From 0c9d8ab0ceb31df5af55b6e63da768519a583689 Mon Sep 17 00:00:00 2001 From: HexLab98 Date: Sat, 8 Aug 2026 15:26:46 +0700 Subject: [PATCH] fix(streaming): join stream worker and drain Relay scopes on interrupt Empty-stream stalls that trip interrupt were raising InterruptedError before the stream worker closed its physical LLM scope, corrupting the Relay LIFO stack and cascading into a CLI EIO redraw storm (#81521). --- agent/chat_completion_helpers.py | 15 +++ agent/relay_runtime.py | 206 ++++++++++++++++++++++--------- cli.py | 79 +++++++++++- 3 files changed, 240 insertions(+), 60 deletions(-) diff --git a/agent/chat_completion_helpers.py b/agent/chat_completion_helpers.py index 46be9f3e83..4a15d2f62d 100644 --- a/agent/chat_completion_helpers.py +++ b/agent/chat_completion_helpers.py @@ -5067,6 +5067,21 @@ def interruptible_streaming_api_call(agent, api_kwargs: dict, *, on_first_delta= _close_request_client_once("stream_interrupt_abort") except Exception: pass + # Wait for the worker to unwind Relay-managed stream scopes + # (physical LLM + deferred logical) before surfacing + # InterruptedError. Raising immediately lets turn teardown + # (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, + ) 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/agent/relay_runtime.py b/agent/relay_runtime.py index a1af9385ae..ec7d016272 100644 --- a/agent/relay_runtime.py +++ b/agent/relay_runtime.py @@ -543,6 +543,114 @@ class RelayRuntime: ) return result if isinstance(result, dict) else args + def _close_scope_handle( + self, + session: RelaySession, + handle: Any, + *, + output: dict[str, Any] | None = None, + allow_closing: bool = False, + failure_label: str = "scope close failed", + drain_limit: int = 32, + ) -> str | None: + """Pop ``handle``, draining orphaned children in the same session context. + + Relay scopes are strict LIFO. Empty-stream retries + interrupt can + abandon a physical LLM scope above TURN/SESSION (#81521). Drain and + close must run inside one ``run_in_session`` callback so ContextVar + stack views stay consistent across pops. + """ + if handle is None: + return None + metadata = { + RUNTIME_SCHEMA_KEY: RUNTIME_SCHEMA_VERSION, + RUNTIME_INSTANCE_KEY: self.runtime_id, + } + close_output = output or {} + session_root = session.handle + drained_holder = {"count": 0} + error_holder: dict[str, BaseException] = {} + + def close_with_drain() -> None: + def current_top() -> Any: + top = self.relay.get_scope_stack() + # Some Relay builds return the live stack (list). Others + # return the top handle directly — including tuple handles + # like ("scope", name, serial) from the test fake. Only + # unwrap real list stacks; never index a handle tuple. + if isinstance(top, list): + return top[-1] if top else None + return top + + try: + self.relay.scope.pop( + handle, + output=close_output, + metadata=metadata, + ) + return + except Exception as first_exc: + error_holder["first"] = first_exc + + for _ in range(drain_limit): + top = current_top() + if top is None or top == handle: + break + # Never pop the session root while draining for a nested handle. + if ( + session_root is not None + and top == session_root + and handle is not session_root + ): + break + try: + self.relay.scope.pop( + top, + output={ + "outcome": "cancelled", + "hermes.orphan_drain": True, + }, + metadata=metadata, + ) + drained_holder["count"] += 1 + except Exception as drain_exc: + error_holder["drain"] = drain_exc + logger.warning( + "Hermes Relay orphaned scope drain failed", + exc_info=True, + ) + break + + if drained_holder["count"]: + logger.warning( + "Hermes Relay drained %d orphaned scope(s) before closing %s", + drained_holder["count"], + handle, + ) + try: + self.relay.scope.pop( + handle, + output=close_output, + metadata=metadata, + ) + error_holder.pop("first", None) + error_holder.pop("drain", None) + except Exception as retry_exc: + error_holder["retry"] = retry_exc + + try: + self.run_in_session( + session, + close_with_drain, + allow_closing=allow_closing, + ) + except Exception as exc: + return f"{failure_label}: {exc}" + retry_exc = error_holder.get("retry") or error_holder.get("first") + if retry_exc is not None: + return f"{failure_label}: {retry_exc}" + return None + def close_session(self, event: dict[str, Any]) -> None: """Close one session scope and remove it from the core registry.""" session_id = _session_id(event) @@ -559,21 +667,15 @@ class RelayRuntime: return session.closing = True if session.handle is not None: - try: - self.run_in_session( - session, - self.relay.scope.pop, - session.handle, - output={}, - metadata={ - RUNTIME_SCHEMA_KEY: RUNTIME_SCHEMA_VERSION, - RUNTIME_INSTANCE_KEY: self.runtime_id, - }, - allow_closing=True, - timeout=_SCOPE_OP_TIMEOUT, - ) - except Exception as exc: - failures.append(f"session scope close failed: {exc}") + failure = self._close_scope_handle( + session, + session.handle, + output={}, + allow_closing=True, + failure_label="session scope close failed", + ) + if failure: + failures.append(failure) try: try: _scope_op_executor().submit( @@ -964,21 +1066,16 @@ class RelaySessionCoordinator: if isinstance(lease.host, RelayRuntime) and lease.session is not None: self._finish_logical_calls(turn, outcome=outcome) if turn.handle is not None: - try: - lease.host.run_in_session( - lease.session, - lease.host.relay.scope.pop, - turn.handle, - output={"outcome": outcome}, - metadata={ - RUNTIME_SCHEMA_KEY: RUNTIME_SCHEMA_VERSION, - RUNTIME_INSTANCE_KEY: lease.host.runtime_id, - }, - timeout=_SCOPE_OP_TIMEOUT, - ) - except Exception: + failure = lease.host._close_scope_handle( + lease.session, + turn.handle, + output={"outcome": outcome}, + failure_label="turn scope close failed", + ) + if failure: logger.warning( - "Hermes Relay turn finalization failed", exc_info=True + "Hermes Relay turn finalization failed: %s", + failure, ) finally: try: @@ -1153,35 +1250,28 @@ class RelaySessionCoordinator: turn.logical_llm_calls.clear() for index in range(len(logical_calls) - 1, -1, -1): request_id, logical_handle = logical_calls[index] - try: - lease.host.run_in_session( - lease.session, - lease.host.relay.scope.pop, - logical_handle, - output={"outcome": outcome}, - metadata={ - RUNTIME_SCHEMA_KEY: RUNTIME_SCHEMA_VERSION, - RUNTIME_INSTANCE_KEY: lease.host.runtime_id, - }, - timeout=_SCOPE_OP_TIMEOUT, - ) - except Exception: - with turn.logical_llm_lock: - # Relay scopes are stack-owned. If the newest remaining - # handle cannot close, older handles cannot close safely - # either, so retain the unclosed prefix for diagnostics. - for pending_request_id, pending_handle in logical_calls[ - : index + 1 - ]: - turn.logical_llm_calls.setdefault( - pending_request_id, - pending_handle, - ) - logger.warning( - "Hermes Relay logical LLM finalization failed", - exc_info=True, - ) - break + failure = lease.host._close_scope_handle( + lease.session, + logical_handle, + output={"outcome": outcome}, + failure_label="logical LLM scope close failed", + ) + if failure is None: + continue + with turn.logical_llm_lock: + # Relay scopes are stack-owned. If the newest remaining + # handle cannot close even after orphan drain, older + # handles cannot close safely either — retain the + # unclosed prefix for diagnostics (#81521). + for pending_request_id, pending_handle in logical_calls[ + : index + 1 + ]: + turn.logical_llm_calls.setdefault( + pending_request_id, + pending_handle, + ) + logger.warning("Hermes Relay logical LLM finalization failed: %s", failure) + break @staticmethod def _reset_turn_context(turn: RelayTurnContext) -> None: diff --git a/cli.py b/cli.py index 0423450370..8b9661f023 100644 --- a/cli.py +++ b/cli.py @@ -4859,6 +4859,10 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin): # don't auto-queue another continuation on top of a user-cancelled # turn (which would make Ctrl+C feel like it did nothing). self._last_turn_interrupted = False + # When stdout/PTY raises EIO (broken pipe after a stream-stall + # interrupt), freeze further UI paints so we don't spin the main + # thread at hundreds of escape-sequence writes/sec (#81521). + self._terminal_io_broken = False self._should_exit = False # /exit --delete: when True, the current session's SQLite history and # on-disk transcripts are deleted during shutdown. Set by @@ -5024,6 +5028,20 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin): finally: self._active_session_lease = None + def _mark_terminal_io_broken(self, reason: str = "") -> None: + """Stop UI paints after the PTY/stdout becomes unusable (#81521).""" + if getattr(self, "_terminal_io_broken", False): + return + self._terminal_io_broken = True + try: + self._pet_stop_anim() + except Exception: + pass + logger.warning( + "Terminal I/O broken%s — freezing UI paints to avoid redraw storm (#81521)", + f" ({reason})" if reason else "", + ) + def _invalidate(self, min_interval: float = 0.25) -> None: """Throttled UI repaint for high-frequency background updates. @@ -5041,12 +5059,20 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin): within the 250ms window — or an in-flight resize — silently drop it, so the prompt never renders and times out unseen (#41098). """ + if getattr(self, "_terminal_io_broken", False): + return if getattr(self, "_resize_recovery_pending", False): return now = time.monotonic() if hasattr(self, "_app") and self._app and (now - getattr(self, "_last_invalidate", 0.0)) >= min_interval: self._last_invalidate = now - self._app.invalidate() + try: + self._app.invalidate() + except OSError as exc: + if getattr(exc, "errno", None) == errno.EIO: + self._mark_terminal_io_broken("invalidate") + return + raise def _paint_now(self) -> None: """Immediate, unthrottled repaint for user-blocking modal prompts. @@ -5059,10 +5085,17 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin): already use. See ``_invalidate`` for why the throttle must not gate these paints (#41098). """ + if getattr(self, "_terminal_io_broken", False): + return app = getattr(self, "_app", None) if app is not None: try: app.invalidate() + except OSError as exc: + if getattr(exc, "errno", None) == errno.EIO: + self._mark_terminal_io_broken("paint_now") + return + raise except Exception: pass @@ -5081,6 +5114,8 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin): matching the standard terminal-UX convention (bash, zsh, fish, vim, htop). """ + if getattr(self, "_terminal_io_broken", False): + return app = getattr(self, "_app", None) if not app: return @@ -5088,9 +5123,16 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin): app, rebuild_scrollback=self._redraw_rebuilds_scrollback(), ) + if getattr(self, "_terminal_io_broken", False): + return _replay_output_history() try: app.invalidate() + except OSError as exc: + if getattr(exc, "errno", None) == errno.EIO: + self._mark_terminal_io_broken("force_full_redraw") + return + raise except Exception: pass @@ -5155,8 +5197,11 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin): screen/cursor state and forces a clean repaint. Both steps are independently safe and self-guard, so a failure of one - never prevents the other. + never prevents the other. If the PTY is already dead (EIO), skip the + redraw entirely — painting a broken fd is the #81521 redraw storm. """ + if getattr(self, "_terminal_io_broken", False): + return try: from hermes_cli.curses_ui import flush_stdin flush_stdin() @@ -5173,6 +5218,8 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin): def _clear_prompt_toolkit_screen(self, app, *, rebuild_scrollback: bool = False) -> None: """Clear the terminal and reset prompt_toolkit renderer state.""" + if getattr(self, "_terminal_io_broken", False): + return try: renderer = app.renderer out = renderer.output @@ -5189,6 +5236,11 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin): # next _redraw() starts from a known (0, 0) origin and # re-renders every cell rather than diffing against stale. renderer.reset(leave_alternate_screen=False) + except OSError as exc: + if getattr(exc, "errno", None) == errno.EIO: + self._mark_terminal_io_broken("clear_screen") + return + pass except Exception: pass @@ -6224,6 +6276,9 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin): """Advance the frame + invalidate on a timer while a pet is enabled.""" while self._pet_anim_running: time.sleep(self._PET_FRAME_INTERVAL) + if getattr(self, "_terminal_io_broken", False): + self._pet_anim_running = False + break now = time.monotonic() if now - self._pet_cfg_checked >= self._PET_CFG_INTERVAL: self._pet_cfg_checked = now @@ -6236,6 +6291,10 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin): if app is not None: try: app.invalidate() + except OSError as exc: + if getattr(exc, "errno", None) == errno.EIO: + self._mark_terminal_io_broken("pet_anim") + break except Exception: pass @@ -18668,7 +18727,23 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin): except Exception: pass # Non-fatal — don't break the main loop + except OSError as e: + if getattr(e, "errno", None) == errno.EIO: + self._mark_terminal_io_broken("process_loop") + logger.warning( + "process_loop EIO — freezing UI paints (#81521): %s", + e, + ) + continue + logger.warning("process_loop unhandled error (msg may be lost): %s", e) except Exception as e: + if isinstance(e, OSError) and getattr(e, "errno", None) == errno.EIO: + self._mark_terminal_io_broken("process_loop") + logger.warning( + "process_loop EIO — freezing UI paints (#81521): %s", + e, + ) + continue logger.warning("process_loop unhandled error (msg may be lost): %s", e) # Start processing thread