diff --git a/gateway/run.py b/gateway/run.py index 0304da3ddd..60a67d88b6 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -2655,15 +2655,27 @@ def _watch_gateway_turn_inactivity( *, agent_holder, task_id: str, process_baseline, timeout: float, worker_done: threading.Event, timeout_fired: threading.Event, cleanup_lock: threading.Lock, poll_interval: float = 5.0, is_still_current: Optional[Callable[[], bool]] = None) -> None: - """Thread watchdog that remains runnable when gateway asyncio is starved.""" + """Thread watchdog that remains runnable when gateway asyncio is starved. + + Until an agent publishes a usable activity snapshot, elapsed worker time is the + liveness clock. Otherwise a provider hang before activity initialization can + retain the session turn lease forever because every watchdog poll just skips it. + """ + activity_origin = time.monotonic() while not worker_done.wait(max(0.01, poll_interval)): + now = time.monotonic() + idle_seconds = now - activity_origin agent = agent_holder[0] if agent_holder else None - if agent is None or not hasattr(agent, "get_activity_summary"): - continue - try: - idle_seconds = float(agent.get_activity_summary().get("seconds_since_activity", 0.0)) - except Exception: - continue + if agent is not None and hasattr(agent, "get_activity_summary"): + try: + reported_idle = agent.get_activity_summary().get("seconds_since_activity") + if reported_idle is not None: + idle_seconds = max(0.0, float(reported_idle)) + # Preserve the most recent usable activity clock as the fallback if + # a later provider-side diagnostic read raises or returns None. + activity_origin = now - idle_seconds + except Exception: + pass if idle_seconds < timeout: continue _abandon_timed_out_gateway_turn( diff --git a/tests/gateway/test_abandoned_turn_process_cleanup.py b/tests/gateway/test_abandoned_turn_process_cleanup.py index da0fcdaa0c..3f37453090 100644 --- a/tests/gateway/test_abandoned_turn_process_cleanup.py +++ b/tests/gateway/test_abandoned_turn_process_cleanup.py @@ -22,10 +22,65 @@ class _IdleAgent: self.interrupts.append(reason) +class _UnavailableActivityAgent: + def __init__(self): + self.interrupts = [] + + def get_activity_summary(self): + return {"seconds_since_activity": None} + + def interrupt(self, reason): + self.interrupts.append(reason) + + def _state(): return threading.Event(), threading.Event(), threading.Lock() +def test_thread_watchdog_times_out_when_agent_activity_is_unavailable(monkeypatch): + agent = _UnavailableActivityAgent() + worker_done, timeout_fired, cleanup_lock = _state() + calls = [] + monkeypatch.setattr( + process_registry, + "kill_started_since", + lambda task_id, baseline, *, source: calls.append( + (task_id, baseline, source) + ) + or 0, + ) + + watchdog = threading.Thread( + target=_watch_gateway_turn_inactivity, + kwargs={ + "agent_holder": [agent], + "task_id": "session-without-agent", + "process_baseline": frozenset(), + "timeout": 0.03, + "worker_done": worker_done, + "timeout_fired": timeout_fired, + "cleanup_lock": cleanup_lock, + "poll_interval": 0.01, + }, + ) + watchdog.start() + watchdog.join(timeout=1) + if watchdog.is_alive(): + worker_done.set() + watchdog.join(timeout=1) + + assert not watchdog.is_alive() + assert timeout_fired.is_set() + assert agent.interrupts == ["Execution timed out (inactivity)"] + assert calls == [ + ( + "session-without-agent", + frozenset(), + "gateway_turn_timeout", + ) + ] + + def test_thread_watchdog_reaps_only_processes_created_by_timed_out_turn(monkeypatch): agent = _IdleAgent() worker_done, timeout_fired, cleanup_lock = _state()