From 17ec0a3702f5eb26f1dcde5176225139954d20a9 Mon Sep 17 00:00:00 2001 From: Josh Dow Date: Mon, 8 Jun 2026 07:58:41 -0600 Subject: [PATCH] fix(tui): avoid stdio fallback for detached websocket sessions --- tests/test_tui_gateway_server.py | 6 +-- tests/tui_gateway/test_protocol.py | 53 ++++++++++++++++++++++++++ tui_gateway/server.py | 20 +++++++--- tui_gateway/transport.py | 19 ++++++++++ tui_gateway/ws.py | 60 +++++++++++++++++++----------- 5 files changed, 127 insertions(+), 31 deletions(-) diff --git a/tests/test_tui_gateway_server.py b/tests/test_tui_gateway_server.py index a25d713b80a45..26b7c20fff5c8 100644 --- a/tests/test_tui_gateway_server.py +++ b/tests/test_tui_gateway_server.py @@ -933,7 +933,7 @@ def close(self): closed["worker"] = True server._sessions["orphan-sid"] = _session( - transport=server._stdio_transport, + transport=server._detached_ws_transport, slash_worker=_FakeWorker(), running=False, ) @@ -961,11 +961,11 @@ def write(self, *a, **k): assert server._ws_session_is_orphaned(reattached) is False # Mid-turn sessions are also spared even if detached. - mid_turn = _session(transport=server._stdio_transport, running=True) + mid_turn = _session(transport=server._detached_ws_transport, running=True) assert server._ws_session_is_orphaned(mid_turn) is False # Already finalized sessions are spared (idempotency). - done = _session(transport=server._stdio_transport, running=False, _finalized=True) + done = _session(transport=server._detached_ws_transport, running=False, _finalized=True) assert server._ws_session_is_orphaned(done) is False diff --git a/tests/tui_gateway/test_protocol.py b/tests/tui_gateway/test_protocol.py index 2fec6617d62d1..e297572c53c60 100644 --- a/tests/tui_gateway/test_protocol.py +++ b/tests/tui_gateway/test_protocol.py @@ -514,6 +514,59 @@ def resume_second(): assert all(sid == winner for sid in server._sessions) +def test_ws_detach_drops_running_session_events_instead_of_falling_back_to_stdio( + server, monkeypatch +): + """A detached WS-owned running session must not stream raw events to stdout. + + The dashboard service runs the gateway in-process with stdout connected to + journald, not a TUI reader. Falling back a disconnected WS session to + ``_stdio_transport`` leaks live stream frames into the service log and leaves + Desktop out of sync until resume rebinds the session. + """ + + import importlib + + ws_mod = importlib.import_module("tui_gateway.ws") + + class _ClosedWsTransport: + def write(self, _obj): + raise AssertionError("closed websocket transport should not be used") + + def close(self): + pass + + closed_transport = _ClosedWsTransport() + stdout = io.StringIO() + server._real_stdout = stdout + scheduled = [] + monkeypatch.setattr(server, "_schedule_ws_orphan_reap", scheduled.append) + + server._sessions["live"] = { + "history_lock": threading.RLock(), + "running": True, + "transport": closed_transport, + } + + detached, reaped = ws_mod._detach_sessions_for_transport( + closed_transport, + peer="test-peer", + ) + + assert detached == 1 + assert reaped == 1 + assert scheduled == ["live"] + assert server._sessions["live"]["transport"] is not server._stdio_transport + + server._emit("reasoning.delta", "live", {"text": "must-not-hit-journal"}) + + assert stdout.getvalue() == "" + assert not server._ws_session_is_orphaned(server._sessions["live"]) + + server._sessions["live"]["running"] = False + assert server._ws_session_is_orphaned(server._sessions["live"]) + + def test_session_resume_live_payload_uses_current_history_with_ancestors(server, monkeypatch): """Live resume should not reuse a stale ancestor-inclusive snapshot.""" diff --git a/tui_gateway/server.py b/tui_gateway/server.py index 4dd9592a77b58..08311ed1a590d 100644 --- a/tui_gateway/server.py +++ b/tui_gateway/server.py @@ -25,6 +25,7 @@ from hermes_cli.env_loader import load_hermes_dotenv from utils import is_truthy_value from tui_gateway.transport import ( + DroppingTransport, StdioTransport, Transport, bind_transport, @@ -207,6 +208,12 @@ def _thread_panic_hook(args): # patches of `_real_stdout` (used extensively in tests) still land correctly. _stdio_transport = StdioTransport(lambda: _real_stdout, _stdout_lock) +# Detached WebSocket sessions are parked here until a reconnect/resume rebinds +# them to a live WS transport. Do not use _stdio_transport for these sessions: +# in the dashboard service stdout is journald, so stream frames would leak to +# service logs and Desktop would remain out of sync. +_detached_ws_transport = DroppingTransport() + class _SlashWorker: """Persistent HermesCLI subprocess for slash commands.""" @@ -383,16 +390,17 @@ def _teardown_session(session: dict | None) -> None: def _ws_session_is_orphaned(session: dict | None) -> bool: """True if a WS session has no live transport and no in-flight turn. - After ``handle_ws`` detaches a disconnected client it points the session - at ``_stdio_transport``. In the dashboard's in-process gateway there is no - real stdio peer reading those frames, so a session left on the stdio - transport (and not mid-turn) is genuinely orphaned and safe to reap. + After ``handle_ws`` detaches a disconnected client it points the session at + ``_detached_ws_transport``. That transport intentionally drops stream frames + instead of falling back to stdout/journald. Once the session is idle, it is + genuinely orphaned and safe to reap unless a reconnect/resume has rebound it + to a live client transport. """ if not session or session.get("_finalized"): return False if session.get("running"): return False - return session.get("transport") is _stdio_transport + return session.get("transport") is _detached_ws_transport def _schedule_ws_orphan_reap(sid: str) -> None: @@ -4348,7 +4356,7 @@ def _(rid, params: dict) -> dict: return err # Re-bind to the current client transport for this request. This keeps # streaming events on the active websocket even if an earlier disconnect - # or fallback moved the session transport to stdio. + # parked the session on the detached-WS dropping transport. if (t := current_transport()) is not None: session["transport"] = t with session["history_lock"]: diff --git a/tui_gateway/transport.py b/tui_gateway/transport.py index ce93e518a3d52..0090c8f1a70cb 100644 --- a/tui_gateway/transport.py +++ b/tui_gateway/transport.py @@ -183,6 +183,25 @@ def close(self) -> None: return None +class DroppingTransport: + """Transport that intentionally drops frames while reporting success. + + Used for WebSocket-owned sessions after their client disconnects. In the + dashboard process, falling back those sessions to stdio writes raw JSON-RPC + stream frames to journald instead of a real TUI reader. Dropping keeps the + running agent safe until a later ``session.resume``/``prompt.submit`` + rebinds the session to a live client transport. + """ + + __slots__ = () + + def write(self, obj: dict) -> bool: + return True + + def close(self) -> None: + return None + + class TeeTransport: """Mirrors writes to one primary plus N best-effort secondaries. diff --git a/tui_gateway/ws.py b/tui_gateway/ws.py index 1babfc1d3c2e6..7421283fad2bf 100644 --- a/tui_gateway/ws.py +++ b/tui_gateway/ws.py @@ -156,6 +156,35 @@ def _disable_nagle(ws: Any) -> None: _log.debug("ws TCP_NODELAY skip: %s", exc) +def _detach_sessions_for_transport(transport: WSTransport, *, peer: str) -> tuple[int, int]: + """Park sessions owned by a dead WS transport without falling back to stdio. + + In stdio/Ink mode the stdio transport is a real JSON-RPC peer. In the + dashboard's in-process WebSocket mode, stdout is the service log. A WS-owned + session that falls back to stdio after disconnect will stream raw event + frames to journald and Desktop will miss them. Point detached sessions at a + dropping transport until ``session.resume`` or ``prompt.submit`` rebinds a + live client transport. + """ + + detached_sessions = 0 + reaped_scheduled = 0 + for _sid, sess in list(server._sessions.items()): + if sess.get("transport") is transport: + sess["transport"] = server._detached_ws_transport + detached_sessions += 1 + try: + server._schedule_ws_orphan_reap(_sid) + reaped_scheduled += 1 + except Exception: + _log.exception( + "ws orphan-reap schedule failed peer=%s sid=%s", + peer, + _sid, + ) + return detached_sessions, reaped_scheduled + + async def handle_ws(ws: Any) -> None: """Run one WebSocket session. Wire-compatible with ``tui_gateway.entry``.""" peer = _ws_peer_label(ws) @@ -288,28 +317,15 @@ async def handle_ws(ws: Any) -> None: if transport is not None: transport.close() - # Detach the transport from any sessions it owned so later emits - # fall back to stdio instead of crashing into a closed socket. - # - # In the dashboard's in-process gateway that stdio fallback has no - # real reader, so a detached session would otherwise sit forever - # holding its _SlashWorker subprocess open (one leaked python proc - # per browser refresh — #38591 fallout). Schedule a grace-delayed - # reap; a quick reconnect / session.resume re-binds a live - # transport and cancels it (see _ws_session_is_orphaned). - for _sid, sess in list(server._sessions.items()): - if sess.get("transport") is transport: - sess["transport"] = server._stdio_transport - detached_sessions += 1 - try: - server._schedule_ws_orphan_reap(_sid) - reaped_scheduled += 1 - except Exception: - _log.exception( - "ws orphan-reap schedule failed peer=%s sid=%s", - peer, - _sid, - ) + # Detach the transport from any sessions it owned so later emits do + # not crash into a closed socket or fall through to dashboard + # stdout/journald. A quick reconnect / session.resume re-binds a + # live transport and cancels the orphan reap (see + # _ws_session_is_orphaned). + detached_sessions, reaped_scheduled = _detach_sessions_for_transport( + transport, + peer=peer, + ) try: await ws.close() except Exception as exc: