From 7f2d6c4bc3463a760f5b22224721a0209208b57d Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Thu, 3 Sep 2026 00:10:18 -0700 Subject: [PATCH] refactor(tools): kanban _delegation_ctx/_own_task_env/_existing_task helpers, None-safe _fields; live-log _best_effort ctx manager; docstring compaction --- tools/debug_helpers.py | 9 +- tools/delegation_live_log.py | 66 ++++++----- tools/delegation_output_schema.py | 25 ++--- tools/desktop_ui.py | 22 ++-- tools/focus_pane_tool.py | 8 +- tools/interrupt.py | 36 +++--- tools/kanban_tools.py | 181 +++++++++++++----------------- 7 files changed, 146 insertions(+), 201 deletions(-) diff --git a/tools/debug_helpers.py b/tools/debug_helpers.py index 516ad29ed9..1fd06f8182 100644 --- a/tools/debug_helpers.py +++ b/tools/debug_helpers.py @@ -47,13 +47,8 @@ class DebugSession: try: filepath = self.log_dir / f"{self.tool_name}_debug_{self.session_id}.json" payload = { - "session_id": self.session_id, - "start_time": self._start_time, - "end_time": _now(), - "debug_enabled": True, - "total_calls": len(self._calls), - "tool_calls": self._calls, - } + "session_id": self.session_id, "start_time": self._start_time, "end_time": _now(), + "debug_enabled": True, "total_calls": len(self._calls), "tool_calls": self._calls} with open(filepath, "w", encoding="utf-8") as f: json.dump(payload, f, indent=2, ensure_ascii=False) logger.debug("%s debug log saved: %s", self.tool_name, filepath) diff --git a/tools/delegation_live_log.py b/tools/delegation_live_log.py index 55d9409a83..40965d9764 100644 --- a/tools/delegation_live_log.py +++ b/tools/delegation_live_log.py @@ -16,6 +16,7 @@ import shutil import threading import time import uuid +from contextlib import contextmanager from pathlib import Path from typing import Any, Dict, List, Optional @@ -42,6 +43,15 @@ def live_transcript_root() -> Path: return get_hermes_dir("cache/delegation", "delegation_cache") / "live" +@contextmanager +def _best_effort(what: str): + """Swallow and debug-log any failure: nothing here may reach the agent loop.""" + try: + yield + except Exception as exc: # noqa: BLE001 + logger.debug("Live transcript %s failed: %s", what, exc) + + 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()) @@ -62,6 +72,10 @@ def _redact(text: str) -> str: return "[line withheld: redaction unavailable]" +def _joined(*parts: str) -> str: + return " ".join(filter(None, parts)) + + def _dump_json(path: Path, payload: Dict[str, Any]) -> None: path.write_text(json.dumps(payload, indent=2, ensure_ascii=False), encoding="utf-8") @@ -74,29 +88,26 @@ class LiveTranscriptWriter: context: Optional[str] = None, root: Optional[Path] = None): self.delegation_id = delegation_id self.task_index = task_index - self._ok = True + self._ok = False self._lock = threading.Lock() self._stream_buf: List[str] = [] self._stream_len = 0 - try: + self.path: Optional[Path] = None + with _best_effort(f"init ({delegation_id} task {task_index})"): goal_line = _one_line(goal, _KICKOFF_MAX) d = (root if root is not None else live_transcript_root()) / delegation_id d.mkdir(parents=True, exist_ok=True) - self.path: Optional[Path] = d / f"task-{task_index}.log" - header = ( + path = d / f"task-{task_index}.log" + path.write_text( "=== Hermes subagent live transcript ===\n" f"delegation: {delegation_id} task: {task_index}\n" 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") - self.path.write_text(header, encoding="utf-8") + + "=" * 40 + "\n", encoding="utf-8") + self.path, self._ok = path, True self.event("user", "kickoff: " + goal_line + (f" | context: {_one_line(context, _KICKOFF_MAX)}" if context else "")) - except Exception as exc: - logger.debug("Live transcript init failed (%s task %s): %s", delegation_id, task_index, exc) - self._ok = False - self.path = None def event(self, role: str, text: str) -> None: """Append one ``HH:MM:SS role | text`` line. Single choke point: every typed @@ -154,13 +165,12 @@ class LiveTranscriptWriter: self.assistant_text(text) def _on_complete(self, tool_name, preview, args, kwargs): - self.flush_stream() dur = kwargs.get("duration_seconds") summary = kwargs.get("summary") or preview - self.marker(" ".join(filter(None, [ + self.marker(_joined( 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 = { @@ -187,11 +197,11 @@ class LiveTranscriptWriter: def finalize(self, entry: Dict[str, Any]) -> None: """Terminal marker with exit-reason detail subagent.complete lacks.""" exit_reason = entry.get("exit_reason") - self.marker(" ".join(filter(None, [ + self.marker(_joined( 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): @@ -199,18 +209,14 @@ def wrap_progress_callback(inner_cb, writer: LiveTranscriptWriter): the log; writer failures never propagate. Preserves the ``_flush`` contract.""" def _cb(event_type, tool_name=None, preview=None, args=None, **kwargs): - try: + with _best_effort("observe"): writer.observe(event_type, tool_name, preview, args, **kwargs) - except Exception as exc: # noqa: BLE001 — must never hit the agent loop - logger.debug("Live transcript observe failed: %s", exc) if inner_cb is not None: inner_cb(event_type, tool_name, preview, args, **kwargs) def _flush(): - try: + with _best_effort("flush"): writer.flush_stream() - except Exception: - pass if callable(getattr(inner_cb, "_flush", None)): inner_cb._flush() @@ -228,7 +234,7 @@ def create_live_transcripts( ``(None, [None]*n, [])`` so delegation proceeds untouched.""" n = len(task_list) prune_stale_live_dirs() # best-effort; never raises - try: + with _best_effort("creation"): # 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) @@ -239,9 +245,7 @@ def create_live_transcripts( return None, [None] * n, [] _write_manifest(deleg_id, task_list, paths, model=model, provider=provider) return deleg_id, writers, paths - except Exception as exc: - logger.debug("Live transcript creation failed: %s", exc) - return None, [None] * n, [] + return None, [None] * n, [] def _manifest_path(delegation_id: str) -> Path: @@ -251,7 +255,7 @@ def _manifest_path(delegation_id: str) -> Path: def _write_manifest(delegation_id: str, task_list: List[Dict[str, Any]], paths: List[str], model: Optional[str] = None, provider: Optional[str] = None) -> None: - try: + with _best_effort("manifest write"): _dump_json(_manifest_path(delegation_id), { "delegation_id": delegation_id, "started": time.strftime(_TIME_FMT), "task_count": len(task_list), "model": model, "provider": provider, @@ -261,8 +265,6 @@ 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)]}) - except Exception as exc: - logger.debug("Live transcript manifest write failed: %s", exc) def update_manifest_statuses(delegation_id: Optional[str], @@ -270,7 +272,7 @@ def update_manifest_statuses(delegation_id: Optional[str], """Best-effort per-task status update once the batch has aggregated.""" if not delegation_id: return - try: + with _best_effort("manifest update"): mp = _manifest_path(delegation_id) manifest = json.loads(mp.read_text(encoding="utf-8")) by_index = {r.get("task_index"): r for r in results if isinstance(r, dict)} @@ -282,14 +284,12 @@ def update_manifest_statuses(delegation_id: Optional[str], task["exit_reason"] = r["exit_reason"] manifest["completed"] = time.strftime(_TIME_FMT) _dump_json(mp, manifest) - except Exception as exc: - logger.debug("Live transcript manifest update failed: %s", exc) def prune_stale_live_dirs(max_age_days: int = LIVE_RETENTION_DAYS) -> int: """Remove live/ dirs older than the retention window. Best-effort.""" removed = 0 - try: + with _best_effort("pruning"): root = live_transcript_root() if not root.is_dir(): return 0 @@ -301,6 +301,4 @@ def prune_stale_live_dirs(max_age_days: int = LIVE_RETENTION_DAYS) -> int: removed += 1 except OSError: continue - except Exception as exc: - logger.debug("Live transcript pruning failed: %s", exc) return removed diff --git a/tools/delegation_output_schema.py b/tools/delegation_output_schema.py index eae47ce8a9..1540520792 100644 --- a/tools/delegation_output_schema.py +++ b/tools/delegation_output_schema.py @@ -24,12 +24,11 @@ def coerce_output_schema(raw: Any) -> Tuple[Optional[Dict[str, Any]], Optional[s if isinstance(raw, str): # Models sometimes double-encode the schema as a JSON string. try: - parsed = json.loads(raw) + raw = json.loads(raw) except (ValueError, TypeError): return None, "output_schema must be a JSON Schema object, got a non-JSON string." - if not isinstance(parsed, dict): + if not isinstance(raw, dict): return None, "output_schema must be a JSON Schema object." - raw = parsed if not isinstance(raw, dict): return None, f"output_schema must be a JSON Schema object, got {type(raw).__name__}." try: @@ -49,12 +48,10 @@ def append_output_contract(context: Optional[str], schema: Dict[str, Any]) -> st schema_text = json.dumps(schema, indent=2, ensure_ascii=False) except (TypeError, ValueError): schema_text = str(schema) - block = ( - "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}") + block = ("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}") base = (context or "").rstrip() return f"{base}\n\n{block}" if base else block @@ -104,9 +101,7 @@ def validate_output(text: str, schema: Dict[str, Any]) -> Tuple[bool, List[str]] def build_retry_message(errors: List[str]) -> str: """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.") + 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.") diff --git a/tools/desktop_ui.py b/tools/desktop_ui.py index 0c7c23db92..6f9b013430 100644 --- a/tools/desktop_ui.py +++ b/tools/desktop_ui.py @@ -28,14 +28,11 @@ def available() -> bool: 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 - 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. - """ + """Read a desktop Appearance switch from ``display.``. The renderer mirrors + these toggles onto the CONNECTED gateway's config, so this 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. 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") @@ -48,10 +45,9 @@ def user_enabled(setting: str, default: bool) -> bool: def emit(event: str, payload: dict) -> bool: """Route ``event`` to the window owning the current turn; False when no emitter.""" - fn = _emit - if fn is None: + if _emit is None: return False - fn(get_session_env("HERMES_UI_SESSION_ID", ""), event, payload) + _emit(get_session_env("HERMES_UI_SESSION_ID", ""), event, payload) return True @@ -63,9 +59,7 @@ def emit_or_error(event: str, payload: dict, fail_prefix: str, desktop_only: str ok = emit(event, payload) except Exception as exc: return tool_error(f"{fail_prefix}{exc}") - if not ok: - return tool_error(desktop_only) - return json.dumps(result, ensure_ascii=False) + return json.dumps(result, ensure_ascii=False) if ok else tool_error(desktop_only) def passthrough_json(raw) -> str: diff --git a/tools/focus_pane_tool.py b/tools/focus_pane_tool.py index 5226879873..2ff5477aa8 100644 --- a/tools/focus_pane_tool.py +++ b/tools/focus_pane_tool.py @@ -17,12 +17,8 @@ def focus_pane_tool(pane: str) -> str: if name not in PANES: return tool_error(f"pane must be one of: {', '.join(PANES)}.") return desktop_ui.emit_or_error( - "pane.reveal", - {"pane": name}, - f"Failed to focus the {name} pane: ", - "Pane focus is only available in the Hermes desktop app.", - {"success": True, "pane": name}, - ) + "pane.reveal", {"pane": name}, f"Failed to focus the {name} pane: ", + "Pane focus is only available in the Hermes desktop app.", {"success": True, "pane": name}) registry.register( diff --git a/tools/interrupt.py b/tools/interrupt.py index f779ecdc4b..814fef6b77 100644 --- a/tools/interrupt.py +++ b/tools/interrupt.py @@ -1,10 +1,7 @@ -"""Per-thread interrupt signaling for all tools. - -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. -""" +"""Per-thread interrupt signaling for all tools: thread-scoped so interrupting one +agent session does not kill tools in other sessions (the gateway runs many agents in one +process). The agent passes its execution thread id to set_interrupt(); tools call +is_interrupted(), which checks the CURRENT thread.""" import logging import os @@ -26,8 +23,8 @@ _lock = threading.Lock() 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.""" + """Set or clear the interrupt for *thread_id* (default: current thread); ``reason`` is + an optional user-safe cause.""" tid = thread_id if thread_id is not None else threading.current_thread().ident with _lock: (_interrupted_threads.add if active else _interrupted_threads.discard)(tid) @@ -44,7 +41,6 @@ def set_interrupt(active: bool, thread_id: int | None = None, *, reason: str | N def is_interrupted() -> bool: - """Check if an interrupt has been requested for the current thread.""" return is_thread_interrupted(threading.current_thread().ident) @@ -65,21 +61,18 @@ def get_interrupt_reason() -> str | None: def clear_current_thread_interrupt() -> None: - """Clear any interrupt bit on the CURRENT thread. - - 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). - """ + """Clear any interrupt bit on the CURRENT thread: 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. 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) class _ThreadAwareEventProxy: - """Backward-compatible ``_interrupt_event``: legacy call sites call - .is_set()/.set()/.clear(); the shim maps those to the per-thread API.""" + """Backward-compatible ``_interrupt_event``: legacy .is_set()/.set()/.clear()/.wait() + call sites mapped onto the per-thread API (``wait`` returns the current state at once).""" def is_set(self) -> bool: return is_interrupted() @@ -91,7 +84,6 @@ class _ThreadAwareEventProxy: set_interrupt(False) def wait(self, timeout: float | None = None) -> bool: - """Not truly supported — returns current state immediately.""" return self.is_set() diff --git a/tools/kanban_tools.py b/tools/kanban_tools.py index 3a9ae789e7..8b69793c30 100644 --- a/tools/kanban_tools.py +++ b/tools/kanban_tools.py @@ -42,34 +42,31 @@ def _profile_has_kanban_toolset() -> bool: return False -def _is_delegated_child_context() -> bool: +def _delegation_ctx(predicate: str, default: bool) -> bool: + """``agent.delegation_context.()``; ``default`` when it cannot be evaluated.""" try: - from agent.delegation_context import is_delegated_child_context - return is_delegated_child_context() + from agent import delegation_context + return getattr(delegation_context, predicate)() except Exception: - return False + return default + + +def _is_delegated_child_context() -> bool: + return _delegation_ctx("is_delegated_child_context", False) def _is_dispatcher_owned_worker() -> bool: """False for delegate_task children AND for cron jobs fired in-process from a worker — i.e. whenever HERMES_KANBAN_* is present but not ours.""" - try: - from agent.delegation_context import is_dispatcher_owned_worker_context - return is_dispatcher_owned_worker_context() - except Exception: - return True - - -def _is_env_worker() -> bool: - """True only for a dispatcher-spawned worker scoped to HERMES_KANBAN_TASK.""" - return bool(os.environ.get("HERMES_KANBAN_TASK")) and _is_dispatcher_owned_worker() + return _delegation_ctx("is_dispatcher_owned_worker_context", True) def _visible(*, to_env_worker: bool) -> bool: - """check_fn core: never for delegate children; env workers per flag; else profile toolset.""" + """check_fn core: never for delegate children; dispatcher-spawned env workers + (HERMES_KANBAN_TASK) per flag; else the profile toolset decides.""" if _is_delegated_child_context(): return False - if _is_env_worker(): + if os.environ.get("HERMES_KANBAN_TASK") and _is_dispatcher_owned_worker(): return to_env_worker return _profile_has_kanban_toolset() @@ -80,16 +77,12 @@ def _check_kanban_mode() -> bool: def _check_kanban_orchestrator_mode() -> bool: - """Board-routing tools (kanban_list, kanban_unblock): hidden from task workers, - who close their own task via complete/block/heartbeat.""" + """Board-routing tools (kanban_list, kanban_unblock): hidden from task workers.""" return _visible(to_env_worker=False) # --- 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)" - - class _Reject(Exception): """Carries a finished ``tool_error`` payload out of a validation helper.""" @@ -123,9 +116,8 @@ def _kanban_handler(tool_name: str) -> Callable: def _reject_delegated_child_mutation(tool_name: str) -> None: - """A delegate_task child shares the parent's process, so inherited - HERMES_KANBAN_* env is not proof of ownership: it may report findings but - must not mutate board state.""" + """A delegate_task child shares the parent's process, so inherited HERMES_KANBAN_* + env is not proof of ownership: it may report findings but must not mutate.""" if _is_delegated_child_context(): raise _Reject( f"{tool_name} refused: delegate_task child agents are not Kanban run owners. " @@ -145,27 +137,28 @@ def _default_task_id(arg: Optional[str]) -> Optional[str]: def _require_task_id(args: dict) -> str: tid = _default_task_id(args.get("task_id")) - _check(tid, _TASK_ID_REQUIRED) + _check(tid, "task_id is required (or set HERMES_KANBAN_TASK in the env)") return tid +def _own_task_env(task_id: str, var: str) -> Optional[str]: + """``$var`` only when this worker is scoped to ``task_id``; else None.""" + return os.environ.get(var) if os.environ.get("HERMES_KANBAN_TASK") == task_id else None + + def _worker_run_id(task_id: str) -> Optional[int]: """This worker's dispatcher run id when it is scoped to task_id.""" - raw = os.environ.get("HERMES_KANBAN_RUN_ID") - if os.environ.get("HERMES_KANBAN_TASK") != task_id or not raw: - return None + raw = _own_task_env(task_id, "HERMES_KANBAN_RUN_ID") try: - return int(raw) + return int(raw) if raw else None except ValueError: return None def _stamp_worker_session_metadata(task_id: str, metadata: Optional[dict]) -> Optional[dict]: """Add trusted worker session id metadata for this worker's own task.""" - session_id = os.environ.get("HERMES_SESSION_ID") - if os.environ.get("HERMES_KANBAN_TASK") != task_id or not session_id: - return metadata - return {**(metadata or {}), "worker_session_id": session_id} + session_id = _own_task_env(task_id, "HERMES_SESSION_ID") + return {**(metadata or {}), "worker_session_id": session_id} if session_id else metadata def _enforce_worker_task_ownership(tid: str) -> None: @@ -215,6 +208,12 @@ def _board(board: Optional[str], *, quiet_close: bool = False): raise +def _existing_task(kb, conn, tid: str): + task = kb.get_task(conn, tid) + _check(task is not None, f"task {tid} not found") + return task + + def _ok(**fields: Any) -> str: return json.dumps({"ok": True, **fields}) @@ -264,10 +263,9 @@ def _require_dict_metadata(metadata: Any) -> None: def _merge_artifacts(metadata: Any, artifacts: list[str]) -> dict: - """Fold ``artifacts`` into ``metadata["artifacts"]`` (merged with, never - overwriting, a list the worker passed manually). Artifacts ride inside - metadata so the completed-event payload needs no DB schema change; the - gateway notifier uploads each path as a native attachment.""" + """Fold ``artifacts`` into ``metadata["artifacts"]`` (merged with, never overwriting, a + list the worker passed manually). Artifacts ride inside metadata so the completed-event + payload needs no DB schema change; the gateway notifier uploads each as an attachment.""" _require_dict_metadata(metadata) metadata = {} if metadata is None else metadata existing = metadata.get("artifacts") @@ -314,10 +312,12 @@ _COMMENT_FIELDS = ("author", "body", "created_at") _EVENT_FIELDS = ("kind", "payload", "created_at", "run_id") _ATTACHMENT_FIELDS = tuple( "id filename content_type size uploaded_by stored_path created_at".split()) +_CREATED_FIELDS = ("status", "workspace_kind", "workspace_path", "project_id") def _fields(obj: Any, names: tuple[str, ...]) -> dict[str, Any]: - return {n: getattr(obj, n) for n in names} + """``{name: getattr(obj, name)}``; every value None when ``obj`` is None.""" + return {n: getattr(obj, n) if obj is not None else None for n in names} def _task_summary_dict(kb, conn, task) -> dict[str, Any]: @@ -474,8 +474,7 @@ def _handle_show(args: dict, **kw) -> str: """Full task state: row, parents, children, comments, runs, last 50 events.""" tid = _require_task_id(args) with _board(args.get("board")) as (kb, conn): - task = kb.get_task(conn, tid) - _check(task is not None, f"task {tid} not found") + task = _existing_task(kb, conn, tid) return json.dumps({ "task": _fields(task, _TASK_FIELDS), "parents": kb.parent_ids(conn, tid), @@ -578,12 +577,11 @@ def _handle_block(args: dict, **kw) -> str: # would be an escape hatch around the completion judge: goal_mode tasks # may only block on genuine external blockers. task = kb.get_task(conn, tid) - 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}). If the task is actually " - f"finished or cannot proceed for another reason, call kanban_complete instead — " - f"the completion judge will evaluate it.") + _check(not (task and task.goal_mode and kind not in _GOAL_MODE_BLOCK_ALLOWED_KINDS), + f"goal_mode tasks can only block with kind in " + 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)) _check(ok, f"could not block {tid} (unknown id or not in running/ready)") return _ok_landed(kb, conn, tid, "blocked", block_kind=kind) @@ -602,10 +600,8 @@ def _handle_request_review(args: dict, **kw) -> str: metadata = _redact_metadata(metadata) _check(metadata is not None, "metadata could not be safely serialized") metadata = _stamp_worker_session_metadata(tid, metadata) - reviewer = args.get("reviewer") or None - if reviewer: - # Model-supplied free text stored durably on the event payload. - reviewer = _redact(reviewer) + # Reviewer is model-supplied free text stored durably on the event payload. + reviewer = _redact_opt(args.get("reviewer") or None) 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( @@ -693,12 +689,10 @@ _MAX_ATTACH_URL_REDIRECTS = 5 def _download_url_with_cap(url: str, max_bytes: int) -> tuple[bytes, Optional[str]]: """Fetch ``url`` over http(s) capped at ``max_bytes`` -> ``(data, content_type)``. - - Every hop is SSRF-checked (redirects followed manually) so a model-controlled URL, - or a public host 302ing, cannot reach loopback/private/cloud-metadata ranges. - ``ValueError`` for bad scheme, blocked target, too many redirects, or a body over - the cap (checked while streaming, so nothing oversize is buffered). - """ + Every hop is SSRF-checked (redirects followed manually) so a model-controlled URL, or a + public host 302ing, cannot reach loopback/private/cloud-metadata ranges. ``ValueError`` + for bad scheme, blocked target, too many redirects, or a body over the cap (checked + while streaming, so nothing oversize is buffered).""" from urllib.parse import urljoin, urlparse import httpx from tools.url_safety import is_safe_url @@ -741,8 +735,7 @@ def _handle_attach_url(args: dict, **kw) -> str: if not filename or not str(filename).strip(): # Derive a name from the URL path's leaf component. from urllib.parse import unquote, urlparse - leaf = unquote(urlparse(url).path.rsplit("/", 1)[-1]).strip() - filename = leaf or "download" + filename = unquote(urlparse(url).path.rsplit("/", 1)[-1]).strip() or "download" try: data, fetched_ct = _download_url_with_cap(url, kb.KANBAN_ATTACHMENT_MAX_BYTES) except ValueError as e: @@ -759,7 +752,7 @@ def _handle_attachments(args: dict, **kw) -> str: """List a task's attachments (read-only; no ownership restriction).""" tid = _require_task_id(args) with _board(args.get("board")) as (kb, conn): - _check(kb.get_task(conn, tid) is not None, f"task {tid} not found") + _existing_task(kb, conn, tid) return json.dumps({ "ok": True, "task_id": tid, "attachments": [ @@ -784,25 +777,21 @@ def _handle_create(args: dict, **kw) -> str: # even for a dispatcher-spawned creator (reusing the parent's path would let a child # mutate review evidence or race its checkout). Project identity is the one safe thing # to inherit implicitly (the DB turns it into a fresh per-task worktree). - workspace_kind = args.get("workspace_kind") - workspace_path = args.get("workspace_path") + workspace_kind, workspace_path = args.get("workspace_kind"), args.get("workspace_path") project_id = args.get("project") or args.get("project_id") project_source_task_id = None - inherit_project = workspace_kind is None and workspace_path is None - triage = _parse_bool_arg(args, "triage") - skills = _coerce_str_list(args.get("skills"), "skills", "skill names") - goal_mode = _parse_bool_arg(args, "goal_mode") - model_override = args.get("model") - provider_override = args.get("provider") + triage, skills, goal_mode = ( + _parse_bool_arg(args, "triage"), _coerce_str_list(args.get("skills"), "skills", "skill names"), + _parse_bool_arg(args, "goal_mode")) + model_override, provider_override = args.get("model"), args.get("provider") _check(model_override or not provider_override, "'provider' requires 'model' to be set as well") parents = _coerce_str_list(args.get("parents") or [], "parents", "task ids") with _board(args.get("board")) as (kb, conn): - if inherit_project and project_id is None: + if project_id is None and workspace_kind is None and workspace_path is None: self_tid = os.environ.get("HERMES_KANBAN_TASK") self_task = kb.get_task(conn, self_tid) if self_tid else None if self_task is not None and self_task.project_id: - project_id = self_task.project_id - project_source_task_id = self_task.id + project_id, project_source_task_id = self_task.project_id, self_task.id new_tid = kb.create_task( conn, title=str(title).strip(), body=args.get("body"), assignee=str(assignee), parents=tuple(parents), tenant=args.get("tenant") or os.environ.get("HERMES_TENANT"), @@ -816,48 +805,35 @@ def _handle_create(args: dict, **kw) -> str: goal_mode=goal_mode, 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) - 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, - 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) + landed = _fields(kb.get_task(conn, new_tid), _CREATED_FIELDS) + return _ok(task_id=new_tid, **landed, subscribed=_maybe_auto_subscribe(conn, new_tid)) def _resolve_notify_target() -> Optional[dict[str, Any]]: """``kanban_db.add_notify_sub`` kwargs for the calling session, or None (CLI/cron/tests). - Gateway sessions: ``HERMES_SESSION_PLATFORM``/``CHAT_ID`` ContextVars. TUI/desktop: those are cleared but the subprocess inherits ``HERMES_SESSION_KEY`` -> ``platform="tui"`` for the TUI poller. ``HERMES_SESSION_ID`` is deliberately NOT a fallback: it is set for - every CLI/ACP invocation and would auto-subscribe every CLI run. - """ - from gateway.session_context import get_session_env - platform = get_session_env("HERMES_SESSION_PLATFORM", "") - chat_id = get_session_env("HERMES_SESSION_CHAT_ID", "") + every CLI/ACP invocation and would auto-subscribe every CLI run.""" + from gateway.session_context import get_session_env as env + platform, chat_id = env("HERMES_SESSION_PLATFORM", ""), 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", "")) + session_key = env("HERMES_SESSION_KEY", "") or os.environ.get("HERMES_SESSION_KEY", "") if not session_key: return None platform, chat_id = "tui", session_key - chat_type = get_session_env("HERMES_SESSION_CHAT_TYPE", "") or None - 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")) + chat_type = env("HERMES_SESSION_CHAT_TYPE", "") or None + thread_id = env("HERMES_SESSION_THREAD_ID", "") or None + message_id = env("HERMES_SESSION_MESSAGE_ID", "") or "" + notifier_profile = env("HERMES_SESSION_PROFILE", "") or os.environ.get("HERMES_PROFILE") if not notifier_profile: try: from hermes_cli.profiles import get_active_profile_name notifier_profile = get_active_profile_name() or "default" except Exception: notifier_profile = "default" - delivery_metadata: dict[str, Any] = {} - if thread_id: - delivery_metadata["thread_id"] = thread_id - if chat_type: - delivery_metadata["chat_type"] = chat_type + delivery_metadata: dict[str, Any] = { + k: v for k, v in (("thread_id", thread_id), ("chat_type", chat_type)) if v} if (platform.lower() == "telegram" and thread_id and (chat_type or "").lower() in {"dm", "direct", "private"}): delivery_metadata["telegram_dm_topic_reply_fallback"] = True @@ -867,8 +843,8 @@ def _resolve_notify_target() -> Optional[dict[str, Any]]: 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, - user_id=get_session_env("HERMES_SESSION_USER_ID", "") or None, - user_id_alt=get_session_env("HERMES_SESSION_USER_ID_ALT", "") or None, + user_id=env("HERMES_SESSION_USER_ID", "") or None, + user_id_alt=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) @@ -880,8 +856,7 @@ def _maybe_auto_subscribe(conn: Any, task_id: str) -> bool: ``kanban_notify-subscribe``). Gated by ``kanban.auto_subscribe_on_create`` (default True). Failures are logged and swallowed: bookkeeping must never fail kanban_create.""" try: - cfg = load_config() - if not cfg_get(cfg, "kanban", "auto_subscribe_on_create", default=True): + if not cfg_get(load_config(), "kanban", "auto_subscribe_on_create", default=True): return False except Exception: pass # unreadable config keeps the user-friendly default (True) @@ -907,11 +882,11 @@ def _handle_unblock(args: dict, **kw) -> str: _require_orchestrator_tool("kanban_unblock") tid = args.get("task_id") _check(tid, "task_id is required") - _enforce_worker_task_ownership(str(tid)) + tid = str(tid) + _enforce_worker_task_ownership(tid) with _board(args.get("board")) as (kb, conn): - _check(kb.unblock_task(conn, str(tid)), f"could not unblock {tid} (not blocked or unknown)") - task = kb.get_task(conn, str(tid)) - return _ok(task_id=str(tid), status=task.status if task else None) + _check(kb.unblock_task(conn, tid), f"could not unblock {tid} (not blocked or unknown)") + return _ok(task_id=tid, **_fields(kb.get_task(conn, tid), ("status",))) @_kanban_handler("kanban_link")