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.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
Reference in New Issue
Block a user