diff --git a/tests/test_tui_gateway_server.py b/tests/test_tui_gateway_server.py index 64d9157c16..deba3bd53e 100644 --- a/tests/test_tui_gateway_server.py +++ b/tests/test_tui_gateway_server.py @@ -4881,8 +4881,8 @@ def test_ws_disconnect_running_sidecar_still_closes_without_orphan_timer(monkeyp ) monkeypatch.setattr( server, - "_close_session_by_id", - lambda sid, *, end_reason: closed.append((sid, end_reason)) or True, + "_teardown_popped_session", + lambda session, *, end_reason: closed.append((session["_sid"], end_reason)) or True, ) monkeypatch.setattr( server, "_schedule_ws_orphan_reap", lambda sid: scheduled.append(sid) @@ -17707,8 +17707,9 @@ def test_session_close_rpc_claims_then_tears_down(monkeypatch): def test_close_sessions_for_transport_closes_flagged_repoints_rest(monkeypatch): seen = [] monkeypatch.setattr( - server, "_close_session_by_id", - lambda sid, *, end_reason: bool(seen.append((sid, end_reason))) or True, + server, + "_teardown_popped_session", + lambda session, *, end_reason: seen.append((session["_sid"], end_reason)) or True, ) # Detached session "b" would schedule a real grace-reap threading.Timer that # outlives the test; grace=0 short-circuits it so no thread lingers. @@ -17725,46 +17726,64 @@ def test_close_sessions_for_transport_closes_flagged_repoints_rest(monkeypatch): server._sessions.clear() -def test_close_sessions_for_transport_skips_rebound_session(monkeypatch): - """Rebind-between-snapshot-and-stomp (#77129 concept salvage). - - _close_sessions_for_transport snapshots owned sessions under - _sessions_lock, then parks each on the drop sentinel. A concurrent - session.resume that rebinds the session to a NEW live transport in - between must NOT be stomped back onto the sentinel — that knocks an - attached client into detached state and arms an orphan reap against a - session with a live owner. The stomp must revalidate ownership under - the lock and skip (park AND reap) when the transport already moved on. - """ +@pytest.mark.parametrize("close_on_disconnect", [True, False]) +def test_close_sessions_for_transport_skips_session_rebound_before_claim( + monkeypatch, close_on_disconnect +): + """A resume between snapshot and claim keeps either session type alive.""" reaps = [] + teardowns = [] monkeypatch.setattr( server, "_schedule_ws_orphan_reap", lambda sid: reaps.append(sid) ) + monkeypatch.setattr( + server, + "_teardown_popped_session", + lambda session, *, end_reason: teardowns.append((session, end_reason)) or True, + ) old_transport = object() # the disconnecting transport new_transport = object() # live rebind target (no _closed attr → alive) + session = {"transport": old_transport, "close_on_disconnect": close_on_disconnect} + original_sessions_lock = server._sessions_lock + rebound = threading.Event() - class _RebindsOnStomp(dict): - """Simulates a session.resume landing between snapshot and stomp: - the first 'viewers' read inside the stomp loop (i.e. after the - snapshot already selected this session) rebinds the transport.""" + class _SnapshotInterlock: + """Rebind in a second thread immediately after the ownership snapshot.""" - def get(self, key, default=None): - if key == "viewers" and not self.get("_rebound_flag"): - dict.__setitem__(self, "_rebound_flag", True) - dict.__setitem__(self, "transport", new_transport) - return dict.get(self, key, default) + def __init__(self): + self._snapshot_released = False - session = _RebindsOnStomp( - {"transport": old_transport, "close_on_disconnect": False} - ) + def __enter__(self): + original_sessions_lock.acquire() + return self + + def __exit__(self, exc_type, exc, traceback): + original_sessions_lock.release() + if not self._snapshot_released: + self._snapshot_released = True + + def _resume_rebind(): + with server._session_resume_lock: + session["transport"] = new_transport + rebound.set() + + thread = threading.Thread(target=_resume_rebind) + thread.start() + assert rebound.wait(timeout=1) + thread.join(timeout=1) + return False + + monkeypatch.setattr(server, "_sessions_lock", _SnapshotInterlock()) server._sessions.clear() server._sessions["rebound"] = session try: reaped, detached = server._close_sessions_for_transport(old_transport) assert reaped == 0 - assert detached == 0 # skipped, not parked - assert session["transport"] is new_transport # rebind preserved - assert reaps == [] # no orphan reap armed against the live owner + assert detached == 0 + assert server._sessions["rebound"] is session + assert session["transport"] is new_transport + assert teardowns == [] + assert reaps == [] finally: server._sessions.clear() diff --git a/tui_gateway/server.py b/tui_gateway/server.py index f110cc0d8f..8486dcd99e 100644 --- a/tui_gateway/server.py +++ b/tui_gateway/server.py @@ -1374,9 +1374,8 @@ def _close_sessions_for_transport( transport, *, end_reason: str = "ws_disconnect" ) -> tuple[int, int]: """On transport disconnect, reap the sessions that opted into - close_on_disconnect (sidecar/dashboard) immediately via the unified - ``_close_session_by_id`` path, and re-point the rest back to stdio so later - emits don't hit a dead socket. + close_on_disconnect (sidecar/dashboard) immediately and re-point the rest + at the detached transport so later emits don't hit a dead socket. Non-flagged detached sessions are handed to the grace-windowed WS-orphan reaper (``_schedule_ws_orphan_reap``): a quick reconnect / session.resume @@ -1391,47 +1390,57 @@ def _close_sessions_for_transport( reaped = 0 detached = 0 for sid, session in owned: - if session.get("close_on_disconnect"): - _close_session_by_id(sid, end_reason=end_reason) - reaped += 1 - else: - # Point detached sessions at the drop sentinel (NOT real stdio) so - # _ws_session_is_orphaned recognizes them and the grace-reap can - # actually fire; a standalone `hermes --tui` keeps real _stdio. - # UNLESS another window still shows the session: multi-window - # pop-outs all register as viewers, so on disconnect re-bind the - # session to the most recent surviving viewer instead of - # stranding the original window on the sentinel (#83716). - viewers = session.get("viewers") - if viewers: - viewers.pop(transport, None) - # Revalidate under the sessions lock before stomping (#77129): - # between the owned-sessions snapshot above and this write, a - # concurrent session.resume can rebind the session to a NEW live - # transport. Stomping it back onto the drop sentinel here would - # knock an attached client into detached state and arm an orphan - # reap against a session that has a live owner. If the transport - # already moved on to a different live transport, this disconnect - # has nothing left to tear down — skip the park AND the reap. + claimed_for_teardown = None + should_schedule_reap = False + # A session.resume fast-path rebinds its live session while holding + # _session_resume_lock. Take that lock before re-checking the snapshot + # so a reconnect cannot move the transport between this check and the + # close/detach ownership claim. Keep the slow teardown below both locks. + with _session_resume_lock: with _sessions_lock: - current = session.get("transport") - if ( - current is not transport - and current is not None - and not _transport_is_dead(current) - ): + current = _sessions.get(sid) + if current is not session: continue - remaining = [ - (ts, v) - for v, ts in (viewers or {}).items() - if v is not transport and not _transport_is_dead(v) - ] - if remaining: - remaining.sort(key=lambda kv: kv[0]) - session["transport"] = remaining[-1][1] + if current.get("transport") is not transport: + # The reconnect owns this session now. Drop only the old + # viewer registration; it must not affect the new owner. + viewers = current.get("viewers") + if viewers: + viewers.pop(transport, None) continue - session["transport"] = _detached_ws_transport - session.pop("_client_gone_interrupt_requested", None) + if current.get("close_on_disconnect"): + claimed_for_teardown = _pop_session_by_id(sid) + else: + # Point detached sessions at the drop sentinel (NOT real + # stdio) so _ws_session_is_orphaned recognizes them and + # the grace-reap can actually fire; a standalone + # `hermes --tui` keeps real _stdio. UNLESS another window + # still shows the session: multi-window pop-outs all + # register as viewers, so on disconnect re-bind the + # session to the most recent surviving viewer instead of + # stranding the original window on the sentinel (#83716). + viewers = current.get("viewers") + if viewers: + viewers.pop(transport, None) + remaining = [ + (ts, viewer_transport) + for viewer_transport, ts in (viewers or {}).items() + if ( + viewer_transport is not transport + and not _transport_is_dead(viewer_transport) + ) + ] + if remaining: + remaining.sort(key=lambda item: item[0]) + current["transport"] = remaining[-1][1] + else: + current["transport"] = _detached_ws_transport + current.pop("_client_gone_interrupt_requested", None) + should_schedule_reap = True + if claimed_for_teardown is not None: + if _teardown_popped_session(claimed_for_teardown, end_reason=end_reason): + reaped += 1 + elif should_schedule_reap: detached += 1 try: _schedule_ws_orphan_reap(sid)