diff --git a/tools/delegate_tool.py b/tools/delegate_tool.py index 915d27d630..fcdc375bd9 100644 --- a/tools/delegate_tool.py +++ b/tools/delegate_tool.py @@ -21,9 +21,8 @@ from utils import is_truthy_value logger = logging.getLogger(__name__) -# The delegate_tool_* siblings hold the pieces split out of this module; every -# name callers or patching tests reach as ``tools.delegate_tool.`` is -# re-imported here. Mutable flag globals live only in their owning module. +# The delegate_tool_* siblings hold the pieces split out of this module; every name callers or patching tests reach as +# ``tools.delegate_tool.`` is re-imported here. Mutable flag globals live only in their owning module. from tools.delegate_tool_child_run import ( # noqa: F401 _ChildRun, _attach_child, _build_result_entry, _dump_subagent_timeout_diagnostic, _fabricated_entry, _lease_child_credential, _merge_late_steer, _register_child, _start_heartbeat, _validate_child_output_schema, @@ -68,11 +67,10 @@ def _normalize_role(r: Optional[str]) -> str: DEFAULT_MAX_ITERATIONS = 250 _HEARTBEAT_INTERVAL = 30 # seconds between parent activity heartbeats during delegation -# Stale-heartbeat thresholds (cycles of _HEARTBEAT_INTERVAL with no progress). -# Progress = iteration, current_tool OR last_activity_ts advancing; an in-flight -# model wait refreshes last_activity_ts, so slow models are not "idle". Idle -# stays tight so a truly wedged child doesn't mask the gateway timeout; in-tool -# is much higher so legitimately long tools can finish. +# Stale-heartbeat thresholds (cycles of _HEARTBEAT_INTERVAL with no progress). Progress = iteration, current_tool OR +# last_activity_ts advancing; an in-flight model wait refreshes last_activity_ts, so slow models are not "idle". Idle +# stays tight so a truly wedged child doesn't mask the gateway timeout; in-tool is much higher so legitimately long +# tools can finish. _HEARTBEAT_STALE_CYCLES_IDLE = 15 # 450s idle between turns → stale _HEARTBEAT_STALE_CYCLES_IN_TOOL = 40 # 1200s stuck on same tool → stale @@ -82,11 +80,10 @@ 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 @@ -118,9 +115,8 @@ 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 @@ -242,9 +238,8 @@ def _run_single_child( """ child_progress_cb = getattr(child, "tool_progress_callback", None) child_pool, leased_cred_id = _lease_child_credential(child) - # Heartbeat keeps the parent's _last_activity_ts moving so the gateway - # inactivity timeout doesn't fire while the child works; it stops itself - # once the child looks stale (see _HEARTBEAT_STALE_CYCLES_*). + # Heartbeat keeps the parent's _last_activity_ts moving so the gateway inactivity timeout doesn't fire while the + # child works; it stops itself once the child looks stale (see _HEARTBEAT_STALE_CYCLES_*). heartbeat = _start_heartbeat(child, parent_agent, task_index) # TUI/RPC registry entry (kill/pause/status by subagent_id); None for test # doubles without a stable id. Unregistered in the finally block. @@ -343,11 +338,10 @@ 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.") @@ -401,9 +395,8 @@ def delegate_task( return tool_error(err) overall_start = time.monotonic() - # Live transcripts: cache/delegation/live//task-.log per task, a - # side channel with zero effect on message content or prompt caching. - # Best-effort: on failure live_paths is empty and delegation proceeds. + # Live transcripts: cache/delegation/live//task-.log per task, a side channel with zero effect on message + # content or prompt caching. Best-effort: on failure live_paths is empty and delegation proceeds. from tools.delegation_live_log import create_live_transcripts live_deleg_id, live_writers, live_paths = create_live_transcripts( task_list, context, model=creds.get("model"), provider=creds.get("provider") @@ -434,9 +427,8 @@ def _build_top_level_description() -> str: except Exception: orchestration_available = False - # Mention recursion only where it's actually available. send_message is - # deliberately not named (gateway-internal vocabulary); model_tools - # session-filters the list to tools the session has. + # Mention recursion only where it's actually available. send_message is deliberately not named (gateway-internal + # vocabulary); model_tools session-filters the list to tools the session has. if orchestration_available: restrictions_rule = ( "- Children cannot call clarify, memory, or cronjob.\n" @@ -505,11 +497,9 @@ def _p(type_: str, description: str, **extra) -> dict: DELEGATE_TASK_SCHEMA = { "name": "delegate_task", - # description / tasks.description are placeholders: the real text is built per - # get_definitions() call by _build_dynamic_schema_overrides() so the model sees - # the user's actual max_concurrent_children / max_spawn_depth. Lazy (not at - # import) so cli.CLI_CONFIG isn't forced to load before the test conftest - # redirects HERMES_HOME. + # description / tasks.description are placeholders: the real text is built per get_definitions() call by + # _build_dynamic_schema_overrides() so the model sees the user's actual max_concurrent_children / max_spawn_depth. + # Lazy (not at import) so cli.CLI_CONFIG isn't forced to load before the test conftest redirects HERMES_HOME. "description": ( "Spawn one or more subagents in isolated contexts. " "Description is rebuilt at every get_definitions() call to reflect the user's current delegation limits." @@ -517,13 +507,10 @@ DELEGATE_TASK_SCHEMA = { "parameters": { "type": "object", "properties": { - # The handler also accepts the legacy single-goal shape (top-level - # `goal`/`context`/`output_schema`), wrapped into a one-entry batch at - # dispatch, and a per-task `role` (legacy, ignored: capability is - # depth-derived). Both unadvertised on purpose (old transcripts only); - # do not re-add. No maxItems — the runtime limit - # (delegation.max_concurrent_children) is enforced with a clear error in - # delegate_task(). + # The handler also accepts the legacy single-goal shape (top-level `goal`/`context`/`output_schema`), + # wrapped into a one-entry batch at dispatch, and a per-task `role` (legacy, ignored: capability is + # depth-derived). Both unadvertised on purpose (old transcripts only); do not re-add. No maxItems — the + # runtime limit (delegation.max_concurrent_children) is enforced with a clear error in delegate_task(). "tasks": { "type": "array", "minItems": 1, @@ -580,13 +567,11 @@ DELEGATE_TASK_SCHEMA = { 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); 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.""" + """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); 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.""" 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 c11610de3d..38df749385 100644 --- a/tools/delegate_tool_child_run.py +++ b/tools/delegate_tool_child_run.py @@ -123,9 +123,8 @@ def _diag_sizes(child: Any) -> List[str]: return ["## Prompt / schema sizes"] + _diag_section("system_prompt: List[str]: - """Worker stack plus all other live threads (bounded to 40): the worker is - often parked on a helper thread, so a pre-HTTP wedge is indistinguishable - from a slow provider without the full picture.""" + """Worker stack plus all other live threads (bounded to 40): the worker is often parked on a helper thread, so a + pre-HTTP wedge is indistinguishable from a slow provider without the full picture.""" import sys as _sys lines = ["## Worker thread stack at timeout"] frames = _sys._current_frames() @@ -155,10 +154,9 @@ 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]: - """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.""" + """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 @@ -263,9 +261,8 @@ 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 (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 @@ -284,10 +281,9 @@ def _register_child( "delegation_id": _str_or_none(getattr(child, "_delegation_id", None)), "model": _str_or_none(getattr(child, "model", None)), "started_at": time.time(), "status": "running", "tool_count": 0, "agent": child, - # Owning conversation's durable session id (same lineage completion - # delivery routes by), sourced from the child's stamp so it survives - # a parent_agent rebuild between dispatch and run; used for - # list/steer/stop ownership when the weakref chain breaks. + # Owning conversation's durable session id (same lineage completion delivery routes by), sourced from the + # child's stamp so it survives a parent_agent rebuild between dispatch and run; used for list/steer/stop + # ownership when the weakref chain breaks. "owner_agent_session_id": ( str(getattr(child, "_parent_session_id", "") or "") or str(getattr(parent_agent, "session_id", "") or "") or None ), @@ -323,14 +319,12 @@ 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) @@ -360,10 +354,9 @@ def _lease_child_credential(child: Any) -> tuple[Any, Optional[str]]: return child_pool, leased_cred_id def _merge_late_steer(result: Dict[str, Any], subagent_id: Optional[str], child: Any) -> None: - """Linearization boundary for registry steering: from here the child cannot - consume another steer. Closing under the registry lock either rejects a - concurrent caller or drains every accepted exact text into the result - before callbacks/result assembly run.""" + """Linearization boundary for registry steering: from here the child cannot consume another steer. Closing under + the registry lock either rejects a concurrent caller or drains every accepted exact text into the result before + callbacks/result assembly run.""" late = _close_subagent_steering(subagent_id, child) if subagent_id else None if late: existing = result.get("pending_steer") @@ -380,9 +373,8 @@ 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) @@ -451,9 +443,8 @@ def _build_result_entry( child: Any, result: Dict[str, Any], task_index: int, duration: float, schema: _SchemaOutcome, ) -> Dict[str, Any]: """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).""" + ``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. @@ -461,16 +452,14 @@ def _build_result_entry( if result.get("interrupted", False): status, exit_reason = "interrupted", "interrupted" elif result.get("failed") or result.get("error"): - # The loop returns the error text as final_response, which would - # otherwise read as "completed". Never report a provider rejection as - # "max_iterations" — that is only truthful for real budget exhaustion. + # The loop returns the error text as final_response, which would otherwise read as "completed". Never report a + # provider rejection as "max_iterations" — that is only truthful for real budget exhaustion. status, exit_reason = "failed", "error" else: - # exit_reason ("completed" vs "max_iterations") tells the parent HOW - # the task ended; completed=False with no failure = budget exhaustion. - # A declared schema still violated after the bounded retry makes the - # summary unusable under the contract, so status must not say completed - # (orchestrators reading only status/icon would accept an empty verdict). + # exit_reason ("completed" vs "max_iterations") tells the parent HOW the task ended; completed=False with no + # failure = budget exhaustion. A declared schema still violated after the bounded retry makes the summary + # unusable under the contract, so status must not say completed (orchestrators reading only status/icon would + # accept an empty verdict). exit_reason = "completed" if result.get("completed", False) else "max_iterations" status = "completed" if schema.valid is not False and usable_summary else "failed" @@ -493,10 +482,9 @@ def _build_result_entry( "output": _num(getattr(child, "session_completion_tokens", 0)), }, "tool_trace": _build_tool_trace(result.get("messages") or []), - # Captured before the finally block calls child.close() so the parent - # thread can fire subagent_stop with the correct role; stripped before - # the dict is serialised back to the model (as is _child_cost_usd, - # folded into the parent's session cost by the aggregator). + # Captured before the finally block calls child.close() so the parent thread can fire subagent_stop with the + # correct role; stripped before the dict is serialised back to the model (as is _child_cost_usd, folded into + # the parent's session cost by the aggregator). "_child_role": getattr(child, "_delegate_role", None), "_child_cost_usd": float(_cost or 0.0) if isinstance(_cost, (int, float)) else 0.0, } @@ -505,8 +493,7 @@ def _build_result_entry( entry["cost_status"] = _cost_status if isinstance(_cost_status, str) and _cost_status else "unknown" if status == "failed": if schema.valid is False and usable_summary: - # The child DID respond; name the contract violation instead of - # the generic "no response" error. + # The child DID respond; name the contract violation instead of the generic "no response" error. entry["error"] = ( "Final answer does not satisfy the declared output_schema" + (" (after 1 retry)." if schema.retries else ".") ) @@ -560,9 +547,8 @@ class _ChildRun: return round(time.monotonic() - self.child_start, 2) def relay_text(self, delta: str) -> None: - """Stream callback forwarding the child's reply text up the progress relay so - gateway watch windows mirror it live (subagent.text → message.delta). Inert - under CLI/TUI: their progress handlers ignore non-tool events.""" + """Stream callback forwarding the child's reply text up the progress relay so gateway watch windows mirror it + live (subagent.text → message.delta). Inert under CLI/TUI: their progress handlers ignore non-tool events.""" if delta: _safe_progress(self.child_progress_cb, "subagent.text", preview=delta) @@ -610,9 +596,8 @@ class _ChildRun: def finish_failed( self, entry: Dict[str, Any], late_steer: Optional[str], *, preview: str, summary: str = "", status: Optional[str] = None, ) -> Dict[str, Any]: - """Shared tail of every failure path: emit ``subagent.complete`` (``status`` - defaults to the entry's), note the steer text that won the race with the - failure, report the worktree.""" + """Shared tail of every failure path: emit ``subagent.complete`` (``status`` defaults to the entry's), note + the steer text that won the race with the failure, report the worktree.""" _safe_progress( self.child_progress_cb, "subagent.complete", preview=preview, status=status or entry["status"], duration_seconds=entry["duration_seconds"], summary=summary, @@ -625,20 +610,17 @@ class _ChildRun: return _close_subagent_steering(self.subagent_id, self.child) if self.subagent_id else None def await_child(self) -> tuple[Optional[Dict[str, Any]], Optional[Dict[str, Any]], bool]: - """Run the child's conversation on a daemon worker: ``(result, None, False)`` - or ``(None, error_entry, close_deferred)`` on timeout/exception. + """Run the child's conversation on a daemon worker: ``(result, None, False)`` or ``(None, error_entry, + close_deferred)`` on timeout/exception. - 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. + 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) from tools.daemon_pool import DaemonThreadPoolExecutor @@ -726,9 +708,8 @@ class _ChildRun: return None, _error_entry, close_deferred def append_sibling_write_reminder(self, entry: Dict[str, Any]) -> None: - """Warn the parent when this child wrote files the parent had already read. - Checks writes by ANY non-parent task_id (not just this child's) so nested - orchestrator→worker chains are covered too.""" + """Warn the parent when this child wrote files the parent had already read. Checks writes by ANY non-parent + task_id (not just this child's) so nested orchestrator→worker chains are covered too.""" if not (self.parent_task_id and self.parent_reads_snapshot): return with _quiet("file_state sibling-write check failed", exc_info=True): @@ -748,9 +729,8 @@ class _ChildRun: entry["stale_paths"] = mod_paths def emit_complete(self, result: Dict[str, Any], entry: Dict[str, Any], duration: float) -> None: - """Fire ``subagent.complete`` with the per-branch observability payload - (tokens, cost, files touched, tool-output tail); every field is optional - and degrades gracefully on the client.""" + """Fire ``subagent.complete`` with the per-branch observability payload (tokens, cost, files touched, + tool-output tail); every field is optional and degrades gracefully on the client.""" if not self.child_progress_cb: return child = self.child @@ -785,11 +765,10 @@ 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() @@ -818,9 +797,8 @@ class _ChildRun: if not close_deferred: _close_child(child, "Failed to close child agent after delegation") - # The AIAgent turn boundary normally closes the child scope itself. This - # fallback covers failures before that boundary starts, but must not pop - # a scope while a timed-out child worker is still unwinding. + # The AIAgent turn boundary normally closes the child scope itself. This fallback covers failures before that + # boundary starts, but must not pop a scope while a timed-out child worker is still unwinding. with _quiet("Failed to close child Relay session after delegation"): from agent import relay_runtime runtime = relay_runtime.get_runtime(create=False) diff --git a/tools/delegate_tool_config.py b/tools/delegate_tool_config.py index fad6dd5f8d..f43f72a0f0 100644 --- a/tools/delegate_tool_config.py +++ b/tools/delegate_tool_config.py @@ -24,10 +24,9 @@ _HIGH_CONCURRENCY_WARNED = False MAX_DEPTH = 1 # flat by default: parent (0) -> child (1); deeper needs max_spawn_depth _MIN_SPAWN_DEPTH = 1 # floor for the configurable cap; MAX_DEPTH stays the default _LEGACY_MAX_ASYNC_WARNED = False -# No default wall-clock cap on children: legitimate heavy work (deep reviews, -# research fan-outs, slow reasoning models) was being killed mid-task. Stuck-child -# detection is the heartbeat staleness monitor; delegation.child_timeout_seconds -# opts back in. +# No default wall-clock cap on children: legitimate heavy work (deep reviews, research fan-outs, slow reasoning +# models) was being killed mid-task. Stuck-child detection is the heartbeat staleness monitor; +# delegation.child_timeout_seconds opts back in. DEFAULT_CHILD_TIMEOUT: Optional[float] = None def _cfg() -> dict: @@ -63,9 +62,8 @@ def _get_subagent_approval_callback(): return _subagent_auto_deny def _knob(key: str, env_var: Optional[str], parse, default, invalid_msg: str): - """delegation. > > default. A config value that fails ``parse`` - logs ``invalid_msg`` (``%r`` = the value) and yields the default; an env value - that fails is silently ignored.""" + """delegation. > > default. A config value that fails ``parse`` logs ``invalid_msg`` (``%r`` = the + value) and yields the default; an env value that fails is silently ignored.""" val = _cfg().get(key) if val is not None: try: @@ -111,10 +109,9 @@ def _get_worktree_isolation() -> bool: 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 warning.""" + """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 warning.""" from tools.delegate_tool import _get_max_concurrent_children if _cfg().get("max_async_children") is not None: _warn_once( @@ -130,20 +127,18 @@ 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: @@ -173,10 +168,9 @@ 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) @@ -185,9 +179,8 @@ 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) and 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) @@ -211,13 +204,11 @@ 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 @@ -243,11 +234,10 @@ def _resolve_child_credential_pool( return None 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. None when both are empty.""" + """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. 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 @@ -265,9 +255,8 @@ def _merge_request_overrides(runtime_overrides, explicit_overrides): merged["extra_body"] = explicit_extra return merged or None -# Native-SDK providers speak their own wire protocol and can't be reached via -# chat_completions against a base_url: always take the runtime-provider path -# (a configured base_url still flows through it, e.g. a Bedrock region). +# Native-SDK providers speak their own wire protocol and can't be reached via chat_completions against a base_url: +# always take the runtime-provider path (a configured base_url still flows through it, e.g. a Bedrock region). _NATIVE_SDK_PROVIDERS = frozenset({"bedrock", "vertex", "google", "google-genai"}) _EXPLICIT_API_MODES = frozenset({"chat_completions", "codex_responses", "anthropic_messages"}) @@ -287,9 +276,8 @@ def _credential_bundle(model, provider, base_url, api_key, api_mode, request_ove def _direct_endpoint_credentials(v: dict, explicit_request_overrides) -> dict: """``delegation.base_url`` branch: provider/api_mode from URL heuristics.""" - # Shared URL-based api_mode detector so Anthropic-compatible direct - # endpoints (/anthropic suffix: Azure AI Foundry, MiniMax, Zhipu, LiteLLM) - # get the Messages transport instead of 404ing on chat_completions. + # Shared URL-based api_mode detector so Anthropic-compatible direct endpoints (/anthropic suffix: Azure AI + # Foundry, MiniMax, Zhipu, LiteLLM) get the Messages transport instead of 404ing on chat_completions. from hermes_cli.runtime_provider import _detect_api_mode_for_url base_lower = v["base_url"].lower() host = base_url_hostname(v["base_url"]) @@ -305,9 +293,8 @@ def _direct_endpoint_credentials(v: dict, explicit_request_overrides) -> dict: if v["api_mode"] in _EXPLICIT_API_MODES: api_mode = v["api_mode"] - # provider configured ALONGSIDE base_url: pull that provider's request - # personality (request_overrides / max_output_tokens) onto the explicit - # endpoint. Best-effort — a resolution failure only skips the overrides. + # provider configured ALONGSIDE base_url: pull that provider's request personality (request_overrides / + # max_output_tokens) onto the explicit endpoint. Best-effort — a resolution failure only skips the overrides. request_overrides = max_output_tokens = None if v["provider"]: try: @@ -361,12 +348,11 @@ def _runtime_provider_credentials(v: dict, explicit_request_overrides) -> dict: ) def _resolve_delegation_credentials(cfg: dict, parent_agent) -> dict: - """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.""" + """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 @@ -383,12 +369,10 @@ def _resolve_delegation_credentials(cfg: dict, parent_agent) -> dict: return _runtime_provider_credentials(values, explicit_request_overrides) def _load_config() -> dict: - """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.""" + """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.""" if os.environ.get("HERMES_IGNORE_USER_CONFIG") != "1": try: from hermes_cli.config import load_config_readonly @@ -404,9 +388,8 @@ def _load_config() -> dict: except Exception: return {} -# OpenRouter routing filters: inherited from the parent, but reset to these -# defaults under a pinned provider — parent filters (e.g. only=["Anthropic"]) -# would silently force the child back onto the parent's provider. +# OpenRouter routing filters: inherited from the parent, but reset to these defaults under a pinned provider — parent +# filters (e.g. only=["Anthropic"]) would silently force the child back onto the parent's provider. # openrouter_min_coding_score stays inherited: model-gated, no-op elsewhere. _ROUTING_FILTER_DEFAULTS = ( ("providers_allowed", None), ("providers_ignored", None), ("providers_order", None), ("provider_sort", None), @@ -420,18 +403,16 @@ 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]: - """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.""" + """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) - # api_mode: each provider has its own wire, so a different provider re-derives - # (None) instead of inheriting (404s otherwise). Nous Portal is dual-wire - # within one provider (anthropic/* → Messages, else chat_completions), so + # api_mode: each provider has its own wire, so a different provider re-derives (None) instead of inheriting (404s + # otherwise). Nous Portal is dual-wire within one provider (anthropic/* → Messages, else chat_completions), so # same-provider inheritance would pin the child on the wrong wire — re-derive. _parent_provider = getattr(parent_agent, "provider", None) or "" if override_api_mode is not None: @@ -486,9 +467,8 @@ def _resolve_child_runtime( "acp_command": effective_acp_command, "acp_args": effective_acp_args, "reasoning_config": child_reasoning, - # Inherit the parent's fallback chain EXCEPT under a pinned provider: a - # mid-run 429/auth failure must not silently reroute the quiet child onto - # the parent's fallbacks. Predictability > liveness for explicit pins. + # Inherit the parent's fallback chain EXCEPT under a pinned provider: a mid-run 429/auth failure must not + # silently reroute the quiet child onto the parent's fallbacks. Predictability > liveness for explicit pins. "fallback_model": None if override_provider else (getattr(parent_agent, "_fallback_chain", None) or None), "openrouter_min_coding_score": getattr(parent_agent, "openrouter_min_coding_score", None), } diff --git a/tools/delegate_tool_dispatch.py b/tools/delegate_tool_dispatch.py index 690e4634e9..090a81dff9 100644 --- a/tools/delegate_tool_dispatch.py +++ b/tools/delegate_tool_dispatch.py @@ -77,9 +77,8 @@ def _capture_origin() -> tuple[str, str, Any, Any]: return (_origin_wake_sid, _origin_ui_session_id, *_capture_gateway_steer_authority(_origin_ui_session_id)) def _report_child_done(parent_agent, spinner_ref, entry, tag, task_labels, n_tasks, remaining) -> None: - """Print one completion line for a finished child and refresh the spinner text. - Failed/errored/timed-out children say WHY on the same line — a bare ✗ reads - as "silently dropped".""" + """Print one completion line for a finished child and refresh the spinner text. Failed/errored/timed-out children + say WHY on the same line — a bare ✗ reads as "silently dropped".""" idx = entry["task_index"] label = task_labels[idx] if idx < len(task_labels) else f"Task {idx}" status = entry.get("status", "?") @@ -94,15 +93,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`` - (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. + """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. from tools.daemon_pool import DaemonThreadPoolExecutor parent_agent, n_tasks = batch.parent_agent, len(batch.task_list) task_labels = [t["goal"][:40] for t in batch.task_list] @@ -136,11 +132,10 @@ 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: @@ -188,12 +183,11 @@ 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 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. + 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 @@ -213,13 +207,11 @@ def _resolve_async_wake_sid(origin_wake_sid: str) -> Optional[str]: def _resolve_async_session_key(parent_agent: Any, origin_ui_session_id: str) -> tuple[str, str]: """``(session_key, origin_ui_session_id)`` the async registry routes completions by. - Desktop/TUI: the routable key is the durable AIAgent.session_id — compression - 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 on the durable session_id (empty would fail closed), - so stamp the parent's durable id. + Desktop/TUI: the routable key is the durable AIAgent.session_id — compression 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 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="") @@ -235,12 +227,11 @@ def _resolve_async_session_key(parent_agent: Any, origin_ui_session_id: str) -> return session_key or agent_session_id, origin_ui_session_id def _batch_progress_token(child_agents: List[Any]) -> tuple: - """Progress token for the async registry's stale monitor: every child's - (api_call_count, current_tool, last_activity_ts). last_activity_ts ticks on - streamed chunks, tool transitions and API-call start/completion, so a child - streaming a long response counts as alive; a fully frozen token past the - threshold means the batch is wedged. ``in_tool`` is True while ANY child is - inside a tool so slow tools get the higher ceiling (mirrors the sync heartbeat).""" + """Progress token for the async registry's stale monitor: every child's (api_call_count, current_tool, + last_activity_ts). last_activity_ts ticks on streamed chunks, tool transitions and API-call start/completion, + so a child streaming a long response counts as alive; a fully frozen token past the threshold means the batch + is wedged. ``in_tool`` is True while ANY child is inside a tool so slow tools get the higher ceiling (mirrors + the sync heartbeat).""" parts = [] in_tool = False for c in child_agents: @@ -292,11 +283,10 @@ def _dispatched_payload(dispatch: dict, goals: List[str], child_agents: List[Any return payload 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 synchronously (with an explanatory ``note``) when the session - cannot receive detached completions or the async pool is at capacity.""" + """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 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) @@ -307,9 +297,8 @@ def _dispatch_background(batch: _Batch) -> str: parent_agent = batch.parent_agent session_key, origin_ui_session_id = _resolve_async_session_key(parent_agent, batch.origin_ui_session_id) child_agents = [c for (_, _, c) in batch.children] - # The batch's lifecycle is owned by the async registry now: drop the children - # from the parent's interrupt-propagation list (_build_child_agent attached - # them, which is correct for sync runs). + # The batch's lifecycle is owned by the async registry now: drop the children from the parent's + # interrupt-propagation list (_build_child_agent attached them, which is correct for sync runs). for c in child_agents: _detach_child(parent_agent, c) diff --git a/tools/delegate_tool_progress.py b/tools/delegate_tool_progress.py index 9158d97409..28db6b1a77 100644 --- a/tools/delegate_tool_progress.py +++ b/tools/delegate_tool_progress.py @@ -16,16 +16,14 @@ from tools.delegate_tool_registry import _active_subagents, _active_subagents_lo # Log-record parity with the origin module. logger = logging.getLogger("tools.delegate_tool") -# Terminal child statuses that mean "the subagent did NOT deliver a usable -# result". Shared by the CLI spinner echo, the gateway failure notice, and -# the parent-facing failure summary so every surface agrees. +# Terminal child statuses that mean "the subagent did NOT deliver a usable result". Shared by the CLI spinner echo, +# the gateway failure notice, and the parent-facing failure summary so every surface agrees. SUBAGENT_FAILURE_STATUSES = frozenset({"failed", "error", "timeout"}) @contextmanager def _quiet(log_message: Optional[str], *log_args: Any, exc_info: bool = False): - """Best-effort block: any Exception is swallowed (never reaches the run) and, - when ``log_message`` is given, logged at debug — the exception fills a trailing - unsatisfied ``%s``.""" + """Best-effort block: any Exception is swallowed (never reaches the run) and, when ``log_message`` is given, + logged at debug — the exception fills a trailing unsatisfied ``%s``.""" try: yield except Exception as exc: @@ -42,9 +40,8 @@ def _safe_progress(cb: Any, event_type: Any, *args: Any, **kwargs: Any) -> None: cb(event_type, *args, **kwargs) def _clean_error_text(error: Any, max_chars: int = 200) -> str: - """Reduce an error payload (traceback / JSON wall) to one clean line: the - exception message (last line of a traceback) or the first non-empty line, - hard-capped in length.""" + """Reduce an error payload (traceback / JSON wall) to one clean line: the exception message (last line of a + traceback) or the first non-empty line, hard-capped in length.""" lines = [ln.strip() for ln in str(error or "").strip().splitlines() if ln.strip()] if not lines: return "" @@ -72,11 +69,10 @@ def format_subagent_failure_line( class DelegateEvent(str, enum.Enum): - """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.""" + """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" @@ -95,9 +91,8 @@ _LEGACY_EVENT_MAP: Dict[str, DelegateEvent] = { "subagent_progress": DelegateEvent.TASK_PROGRESS, } -# Event → _ChildProgressRelay method name. Lifecycle strings are emitted by the -# orchestrator itself (not DelegateEvent). Any other DelegateEvent -# (TASK_TOOL_STARTED and the reserved TASK_* values) takes the tool-started +# Event → _ChildProgressRelay method name. Lifecycle strings are emitted by the orchestrator itself (not +# DelegateEvent). Any other DelegateEvent (TASK_TOOL_STARTED and the reserved TASK_* values) takes the tool-started # path; None means "recognised but ignored". _LIFECYCLE_EVENTS = frozenset({"subagent.start", "subagent.complete", "subagent.text"}) _EVENT_HANDLERS: Dict[Any, Optional[str]] = { @@ -123,10 +118,9 @@ 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: - """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}") @@ -136,12 +130,10 @@ def _build_child_system_prompt( f"{workspace_path}\n" "Use this exact path for local repository/workdir operations unless the task explicitly says otherwise." ) - # Project context files (AGENTS.md / CLAUDE.md / .cursorrules ...) via - # the SAME discovery/priority/cap logic as the main agent's prompt: - # children are built with skip_context_files=True, so without this a - # subagent works in a repo blind to its conventions. SOUL.md is skipped - # (identity belongs to the parent). workspace_path comes only from - # explicit sources (_resolve_workspace_hint, never bare getcwd), so the + # Project context files (AGENTS.md / CLAUDE.md / .cursorrules ...) via the SAME discovery/priority/cap logic + # as the main agent's prompt: children are built with skip_context_files=True, so without this a subagent + # works in a repo blind to its conventions. SOUL.md is skipped (identity belongs to the parent). + # workspace_path comes only from explicit sources (_resolve_workspace_hint, never bare getcwd), so the # install-tree-fallback leak doesn't apply. Best-effort. _ctx_files = "" with _quiet("subagent: workspace context-files load failed", exc_info=True): @@ -214,12 +206,11 @@ _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: @@ -236,9 +227,8 @@ def _batch_prefix(delegation_id: Optional[str], task_index: int, task_count: int return f"[{tag}] " if tag else "" def _emit_parent_console(parent_agent, line: str) -> None: - """Emit a progress line through ``parent_agent._safe_print`` when available - so headless stdio hosts (ACP, gateway API) can redirect it to stderr; a - bare ``print()`` would land on stdout and corrupt JSON-RPC framing.""" + """Emit a progress line through ``parent_agent._safe_print`` when available so headless stdio hosts (ACP, gateway + API) can redirect it to stderr; a bare ``print()`` would land on stdout and corrupt JSON-RPC framing.""" printer = getattr(parent_agent, "_safe_print", None) if callable(printer): with _quiet(None): @@ -260,12 +250,10 @@ 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 @@ -346,9 +334,8 @@ class _ChildProgressRelay: self._relay("subagent.thinking", preview=text) def _on_progress(self, tool_name, preview, args, kwargs): - # Pre-batched summary from a nested orchestrator's grandchild arrives in - # the tool_name slot: render distinctly (no tool-emoji lookup) and relay - # upward without re-batching. + # Pre-batched summary from a nested orchestrator's grandchild arrives in the tool_name slot: render distinctly + # (no tool-emoji lookup) and relay upward without re-batching. summary_text = tool_name or preview or "" if summary_text: self._tree_line(f"🔀 {summary_text}") @@ -386,9 +373,8 @@ def _build_child_progress_callback( parent_id: Optional[str] = None, depth: Optional[int] = None, model: Optional[str] = None, toolsets: Optional[List[str]] = None, session_ref: Optional[Dict[str, Any]] = None, ) -> Optional[callable]: - """Relay for one child's events (see ``_ChildProgressRelay``), or None when - the parent has neither a spinner nor a progress callback — the child then - runs with no progress callback at all (zero behavior change).""" + """Relay for one child's events (see ``_ChildProgressRelay``), or None when the parent has neither a spinner nor a + progress callback — the child then runs with no progress callback at all (zero behavior change).""" spinner = getattr(parent_agent, "_delegate_spinner", None) parent_cb = getattr(parent_agent, "tool_progress_callback", None) if not spinner and not parent_cb: diff --git a/tools/delegate_tool_registry.py b/tools/delegate_tool_registry.py index 3257df2db5..64991b3c90 100644 --- a/tools/delegate_tool_registry.py +++ b/tools/delegate_tool_registry.py @@ -22,19 +22,16 @@ _active_subagents_lock = threading.Lock() # subagent_id -> mutable record tracking the live child agent. Stays only # for the lifetime of the run; _run_single_child is the owner. _active_subagents: Dict[str, Dict[str, Any]] = {} -# subagent_id -> {goal, delegation_id, owner_agent_session_id} retained AFTER -# the child finishes (bounded FIFO). Child-started background processes -# routinely outlive the child (its npm ci with notify_on_complete=true finishes -# after the summary was delivered); their completion notifications reach the -# parent via the shared completion_queue and need delegation attribution even -# though the live registry entry is gone. +# subagent_id -> {goal, delegation_id, owner_agent_session_id} retained AFTER the child finishes (bounded FIFO). +# Child-started background processes routinely outlive the child (its npm ci with notify_on_complete=true finishes +# after the summary was delivered); their completion notifications reach the parent via the shared completion_queue +# and need delegation attribution even though the live registry entry is gone. _RECENT_SUBAGENTS_CAP = 200 _recent_subagents: Dict[str, Dict[str, Any]] = {} def get_subagent_attribution(task_id: Optional[str]) -> Optional[Dict[str, Any]]: - """``{subagent_id, goal, delegation_id}`` for a process task_id that belongs to a - live or recently-finished child (children run their terminal sessions under - ``task_id == subagent_id``), else None.""" + """``{subagent_id, goal, delegation_id}`` for a process task_id that belongs to a live or recently-finished child + (children run their terminal sessions under ``task_id == subagent_id``), else None.""" if not task_id or not isinstance(task_id, str): return None with _active_subagents_lock: @@ -77,11 +74,10 @@ def _unregister_subagent(subagent_id: str, *, agent: Any = None) -> None: _recent_subagents.pop(next(iter(_recent_subagents)), 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``, 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.""" + """Atomically close steer acceptance and drain its final durable artifact. ``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: @@ -118,14 +114,11 @@ def steer_subagent( ) -> bool: """Queue steering text into a running subagent without stopping it. - 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``. + 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 @@ -171,10 +164,9 @@ def list_active_subagents() -> List[Dict[str, Any]]: return [{k: v for k, v in r.items() if k not in _PRIVATE_RECORD_KEYS} for r in _active_subagents.values()] def _is_descendant_of(child_agent: Any, parent_agent: Any, max_hops: int = 8) -> bool: - """True when *child_agent* sits below *parent_agent* in the spawn tree (walks the - ``_delegate_parent_ref`` weakref chain stamped at build time). Identity only — - a parent may steer/stop its own children and grandchildren, never a sibling - tree owned by another conversation.""" + """True when *child_agent* sits below *parent_agent* in the spawn tree (walks the ``_delegate_parent_ref`` weakref + chain stamped at build time). Identity only — a parent may steer/stop its own children and grandchildren, never + a sibling tree owned by another conversation.""" if child_agent is None or parent_agent is None: return False cur = child_agent @@ -193,9 +185,8 @@ def _is_descendant_of(child_agent: Any, parent_agent: Any, max_hops: int = 8) -> _CONTROL_ACTIONS = frozenset({"list", "steer", "stop"}) def _resolve_session_lineage(session_id: Optional[str], parent_agent: Any) -> str: - """Tip of a session id's compression lineage via the parent's live SessionDB - (best-effort; input unchanged when unavailable) so a delegation dispatched - before a compression rotation still matches the rotated parent.""" + """Tip of a session id's compression lineage via the parent's live SessionDB (best-effort; input unchanged when + unavailable) so a delegation dispatched before a compression rotation still matches the rotated parent.""" sid = str(session_id or "") db = getattr(parent_agent, "_session_db", None) if not sid or db is None: @@ -258,9 +249,8 @@ 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) over the same registry the TUI overlay drives, 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_results.py b/tools/delegate_tool_results.py index 36a84bac61..52daf1eeb5 100644 --- a/tools/delegate_tool_results.py +++ b/tools/delegate_tool_results.py @@ -54,11 +54,10 @@ def _looks_like_error_output(content: Any) -> bool: return first.startswith(("error:", "failed:", "traceback ", "exception:")) def _extract_output_tail(result: Dict[str, Any], *, max_entries: int = 12, max_chars: int = 8000) -> List[Dict[str, Any]]: - """Last N tool-call results ``{tool, preview, is_error}`` from a child's - conversation (the overlay's "Output" section), chronological order. Content - blocks are flattened first so a block-wrapped "Error: ..." is still flagged; - line structure is preserved (capped at ``max_chars``) so the overlay shows - real output rather than a whitespace-collapsed blob.""" + """Last N tool-call results ``{tool, preview, is_error}`` from a child's conversation (the overlay's "Output" + section), chronological order. Content blocks are flattened first so a block-wrapped "Error: ..." is still + flagged; line structure is preserved (capped at ``max_chars``) so the overlay shows real output rather than a + whitespace-collapsed blob.""" messages = result.get("messages") if isinstance(result, dict) else None if not isinstance(messages, list): return [] @@ -100,10 +99,9 @@ def _sanitize_tool_target(key: str, value: Any) -> Any: hostname = parsed.hostname if not hostname: return None - # ``SplitResult.netloc`` includes ``user:password@``. Rebuild - # the authority from parsed host/port so hook-visible history - # cannot carry URL credentials. Bracket IPv6 literals before - # appending a validated port. + # ``SplitResult.netloc`` includes ``user:password@``. Rebuild the authority from parsed host/port so + # hook-visible history cannot carry URL credentials. Bracket IPv6 literals before appending a + # validated port. host = f"[{hostname}]" if ":" in hostname else hostname netloc = f"{host}:{parsed.port}" if parsed.port is not None else host return urlunsplit((parsed.scheme, netloc, parsed.path, "", "")) @@ -166,23 +164,20 @@ def _subagent_stop_tool_call_history(tool_trace: Any) -> List[Dict[str, Any]]: }) return history -# Hard per-summary character ceiling layered on top of the dynamic headroom -# budget (see _apply_summary_budget): belt-and-suspenders for models that -# ignore "be concise". 0 disables the ceiling. +# Hard per-summary character ceiling layered on top of the dynamic headroom budget (see _apply_summary_budget): +# belt-and-suspenders for models that ignore "be concise". 0 disables the ceiling. DEFAULT_MAX_SUMMARY_CHARS = 24000 -# Fraction of the parent's *remaining* context headroom the whole batch of -# summaries may consume, split per summary, so N children can't collectively -# blow the parent's window (the compression/429 death spiral). +# Fraction of the parent's *remaining* context headroom the whole batch of summaries may consume, split per summary, +# so N children can't collectively blow the parent's window (the compression/429 death spiral). _SUMMARY_HEADROOM_FRACTION = 0.5 # Floor so a single summary always gets a usable slice even when the parent is # already nearly full — below this we'd be truncating to noise. _MIN_SUMMARY_CHARS = 2000 def _spill_summary_to_file(task_index: int, summary: str) -> Optional[str]: - """Write the full summary under ``cache/delegation`` (mounted read-only into - remote backends via ``credential_files._CACHE_DIRS``, so the parent's - terminal/``read_file`` can page it on any backend). Absolute path, or None on - failure — the trimmed head+tail is still returned regardless.""" + """Write the full summary under ``cache/delegation`` (mounted read-only into remote backends via + ``credential_files._CACHE_DIRS``, so the parent's terminal/``read_file`` can page it on any backend). Absolute + path, or None on failure — the trimmed head+tail is still returned regardless.""" try: from hermes_constants import get_hermes_dir import datetime as _dt @@ -190,9 +185,8 @@ def _spill_summary_to_file(task_index: int, summary: str) -> Optional[str]: cache_dir.mkdir(parents=True, exist_ok=True) path = cache_dir / f"subagent-summary-{task_index}-{_dt.datetime.now().strftime('%Y%m%d_%H%M%S_%f')}.txt" from tools.spill_safety import write_text_exclusive - # Exclusive symlink-refusing create; not private because cache/delegation - # is bind-mounted read-only into remote backends whose container UID - # must be able to read it. + # Exclusive symlink-refusing create; not private because cache/delegation is bind-mounted read-only into + # remote backends whose container UID must be able to read it. write_text_exclusive(path, summary, private=False) return str(path) except Exception as exc: @@ -200,10 +194,9 @@ def _spill_summary_to_file(task_index: int, summary: str) -> Optional[str]: return None def _trim_summary_with_footer(summary: str, cap: int, task_index: int) -> tuple[str, Optional[str]]: - """``(model_text, spill_path)`` for one over-budget summary: a ~75% head / - ~25% tail window snapped to line boundaries (so the opening AND the closing - outcomes/files-changed/issues both survive), the full text spilled to disk, - and a footer giving the exact ``read_file offset=`` for the omitted middle.""" + """``(model_text, spill_path)`` for one over-budget summary: a ~75% head / ~25% tail window snapped to line + boundaries (so the opening AND the closing outcomes/files-changed/issues both survive), the full text spilled + to disk, and a footer giving the exact ``read_file offset=`` for the omitted middle.""" original_len = len(summary) head_budget = int(cap * 0.75) tail_budget = cap - head_budget @@ -235,10 +228,9 @@ def _trim_summary_with_footer(summary: str, cap: int, task_index: int) -> tuple[ return head + "\n\n[... middle omitted — see footer ...]\n\n" + tail + "\n".join(footer_lines), spill_path def _parent_summary_char_budget(parent_agent, n_summaries: int) -> Optional[int]: - """Per-summary char budget from the parent's *remaining* context headroom - (context length − prompt tokens − the compressor's output reserve), a - fraction of it split across the batch at ~4 chars/token. None when the - parent's context state is unknown — caller then uses the static ceiling only.""" + """Per-summary char budget from the parent's *remaining* context headroom (context length − prompt tokens − the + compressor's output reserve), a fraction of it split across the batch at ~4 chars/token. None when the parent's + context state is unknown — caller then uses the static ceiling only.""" try: compressor = getattr(parent_agent, "context_compressor", None) context_length = getattr(compressor, "context_length", None) @@ -257,10 +249,9 @@ def _parent_summary_char_budget(parent_agent, n_summaries: int) -> Optional[int] return None def _apply_summary_budget(results: List[Dict[str, Any]], parent_agent) -> None: - """Trim subagent summaries in-place so a batch can't overflow the parent's - context window (full text spilled to disk). Per-summary cap = MIN(dynamic - headroom budget, static ``delegation.max_summary_chars`` ceiling; 0 = disabled); - over-cap summaries become head+tail plus a pointer to the spill file.""" + """Trim subagent summaries in-place so a batch can't overflow the parent's context window (full text spilled to + disk). Per-summary cap = MIN(dynamic headroom budget, static ``delegation.max_summary_chars`` ceiling; 0 = + disabled); over-cap summaries become head+tail plus a pointer to the spill file.""" from tools.delegate_tool import _load_config summaries = [r for r in results if isinstance(r, dict) and isinstance(r.get("summary"), str) and r["summary"]] if not summaries: diff --git a/tools/delegate_tool_tasks.py b/tools/delegate_tool_tasks.py index 3e92fffa0d..df4dd6b7ed 100644 --- a/tools/delegate_tool_tasks.py +++ b/tools/delegate_tool_tasks.py @@ -9,12 +9,10 @@ import json import re from typing import Any, Dict, List, Optional -# Placeholder shapes for batch goal validation: bare 'TODO' / 'task N' labels, -# or unexpanded template markers. The marker regex is deliberately NARROW — -# only snake_case / space-separated placeholder identifiers (``, -# `{file path}`, ``), the shape LLM templates leave behind. Bare -# single-word brackets must never be rejected: legitimate goals are full of -# generics (`Vec`), HTML tags (`
`), dict snippets (`{"key": 1}`), glob +# Placeholder shapes for batch goal validation: bare 'TODO' / 'task N' labels, or unexpanded template markers. The +# marker regex is deliberately NARROW — only snake_case / space-separated placeholder identifiers (``, +# `{file path}`, ``), the shape LLM templates leave behind. Bare single-word brackets must never be +# rejected: legitimate goals are full of generics (`Vec`), HTML tags (`
`), dict snippets (`{"key": 1}`), glob # braces (`{a,b}`) and f-string style (`{i}`). _PLACEHOLDER_GOAL_RE = re.compile(r"^(todo|task\s*\d+)$", re.IGNORECASE) _TEMPLATE_MARKER_RE = re.compile( @@ -38,13 +36,11 @@ 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())): @@ -110,9 +106,8 @@ def _normalize_task_list( def _coerce_task_schemas( task_list: List[Dict[str, Any]], output_schema: Optional[Dict[str, Any]] ) -> tuple[List[Optional[Dict[str, Any]]], Optional[str]]: - """Per-task coerced output schemas. A malformed output_schema fails the whole - call before any child spawns; schema-less tasks resolve to None and take no - new code paths downstream.""" + """Per-task coerced output schemas. A malformed output_schema fails the whole call before any child spawns; + schema-less tasks resolve to None and take no new code paths downstream.""" from tools.delegation_output_schema import coerce_output_schema task_schemas: List[Optional[Dict[str, Any]]] = [] for i, task in enumerate(task_list): diff --git a/tools/delegate_tool_toolsets.py b/tools/delegate_tool_toolsets.py index 581de9ad9c..a28eb4b4b1 100644 --- a/tools/delegate_tool_toolsets.py +++ b/tools/delegate_tool_toolsets.py @@ -40,9 +40,8 @@ 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: @@ -53,9 +52,8 @@ def _expand_parent_toolsets(parent_toolsets: set) -> set: return expanded def _strip_blocked_tools(toolsets: List[str]) -> List[str]: - """Remove toolsets whose tools are ALL blocked (derived from DELEGATE_BLOCKED_TOOLS - so the two can't drift) plus composite toolsets children must never get - (``delegation``, ``kanban``).""" + """Remove toolsets whose tools are ALL blocked (derived from DELEGATE_BLOCKED_TOOLS so the two can't drift) plus + composite toolsets children must never get (``delegation``, ``kanban``).""" blocked_toolset_names = {"delegation", "kanban"} | { name for name, defn in TOOLSETS.items() if all(t in DELEGATE_BLOCKED_TOOLS for t in defn.get("tools", [])) } @@ -74,13 +72,11 @@ 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]]: - """``(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 + """``(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)