diff --git a/agent/chat_completion_helpers.py b/agent/chat_completion_helpers.py index 98d5ceb83c..160b65b704 100644 --- a/agent/chat_completion_helpers.py +++ b/agent/chat_completion_helpers.py @@ -71,6 +71,35 @@ def _context_thread_target(callback): return lambda: context.run(callback) +def _join_worker_for_relay_teardown(worker, *, label: str) -> None: + """Bounded worker join before raising InterruptedError (#81521). + + Raising immediately lets turn teardown (finish_logical_calls / + end_turn / close_session) race a still-open Relay physical LLM scope + and corrupt the LIFO stack — "scope handle is not at the top of the + stack" → CLI EIO / redraw storm. Only joins when Relay managed + execution is actually live: when no Relay consumers are registered + there is no scope to unwind, and the join would just delay interrupt + detection (tests/run_agent/test_interrupt_propagation.py). + """ + try: + from agent import relay_runtime + + runtime = relay_runtime.get_runtime(create=False) + if runtime is None or not runtime.managed_execution_enabled(): + return + except Exception: + return + worker.join(timeout=2.0) + if worker.is_alive(): + logger.warning( + "%s worker still alive after interrupt abort (2.0s join " + "timeout); Relay teardown will best-effort drain orphaned " + "scopes (#81521).", + label, + ) + + def _ra(): """Lazy ``run_agent`` reference. @@ -1724,15 +1753,9 @@ def interruptible_api_call(agent, api_kwargs: dict): # #81521 (sibling of the streaming-path fix): wait for the worker # to unwind Relay-managed scopes before surfacing # InterruptedError, so turn teardown cannot race a still-open - # physical scope and corrupt the LIFO stack. - t.join(timeout=2.0) - if t.is_alive(): - logger.warning( - "Non-streaming worker still alive after interrupt abort " - "(%.1fs join timeout); Relay teardown will best-effort " - "drain orphaned scopes (#81521).", - 2.0, - ) + # physical scope and corrupt the LIFO stack. No-op when Relay + # managed execution is not live. + _join_worker_for_relay_teardown(t, label="Non-streaming") raise InterruptedError("Agent interrupted during API call") if result["error"] is not None: raise result["error"] @@ -3426,16 +3449,9 @@ def interruptible_streaming_api_call(agent, api_kwargs: dict, *, on_first_delta= # #81521 (sibling of the main streaming-path fix): give # the Bedrock worker a bounded window to unwind its # Relay-managed stream scopes before surfacing - # InterruptedError, so turn teardown cannot race a - # still-open physical scope and corrupt the LIFO stack. - t.join(timeout=2.0) - if t.is_alive(): - logger.warning( - "Bedrock streaming worker still alive after " - "interrupt (%.1fs join timeout); Relay teardown " - "will best-effort drain orphaned scopes (#81521).", - 2.0, - ) + # InterruptedError. No-op when Relay managed execution + # is not live. + _join_worker_for_relay_teardown(t, label="Bedrock streaming") raise InterruptedError("Agent interrupted during Bedrock API call") # Liveness watchdog: no Bedrock event for longer than the stale # timeout means the stream has wedged (open socket, keep-alives but @@ -5098,15 +5114,9 @@ def interruptible_streaming_api_call(agent, api_kwargs: dict, *, on_first_delta= # (finish_logical_calls / end_turn / close_session) race a # still-open physical scope and corrupt the LIFO stack — # "scope handle is not at the top of the stack" → CLI EIO / - # redraw storm (#81521). - t.join(timeout=2.0) - if t.is_alive(): - logger.warning( - "Streaming worker still alive after interrupt abort " - "(%.1fs join timeout); Relay teardown will best-effort " - "drain orphaned scopes (#81521).", - 2.0, - ) + # redraw storm (#81521). No-op when Relay managed execution + # is not live. + _join_worker_for_relay_teardown(t, label="Streaming") raise InterruptedError("Agent interrupted during streaming API call") # Worker thread exited before the main thread's poll loop could check # the interrupt flag. If the worker returned early due to an interrupt diff --git a/tests/run_agent/test_stream_interrupt_retry.py b/tests/run_agent/test_stream_interrupt_retry.py index acb618175b..96fb9c9cfa 100644 --- a/tests/run_agent/test_stream_interrupt_retry.py +++ b/tests/run_agent/test_stream_interrupt_retry.py @@ -272,6 +272,21 @@ class TestStreamInterruptJoinsWorkerBeforeRaise: monkeypatch.setattr(threading.Thread, "join", spy_join) + # The join is gated on Relay managed execution being live (it is + # pointless — and delays interrupt detection — when no Relay + # consumers are registered). Simulate a live runtime. + from agent import relay_runtime as rr + + monkeypatch.setattr( + rr, + "get_runtime", + lambda *a, **k: SimpleNamespace( + managed_execution_enabled=lambda: True, + get_session=lambda *a2, **k2: None, + ensure_session=lambda *a2, **k2: None, + ), + ) + class HangUntilClosedStream: response = SimpleNamespace(headers={})