refactor(delegate): docstring compaction by hand (every rule/why kept)

This commit is contained in:
Teknium
2026-09-02 19:16:07 -07:00
parent aa31f83995
commit 1283ccd855
8 changed files with 195 additions and 279 deletions
+18 -26
View File
@@ -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"}
+42 -60
View File
@@ -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-<sid>-<ts>.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-<sid>-<ts>.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()
+48 -75
View File
@@ -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:<name>`` 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:<name>`` 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)
+24 -32
View File
@@ -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)
+21 -31
View File
@@ -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
+24 -30
View File
@@ -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)
+7 -10
View File
@@ -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())):
+11 -15
View File
@@ -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: