diff --git a/tools/delegate_tool.py b/tools/delegate_tool.py index b61ee0c28e..dedd7292c4 100644 --- a/tools/delegate_tool.py +++ b/tools/delegate_tool.py @@ -82,13 +82,11 @@ def check_delegate_requirements() -> bool: def _open_child_session_db(parent_agent) -> Any: - """DEDICATED SessionDB handle for the child, or None. - - The parent's handle can be closed by its own lifecycle while a background - child still flushes (transcript silently dropped). It MUST open the same db - FILE as the parent's handle (non-launch profiles), else lineage / - session_search break; released by the child's close() via _owns_session_db. - """ + """DEDICATED SessionDB handle for the child, or None: the parent's handle can be + closed by its own lifecycle while a background child still flushes (transcript + silently dropped). It MUST open the same db FILE as the parent's handle + (non-launch profiles), else lineage / session_search break; released by the + child's close() via _owns_session_db.""" parent_session_db = getattr(parent_agent, "_session_db", None) if parent_session_db is None: return None @@ -120,11 +118,9 @@ def _build_child_agent( # Legacy; accepted for wire compat but ignored (capability is depth-derived). role: str = "leaf", ): - """Build (don't run) a child AIAgent on the main thread. - - override_* (from delegation config) replace parent inheritance so children - can run on a different provider:model pair. - """ + """Build (don't run) a child AIAgent on the main thread. override_* (from + delegation config) replace parent inheritance so children can run on a + different provider:model pair.""" import uuid as _uuid from run_agent import AIAgent from agent.delegation_context import delegated_child_context @@ -347,13 +343,11 @@ def delegate_task( output_schema: Optional[Dict[str, Any]] = None, action: Optional[str] = None, subagent_id: Optional[str] = None, message: Optional[str] = None, parent_agent=None, credentials_cfg: Optional[Dict[str, Any]] = None, ) -> str: - """Spawn child agents (single ``goal`` or ``tasks=[...]`` batch) or control running ones. - - ``action`` list/steer/stop run synchronously and bypass the pause gate, - depth limit and async dispatch. ``role`` is legacy (per-task beats - top-level; capability is depth-derived). Returns JSON with one results - entry per task, or a dispatch handle when running in the background. - """ + """Spawn child agents (single ``goal`` or ``tasks=[...]`` batch) or control + running ones. ``action`` list/steer/stop run synchronously and bypass the pause + gate, depth limit and async dispatch. ``role`` is legacy (per-task beats + top-level; capability is depth-derived). Returns JSON with one results entry + per task, or a dispatch handle when running in the background.""" if parent_agent is None: return tool_error("delegate_task requires a parent agent context.") @@ -597,14 +591,12 @@ from tools.registry import registry, tool_error def _model_background_value(args: dict, parent_agent=None) -> bool: """Background flag for the MODEL-facing dispatch path (registry fallback). - - Top-level delegations always run in the background — the model does not - choose — for single tasks and fan-out batches alike (one async unit, one - consolidated result). The exception is an orchestrator subagent (depth > 0), - which needs its workers' results within its own turn. The live path is + Top-level delegations always run in the background — the model does not choose + — for single tasks and fan-out batches alike (one async unit, one consolidated + result); an orchestrator subagent (depth > 0) is the exception since it needs + its workers' results within its own turn. The live path is ``run_agent._dispatch_delegate_task``; this mirrors it for the rare case the - intercept is bypassed. Direct Python callers keep the synchronous default. - """ + intercept is bypassed. Direct Python callers keep the synchronous default.""" return not getattr(parent_agent, "_delegate_depth", 0) > 0 _MODEL_HIDDEN_TASK_FIELDS = {"acp_command", "acp_args"} diff --git a/tools/delegate_tool_child_run.py b/tools/delegate_tool_child_run.py index f4e7c2e047..1e46868bf1 100644 --- a/tools/delegate_tool_child_run.py +++ b/tools/delegate_tool_child_run.py @@ -155,12 +155,10 @@ def _dump_subagent_timeout_diagnostic( *, child: Any, task_index: int, timeout_seconds: float, duration_seconds: float, worker_thread: Optional[threading.Thread], goal: str, ) -> Optional[str]: - """Write a structured diagnostic for a subagent that timed out before any - API call (users hit "timed out with no response" and 0 API calls with no way - to inspect it). Lands under ``~/.hermes/logs/subagent-timeout--.log`` - with the child's config, prompt/schema sizes, activity snapshot and the - worker thread's stack. Returns the path, or None on failure. - """ + """Structured diagnostic for a subagent that timed out before any API call + (otherwise "timed out with no response", 0 API calls, nothing to inspect): + ``~/.hermes/logs/subagent-timeout--.log`` with the child's config, + prompt/schema sizes, activity snapshot and worker stack. Path, or None on failure.""" try: from hermes_constants import get_hermes_home import datetime as _dt @@ -213,12 +211,9 @@ def _dump_subagent_timeout_diagnostic( # ── Per-run helpers ────────────────────────────────────────────────────────── def _start_heartbeat(child: Any, parent_agent: Any, task_index: int) -> tuple: - """Build the parent-activity heartbeat thread for one child (not started). - - Returns ``(stop_event, thread)``. The caller starts the thread inside its - ``try`` so a failed ``start()`` (OS thread exhaustion) leaves ``ident`` None - and the finally-path join can be skipped safely. - """ + """``(stop_event, thread)`` for one child's parent-activity heartbeat, NOT + started: the caller starts it inside its ``try`` so a failed ``start()`` (OS + thread exhaustion) leaves ``ident`` None and the finally-path join is skipped.""" from tools.delegate_tool import (_HEARTBEAT_INTERVAL, _HEARTBEAT_STALE_CYCLES_IDLE, _HEARTBEAT_STALE_CYCLES_IN_TOOL) _heartbeat_stop = threading.Event() # Stale detection: a cycle counts as stale when (tool, iteration, @@ -269,11 +264,9 @@ def _register_child( child: Any, parent_agent: Any, goal: str, *, owner_session_id: Optional[str], owner_transport: Any, owner_session_record: Any, ) -> Optional[str]: - """Register the live child in the module registry; return its subagent_id. - - Test doubles without a stable string ``_subagent_id`` are not registered - (returns None) and the caller skips every registry interaction for them. - """ + """Register the live child in the module registry; return its subagent_id. Test + doubles without a stable string ``_subagent_id`` are not registered (None) and + the caller skips every registry interaction for them.""" _subagent_id = getattr(child, "_subagent_id", None) if not isinstance(_subagent_id, str) or not _subagent_id: return None @@ -340,14 +333,14 @@ def _create_isolated_worktree(parent_agent: Any, parent_task_id: Any, subagent_i def _defer_close_after_timeout(child: Any, child_future: Any) -> None: """Hand ``child.close()`` to a Future done-callback and drain its transports. - The interrupt is cooperative: the worker still runs its finally path, so - closing now could close SQLite under its final write — the done-callback is - the first safe boundary. The abandoned worker is usually parked in an - OpenSSL read; NEVER hard-close that transport from this thread (cross-thread - FD release under a live SSL read corrupts native state) — shutdown() the - pooled sockets so the read settles with EOF and the worker unwinds. One - immediate sweep + one delayed re-sweep for a connection opened in between; a - worker that still won't settle keeps its resources until process exit. + The interrupt is cooperative: the worker still runs its finally path, so closing + now could close SQLite under its final write — the done-callback is the first + safe boundary. The abandoned worker is usually parked in an OpenSSL read; NEVER + hard-close that transport from this thread (cross-thread FD release under a + live SSL read corrupts native state) — shutdown() the pooled sockets so the read + settles with EOF and the worker unwinds. One immediate sweep + one delayed + re-sweep for a connection opened in between; a worker that still won't settle + keeps its resources until process exit. """ child_future.add_done_callback(lambda _done: _close_child(child, "Failed to close timed-out child after worker exit")) _drain = getattr(child, "_drain_transports_after_abandonment", None) @@ -397,11 +390,9 @@ class _SchemaOutcome: def _validate_child_output_schema( child: Any, result: Dict[str, Any], task_index: int, child_task_id: str, relay_child_text: Any ) -> _SchemaOutcome: - """Validate the final answer against the attached output_schema with ONE bounded retry. - - Schema-less children (no dict on ``child._delegate_output_schema``) take no - branch here so their result entry stays byte-identical. - """ + """Validate the final answer against the attached output_schema with ONE bounded + retry. Schema-less children (no dict on ``child._delegate_output_schema``) take + no branch here so their result entry stays byte-identical.""" _output_schema = getattr(child, "_delegate_output_schema", None) if not isinstance(_output_schema, dict): return _SchemaOutcome(_output_schema, None, [], 0) @@ -469,12 +460,10 @@ def _build_tool_trace(messages: Any) -> list[Dict[str, Any]]: def _build_result_entry( child: Any, result: Dict[str, Any], task_index: int, duration: float, schema: _SchemaOutcome, ) -> Dict[str, Any]: - """Derive the parent-visible result entry (status, exit_reason, tool trace, tokens, cost). - - ``status`` / ``exit_reason`` / ``truncated`` follow the contract in the - ``_run_single_child`` docstring; a structured failure always wins over the - summary-presence heuristic (which is only a fallback for legacy/mock results). - """ + """Parent-visible result entry (status, exit_reason, tool trace, tokens, cost). + ``status``/``exit_reason``/``truncated`` follow the ``_run_single_child`` contract; + a structured failure always wins over the summary-presence heuristic (a fallback + for legacy/mock results only).""" summary = result.get("final_response") or "" # "(empty)" is run_agent's give-up sentinel after repeated empty LLM # responses (usually a transport bug) — a failure, not a success. @@ -561,11 +550,8 @@ def _build_result_entry( @dataclass class _ChildRun: """State of one child run, shared by every phase of ``_run_single_child``. - - ``worktree_info`` stays None until isolation engages, so ``attach_worktree`` - is a no-op on every early error path; ``child_task_id`` / - ``parent_reads_snapshot`` are set by ``seed_workspace``. - """ + ``worktree_info`` stays None until isolation engages (``attach_worktree`` is + then a no-op on every early error path); ``seed_workspace`` sets the rest.""" child: Any parent_agent: Any @@ -652,18 +638,16 @@ class _ChildRun: """Run the child's conversation on a daemon worker: ``(result, None, False)`` or ``(None, error_entry, close_deferred)`` on timeout/exception. - The hard timeout is off by default (``result(timeout=None)`` blocks; stuck - children are the heartbeat's job). Daemon worker: a timed-out child is - abandoned and a non-daemon thread would block interpreter exit at atexit - join. The worker installs a non-interactive approval callback so dangerous - command prompts never fall back to ``input()`` and deadlock the parent TUI - (deny vs approve follows delegation.subagent_auto_approve). - - On failure: steer acceptance closes BEFORE the stop signal (a concurrent - steer is drained into the entry or rejected, never silently lost); a - 0-API-call timeout gets a diagnostic dump; a timed-out worker that still - owns the child gets ``child.close()`` via a Future done-callback - (``close_deferred=True``) because closing from this thread races its + Hard timeout is off by default (``result(timeout=None)``; stuck children are + the heartbeat's job). Daemon worker: an abandoned timed-out child on a + non-daemon thread would block interpreter exit at atexit join. The worker + installs a non-interactive approval callback (deny/approve per + delegation.subagent_auto_approve) so dangerous-command prompts never fall + back to ``input()`` and deadlock the parent TUI. On failure: steer acceptance + closes BEFORE the stop signal (a concurrent steer is drained into the entry + or rejected, never lost); a 0-API-call timeout gets a diagnostic dump; a + worker that still owns the child gets ``child.close()`` via a Future + done-callback (``close_deferred=True``) — closing here would race its still-unwinding finally path. """ from tools.delegate_tool import (_get_child_timeout, _get_subagent_approval_callback, _set_subagent_approval_cb) @@ -811,13 +795,11 @@ class _ChildRun: _safe_progress(self.child_progress_cb, "subagent.complete", **complete_kwargs) def cleanup(self, *, heartbeat: tuple, child_pool: Any, leased_cred_id: Any, close_deferred: bool) -> None: - """Finally-path teardown (idempotent, never raises). - - Order matters: stop heartbeat → drop registry entry → release credential - lease → restore the parent's process-global tool names → detach from the - parent's interrupt list → close the child (unless a timed-out worker still - owns it) → pop the child's Relay scope if no turn is active. - """ + """Finally-path teardown (idempotent, never raises). Order matters: stop + heartbeat → drop registry entry → release credential lease → restore the + parent's process-global tool names → detach from the parent's interrupt list + → close the child (unless a timed-out worker still owns it) → pop the child's + Relay scope if no turn is active.""" child = self.child _heartbeat_stop, _heartbeat_thread = heartbeat _heartbeat_stop.set() diff --git a/tools/delegate_tool_config.py b/tools/delegate_tool_config.py index da144e78a0..fad6dd5f8d 100644 --- a/tools/delegate_tool_config.py +++ b/tools/delegate_tool_config.py @@ -105,22 +105,16 @@ def _get_max_concurrent_children() -> int: return result def _get_worktree_isolation() -> bool: - """delegation.worktree_isolation (bool, default False). - - When enabled each child gets its own git worktree off the parent's HEAD so - parallel children never contend for one working copy. Git-only and - local-backend-only; otherwise silently ignored (shared workspace as before). - """ + """delegation.worktree_isolation (bool, default False): each child gets its own + git worktree off the parent's HEAD so parallel children never contend for one + working copy. Git-only and local-backend-only; otherwise silently ignored.""" return bool(_cfg().get("worktree_isolation", False)) def _get_max_async_children() -> int: """Concurrency cap for background delegations == delegation.max_concurrent_children. - At capacity a new async dispatch is REJECTED (not queued) so a runaway model - can't pile up unbounded background work; the caller then runs synchronously. - A leftover ``delegation.max_async_children`` key is ignored with a one-time - deprecation warning. - """ + can't pile up unbounded background work; the caller then runs synchronously. A + leftover ``delegation.max_async_children`` key is ignored with a one-time warning.""" from tools.delegate_tool import _get_max_concurrent_children if _cfg().get("max_async_children") is not None: _warn_once( @@ -136,24 +130,20 @@ def _parse_timeout(raw: Any) -> Optional[float]: return None if parsed <= 0 else max(30.0, parsed) def _get_child_timeout() -> Optional[float]: - """Hard wall-clock cap for one child, or None (default: no timeout). - - Failures should come from what the child does (API/tool errors, iteration - budget), not a stopwatch; stuck children are caught by the heartbeat - staleness monitor. delegation.child_timeout_seconds > 0 opts in (floor 30 s); - 0 or negative disables. Env fallback: DELEGATION_CHILD_TIMEOUT_SECONDS. - """ + """Hard wall-clock cap for one child, or None (default: no timeout). Failures + should come from what the child does (API/tool errors, iteration budget), not a + stopwatch; stuck children are caught by the heartbeat staleness monitor. + delegation.child_timeout_seconds > 0 opts in (floor 30 s); 0 or negative + disables. Env fallback: DELEGATION_CHILD_TIMEOUT_SECONDS.""" return _knob( "child_timeout_seconds", "DELEGATION_CHILD_TIMEOUT_SECONDS", _parse_timeout, DEFAULT_CHILD_TIMEOUT, "delegation.child_timeout_seconds=%r is not a valid number; using default (no timeout)", ) def _get_max_spawn_depth() -> int: - """delegation.max_spawn_depth floored at 1 (no ceiling). - - Depth 0 is the parent; agents at depths 0..N-1 may spawn, depth N is the - leaf floor. Default 1 is flat. Each extra level multiplies API cost. - """ + """delegation.max_spawn_depth floored at 1 (no ceiling). Depth 0 is the parent; + agents at depths 0..N-1 may spawn, depth N is the leaf floor. Default 1 is + flat. Each extra level multiplies API cost.""" def _floored(v): ival = int(v) if ival < _MIN_SPAWN_DEPTH: @@ -183,12 +173,10 @@ def _normalized_runtime_url(value: Any) -> str: return str(value or "").strip().rstrip("/") def _inherit_parent_capabilities(parent_agent, override_provider, override_base_url) -> Optional[dict]: - """Parent's endpoint-trust capability map for a child, or None. - - ``agent.capabilities`` is a trust decision scoped to one provider+endpoint: - inherited ONLY when the child runs the parent's exact route; any provider or - base_url override stays DEFAULT-DENY (matches the /model switch posture). - """ + """Parent's endpoint-trust capability map for a child, or None. ``agent.capabilities`` + is a trust decision scoped to one provider+endpoint: inherited ONLY when the + child runs the parent's exact route; any provider or base_url override stays + DEFAULT-DENY (matches the /model switch posture).""" if override_provider or override_base_url: return None parent_caps = getattr(parent_agent, "capabilities", None) @@ -197,11 +185,9 @@ def _inherit_parent_capabilities(parent_agent, override_provider, override_base_ return {key: value for key, value in parent_caps.items() if isinstance(key, str) and isinstance(value, bool)} def _inherit_parent_base_url(parent_agent, fallback_base_url: Optional[str]) -> Optional[str]: - """Base URL the parent is actually calling (live client), not a stale attribute. - - ``parent_agent.base_url`` can lag the live client (old OpenRouter URL vs - local Ollama); inheriting the stale one 401s with a dummy/local key. - """ + """Base URL the parent is actually calling (live client), not a stale attribute: + ``parent_agent.base_url`` can lag the live client (old OpenRouter URL vs local + Ollama) and inheriting the stale one 401s with a dummy/local key.""" surface_url = _normalized_runtime_url(fallback_base_url) client_kwargs = getattr(parent_agent, "_client_kwargs", None) client = getattr(parent_agent, "client", None) @@ -225,15 +211,13 @@ def _loaded_pool(key: Any): def _resolve_child_credential_pool( effective_provider: Optional[str], parent_agent, effective_base_url: Optional[str] = None, ): - """Credential pool for the child: parent's pool (same provider), that - provider's own pool, or None (child keeps its fixed credential). - - Custom endpoints all collapse to ``provider="custom"``, so they are matched - by endpoint identity (the ``custom:`` pool key) — sharing the parent's - pool across different custom endpoints would overwrite the child's delegated - base_url on lease. An unregistered custom endpoint (no custom_providers - entry) keeps the child's fixed credential rather than inherit the parent's. - """ + """Credential pool for the child: parent's pool (same provider), that provider's + own pool, or None (child keeps its fixed credential). Custom endpoints all + collapse to ``provider="custom"``, so they are matched by endpoint identity (the + ``custom:`` pool key) — sharing the parent's pool across different custom + endpoints would overwrite the child's delegated base_url on lease; an + unregistered custom endpoint (no custom_providers entry) keeps the child's fixed + credential rather than inherit the parent's.""" parent_pool = getattr(parent_agent, "_credential_pool", None) if not effective_provider: return parent_pool @@ -260,13 +244,10 @@ def _resolve_child_credential_pool( def _merge_request_overrides(runtime_overrides, explicit_overrides): """Merge explicit ``delegation.request_overrides`` OVER runtime-derived ones. - - Explicit top-level keys win; ``extra_body`` is deep-merged ONE level so - provider personality (e.g. ``thinking: {type: disabled}``) survives unless - the explicit dict redefines that exact key. Both sides are deep-copied so - transport-side mutation can't leak into the config/runtime cache. Returns - None when both sides are empty. - """ + Explicit top-level keys win; ``extra_body`` is deep-merged ONE level so provider + personality (e.g. ``thinking: {type: disabled}``) survives unless the explicit + dict redefines that exact key. Both sides are deep-copied so transport-side + mutation can't leak into the config/runtime cache. None when both are empty.""" import copy as _copy runtime_overrides = runtime_overrides if isinstance(runtime_overrides, dict) else None explicit_overrides = explicit_overrides if isinstance(explicit_overrides, dict) else None @@ -380,15 +361,12 @@ def _runtime_provider_credentials(v: dict, explicit_request_overrides) -> dict: ) def _resolve_delegation_credentials(cfg: dict, parent_agent) -> dict: - """Resolve the child credential bundle from the ``delegation`` config section. - - Three branches: ``base_url`` set → direct endpoint (``api_key`` None means - inherit the parent's key, so providers keyed outside OPENAI_API_KEY work); - ``provider`` set → full bundle via the runtime provider system (same path as - CLI/gateway startup); neither → None values, child inherits everything. - ``request_overrides`` is honored on every branch. Raises ValueError with a - user-facing message on credential failure. - """ + """Child credential bundle from the ``delegation`` config section. Three + branches: ``base_url`` set → direct endpoint (``api_key`` None means inherit the + parent's key, so providers keyed outside OPENAI_API_KEY work); ``provider`` set + → full bundle via the runtime provider system (same path as CLI/gateway + startup); neither → None values, child inherits everything. ``request_overrides`` + is honored on every branch. Raises ValueError with a user-facing message.""" values = {k: str(cfg.get(k) or "").strip() or None for k in ("model", "provider", "base_url", "api_key")} values["api_mode"] = str(cfg.get("api_mode") or "").strip().lower() or None explicit_request_overrides = cfg.get("request_overrides") if isinstance(cfg.get("request_overrides"), dict) else None @@ -405,14 +383,12 @@ def _resolve_delegation_credentials(cfg: dict, parent_agent) -> dict: return _runtime_provider_credentials(values, explicit_request_overrides) def _load_config() -> dict: - """Return the ``delegation`` config section (read-only — do NOT mutate). - - Prefers the shared ``load_config_readonly()`` (follows HERMES_HOME/profile; - no deepcopy, since this runs on every get_definitions() rebuild) over the - legacy ``cli.CLI_CONFIG``, which can hide user-set keys. Exception: + """The ``delegation`` config section (read-only — do NOT mutate). Prefers the + shared ``load_config_readonly()`` (follows HERMES_HOME/profile; no deepcopy, + since this runs on every get_definitions() rebuild) over the legacy + ``cli.CLI_CONFIG``, which can hide user-set keys — except that ``HERMES_IGNORE_USER_CONFIG=1`` is only honored by the legacy loader, so it - stays authoritative when that flag is set. - """ + stays authoritative when that flag is set.""" if os.environ.get("HERMES_IGNORE_USER_CONFIG") != "1": try: from hermes_cli.config import load_config_readonly @@ -444,15 +420,12 @@ def _resolve_child_runtime( override_base_url: Optional[str], override_api_key: Optional[str], override_api_mode: Optional[str], override_max_tokens: Optional[int], override_acp_command: Optional[str], override_acp_args: Optional[List[str]], ) -> Dict[str, Any]: - """Resolve the child's credentials, transport and routing (config override > - parent inherit) as ``AIAgent`` keyword arguments. - - Rules that are easy to break: api_mode is re-derived (not inherited) when - the child's provider differs from the parent's or is Nous Portal (dual-wire); - a pinned ``delegation.command`` must exist on PATH or the spawn fails loudly; - ``override_provider`` clears the parent's ACP transport, fallback chain and - OpenRouter routing filters so the pinned provider is actually honoured. - """ + """Child credentials, transport and routing (config override > parent inherit) + as ``AIAgent`` kwargs. Rules that are easy to break: api_mode is re-derived (not + inherited) when the child's provider differs from the parent's or is Nous Portal + (dual-wire); a pinned ``delegation.command`` must exist on PATH or the spawn + fails loudly; ``override_provider`` clears the parent's ACP transport, fallback + chain and OpenRouter routing filters so the pinned provider is actually honoured.""" effective_model = model or parent_agent.model effective_provider = override_provider or getattr(parent_agent, "provider", None) effective_base_url = override_base_url or _inherit_parent_base_url(parent_agent, parent_agent.base_url) diff --git a/tools/delegate_tool_dispatch.py b/tools/delegate_tool_dispatch.py index bfe72e0d1a..690e4634e9 100644 --- a/tools/delegate_tool_dispatch.py +++ b/tools/delegate_tool_dispatch.py @@ -94,14 +94,12 @@ def _report_child_done(parent_agent, spinner_ref, entry, tag, task_labels, n_tas spinner_ref.update_text(f"🔀 {'[' + tag + '] ' if tag else ''}{remaining} task{'s' if remaining != 1 else ''} remaining") def _run_children_parallel(batch: _Batch, results: list, *, honor_parent_interrupt: bool) -> None: - """Run the batch's children in parallel, appending entries to ``results``. - - Polls futures with a short ``wait()`` timeout instead of ``as_completed()`` - so a wedged child cannot block the parent forever after an interrupt; - on parent interrupt the still-pending children are reported as - ``interrupted`` and abandoned (they already got the interrupt signal). - Prints one completion line per child. ``results`` ends sorted by task_index. - """ + """Run the batch's children in parallel, appending entries to ``results`` + (sorted by task_index on return, one completion line printed per child). + Polls futures with a short ``wait()`` timeout instead of ``as_completed()`` so + a wedged child cannot block the parent forever after an interrupt; on parent + interrupt the still-pending children are reported ``interrupted`` and + abandoned (they already got the interrupt signal).""" # Daemon workers (tools.daemon_pool): the `with` block still joins normally, # but if the parent is interrupted while a child is wedged, the abandoned # worker must not block interpreter exit. @@ -138,13 +136,11 @@ def _run_children_parallel(batch: _Batch, results: list, *, honor_parent_interru results.sort(key=lambda r: r["task_index"]) # match input order def _execute_and_aggregate(batch: _Batch, *, honor_parent_interrupt: bool = True) -> dict: - """Run all built children, join, finalize (hooks + cost rollup), return the combined dict. - - Shared by the sync path and the background runner: even in the background - the batch JOINS on itself here so ONE consolidated results block re-enters - the conversation. Live transcripts are finalized but retained as the - full-fidelity record (retention pruning happens on future dispatches). - """ + """Run all built children, join, finalize (hooks + cost rollup), return the + combined dict. Shared by the sync path and the background runner: even in the + background the batch JOINS on itself here so ONE consolidated results block + re-enters the conversation. Live transcripts are finalized but retained as the + full-fidelity record (retention pruning happens on future dispatches).""" from tools.delegation_live_log import update_manifest_statuses results: list = [] if len(batch.task_list) == 1: @@ -192,13 +188,12 @@ def _run_sync_with_note(batch: _Batch, reason: str) -> str: def _resolve_async_wake_sid(origin_wake_sid: str) -> Optional[str]: """Wake target for a detached batch, or None to force synchronous execution. - Finite sessions (stateless HTTP requests, one-shot Kanban workers) cannot - route a detached result back after their turn/process ends. But if a raw - session id is bound (the API server always binds one), gateway.wake can - still reach the session by self-POSTing /v1/chat/completions with that id, - so only fall back to sync when there is truly no session id to wake. Uses - the origin captured BEFORE child construction — HERMES_SESSION_ID here - would be the subagent's internal id. + Finite sessions (stateless HTTP requests, one-shot Kanban workers) cannot route + a detached result back after their turn/process ends — but if a raw session id + is bound (the API server always binds one), gateway.wake can still reach it by + self-POSTing /v1/chat/completions, so only fall back to sync when there is truly + no session id to wake. Uses the origin captured BEFORE child construction — + HERMES_SESSION_ID here would be the subagent's internal id. """ try: from gateway.session_context import async_delivery_supported @@ -222,9 +217,9 @@ def _resolve_async_session_key(parent_agent: Any, origin_ui_session_id: str) -> can rotate it mid-turn before the TUI-side dict is re-anchored, and a stale approval-context key would orphan the completion. Gateway chats keep the platform conversation key (agent:main:...). The CLI has no bound approval - contextvar and no HERMES_SESSION_KEY, so the key resolves empty; its drain - is a positive-ownership filter keyed on the durable session_id, so an empty - key would fail closed — stamp the parent's durable id. + contextvar and no HERMES_SESSION_KEY, so the key resolves empty; its drain is a + positive-ownership filter on the durable session_id (empty would fail closed), + so stamp the parent's durable id. """ from tools.approval import get_current_session_key session_key = get_current_session_key(default="") @@ -298,13 +293,10 @@ def _dispatched_payload(dispatch: dict, goals: List[str], child_agents: List[Any def _dispatch_background(batch: _Batch) -> str: """Dispatch the WHOLE batch as one async unit and return the tool result JSON. - - The runner joins on every child and yields ONE consolidated results block - that re-enters the conversation as a single message when ALL children - finish. Falls back to running it synchronously (with an explanatory - ``note``) when the session cannot receive detached completions or the async - pool is at capacity. - """ + The runner joins on every child and yields ONE consolidated results block that + re-enters the conversation as a single message when ALL children finish. Falls + back to running synchronously (with an explanatory ``note``) when the session + cannot receive detached completions or the async pool is at capacity.""" from tools.delegate_tool import _get_max_async_children from tools.async_delegation import dispatch_async_delegation_batch wake_sid = _resolve_async_wake_sid(batch.origin_wake_sid) diff --git a/tools/delegate_tool_progress.py b/tools/delegate_tool_progress.py index 5c7c56c147..9158d97409 100644 --- a/tools/delegate_tool_progress.py +++ b/tools/delegate_tool_progress.py @@ -72,14 +72,11 @@ def format_subagent_failure_line( class DelegateEvent(str, enum.Enum): - """Formal event types emitted during delegation progress. - - The relay normalises incoming legacy strings (``tool.started``, - ``_thinking``, …) to these values via ``_LEGACY_EVENT_MAP``; external - consumers (gateway SSE, ACP adapter, CLI) still receive the legacy strings - during the deprecation window. TASK_SPAWNED / TASK_COMPLETED / TASK_FAILED - are reserved for future orchestrator lifecycle events, not emitted yet. - """ + """Formal delegation progress event types. The relay normalises incoming legacy + strings (``tool.started``, ``_thinking``, …) to these via ``_LEGACY_EVENT_MAP``; + external consumers (gateway SSE, ACP adapter, CLI) still receive the legacy + strings during the deprecation window. TASK_SPAWNED / TASK_COMPLETED / + TASK_FAILED are reserved for future orchestrator lifecycle events, not emitted yet.""" TASK_SPAWNED = "delegate.task_spawned" TASK_PROGRESS = "delegate.task_progress" @@ -126,12 +123,10 @@ def _build_child_system_prompt( goal: str, context: Optional[str] = None, *, workspace_path: Optional[str] = None, role: str = "leaf", max_spawn_depth: int = 2, child_depth: int = 1, ) -> str: - """Build a focused system prompt for a child agent. - - role='orchestrator' appends a delegation-capability block (modeled on - OpenClaw's buildSubagentSystemPrompt); its depth note is literal truth - grounded in the passed config so the LLM can't confabulate nesting. - """ + """Focused system prompt for a child agent. role='orchestrator' appends a + delegation-capability block (modeled on OpenClaw's buildSubagentSystemPrompt); + its depth note is literal truth grounded in the passed config so the LLM can't + confabulate nesting.""" parts = ["You are a focused subagent working on a specific delegated task.", "", f"YOUR TASK:\n{goal}"] if context and context.strip(): parts.append(f"\nCONTEXT:\n{context}") @@ -219,15 +214,12 @@ _BATCH_ORDINALS: Dict[str, int] = {} _BATCH_ORDINALS_LOCK = threading.Lock() def format_batch_tag(delegation_id: Optional[str]) -> str: - """Short human tag for a delegation batch: ``deleg_6a664903`` → ``set 1`` - (first batch seen in this process), the next distinct id → ``set 2``. - - Several batches (a parent's fan-out plus a child's nested fan-out, or two - concurrent tools) print interleaved ``[n/N]`` lines to one console; without - a tag ``✓ [3/3]`` and ``✓ [3/9]`` are indistinguishable, and a raw hex - slice is unreadable. Empty string when no id is known so callers can - concatenate unconditionally. - """ + """Short human tag for a delegation batch: ``deleg_6a664903`` → ``set 1`` (first + batch seen in this process), the next distinct id → ``set 2``. Several batches + (a parent's fan-out plus a child's nested fan-out, or two concurrent tools) + print interleaved ``[n/N]`` lines to one console; without a tag ``✓ [3/3]`` and + ``✓ [3/9]`` are indistinguishable, and a raw hex slice is unreadable. Empty + string when no id is known so callers can concatenate unconditionally.""" if not isinstance(delegation_id, str) or not delegation_id: return "" with _BATCH_ORDINALS_LOCK: @@ -268,14 +260,12 @@ def _short(text: str, n: int) -> str: class _ChildProgressRelay: - """Callable relaying one child's events to the parent display. - - CLI: prints tree-view lines above the parent's delegation spinner. - Gateway: batches tool names (``_BATCH_SIZE``) and relays to the parent's - progress callback, threading the identity kwargs (subagent_id, parent_id, - depth, model, toolsets) into every event so the TUI can rebuild the live - spawn tree and route per-branch controls back by ``subagent_id``. - """ + """Callable relaying one child's events to the parent display. CLI: prints + tree-view lines above the parent's delegation spinner. Gateway: batches tool + names (``_BATCH_SIZE``) and relays to the parent's progress callback, threading + the identity kwargs (subagent_id, parent_id, depth, model, toolsets) into every + event so the TUI can rebuild the live spawn tree and route per-branch controls + back by ``subagent_id``.""" _BATCH_SIZE = 5 diff --git a/tools/delegate_tool_registry.py b/tools/delegate_tool_registry.py index 11cebb14b9..3257df2db5 100644 --- a/tools/delegate_tool_registry.py +++ b/tools/delegate_tool_registry.py @@ -78,12 +78,10 @@ def _unregister_subagent(subagent_id: str, *, agent: Any = None) -> None: def _close_subagent_steering(subagent_id: str, agent: Any) -> Optional[str]: """Atomically close steer acceptance and drain its final durable artifact. - - ``steer_subagent`` holds the same registry lock through ``agent.steer``. - Therefore either acceptance wins and this drain sees its exact text, or - closure wins and the caller is rejected. Exact agent identity prevents a - finishing child with a recycled public id from closing its replacement. - """ + ``steer_subagent`` holds the same registry lock through ``agent.steer``, so + either acceptance wins and this drain sees its exact text, or closure wins and + the caller is rejected. Exact agent identity prevents a finishing child with a + recycled public id from closing its replacement.""" with _active_subagents_lock: record = _active_subagents.get(subagent_id) if record is None or record.get("agent") is not agent: @@ -120,14 +118,14 @@ def steer_subagent( ) -> bool: """Queue steering text into a running subagent without stopping it. - Calls AIAgent.steer(), which appends the text to the child's last tool result - at its next iteration boundary — the current tool call is never cut. True iff - the text was QUEUED while the child still accepted work; False for - unknown/closed id, ownership mismatch, no live agent, or empty text. - ``owner_session_id=None`` keeps the in-process helper contract; gateway - callers must pass exact authority. Acceptance and completion are linearized - by the registry lock: if acceptance wins but no delivery boundary remains, - the text lands in the entry as ``missed_steer``. + AIAgent.steer() appends the text to the child's last tool result at its next + iteration boundary — the current tool call is never cut. True iff the text was + QUEUED while the child still accepted work; False for unknown/closed id, + ownership mismatch, no live agent, or empty text. ``owner_session_id=None`` + keeps the in-process helper contract; gateway callers must pass exact + authority. Acceptance and completion are linearized by the registry lock: if + acceptance wins but no delivery boundary remains, the text lands in the entry + as ``missed_steer``. """ if not text or not text.strip(): return False @@ -211,16 +209,15 @@ def _resolve_session_lineage(session_id: Optional[str], parent_agent: Any) -> st def _owns_subagent_record(record: Dict[str, Any], parent_agent: Any) -> bool: """True when *parent_agent*'s conversation owns this live-child record. - Tier 1: object identity — the ``_delegate_parent_ref`` weakref chain reaches - *parent_agent* (fast path while the parent AIAgent survives the run). - Tier 2: durable lineage — the record's ``owner_agent_session_id`` matches the - caller's ``session_id``, resolving compression-rotation lineage on both - sides. Tier 2 exists because the identity chain is BRITTLE across parent - rebuilds: the CLI sets ``self.agent = None`` mid-session (route change, - credential refresh, /model, MoA one-shots) and builds a NEW AIAgent while - the child keeps a weakref to the old one. Delivery always routed by durable - session id; control must use the same spine or running children go - invisible/unsteerable. + Tier 1: identity — the ``_delegate_parent_ref`` weakref chain reaches + *parent_agent* (fast path while the parent AIAgent survives the run). Tier 2: + durable lineage — the record's ``owner_agent_session_id`` matches the caller's + ``session_id`` after resolving compression-rotation lineage on both sides. + Tier 2 exists because the identity chain is BRITTLE across parent rebuilds: + the CLI sets ``self.agent = None`` mid-session (route change, credential + refresh, /model, MoA one-shots) and builds a NEW AIAgent while the child keeps + a weakref to the old one. Delivery routes by durable session id; control must + use the same spine or running children go invisible/unsteerable. """ if _is_descendant_of(record.get("agent"), parent_agent): return True @@ -261,12 +258,9 @@ def _list_payload(parent_agent: Any) -> Dict[str, Any]: return payload def _handle_control_action(action: str, subagent_id: Optional[str], message: Optional[str], parent_agent: Any) -> str: - """Synchronous control plane for delegate_task: list/steer/stop. - - Runs in-turn (never backgrounded) and only over subagents descended from - *parent_agent* — the same registry the TUI overlay drives, but scoped so - a conversation can only control its own spawn tree. - """ + """Synchronous control plane for delegate_task: list/steer/stop. Runs in-turn + (never backgrounded) over the same registry the TUI overlay drives, scoped so a + conversation can only control its own spawn tree.""" if action == "list": return json.dumps(_list_payload(parent_agent), ensure_ascii=False) diff --git a/tools/delegate_tool_tasks.py b/tools/delegate_tool_tasks.py index 3e658d6a2d..3e92fffa0d 100644 --- a/tools/delegate_tool_tasks.py +++ b/tools/delegate_tool_tasks.py @@ -38,16 +38,13 @@ def _recover_tasks_from_json_string(tasks: Any) -> tuple[Optional[List[Dict[str, return parsed, None def _validate_batch_tasks(task_list: List[Dict[str, Any]]) -> Optional[str]: - """Batch-only quality gate beyond per-task goal presence; actionable error or None. - - No minimum count: a one-entry array is the canonical single-task shape - (legacy top-level `goal` is wrapped into one). Duplicate goals are - deliberately NOT rejected — identical-goal fan-outs (best-of-N / ensemble - sampling) are legitimate and blocking them broke real workflows. The - too-short check applies only to multi-task fan-outs (terse goals there are - usually unexpanded templates); a SINGLE task legitimately uses short goals - ("Fix the tests"). - """ + """Batch-only quality gate beyond per-task goal presence; actionable error or + None. No minimum count: a one-entry array is the canonical single-task shape + (legacy top-level `goal` is wrapped into one). Duplicate goals are deliberately + NOT rejected — identical-goal fan-outs (best-of-N / ensemble sampling) are + legitimate and blocking them broke real workflows. The too-short check applies + only to multi-task fan-outs (terse goals there are usually unexpanded + templates); a SINGLE task legitimately uses short goals ("Fix the tests").""" for i, task in enumerate(task_list): goal = str(task.get("goal", "")).strip() if _PLACEHOLDER_GOAL_RE.match(" ".join(goal.lower().split())): diff --git a/tools/delegate_tool_toolsets.py b/tools/delegate_tool_toolsets.py index f49ad87d09..581de9ad9c 100644 --- a/tools/delegate_tool_toolsets.py +++ b/tools/delegate_tool_toolsets.py @@ -40,11 +40,9 @@ def _is_mcp_toolset_name(name: str) -> bool: return bool(target and str(target).startswith("mcp-")) def _expand_parent_toolsets(parent_toolsets: set) -> set: - """Add every toolset whose tools are a subset of the parent's tools. - - A parent on a composite like ``hermes-cli`` must still let a child request - ``web``/``terminal``; bare name intersection would reject them. - """ + """Add every toolset whose tools are a subset of the parent's tools: a parent on + a composite like ``hermes-cli`` must still let a child request ``web``/``terminal``; + bare name intersection would reject them.""" parent_tool_names = {t for ts_name in parent_toolsets for t in (TOOLSETS.get(ts_name) or {}).get("tools", [])} expanded = set(parent_toolsets) if parent_tool_names: @@ -76,16 +74,14 @@ def _blocked_toolsets_for_role(role: str) -> List[str]: def _resolve_child_toolsets( parent_agent, toolsets: Optional[List[str]], effective_role: str ) -> tuple[List[str], List[str]]: - """Return ``(enabled_toolsets, disabled_toolsets)`` for a child. - - Children never gain tools the parent lacks: explicit ``toolsets`` are - intersected with the parent's (composite-expanded) set, else the parent's - enabled set is inherited. Blocked tools are stripped twice — whole blocked - toolsets here, and exact one-tool deny toolsets via ``disabled_toolsets`` so - blocked names inside mixed bundles (hermes-cli) are subtracted AFTER - composite expansion and survive registry refreshes. Orchestrators get - ``delegation`` re-added unconditionally (role-granted, not inherited). - """ + """``(enabled_toolsets, disabled_toolsets)`` for a child. Children never gain + tools the parent lacks: explicit ``toolsets`` are intersected with the parent's + (composite-expanded) set, else the parent's enabled set is inherited. Blocked + tools are stripped twice — whole blocked toolsets here, and exact one-tool deny + toolsets via ``disabled_toolsets`` so blocked names inside mixed bundles + (hermes-cli) are subtracted AFTER composite expansion and survive registry + refreshes. Orchestrators get ``delegation`` re-added unconditionally + (role-granted, not inherited).""" # enabled_toolsets=None means "all tools", so derive from loaded tool names. parent_enabled = getattr(parent_agent, "enabled_toolsets", None) if parent_enabled is not None: