feat(delegation): subagents hand background processes to the parent; leftovers are named, not trusted
A child's background processes are killed at its teardown and their notify_on_complete notices are suppressed in the parent, yet the child's terminal result still said `notify_on_complete: true` and the parent's delegation notice said nothing about processes left behind. Orchestrators believed "CI watcher running" and waited on a completion that could never arrive (recurring in the Sep 7 campaign sessions). - process_manage(action="handoff", session_id, data="<purpose>"), children only: process_registry.transfer_ownership flips owner_task_id/task_id/ session_key to the parent under the registry lock, so the completion is stamped with the parent's owner at exit, passes the parent's sa- filter, and is reaped by the parent, not the child. Cap 3 per child; an exited, foreign, or non-child request is a tool error. The purpose rides the event as handoff_note and renders in the parent's notice. - Child terminal(background=True, notify=True) now returns notify_on_complete=false plus a note: wait, kill, or hand off. - _ChildRun.account_background_processes records handed_off_processes and orphaned_processes on the result before cleanup kills the leftovers; the parent's delegation block renders both.
This commit is contained in:
@@ -0,0 +1,120 @@
|
||||
"""Subagent → parent background-process handoff (process_manage action='handoff').
|
||||
|
||||
Ownership is the process's ``owner_task_id``: completion notices are stamped from it at exit and the parent's
|
||||
``drain_notifications`` suppresses ``sa-`` owners, so a child's watcher never reaches the parent and is killed at child
|
||||
teardown. A validated handoff flips the owner so the completion lands in the parent's chat; anything not handed off is
|
||||
named on the child's result as orphaned before teardown kills it.
|
||||
"""
|
||||
|
||||
import json
|
||||
import time
|
||||
import weakref
|
||||
|
||||
import pytest
|
||||
|
||||
from tools.delegate_tool import _register_subagent, _unregister_subagent
|
||||
from tools.process_registry import ProcessRegistry, process_registry, _handle_process
|
||||
from tools.process_registry_notifications import format_process_notification
|
||||
|
||||
|
||||
class _Parent:
|
||||
def __init__(self):
|
||||
self.session_id = "sess-handoff"
|
||||
self._current_task_id = "parent-turn-1"
|
||||
|
||||
|
||||
class _Child:
|
||||
def __init__(self, parent):
|
||||
self._delegate_parent_ref = weakref.ref(parent)
|
||||
|
||||
|
||||
def _register(sid, child):
|
||||
_register_subagent({"subagent_id": sid, "parent_id": None, "depth": 0, "goal": "g", "model": "m",
|
||||
"started_at": time.time(), "status": "running", "tool_count": 0, "agent": child})
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _plain_spawn(monkeypatch):
|
||||
"""Spawn plain children: the systemd-run --user --scope wrapper is irrelevant here and stalls under pytest."""
|
||||
import tools.process_registry as _pr
|
||||
monkeypatch.setattr(_pr, "_SYSTEMD_SCOPE_AVAILABLE", False)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def clean_queue():
|
||||
while not process_registry.completion_queue.empty():
|
||||
process_registry.completion_queue.get_nowait()
|
||||
yield
|
||||
while not process_registry.completion_queue.empty():
|
||||
process_registry.completion_queue.get_nowait()
|
||||
|
||||
|
||||
def test_handed_off_process_completion_reaches_parent_and_leftover_is_reported(clean_queue):
|
||||
"""A real child-owned process handed off carries the parent's owner id (so the parent's drain accepts it, with the
|
||||
handoff purpose), while a sibling the child did not hand off is still owned by the child and is listed as orphaned
|
||||
on the child's result."""
|
||||
sid = "sa-0-handoff01"
|
||||
parent = _Parent()
|
||||
child = _Child(parent)
|
||||
_register(sid, child)
|
||||
try:
|
||||
handed = process_registry.spawn_local("sleep 0.4; echo ci-green", task_id=sid, owner_task_id=sid)
|
||||
handed.notify_on_complete = True
|
||||
leftover = process_registry.spawn_local("sleep 30", task_id=sid, owner_task_id=sid)
|
||||
|
||||
out = json.loads(_handle_process(
|
||||
{"action": "handoff", "session_id": handed.id, "data": "CI watcher for PR 1"}, task_id=sid))
|
||||
assert out["status"] == "handed_off"
|
||||
assert handed.owner_task_id == "parent-turn-1" and handed.session_key == "sess-handoff"
|
||||
assert child._handed_off_processes[0]["session_id"] == handed.id
|
||||
|
||||
# Only the un-handed sibling is left in the child's name — this is what the result entry reports as orphaned.
|
||||
assert [s.id for s in process_registry.running_owned_by(sid)] == [leftover.id]
|
||||
|
||||
# Let it exit on its own — a registry wait() would mark the completion consumed (that is the parent-observed path).
|
||||
deadline = time.time() + 10
|
||||
while process_registry.completion_queue.empty() and time.time() < deadline:
|
||||
time.sleep(0.05)
|
||||
assert handed.exited
|
||||
# The parent drains with the default suppression of sa- owners: the handed-off completion passes it.
|
||||
events = process_registry.drain_notifications(owns_event=lambda e: True)
|
||||
mine = [(e, text) for e, text in events if e.get("session_id") == handed.id]
|
||||
assert len(mine) == 1
|
||||
evt, text = mine[0]
|
||||
assert evt["owner_task_id"] == "parent-turn-1"
|
||||
assert "Handed off to you by a subagent" in text and "CI watcher for PR 1" in text
|
||||
finally:
|
||||
process_registry.kill_all(source="test")
|
||||
_unregister_subagent(sid)
|
||||
|
||||
|
||||
def test_handoff_refuses_exited_foreign_or_non_child_callers(clean_queue):
|
||||
"""A PID in prose is not a transfer: handoff is an error for a process that already exited, one the caller does
|
||||
not own, or a caller that is not a registered subagent — never a silent no-op."""
|
||||
sid, other = "sa-0-handoff02", "sa-1-handoff03"
|
||||
parent = _Parent()
|
||||
_register(sid, _Child(parent))
|
||||
try:
|
||||
done = process_registry.spawn_local("true", task_id=sid, owner_task_id=sid)
|
||||
process_registry.wait(done.id, timeout=10)
|
||||
assert "error" in json.loads(_handle_process(
|
||||
{"action": "handoff", "session_id": done.id, "data": "x"}, task_id=sid))
|
||||
|
||||
foreign = process_registry.spawn_local("sleep 30", task_id=other, owner_task_id=other)
|
||||
assert "error" in json.loads(_handle_process(
|
||||
{"action": "handoff", "session_id": foreign.id, "data": "x"}, task_id=sid))
|
||||
assert foreign.owner_task_id == other
|
||||
|
||||
assert "error" in json.loads(_handle_process(
|
||||
{"action": "handoff", "session_id": foreign.id, "data": "x"}, task_id="not-a-subagent"))
|
||||
|
||||
# A completion for an un-handed child process is still suppressed in the parent.
|
||||
reg = ProcessRegistry()
|
||||
reg.completion_queue.put({"type": "completion", "session_id": "proc_x", "task_id": other,
|
||||
"owner_task_id": other, "command": "sleep", "exit_code": 0, "output": ""})
|
||||
assert reg.drain_notifications() == []
|
||||
assert format_process_notification({"type": "completion", "session_id": "p", "command": "c", "exit_code": 0,
|
||||
"output": "", "handoff_note": "why"}).count("Handed off to you") == 1
|
||||
finally:
|
||||
process_registry.kill_all(source="test")
|
||||
_unregister_subagent(sid)
|
||||
+5
-1
@@ -89,7 +89,11 @@ no `delegate_task`, `clarify`, `memory`, `send_message`, `cronjob`; keeps `execu
|
||||
`orchestrator` (keeps `delegate_task`; gated by `delegation.orchestrator_enabled`, bounded by
|
||||
`delegation.max_spawn_depth`, default 2). Config knobs under `delegation:`:
|
||||
`max_concurrent_children, independent_completions, max_spawn_depth, child_timeout_seconds, orchestrator_enabled,
|
||||
subagent_auto_approve, inherit_mcp_toolsets, max_iterations`. **Durability:** background
|
||||
subagent_auto_approve, inherit_mcp_toolsets, max_iterations`. **Child processes:** a child's background
|
||||
processes are killed at its teardown and their notices are suppressed in the parent; `process_manage(action="handoff")`
|
||||
(children only) flips `ProcessSession.owner_task_id` to the parent under the registry lock
|
||||
(`process_registry.transfer_ownership`) so the completion routes and reaps by the new owner; un-handed leftovers land on
|
||||
the result as `orphaned_processes` (`_ChildRun.account_background_processes`, before `cleanup` kills them). **Durability:** background
|
||||
delegation is process-local; work that must survive restart uses `cronjob` or
|
||||
`terminal(background=True, notify_on_complete=True)`. API: `website/docs/developer-guide/subagent-lifecycle-api.md`.
|
||||
|
||||
|
||||
@@ -332,6 +332,7 @@ def _run_single_child(
|
||||
duration = run.elapsed()
|
||||
entry = _build_result_entry(child, result, task_index, duration, schema)
|
||||
run.append_sibling_write_reminder(entry)
|
||||
run.account_background_processes(entry)
|
||||
run.emit_complete(result, entry, duration)
|
||||
return run.attach_worktree(entry)
|
||||
except Exception as exc:
|
||||
|
||||
@@ -741,6 +741,21 @@ class _ChildRun:
|
||||
else:
|
||||
entry["stale_paths"] = mod_paths
|
||||
|
||||
def account_background_processes(self, entry: Dict[str, Any]) -> None:
|
||||
"""Name the child's background processes on the result BEFORE ``cleanup`` kills them: handed-off ones now
|
||||
belong to the parent (their completion lands in the parent's chat); anything else still running is about to be
|
||||
terminated, and the parent must hear that from the runtime rather than trust a child's "watcher running"."""
|
||||
handed = list(getattr(self.child, "_handed_off_processes", None) or [])
|
||||
if handed:
|
||||
entry["handed_off_processes"] = handed
|
||||
with _quiet(None):
|
||||
from tools.process_registry import process_registry
|
||||
leftover = process_registry.running_owned_by(self.child_task_id)
|
||||
if leftover:
|
||||
entry["orphaned_processes"] = [
|
||||
{"session_id": s.id, "command": s.command[:200], "runtime_seconds": round(time.time() - s.started_at)}
|
||||
for s in leftover]
|
||||
|
||||
def emit_complete(self, result: Dict[str, Any], entry: Dict[str, Any], duration: float) -> None:
|
||||
"""Fire ``subagent.complete`` with the per-branch observability payload (tokens, cost, files touched,
|
||||
tool-output tail); every field is optional and degrades gracefully on the client."""
|
||||
|
||||
@@ -324,6 +324,7 @@ class ProcessSession:
|
||||
detached: bool = False # Recovered from checkpoint (no pipe)
|
||||
pid_scope: str = "host" # "host" for local/PTY PIDs, "sandbox" for env-local PIDs
|
||||
systemd_unit: str = "" # transient scope unit name when spawned under systemd-run
|
||||
handoff_note: str = "" # why a subagent handed this process to its parent (rides the notice)
|
||||
# Watcher/notification routing (persisted for crash recovery)
|
||||
# systemd_unit: str = "" # transient scope unit name when spawned under systemd-run
|
||||
# (#70716)
|
||||
@@ -1197,6 +1198,7 @@ class ProcessRegistry(ProcessCheckpointMixin):
|
||||
"task_id": session.task_id,
|
||||
"owner_task_id": session.owner_task_id or session.task_id,
|
||||
"command": session.command,
|
||||
**({"handoff_note": session.handoff_note} if session.handoff_note else {}),
|
||||
**self._exit_fields(session),
|
||||
"output": _output_tail(session, 2000),
|
||||
# Stable producer identity across checkpoint recovery (unlike a
|
||||
@@ -1857,6 +1859,27 @@ class ProcessRegistry(ProcessCheckpointMixin):
|
||||
"""Whether any process for ``task_id`` is still running."""
|
||||
return self._any_running(lambda s: s.task_id == task_id)
|
||||
|
||||
def running_owned_by(self, owner_task_id: str) -> List[ProcessSession]:
|
||||
"""Running processes whose RAW spawning owner is ``owner_task_id``."""
|
||||
with self._lock:
|
||||
return [s for s in self._running.values() if s.owner_task_id == owner_task_id and not s.exited]
|
||||
|
||||
def transfer_ownership(self, session_id: str, *, from_owner: str, to_owner: str, to_task_id: str,
|
||||
to_session_key: str, note: str = "") -> Optional[ProcessSession]:
|
||||
"""Move a RUNNING process from one owner to another under the registry lock. Ownership is the ``owner_task_id``
|
||||
field: completion notices are stamped from it at exit time and teardown kills by it, so flipping it here is the
|
||||
whole transfer. Returns the session, or None when it is unknown, already exited, or not owned by ``from_owner``
|
||||
(the caller must not report a transfer that did not happen)."""
|
||||
session = self.get(session_id)
|
||||
with self._lock:
|
||||
if session is None or session.exited or session.owner_task_id != from_owner:
|
||||
return None
|
||||
session.owner_task_id = to_owner
|
||||
session.task_id = to_task_id
|
||||
session.session_key = to_session_key
|
||||
session.handoff_note = note
|
||||
return session
|
||||
|
||||
def has_active_for_session(self, session_key: str, max_active_age: Optional[float] = None) -> bool:
|
||||
"""Active processes for a gateway session key. Processes older than
|
||||
``max_active_age`` seconds are ignored as stale so a forgotten ``http.server``
|
||||
@@ -1938,14 +1961,17 @@ PROCESS_SCHEMA = {
|
||||
"poll: status + new output. log: full output, paged. wait: block "
|
||||
"until exit or timeout (partial output on timeout). write vs "
|
||||
"submit: submit appends Enter — use it to answer prompts; write "
|
||||
"sends raw bytes, no newline. close: EOF stdin. kill: terminate."
|
||||
"sends raw bytes, no newline. close: EOF stdin. kill: terminate. "
|
||||
"handoff (subagents only): transfer a running process you started to your parent agent, which then "
|
||||
"receives its completion; `data` = one sentence on its purpose. Subagent-owned processes are otherwise "
|
||||
"killed when the subagent finishes and their notifications never reach the parent."
|
||||
),
|
||||
"parameters": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"action": {
|
||||
"type": "string",
|
||||
"enum": ["list", "poll", "log", "wait", "kill", "write", "submit", "close"]
|
||||
"enum": ["list", "poll", "log", "wait", "kill", "write", "submit", "close", "handoff"]
|
||||
},
|
||||
"session_id": {
|
||||
"type": "string",
|
||||
@@ -1953,7 +1979,7 @@ PROCESS_SCHEMA = {
|
||||
},
|
||||
"data": {
|
||||
"type": "string",
|
||||
"description": "Stdin text for write/submit."
|
||||
"description": "Stdin text for write/submit; purpose sentence for handoff."
|
||||
},
|
||||
"timeout": {
|
||||
"type": "integer",
|
||||
@@ -2023,19 +2049,64 @@ _SESSION_ACTIONS = {
|
||||
}
|
||||
|
||||
|
||||
def _handoff_process(session_id: str, args: dict, task_id: Optional[str]) -> dict:
|
||||
"""Subagent-only: transfer a running background process to the parent agent so its completion is delivered THERE
|
||||
(child-owned process notices are suppressed and child teardown kills what it owns). Validated against the live spawn
|
||||
tree: the caller must be a registered child and must own the process; anything else is an error, never a silent
|
||||
no-op, so a PID mentioned in prose can't masquerade as a transfer."""
|
||||
from tools.delegate_tool_registry import _active_subagents, _active_subagents_lock
|
||||
from tools.terminal_tool import _resolve_container_task_id
|
||||
with _active_subagents_lock:
|
||||
record = _active_subagents.get(str(task_id or ""))
|
||||
child = record.get("agent") if record else None
|
||||
parent_ref = getattr(child, "_delegate_parent_ref", None)
|
||||
parent = parent_ref() if callable(parent_ref) else None
|
||||
if parent is None:
|
||||
return {"error": "handoff is only available to a running subagent with a live parent; you are not one."}
|
||||
parent_owner = str(getattr(parent, "_current_task_id", "") or getattr(parent, "session_id", "") or "")
|
||||
if not parent_owner:
|
||||
return {"error": "parent has no process owner id yet; retry after the parent's turn has started."}
|
||||
handed = getattr(child, "_handed_off_processes", None)
|
||||
if handed is None:
|
||||
handed = child._handed_off_processes = []
|
||||
if len(handed) >= _MAX_HANDOFFS_PER_CHILD:
|
||||
return {"error": f"handoff cap reached ({_MAX_HANDOFFS_PER_CHILD} per subagent); wait on or kill the rest yourself."}
|
||||
note = str(args.get("data") or "").strip()
|
||||
if not note:
|
||||
return {"error": "handoff requires `data`: one sentence saying what the process is for and what the parent should do with its result."}
|
||||
session = process_registry.transfer_ownership(
|
||||
session_id, from_owner=str(task_id or ""), to_owner=parent_owner,
|
||||
to_task_id=_resolve_container_task_id(parent_owner),
|
||||
to_session_key=str(getattr(parent, "session_id", "") or ""), note=note)
|
||||
if session is None:
|
||||
return {"error": f"cannot hand off {session_id}: not a running process you own (already exited? read its result "
|
||||
"with poll/log and report it instead)."}
|
||||
handed.append({"session_id": session.id, "command": session.command, "note": note})
|
||||
return {"status": "handed_off", "session_id": session.id, "command": session.command,
|
||||
"note": "Your parent now owns this process and will receive its completion; you will not. Mention the handoff "
|
||||
"in your final answer."}
|
||||
|
||||
|
||||
_MAX_HANDOFFS_PER_CHILD = 3
|
||||
|
||||
|
||||
def _handle_process(args, **kw):
|
||||
action = args.get("action", "")
|
||||
# Coerce to string — some models send session_id as an integer
|
||||
session_id = str(args.get("session_id", "")) if args.get("session_id") is not None else ""
|
||||
if action == "list":
|
||||
return json.dumps(_list_processes(kw.get("task_id")), ensure_ascii=False)
|
||||
if action == "handoff":
|
||||
if not session_id:
|
||||
return tool_error("session_id is required for handoff")
|
||||
return json.dumps(_handoff_process(session_id, args, kw.get("task_id")), ensure_ascii=False)
|
||||
if action in _SESSION_ACTIONS:
|
||||
if not session_id:
|
||||
return tool_error(f"session_id is required for {action}")
|
||||
handler, redact = _SESSION_ACTIONS[action]
|
||||
result = handler(session_id, args)
|
||||
return json.dumps(_redact_process_result(result) if redact else result, ensure_ascii=False)
|
||||
return tool_error(f"Unknown process action: {action}. Use: list, poll, log, wait, kill, write, submit, close")
|
||||
return tool_error(f"Unknown process action: {action}. Use: list, poll, log, wait, kill, write, submit, close, handoff")
|
||||
|
||||
|
||||
registry.register(
|
||||
|
||||
@@ -212,9 +212,26 @@ def _format_batch_delegation(evt: dict, deleg_id: str, completed_at: float) -> s
|
||||
lines.append(f"(no summary — status={r_status}" + (f": {r_error}" if r_error else "") + ")")
|
||||
if r.get("live_transcript"):
|
||||
lines.append(f"Full live transcript (complete tool/assistant trace): {r['live_transcript']}")
|
||||
lines += _process_accounting_lines(r)
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
def _process_accounting_lines(r: dict) -> list:
|
||||
"""Runtime-truth lines about a child's background processes: what it handed to you (you own it now, its
|
||||
completion lands here) and what it left running (terminated at teardown — never trust a child's "watcher running")."""
|
||||
lines = []
|
||||
for h in r.get("handed_off_processes") or []:
|
||||
lines.append(f"Handed off to you: {h.get('session_id')} ({h.get('command', '')[:120]}) — {h.get('note', '')}. "
|
||||
"You own it now; its completion notice will arrive here.")
|
||||
orphans = r.get("orphaned_processes") or []
|
||||
if orphans:
|
||||
lines.append(f"Child left {len(orphans)} background process(es) running that were TERMINATED with it "
|
||||
"(subagent process notices never reach you): "
|
||||
+ "; ".join(f"{o.get('session_id')} `{o.get('command', '')[:100]}` ({o.get('runtime_seconds')}s)" for o in orphans)
|
||||
+ ". Re-launch in this session anything you still need.")
|
||||
return lines
|
||||
|
||||
|
||||
def _format_async_delegation(evt: dict) -> str:
|
||||
"""Self-contained re-injection for an async-delegation completion: the FULL
|
||||
original task source (goal, context, toolsets, role, model), dispatch time, status
|
||||
@@ -331,6 +348,8 @@ def format_process_notification(evt: dict) -> "str | None":
|
||||
return _format_async_delegation(evt)
|
||||
_sid, _cmd = evt.get("session_id", "unknown"), evt.get("command", "unknown")
|
||||
_attribution = _delegation_attribution_line(evt)
|
||||
if evt.get("handoff_note"):
|
||||
_attribution = f"Handed off to you by a subagent before it finished. Purpose: {evt['handoff_note']}"
|
||||
attribution = f"{_attribution}\n" if _attribution else ""
|
||||
if evt_type == "watch_match":
|
||||
_sup = evt.get("suppressed", 0)
|
||||
|
||||
@@ -177,6 +177,10 @@ def spawn_background_process(
|
||||
result_data["notify_on_complete"] = True
|
||||
if proc_session.watcher_platform:
|
||||
_register_completion_watcher(process_registry, proc_session, session_key)
|
||||
from agent.delegation_context import is_delegated_child_context
|
||||
if is_delegated_child_context():
|
||||
result_data["notify_on_complete"] = False
|
||||
result_data["subagent_note"] = _SUBAGENT_NOTIFY_NOTE
|
||||
if watch_patterns:
|
||||
proc_session.watch_patterns = list(watch_patterns)
|
||||
result_data["watch_patterns"] = proc_session.watch_patterns
|
||||
@@ -188,6 +192,13 @@ def spawn_background_process(
|
||||
}, ensure_ascii=False)
|
||||
|
||||
|
||||
_SUBAGENT_NOTIFY_NOTE = (
|
||||
"You are a subagent: this process's completion notice will NOT reach your parent, and the process is killed when "
|
||||
"you finish. Before you finish, either wait for it (process_manage wait), kill it, or hand it to your parent with "
|
||||
"process_manage(action='handoff', session_id=..., data='<purpose>') so the parent receives its completion. For CI "
|
||||
"watchers prefer returning the fact (PR number, SHA) and letting the parent watch."
|
||||
)
|
||||
|
||||
_YIELDED_NOTE = (
|
||||
"The user sent a message while this command was running, so it was moved to the "
|
||||
"background WITHOUT being killed and is still running. You will be notified when it "
|
||||
|
||||
@@ -189,6 +189,16 @@ delegation:
|
||||
surface_child_process_notifications: true # default: false
|
||||
```
|
||||
|
||||
### Handing a process to the parent
|
||||
|
||||
A subagent's background processes are also **killed when the subagent finishes**, so a CI watcher or build a child starts with `notify=true` never reports to anyone. The child's `terminal` result says so (`notify_on_complete: false` plus a `subagent_note`), and the child has three honest options before it finishes:
|
||||
|
||||
- **wait** — `process_manage(action="wait", session_id=...)` and report the result itself;
|
||||
- **kill** — `process_manage(action="kill", ...)`;
|
||||
- **hand off** — `process_manage(action="handoff", session_id=..., data="<one sentence: what it is for>")`. The runtime transfers ownership to the parent under the registry lock (up to 3 per child; only a running process the child owns is accepted, anything else is a tool error). The parent's completion notice then arrives in the parent chat with `Handed off to you by a subagent… Purpose: …`, and the parent can poll/log/kill it like its own.
|
||||
|
||||
Whatever the child neither waited on, killed, nor handed off is named on its result (`orphaned_processes`) and in the parent's delegation notice as terminated, so the parent hears from the runtime, never from the child's prose, that "the watcher is running" is no longer true. For CI watchers the better pattern is still: the child returns the fact (PR number, SHA) and the parent launches its own watcher.
|
||||
|
||||
## Model Override
|
||||
|
||||
You can configure a different model for subagents via `config.yaml` — useful for delegating simple tasks to cheaper/faster models:
|
||||
|
||||
Reference in New Issue
Block a user