refactor(delegate): child_run best-effort blocks via _quiet; shared failure-entry tail

This commit is contained in:
Teknium
2026-09-02 17:54:08 -07:00
parent 64957b7e02
commit 901ba7db6f
2 changed files with 70 additions and 130 deletions

View File

@@ -412,16 +412,12 @@ def _run_single_child(
return entry
except Exception as exc:
# Close steer acceptance before any completion callback (see _merge_late_steer).
_late_pending_steer = (_close_subagent_steering(_subagent_id, child) if _subagent_id else None)
duration = round(time.monotonic() - child_start, 2)
logging.exception(f"[subagent-{task_index}] failed")
_safe_progress(
child_progress_cb,
"subagent.complete",
preview=str(exc),
status="failed",
duration_seconds=duration,
summary=str(exc),
child_progress_cb, "subagent.complete", preview=str(exc), status="failed", duration_seconds=duration, summary=str(exc),
)
_error_entry = _fabricated_entry(task_index, "error", str(exc), child, duration)
_append_missed_steer(_error_entry, _late_pending_steer)

View File

@@ -16,7 +16,7 @@ from typing import Any, Dict, List, Optional
from agent.interrupt_compat import request_hard_interrupt
from dataclasses import dataclass
from tools import file_state
from tools.delegate_tool_progress import _safe_progress
from tools.delegate_tool_progress import _quiet, _safe_progress
from tools.delegate_tool_registry import (
_capture_gateway_steer_authority, _close_subagent_steering, _register_subagent, _unregister_subagent,
)
@@ -31,6 +31,9 @@ def _num(value: Any, default: int = 0) -> int:
"""int() for counters that may be mocks/None on test doubles."""
return int(value) if isinstance(value, (int, float)) else default
def _str_or_none(value: Any) -> Optional[str]:
return value if isinstance(value, str) else None
def _fabricated_entry(idx: int, status: str, error: str, child: Any, duration: float = 0) -> Dict[str, Any]:
"""Result entry for a child that raised, never finished, or was abandoned."""
return {
@@ -51,54 +54,44 @@ def _append_missed_steer(entry: Dict[str, Any], late_steer: Optional[str]) -> No
def _close_child(child: Any, log_message: str) -> None:
"""Best-effort ``child.close()`` (tool sandboxes, browser daemons, httpx clients)."""
try:
with _quiet(log_message, exc_info=True):
close = getattr(child, "close", None)
if callable(close):
close()
except Exception:
logger.debug(log_message, exc_info=True)
def _attach_child(parent_agent: Any, child: Any) -> None:
"""Register the child for parent interrupt propagation."""
if not hasattr(parent_agent, "_active_children"):
return
def _with_children_lock(parent_agent: Any, op: str, child: Any) -> None:
"""``parent_agent._active_children.<op>(child)`` under the parent's lock when it has one."""
lock = getattr(parent_agent, "_active_children_lock", None)
if lock:
with lock:
parent_agent._active_children.append(child)
getattr(parent_agent._active_children, op)(child)
else:
parent_agent._active_children.append(child)
getattr(parent_agent._active_children, op)(child)
def _attach_child(parent_agent: Any, child: Any) -> None:
"""Register the child for parent interrupt propagation."""
if hasattr(parent_agent, "_active_children"):
_with_children_lock(parent_agent, "append", child)
def _detach_child(parent_agent: Any, child: Any) -> None:
"""Remove the child from parent interrupt propagation (no-op if absent)."""
if not hasattr(parent_agent, "_active_children"):
return
try:
lock = getattr(parent_agent, "_active_children_lock", None)
if lock:
with lock:
parent_agent._active_children.remove(child)
else:
parent_agent._active_children.remove(child)
_with_children_lock(parent_agent, "remove", child)
except (ValueError, UnboundLocalError) as e:
logger.debug("Could not remove child from active_children: %s", e)
def _signal_child_stop(child: Any, *reason: str) -> None:
"""Cooperative interrupt so the child's worker thread can exit cleanly."""
try:
with _quiet(None):
if child is not None and not request_hard_interrupt(child, *reason) and hasattr(child, "_interrupt_requested"):
child._interrupt_requested = True
except Exception:
pass
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")
]
return [f"{indent}{sub}" for frame_line in _traceback.format_stack(frame) for sub in frame_line.rstrip().split("\n")]
_DIAG_CHILD_ATTRS = (
"model", "provider", "api_mode", "base_url", "max_iterations",
@@ -128,16 +121,12 @@ 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
import threading as _threading
lines = ["## Worker thread stack at timeout"]
frames = _sys._current_frames()
if worker_thread is not None and worker_thread.is_alive():
worker_frame = frames.get(worker_thread.ident)
lines.extend(
_format_thread_stack(worker_frame, " ") if worker_frame is not None
else [" <worker frame not available>"]
)
lines.extend(_format_thread_stack(worker_frame, " ") if worker_frame is not None else [" <worker frame not available>"])
elif worker_thread is None:
lines.append(" <no worker thread handle>")
else:
@@ -145,7 +134,7 @@ def _diag_threads(worker_thread: Optional[threading.Thread]) -> List[str]:
lines += ["", "## All thread stacks at timeout"]
try:
frames = _sys._current_frames()
by_ident = {th.ident: th for th in _threading.enumerate() if th.ident}
by_ident = {th.ident: th for th in threading.enumerate() if th.ident}
worker_ident = worker_thread.ident if worker_thread else None
dumped = 0
for ident, frame in frames.items():
@@ -217,10 +206,8 @@ def _dump_subagent_timeout_diagnostic(
tool_names = getattr(child, "valid_tool_names", None)
if tool_names:
lines.append(f" loaded tool count: {len(tool_names)}")
try:
with _quiet(None):
lines.append(f" loaded tools: {sorted(tool_names)}")
except Exception:
pass
lines += [""] + _diag_sizes(child) + ["", "## Activity summary"]
try:
lines += [f" {k}: {v!r}" for k, v in child.get_activity_summary().items()]
@@ -282,11 +269,8 @@ def _start_heartbeat(child: Any, parent_agent: Any, task_index: int) -> tuple:
stale_limit = _HEARTBEAT_STALE_CYCLES_IN_TOOL if child_tool else _HEARTBEAT_STALE_CYCLES_IDLE
if last_seen["stale"] >= stale_limit:
logger.warning(
"Subagent %d appears stale (no progress for %d "
"heartbeat cycles, tool=%s) — stopping heartbeat",
task_index,
last_seen["stale"],
child_tool or "<none>",
"Subagent %d appears stale (no progress for %d heartbeat cycles, tool=%s) — stopping heartbeat",
task_index, last_seen["stale"], child_tool or "<none>",
)
break # stop touching parent, let gateway timeout fire
@@ -299,10 +283,8 @@ def _start_heartbeat(child: Any, parent_agent: Any, task_index: int) -> tuple:
)
except Exception:
pass
try:
with _quiet(None):
touch(desc)
except Exception:
pass
return _heartbeat_stop, threading.Thread(target=_heartbeat_loop, daemon=True)
@@ -324,26 +306,21 @@ def _register_child(
if not isinstance(_subagent_id, str) or not _subagent_id:
return None
if owner_session_id is None:
try:
with _quiet(None):
from gateway.session_context import get_session_env
owner_session_id = get_session_env("HERMES_UI_SESSION_ID", "") or None
except Exception:
owner_session_id = 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)
_raw_depth = getattr(child, "_delegate_depth", 1)
_parent_sid = getattr(child, "_parent_subagent_id", None)
_delegation_id = getattr(child, "_delegation_id", None)
_model = getattr(child, "model", None)
_register_subagent(
{
"subagent_id": _subagent_id,
"parent_id": _parent_sid if isinstance(_parent_sid, str) else None,
"parent_id": _str_or_none(getattr(child, "_parent_subagent_id", None)),
"depth": max(0, _raw_depth - 1) if isinstance(_raw_depth, int) else 0,
"goal": goal,
"delegation_id": _delegation_id if isinstance(_delegation_id, str) else None,
"model": _model if isinstance(_model, str) else None,
"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,
@@ -416,12 +393,10 @@ def _create_isolated_worktree(parent_agent: Any, parent_task_id: Any, subagent_i
logger.debug("worktree isolation skipped: non-local terminal backend")
return None
_parent_cwd = None
try:
with _quiet(None):
from tools.terminal_tool import get_session_cwd as _gsc
_parent_cwd = _gsc(parent_task_id)
except Exception:
pass
return subagent_worktree.create_subagent_worktree(
_parent_cwd or _resolve_workspace_hint(parent_agent), subagent_id=subagent_id,
)
@@ -449,22 +424,18 @@ def _seed_child_workspace(
# Seed the child's cwd record from the parent's: same starting directory,
# but the child's later `cd`s stay in its own record. Per-session container
# isolation keys containers by task_id; the child must share the PARENT's.
try:
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)
except Exception as e:
logger.debug("Child cwd seed failed: %s", e)
_worktree_info = _create_isolated_worktree(parent_agent, parent_task_id, subagent_id)
if _worktree_info is not None:
try:
with _quiet("worktree cwd seed failed: %s"):
from tools.terminal_tool import record_session_cwd as _rsc
_rsc(child_task_id, _worktree_info["path"])
except Exception as e:
logger.debug("worktree cwd seed failed: %s", e)
# 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
@@ -493,18 +464,14 @@ def _defer_close_after_timeout(child: Any, child_future: Any) -> None:
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")
)
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)
if not callable(_drain):
return
def _drain_once(phase: str) -> None:
try:
with _quiet("Timed-out child transport drain (%s) failed", phase, exc_info=True):
_drain(reason=f"delegate_timeout_{phase}")
except Exception:
logger.debug("Timed-out child transport drain (%s) failed", phase, exc_info=True)
_drain_once("immediate")
@@ -523,12 +490,10 @@ def _lease_child_credential(child: Any) -> tuple[Any, Optional[str]]:
return None, None
leased_cred_id = child_pool.acquire_lease()
if leased_cred_id is not None:
try:
with _quiet("Failed to bind child to leased credential: %s"):
leased_entry = child_pool.current()
if leased_entry is not None and hasattr(child, "_swap_credential"):
child._swap_credential(leased_entry)
except Exception as exc:
logger.debug("Failed to bind child to leased credential: %s", exc)
return child_pool, leased_cred_id
def _make_text_relay(child_progress_cb: Any):
@@ -612,6 +577,19 @@ def _merge_late_steer(result: Dict[str, Any], subagent_id: Optional[str], child:
existing = result.get("pending_steer")
result["pending_steer"] = f"{existing}\n{late}" if isinstance(existing, str) and existing else late
def _finish_failed_entry(
entry: Dict[str, Any], late_steer: Optional[str], child_progress_cb: Any, worktree: _WorktreeReporter, *, preview: str,
) -> Dict[str, Any]:
"""Shared tail of every child failure path: emit ``subagent.complete``, note
the steer text that won the race with the failure, report the worktree."""
_safe_progress(
child_progress_cb, "subagent.complete", preview=preview, status=entry["status"],
duration_seconds=entry["duration_seconds"], summary=entry["summary"] or "",
)
_append_missed_steer(entry, late_steer)
worktree.attach(entry) # no-op when isolation never engaged
return entry
def _handle_child_wait_failure(
exc: BaseException,
*,
@@ -633,28 +611,20 @@ def _handle_child_wait_failure(
``close_deferred=True``) because closing from this thread races the
still-unwinding worker's finally path.
"""
# No consumer boundary remains once this owner stops waiting for the
# child. Close acceptance before any completion callback and retain steer
# text that won the race with this failure/timeout.
# Steer acceptance must close BEFORE the stop signal so a concurrent steer
# is either drained into the entry or rejected — never silently lost.
_late_pending_steer = _close_subagent_steering(subagent_id, child) if subagent_id else None
_signal_child_stop(child)
is_timeout = isinstance(exc, (FuturesTimeoutError, TimeoutError))
duration = round(time.monotonic() - child_start, 2)
logger.warning(
"Subagent %d %s after %.1fs",
task_index,
"timed out" if is_timeout else f"raised {type(exc).__name__}",
duration,
)
logger.warning("Subagent %d %s after %.1fs", task_index, "timed out" if is_timeout else f"raised {type(exc).__name__}", duration)
# A timeout BEFORE any API call is a black box without a diagnostic dump.
diagnostic_path: Optional[str] = None
child_api_calls = 0
try:
with _quiet(None):
child_api_calls = int(child.get_activity_summary().get("api_call_count", 0) or 0)
except Exception:
pass
if is_timeout and child_api_calls == 0:
diagnostic_path = _dump_subagent_timeout_diagnostic(
child=child,
@@ -670,15 +640,6 @@ def _handle_child_wait_failure(
logger.warning("Subagent %d 0-API-call timeout — diagnostic written to %s", task_index, diagnostic_path)
status = "timeout" if is_timeout else "error"
_safe_progress(
child_progress_cb,
"subagent.complete",
preview=f"Timed out after {duration}s" if is_timeout else str(exc),
status=status,
duration_seconds=duration,
summary="",
)
if not is_timeout:
_err = str(exc)
elif child_api_calls == 0:
@@ -716,8 +677,10 @@ def _handle_child_wait_failure(
"_child_role": getattr(child, "_delegate_role", None),
"diagnostic_path": diagnostic_path,
}
_append_missed_steer(_error_entry, _late_pending_steer)
worktree.attach(_error_entry)
_finish_failed_entry(
_error_entry, _late_pending_steer, child_progress_cb, worktree,
preview=f"Timed out after {duration}s" if is_timeout else str(exc),
)
close_deferred = is_timeout and not child_future.done()
if close_deferred:
_defer_close_after_timeout(child, child_future)
@@ -754,9 +717,7 @@ def _validate_child_output_schema(
_retry_result = None
try:
_retry_result = child.run_conversation(
user_message=build_retry_message(_schema_errors),
task_id=child_task_id,
stream_callback=relay_child_text,
user_message=build_retry_message(_schema_errors), task_id=child_task_id, stream_callback=relay_child_text,
)
except Exception as _retry_exc:
logger.warning("Subagent %d schema-retry turn failed: %s", task_index, _retry_exc)
@@ -798,10 +759,7 @@ def _build_tool_trace(messages: Any) -> list[Dict[str, Any]]:
trace_by_id[tc["id"]] = entry_t
elif msg.get("role") == "tool":
content = _stringify_tool_content(msg.get("content", ""))
result_meta = {
"result_bytes": len(content),
"status": "error" if _looks_like_error_output(content) else "ok",
}
result_meta = {"result_bytes": len(content), "status": "error" if _looks_like_error_output(content) else "ok"}
tc_id = msg.get("tool_call_id")
target = trace_by_id.get(tc_id) if tc_id else None
if target is not None:
@@ -828,8 +786,7 @@ def _build_result_entry(
structured_failure = bool(result.get("failed") or result.get("error"))
# "(empty)" is run_agent's give-up sentinel after repeated empty LLM
# responses (usually a transport bug) — a failure, not a success.
_empty_sentinel = summary.strip() == "(empty)"
usable_summary = bool(summary) and not _empty_sentinel
usable_summary = bool(summary) and summary.strip() != "(empty)"
if interrupted:
status, exit_reason = "interrupted", "interrupted"
@@ -847,7 +804,6 @@ def _build_result_entry(
# (orchestrators reading only status/icon would accept an empty verdict).
status = "completed" if schema.valid is not False and usable_summary else "failed"
_model = getattr(child, "model", None)
_cost = getattr(child, "session_estimated_cost_usd", 0.0)
_cost_status = getattr(child, "session_cost_status", None)
# Result entry contract: see the _run_single_child docstring.
@@ -857,7 +813,7 @@ def _build_result_entry(
"summary": summary,
"api_calls": result.get("api_calls", 0),
"duration_seconds": duration,
"model": _model if isinstance(_model, str) else None,
"model": _str_or_none(getattr(child, "model", None)),
"exit_reason": exit_reason,
# A budget-exhausted child still returns a summary (status stays
# "completed"), so the parent needs this explicit flag.
@@ -882,8 +838,7 @@ def _build_result_entry(
# 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 ".")
"Final answer does not satisfy the declared output_schema" + (" (after 1 retry)." if schema.retries else ".")
)
else:
entry["error"] = result.get("error", "Subagent did not produce a response.")
@@ -919,7 +874,7 @@ def _append_sibling_write_reminder(entry: Dict[str, Any], ws: _ChildWorkspace) -
"""
if not (ws.parent_task_id and ws.parent_reads_snapshot):
return
try:
with _quiet("file_state sibling-write check failed", exc_info=True):
sibling_writes = file_state.writes_since(ws.parent_task_id, ws.wall_start, ws.parent_reads_snapshot)
mod_paths = sorted({p for paths in sibling_writes.values() for p in paths}) if sibling_writes else []
if not mod_paths:
@@ -935,8 +890,6 @@ def _append_sibling_write_reminder(entry: Dict[str, Any], ws: _ChildWorkspace) -
entry["summary"] = entry["summary"] + reminder
else:
entry["stale_paths"] = mod_paths
except Exception:
logger.debug("file_state sibling-write check failed", exc_info=True)
def _emit_child_complete(
child: Any,
@@ -954,17 +907,13 @@ def _emit_child_complete(
if not child_progress_cb:
return
summary, status = entry["summary"], entry["status"]
try:
_files_read: list = []
with _quiet(None):
_files_read = list(file_state.known_reads(ws.child_task_id))[:40]
except Exception:
_files_read = []
try:
_files_written_map: dict = {}
with _quiet(None):
_files_written_map = file_state.writes_since("", ws.wall_start, []) # all writes since wall_start
except Exception:
_files_written_map = {}
_files_written = sorted(
{p for tid, paths in _files_written_map.items() if tid == ws.child_task_id for p in paths}
)[:40]
_files_written = sorted({p for tid, paths in _files_written_map.items() if tid == ws.child_task_id for p in paths})[:40]
complete_kwargs: Dict[str, Any] = {
"preview": summary[:160] if summary else entry.get("error", ""),
@@ -1014,10 +963,8 @@ def _cleanup_child_run(
_unregister_subagent(subagent_id, agent=child)
if child_pool is not None and leased_cred_id is not None:
try:
with _quiet("Failed to release credential lease: %s"):
child_pool.release_lease(leased_cred_id)
except Exception as exc:
logger.debug("Failed to release credential lease: %s", exc)
# Restore the parent's tool names so the process-global is correct for
# any subsequent execute_code calls or other consumers.
@@ -1037,16 +984,13 @@ def _cleanup_child_run(
# 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.
try:
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(
profile_key=relay_runtime.current_profile_key(),
session_id=child_session_id,
profile_key=relay_runtime.current_profile_key(), session_id=child_session_id,
)
if runtime is not None and child_session_id and not child_turn_is_active:
runtime.unregister_subagent({"child_session_id": child_session_id})
except Exception:
logger.debug("Failed to close child Relay session after delegation")