From 4f3e68c7ac43f6b3a07af08905200214d37dda23 Mon Sep 17 00:00:00 2001 From: ciregenz Date: Tue, 18 Aug 2026 21:13:47 -0700 Subject: [PATCH] [eric] agents: the 25s unwedge verifies the sidecar is actually wedged (heartbeat file) before shooting it; alive-but-slow gets 120s (ENG-353) --- .../apps/agents/combined_meta_mcp_server.py | 25 ++++++++++ .../manager/streaming/unwedge_sidecar.py | 40 +++++++++++++++- backend/tests/test_wedge_verdict.py | 46 +++++++++++++++++++ 3 files changed, 110 insertions(+), 1 deletion(-) create mode 100644 backend/tests/test_wedge_verdict.py diff --git a/backend/apps/agents/combined_meta_mcp_server.py b/backend/apps/agents/combined_meta_mcp_server.py index 216b3ff9..5c7825d4 100644 --- a/backend/apps/agents/combined_meta_mcp_server.py +++ b/backend/apps/agents/combined_meta_mcp_server.py @@ -95,7 +95,32 @@ def p_call_async(id_, tool_name: str, arguments: dict) -> None: send_response(id_, error={"code": -32000, "message": str(e)}) +def start_heartbeat(): + """Touch a per-session file every 5s from a daemon thread: proof this process is scheduled and + alive. The backend's wedge watchdog reads the mtime to tell a WEDGED sidecar (SIGSTOP, dead + process: heartbeat stops) from a merely SLOW tool call (threads fine, heartbeat keeps beating), + because killing the second kind is exactly the "MCP disconnected" a user reports (ENG-353).""" + import tempfile + import time as p_time + session = os.environ.get("OPENSWARM_PARENT_SESSION_ID", "") + if not session: + return + path = os.path.join(tempfile.gettempdir(), f"osw-mcp-hb-{session}") + + def p_beat(): + while True: + try: + with open(path, "a"): + os.utime(path, None) + except Exception: + pass + p_time.sleep(5) + + threading.Thread(target=p_beat, daemon=True, name="hb").start() + + def main(): + start_heartbeat() for line in sys.stdin: line = line.strip() if not line: diff --git a/backend/apps/agents/manager/streaming/unwedge_sidecar.py b/backend/apps/agents/manager/streaming/unwedge_sidecar.py index 0811b937..4e37f6a2 100644 --- a/backend/apps/agents/manager/streaming/unwedge_sidecar.py +++ b/backend/apps/agents/manager/streaming/unwedge_sidecar.py @@ -10,6 +10,7 @@ block (a human answering AskUI, a delegated browser run) are exempt by name, nev import asyncio import logging +import os import subprocess import threading import time @@ -24,6 +25,11 @@ logger = logging.getLogger(__name__) # the recovery as a hiccup rather than a hang. Anything that legitimately blocks (a human, a # delegated run) is exempt by name below, so this deadline never races real work. WEDGE_SECONDS = 25.0 +# A sidecar whose heartbeat still beats is ALIVE with one slow tool, not wedged; give it this long +# before concluding the call is hung anyway (measured: 5 healthy-sidecar kills in one loaded evening +# were every one of Haik's "MCP disconnected" reports, ENG-353). +LATE_WEDGE_SECONDS = 120.0 +HEARTBEAT_FRESH_S = 12.0 P_CORE_PREFIX = "mcp__openswarm-core__" @@ -72,6 +78,30 @@ def find_sidecar_pids(session_id: str) -> list: return pids +@typechecked +def heartbeat_age(session_id: str) -> float: + """Seconds since the session's sidecar last proved its process alive, or a huge number when no + heartbeat exists (old sidecar builds have none: treat as wedged-on-timeout, the old behavior).""" + import tempfile + path = os.path.join(tempfile.gettempdir(), f"osw-mcp-hb-{session_id}") + try: + return max(0.0, time.time() - os.path.getmtime(path)) + except OSError: + return 1e9 + + +@typechecked +def wedge_verdict(outstanding_s: float, hb_age: float) -> str: + """kill | extend | wait. Stale heartbeat = the PROCESS is wedged, kill at the first deadline. + Fresh heartbeat = alive with a slow call: extend once, and only a call still outstanding at the + late deadline dies (a hung per-call thread must not hang the session forever).""" + if outstanding_s >= LATE_WEDGE_SECONDS: + return "kill" + if hb_age > HEARTBEAT_FRESH_S: + return "kill" + return "extend" + + RETRY_PROMPT = ( "Your last tool call never returned because its server had frozen; that server has been " "restarted and works now. Retry that one step, then carry on where you left off." @@ -143,10 +173,18 @@ def arm_wedge_watchdog(ctx: object, tool_use_id: str, tool_name: str) -> None: session_id = getattr(ctx, "session_id", "") if not session_id: return + outstanding = time.time() - started + verdict = wedge_verdict(outstanding, heartbeat_age(session_id)) + if verdict == "extend": + logger.info( + f"Agent {session_id}: core tool {tool_name} outstanding {outstanding:.0f}s but the " + f"sidecar heartbeat is fresh (alive, slow); re-checking at {LATE_WEDGE_SECONDS:.0f}s") + loop.call_later(LATE_WEDGE_SECONDS - outstanding, p_check) + return # ps + kill are blocking; keep them off the event loop. A daemon thread, not the loop's # default executor: executor workers are non-daemon and a per-test loop that closes without # shutdown leaks them parked forever (the suite's flaky hang at interpreter exit). - threading.Thread(target=unwedge, args=(session_id, tool_name, time.time() - started), daemon=True, name="unwedge").start() + threading.Thread(target=unwedge, args=(session_id, tool_name, outstanding), daemon=True, name="unwedge").start() arm_retry(getattr(ctx, "session", None)) loop.call_later(WEDGE_SECONDS, p_check) diff --git a/backend/tests/test_wedge_verdict.py b/backend/tests/test_wedge_verdict.py new file mode 100644 index 00000000..ae9f7a6a --- /dev/null +++ b/backend/tests/test_wedge_verdict.py @@ -0,0 +1,46 @@ +"""Pins the ENG-353 two-stage wedge contract: a fresh heartbeat (alive sidecar, slow tool) must not +be shot at the first deadline; a stale one must; nothing survives the late deadline. Five healthy +kills in one loaded evening were every one of the "MCP disconnected" reports.""" +import os +import tempfile +import time +from backend.apps.agents.manager.streaming.unwedge_sidecar import ( + HEARTBEAT_FRESH_S, + LATE_WEDGE_SECONDS, + WEDGE_SECONDS, + heartbeat_age, + wedge_verdict, +) + + +def test_stale_heartbeat_kills_at_first_deadline(): + assert wedge_verdict(WEDGE_SECONDS + 1, HEARTBEAT_FRESH_S + 1) == "kill" + + +def test_fresh_heartbeat_extends_instead_of_killing(): + assert wedge_verdict(WEDGE_SECONDS + 1, 2.0) == "extend" + + +def test_late_deadline_kills_even_with_fresh_heartbeat(): + assert wedge_verdict(LATE_WEDGE_SECONDS + 1, 0.5) == "kill" + + +def test_missing_heartbeat_reads_as_wedged(): + assert heartbeat_age("no-such-session-anywhere") > 1e8 + + +def test_real_heartbeat_file_reads_fresh(): + sid = "wedge-verdict-test" + path = os.path.join(tempfile.gettempdir(), f"osw-mcp-hb-{sid}") + with open(path, "a"): + os.utime(path, None) + try: + assert heartbeat_age(sid) < 5.0 + assert wedge_verdict(WEDGE_SECONDS + 1, heartbeat_age(sid)) == "extend" + finally: + os.unlink(path) + + +def test_negative_control_old_behavior_without_heartbeat(): + age = heartbeat_age("no-such-session-anywhere") + assert wedge_verdict(WEDGE_SECONDS + 1, age) == "kill"