From c548bdd6ada1ffb6a0017c028c3055c58fee713b Mon Sep 17 00:00:00 2001 From: John Paul Soliva Date: Sun, 20 Sep 2026 02:29:36 +0900 Subject: [PATCH] fix(tui-gateway): a notification that loses its delivery claim hands the session's turn back MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The Desktop/TUI poller claims the session's turn (running=True) before it claims the event's durable delivery row. When that row is held by another consumer — a gateway sharing the home claims it before verifying the target — or the ledger read raises, _notif_dispatch_event returned with no turn started and running still set. Nothing else clears it: a running session is exempt from the reaper, keeps its active-session lease, diverts every prompt.submit into a queue only a finishing turn drains, and stops polling its bot mailbox, so the chat and every DM to that bot stay stuck until Stop or a backend restart. The same release was missing around the completion batch's claim/render and in the /loop slash wakeup's swallowed send, and an exception out of any dispatch ended the poller thread — the session's only path to notifications, /loop, /heartbeat and its mailbox. All four now hand the turn (and any claims already taken) back, and the loop logs a bad event instead of dying on it. Fixes #116255 --- .../test_notification_turn_release.py | 134 ++++++++++++++++++ tui_gateway/session_notifications.py | 43 ++++-- 2 files changed, 169 insertions(+), 8 deletions(-) create mode 100644 tests/tui_gateway/test_notification_turn_release.py diff --git a/tests/tui_gateway/test_notification_turn_release.py b/tests/tui_gateway/test_notification_turn_release.py new file mode 100644 index 0000000000..d35edbc833 --- /dev/null +++ b/tests/tui_gateway/test_notification_turn_release.py @@ -0,0 +1,134 @@ +"""A notification that claims the session's turn must hand it back whenever no turn runs. + +The poller claims the idle session (``running = True``) and only then claims the event's durable +delivery row. Nothing else clears ``running``: a busy session is exempt from the reaper, keeps its +active-session lease, diverts every ``prompt.submit`` into a queue that only a finishing turn +drains, and never polls its bot mailbox again — so a dispatch that returns or raises without +starting a turn leaves that session unusable for the life of the backend. +""" + +from __future__ import annotations + +import queue +import sqlite3 +import threading +from types import SimpleNamespace + +import pytest + +from tui_gateway import server + +DELEGATION = {"type": "async_delegation", "delegation_id": "deleg-1", "session_key": "stored"} + + +def _claimed_session() -> dict: + session = {"history_lock": threading.RLock(), "running": False, "history": []} + assert server._notif_claim_turn(session) is True + return session + + +def _no_turn(monkeypatch) -> list: + started: list = [] + monkeypatch.setattr(server, "_run_prompt_submit", lambda *a, **k: started.append(a)) + monkeypatch.setattr(server, "_emit", lambda *a, **k: None) + return started + + +@pytest.mark.parametrize("claim", [None, sqlite3.OperationalError("database is locked")], + ids=["row-held-by-another-consumer", "ledger-unreadable"]) +def test_a_lost_delivery_claim_hands_the_turn_back(monkeypatch, claim): + """A gateway sharing this home claims the durable row before it verifies the target, so the + poller holding the live copy of the same event gets ``None`` — after it already took the turn.""" + def _claim(evt, consumer): + if isinstance(claim, Exception): + raise claim + return claim + + monkeypatch.setattr("tools.async_delegation.claim_event_delivery", _claim) + started = _no_turn(monkeypatch) + session = _claimed_session() + + server._notif_dispatch_event("sid", session, dict(DELEGATION), "text") + + assert session["running"] is False + assert started == [] + assert server._notif_claim_turn(session) is True, "the session must be claimable again" + + +@pytest.mark.parametrize("fail_at", ["claim", "render"]) +def test_a_completion_batch_that_cannot_be_prepared_hands_the_turn_back(monkeypatch, fail_at): + events = [{"type": "completion", "session_id": "proc_a"}, {"type": "completion", "session_id": "proc_b"}] + released: list = [] + claims = iter(["claim-a", sqlite3.OperationalError("database is locked")] if fail_at == "claim" + else ["claim-a", "claim-b"]) + + def _claim(evt, consumer): + value = next(claims) + if isinstance(value, Exception): + raise value + return value + + monkeypatch.setattr("tools.async_delegation.claim_event_delivery", _claim) + monkeypatch.setattr("tools.async_delegation.release_event_delivery", lambda evt, c: released.append(c)) + if fail_at == "render": + monkeypatch.setattr("tools.process_registry_notifications.ProcessNotificationBatch.render", + lambda self, registry: (_ for _ in ()).throw(ValueError("bad payload"))) + started = _no_turn(monkeypatch) + session = {"history_lock": threading.RLock(), "running": False, "history": []} + + server._notif_dispatch_completions("sid", session, [(e, "t") for e in events], + SimpleNamespace(completion_queue=queue.Queue()), None) + + assert session["running"] is False and started == [] + # Claims already taken go back too, or those completions are lost to every consumer for 300 s. + assert released == (["claim-a"] if fail_at == "claim" else ["claim-a", "claim-b"]) + + +def test_a_loop_wakeup_whose_send_cannot_start_hands_the_turn_back(monkeypatch): + """The /loop slash wakeup re-claims the turn for a command that resolves to a prompt, then runs + the send under ``except Exception: pass`` — which swallowed the only signal that no turn started.""" + monkeypatch.setitem(server._methods, "command.dispatch", + lambda rid, params: {"result": {"type": "send", "message": "run the skill"}}) + monkeypatch.setattr(server, "_emit", lambda *a, **k: None) + monkeypatch.setattr(server, "_run_prompt_submit", + lambda *a, **k: (_ for _ in ()).throw(RuntimeError("no free worker"))) + ticks: list = [] + mgr = SimpleNamespace(abandon_tick=lambda: ticks.append("abandoned"), + complete_tick=lambda text: ticks.append("completed") or {}) + session = _claimed_session() + + server._notif_slash_loop_tick("rid", "sid", session, mgr, "/skill go") + + assert session["running"] is False + assert ticks == ["completed"] + + +def test_the_poller_thread_survives_a_dispatch_that_raises(monkeypatch): + """The poller is the session's only path to notifications, /loop, /heartbeat and its bot + mailbox; an exception out of one event's dispatch used to end the thread for good.""" + events: queue.Queue = queue.Queue() + events.put({"type": "completion", "session_id": "proc_a"}) + monkeypatch.setattr("tools.process_registry.process_registry", SimpleNamespace(completion_queue=events)) + for name in ("_poll_bot_live_delivery_guarded", "_maybe_fire_tui_loop_tick", + "_maybe_fire_tui_heartbeat_tick", "_notif_poll_kanban"): + monkeypatch.setattr(server, name, lambda *a, **k: None) + stop = threading.Event() + handled: list = [] + + def _handle_ready(sid, session, ready, emitted, registry, fmt, deferred, **kwargs): + if deferred is not None: # the post-stop drain + return + handled.append(len(ready)) + if len(handled) == 1: + events.put({"type": "completion", "session_id": "proc_b"}) + raise RuntimeError("one bad event") + stop.set() + + monkeypatch.setattr(server, "_notif_handle_ready", _handle_ready) + session = {"history_lock": threading.RLock(), "running": False, "history": []} + worker = threading.Thread(target=server._notification_poller_scoped_loop, args=(stop, "sid", session), daemon=True) + worker.start() + worker.join(timeout=10) + + assert not worker.is_alive() + assert handled == [1, 1], "the second event must still be dispatched after the first one raised" diff --git a/tui_gateway/session_notifications.py b/tui_gateway/session_notifications.py index a261db1fb6..22e504a7ae 100644 --- a/tui_gateway/session_notifications.py +++ b/tui_gateway/session_notifications.py @@ -189,8 +189,12 @@ def _notif_slash_loop_tick(rid: str, sid: str, session: dict, mgr, wakeup: str) if not _notif_claim_turn(session): mgr.abandon_tick() return - _emit("message.start", sid) - _run_prompt_submit(rid, sid, session, payload["message"]) + try: + _emit("message.start", sid) + _run_prompt_submit(rid, sid, session, payload["message"]) + except Exception: + _notif_release_turn(session) # the swallow below would otherwise leave the session busy for good + raise return except Exception: pass @@ -456,7 +460,16 @@ def _notif_poll_kanban_scoped(sid: str, session: dict) -> None: def _notif_dispatch_event(sid: str, session: dict, evt: dict, text: str) -> None: """Run the claimed (running=True) agent turn for one notification event.""" from tools.async_delegation import claim_event_delivery, complete_event_delivery, release_event_delivery - if (claim := claim_event_delivery(evt, "tui-poller")) is None: + try: + claim = claim_event_delivery(evt, "tui-poller") + except Exception as exc: # shared ledger busy/unreadable: the durable row stays pending and replays + _notif_log_failure("notification delivery claim failed", exc) + claim = None + if claim is None: + # Another consumer holds the durable row — a gateway sharing this home claims before it verifies + # the target. No turn will run, and nothing else clears ``running``: a busy session is exempt + # from the reaper, keeps its lease, and never reaches its bot mailbox again. + _notif_release_turn(session) return kwargs = ({"display_kind": "async_delegation_complete", "display_metadata": _async_delegation_display_metadata(evt)} if evt.get("type") == "async_delegation" else {}) @@ -539,10 +552,19 @@ def _notif_dispatch_completions(sid, session, notifications, registry, deferred) if deferred is None: time.sleep(0.25) return - claimed = [(event, text, claim) for event, text in notifications - if (claim := claim_event_delivery(event, "tui-completion-batch")) is not None] - batch = ProcessNotificationBatch(tuple((event, text) for event, text, _claim in claimed)) - text = batch.render(registry) + claimed: list = [] + try: + for event, event_text in notifications: + if (claim := claim_event_delivery(event, "tui-completion-batch")) is not None: + claimed.append((event, event_text, claim)) + batch = ProcessNotificationBatch(tuple((event, event_text) for event, event_text, _claim in claimed)) + text = batch.render(registry) + except Exception as exc: + _notif_log_failure("completion batch preparation failed", exc) + _notif_release_turn(session) + for event, _text, claim in claimed: + release_event_delivery(event, claim) + return if text is None: _notif_release_turn(session) try: @@ -705,7 +727,12 @@ def _notification_poller_scoped_loop(stop_event: threading.Event, sid: str, sess ready.append(queue.get_nowait()) except Exception: break - handle(ready, None) + try: + handle(ready, None) + except Exception as exc: + # This thread is the session's only path to notifications, /loop, /heartbeat and its + # bot mailbox; one bad event must not end all four. + _notif_log_failure("notification dispatch failed", exc) # Drain remaining events after the stop signal so nothing is lost on shutdown; foreign and orphaned-delegation # events are handed back to the shared queue afterwards. deferred: list = []