fix(tui-gateway): claim disconnect sessions before teardown

This commit is contained in:
konsisumer
2026-08-26 04:54:55 +02:00
committed by Teknium
parent 9dfbde19db
commit c760143935
2 changed files with 99 additions and 71 deletions
+49 -30
View File
@@ -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()
+50 -41
View File
@@ -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)