refactor(tools): kanban _delegation_ctx/_own_task_env/_existing_task helpers, None-safe _fields; live-log _best_effort ctx manager; docstring compaction
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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/<delegation_id> 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
|
||||
|
||||
@@ -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.")
|
||||
|
||||
+8
-14
@@ -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.<setting>``.
|
||||
|
||||
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.<setting>``. 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:
|
||||
|
||||
@@ -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(
|
||||
|
||||
+14
-22
@@ -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()
|
||||
|
||||
|
||||
|
||||
+78
-103
@@ -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.<predicate>()``; ``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")
|
||||
|
||||
Reference in New Issue
Block a user