From 8109c400a2e46220f424ce84febcb67a5d7637b5 Mon Sep 17 00:00:00 2001 From: teknium1 <127238744+teknium1@users.noreply.github.com> Date: Sun, 20 Sep 2026 00:02:13 -0700 Subject: [PATCH] fix: late resume/submit from a closed socket keeps the ws-orphan reap armed A WebSocket that drops right after sending session.resume or prompt.submit has its disconnect cleanup (_close_sessions_for_transport) run before the RPC's rebind. The rebind then registered the already-closed socket and cancelled the pending orphan reap; nothing ever detaches that socket again, so the detached session kept its active-session lease (surface=ios Bot Chat) with no Timer until the 300 s idle-reaper repair sweep, and the canonical Bot Chat stayed "connecting" for the next client (#116464). _rebind_live_transport now treats a dead rebinding transport as "client not back": it leaves the reap armed (re-arming it when the caller already cancelled, as the reuse fast path does) and does not register the dead socket as a viewer. prompt.submit goes through the same seam instead of its own attach + unconditional cancel. Live: loopback dashboard rig, grace 3 s, ios session, owning socket aborted mid-submit + 12 resume-then-drop sockets. Before: lease still held 15 s later, no reap logged. After: lease released within the grace, reap logged, fresh resume returns the stored history. Control (plain 1001 close + reply- awaiting resume/close burst) still releases. --- tests/tui_gateway/test_ws_orphan_races.py | 60 +++++++++++++++++++++++ tui_gateway/methods_prompt.py | 6 +-- tui_gateway/session_lifecycle.py | 8 +++ 3 files changed, 71 insertions(+), 3 deletions(-) diff --git a/tests/tui_gateway/test_ws_orphan_races.py b/tests/tui_gateway/test_ws_orphan_races.py index b73af43d80..30a3870971 100644 --- a/tests/tui_gateway/test_ws_orphan_races.py +++ b/tests/tui_gateway/test_ws_orphan_races.py @@ -207,3 +207,63 @@ def test_reconnect_cannot_cross_orphan_interrupt_claim(monkeypatch, path, claim) assert session["transport"] is server._detached_ws_transport assert sid in server._pending_ws_reaps assert session["queued_prompt"] is None + + +@pytest.mark.parametrize("path", ["unpersisted", "reuse", "prompt"]) +@pytest.mark.parametrize("socket", ["closed", "live"]) +def test_late_rpc_from_closed_socket_keeps_orphan_reap_armed(monkeypatch, path, socket): + """A resume/submit whose socket already closed (disconnect cleanup ran first, so nothing detaches it + again) must not cancel the orphan reap: the client is not back. A live socket still cancels it.""" + sid = "late-rebind" + session = dict(transport=server._detached_ws_transport, running=True, history_lock=threading.Lock(), + history=[], session_key="stored", agent=SimpleNamespace(model="test"), queued_prompt=None) + timers = [] + + class Timer: + def __init__(self, delay, callback): + timers.append(self) + + def start(self): + pass + + def cancel(self): + pass + + class Socket: + _closed = socket == "closed" + + def send(self, *a, **kw): + pass + + transport = Socket() + monkeypatch.setattr(server, "_sessions", {sid: session}) + monkeypatch.setattr(server, "_pending_ws_reaps", {sid: Mock()}) + monkeypatch.setattr(server.threading, "Timer", Timer) + monkeypatch.setattr(server, "_WS_ORPHAN_REAP_GRACE_S", 20) + monkeypatch.setattr(server, "current_transport", lambda: transport) + monkeypatch.setattr(server, "_resolve_model", lambda: "test") + monkeypatch.setattr(server, "_ensure_active_session_slot", lambda *a: None) + monkeypatch.setattr(server, "_legacy_group_fence_error", lambda *a: None) + monkeypatch.setattr(server, "_session_uses_compute_host", lambda *a: False) + monkeypatch.setattr(server, "_load_dashboard_process_isolation_config", lambda: {}) + monkeypatch.setattr(server, "_handle_busy_submit", lambda *a, **kw: {"result": {"status": "queued"}}) + monkeypatch.setattr(server, "_sess", lambda *a: (session, None)) + ctx = SimpleNamespace(rid=1, owns_db=False, db=None, cols=80, omit_messages=True, + defer_history=False, target="stored", profile=None, + profile_home=None, profile_resume_cwd=None, found={}, + messages=lambda history: [], mint=lambda: ("unused", "tui", "."), + restore=lambda: ([], [], []), display_prefix=lambda: []) + if path == "unpersisted": + response = server._resume_live_unpersisted(ctx, sid, session) + elif path == "reuse": + response = server._resume_reuse_live(ctx, sid, session) + else: + response = server.handle_request({"jsonrpc": "2.0", "id": 1, "method": "prompt.submit", + "params": {"session_id": sid, "text": "continue"}}) + assert "error" not in response + if socket == "closed": + assert session["transport"] is server._detached_ws_transport + assert sid in server._pending_ws_reaps # left armed (unpersisted/prompt) or re-armed (reuse cancels first) + else: + assert session["transport"] is transport + assert sid not in server._pending_ws_reaps diff --git a/tui_gateway/methods_prompt.py b/tui_gateway/methods_prompt.py index 093aa906a8..18b21cbc2b 100644 --- a/tui_gateway/methods_prompt.py +++ b/tui_gateway/methods_prompt.py @@ -621,13 +621,13 @@ def _(rid, params: dict) -> dict: if internal_hosted_submit and turn_isolation: return _err(rid, 4121, "hosted room turns do not support isolated compute workers yet") # Re-bind to the current transport: streaming must stay on the active websocket even - # if a disconnect/fallback moved the session to stdio. + # if a disconnect/fallback moved the session to stdio. Through _rebind_live_transport so a + # socket that already closed cannot cancel the orphan reap without coming back (#116464). with _session_resume_lock: if (refusal := _reattach_refusal(rid, sid, session)) is not None: return refusal if (t := current_transport()) is not None: - _attach_session_transport(session, t) - _cancel_ws_orphan_reap(sid) + _rebind_live_transport(sid, session, t) # Claim the turn against a possibly-running session (busy/queued reply, else fall # through once ``running`` is observed False). The provider interrupt happens after # history_lock is released (a non-interruptible tool may hold it); if the old turn diff --git a/tui_gateway/session_lifecycle.py b/tui_gateway/session_lifecycle.py index a200af8991..0beeb087b6 100644 --- a/tui_gateway/session_lifecycle.py +++ b/tui_gateway/session_lifecycle.py @@ -703,6 +703,14 @@ def _rebind_live_transport(sid: str, session: dict, transport: Transport) -> Non """Attach a live peer without displacing existing subscribers (caller holds ``history_lock``). Subagent control authority needs no bookkeeping here: it resolves against ``session["transport"]`` at RPC time (``tools.delegate_tool_registry._subagent_transport_matches``).""" + if transport is not _detached_ws_transport and _transport_is_dead(transport): + # The rebinding socket already closed: its disconnect cleanup ran before this late RPC (a + # resume-then-drop burst), so nothing will detach it again. The client is NOT back — re-arm the + # reap the caller cancelled instead of leaving a detached session with no Timer (#116464). + with _sessions_lock: + if _ws_session_is_detached(session) and sid not in _pending_ws_reaps: + _schedule_ws_orphan_reap(sid) + return _attach_session_transport(session, transport) # Every transport that showed this session (pop-outs resume the same sid); on disconnect the last # viewer becomes the transport instead of the drop sentinel.