fix(tui-gateway): interrupt turns after websocket disconnect

After the existing reconnect grace, route a still-detached running session through the same interrupt mechanism as session.interrupt. Preserve delegation deferral, sidecar teardown, partial history, and single-owner reap semantics.

Verified on upstream main: RED 4 failed/2 passed without production changes; GREEN 603 related gateway/compute-host tests. Ruff and py_compile passed. Momus pass 2: APPROVE.
This commit is contained in:
Kyzcreig
2026-08-19 16:39:21 -07:00
committed by Teknium
parent 30a37668d9
commit 14b50f5edd
4 changed files with 417 additions and 100 deletions
+273
View File
@@ -4585,6 +4585,279 @@ def test_session_close_settles_active_turn_before_teardown(monkeypatch):
assert response["result"] == {"closed": True}
def test_ws_orphan_reap_interrupts_isolated_turn_then_reaps(monkeypatch):
callbacks = []
interrupted = []
torn_down = []
class _Timer:
def __init__(self, _delay, callback):
callbacks.append(callback)
self.daemon = False
def start(self):
return None
class _Supervisor:
def interrupt(self, sid, *, request_id=None):
interrupted.append((sid, request_id))
session = _session(
agent=None,
agent_ready=threading.Event(),
transport=server._detached_ws_transport,
running=True,
_compute_host_active=True,
history=[{"role": "assistant", "content": "partial"}],
queued_prompt={"text": "must not run"},
)
server._sessions["isolated-sid"] = session
monkeypatch.setattr(server, "_WS_ORPHAN_REAP_GRACE_S", 0.01)
monkeypatch.setattr(server.threading, "Timer", _Timer)
monkeypatch.setattr(
server, "_load_cfg", lambda: {"dashboard": {"turn_isolation": True}}
)
monkeypatch.setattr(
server, "_get_compute_host_supervisor", lambda _cfg=None: _Supervisor()
)
monkeypatch.setattr(
server,
"_teardown_popped_session",
lambda claimed, *, end_reason: torn_down.append((claimed, end_reason)) or True,
)
try:
server._schedule_ws_orphan_reap("isolated-sid")
callbacks.pop(0)()
assert interrupted == [("isolated-sid", "client-gone-isolated-sid")]
assert session["_turn_cancel_requested"] is True
assert session["queued_prompt"] is None
assert session["history"] == [{"role": "assistant", "content": "partial"}]
assert len(callbacks) == 1
callbacks.pop(0)()
assert interrupted == [("isolated-sid", "client-gone-isolated-sid")]
assert len(callbacks) == 1
session["running"] = False
callbacks.pop(0)()
assert "isolated-sid" not in server._sessions
assert torn_down == [(session, "ws_orphan_reap")]
finally:
server._sessions.pop("isolated-sid", None)
def test_ws_orphan_reap_spares_turn_reattached_within_grace(monkeypatch):
callbacks = []
interrupted = []
class _Timer:
def __init__(self, _delay, callback):
callbacks.append(callback)
def start(self):
return None
class _LiveThread:
def is_alive(self):
return True
class _LiveTransport:
def write(self, *_args, **_kwargs):
return True
disconnecting_transport = _LiveTransport()
session = _session(
agent=types.SimpleNamespace(
interrupt=lambda: interrupted.append("interrupted")
),
transport=disconnecting_transport,
running=True,
_run_thread=_LiveThread(),
)
server._sessions["reattached-sid"] = session
monkeypatch.setattr(server, "_WS_ORPHAN_REAP_GRACE_S", 0.01)
monkeypatch.setattr(server.threading, "Timer", _Timer)
monkeypatch.setattr(server, "_load_cfg", lambda: {})
try:
server._close_sessions_for_transport(disconnecting_transport)
assert session["transport"] is server._detached_ws_transport
session["transport"] = _LiveTransport()
callbacks.pop(0)()
assert interrupted == []
assert "reattached-sid" in server._sessions
assert callbacks == []
finally:
server._sessions.pop("reattached-sid", None)
def test_session_resume_does_not_rebind_after_client_gone_interrupt_claim(monkeypatch):
class _DB:
def get_session(self, session_id):
assert session_id == "stored-sid"
return {"id": session_id, "cwd": "/tmp"}
def resolve_resume_session_id(self, session_id):
return session_id
live_transport = object()
session = _session(
session_key="stored-sid",
transport=server._detached_ws_transport,
running=True,
_client_gone_interrupt_requested=True,
)
server._sessions["live-sid"] = session
monkeypatch.setattr(server, "_get_db", lambda: _DB())
monkeypatch.setattr(server, "current_transport", lambda: live_transport)
try:
response = server.handle_request(
{
"id": "resume-after-claim",
"method": "session.resume",
"params": {"session_id": "stored-sid"},
}
)
assert response is not None
assert response["error"]["code"] == 4009
assert response["error"]["message"] == "session disconnect interrupt settling"
assert session["transport"] is server._detached_ws_transport
finally:
server._sessions.pop("live-sid", None)
def test_ws_orphan_reap_defers_running_turn_for_active_delegation(monkeypatch):
callbacks = []
interrupted = []
delegation_active = iter((True, False, False))
class _Timer:
def __init__(self, _delay, callback):
callbacks.append(callback)
def start(self):
return None
class _LiveThread:
def is_alive(self):
return True
def _interrupt():
interrupted.append("interrupted")
session["running"] = False
session = _session(
agent=types.SimpleNamespace(interrupt=_interrupt),
transport=server._detached_ws_transport,
running=True,
_run_thread=_LiveThread(),
)
server._sessions["delegating-turn"] = session
monkeypatch.setattr(server, "_WS_ORPHAN_REAP_GRACE_S", 0.01)
monkeypatch.setattr(server.threading, "Timer", _Timer)
monkeypatch.setattr(server, "_load_cfg", lambda: {})
monkeypatch.setattr(
server,
"_session_has_active_delegations",
lambda *_args, **_kwargs: next(delegation_active),
)
monkeypatch.setattr(server, "_teardown_popped_session", lambda *_args, **_kwargs: True)
try:
server._schedule_ws_orphan_reap("delegating-turn")
callbacks.pop(0)()
assert interrupted == []
assert len(callbacks) == 1
callbacks.pop(0)()
assert interrupted == ["interrupted"]
assert len(callbacks) == 1
callbacks.pop(0)()
assert "delegating-turn" not in server._sessions
finally:
server._sessions.pop("delegating-turn", None)
def test_ws_orphan_reap_interrupts_in_process_turn(monkeypatch):
callbacks = []
interrupted = []
class _Timer:
def __init__(self, _delay, callback):
callbacks.append(callback)
def start(self):
return None
class _LiveThread:
def is_alive(self):
return True
def _interrupt():
interrupted.append("interrupted")
session["running"] = False
session = _session(
agent=types.SimpleNamespace(interrupt=_interrupt),
transport=server._detached_ws_transport,
running=True,
_run_thread=_LiveThread(),
)
server._sessions["inline-sid"] = session
monkeypatch.setattr(server, "_WS_ORPHAN_REAP_GRACE_S", 0.01)
monkeypatch.setattr(server.threading, "Timer", _Timer)
monkeypatch.setattr(server, "_load_cfg", lambda: {})
try:
server._schedule_ws_orphan_reap("inline-sid")
callbacks.pop(0)()
assert interrupted == ["interrupted"]
assert session["_turn_cancel_requested"] is True
assert len(callbacks) == 1
finally:
server._sessions.pop("inline-sid", None)
def test_ws_disconnect_running_sidecar_still_closes_without_orphan_timer(monkeypatch):
closed = []
scheduled = []
transport = object()
server._sessions["sidecar-sid"] = _session(
transport=transport,
running=True,
close_on_disconnect=True,
)
monkeypatch.setattr(
server,
"_close_session_by_id",
lambda sid, *, end_reason: closed.append((sid, end_reason)) or True,
)
monkeypatch.setattr(
server, "_schedule_ws_orphan_reap", lambda sid: scheduled.append(sid)
)
try:
reaped, detached = server._close_sessions_for_transport(transport)
assert (reaped, detached) == (1, 0)
assert closed == [("sidecar-sid", "ws_disconnect")]
assert scheduled == []
finally:
server._sessions.pop("sidecar-sid", None)
def test_ws_orphan_reap_closes_worker_when_session_stays_detached(monkeypatch):
"""A detached WS session past its grace window has its slash_worker closed.
@@ -159,7 +159,14 @@ def test_resume_closes_profile_db_on_live_session_fast_path(profile_dbs, monkeyp
return db
monkeypatch.setattr("hermes_state.SessionDB", _factory)
monkeypatch.setattr(server, "_find_live_session_by_key", lambda _key: ("live-sid", {}))
live_session = {}
with server._sessions_lock:
server._sessions["live-sid"] = live_session
monkeypatch.setattr(
server,
"_find_live_session_by_key",
lambda _key: ("live-sid", live_session),
)
monkeypatch.setattr(
server,
"_live_session_payload",
+22 -72
View File
@@ -552,11 +552,23 @@ def _(rid, params: dict) -> dict:
payload["status"] = "streaming"
return payload
def _reuse_live_response(sid: str, session: dict) -> dict:
# The helper owns the resume lock because slow-path claim races can
# discover a live winner and return it after releasing their own lock.
# Keeping the client-gone check and transport rebind in one critical
# section makes grace expiry atomic across every reuse path.
with _session_resume_lock:
if _sessions.get(sid) is not session:
return _err(rid, 4007, "session no longer live; retry resume")
if session.get("_client_gone_interrupt_requested"):
return _err(rid, 4009, "session disconnect interrupt settling")
return _ok(rid, _reuse_live_payload(sid, session))
# Fast path: if the session is already live, reuse it under the lock.
with _session_resume_lock:
live = _find_live_session_by_key(target)
if live is not None:
return _ok(rid, _reuse_live_payload(*live))
if live is not None:
return _reuse_live_response(*live)
# Lazy/watch resume: register the live session WITHOUT building an agent.
# Used by the desktop's subagent windows — the child runs inside the
@@ -597,7 +609,7 @@ def _(rid, params: dict) -> dict:
lazy=True,
)
if (live := _claim_or_reuse_live(sid, target, record, lease)) is not None:
return _ok(rid, _reuse_live_payload(*live))
return _reuse_live_response(*live)
# A delegated child mid-run emits no session events of its own — report
# its liveness from the relay registry so the window shows a busy turn.
child_running = _child_run_active(target)
@@ -668,7 +680,7 @@ def _(rid, params: dict) -> dict:
record["resume_hydrating"] = True
record["resume_message_count"] = int(found.get("message_count") or 0)
if (live := _claim_or_reuse_live(sid, target, record, lease)) is not None:
return _ok(rid, _reuse_live_payload(*live))
return _reuse_live_response(*live)
_schedule_resume_hydration(sid, target, db, close_db=owns_db)
# The hydration worker now owns a profile-scoped handle and closes it
@@ -761,7 +773,7 @@ def _(rid, params: dict) -> dict:
resume_runtime_overrides=overrides or None,
)
if (live := _claim_or_reuse_live(sid, target, record, lease)) is not None:
return _ok(rid, _reuse_live_payload(*live))
return _reuse_live_response(*live)
_schedule_agent_build(sid)
_schedule_session_cap_enforcement() # trim detached idle sessions over the cap
@@ -871,17 +883,7 @@ def _(rid, params: dict) -> dict:
pass
if lease is not None:
lease.release()
other_sid, other_session = live
payload = _live_session_payload(
other_sid,
other_session,
cols=cols,
touch=True,
transport=current_transport() or _stdio_transport,
omit_messages=omit_messages,
)
payload["resumed"] = target
return _ok(rid, payload)
return _reuse_live_response(*live)
try:
init_home_token = (
set_hermes_home_override(str(profile_home))
@@ -3250,67 +3252,15 @@ def _(rid, params: dict) -> dict:
return err
if _session_uses_compute_host(session):
sid = str(params.get("session_id") or "")
if session.get("running"):
try:
_get_compute_host_supervisor().interrupt(sid, request_id=f"interrupt-{rid}")
except Exception as exc:
return _err(rid, 5019, f"compute-host interrupt failed: {exc}")
with session["history_lock"]:
session["_turn_cancel_requested"] = True
session["queued_prompt"] = None
session.pop("queued_prompts", None)
session["_queued_prompt_generation"] = int(session.get("_queued_prompt_generation", 0)) + 1
_clear_pending(sid)
try:
from tools.approval import resolve_gateway_approval
resolve_gateway_approval(session["session_key"], "deny", resolve_all=True)
except Exception:
pass
_interrupt_session_turn(sid, session, request_id=f"interrupt-{rid}")
except Exception as exc:
return _err(rid, 5019, f"compute-host interrupt failed: {exc}")
return _ok(rid, {"status": "interrupted", "turn_isolation": True})
session, err = _sess(params, rid)
if err:
return err
# Safety net: if the turn's run thread is already gone but `running` stayed
# stuck (a crash/desync that skipped the run loop's `finally`), force-clear it
# so the session can't be permanently bricked at 4009 "session busy" — every
# send/restore/resume would otherwise reject until a full backend restart.
# Always tell the agent to interrupt when the session claims a run is active:
# stale flags are cleared below, and fresh turns clear the interrupt flag at
# entry. This keeps a stale/missing thread handle from making Stop a no-op.
run_thread = session.get("_run_thread")
run_thread_alive = run_thread is not None and run_thread.is_alive()
should_interrupt = bool(session.get("running"))
with session["history_lock"]:
session["_turn_cancel_requested"] = True
session["queued_prompt"] = None
session.pop("queued_prompts", None)
session["_queued_prompt_generation"] = int(session.get("_queued_prompt_generation", 0)) + 1
if should_interrupt:
from agent.interrupt_compat import request_hard_interrupt
request_hard_interrupt(session["agent"])
if not run_thread_alive:
with session["history_lock"]:
if session.get("running"):
session["running"] = False
_clear_inflight_turn(session)
# Stop = stop the TURN (cooperative interrupt above also kills the in-flight
# foreground subprocess). Background processes the agent started (dev servers,
# watchers) are intentionally left running — kill those individually with the
# "x" on the task row (process.kill). Don't reap them here.
# Scope the pending-prompt release to THIS session. A global
# _clear_pending() would collaterally cancel clarify/sudo/secret
# prompts on unrelated sessions sharing the same tui_gateway
# process, silently resolving them to empty strings.
_clear_pending(params.get("session_id", ""))
try:
from tools.approval import resolve_gateway_approval
resolve_gateway_approval(session["session_key"], "deny", resolve_all=True)
except Exception:
pass
_interrupt_session_turn(str(params.get("session_id") or ""), session)
return _ok(rid, {"status": "interrupted"})
+114 -27
View File
@@ -177,8 +177,8 @@ _SLASH_WORKER_TIMEOUT_S = max(5.0, _slash_timeout)
# ``session.create`` (new sid + a fresh _SlashWorker via _deferred_build) and
# never reattaches the OLD sid, so the old session's slash-worker subprocess
# lingers forever — one leaked python process per refresh (#38591 fallout).
# After this grace window, an orphaned (transport-detached, not-running) WS
# session is reaped: its _SlashWorker is closed and the session finalized.
# After this grace window, an orphaned WS session is interrupted if it is still
# running, then reaped once the normal turn-finalization path settles.
# Set to 0 to disable (park forever, pre-fix behaviour).
try:
_ws_orphan_reap_grace = float(
@@ -187,6 +187,7 @@ try:
except (ValueError, TypeError):
_ws_orphan_reap_grace = 20.0
_WS_ORPHAN_REAP_GRACE_S = max(0.0, _ws_orphan_reap_grace)
_WS_ORPHAN_INTERRUPT_REAP_POLL_S = 1.0
_TURN_SETTLE_BEFORE_CLOSE_SECONDS = 5.0
_DETAIL_SECTION_NAMES = ("thinking", "tools", "subagents", "activity")
_DETAIL_MODES = frozenset({"hidden", "collapsed", "expanded"})
@@ -1050,6 +1051,15 @@ def _close_session_by_id(
return _teardown_popped_session(session, end_reason=end_reason)
def _ws_session_is_detached(session: dict | None) -> bool:
"""True if a live session is still bound to the disconnected-WS sentinel."""
return bool(
session
and not session.get("_finalized")
and session.get("transport") is _detached_ws_transport
)
def _ws_session_is_orphaned(session: dict | None) -> bool:
"""True if a WS session has no live transport and no in-flight turn.
@@ -1057,11 +1067,61 @@ def _ws_session_is_orphaned(session: dict | None) -> bool:
``_detached_ws_transport``. A session left on that transport (and not
mid-turn) is genuinely orphaned and safe to reap.
"""
if not session or session.get("_finalized"):
return False
if session.get("running"):
return False
return session.get("transport") is _detached_ws_transport
return bool(
_ws_session_is_detached(session)
and session is not None
and not session.get("running")
)
def _interrupt_session_turn(
sid: str, session: dict, *, request_id: str | None = None
) -> bool:
"""Apply the shared ``session.interrupt`` contract to one claimed session.
Returns whether the interrupt used the compute-host control channel. The WS
orphan reaper calls this same helper after its reconnect grace expires, so a
dead client gets the same partial-history and queued-prompt semantics as an
explicit user interrupt.
"""
use_compute_host = _session_uses_compute_host(session)
should_interrupt = bool(session.get("running"))
run_thread_alive = False
if use_compute_host:
if should_interrupt:
_get_compute_host_supervisor().interrupt(sid, request_id=request_id)
else:
run_thread = session.get("_run_thread")
run_thread_alive = run_thread is not None and run_thread.is_alive()
with session["history_lock"]:
session["_turn_cancel_requested"] = True
session["queued_prompt"] = None
session.pop("queued_prompts", None)
session["_queued_prompt_generation"] = int(
session.get("_queued_prompt_generation", 0)
) + 1
if not use_compute_host:
if should_interrupt:
from agent.interrupt_compat import request_hard_interrupt
request_hard_interrupt(session.get("agent"))
if not run_thread_alive:
with session["history_lock"]:
if session.get("running"):
session["running"] = False
_clear_inflight_turn(session)
_clear_pending(sid)
try:
from tools.approval import resolve_gateway_approval
resolve_gateway_approval(session["session_key"], "deny", resolve_all=True)
except Exception:
pass
return use_compute_host
def _session_owns_durable_lifecycle(session_id: str | None) -> bool:
@@ -1139,7 +1199,7 @@ def _session_has_active_delegations(sid: str, session: dict | None = None) -> bo
return True
def _schedule_ws_orphan_reap(sid: str) -> None:
def _schedule_ws_orphan_reap(sid: str, *, delay_s: float | None = None) -> None:
"""After a grace window, reap session ``sid`` iff it's still orphaned.
Called from the WS-disconnect path. The grace window lets a transient
@@ -1160,34 +1220,60 @@ def _schedule_ws_orphan_reap(sid: str) -> None:
# mutual exclusion against _init_session / _close_session_by_id, which
# guard with _sessions_lock). _sessions_lock is an RLock and the global
# ordering is always resume_lock -> sessions_lock, so nesting is safe.
reschedule = False
reschedule_delay = None
interrupt_session = None
session = None
with _session_resume_lock:
current = _sessions.get(sid)
# Mid-turn detached sessions are not yet orphaned (_running
# short-circuits _ws_session_is_orphaned). If we return here
# the single Timer is gone and the session is immortal until
# the 6h TTL — which also skips running (#85578). Reschedule
# like the active-delegation branch below.
if (
current
and not current.get("_finalized")
and current.get("running")
and current.get("transport") is _detached_ws_transport
):
reschedule = True
elif not _ws_session_is_orphaned(current):
if current is None or not _ws_session_is_detached(current):
return
elif _session_has_active_delegations(sid, current):
reschedule = True
if _session_has_active_delegations(sid, current):
reschedule_delay = _WS_ORPHAN_REAP_GRACE_S
elif current.get("running"):
# Mid-turn detached sessions must never drop the single
# Timer (#85578): after the reconnect grace the turn is
# interrupted once, then the reap keeps polling until the
# normal turn-finalization path settles.
if not current.get("_client_gone_interrupt_requested"):
current["_client_gone_interrupt_requested"] = True
interrupt_session = current
reschedule_delay = _WS_ORPHAN_INTERRUPT_REAP_POLL_S
else:
session = _pop_session_by_id(sid)
if reschedule:
_schedule_ws_orphan_reap(sid)
if interrupt_session is not None:
try:
isolated = _interrupt_session_turn(
sid,
interrupt_session,
request_id=f"client-gone-{sid}",
)
logger.info(
"client_gone sid=%s action=interrupt turn_isolation=%s",
sid,
isolated,
)
except Exception:
logger.exception("client_gone interrupt failed sid=%s", sid)
with _sessions_lock:
if _sessions.get(sid) is interrupt_session:
interrupt_session.pop(
"_client_gone_interrupt_requested", None
)
if reschedule_delay is not None:
_schedule_ws_orphan_reap(sid, delay_s=reschedule_delay)
return
if session is not None and session.get(
"_client_gone_interrupt_requested"
):
logger.info("client_gone sid=%s action=reap", sid)
_teardown_popped_session(session, end_reason="ws_orphan_reap")
timer = threading.Timer(_WS_ORPHAN_REAP_GRACE_S, _reap)
timer = threading.Timer(
_WS_ORPHAN_REAP_GRACE_S if delay_s is None else max(0.0, delay_s),
_reap,
)
timer.daemon = True
timer.start()
@@ -1221,6 +1307,7 @@ def _close_sessions_for_transport(
# _ws_session_is_orphaned recognizes them and the grace-reap can
# actually fire; a standalone `hermes --tui` keeps real _stdio.
session["transport"] = _detached_ws_transport
session.pop("_client_gone_interrupt_requested", None)
detached += 1
try:
_schedule_ws_orphan_reap(sid)