refactor(delegate): AST-identical re-layout (pack signatures/call args, join split literals)

This commit is contained in:
Teknium
2026-09-02 17:59:28 -07:00
parent 901ba7db6f
commit b18228f377
9 changed files with 131 additions and 442 deletions

View File

@@ -107,7 +107,6 @@ def _open_child_session_db(parent_agent) -> Any:
return None
try:
from hermes_state import get_shared_session_db
_parent_db_path = getattr(parent_session_db, "db_path", None)
return get_shared_session_db(_parent_db_path) if _parent_db_path is not None else get_shared_session_db()
except Exception:
@@ -115,24 +114,14 @@ def _open_child_session_db(parent_agent) -> Any:
return None
def _construct_child_agent(
rt: Dict[str, Any],
*,
task_index: int,
max_iterations: int,
parent_agent,
child_toolsets: List[str],
child_disabled_toolsets: List[str],
child_prompt: str,
child_progress_cb: Any,
child_session_db: Any,
override_provider: Optional[str],
override_request_overrides: Optional[Dict[str, Any]],
rt: Dict[str, Any], *, task_index: int, max_iterations: int, parent_agent, child_toolsets: List[str],
child_disabled_toolsets: List[str], child_prompt: str, child_progress_cb: Any, child_session_db: Any,
override_provider: Optional[str], override_request_overrides: Optional[Dict[str, Any]],
):
"""Instantiate the child AIAgent; releases the dedicated SessionDB handle on
a construction failure (no child close() will ever run)."""
from run_agent import AIAgent
from agent.delegation_context import delegated_child_context
child_thinking_cb = None
if child_progress_cb:
@@ -186,13 +175,9 @@ def _announce_child_spawn(child, parent_agent, child_progress_cb, *, goal, subag
try:
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
_invoke_hook(
"subagent_start",
parent_session_id=getattr(parent_agent, "session_id", None),
parent_turn_id=getattr(parent_agent, "_current_turn_id", "") or "",
parent_subagent_id=parent_subagent_id,
child_session_id=getattr(child, "session_id", None),
child_subagent_id=subagent_id,
child_role=role,
"subagent_start", parent_session_id=getattr(parent_agent, "session_id", None),
parent_turn_id=getattr(parent_agent, "_current_turn_id", "") or "", parent_subagent_id=parent_subagent_id,
child_session_id=getattr(child, "session_id", None), child_subagent_id=subagent_id, child_role=role,
child_goal=goal,
)
except Exception:
@@ -226,7 +211,6 @@ def _build_child_agent(
can run on a different provider:model pair.
"""
import uuid as _uuid
# Role is depth-derived: a child may delegate iff the kill switch is on and
# depth budget remains below max_spawn_depth. The `role` arg is ignored.
child_depth = getattr(parent_agent, "_delegate_depth", 0) + 1
@@ -241,12 +225,8 @@ def _build_child_agent(
delegation_cfg = _load_config()
child_toolsets, child_disabled_toolsets = _resolve_child_toolsets(parent_agent, toolsets, effective_role)
child_prompt = _build_child_system_prompt(
goal,
context,
workspace_path=_resolve_workspace_hint(parent_agent),
role=effective_role,
max_spawn_depth=max_spawn,
child_depth=child_depth,
goal, context, workspace_path=_resolve_workspace_hint(parent_agent), role=effective_role,
max_spawn_depth=max_spawn, child_depth=child_depth,
)
parent_api_key = getattr(parent_agent, "api_key", None)
if (not parent_api_key) and hasattr(parent_agent, "_client_kwargs"):
@@ -266,30 +246,16 @@ def _build_child_agent(
session_ref=child_session_ref,
)
rt = _resolve_child_runtime(
parent_agent,
delegation_cfg,
parent_api_key,
model=model,
override_provider=override_provider,
override_base_url=override_base_url,
override_api_key=override_api_key,
override_api_mode=override_api_mode,
override_max_tokens=override_max_tokens,
override_acp_command=override_acp_command,
parent_agent, delegation_cfg, parent_api_key, model=model, override_provider=override_provider,
override_base_url=override_base_url, override_api_key=override_api_key, override_api_mode=override_api_mode,
override_max_tokens=override_max_tokens, override_acp_command=override_acp_command,
override_acp_args=override_acp_args,
)
child_session_db = _open_child_session_db(parent_agent)
child = _construct_child_agent(
rt,
task_index=task_index,
max_iterations=max_iterations,
parent_agent=parent_agent,
child_toolsets=child_toolsets,
child_disabled_toolsets=child_disabled_toolsets,
child_prompt=child_prompt,
child_progress_cb=child_progress_cb,
child_session_db=child_session_db,
override_provider=override_provider,
rt, task_index=task_index, max_iterations=max_iterations, parent_agent=parent_agent,
child_toolsets=child_toolsets, child_disabled_toolsets=child_disabled_toolsets, child_prompt=child_prompt,
child_progress_cb=child_progress_cb, child_session_db=child_session_db, override_provider=override_provider,
override_request_overrides=override_request_overrides,
)
child._print_fn = getattr(parent_agent, "_print_fn", None)
@@ -328,15 +294,8 @@ def _build_child_agent(
return child
def _run_single_child(
task_index: int,
goal: str,
child=None,
parent_agent=None,
*,
owner_session_id: Optional[str] = None,
owner_transport: Any = None,
owner_session_record: Any = None,
**_kwargs,
task_index: int, goal: str, child=None, parent_agent=None, *, owner_session_id: Optional[str] = None,
owner_transport: Any = None, owner_session_record: Any = None, **_kwargs,
) -> Dict[str, Any]:
"""Run a pre-built child agent (called from a worker thread) and return its result entry.
@@ -363,11 +322,7 @@ def _run_single_child(
# TUI/RPC registry entry (kill/pause/status by subagent_id); None for test
# doubles without a stable id. Unregistered in the finally block.
_subagent_id = _register_child(
child,
parent_agent,
goal,
owner_session_id=owner_session_id,
owner_transport=owner_transport,
child, parent_agent, goal, owner_session_id=owner_session_id, owner_transport=owner_transport,
owner_session_record=owner_session_record,
)
worktree = _WorktreeReporter()
@@ -380,15 +335,8 @@ def _run_single_child(
goal = ws.goal
_relay_child_text = _make_text_relay(child_progress_cb)
result, failure = _await_child(
child,
goal,
ws,
_relay_child_text,
task_index=task_index,
subagent_id=_subagent_id,
child_start=child_start,
child_progress_cb=child_progress_cb,
worktree=worktree,
child, goal, ws, _relay_child_text, task_index=task_index, subagent_id=_subagent_id,
child_start=child_start, child_progress_cb=child_progress_cb, worktree=worktree,
)
if failure is not None:
_child_close_deferred = failure.close_deferred
@@ -426,31 +374,18 @@ def _run_single_child(
finally:
_cleanup_child_run(
child,
parent_agent,
subagent_id=_subagent_id,
heartbeat=heartbeat,
child_pool=child_pool,
leased_cred_id=leased_cred_id,
close_deferred=_child_close_deferred,
child, parent_agent, subagent_id=_subagent_id, heartbeat=heartbeat, child_pool=child_pool,
leased_cred_id=leased_cred_id, close_deferred=_child_close_deferred,
)
def _build_children(
task_list: List[Dict[str, Any]],
task_schemas: List[Optional[Dict[str, Any]]],
creds: Dict[str, Any],
*,
top_role: str,
max_iterations: int,
parent_agent,
live_deleg_id: Optional[str],
live_writers: list,
task_list: List[Dict[str, Any]], task_schemas: List[Optional[Dict[str, Any]]], creds: Dict[str, Any], *,
top_role: str, max_iterations: int, parent_agent, live_deleg_id: Optional[str], live_writers: list,
) -> tuple[List[tuple], Optional[str]]:
"""Build every child on the main thread (construction is not thread-safe);
``(children, None)`` or ``([], error)`` on an explicit-pin preflight failure."""
from tools.delegation_live_log import wrap_progress_callback
children = []
for i, t in enumerate(task_list):
effective_role = _normalize_role(t.get("role") or top_role)
@@ -458,7 +393,6 @@ def _build_children(
_child_context = t.get("context")
if _task_schema is not None:
from tools.delegation_output_schema import append_output_contract
_child_context = append_output_contract(_child_context, _task_schema)
try:
child = _build_child_preserving_parent_tools(
@@ -505,18 +439,10 @@ def _build_children(
def delegate_task(
goal: Optional[str] = None,
context: Optional[str] = None,
tasks: Optional[List[Dict[str, Any]]] = None,
max_iterations: Optional[int] = None,
role: Optional[str] = None,
background: Optional[bool] = None,
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,
goal: Optional[str] = None, context: Optional[str] = None, tasks: Optional[List[Dict[str, Any]]] = None,
max_iterations: Optional[int] = None, role: Optional[str] = None, background: Optional[bool] = None,
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.
@@ -550,10 +476,8 @@ def delegate_task(
max_spawn = _get_max_spawn_depth()
if depth >= max_spawn:
return tool_error(
f"Delegation depth limit reached (depth={depth}, "
f"max_spawn_depth={max_spawn}). Raise "
f"delegation.max_spawn_depth in config.yaml if deeper "
f"nesting is required (no hard ceiling, but each level "
f"Delegation depth limit reached (depth={depth}, max_spawn_depth={max_spawn}). Raise "
f"delegation.max_spawn_depth in config.yaml if deeper nesting is required (no hard ceiling, but each level "
f"multiplies API cost)."
)
@@ -563,8 +487,7 @@ def delegate_task(
# so budgets stay predictable (kwarg kept for internal callers/tests).
if max_iterations is not None and max_iterations != default_max_iter:
logger.debug(
"delegate_task: ignoring caller-supplied max_iterations=%s; "
"using delegation.max_iterations=%s from config",
"delegate_task: ignoring caller-supplied max_iterations=%s; using delegation.max_iterations=%s from config",
max_iterations, default_max_iter,
)
@@ -588,7 +511,6 @@ def delegate_task(
# 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")
)
@@ -596,14 +518,8 @@ def delegate_task(
origin = _capture_origin()
children, err = _build_children(
task_list,
task_schemas,
creds,
top_role=top_role,
max_iterations=default_max_iter,
parent_agent=parent_agent,
live_deleg_id=live_deleg_id,
live_writers=live_writers,
task_list, task_schemas, creds, top_role=top_role, max_iterations=default_max_iter, parent_agent=parent_agent,
live_deleg_id=live_deleg_id, live_writers=live_writers,
)
if err:
return tool_error(err)
@@ -667,8 +583,7 @@ def _build_top_level_description() -> str:
"require a verifiable handle (URL, ID, absolute path) and verify it "
"yourself before telling the user the operation succeeded.\n"
+ restrictions_rule +
"- Children inherit the parent model unless pinned via "
"delegation.provider / delegation.model in config.yaml."
"- Children inherit the parent model unless pinned via delegation.provider / delegation.model in config.yaml."
)
def _build_tasks_param_description() -> str:
@@ -703,8 +618,7 @@ DELEGATE_TASK_SCHEMA = {
# 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."
"Description is rebuilt at every get_definitions() call to reflect the user's current delegation limits."
),
"parameters": {
"type": "object",
@@ -721,29 +635,23 @@ DELEGATE_TASK_SCHEMA = {
"goal": {
"type": "string",
"description": (
"What this subagent should accomplish. Be "
"specific and self-contained — it knows "
"What this subagent should accomplish. Be specific and self-contained — it knows "
"nothing about your conversation history."
),
},
"context": {
"type": "string",
"description": (
"Background THIS child needs: file paths, "
"error messages, constraints. Each child "
"sees only its own context — repeat shared "
"background in every task that needs it."
"Background THIS child needs: file paths, error messages, constraints. Each child "
"sees only its own context — repeat shared background in every task that needs it."
),
},
"output_schema": {
"type": "object",
"description": (
"Optional JSON Schema this child's final "
"answer must validate against (told to the "
"child up front; parent validates with one "
"bounded correction retry; result gains "
"schema_valid, plus schema_errors on "
"failure). Keep it forgiving — require only "
"Optional JSON Schema this child's final answer must validate against (told to the "
"child up front; parent validates with one bounded correction retry; result gains "
"schema_valid, plus schema_errors on failure). Keep it forgiving — require only "
"fields you will read."
),
},
@@ -767,23 +675,18 @@ DELEGATE_TASK_SCHEMA = {
"course-correction text into one child (subagent_id + "
"message) without stopping it; 'stop' = end one child "
"early (subagent_id; partial result still returns). "
"Control actions return immediately; goal/tasks are "
"ignored unless spawning."
"Control actions return immediately; goal/tasks are ignored unless spawning."
),
},
"subagent_id": {
"type": "string",
"description": (
"Target for action='steer'/'stop' (ids from the spawn "
"response or action='list')."
),
"description": ("Target for action='steer'/'stop' (ids from the spawn response or action='list')."),
},
"message": {
"type": "string",
"description": (
"For action='steer': the course correction, appended to "
"the child's next tool result mid-run. Be directive and "
"specific."
"the child's next tool result mid-run. Be directive and specific."
),
},
},
@@ -829,16 +732,10 @@ registry.register(
toolset="delegation",
schema=DELEGATE_TASK_SCHEMA,
handler=lambda args, **kw: delegate_task(
goal=args.get("goal"),
context=args.get("context"),
tasks=_strip_model_hidden_task_fields(args.get("tasks")),
max_iterations=args.get("max_iterations"),
role=args.get("role"),
background=_model_background_value(args, kw.get("parent_agent")),
output_schema=args.get("output_schema"),
action=args.get("action"),
subagent_id=args.get("subagent_id"),
message=args.get("message"),
goal=args.get("goal"), context=args.get("context"), tasks=_strip_model_hidden_task_fields(args.get("tasks")),
max_iterations=args.get("max_iterations"), role=args.get("role"),
background=_model_background_value(args, kw.get("parent_agent")), output_schema=args.get("output_schema"),
action=args.get("action"), subagent_id=args.get("subagent_id"), message=args.get("message"),
parent_agent=kw.get("parent_agent"),
),
check_fn=check_delegate_requirements,

View File

@@ -90,7 +90,6 @@ def _signal_child_stop(child: Any, *reason: str) -> None:
def _format_thread_stack(frame: Any, indent: str) -> List[str]:
import traceback as _traceback
return [f"{indent}{sub}" for frame_line in _traceback.format_stack(frame) for sub in frame_line.rstrip().split("\n")]
_DIAG_CHILD_ATTRS = (
@@ -121,7 +120,6 @@ def _diag_threads(worker_thread: Optional[threading.Thread]) -> List[str]:
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()
if worker_thread is not None and worker_thread.is_alive():
@@ -153,13 +151,8 @@ def _diag_threads(worker_thread: Optional[threading.Thread]) -> List[str]:
return lines
def _dump_subagent_timeout_diagnostic(
*,
child: Any,
task_index: int,
timeout_seconds: float,
duration_seconds: float,
worker_thread: Optional[threading.Thread],
goal: str,
*, 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
@@ -170,7 +163,6 @@ def _dump_subagent_timeout_diagnostic(
try:
from hermes_constants import get_hermes_home
import datetime as _dt
logs_dir = get_hermes_home() / "logs"
try:
logs_dir.mkdir(parents=True, exist_ok=True)
@@ -183,19 +175,10 @@ def _dump_subagent_timeout_diagnostic(
if len(_goal_preview) > 1000:
_goal_preview = _goal_preview[:1000] + " ...[truncated]"
lines: List[str] = [
"# Subagent timeout diagnostic — issue #14726",
f"# Generated: {_dt.datetime.now().isoformat()}",
"",
"## Timeout",
f" task_index: {task_index}",
f" subagent_id: {subagent_id}",
f" configured_timeout: {timeout_seconds}s",
f" actual_duration: {duration_seconds:.2f}s",
"",
"## Goal",
_goal_preview or "(empty)",
"",
"## Child config",
"# Subagent timeout diagnostic — issue #14726", f"# Generated: {_dt.datetime.now().isoformat()}", "",
"## Timeout", f" task_index: {task_index}", f" subagent_id: {subagent_id}",
f" configured_timeout: {timeout_seconds}s", f" actual_duration: {duration_seconds:.2f}s", "", "## Goal",
_goal_preview or "(empty)", "", "## Child config",
]
for attr in _DIAG_CHILD_ATTRS:
try:
@@ -235,7 +218,6 @@ def _start_heartbeat(child: Any, parent_agent: Any, task_index: int) -> tuple:
and the finally-path join can be skipped safely.
"""
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,
# activity_ts) all froze; thresholds differ idle vs in-tool.
@@ -289,12 +271,7 @@ def _start_heartbeat(child: Any, parent_agent: Any, task_index: int) -> tuple:
return _heartbeat_stop, threading.Thread(target=_heartbeat_loop, daemon=True)
def _register_child(
child: Any,
parent_agent: Any,
goal: str,
*,
owner_session_id: Optional[str],
owner_transport: Any,
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.
@@ -308,7 +285,6 @@ def _register_child(
if owner_session_id is None:
with _quiet(None):
from gateway.session_context import get_session_env
owner_session_id = get_session_env("HERMES_UI_SESSION_ID", "") or None
if owner_session_id and (owner_transport is None or owner_session_record is None):
owner_transport, owner_session_record = _capture_gateway_steer_authority(owner_session_id)
@@ -360,7 +336,6 @@ class _WorktreeReporter:
if info is None:
return
from tools import subagent_worktree
try:
entry_dict["worktree"] = subagent_worktree.finalize_subagent_worktree(info)
except Exception as e:
@@ -383,19 +358,16 @@ def _create_isolated_worktree(parent_agent: Any, parent_task_id: Any, subagent_i
child's terminal starts there). Git-only, local-backend-only; failures
degrade silently to the shared workspace. Returns the worktree info or None."""
from tools.delegate_tool import _get_worktree_isolation, _resolve_workspace_hint
if not _get_worktree_isolation():
return None
try:
from tools import subagent_worktree
if not subagent_worktree.local_backend_active():
logger.debug("worktree isolation skipped: non-local terminal backend")
return None
_parent_cwd = None
with _quiet(None):
from tools.terminal_tool import get_session_cwd as _gsc
_parent_cwd = _gsc(parent_task_id)
return subagent_worktree.create_subagent_worktree(
_parent_cwd or _resolve_workspace_hint(parent_agent), subagent_id=subagent_id,
@@ -405,12 +377,7 @@ def _create_isolated_worktree(parent_agent: Any, parent_task_id: Any, subagent_i
return None
def _seed_child_workspace(
child: Any,
parent_agent: Any,
goal: str,
task_index: int,
subagent_id: Optional[str],
worktree: _WorktreeReporter,
child: Any, parent_agent: Any, goal: str, task_index: int, subagent_id: Optional[str], worktree: _WorktreeReporter,
) -> _ChildWorkspace:
"""Seed cwd/container aliases and optional worktree isolation for the child.
@@ -418,7 +385,6 @@ def _seed_child_workspace(
worktree contract note when isolation engaged.
"""
import uuid as _uuid
child_task_id = subagent_id or f"subagent-{task_index}-{_uuid.uuid4().hex[:8]}"
parent_task_id = getattr(parent_agent, "_current_task_id", None)
# Seed the child's cwd record from the parent's: same starting directory,
@@ -426,7 +392,6 @@ def _seed_child_workspace(
# isolation keys containers by task_id; the child must share the PARENT's.
with _quiet("Child cwd seed failed: %s"):
from tools.terminal_tool import get_session_cwd, record_session_cwd, register_container_alias
record_session_cwd(child_task_id, get_session_cwd(parent_task_id))
register_container_alias(child_task_id, parent_task_id)
@@ -434,12 +399,10 @@ def _seed_child_workspace(
if _worktree_info is not None:
with _quiet("worktree cwd seed failed: %s"):
from tools.terminal_tool import record_session_cwd as _rsc
_rsc(child_task_id, _worktree_info["path"])
# The child's context is already built; carry the isolation contract on
# the goal message instead (same turn, no system-prompt mutation).
from tools.subagent_worktree import build_worktree_context_note
goal = goal + build_worktree_context_note(_worktree_info)
worktree.info = _worktree_info
@@ -508,16 +471,8 @@ def _make_text_relay(child_progress_cb: Any):
return _relay_child_text
def _await_child(
child: Any,
goal: str,
ws: "_ChildWorkspace",
relay_child_text: Any,
*,
task_index: int,
subagent_id: Optional[str],
child_start: float,
child_progress_cb: Any,
worktree: _WorktreeReporter,
child: Any, goal: str, ws: "_ChildWorkspace", relay_child_text: Any, *, task_index: int, subagent_id: Optional[str],
child_start: float, child_progress_cb: Any, worktree: _WorktreeReporter,
) -> tuple[Optional[Dict[str, Any]], Optional[_ChildFailure]]:
"""Run the child's conversation on a daemon worker; ``(result, None)`` or
``(None, failure)`` on timeout/exception.
@@ -531,7 +486,6 @@ def _await_child(
"""
from tools.delegate_tool import (_get_child_timeout, _get_subagent_approval_callback, _set_subagent_approval_cb)
from tools.daemon_pool import DaemonThreadPoolExecutor
child_timeout = _get_child_timeout()
executor = DaemonThreadPoolExecutor(
max_workers=1, initializer=_set_subagent_approval_cb, initargs=(_get_subagent_approval_callback(),),
@@ -542,7 +496,6 @@ def _await_child(
def _run_with_thread_capture():
worker_thread_holder["t"] = threading.current_thread()
from agent.delegation_context import delegated_child_context
with delegated_child_context(str(getattr(child, "session_id", "") or "")):
return child.run_conversation(user_message=goal, task_id=ws.child_task_id, stream_callback=relay_child_text)
@@ -551,17 +504,9 @@ def _await_child(
return future.result(timeout=child_timeout), None
except Exception as exc:
return None, _handle_child_wait_failure(
exc,
child=child,
task_index=task_index,
goal=goal,
subagent_id=subagent_id,
child_future=future,
child_timeout=child_timeout,
child_start=child_start,
child_progress_cb=child_progress_cb,
worker_thread_holder=worker_thread_holder,
worktree=worktree,
exc, child=child, task_index=task_index, goal=goal, subagent_id=subagent_id, child_future=future,
child_timeout=child_timeout, child_start=child_start, child_progress_cb=child_progress_cb,
worker_thread_holder=worker_thread_holder, worktree=worktree,
)
finally:
# Shut down without waiting — a child stuck on blocking I/O would hang wait=True forever.
@@ -591,18 +536,9 @@ def _finish_failed_entry(
return entry
def _handle_child_wait_failure(
exc: BaseException,
*,
child: Any,
task_index: int,
goal: str,
subagent_id: Optional[str],
child_future: Any,
child_timeout: Optional[float],
child_start: float,
child_progress_cb: Any,
worker_thread_holder: Dict[str, Optional[threading.Thread]],
worktree: _WorktreeReporter,
exc: BaseException, *, child: Any, task_index: int, goal: str, subagent_id: Optional[str], child_future: Any,
child_timeout: Optional[float], child_start: float, child_progress_cb: Any,
worker_thread_holder: Dict[str, Optional[threading.Thread]], worktree: _WorktreeReporter,
) -> _ChildFailure:
"""Build the error entry for a child whose Future timed out or raised.
@@ -644,17 +580,13 @@ def _handle_child_wait_failure(
_err = str(exc)
elif child_api_calls == 0:
_err = (
f"Subagent timed out after {child_timeout}s without "
f"making any API call — the child never reached its "
f"first LLM request (prompt construction, credential "
f"resolution, or transport may be stuck)."
f"Subagent timed out after {child_timeout}s without making any API call — the child never reached its "
f"first LLM request (prompt construction, credential resolution, or transport may be stuck)."
)
else:
_err = (
f"Subagent timed out after {child_timeout}s with "
f"{child_api_calls} API call(s) completed — likely "
f"stuck on a slow API call, tool call, or unresponsive "
f"network request."
f"Subagent timed out after {child_timeout}s with {child_api_calls} API call(s) completed — likely "
f"stuck on a slow API call, tool call, or unresponsive network request."
)
if is_timeout and diagnostic_path:
_err += f" Diagnostic: {diagnostic_path}"
@@ -706,7 +638,6 @@ def _validate_child_output_schema(
if not isinstance(_output_schema, dict):
return _SchemaOutcome(_output_schema, None, [], 0)
from tools.delegation_output_schema import build_retry_message, validate_output
_first_text = result.get("final_response") or ""
_schema_valid, _schema_errors = validate_output(_first_text, _output_schema)
if _schema_valid or not _first_text.strip() or result.get("interrupted", False):
@@ -750,8 +681,7 @@ def _build_tool_trace(messages: Any) -> list[Dict[str, Any]]:
fn = tc.get("function", {})
arguments = fn.get("arguments", "")
entry_t = {
"tool": fn.get("name", "unknown"),
"args_bytes": len(arguments),
"tool": fn.get("name", "unknown"), "args_bytes": len(arguments),
"input_summary": _summarize_tool_arguments(arguments),
}
tool_trace.append(entry_t)
@@ -769,11 +699,7 @@ def _build_tool_trace(messages: Any) -> list[Dict[str, Any]]:
return tool_trace
def _build_result_entry(
child: Any,
result: Dict[str, Any],
task_index: int,
duration: float,
schema: _SchemaOutcome,
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).
@@ -892,11 +818,7 @@ def _append_sibling_write_reminder(entry: Dict[str, Any], ws: _ChildWorkspace) -
entry["stale_paths"] = mod_paths
def _emit_child_complete(
child: Any,
result: Dict[str, Any],
entry: Dict[str, Any],
ws: _ChildWorkspace,
duration: float,
child: Any, result: Dict[str, Any], entry: Dict[str, Any], ws: _ChildWorkspace, duration: float,
child_progress_cb: Any,
) -> None:
"""Fire ``subagent.complete`` with the per-branch observability payload.
@@ -937,14 +859,8 @@ def _emit_child_complete(
_safe_progress(child_progress_cb, "subagent.complete", **complete_kwargs)
def _cleanup_child_run(
child: Any,
parent_agent: Any,
*,
subagent_id: Optional[str],
heartbeat: tuple,
child_pool: Any,
leased_cred_id: Any,
close_deferred: bool,
child: Any, parent_agent: Any, *, subagent_id: Optional[str], heartbeat: tuple, child_pool: Any,
leased_cred_id: Any, close_deferred: bool,
) -> None:
"""Finally-path teardown for one child run (idempotent, never raises).
@@ -969,7 +885,6 @@ def _cleanup_child_run(
# Restore the parent's tool names so the process-global is correct for
# any subsequent execute_code calls or other consumers.
import model_tools
saved_tool_names = getattr(child, "_delegate_saved_tool_names", None)
if isinstance(saved_tool_names, list):
model_tools._last_resolved_tool_names = list(saved_tool_names)
@@ -986,7 +901,6 @@ def _cleanup_child_run(
# 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)
child_session_id = str(getattr(child, "session_id", "") or "")
child_turn_is_active = relay_runtime.SESSION_COORDINATOR.has_active_turn(

View File

@@ -33,7 +33,6 @@ DEFAULT_CHILD_TIMEOUT: Optional[float] = None
def _cfg() -> dict:
"""The ``delegation`` section, read through the origin so tests can patch it."""
from tools.delegate_tool import _load_config
return _load_config()
@@ -47,8 +46,7 @@ def _cfg() -> dict:
def _subagent_auto_deny(command: str, description: str, **kwargs) -> str:
"""Auto-deny (safe default): returns 'deny' so the child sees a recoverable refusal."""
logger.warning(
"Subagent auto-denied dangerous command: %s (%s). "
"Set delegation.subagent_auto_approve: true to allow.",
"Subagent auto-denied dangerous command: %s (%s). Set delegation.subagent_auto_approve: true to allow.",
command, description,
)
return "deny"
@@ -102,8 +100,7 @@ def _get_max_concurrent_children() -> int:
_HIGH_CONCURRENCY_WARNED = True
logger.warning(
"delegation.max_concurrent_children=%d: each child consumes API tokens "
"independently. High values multiply cost linearly.",
result,
"independently. High values multiply cost linearly.", result,
)
return result
@@ -228,14 +225,11 @@ def _inherit_parent_base_url(parent_agent, fallback_base_url: Optional[str]) ->
def _loaded_pool(key: Any):
"""``load_pool(key)`` when it holds credentials, else None."""
from agent.credential_pool import load_pool
pool = load_pool(key)
return pool if pool is not None and pool.has_credentials() else None
def _resolve_child_credential_pool(
effective_provider: Optional[str],
parent_agent,
effective_base_url: Optional[str] = None,
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).
@@ -253,7 +247,6 @@ def _resolve_child_credential_pool(
if effective_provider == "custom":
try:
from agent.credential_pool import get_custom_provider_pool_key
child_key = get_custom_provider_pool_key(effective_base_url)
if child_key is None:
# Unregistered endpoint (no custom_providers entry): keep the
@@ -294,7 +287,6 @@ def _merge_request_overrides(runtime_overrides, explicit_overrides):
None when both sides 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
if not runtime_overrides and not explicit_overrides:
@@ -323,7 +315,6 @@ def _require_pinned_command(command: Optional[str], message: str) -> None:
if not command:
return
import shutil as _shutil
if not _shutil.which(command):
raise ValueError(message)
@@ -337,7 +328,6 @@ def _direct_endpoint_credentials(cfg_values: dict, explicit_request_overrides) -
# 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 = configured_base_url.lower()
host = base_url_hostname(configured_base_url)
provider = "custom"
@@ -360,7 +350,6 @@ def _direct_endpoint_credentials(cfg_values: dict, explicit_request_overrides) -
if configured_provider:
try:
from hermes_cli.runtime_provider import resolve_runtime_provider
runtime = resolve_runtime_provider(requested=configured_provider, target_model=configured_model)
request_overrides = dict(runtime.get("request_overrides") or {}) or None
max_output_tokens = runtime.get("max_output_tokens")
@@ -386,7 +375,6 @@ def _runtime_provider_credentials(cfg_values: dict, explicit_request_overrides)
configured_model, configured_provider = cfg_values["model"], cfg_values["provider"]
try:
from hermes_cli.runtime_provider import resolve_runtime_provider
runtime = resolve_runtime_provider(requested=configured_provider, target_model=configured_model)
except Exception as exc:
raise ValueError(
@@ -404,8 +392,7 @@ def _runtime_provider_credentials(cfg_values: dict, explicit_request_overrides)
)
pinned_command = runtime.get("command")
_require_pinned_command(
pinned_command,
f"Delegation provider '{configured_provider}' is pinned to the "
pinned_command, f"Delegation provider '{configured_provider}' is pinned to the "
f"'{pinned_command}' command, which was not found on PATH. "
f"Install it or choose a different delegation provider.",
)
@@ -415,9 +402,7 @@ def _runtime_provider_credentials(cfg_values: dict, explicit_request_overrides)
"base_url": runtime.get("base_url"),
"api_key": api_key,
"api_mode": runtime.get("api_mode"),
"request_overrides": _merge_request_overrides(
runtime.get("request_overrides"), explicit_request_overrides
)
"request_overrides": _merge_request_overrides(runtime.get("request_overrides"), explicit_request_overrides)
or {},
"max_output_tokens": runtime.get("max_output_tokens"),
"command": runtime.get("command"),
@@ -454,8 +439,7 @@ def _resolve_delegation_credentials(cfg: dict, parent_agent) -> dict:
"api_key": None,
"api_mode": None,
"request_overrides": _merge_request_overrides(
getattr(parent_agent, "request_overrides", None),
explicit_request_overrides,
getattr(parent_agent, "request_overrides", None), explicit_request_overrides,
),
"max_output_tokens": None,
}
@@ -473,7 +457,6 @@ def _load_config() -> dict:
if os.environ.get("HERMES_IGNORE_USER_CONFIG") != "1":
try:
from hermes_cli.config import load_config_readonly
cfg = load_config_readonly().get("delegation") or {}
if isinstance(cfg, dict):
return cfg
@@ -481,7 +464,6 @@ def _load_config() -> dict:
pass
try:
from cli import CLI_CONFIG
cfg = CLI_CONFIG.get("delegation") or {}
return cfg if isinstance(cfg, dict) else {}
except Exception:
@@ -492,29 +474,16 @@ def _load_config() -> dict:
# 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),
("provider_require_parameters", False),
("provider_data_collection", ""),
("providers_allowed", None), ("providers_ignored", None), ("providers_order", None), ("provider_sort", None),
("provider_require_parameters", False), ("provider_data_collection", ""),
)
_NOUS_PROVIDERS = frozenset({"nous", "nous-portal", "nousresearch"})
def _resolve_child_runtime(
parent_agent,
delegation_cfg: dict,
parent_api_key: Any,
*,
model: Optional[str],
override_provider: Optional[str],
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]],
parent_agent, delegation_cfg: dict, parent_api_key: Any, *, model: Optional[str], override_provider: Optional[str],
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.
@@ -537,7 +506,6 @@ def _resolve_child_runtime(
effective_api_mode = override_api_mode
elif (effective_provider or "").strip().lower() in _NOUS_PROVIDERS:
from hermes_cli.providers import nous_api_mode
effective_api_mode = nous_api_mode(effective_model)
elif effective_provider != _parent_provider:
effective_api_mode = None # force re-derivation from provider's defaults
@@ -546,10 +514,8 @@ def _resolve_child_runtime(
# A pinned transport that cannot run must fail the spawn loudly, never fall
# back silently (delegate_task pre-validates; this covers direct callers).
_require_pinned_command(
override_acp_command,
f"Pinned delegation command '{override_acp_command}' was not "
f"found on PATH. Install it or remove delegation.command from "
f"config.yaml.",
override_acp_command, f"Pinned delegation command '{override_acp_command}' was not "
f"found on PATH. Install it or remove delegation.command from config.yaml.",
)
effective_acp_command = override_acp_command or getattr(parent_agent, "acp_command", None)
effective_acp_args = list(
@@ -570,7 +536,6 @@ def _resolve_child_runtime(
delegation_effort = delegation_cfg.get("reasoning_effort")
if delegation_effort or delegation_effort is False:
from hermes_constants import parse_reasoning_effort
parsed = parse_reasoning_effort(delegation_effort)
if parsed is not None:
child_reasoning = parsed

View File

@@ -49,8 +49,7 @@ class _Batch:
def owner_kwargs(self) -> Dict[str, Any]:
"""Steer/stop authority of the originating session, passed to every child run."""
return {
"owner_session_id": self.origin_ui_session_id or None,
"owner_transport": self.origin_owner_transport,
"owner_session_id": self.origin_ui_session_id or None, "owner_transport": self.origin_owner_transport,
"owner_session_record": self.origin_owner_session_record,
}
@@ -68,12 +67,10 @@ def _capture_origin() -> tuple[str, str, Any, Any]:
ORIGINATING session, captured BEFORE building any child: AIAgent construction
clobbers the HERMES_SESSION_ID ContextVar/os.environ with the subagent's id."""
from tools.async_delegation import _current_origin_session_id
_origin_wake_sid = _current_origin_session_id()
_origin_ui_session_id = ""
with _quiet(None):
from gateway.session_context import get_session_env
_origin_ui_session_id = get_session_env("HERMES_UI_SESSION_ID", "")
transport, record = _capture_gateway_steer_authority(_origin_ui_session_id)
return _origin_wake_sid, _origin_ui_session_id, transport, record
@@ -120,7 +117,6 @@ def _run_children_parallel(batch: _Batch, results: list, *, honor_parent_interru
# 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]
completed_count = 0
@@ -185,7 +181,6 @@ def _execute_and_aggregate(batch: _Batch, *, honor_parent_interrupt: bool = True
"""
from tools.delegate_tool import _run_single_child
from tools.delegation_live_log import update_manifest_statuses
results: list = []
if len(batch.task_list) == 1:
_i, _t, child = batch.children[0]
@@ -209,15 +204,12 @@ _SYNC_FALLBACK_NOTES = {
"background=true is not available in this session — it cannot "
"receive a detached subagent result after the turn ends (a "
"one-shot runner such as `hermes -z`, a cron job, a Kanban "
"worker, or a stateless HTTP endpoint). The subagent(s) ran "
"SYNCHRONOUSLY and the result is included above."
"worker, or a stateless HTTP endpoint). The subagent(s) ran SYNCHRONOUSLY and the result is included above."
),
"at_capacity": (
"The background delegation pool was at capacity "
"(delegation.max_concurrent_children), so the subagent(s) ran "
"The background delegation pool was at capacity (delegation.max_concurrent_children), so the subagent(s) ran "
"SYNCHRONOUSLY and the result is included above. Raise "
"delegation.max_concurrent_children in config.yaml to allow "
"more concurrent background delegations."
"delegation.max_concurrent_children in config.yaml to allow more concurrent background delegations."
),
}
@@ -249,12 +241,9 @@ def _resolve_async_wake_sid(origin_wake_sid: str) -> Optional[str]:
return ""
if origin_wake_sid:
logger.info(
"delegate_task: async delivery unsupported on this "
"session, but a session id is bound (%s) — dispatching "
"in the background and waking the session via self-post "
"when it completes instead of forcing synchronous "
"execution.",
origin_wake_sid,
"delegate_task: async delivery unsupported on this session, but a session id is bound (%s) — dispatching "
"in the background and waking the session via self-post when it completes instead of forcing synchronous "
"execution.", origin_wake_sid,
)
return origin_wake_sid
return None
@@ -272,12 +261,10 @@ def _resolve_async_session_key(parent_agent: Any, origin_ui_session_id: str) ->
key would fail closed — stamp the parent's durable id.
"""
from tools.approval import get_current_session_key
session_key = get_current_session_key(default="")
agent_session_id = str(getattr(parent_agent, "session_id", "") or "")
with _quiet(None):
from gateway.session_context import get_session_env
source = get_session_env("HERMES_SESSION_SOURCE", "")
# Refresh from the task-local source when available, else retain the
# immutable value captured before child construction.
@@ -321,14 +308,12 @@ def _dispatched_payload(dispatch: dict, goals: List[str], child_agents: List[Any
"note": (
"Subagent is running in the background. You and the user can "
"keep working; its full result re-enters the conversation as a "
"new message when it finishes. Do not wait or poll — just "
"continue."
"new message when it finishes. Do not wait or poll — just continue."
if n == 1 else
f"{n} subagents are running in parallel in the background. You "
f"and the user can keep working; they wait on each other and "
f"their consolidated results re-enter the conversation as a "
f"single message once ALL of them finish. Do not wait or poll "
f"— just continue."
f"single message once ALL of them finish. Do not wait or poll — just continue."
),
}
sids = [getattr(c, "_subagent_id", None) for c in child_agents]
@@ -338,16 +323,14 @@ def _dispatched_payload(dispatch: dict, goals: List[str], child_agents: List[Any
"While a child runs you can orchestrate it live with this "
"same tool: delegate_task(action='list') to see live "
"children, action='steer' with subagent_id + message to "
"redirect one, action='stop' with subagent_id to end one "
"early."
"redirect one, action='stop' with subagent_id to end one early."
)
if live_paths:
payload["live_transcripts"] = list(live_paths)
payload["live_transcripts_hint"] = (
"Each subagent streams a human-readable transcript of its "
"operations to the file listed above (append-only, one per "
"task). Read or `tail -f` these paths at any time to watch "
"a child work while it runs."
"task). Read or `tail -f` these paths at any time to watch a child work while it runs."
)
return payload
@@ -363,7 +346,6 @@ def _dispatch_background(batch: _Batch) -> str:
"""
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)
if wake_sid is None:
logger.info("delegate_task: async delivery unsupported on this session runtime; running the batch synchronously instead.")

View File

@@ -52,10 +52,7 @@ def _clean_error_text(error: Any, max_chars: int = 200) -> str:
return line[: max_chars - 3] + "..." if len(line) > max_chars else line
def format_subagent_failure_line(
goal: Optional[str],
status: Optional[str],
error: Any = None,
duration_seconds: Any = None,
goal: Optional[str], status: Optional[str], error: Any = None, duration_seconds: Any = None,
) -> str:
"""One clean, human-readable line describing a failed subagent, rendered
directly to the user (CLI spinner echo, gateway platform notice), e.g.
@@ -129,13 +126,8 @@ def _normalize_event(event_type: Any) -> Any:
return None
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,
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.
@@ -162,7 +154,6 @@ def _build_child_system_prompt(
_ctx_files = ""
with _quiet("subagent: workspace context-files load failed", exc_info=True):
from agent.prompt_builder import build_context_files_prompt
_ctx_files = build_context_files_prompt(cwd=str(workspace_path), skip_soul=True)
if _ctx_files.strip():
parts.append(
@@ -192,8 +183,7 @@ def _build_child_system_prompt(
if child_depth + 1 >= max_spawn_depth
else "Your own children can themselves be orchestrators or leaves, "
"depending on the `role` you pass to delegate_task. Default is "
"'leaf'; pass role='orchestrator' explicitly when a child "
"needs to further decompose its work."
"'leaf'; pass role='orchestrator' explicitly when a child needs to further decompose its work."
)
parts.append(
"\n## Subagent Spawning (Orchestrator Role)\n"
@@ -221,12 +211,8 @@ def _resolve_workspace_hint(parent_agent) -> Optional[str]:
"""Best-effort local workspace hint for child prompts: only a concrete
absolute directory is ever injected (never a fake container path)."""
candidates = [
os.getenv("TERMINAL_CWD"),
getattr(
getattr(parent_agent, "_subdirectory_hints", None), "working_dir", None
),
getattr(parent_agent, "terminal_cwd", None),
getattr(parent_agent, "cwd", None),
os.getenv("TERMINAL_CWD"), getattr(getattr(parent_agent, "_subdirectory_hints", None), "working_dir", None),
getattr(parent_agent, "terminal_cwd", None), getattr(parent_agent, "cwd", None),
]
for candidate in candidates:
if not candidate:
@@ -367,8 +353,7 @@ class _ChildProgressRelay:
# sees WHY, not just a vanished branch (gateway renders off the relayed event).
if kwargs.get("status") in SUBAGENT_FAILURE_STATUSES:
self._tree_line(format_subagent_failure_line(
self.goal_label, kwargs.get("status"),
error=kwargs.get("summary") or preview,
self.goal_label, kwargs.get("status"), error=kwargs.get("summary") or preview,
duration_seconds=kwargs.get("duration_seconds"),
))
self._relay("subagent.complete", preview=preview, **kwargs)
@@ -405,7 +390,6 @@ class _ChildProgressRelay:
rec["last_tool"] = tool_name or ""
if self.spinner:
from agent.display import get_tool_emoji
line = f"{get_tool_emoji(tool_name or '')} {tool_name}"
short = _short(preview, 35) if preview else ""
self._tree_line(f'{line} "{short}"' if short else line)
@@ -424,17 +408,9 @@ class _ChildProgressRelay:
getattr(self, method)(tool_name, preview, args, kwargs)
def _build_child_progress_callback(
task_index: int,
goal: str,
parent_agent,
task_count: int = 1,
*,
subagent_id: Optional[str] = None,
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,
task_index: int, goal: str, parent_agent, task_count: int = 1, *, subagent_id: Optional[str] = None,
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
@@ -444,6 +420,5 @@ def _build_child_progress_callback(
if not spinner and not parent_cb:
return None
return _ChildProgressRelay(
task_index, goal, spinner, parent_cb, task_count,
subagent_id, parent_id, depth, model, toolsets, session_ref,
task_index, goal, spinner, parent_cb, task_count, subagent_id, parent_id, depth, model, toolsets, session_ref,
)

View File

@@ -82,8 +82,7 @@ def _retain_recent_subagent(record: Dict[str, Any]) -> None:
if not sid:
return
_recent_subagents[sid] = {
"goal": record.get("goal"),
"delegation_id": record.get("delegation_id"),
"goal": record.get("goal"), "delegation_id": record.get("delegation_id"),
"owner_agent_session_id": record.get("owner_agent_session_id"),
}
while len(_recent_subagents) > _RECENT_SUBAGENTS_CAP:
@@ -143,11 +142,7 @@ def interrupt_subagent(subagent_id: str) -> bool:
return True
def steer_subagent(
subagent_id: str,
text: str,
*,
owner_session_id: Optional[str] = None,
owner_transport: Any = None,
subagent_id: str, text: str, *, owner_session_id: Optional[str] = None, owner_transport: Any = None,
owner_session_record: Any = None,
) -> bool:
"""Queue steering text into a running subagent without stopping it.
@@ -194,7 +189,6 @@ def _capture_gateway_steer_authority(owner_session_id: Optional[str]) -> tuple[A
return None, None
try:
from tui_gateway.server import _current_session_steer_authority
return _current_session_steer_authority(owner_session_id)
except Exception:
return None, None
@@ -281,8 +275,7 @@ def _owns_subagent_record(record: Dict[str, Any], parent_agent: Any) -> bool:
return True
# Compression rotation on either side: compare lineage tips.
return _resolve_session_lineage(owner_sid, parent_agent) in {
parent_sid,
_resolve_session_lineage(parent_sid, parent_agent),
parent_sid, _resolve_session_lineage(parent_sid, parent_agent),
}
def _handle_control_action(action: str, subagent_id: Optional[str], message: Optional[str], parent_agent: Any) -> str:
@@ -308,11 +301,7 @@ def _handle_control_action(action: str, subagent_id: Optional[str], message: Opt
"goal": r.get("goal"),
"model": r.get("model"),
"status": r.get("status"),
"running_seconds": (
round(time.time() - started, 1)
if isinstance(started, (int, float))
else None
),
"running_seconds": (round(time.time() - started, 1) if isinstance(started, (int, float)) else None),
"accepting_steer": bool(r.get("accepting_steer", False)),
"live_transcript": getattr(agent, "_live_transcript_path", None),
}
@@ -357,21 +346,16 @@ def _handle_control_action(action: str, subagent_id: Optional[str], message: Opt
_CONTROL_OUTCOMES = {
"stop": (
"interrupt_requested",
"The subagent stops at its next iteration boundary "
"(in-flight tool calls are asked to cancel). Its "
"partial result still re-enters the conversation as a "
"completion message — do not wait or poll.",
"The subagent stops at its next iteration boundary (in-flight tool calls are asked to cancel). Its "
"partial result still re-enters the conversation as a completion message — do not wait or poll.",
"Could not interrupt '{sid}' — it likely finished in the last "
"moment. Its result arrives as a normal completion message.",
),
"steer": (
"queued",
"Steering text queued. The subagent sees it appended "
"to its next tool result — the current tool call is "
"never cut. If the child finishes before a delivery "
"boundary remains, the text is reported back as "
"missed_steer in its completion entry.",
"Subagent '{sid}' is no longer accepting steering (finishing or "
"Steering text queued. The subagent sees it appended to its next tool result — the current tool call is "
"never cut. If the child finishes before a delivery boundary remains, the text is reported back as "
"missed_steer in its completion entry.", "Subagent '{sid}' is no longer accepting steering (finishing or "
"already finished). Its result arrives as a normal completion "
"message; re-delegate a follow-up task if more work is needed.",
),

View File

@@ -15,10 +15,7 @@ from urllib.parse import urlsplit, urlunsplit
logger = logging.getLogger("tools.delegate_tool")
def _extract_output_tail(
result: Dict[str, Any],
*,
max_entries: int = 12,
max_chars: int = 8000,
result: Dict[str, Any], *, max_entries: int = 12, max_chars: int = 8000,
) -> List[Dict[str, Any]]:
"""Pull the last N tool-call results from a child's conversation.
@@ -260,13 +257,11 @@ def _spill_summary_to_file(task_index: int, summary: str) -> Optional[str]:
try:
from hermes_constants import get_hermes_dir
import datetime as _dt
cache_dir = get_hermes_dir("cache/delegation", "delegation_cache")
cache_dir.mkdir(parents=True, exist_ok=True)
ts = _dt.datetime.now().strftime("%Y%m%d_%H%M%S_%f")
path = cache_dir / f"subagent-summary-{task_index}-{ts}.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.
@@ -304,8 +299,7 @@ def _trim_summary_with_footer(summary: str, cap: int, task_index: int) -> tuple[
spill_path = _spill_summary_to_file(task_index, summary)
footer_lines = [
"",
"─" * 8 + " [SUMMARY TRUNCATED] " + "─" * 8,
"", "─" * 8 + " [SUMMARY TRUNCATED] " + "─" * 8,
f"Showing {len(head):,} chars (head) + {len(tail):,} chars (tail) "
f"of {original_len:,} total — trimmed to protect the parent's context window.",
]
@@ -319,10 +313,7 @@ def _trim_summary_with_footer(summary: str, cap: int, task_index: int) -> tuple[
f"summary; raise/lower offset to page through it)."
)
else:
footer_lines.append(
"Full output could not be stored to disk; the head+tail above is "
"all that was preserved."
)
footer_lines.append("Full output could not be stored to disk; the head+tail above is all that was preserved.")
footer_lines.append("─" * 37)
model_text = head + "\n\n[... middle omitted — see footer ...]\n\n" + tail + "\n".join(footer_lines)
@@ -400,10 +391,7 @@ def _apply_summary_budget(results: List[Dict[str, Any]], parent_agent) -> None:
if spill_path:
entry["summary_full_path"] = spill_path
logger.debug(
"[subagent-%s] summary trimmed %d → ~%d chars (spill=%s)",
entry.get("task_index", "?"),
original_len,
cap,
"[subagent-%s] summary trimmed %d → ~%d chars (spill=%s)", entry.get("task_index", "?"), original_len, cap,
spill_path or "none",
)
@@ -417,7 +405,6 @@ def _build_child_preserving_parent_tools(**kwargs):
"""Build a child without leaking its resolved toolset into the parent."""
from tools.delegate_tool import _build_child_agent
import model_tools
with _CHILD_CONSTRUCTION_LOCK:
parent_tool_names = list(model_tools._last_resolved_tool_names)
try:
@@ -453,8 +440,7 @@ def _notify_memory_manager(results, task_list, child_by_index, parent_agent) ->
task_index = entry.get("task_index", -1)
in_range = isinstance(task_index, int) and 0 <= task_index < len(task_list)
memory.on_delegation(
task=task_list[task_index]["goal"] if in_range else "",
result=entry.get("summary", "") or "",
task=task_list[task_index]["goal"] if in_range else "", result=entry.get("summary", "") or "",
child_session_id=getattr(child_by_index.get(task_index), "session_id", ""),
)
except Exception:
@@ -482,13 +468,10 @@ def _fire_subagent_stop_hooks(results, child_by_index, parent_agent) -> float:
try:
child = child_by_index.get(entry.get("task_index", -1))
invoke_hook(
"subagent_stop",
parent_session_id=getattr(parent_agent, "session_id", None),
"subagent_stop", parent_session_id=getattr(parent_agent, "session_id", None),
parent_turn_id=getattr(parent_agent, "_current_turn_id", "") or "",
child_session_id=getattr(child, "session_id", None),
child_role=child_role,
child_summary=entry.get("summary"),
child_status=entry.get("status"),
child_session_id=getattr(child, "session_id", None), child_role=child_role,
child_summary=entry.get("summary"), child_status=entry.get("status"),
tool_call_history=_subagent_stop_tool_call_history(entry.get("tool_trace")),
duration_ms=int((entry.get("duration_seconds") or 0) * 1000),
)
@@ -512,9 +495,7 @@ def _rollup_children_cost(parent_agent, children_cost_total: float) -> None:
logger.debug("Subagent cost rollup failed", exc_info=True)
def _finalize_child_results(
results: List[Dict[str, Any]],
task_list: List[Dict[str, Any]],
children: List[tuple[int, Dict[str, Any], Any]],
results: List[Dict[str, Any]], task_list: List[Dict[str, Any]], children: List[tuple[int, Dict[str, Any], Any]],
parent_agent,
) -> None:
"""Apply host-owned summary, memory, hook, and cost contracts once."""

View File

@@ -36,8 +36,7 @@ def _recover_tasks_from_json_string(tasks: Any) -> tuple[Optional[List[Dict[str,
_PLACEHOLDER_GOAL_RE = re.compile(r"^(todo|task\s*\d+)$", re.IGNORECASE)
_TEMPLATE_MARKER_RE = re.compile(
r"<[A-Za-z][A-Za-z0-9]*(?:[ _-][A-Za-z0-9]+)+>"
r"|\{[A-Za-z][A-Za-z0-9]*(?:[ _-][A-Za-z0-9]+)+\}"
r"<[A-Za-z][A-Za-z0-9]*(?:[ _-][A-Za-z0-9]+)+>|\{[A-Za-z][A-Za-z0-9]*(?:[ _-][A-Za-z0-9]+)+\}"
)
_MIN_BATCH_GOAL_LEN = 10
@@ -58,16 +57,14 @@ def _validate_batch_tasks(task_list: List[Dict[str, Any]]) -> Optional[str]:
if _PLACEHOLDER_GOAL_RE.match(normalized):
return (
f"Task {i} has a placeholder goal ({goal!r}). Replace it "
"with a specific, self-contained description of what the "
"subagent should accomplish."
"with a specific, self-contained description of what the subagent should accomplish."
)
marker = _TEMPLATE_MARKER_RE.search(goal)
if marker:
return (
f"Task {i} goal contains an unexpanded template marker "
f"({marker.group(0)!r}). Substitute the real value before "
"calling delegate_task — subagents cannot resolve "
"placeholders."
"calling delegate_task — subagents cannot resolve placeholders."
)
if len(goal) < _MIN_BATCH_GOAL_LEN and len(task_list) >= 2:
# Multi-task fan-outs with terse goals are usually unexpanded
@@ -98,10 +95,8 @@ def _normalize_task_list(
if tasks and isinstance(tasks, list):
if len(tasks) > max_children:
return None, (
f"Too many tasks: {len(tasks)} provided, but "
f"max_concurrent_children is {max_children}. "
f"Either reduce the task count, split into multiple "
f"delegate_task calls, or increase "
f"Too many tasks: {len(tasks)} provided, but max_concurrent_children is {max_children}. "
f"Either reduce the task count, split into multiple delegate_task calls, or increase "
f"delegation.max_concurrent_children in config.yaml."
)
task_list = tasks
@@ -113,8 +108,7 @@ def _normalize_task_list(
else:
return None, (
"No tasks provided. Pass tasks=[{goal: '...', context: '...'}, "
"...] — one entry per subagent (a single task is a one-entry "
"array)."
"...] — one entry per subagent (a single task is a one-entry array)."
)
for i, task in enumerate(task_list):
@@ -138,7 +132,6 @@ def _coerce_task_schemas(
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):
raw_schema = task.get("output_schema")

View File

@@ -33,7 +33,6 @@ def _is_mcp_toolset_name(name: str) -> bool:
return True
try:
from tools.registry import registry
target = registry.get_toolset_alias_target(str(name))
except Exception:
target = None
@@ -118,7 +117,6 @@ def _resolve_child_toolsets(
parent_toolsets = set(parent_enabled)
elif parent_agent and hasattr(parent_agent, "valid_tool_names"):
import model_tools
parent_toolsets = {
ts
for name in parent_agent.valid_tool_names