diff --git a/tools/delegate_tool_registry.py b/tools/delegate_tool_registry.py index 8b8aa77d95..cf9a99f4c8 100644 --- a/tools/delegate_tool_registry.py +++ b/tools/delegate_tool_registry.py @@ -17,34 +17,24 @@ from tools.registry import tool_error logger = logging.getLogger("tools.delegate_tool") _spawn_pause_lock = threading.Lock() - _spawn_paused: bool = False - _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, parent_session_id} retained AFTER the -# child finishes (bounded FIFO). Child-started background processes routinely -# outlive the child itself (its npm ci with notify_on_complete=true finishes -# after the child's summary was delivered); their completion notifications -# reach the parent conversation 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]]: - """Resolve a process task_id to its originating delegation, if any. - - Children run their terminal sessions under ``task_id == subagent_id`` - (see _run_single_child's child_task_id), so a background process spawned - by a subagent carries that id in ``ProcessSession.task_id``. Returns - ``{subagent_id, goal, delegation_id}`` for live AND recently-finished - children, or None when the task_id is not a known subagent. - """ + """``{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: @@ -54,11 +44,7 @@ def get_subagent_attribution(task_id: Optional[str]) -> Optional[Dict[str, Any]] return {"subagent_id": task_id, "goal": record.get("goal"), "delegation_id": record.get("delegation_id")} def set_spawn_paused(paused: bool) -> bool: - """Globally block/unblock new delegate_task spawns. - - Active children keep running; only NEW calls to delegate_task fail fast - with a "spawning paused" error until unblocked. Returns the new state. - """ + """Globally block/unblock NEW delegate_task spawns (active children keep running). Returns the new state.""" global _spawn_paused with _spawn_pause_lock: _spawn_paused = bool(paused) @@ -76,24 +62,21 @@ def _register_subagent(record: Dict[str, Any]) -> None: with _active_subagents_lock: _active_subagents[sid] = record -def _retain_recent_subagent(record: Dict[str, Any]) -> None: - """Keep a bounded attribution stub after a child finishes (lock held).""" - sid = record.get("subagent_id") - if not sid: - return - _recent_subagents[sid] = { - "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: - _recent_subagents.pop(next(iter(_recent_subagents)), None) - def _unregister_subagent(subagent_id: str, *, agent: Any = None) -> None: + """Drop the live record (exact agent identity when given) and keep a bounded attribution stub.""" with _active_subagents_lock: record = _active_subagents.get(subagent_id) - if record is not None and (agent is None or record.get("agent") is agent): - _active_subagents.pop(subagent_id, None) - _retain_recent_subagent(record) + if record is None or not (agent is None or record.get("agent") is agent): + return + _active_subagents.pop(subagent_id, None) + sid = record.get("subagent_id") + if not sid: + return + _recent_subagents[sid] = { + k: record.get(k) for k in ("goal", "delegation_id", "owner_agent_session_id") + } + while len(_recent_subagents) > _RECENT_SUBAGENTS_CAP: + _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. @@ -119,27 +102,19 @@ def _close_subagent_steering(subagent_id: str, agent: Any) -> Optional[str]: return pending if isinstance(pending, str) and pending.strip() else None def interrupt_subagent(subagent_id: str) -> bool: - """Request that a single running subagent stop at its next iteration boundary. - - Does not hard-kill the worker thread (Python can't); sets the child's - interrupt flag which propagates to in-flight tools and recurses into - grandchildren via AIAgent.interrupt(). Returns True if a matching - subagent was found. - """ + """Request that one running subagent stop at its next iteration boundary + (cooperative: the flag propagates to in-flight tools and recurses into + grandchildren via AIAgent.interrupt()). True iff a matching subagent was found.""" with _active_subagents_lock: record = _active_subagents.get(subagent_id) - if not record: - return False - agent = record.get("agent") + agent = record.get("agent") if record else None if agent is None: return False try: - if not request_hard_interrupt(agent, f"Interrupted via TUI ({subagent_id})"): - return False + return bool(request_hard_interrupt(agent, f"Interrupted via TUI ({subagent_id})")) except Exception as exc: logger.debug("interrupt_subagent(%s) failed: %s", subagent_id, exc) return False - return True def steer_subagent( subagent_id: str, text: str, *, owner_session_id: Optional[str] = None, owner_transport: Any = None, @@ -147,14 +122,14 @@ def steer_subagent( ) -> bool: """Queue steering text into a running subagent without stopping it. - Mirror of interrupt_subagent(): calls AIAgent.steer(), which appends the - text to the child's last tool result at its next iteration boundary — the - current tool call is never cut. True iff the text was QUEUED while the child - still accepted work; False for unknown/closed id, ownership mismatch, no - live agent, or empty text. ``owner_session_id=None`` keeps the in-process - helper contract; gateway callers must pass exact authority. Acceptance and - completion are linearized by the registry lock: if acceptance wins but no - delivery boundary remains, the text lands in the entry as ``missed_steer``. + Calls AIAgent.steer(), which appends the text to the child's last tool result + at its next iteration boundary — the current tool call is never cut. True iff + the text was QUEUED while the child still accepted work; False for + unknown/closed id, ownership mismatch, no live agent, or empty text. + ``owner_session_id=None`` keeps the in-process helper contract; gateway + callers must pass exact authority. Acceptance and completion are linearized + by the registry lock: if acceptance wins but no delivery boundary remains, + the text lands in the entry as ``missed_steer``. """ if not text or not text.strip(): return False @@ -180,11 +155,8 @@ def steer_subagent( return False def _capture_gateway_steer_authority(owner_session_id: Optional[str]) -> tuple[Any, Any]: - """Capture exact request transport + live session generation, if any. - - This is intentionally an in-process bridge, not a serializable capability. - Non-gateway hosts (including the CLI helper path) receive ``(None, None)``. - """ + """Exact request transport + live session generation, if any — an in-process + bridge, not a serializable capability. Non-gateway hosts get ``(None, None)``.""" if not owner_session_id: return None, None try: @@ -197,21 +169,16 @@ def _capture_gateway_steer_authority(owner_session_id: Optional[str]) -> tuple[A _PRIVATE_RECORD_KEYS = frozenset({"agent", "owner_session_id", "owner_transport", "owner_session_record", "accepting_steer"}) def list_active_subagents() -> List[Dict[str, Any]]: - """Snapshot of the currently running subagent tree. - - Each record: {subagent_id, parent_id, depth, goal, model, started_at, - tool_count, status}. Safe to call from any thread — returns a copy. - """ + """Copy of the running subagent tree ({subagent_id, parent_id, depth, goal, model, + started_at, tool_count, status, ...}); safe from any thread.""" with _active_subagents_lock: 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 comparison 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 @@ -230,17 +197,12 @@ 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: - """Resolve a session id to the tip of its compression lineage. - - Best-effort: uses the parent's live SessionDB handle when present so a - delegation dispatched before a compression rotation still matches the - rotated parent. Returns the input unchanged when resolution fails. - """ + """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 "") - if not sid: - return "" db = getattr(parent_agent, "_session_db", None) - if db is None: + if not sid or db is None: return sid try: resolved = db.resolve_resume_session_id(sid) @@ -262,21 +224,43 @@ def _owns_subagent_record(record: Dict[str, Any], parent_agent: Any) -> bool: session id; control must use the same spine or running children go invisible/unsteerable. """ - agent = record.get("agent") - if _is_descendant_of(agent, parent_agent): + if _is_descendant_of(record.get("agent"), parent_agent): return True owner_sid = str(record.get("owner_agent_session_id") or "") - if not owner_sid: - return False parent_sid = str(getattr(parent_agent, "session_id", "") or "") - if not parent_sid: + if not owner_sid or not parent_sid: return False if owner_sid == parent_sid: 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), - } + return _resolve_session_lineage(owner_sid, parent_agent) in {parent_sid, _resolve_session_lineage(parent_sid, parent_agent)} + +def _list_payload(parent_agent: Any) -> Dict[str, Any]: + with _active_subagents_lock: + records = list(_active_subagents.values()) + entries = [] + for r in records: + if not _owns_subagent_record(r, parent_agent): + continue + started = r.get("started_at") + entries.append({ + "subagent_id": r.get("subagent_id"), + "parent_id": r.get("parent_id"), + "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, + "accepting_steer": bool(r.get("accepting_steer", False)), + "live_transcript": getattr(r.get("agent"), "_live_transcript_path", None), + }) + payload: Dict[str, Any] = {"action": "list", "count": len(entries), "subagents": entries} + if not entries: + payload["note"] = ( + "No live subagents right now. Children that already finished " + "have delivered (or will deliver) their results as normal " + "completion messages — there is nothing to steer or stop." + ) + 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. @@ -286,42 +270,12 @@ def _handle_control_action(action: str, subagent_id: Optional[str], message: Opt a conversation can only control its own spawn tree. """ if action == "list": - with _active_subagents_lock: - records = list(_active_subagents.values()) - entries = [] - for r in records: - agent = r.get("agent") - if not _owns_subagent_record(r, parent_agent): - continue - started = r.get("started_at") - entries.append( - { - "subagent_id": r.get("subagent_id"), - "parent_id": r.get("parent_id"), - "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), - "accepting_steer": bool(r.get("accepting_steer", False)), - "live_transcript": getattr(agent, "_live_transcript_path", None), - } - ) - payload: Dict[str, Any] = {"action": "list", "count": len(entries), "subagents": entries} - if not entries: - payload["note"] = ( - "No live subagents right now. Children that already finished " - "have delivered (or will deliver) their results as normal " - "completion messages — there is nothing to steer or stop." - ) - return json.dumps(payload, ensure_ascii=False) + return json.dumps(_list_payload(parent_agent), ensure_ascii=False) # steer / stop need a resolvable, owned target. sid = (subagent_id or "").strip() if not sid: - return tool_error( - f"action='{action}' requires subagent_id (from the spawn dispatch " - "response or action='list')." - ) + return tool_error(f"action='{action}' requires subagent_id (from the spawn dispatch response or action='list').") with _active_subagents_lock: record = _active_subagents.get(sid) if record is None or not _owns_subagent_record(record, parent_agent): @@ -330,9 +284,8 @@ def _handle_control_action(action: str, subagent_id: Optional[str], message: Opt "may have already finished (its result arrives as a normal " "completion message). Use action='list' to see live children." ) - if action == "steer" and not (message or "").strip(): - return tool_error("action='steer' requires a non-empty 'message' describing the " "course correction.") + return tool_error("action='steer' requires a non-empty 'message' describing the course correction.") outcome = _CONTROL_OUTCOMES.get(action) if outcome is None: return tool_error(f"Unknown action '{action}'. Use spawn, list, steer, or stop.") @@ -346,16 +299,21 @@ 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.", ),