mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-09-30 05:24:50 +02:00
[eric] agents: the unwedge envelope carries each child's status and age, so a lost result can be told from a dead child
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012G8kyALnPjsA7aJFmMBq3R
This commit is contained in:
co-authored by
Claude Fable 5.1
parent
d2f1523004
commit
e6b3865851
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user