refactor(tools): compact kanban/delegation/interrupt/desktop tool modules; drop dead helpers

This commit is contained in:
Teknium
2026-09-02 22:06:11 -07:00
parent 4ab6bdb94b
commit 26186cb0b8
8 changed files with 245 additions and 469 deletions

View File

@@ -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,

View File

@@ -1,15 +1,11 @@
"""Live, tail-able transcripts for delegated subagents.
Each ``delegate_task`` dispatch creates one append-only log per child under
``<hermes_home>/cache/delegation/live/<delegation_id>/task-<n>.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 ``<hermes_home>/cache/delegation/live/
<delegation_id>/task-<n>.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>) — 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", <text>, 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)

View File

@@ -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.")

View File

@@ -31,15 +31,13 @@ def user_enabled(setting: str, default: bool) -> bool:
"""Read one of the desktop's Appearance switches from ``display.<setting>``.
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:

View File

@@ -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="🪟",
)

View File

@@ -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()

View File

@@ -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()

View File

@@ -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)