diff --git a/backend/apps/agents/manager/streaming/delegation_watchdog.py b/backend/apps/agents/manager/streaming/delegation_watchdog.py index 6797469f..7a106401 100644 --- a/backend/apps/agents/manager/streaming/delegation_watchdog.py +++ b/backend/apps/agents/manager/streaming/delegation_watchdog.py @@ -14,7 +14,7 @@ from typing import Set from typeguard import typechecked from backend.apps.agents.manager.streaming.unwedge_sidecar import ( - HEARTBEAT_FRESH_S, CORE_PREFIX, RETRY_PROMPT, arm_retry, heartbeat_age, unwedge, + HEARTBEAT_FRESH_S, CORE_PREFIX, RETRY_PROMPT, arm_retry, delegation_children_born_after, heartbeat_age, unwedge, ) logger = logging.getLogger(__name__) @@ -49,17 +49,7 @@ def delegation_children_settled(session_id: str, since: float) -> bool: p_parent = agent_manager.sessions.get(session_id) if p_parent is not None and getattr(p_parent, "ended_by_user", False): return False - kids = [] - for s in agent_manager.sessions.values(): - if getattr(s, "parent_session_id", None) != session_id or getattr(s, "mode", "") != "browser-agent": - continue - born = getattr(s, "created_at", None) - try: - if born is None or born.timestamp() < since: - continue - except Exception: - continue - kids.append(s) + kids = delegation_children_born_after(session_id, since) if not kids: return False return all(getattr(s, "status", "") in ("completed", "error", "failed", "stopped") for s in kids) diff --git a/backend/apps/agents/manager/streaming/unwedge_sidecar.py b/backend/apps/agents/manager/streaming/unwedge_sidecar.py index 86332d73..decb1a32 100644 --- a/backend/apps/agents/manager/streaming/unwedge_sidecar.py +++ b/backend/apps/agents/manager/streaming/unwedge_sidecar.py @@ -14,7 +14,7 @@ import os import subprocess import threading import time -from typing import Set +from typing import Dict, List, Set from typeguard import typechecked @@ -146,6 +146,42 @@ def arm_retry(session: object) -> bool: return False +@typechecked +def delegation_children_born_after(session_id: str, since: float) -> List[object]: + """The browser/app children this session spawned after `since`. One definition, shared by the + watchdog's settled test and the unwedge envelope, so the two can never count different children.""" + from backend.apps.agents.agent_manager import agent_manager + kids: List[object] = [] + for s in agent_manager.sessions.values(): + if getattr(s, "parent_session_id", None) != session_id or getattr(s, "mode", "") != "browser-agent": + continue + born = getattr(s, "created_at", None) + try: + if born is None or born.timestamp() < since: + continue + except Exception: + continue + kids.append(s) + return kids + + +@typechecked +def children_summary(session_id: str, since: float) -> List[Dict[str, object]]: + """What the children looked like at the moment of a kill: status plus how long ago each was born. + Haik's 154 CreateBrowserAgent kills could not be split into "child finished, result lost" versus + "child died first" because the envelope carried only tool, seconds and pids (read 2026-09-01).""" + now = time.time() + out: List[Dict[str, object]] = [] + for s in delegation_children_born_after(session_id, since): + born = getattr(s, "created_at", None) + try: + age = round(now - born.timestamp(), 1) if born is not None else None + except Exception: + age = None + out.append({"status": str(getattr(s, "status", "") or ""), "age_s": age, "tools": len(getattr(s, "tool_latencies", {}) or {})}) + return out + + @typechecked def unwedge(session_id: str, tool_name: str, outstanding_s: float) -> int: """CONT first: a STOPPED process queues TERM forever (the ghost-reaper lesson, ENG-196), so @@ -169,6 +205,7 @@ def unwedge(session_id: str, tool_name: str, outstanding_s: float) -> int: "tool": tool_name, "outstanding_s": round(outstanding_s, 1), "pids": pids, + "children": children_summary(session_id, time.time() - outstanding_s), }) except Exception: pass diff --git a/backend/tests/test_unwedge_envelope_children.py b/backend/tests/test_unwedge_envelope_children.py new file mode 100644 index 00000000..ae7f4be9 --- /dev/null +++ b/backend/tests/test_unwedge_envelope_children.py @@ -0,0 +1,48 @@ +"""The unwedge envelope must say what the children looked like at the kill. + +Haik's install: 154 CreateBrowserAgent kills in 14 days at a 150 s floor, and the envelope carried +only tool, seconds and pids, so nobody could tell "child finished, result lost" (recovery by design) +from "child died first" (a different bug). One shared helper feeds the watchdog's settled test AND the +envelope, so the two can never disagree about which children count.""" +import datetime +from backend.apps.agents.agent_manager import agent_manager +from backend.apps.agents.core.models import AgentSession +from backend.apps.agents.manager.streaming import unwedge_sidecar as us +from backend.apps.agents.manager.streaming.delegation_watchdog import delegation_children_settled +from backend.apps.service import client as svc + + +def p_child(parent_id: str, status: str, born_ago_s: float) -> AgentSession: + s = AgentSession(name="Browser Agent", model="sonnet", mode="browser-agent", status=status) + s.parent_session_id = parent_id + s.created_at = datetime.datetime.now() - datetime.timedelta(seconds=born_ago_s) + return s + + +def test_the_envelope_carries_each_child_status_and_age(monkeypatch) -> None: + parent = AgentSession(name="p", model="sonnet") + agent_manager.sessions.clear(); agent_manager.sessions[parent.id] = parent + for kid in (p_child(parent.id, "completed", 30), p_child(parent.id, "error", 10), p_child(parent.id, "running", 5), p_child(parent.id, "completed", 900)): + agent_manager.sessions[kid.id] = kid + sent: list = [] + monkeypatch.setattr(svc, "submit_diagnostic", lambda d: sent.append(d)) + monkeypatch.setattr(us, "find_sidecar_pids", lambda sid: [424242]) + monkeypatch.setattr(us.subprocess, "run", lambda *a, **k: None) + monkeypatch.setattr(us.time, "sleep", lambda s: None) + us.unwedge(parent.id, "mcp__openswarm-core__CreateBrowserAgent", 150.0) + env = [d for d in sent if isinstance(d, dict) and d.get("kind") == "mcp_sidecar_unwedged"] + assert len(env) == 1 + kids = env[0]["children"] + assert sorted(k["status"] for k in kids) == ["completed", "error", "running"], "only children born after the call started; the 900 s one is an earlier delegation's" + assert all(isinstance(k["age_s"], float) for k in kids) + + +def test_the_watchdog_and_the_envelope_count_the_same_children() -> None: + parent = AgentSession(name="p", model="sonnet") + agent_manager.sessions.clear(); agent_manager.sessions[parent.id] = parent + agent_manager.sessions.update({k.id: k for k in (p_child(parent.id, "completed", 20), p_child(parent.id, "completed", 600))}) + import time + since = time.time() - 100 + assert len(us.delegation_children_born_after(parent.id, since)) == 1 + assert delegation_children_settled(parent.id, since) is True + assert len(us.children_summary(parent.id, since)) == 1