From 26186cb0b88ccddddd03251ada9a8a672fcadce1 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 22:06:11 -0700 Subject: [PATCH] refactor(tools): compact kanban/delegation/interrupt/desktop tool modules; drop dead helpers --- tools/debug_helpers.py | 21 +- tools/delegation_live_log.py | 90 +++---- tools/delegation_output_schema.py | 42 +--- tools/desktop_ui.py | 22 +- tools/focus_pane_tool.py | 41 ++-- tools/interpreter_shutdown.py | 47 +--- tools/interrupt.py | 73 ++---- tools/kanban_tools.py | 378 ++++++++++-------------------- 8 files changed, 245 insertions(+), 469 deletions(-) diff --git a/tools/debug_helpers.py b/tools/debug_helpers.py index 16a97dad42..516ad29ed9 100644 --- a/tools/debug_helpers.py +++ b/tools/debug_helpers.py @@ -13,6 +13,10 @@ from hermes_constants import get_hermes_home logger = logging.getLogger(__name__) +def _now() -> str: + return datetime.datetime.now().isoformat() + + class DebugSession: """Per-tool debug session that records tool calls to a JSON log file.""" @@ -22,12 +26,10 @@ class DebugSession: self.session_id = str(uuid.uuid4()) if self.enabled else "" self.log_dir = get_hermes_home() / "logs" self._calls: list[Dict[str, Any]] = [] - self._start_time = datetime.datetime.now().isoformat() if self.enabled else "" - + self._start_time = _now() if self.enabled else "" if self.enabled: self.log_dir.mkdir(parents=True, exist_ok=True) - logger.debug("%s debug mode enabled - Session ID: %s", - tool_name, self.session_id) + logger.debug("%s debug mode enabled - Session ID: %s", tool_name, self.session_id) @property def active(self) -> bool: @@ -35,13 +37,8 @@ class DebugSession: def log_call(self, call_name: str, call_data: Dict[str, Any]) -> None: """Append a tool-call entry to the in-memory log.""" - if not self.enabled: - return - self._calls.append({ - "timestamp": datetime.datetime.now().isoformat(), - "tool_name": call_name, - **call_data, - }) + if self.enabled: + self._calls.append({"timestamp": _now(), "tool_name": call_name, **call_data}) def save(self) -> None: """Flush the in-memory log to a JSON file in the logs directory.""" @@ -52,7 +49,7 @@ class DebugSession: payload = { "session_id": self.session_id, "start_time": self._start_time, - "end_time": datetime.datetime.now().isoformat(), + "end_time": _now(), "debug_enabled": True, "total_calls": len(self._calls), "tool_calls": self._calls, diff --git a/tools/delegation_live_log.py b/tools/delegation_live_log.py index b8ecce33f4..13638025ac 100644 --- a/tools/delegation_live_log.py +++ b/tools/delegation_live_log.py @@ -1,15 +1,11 @@ """Live, tail-able transcripts for delegated subagents. -Each ``delegate_task`` dispatch creates one append-only log per child under -``/cache/delegation/live//task-.log``, pre-created -with a header at dispatch (so ``tail -f`` attaches immediately), streaming one -line per child event; paths are returned from ``delegate_task``. -``cache/delegation`` is mounted read-only into remote terminal backends -(``credential_files._CACHE_DIRS``), so every line written here must be -credential-redacted. Constraints: never raise into the agent loop (first write -failure disables the writer); append mode per write (no handle to lose on a -child crash, close() is the flush); side-channel only (prompt cache unaffected); -no config knobs (7-day retention constant, pruned on each dispatch). +One append-only log per child under ``/cache/delegation/live/ +/task-.log``, pre-created with a header at dispatch (so +``tail -f`` attaches immediately); paths are returned from ``delegate_task``. +``cache/delegation`` is mounted read-only into remote terminal backends, so +every line written here must be credential-redacted. Never raises into the +agent loop; append mode per write (close() is the flush); 7-day retention. """ from __future__ import annotations @@ -27,33 +23,25 @@ logger = logging.getLogger(__name__) LIVE_RETENTION_DAYS = 7 -# Per-line truncation budgets (chars). The .log is a compact operational view; +# Per-line truncation budgets (chars): the .log is a compact operational view; # the child's SessionDB transcript and summary spill files carry full text. _ASSISTANT_MAX = 600 _THINKING_MAX = 300 _ARGS_MAX = 220 _RESULT_MAX = 400 _KICKOFF_MAX = 500 - # Stream deltas are buffered and flushed as one assistant line when another # event type arrives (or on completion); capped so a huge reply can't hold memory. _STREAM_BUFFER_FLUSH_CHARS = 4000 - _TIME_FMT = "%Y-%m-%d %H:%M:%S" def live_transcript_root() -> Path: """Root directory for live transcripts (profile-safe, never ~/.hermes).""" from hermes_constants import get_hermes_dir - return get_hermes_dir("cache/delegation", "delegation_cache") / "live" -def new_live_delegation_id() -> str: - """Same shape as async_delegation's ids so the dir name matches the handle.""" - return f"deleg_{uuid.uuid4().hex[:8]}" - - def _one_line(text: Any, limit: int) -> str: """Collapse to a single line and truncate with an elided-chars note.""" s = " ".join(str(text or "").split()) @@ -63,14 +51,12 @@ def _one_line(text: Any, limit: int) -> str: def _redact(text: str) -> str: - """Mask credentials before anything reaches the sandbox-readable transcript. - ``force=True``: safety boundary, redact even when the global toggle is off; - if the redactor is unavailable, withhold the line rather than leak.""" + """Mask credentials (``force=True``: safety boundary, even when the global + toggle is off); if the redactor is unavailable, withhold rather than leak.""" if not text: return text try: from agent.redact import redact_sensitive_text - return redact_sensitive_text(text, force=True) or "" except Exception: # pragma: no cover - core module; never leak on failure return "[line withheld: redaction unavailable]" @@ -81,9 +67,8 @@ def _dump_json(path: Path, payload: Dict[str, Any]) -> None: class LiveTranscriptWriter: - """Append-only human-readable event log for ONE subagent task. Best-effort: - the first write failure flips ``_ok`` off and later calls become - debug-logged no-ops. Never raises.""" + """Append-only event log for ONE subagent task. Best-effort: the first write + failure flips ``_ok`` off and later calls become debug-logged no-ops.""" def __init__(self, delegation_id: str, task_index: int, goal: str, context: Optional[str] = None, root: Optional[Path] = None): @@ -104,8 +89,7 @@ class LiveTranscriptWriter: f"goal: {_redact(goal_line)}\n" # header bypasses event(), so redact here too f"started: {time.strftime(_TIME_FMT)}\n" "(append-only; streams while the subagent runs — tail -f me)\n" - + "=" * 40 + "\n" - ) + + "=" * 40 + "\n") self.path.write_text(header, encoding="utf-8") self.event("user", "kickoff: " + goal_line + (f" | context: {_one_line(context, _KICKOFF_MAX)}" if context else "")) @@ -115,9 +99,8 @@ class LiveTranscriptWriter: self.path = None def event(self, role: str, text: str) -> None: - """Append one ``HH:MM:SS role | text`` line, flushed per event. Single - choke point: every typed helper funnels through here so one redaction - covers args, results, thinking and streamed text.""" + """Append one ``HH:MM:SS role | text`` line. Single choke point: every typed + helper funnels through here so one redaction covers everything.""" if not self._ok or self.path is None: return line = f"{time.strftime('%H:%M:%S')} {role:<9}| {_redact(text)}\n" @@ -171,10 +154,6 @@ class LiveTranscriptWriter: text, self._stream_buf, self._stream_len = "".join(self._stream_buf), [], 0 self.assistant_text(text) - def _on_tool_completed(self, tool_name, preview, args, kwargs): - self.tool_result(str(tool_name or ""), result=kwargs.get("result"), - duration=kwargs.get("duration"), is_error=bool(kwargs.get("is_error"))) - def _on_complete(self, tool_name, preview, args, kwargs): self.flush_stream() dur = kwargs.get("duration_seconds") @@ -182,27 +161,26 @@ class LiveTranscriptWriter: self.marker(" ".join(filter(None, [ f"status={kwargs.get('status', '?')}", f"duration={dur}s" if dur is not None else "", - f"summary: {_one_line(summary, _RESULT_MAX)}" if summary else "", - ]))) + f"summary: {_one_line(summary, _RESULT_MAX)}" if summary else ""]))) # Event demux (the tool_progress_callback surface): handler(self, tool_name, preview, args, kwargs). _OBSERVERS = { "tool.started": lambda s, n, p, a, kw: s.tool_start(str(n or ""), p if p else a), - "tool.completed": _on_tool_completed, + "tool.completed": lambda s, n, p, a, kw: s.tool_result( + str(n or ""), result=kw.get("result"), duration=kw.get("duration"), + is_error=bool(kw.get("is_error"))), # Fired as cb("_thinking", ) — text rides in the tool_name slot. "_thinking": lambda s, n, p, a, kw: s.thinking(str(n or p or "")), # cb("reasoning.available", "_thinking", , None) "reasoning.available": lambda s, n, p, a, kw: s.thinking(str(p or "")), "subagent.text": lambda s, n, p, a, kw: s.add_stream_delta(str(p or "")), "subagent.start": lambda s, n, p, a, kw: s.event("start", _one_line(p, _KICKOFF_MAX)), - "subagent.complete": _on_complete, - } + "subagent.complete": _on_complete} def observe(self, event_type: Any, tool_name: Any = None, preview: Any = None, args: Any = None, **kwargs: Any) -> None: - """Map a child tool_progress_callback event (shapes from agent/tool_executor.py, - agent/conversation_loop.py, delegate_tool._run_single_child) onto transcript - lines. Unknown events are ignored. Never raises (event() swallows I/O).""" + """Map a child tool_progress_callback event onto transcript lines. + Unknown events are ignored. Never raises (event() swallows I/O).""" handler = self._OBSERVERS.get(str(event_type or "")) if handler is not None: handler(self, tool_name, preview, args, kwargs) @@ -214,14 +192,12 @@ class LiveTranscriptWriter: f"end status={entry.get('status', '?')}", f"exit_reason={exit_reason}" if exit_reason else "", "(iteration budget exhausted)" if exit_reason == "max_iterations" else "", - f"error: {_one_line(entry['error'], _RESULT_MAX)}" if entry.get("error") else "", - ]))) + f"error: {_one_line(entry['error'], _RESULT_MAX)}" if entry.get("error") else ""]))) def wrap_progress_callback(inner_cb, writer: LiveTranscriptWriter): - """Wrap a child's tool_progress_callback so events also land in the log. - ``inner_cb`` may be None; writer failures never propagate and the inner - callback is unchanged. Preserves the ``_flush`` attribute contract.""" + """Wrap a child's tool_progress_callback (may be None) so events also land in + the log; writer failures never propagate. Preserves the ``_flush`` contract.""" def _cb(event_type, tool_name=None, preview=None, args=None, **kwargs): try: @@ -244,22 +220,21 @@ def wrap_progress_callback(inner_cb, writer: LiveTranscriptWriter): def create_live_transcripts( - task_list: List[Dict[str, Any]], - context: Optional[str] = None, - delegation_id: Optional[str] = None, - model: Optional[str] = None, + task_list: List[Dict[str, Any]], context: Optional[str] = None, + delegation_id: Optional[str] = None, model: Optional[str] = None, provider: Optional[str] = None, ) -> tuple[Optional[str], List[Optional[LiveTranscriptWriter]], List[str]]: - """Create one pre-headered writer per task + a manifest.json; also prunes - stale live dirs. Returns ``(delegation_id, writers, paths)``; on any - top-level failure ``(None, [None]*n, [])`` so delegation proceeds untouched.""" + """One pre-headered writer per task + a manifest.json; prunes stale dirs. + Returns ``(delegation_id, writers, paths)``; on any top-level failure + ``(None, [None]*n, [])`` so delegation proceeds untouched.""" n = len(task_list) try: prune_stale_live_dirs() except Exception: pass try: - deleg_id = delegation_id or new_live_delegation_id() + # Same id shape as async_delegation's so the dir name matches the handle. + deleg_id = delegation_id or f"deleg_{uuid.uuid4().hex[:8]}" made = [LiveTranscriptWriter(deleg_id, i, str(t.get("goal", "")), context=t.get("context") or context) for i, t in enumerate(task_list)] writers: List[Optional[LiveTranscriptWriter]] = [w if w.path is not None else None for w in made] @@ -293,8 +268,7 @@ def _write_manifest(delegation_id: str, task_list: List[Dict[str, Any]], "goal": _redact(str(t.get("goal", ""))[:500]), "log": paths[i] if i < len(paths) else None, "status": "running", - } for i, t in enumerate(task_list)], - }) + } for i, t in enumerate(task_list)]}) except Exception as exc: logger.debug("Live transcript manifest write failed: %s", exc) diff --git a/tools/delegation_output_schema.py b/tools/delegation_output_schema.py index 1f8d0abcda..eae47ce8a9 100644 --- a/tools/delegation_output_schema.py +++ b/tools/delegation_output_schema.py @@ -3,8 +3,8 @@ Optional per-task ``output_schema`` (a JSON Schema object): the child gets an OUTPUT CONTRACT block appended to its context, the parent validates the final answer with jsonschema, and on failure sends exactly ONE bounded retry turn -carrying the validation errors verbatim (max 1 retry — more retries make -frontier models drop fields that were right the first time; no schema re-paste). +carrying the validation errors verbatim (more retries make frontier models +drop fields that were right the first time; the schema is never re-pasted). """ from __future__ import annotations @@ -15,15 +15,10 @@ from typing import Any, Dict, List, Optional, Tuple logger = logging.getLogger(__name__) -_CONTRACT_HEADER = "OUTPUT CONTRACT (machine-validated)" - def coerce_output_schema(raw: Any) -> Tuple[Optional[Dict[str, Any]], Optional[str]]: - """Validate a model/caller-supplied output_schema value. - - Returns ``(schema, None)`` when usable, ``(None, error)`` when not; - ``None`` input passes through as ``(None, None)`` (no schema requested). - """ + """``(schema, None)`` when usable, ``(None, error)`` when not; ``None`` input + passes through as ``(None, None)`` (no schema requested).""" if raw is None: return None, None if isinstance(raw, str): @@ -39,7 +34,6 @@ def coerce_output_schema(raw: Any) -> Tuple[Optional[Dict[str, Any]], Optional[s return None, f"output_schema must be a JSON Schema object, got {type(raw).__name__}." try: from jsonschema.validators import validator_for # type: ignore[import-untyped] - validator_for(raw).check_schema(raw) except ImportError: # Degrade to accepting the dict as-is so delegation still works without jsonschema. @@ -56,20 +50,17 @@ def append_output_contract(context: Optional[str], schema: Dict[str, Any]) -> st except (TypeError, ValueError): schema_text = str(schema) block = ( - f"{_CONTRACT_HEADER}:\n" + "OUTPUT CONTRACT (machine-validated):\n" "Your FINAL response must be a single JSON object that validates " "against this JSON Schema. No prose before or after the JSON; a " "```json code fence is acceptable but not required.\n" - f"{schema_text}" - ) + f"{schema_text}") base = (context or "").rstrip() return f"{base}\n\n{block}" if base else block -def extract_json_candidate(text: str) -> str: - """Strip markdown fences and prose around the outermost ``{...}``/``[...]`` - span. Returns the (possibly unchanged) candidate; parse errors are - reported by validate_output.""" +def _extract_json_candidate(text: str) -> str: + """Strip markdown fences and prose around the outermost ``{...}``/``[...]``.""" raw = (text or "").strip() if raw.startswith("```"): raw = raw.split("\n", 1)[-1] @@ -89,12 +80,8 @@ def extract_json_candidate(text: str) -> str: def validate_output(text: str, schema: Dict[str, Any]) -> Tuple[bool, List[str]]: - """Validate a child's final answer against ``schema``. - - Returns ``(True, [])`` or ``(False, errors)`` with human-readable strings - suitable for the retry turn. - """ - candidate = extract_json_candidate(text or "") + """``(True, [])`` or ``(False, errors)`` with strings suitable for the retry turn.""" + candidate = _extract_json_candidate(text or "") if not candidate.strip(): return False, ["Response was empty — expected a JSON object matching the schema."] try: @@ -110,19 +97,16 @@ def validate_output(text: str, schema: Dict[str, Any]) -> Tuple[bool, List[str]] errors = sorted(validator.iter_errors(parsed), key=lambda e: list(e.absolute_path)) rendered = [ # bound error volume for the retry prompt "$" + "".join(f"[{p}]" if isinstance(p, int) else f".{p}" for p in err.absolute_path) + f": {err.message}" - for err in errors[:10] - ] + for err in errors[:10]] return not rendered, rendered def build_retry_message(errors: List[str]) -> str: - """Single bounded retry turn: errors verbatim, deliberately NOT re-pasting - the schema (the child already has it in its context).""" + """Single bounded retry turn: errors verbatim, schema deliberately NOT re-pasted.""" error_block = "\n".join(f"- {e}" for e in errors) return ( "Your previous final response was rejected by the output contract " "validator. Validation errors:\n" f"{error_block}\n\n" "Reply with ONLY the corrected JSON object matching the OUTPUT " - "CONTRACT schema from your task context. No prose, no explanations." - ) + "CONTRACT schema from your task context. No prose, no explanations.") diff --git a/tools/desktop_ui.py b/tools/desktop_ui.py index 46f5b2db71..0c7c23db92 100644 --- a/tools/desktop_ui.py +++ b/tools/desktop_ui.py @@ -31,15 +31,13 @@ def user_enabled(setting: str, default: bool) -> bool: """Read one of the desktop's Appearance switches from ``display.``. The renderer mirrors these toggles onto the CONNECTED gateway's config, so this - reads the user's real answer for local/SSH/URL/cloud gateways alike (an env var - would only describe the process). ``check_fn``s use it to withdraw a tool from the - schema when the user switched the feature off — Hermes should not be told about - a surface it may not use. Unreadable config -> ``default`` so a shipped-on - feature does not vanish on a transient read error. + is the user's real answer for local/SSH/URL/cloud gateways alike. ``check_fn``s + use it to withdraw a tool from the schema when the feature is switched off. + Unreadable config -> ``default`` so a shipped-on feature does not vanish on a + transient read error. """ try: from hermes_cli.config import load_config_readonly - display = load_config_readonly().get("display") except Exception: return default @@ -49,9 +47,7 @@ def user_enabled(setting: str, default: bool) -> bool: def emit(event: str, payload: dict) -> bool: - """Route ``event`` to the window that owns the current turn. - - Returns ``False`` when no emitter is wired (i.e. not the desktop app).""" + """Route ``event`` to the window owning the current turn; False when no emitter.""" fn = _emit if fn is None: return False @@ -60,11 +56,9 @@ def emit(event: str, payload: dict) -> bool: def emit_or_error(event: str, payload: dict, fail_prefix: str, desktop_only: str, result: dict) -> str: - """Emit ``event``; return ``tool_error`` text on failure, else ``result`` as JSON. - - ``fail_prefix`` is prepended to the exception text; ``desktop_only`` is the error - when no emitter is wired. Looked up as ``desktop_ui.emit`` so tests can patch it. - """ + """Emit ``event``; ``tool_error`` text on failure (``fail_prefix`` + exception, or + ``desktop_only`` when no emitter), else ``result`` as JSON. Calls ``emit`` via the + module attribute so tests can patch it.""" try: ok = emit(event, payload) except Exception as exc: diff --git a/tools/focus_pane_tool.py b/tools/focus_pane_tool.py index 5e18decc62..5226879873 100644 --- a/tools/focus_pane_tool.py +++ b/tools/focus_pane_tool.py @@ -25,31 +25,28 @@ def focus_pane_tool(pane: str) -> str: ) -FOCUS_PANE_SCHEMA = { - "name": "focus_pane", - "description": ( - "Reveal and focus a Hermes desktop pane when the user asks to see it: " - "chat, files, terminal, review (git diff), or sessions. For URLs/" - "files use the desktop_preview tool instead." - ), - "parameters": { - "type": "object", - "properties": { - "pane": { - "type": "string", - "enum": list(PANES), - "description": "Which pane to reveal.", - }, - }, - "required": ["pane"], - }, -} - - registry.register( name="focus_pane", toolset="desktop_ui", - schema=FOCUS_PANE_SCHEMA, + schema={ + "name": "focus_pane", + "description": ( + "Reveal and focus a Hermes desktop pane when the user asks to see it: " + "chat, files, terminal, review (git diff), or sessions. For URLs/" + "files use the desktop_preview tool instead." + ), + "parameters": { + "type": "object", + "properties": { + "pane": { + "type": "string", + "enum": list(PANES), + "description": "Which pane to reveal.", + }, + }, + "required": ["pane"], + }, + }, handler=lambda args, **kw: focus_pane_tool(pane=args.get("pane", "")), emoji="🪟", ) diff --git a/tools/interpreter_shutdown.py b/tools/interpreter_shutdown.py index df00316e59..770010fad1 100644 --- a/tools/interpreter_shutdown.py +++ b/tools/interpreter_shutdown.py @@ -1,26 +1,13 @@ -"""Shared interpreter-shutdown detection. +"""Shared interpreter-shutdown predicate for background threads that can outlive +process teardown (cron delivery, concurrent tool submission, retry paths). -Single home for the "is the Python interpreter finalizing?" predicate used by -every subsystem whose background threads can outlive process teardown (cron -delivery, concurrent tool submission, the conversation loop's retry path, -background review forks). Once finalization starts, ``concurrent.futures`` -refuses new work and asyncio's default executor is gone — any further attempt -to schedule work (an API retry, a pool submit, ``asyncio.run``) is doomed and -only produces noise: stray ``❌`` prints after the TUI exited, tracebacks in -``errors.log``, and futile retry loops against a dying process. - -CPython emits two message variants depending on the failing site: - -- ``cannot schedule new futures after interpreter shutdown`` — the - module-global finalization flag (asyncio.run_coroutine_threadsafe, a - torn-down default executor, ThreadPoolExecutor.submit during teardown). -- ``cannot schedule new futures after shutdown`` — a plain - ``ThreadPoolExecutor`` whose ``shutdown()`` ran. - -The common short prefix catches both. Matching the second variant is safe -for shutdown detection at every current call site: the pools involved are -either module-global daemons or ``with``-scoped locals that cannot be shut -down mid-use by anything except interpreter finalization. +Once finalization starts, ``concurrent.futures`` refuses new work and asyncio's +default executor is gone, so any further scheduling only produces noise (stray +prints after the TUI exited, tracebacks in ``errors.log``, futile retries). +CPython raises ``cannot schedule new futures after interpreter shutdown`` +(module-global flag) or ``... after shutdown`` (a pool whose ``shutdown()`` ran); +the short prefix catches both — safe here because every pool involved is a +module-global daemon or a ``with``-scoped local only finalization can shut down. """ from __future__ import annotations @@ -28,19 +15,11 @@ from __future__ import annotations import sys from typing import Optional -_SHUTDOWN_SUBMIT_ERROR_PREFIX = "cannot schedule new futures" - def interpreter_shutting_down(exc: Optional[BaseException] = None) -> bool: - """Return True when the Python interpreter is finalizing. - - ``exc`` lets a caller also treat an already-raised scheduling error as a - shutdown signal: the ``concurrent.futures`` module-global flag can be set - a hair before ``sys.is_finalizing()`` flips, so matching the error text - is a safe fallback for that race. - """ + """True when the interpreter is finalizing. ``exc`` lets a caller treat an + already-raised scheduling error as a shutdown signal: the ``concurrent.futures`` + flag can be set a hair before ``sys.is_finalizing()`` flips.""" if sys.is_finalizing(): return True - if exc is not None: - return _SHUTDOWN_SUBMIT_ERROR_PREFIX in str(exc).lower() - return False + return exc is not None and "cannot schedule new futures" in str(exc).lower() diff --git a/tools/interrupt.py b/tools/interrupt.py index f0d012643d..8a15b4946f 100644 --- a/tools/interrupt.py +++ b/tools/interrupt.py @@ -1,15 +1,9 @@ """Per-thread interrupt signaling for all tools. -Thread-scoped interrupt tracking so that interrupting one agent session does -not kill tools running in other sessions — critical in the gateway where -multiple agents run concurrently in one process. The agent stores its -execution thread ID at the start of run_conversation() and passes it to +Thread-scoped so interrupting one agent session does not kill tools running in +other sessions (the gateway runs many agents in one process). The agent stores +its execution thread id at the start of run_conversation() and passes it to set_interrupt(); tools call is_interrupted(), which checks the CURRENT thread. - -Usage in tools: - from tools.interrupt import is_interrupted - if is_interrupted(): - return {"output": "[interrupted]", "returncode": 130} """ import logging @@ -19,13 +13,12 @@ import threading logger = logging.getLogger(__name__) # Opt-in debug tracing — pairs with HERMES_DEBUG_INTERRUPT in -# tools/environments/base.py. Logs caller thread, target thread, and current -# state per set/check for "interrupt signaled but tool never saw it" reports. +# tools/environments/base.py; logs caller/target thread and state per set/check. _DEBUG_INTERRUPT = bool(os.getenv("HERMES_DEBUG_INTERRUPT")) if _DEBUG_INTERRUPT: - # AIAgent's quiet_mode path forces the `tools` logger to ERROR on CLI - # startup; force ours back to INFO so the trace is visible in agent.log. + # AIAgent's quiet_mode forces the `tools` logger to ERROR on CLI startup; + # force ours back to INFO so the trace is visible in agent.log. logger.setLevel(logging.INFO) # Interrupted thread idents, plus an optional user-safe cause per signal. The @@ -35,15 +28,9 @@ _interrupt_reasons: dict[int, str] = {} _lock = threading.Lock() -def set_interrupt( - active: bool, - thread_id: int | None = None, - *, - reason: str | None = None, -) -> None: - """Set (``active=True``) or clear the interrupt for *thread_id*, defaulting - to the current thread (backward compat for CLI/tests). ``reason`` is an - optional user-safe cause.""" +def set_interrupt(active: bool, thread_id: int | None = None, *, reason: str | None = None) -> None: + """Set or clear the interrupt for *thread_id* (default: current thread, for + CLI/tests). ``reason`` is an optional user-safe cause.""" tid = thread_id if thread_id is not None else threading.current_thread().ident with _lock: if active: @@ -60,8 +47,7 @@ def set_interrupt( logger.info( "[interrupt-debug] set_interrupt(active=%s, target_tid=%s) " "called_from_tid=%s current_set=%s", - active, tid, threading.current_thread().ident, _snapshot, - ) + active, tid, threading.current_thread().ident, _snapshot) def is_interrupted() -> bool: @@ -70,12 +56,9 @@ def is_interrupted() -> bool: def is_thread_interrupted(thread_id: int | None) -> bool: - """Check whether *thread_id* has an interrupt bit set (``None`` never is). - - Used when a wait is moved onto a deadline worker (``run_bounded_sync``) - so ``/stop`` targeting the original tool-worker tid still kills the - subprocess. - """ + """Whether *thread_id* has an interrupt bit set (``None`` never is). Used when + a wait moves onto a deadline worker (``run_bounded_sync``) so ``/stop`` + targeting the original tool-worker tid still kills the subprocess.""" if thread_id is None: return False with _lock: @@ -83,35 +66,27 @@ def is_thread_interrupted(thread_id: int | None) -> bool: def get_interrupt_reason() -> str | None: - """Return the user-safe interrupt cause for the current thread, if known.""" - tid = threading.current_thread().ident + """User-safe interrupt cause for the current thread, if known.""" with _lock: - return _interrupt_reasons.get(tid) + return _interrupt_reasons.get(threading.current_thread().ident) def clear_current_thread_interrupt() -> None: """Clear any interrupt bit on the CURRENT thread. - Gives a user-approved command a clean interrupt slate immediately before - it spawns its child process, so a stale bit that landed on this thread - during the blocking approval-wait cannot SIGINT the just-approved run - (exit 130 + "[Command interrupted]"). Single-thread ordering on this tid - keeps the DO-NOT-BREAK invariant intact: a *genuine* interrupt arriving - after this call re-sets the bit on the same thread and is still observed by - the executor's poll loop. Call this directly, never via the - _interrupt_event proxy (its .clear() binds to whatever thread runs it). + Gives a user-approved command a clean slate right before it spawns its child, + so a stale bit that landed during the blocking approval-wait cannot SIGINT the + just-approved run. Single-thread ordering keeps the invariant: a *genuine* + interrupt arriving after this call re-sets the bit and is still observed by the + executor's poll loop. Call directly, never via the _interrupt_event proxy (its + .clear() binds to whatever thread runs it). """ - set_interrupt(False) # thread_id=None -> current thread + set_interrupt(False) -# --------------------------------------------------------------------------- -# Backward-compatible _interrupt_event proxy: legacy call sites -# (code_execution_tool, process_registry, tests) import it and call -# .is_set() / .set() / .clear(); the shim maps those to the per-thread API. -# --------------------------------------------------------------------------- - class _ThreadAwareEventProxy: - """Drop-in proxy that maps threading.Event methods to per-thread state.""" + """Backward-compatible ``_interrupt_event``: legacy call sites call + .is_set()/.set()/.clear(); the shim maps those to the per-thread API.""" def is_set(self) -> bool: return is_interrupted() diff --git a/tools/kanban_tools.py b/tools/kanban_tools.py index 9d2c0cd49e..e8cdd906fa 100644 --- a/tools/kanban_tools.py +++ b/tools/kanban_tools.py @@ -22,24 +22,11 @@ from hermes_cli.goals import judge_goal from tools.registry import registry, tool_error from hermes_cli.config import cfg_get, load_config from tools.kanban_tools_schemas import ( # noqa: F401 - re-exported for callers/tests - _DESC_BOARD, - _DESC_TASK_ID_DEFAULT, - _board_schema_prop, - KANBAN_ATTACH_SCHEMA, - KANBAN_ATTACH_URL_SCHEMA, - KANBAN_ATTACHMENTS_SCHEMA, - KANBAN_BLOCK_SCHEMA, - KANBAN_COMMENT_SCHEMA, - KANBAN_COMPLETE_SCHEMA, - KANBAN_CREATE_SCHEMA, - KANBAN_HEARTBEAT_SCHEMA, - KANBAN_LINK_SCHEMA, - KANBAN_LIST_SCHEMA, - KANBAN_REQUEST_CHANGES_SCHEMA, - KANBAN_REQUEST_REVIEW_SCHEMA, - KANBAN_SHOW_SCHEMA, - KANBAN_UNBLOCK_SCHEMA, -) + _DESC_BOARD, _DESC_TASK_ID_DEFAULT, _board_schema_prop, KANBAN_ATTACH_SCHEMA, + KANBAN_ATTACH_URL_SCHEMA, KANBAN_ATTACHMENTS_SCHEMA, KANBAN_BLOCK_SCHEMA, KANBAN_COMMENT_SCHEMA, + KANBAN_COMPLETE_SCHEMA, KANBAN_CREATE_SCHEMA, KANBAN_HEARTBEAT_SCHEMA, KANBAN_LINK_SCHEMA, + KANBAN_LIST_SCHEMA, KANBAN_REQUEST_CHANGES_SCHEMA, KANBAN_REQUEST_REVIEW_SCHEMA, + KANBAN_SHOW_SCHEMA, KANBAN_UNBLOCK_SCHEMA) logger = logging.getLogger(__name__) @@ -47,9 +34,7 @@ KANBAN_LIST_DEFAULT_LIMIT = 50 KANBAN_LIST_MAX_LIMIT = 200 -# --------------------------------------------------------------------------- -# Gating -# --------------------------------------------------------------------------- +# --- Gating --- def _profile_has_kanban_toolset() -> bool: # load_config() is mtime-cached and check_fn results are TTL-cached (~30s). @@ -102,9 +87,7 @@ def _check_kanban_orchestrator_mode() -> bool: return _visible(to_env_worker=False) -# --------------------------------------------------------------------------- -# Shared helpers — validation failures raise _Reject; _kanban_handler renders it -# --------------------------------------------------------------------------- +# --- Shared helpers: validation failures raise _Reject; _kanban_handler renders it --- _TASK_ID_REQUIRED = "task_id is required (or set HERMES_KANBAN_TASK in the env)" @@ -145,11 +128,9 @@ def _reject_delegated_child_mutation(tool_name: str) -> None: must not mutate board state.""" if _is_delegated_child_context(): raise _Reject( - f"{tool_name} refused: delegate_task child agents are not Kanban " - "run owners. Return findings to the parent agent; the dispatcher " - "worker or an explicitly configured Kanban orchestrator must perform " - "board mutations." - ) + f"{tool_name} refused: delegate_task child agents are not Kanban run owners. " + "Return findings to the parent agent; the dispatcher worker or an explicitly " + "configured Kanban orchestrator must perform board mutations.") def _default_task_id(arg: Optional[str]) -> Optional[str]: @@ -195,10 +176,8 @@ def _enforce_worker_task_ownership(tid: str) -> None: env_tid = os.environ.get("HERMES_KANBAN_TASK") if env_tid and tid != env_tid: raise _Reject( - f"worker is scoped to task {env_tid}; refusing to mutate " - f"{tid}. Use kanban_comment to hand off information to other " - f"tasks, or kanban_create to spawn follow-up work." - ) + f"worker is scoped to task {env_tid}; refusing to mutate {tid}. Use kanban_comment " + f"to hand off information to other tasks, or kanban_create to spawn follow-up work.") def _worker_guard(tool_name: str, args: dict) -> str: @@ -215,10 +194,9 @@ def _require_orchestrator_tool(tool_name: str) -> None: a stale registration or test harness routing a worker here anyway.""" if os.environ.get("HERMES_KANBAN_TASK"): raise _Reject( - f"{tool_name} is orchestrator-only; dispatcher-spawned workers " - "must use kanban_complete, kanban_block, kanban_heartbeat, or " - "kanban_comment for their assigned task." - ) + f"{tool_name} is orchestrator-only; dispatcher-spawned workers must use " + "kanban_complete, kanban_block, kanban_heartbeat, or kanban_comment for their " + "assigned task.") @contextmanager @@ -317,27 +295,19 @@ def _opt_int(value: Any, default: Optional[int] = None) -> Optional[int]: _TASK_FIELDS = ( - "id", "title", "body", "assignee", "status", "tenant", "priority", - "workspace_kind", "workspace_path", "created_by", "created_at", - "started_at", "completed_at", "result", "current_run_id", - "model_override", "provider_override", -) + "id", "title", "body", "assignee", "status", "tenant", "priority", "workspace_kind", + "workspace_path", "created_by", "created_at", "started_at", "completed_at", "result", + "current_run_id", "model_override", "provider_override") _TASK_SUMMARY_FIELDS = ( - "id", "title", "assignee", "status", "priority", "tenant", - "workspace_kind", "workspace_path", "project_id", "created_by", - "created_at", "started_at", "completed_at", "current_run_id", - "model_override", "provider_override", -) + "id", "title", "assignee", "status", "priority", "tenant", "workspace_kind", "workspace_path", + "project_id", "created_by", "created_at", "started_at", "completed_at", "current_run_id", + "model_override", "provider_override") _RUN_FIELDS = ( - "id", "profile", "status", "outcome", "summary", "error", "metadata", - "started_at", "ended_at", -) + "id", "profile", "status", "outcome", "summary", "error", "metadata", "started_at", "ended_at") _COMMENT_FIELDS = ("author", "body", "created_at") _EVENT_FIELDS = ("kind", "payload", "created_at", "run_id") _ATTACHMENT_FIELDS = ( - "id", "filename", "content_type", "size", "uploaded_by", "stored_path", - "created_at", -) + "id", "filename", "content_type", "size", "uploaded_by", "stored_path", "created_at") def _fields(obj: Any, names: tuple[str, ...]) -> dict[str, Any]: @@ -348,17 +318,11 @@ def _task_summary_dict(kb, conn, task) -> dict[str, Any]: parents = kb.parent_ids(conn, task.id) children = kb.child_ids(conn, task.id) return { - **_fields(task, _TASK_SUMMARY_FIELDS), - "parents": parents, - "children": children, - "parent_count": len(parents), - "child_count": len(children), - } + **_fields(task, _TASK_SUMMARY_FIELDS), "parents": parents, "children": children, + "parent_count": len(parents), "child_count": len(children)} -# --------------------------------------------------------------------------- -# Goal-mode judge gate -# --------------------------------------------------------------------------- +# --- Goal-mode judge gate --- _GOAL_MODE_BLOCK_ALLOWED_KINDS = frozenset({"dependency", "needs_input"}) @@ -379,33 +343,20 @@ def _goal_judge_available() -> bool: _GOAL_GATE_MESSAGES = { "kanban_complete": { "blocked": ( - "Goal completion rejected: judge ruled the goal " - "unachievable — {reason}. The task will NOT complete " - "silently. Either re-scope the task with kanban_edit, " - "or record the block with kanban_block and hand the " - "decision to a human / reviewer." - ), + "Goal completion rejected: judge ruled the goal unachievable — {reason}. The task " + "will NOT complete silently. Either re-scope the task with kanban_edit, or record " + "the block with kanban_block and hand the decision to a human / reviewer."), "continue": ( - "Goal completion rejected by judge: {reason}. " - "To proceed, either: (1) provide explicit acceptance " - "evidence in your summary matching the task's criteria, " - "or (2) create continuation tasks with parents=[{tid}] " - "and keep this task alive." - ), - }, + "Goal completion rejected by judge: {reason}. To proceed, either: (1) provide " + "explicit acceptance evidence in your summary matching the task's criteria, or (2) " + "create continuation tasks with parents=[{tid}] and keep this task alive.")}, "kanban_request_review": { "blocked": ( - "Goal review handoff rejected: judge ruled the goal " - "unachievable — {reason}. Record the block with " - "kanban_block instead of requesting review." - ), + "Goal review handoff rejected: judge ruled the goal unachievable — {reason}. " + "Record the block with kanban_block instead of requesting review."), "continue": ( - "Goal review handoff rejected by judge: {reason}. " - "Provide acceptance evidence matching the card before " - "requesting review." - ), - }, -} + "Goal review handoff rejected by judge: {reason}. Provide acceptance evidence " + "matching the card before requesting review.")}} def _goal_gate(tool_name: str, task, tid: str, evidence: str) -> None: @@ -417,13 +368,10 @@ def _goal_gate(tool_name: str, task, tid: str, evidence: str) -> None: return try: verdict, reason, _, _, _ = judge_goal( - goal=f"{task.title}\n\n{task.body or ''}".strip(), - last_response=evidence.strip(), - ) + goal=f"{task.title}\n\n{task.body or ''}".strip(), last_response=evidence.strip()) except Exception as judge_exc: logger.warning( - "goal judge check failed, allowing lifecycle handoff: %s", judge_exc, exc_info=True, - ) + "goal judge check failed, allowing lifecycle handoff: %s", judge_exc, exc_info=True) return if verdict == "done": return @@ -431,9 +379,7 @@ def _goal_gate(tool_name: str, task, tid: str, evidence: str) -> None: raise _Reject(_GOAL_GATE_MESSAGES[tool_name][key].format(reason=reason, tid=tid)) -# --------------------------------------------------------------------------- -# Runtime-activity → board bridges (auto-heartbeat, live comment injection) -# --------------------------------------------------------------------------- +# --- Runtime-activity → board bridges (auto-heartbeat, live comment injection) --- # The dispatcher watchdog reads ``tasks.last_heartbeat_at``, not the agent's # in-process activity timestamp, so normal work is mirrored onto the board here; # the explicit ``kanban_heartbeat`` tool stays for notes / pre-extending a claim. @@ -530,8 +476,7 @@ def inject_new_comments_from_env(agent: Any) -> bool: + ("s" if len(fresh) > 1 else "") + " on your kanban task from the operator (delivered mid-run). " + "Take it into account for the work you're doing right now:\n" - + "\n".join(lines) - ) + + "\n".join(lines)) try: return bool(agent.steer(note)) except Exception: @@ -539,9 +484,7 @@ def inject_new_comments_from_env(agent: Any) -> bool: return False -# --------------------------------------------------------------------------- -# Handlers -# --------------------------------------------------------------------------- +# --- Handlers --- @_kanban_handler("kanban_show") def _handle_show(args: dict, **kw) -> str: @@ -560,8 +503,7 @@ def _handle_show(args: dict, **kw) -> str: "events": [_fields(e, _EVENT_FIELDS) for e in kb.list_events(conn, tid)[-50:]], "runs": [_fields(r, _RUN_FIELDS) for r in kb.list_runs(conn, tid)], # Same string build_worker_context hands the dispatcher at spawn time. - "worker_context": kb.build_worker_context(conn, tid), - }) + "worker_context": kb.build_worker_context(conn, tid)}) @_kanban_handler("kanban_list") @@ -586,26 +528,16 @@ def _handle_list(args: dict, **kw) -> str: promoted = kb.recompute_ready(conn) # One extra row lets the output report truncation without dumping the board. rows = kb.list_tasks( - conn, - assignee=args.get("assignee"), - status=args.get("status"), - tenant=args.get("tenant"), - include_archived=include_archived, - limit=limit + 1, - ) + conn, assignee=args.get("assignee"), status=args.get("status"), + tenant=args.get("tenant"), include_archived=include_archived, limit=limit + 1) truncated = len(rows) > limit tasks = rows[:limit] return json.dumps({ "tasks": [_task_summary_dict(kb, conn, t) for t in tasks], - "count": len(tasks), - "limit": limit, - "truncated": truncated, - "next_limit": ( - min(limit * 2, KANBAN_LIST_MAX_LIMIT) - if truncated and limit < KANBAN_LIST_MAX_LIMIT else None - ), - "promoted": promoted, - }) + "count": len(tasks), "limit": limit, "truncated": truncated, + "next_limit": (min(limit * 2, KANBAN_LIST_MAX_LIMIT) + if truncated and limit < KANBAN_LIST_MAX_LIMIT else None), + "promoted": promoted}) @_kanban_handler("kanban_complete") @@ -624,7 +556,8 @@ def _handle_complete(args: dict, **kw) -> str: redacted = _redact_metadata(metadata) if redacted is not None: metadata = redacted - created_cards = _coerce_str_list(args.get("created_cards"), "created_cards", "task ids", strip=True) + created_cards = _coerce_str_list( + args.get("created_cards"), "created_cards", "task ids", strip=True) artifacts = _coerce_str_list(args.get("artifacts"), "artifacts", "file paths", strip=True) if artifacts: metadata = _merge_artifacts(metadata, artifacts) @@ -637,31 +570,24 @@ def _handle_complete(args: dict, **kw) -> str: _goal_gate("kanban_complete", task, tid, (summary or result or "").strip()) try: ok = kb.complete_task( - conn, tid, - result=result, summary=summary, metadata=metadata, - created_cards=created_cards, - expected_run_id=_worker_run_id(tid), - ) + conn, tid, result=result, summary=summary, metadata=metadata, + created_cards=created_cards, expected_run_id=_worker_run_id(tid)) except kb.ArtifactPreservationError as artifact_err: return tool_error( - f"kanban_complete could not preserve the declared artifacts: " - f"{artifact_err}. Your task is still in-flight and its " - f"scratch workspace was kept. Fix the artifact path or " - f"storage error, then retry kanban_complete with the same handoff." - ) + f"kanban_complete could not preserve the declared artifacts: {artifact_err}. " + f"Your task is still in-flight and its scratch workspace was kept. Fix the " + f"artifact path or storage error, then retry kanban_complete with the same " + f"handoff.") except kb.HallucinatedCardsError as hall_err: # The gate runs before the write txn, so the task was NOT mutated; # say so explicitly or the model treats the error as terminal and # blocks/crashes instead of retrying. Audit event already landed. return tool_error( - f"kanban_complete blocked: the following created_cards " - f"do not exist or were not created by this worker: " - f"{', '.join(hall_err.phantom)}. " - f"Your task is still in-flight (no state change). " - f"Retry kanban_complete with the same summary/metadata " - f"and either drop these ids from created_cards, or pass " - f"created_cards=[] to skip the card-claim check entirely." - ) + f"kanban_complete blocked: the following created_cards do not exist or were not " + f"created by this worker: {', '.join(hall_err.phantom)}. Your task is still " + f"in-flight (no state change). Retry kanban_complete with the same " + f"summary/metadata and either drop these ids from created_cards, or pass " + f"created_cards=[] to skip the card-claim check entirely.") if not ok: return tool_error(f"could not complete {tid} (unknown id or already terminal)") run = kb.latest_run(conn, tid) @@ -672,7 +598,8 @@ def _handle_complete(args: dict, **kw) -> str: def _handle_block(args: dict, **kw) -> str: """Transition the task to blocked with a reason a human will read.""" tid = _worker_guard("kanban_block", args) - reason = _redact(_require_text(args, "reason", "reason is required — explain what input you need")) + reason = _redact( + _require_text(args, "reason", "reason is required — explain what input you need")) kind = args.get("kind") with _board(args.get("board")) as (kb, conn): if kind is not None and kind not in kb.VALID_BLOCK_KINDS: @@ -684,23 +611,17 @@ def _handle_block(args: dict, **kw) -> str: if task and task.goal_mode and kind not in _GOAL_MODE_BLOCK_ALLOWED_KINDS: return tool_error( f"goal_mode tasks can only block with kind in " - f"{sorted(_GOAL_MODE_BLOCK_ALLOWED_KINDS)} (got {kind!r}). " - f"If the task is actually finished or cannot proceed for " - f"another reason, call kanban_complete instead — the " - f"completion judge will evaluate it." - ) + f"{sorted(_GOAL_MODE_BLOCK_ALLOWED_KINDS)} (got {kind!r}). If the task is actually " + f"finished or cannot proceed for another reason, call kanban_complete instead — " + f"the completion judge will evaluate it.") ok = kb.block_task(conn, tid, reason=reason, kind=kind, expected_run_id=_worker_run_id(tid)) if not ok: return tool_error(f"could not block {tid} (unknown id or not in running/ready)") run = kb.latest_run(conn, tid) # Report where the task actually landed; routing may not leave it in 'blocked'. landed = kb.get_task(conn, tid) - return _ok( - task_id=tid, - run_id=run.id if run else None, - status=landed.status if landed else "blocked", - block_kind=kind, - ) + return _ok(task_id=tid, run_id=run.id if run else None, + status=landed.status if landed else "blocked", block_kind=kind) @_kanban_handler("kanban_request_review") @@ -708,10 +629,8 @@ def _handle_request_review(args: dict, **kw) -> str: """Move implementation into the first-class review phase.""" tid = _worker_guard("kanban_request_review", args) summary = _redact(_require_text( - args, "summary", - "summary is required — describe what was implemented and how it " - "was verified so the reviewer has context", - )) + args, "summary", "summary is required — describe what was implemented and how it " + "was verified so the reviewer has context")) metadata = args.get("metadata") _require_dict_metadata(metadata) if metadata is not None: @@ -726,44 +645,33 @@ def _handle_request_review(args: dict, **kw) -> str: with _board(args.get("board")) as (kb, conn): _goal_gate("kanban_request_review", kb.get_task(conn, tid), tid, summary) ok, fail_reason = kb.request_review( - conn, tid, - summary=summary, - metadata=metadata, - reviewer=reviewer, - expected_run_id=_worker_run_id(tid), - with_reason=True, - ) + conn, tid, summary=summary, metadata=metadata, reviewer=reviewer, + expected_run_id=_worker_run_id(tid), with_reason=True) if not ok: detail = fail_reason or "unknown id or not in running/ready" return tool_error(f"could not request review for {tid}: {detail}") run = kb.latest_run(conn, tid) landed = kb.get_task(conn, tid) - return _ok( - task_id=tid, - run_id=run.id if run else None, - status=landed.status if landed else "review", - ) + return _ok(task_id=tid, run_id=run.id if run else None, + status=landed.status if landed else "review") @_kanban_handler("kanban_request_changes") def _handle_request_changes(args: dict, **kw) -> str: """Return a reviewer-owned running task to its implementer.""" tid = _worker_guard("kanban_request_changes", args) - reason = _redact(_require_text(args, "reason", "reason is required — describe the changes needed")) + reason = _redact( + _require_text(args, "reason", "reason is required — describe the changes needed")) with _board(args.get("board")) as (kb, conn): - ok, detail = kb.request_changes(conn, tid, reason=reason, expected_run_id=_worker_run_id(tid)) + ok, detail = kb.request_changes( + conn, tid, reason=reason, expected_run_id=_worker_run_id(tid)) if not ok: return tool_error( - f"could not request changes for {tid}: {detail or 'invalid review state'}" - ) + f"could not request changes for {tid}: {detail or 'invalid review state'}") landed = kb.get_task(conn, tid) run = kb.latest_run(conn, tid) - return _ok( - task_id=tid, - run_id=run.id if run else None, - status=landed.status if landed else "ready", - implementer=detail, - ) + return _ok(task_id=tid, run_id=run.id if run else None, + status=landed.status if landed else "ready", implementer=detail) @_kanban_handler("kanban_heartbeat") @@ -776,7 +684,8 @@ def _handle_heartbeat(args: dict, **kw) -> str: # The dispatcher pins HERMES_KANBAN_CLAIM_LOCK at spawn; the default # claimer covers locally-driven workers that bypassed the dispatcher. kb.heartbeat_claim(conn, tid, claimer=os.environ.get("HERMES_KANBAN_CLAIM_LOCK")) - ok = kb.heartbeat_worker(conn, tid, note=args.get("note"), expected_run_id=_worker_run_id(tid)) + ok = kb.heartbeat_worker( + conn, tid, note=args.get("note"), expected_run_id=_worker_run_id(tid)) if not ok: return tool_error(f"could not heartbeat {tid} (unknown id or not running)") return _ok(task_id=tid) @@ -788,10 +697,8 @@ def _handle_comment(args: dict, **kw) -> str: _reject_delegated_child_mutation("kanban_comment") tid = args.get("task_id") if not tid: - return tool_error( - "task_id is required (use the current task id if that's what " - "you mean — pulls from env but kept explicit here)" - ) + return tool_error("task_id is required (use the current task id if that's what " + "you mean — pulls from env but kept explicit here)") body = _redact(_require_text(args, "body")) # Author comes from the worker's runtime identity, never caller args: # comments are injected into future workers' system prompts as @@ -810,8 +717,7 @@ def _store_attachment(board, tid, filename, data, content_type) -> str: with _board(board) as (kb, conn): att_id = kb.store_attachment_bytes( conn, tid, str(filename), data, - content_type=content_type, uploaded_by="agent", board=board, - ) + content_type=content_type, uploaded_by="agent", board=board) return _ok(task_id=tid, attachment_id=att_id, size=len(data)) @@ -845,9 +751,7 @@ def _download_url_with_cap(url: str, max_bytes: int) -> tuple[bytes, Optional[st streaming, so nothing oversize is buffered). """ from urllib.parse import urljoin, urlparse - import httpx - from tools.url_safety import is_safe_url current_url = url @@ -857,17 +761,11 @@ def _download_url_with_cap(url: str, max_bytes: int) -> tuple[bytes, Optional[st raise ValueError(f"unsupported URL scheme {scheme!r}; only http/https are allowed") if not is_safe_url(current_url): raise ValueError( - f"URL blocked by SSRF protection (private/internal address): {current_url}" - ) + f"URL blocked by SSRF protection (private/internal address): {current_url}") chunks: list[bytes] = [] total = 0 - with httpx.stream( - "GET", - current_url, - headers={"User-Agent": "hermes-kanban/attach"}, - timeout=30, - follow_redirects=False, - ) as resp: + with httpx.stream("GET", current_url, headers={"User-Agent": "hermes-kanban/attach"}, + timeout=30, follow_redirects=False) as resp: if resp.is_redirect: location = resp.headers.get("location") if not location: @@ -906,8 +804,7 @@ def _handle_attach_url(args: dict, **kw) -> str: logger.exception("kanban_attach_url download failed") return tool_error(f"kanban_attach_url: failed to fetch {url}: {e}") return _store_attachment( - args.get("board"), tid, filename, data, args.get("content_type") or fetched_ct - ) + args.get("board"), tid, filename, data, args.get("content_type") or fetched_ct) @_kanban_handler("kanban_attachments") @@ -918,10 +815,9 @@ def _handle_attachments(args: dict, **kw) -> str: if kb.get_task(conn, tid) is None: return tool_error(f"task {tid} not found") return json.dumps({ - "ok": True, - "task_id": tid, - "attachments": [_fields(a, _ATTACHMENT_FIELDS) for a in kb.list_attachments(conn, tid)], - }) + "ok": True, "task_id": tid, + "attachments": [ + _fields(a, _ATTACHMENT_FIELDS) for a in kb.list_attachments(conn, tid)]}) @_kanban_handler("kanban_create") @@ -931,21 +827,16 @@ def _handle_create(args: dict, **kw) -> str: title = _require_text(args, "title") assignee = args.get("assignee") if not assignee: - return tool_error( - "assignee is required — name the profile that should execute this " - "task (the dispatcher will only spawn tasks with an assignee)" - ) + return tool_error("assignee is required — name the profile that should execute this " + "task (the dispatcher will only spawn tasks with an assignee)") # Prefer the request-scoped api_server origin binding over HERMES_SESSION_ID: # the env var is clobbered with a subagent's internal id whenever a child # agent is constructed in-process, which would stamp — and later wake — # the wrong session. NULL on CLI/dashboard paths that set neither. from tools.async_delegation import _current_origin_session_id - session_id = ( - args.get("session_id") - or _current_origin_session_id() - or os.environ.get("HERMES_SESSION_ID") - ) + session_id = (args.get("session_id") or _current_origin_session_id() + or os.environ.get("HERMES_SESSION_ID")) # Workspace sharing is always explicit: omitted fields mean a fresh scratch # workspace even for a dispatcher-spawned creator — reusing the parent's # literal path would let a child mutate review evidence or race its @@ -993,18 +884,14 @@ def _handle_create(args: dict, **kw) -> str: goal_max_turns=_opt_int(args.get("goal_max_turns")), initial_status=str(args.get("initial_status") or "running"), created_by=os.environ.get("HERMES_PROFILE") or "worker", - session_id=session_id, - ) + session_id=session_id) new_task = kb.get_task(conn, new_tid) subscribed = _maybe_auto_subscribe(conn, new_tid) return _ok( - task_id=new_tid, - status=new_task.status if new_task else None, + task_id=new_tid, status=new_task.status if new_task else None, workspace_kind=new_task.workspace_kind if new_task else None, workspace_path=new_task.workspace_path if new_task else None, - project_id=new_task.project_id if new_task else None, - subscribed=subscribed, - ) + project_id=new_task.project_id if new_task else None, subscribed=subscribed) def _resolve_notify_target() -> Optional[dict[str, Any]]: @@ -1025,9 +912,7 @@ def _resolve_notify_target() -> Optional[dict[str, Any]]: chat_id = get_session_env("HERMES_SESSION_CHAT_ID", "") if not platform or not chat_id: session_key = ( - get_session_env("HERMES_SESSION_KEY", "") - or os.environ.get("HERMES_SESSION_KEY", "") - ) + get_session_env("HERMES_SESSION_KEY", "") or os.environ.get("HERMES_SESSION_KEY", "")) if not session_key: return None platform, chat_id = "tui", session_key @@ -1035,9 +920,7 @@ def _resolve_notify_target() -> Optional[dict[str, Any]]: thread_id = get_session_env("HERMES_SESSION_THREAD_ID", "") or None message_id = get_session_env("HERMES_SESSION_MESSAGE_ID", "") or "" notifier_profile = ( - get_session_env("HERMES_SESSION_PROFILE", "") - or os.environ.get("HERMES_PROFILE") - ) + get_session_env("HERMES_SESSION_PROFILE", "") or os.environ.get("HERMES_PROFILE")) if not notifier_profile: try: from hermes_cli.profiles import get_active_profile_name @@ -1052,24 +935,19 @@ def _resolve_notify_target() -> Optional[dict[str, Any]]: if ( platform.lower() == "telegram" and thread_id - and (chat_type or "").lower() in {"dm", "direct", "private"} - ): + and (chat_type or "").lower() in {"dm", "direct", "private"}): delivery_metadata["telegram_dm_topic_reply_fallback"] = True if str(thread_id) not in {"", "1"}: delivery_metadata["direct_messages_topic_id"] = str(thread_id) if message_id: delivery_metadata["telegram_reply_to_message_id"] = str(message_id) return dict( - platform=platform, - chat_id=chat_id, - chat_type=chat_type, - thread_id=thread_id, + platform=platform, chat_id=chat_id, chat_type=chat_type, thread_id=thread_id, user_id=get_session_env("HERMES_SESSION_USER_ID", "") or None, user_id_alt=get_session_env("HERMES_SESSION_USER_ID_ALT", "") or None, notifier_profile=notifier_profile, delivery_mode="notify+wake" if platform != "tui" else None, - delivery_metadata=delivery_metadata or None, - ) + delivery_metadata=delivery_metadata or None) def _maybe_auto_subscribe(conn: Any, task_id: str) -> bool: @@ -1100,8 +978,7 @@ def _maybe_auto_subscribe(conn: Any, task_id: str) -> bool: except Exception as _exc: logger.warning( "_maybe_auto_subscribe failed: %r (platform=%r key_set=%r)", - _exc, target["platform"] if target else "", bool(target and target["chat_id"]), - ) + _exc, target["platform"] if target else "", bool(target and target["chat_id"])) return False @@ -1135,28 +1012,27 @@ def _handle_link(args: dict, **kw) -> str: return _ok(parent_id=parent_id, child_id=child_id) -# --------------------------------------------------------------------------- -# Registration (order preserved: it is the order tools appear in the schema) -# --------------------------------------------------------------------------- +# --- Registration (order preserved: it is the order tools appear in the schema) --- +# kanban_list / kanban_unblock route the board and are hidden from task workers. +_ORCHESTRATOR_TOOLS = frozenset({"kanban_list", "kanban_unblock"}) _TOOLS = ( - ("kanban_show", KANBAN_SHOW_SCHEMA, _handle_show, _check_kanban_mode, "📋"), - ("kanban_list", KANBAN_LIST_SCHEMA, _handle_list, _check_kanban_orchestrator_mode, "📋"), - ("kanban_complete", KANBAN_COMPLETE_SCHEMA, _handle_complete, _check_kanban_mode, "✔"), - ("kanban_block", KANBAN_BLOCK_SCHEMA, _handle_block, _check_kanban_mode, "⏸"), - ("kanban_request_review", KANBAN_REQUEST_REVIEW_SCHEMA, _handle_request_review, _check_kanban_mode, "👀"), - ("kanban_request_changes", KANBAN_REQUEST_CHANGES_SCHEMA, _handle_request_changes, _check_kanban_mode, "↩"), - ("kanban_heartbeat", KANBAN_HEARTBEAT_SCHEMA, _handle_heartbeat, _check_kanban_mode, "💓"), - ("kanban_comment", KANBAN_COMMENT_SCHEMA, _handle_comment, _check_kanban_mode, "💬"), - ("kanban_attach", KANBAN_ATTACH_SCHEMA, _handle_attach, _check_kanban_mode, "📎"), - ("kanban_attach_url", KANBAN_ATTACH_URL_SCHEMA, _handle_attach_url, _check_kanban_mode, "📎"), - ("kanban_attachments", KANBAN_ATTACHMENTS_SCHEMA, _handle_attachments, _check_kanban_mode, "📎"), - ("kanban_create", KANBAN_CREATE_SCHEMA, _handle_create, _check_kanban_mode, "➕"), - ("kanban_unblock", KANBAN_UNBLOCK_SCHEMA, _handle_unblock, _check_kanban_orchestrator_mode, "▶"), - ("kanban_link", KANBAN_LINK_SCHEMA, _handle_link, _check_kanban_mode, "🔗"), -) + ("kanban_show", KANBAN_SHOW_SCHEMA, _handle_show, "📋"), + ("kanban_list", KANBAN_LIST_SCHEMA, _handle_list, "📋"), + ("kanban_complete", KANBAN_COMPLETE_SCHEMA, _handle_complete, "✔"), + ("kanban_block", KANBAN_BLOCK_SCHEMA, _handle_block, "⏸"), + ("kanban_request_review", KANBAN_REQUEST_REVIEW_SCHEMA, _handle_request_review, "👀"), + ("kanban_request_changes", KANBAN_REQUEST_CHANGES_SCHEMA, _handle_request_changes, "↩"), + ("kanban_heartbeat", KANBAN_HEARTBEAT_SCHEMA, _handle_heartbeat, "💓"), + ("kanban_comment", KANBAN_COMMENT_SCHEMA, _handle_comment, "💬"), + ("kanban_attach", KANBAN_ATTACH_SCHEMA, _handle_attach, "📎"), + ("kanban_attach_url", KANBAN_ATTACH_URL_SCHEMA, _handle_attach_url, "📎"), + ("kanban_attachments", KANBAN_ATTACHMENTS_SCHEMA, _handle_attachments, "📎"), + ("kanban_create", KANBAN_CREATE_SCHEMA, _handle_create, "➕"), + ("kanban_unblock", KANBAN_UNBLOCK_SCHEMA, _handle_unblock, "▶"), + ("kanban_link", KANBAN_LINK_SCHEMA, _handle_link, "🔗")) -for _name, _sch, _handler, _check_fn, _emoji in _TOOLS: - registry.register( - name=_name, toolset="kanban", schema=_sch, handler=_handler, check_fn=_check_fn, emoji=_emoji, - ) +for _name, _sch, _handler, _emoji in _TOOLS: + _check = _check_kanban_orchestrator_mode if _name in _ORCHESTRATOR_TOOLS else _check_kanban_mode + registry.register(name=_name, toolset="kanban", schema=_sch, handler=_handler, emoji=_emoji, + check_fn=_check)