refactor(delegate): reflow comments/docstrings to 118 cols (word-identical, AST-identical)
This commit is contained in:
@@ -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.<name>`` 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.<name>`` 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/<id>/task-<n>.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/<id>/task-<n>.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"}
|
||||
|
||||
@@ -123,9 +123,8 @@ def _diag_sizes(child: Any) -> List[str]:
|
||||
return ["## Prompt / schema sizes"] + _diag_section("system_prompt: <error: ", _prompt) + _diag_section("tool_schema: <error: ", _tools)
|
||||
|
||||
def _diag_threads(worker_thread: Optional[threading.Thread]) -> 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-<sid>-<ts>.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-<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
|
||||
@@ -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)
|
||||
|
||||
@@ -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.<key> > <env_var> > 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.<key> > <env_var> > 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:<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
|
||||
@@ -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),
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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 (`<feature_name>`,
|
||||
# `{file path}`, `<FEATURE-NAME>`), the shape LLM templates leave behind. Bare
|
||||
# single-word brackets must never be rejected: legitimate goals are full of
|
||||
# generics (`Vec<T>`), HTML tags (`<div>`), 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 (`<feature_name>`,
|
||||
# `{file path}`, `<FEATURE-NAME>`), the shape LLM templates leave behind. Bare single-word brackets must never be
|
||||
# rejected: legitimate goals are full of generics (`Vec<T>`), HTML tags (`<div>`), 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):
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user