Files
6krrt/plans/tui-live-routing-panel-sse-fix-review.md
adlee-was-taken 3523dcf93e docs(plans): give every plan a Status line so the queue is greppable
plans/ held 58 documents and exactly one said whether it was open. The rest
mixed finished work, reviews of shipped work, parked specs and genuinely
pending ones, with nothing distinguishing them, so "how many plans are in
the queue" had no answer short of reading all 58.

Now `grep -H '^Status:' plans/*.md` is the answer:

    50 done   3 in progress   2 planned   2 reference   1 parked

Statuses were derived rather than guessed: CLAUDE.md's own built list and
"What's NOT built yet" section, plus checking the subject exists in the
code. A review of work that shipped counts as done -- it records what was
found, it is not a request for anything. `reference` separates the two docs
that are conventions rather than work items (admin-design-standards,
admin-work-framework), which otherwise read as permanently-open plans.

The vocabulary is deliberately five words. A larger one invites "mostly
done" and "blocked-ish", which is how the directory became unreadable.

test_plans_declare_status.py keeps it from rotting: a new plan without a
marker fails, as does an unknown status, one buried below the eighth line,
or an open status with no reason -- "planned" alone is the state that rots,
since nobody can tell later whether it waits on a decision, a dependency,
or just nobody's turn.

Also updates the sweep plan with what landed and what did not, including
that #9 was not a defect.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01VRQXz5SYZYVWscxS1QqF6U
2026-09-08 18:55:16 -04:00

11 KiB

Review: re-fix for finding #5 (async SSE) + finding #9 test gap

Status: done -- review of shipped work

What it was reviewing: the uncommitted working tree on top of 943b2d5, per opencode's summary claiming to (1) make _decision_event_stream a genuine async def generator against asyncio.Queue — the fix tui-live-routing-panel-fixes-review.md said finding #5 still needed — and (2) close the test-coverage gap that same report opened around the finding #9 fix. 7 files touched (events.py, dispatcher.py, tui_screens.py, README.md, three test files), +/- not yet committed. Reviewed via /code-review high (8 parallel finder angles) plus independent verification of every finding against the actual files and, for the two claims below, against live reproductions rather than static reading alone. Full suite: 622/622 passing (tests/test_tui.py + tests/test_events.py + tests/test_metrics_endpoint.py: 54/54).

Verdict: #9 fixed and verified; #5 not fixed — the new mechanism is broken, worse than the sync-threadpool problem it replaced

#9 — fixed and verified

test_live_decision_existing_bucket_defers_count_update (tests/test_tui.py:805) does exactly what the prior report asked for: posts a live decision into a bucket that already has a row, asserts the count is unchanged immediately after, then asserts it updates after the next _on_interval (/metrics) poll. Ran it in isolation — passes, and fails as expected if the deferred-update guard is bypassed by hand. No further action needed here.

#5 — not fixed: the new async bridge silently drops every live decision in production

Confirmed by direct reproduction, not just reading. events.publish_decision (events.py:79-84) bridges to _sse_subscribers via:

try:
    loop = asyncio.get_event_loop()
    for q in _sse_subscribers:
        loop.call_soon_threadsafe(q.put_nowait, decision)
except RuntimeError:
    pass

asyncio.get_event_loop() only returns a usable loop when called from the thread that has one running or set. Every real caller of publish_decision is persist_route_decision (dispatcher.py:820), called from route_endpoint, dispatch_endpoint, and chat_completions — all plain def, not async def (confirmed at dispatcher.py:1543, 1968, 2530), so FastAPI/Starlette runs them on an anyio worker thread, not the thread serving /events/decisions. Reproduced directly on this repo's Python 3.14 venv:

$ python3 -c "
import asyncio, threading
def worker():
    try: asyncio.get_event_loop()
    except RuntimeError as e: print('RuntimeError:', e)
threading.Thread(target=worker).start()
"
RuntimeError: There is no current event loop in thread 'Thread-1 (worker)'.

And end-to-end, simulating the actual production shape (event loop on the main thread serving a subscriber queue, publish_decision called from a separate thread exactly as persist_route_decision does it):

async def main():
    q = asyncio.Queue()
    events.subscribe_sse(q, replay=False)
    threading.Thread(target=lambda: events.publish_decision({"id": 1})).start()
    ...join...
    await asyncio.wait_for(q.get(), timeout=1.0)   # -> asyncio.TimeoutError

Times out every time — the decision never arrives. The except RuntimeError: pass swallows it with no log, so this fails silently: the dashboard shows the initial replay burst and then goes permanently idle (heartbeats only) for the rest of the process's life. This is strictly worse than the finding-#5 report's original complaint (thread-pool pinning) — that version at least delivered decisions.

Why the tests don't catch it: no test calls publish_decision from a plain thread while checking delivery into an asyncio.Queue. test_concurrent_publish_subscribe_unsubscribe_raises_no_error (tests/test_events.py:90) does call publish_decision from real threading.Threads, but never subscribes anything to _sse_subscribers and never checks delivery — it only asserts no exception escaped, and the RuntimeError is already swallowed inside publish_decision before it could. test_full_subscriber_eviction_terminates_decision_event_stream (tests/test_events.py:118) is @pytest.mark.anyio and bypasses publish_decision entirely (monkeypatch.setattr(events, "subscribe_sse", fake_subscribe_sse), then await q.put(...) directly) — the one context where get_event_loop() would have worked is also the one test that never calls the function under test.

Fix direction: capture the loop once at subscribe time (e.g. asyncio.get_running_loop() inside subscribe_sse, called from the async def endpoint where it's valid, stored alongside the queue) rather than calling get_event_loop() from the publisher's thread. call_soon_threadsafe already needs a loop reference callable from any thread — it just needs to be the right loop's reference, obtained once from a context that has one.

Secondary issues found in the same code (real, but currently masked by #5)

  1. Unbounded asyncio.Queue, no backpressure — dispatcher.py:1308 creates asyncio.Queue() with no maxsize, and nothing in the SSE path ever hits QueueFull. The old queue.Queue-based design this replaced bounded each subscriber at DEFAULT_QUEUE_SIZE=100 and evicted slow consumers (still true for the now-unused subscribe/_subscribers path). A stalled dashboard connection has no cap on the SSE side — unbounded per-connection memory growth once #5 is fixed and decisions actually flow.

  2. _sse_subscribers iterated without the lock that guards its mutation — events.py:80, for q in _sse_subscribers:, runs outside _subscribers_lock, while subscribe_sse/unsubscribe_sse/clear all mutate the same set under that lock. A connect/disconnect racing a publish can raise RuntimeError: Set changed size during iteration — currently unreachable in production only because get_event_loop() already raises before this loop runs, so this needs fixing in the same pass as #5, not after.

  3. Old thread-based subscriber path (subscribe/unsubscribe/_subscribers) is now dead in production — dispatcher.py no longer calls events.subscribe/events.unsubscribe anywhere (confirmed by grep); only tests still exercise it. Two independent fan-out implementations now have to be kept in sync by hand, and they've already diverged: the dead one has real bounding/eviction, the live one (once #5 is fixed) doesn't (see #1 above).

  4. Dead sentinel write in _decision_event_stream's finally — dispatcher.py:1326-1329 does queue.put_nowait(events._EVICTED) immediately before unsubscribe_sse(queue). By the time finally runs, this generator's own while True loop has already exited, and no other code ever reads this particular queue — nothing will ever consume the sentinel. Leftover from the old design where eviction needed to wake a different consumer; harmless but confusing cruft.

  5. queue (the stdlib module) is shadowed by a same-named local variable in both dispatcher.py:1308 (_decision_event_stream) and events.py:110,134 (subscribe_sse/unsubscribe_sse's parameter). In dispatcher.py, import queue (line 40) is now unused everywhere else in the file — confirmed via grep, no remaining queue.Full/queue.Empty/ queue.Queue( call sites. Low severity, but a future edit adding queue.Full-style handling inside either function would silently resolve to the local asyncio.Queue instead and produce a confusing AttributeError.

  6. Docstring guarantee dropped — events_decisions()'s docstring lost the line "Unauthenticated and loopback-bound like /metrics; each event carries only the fields already on a route_decisions row — no conversation text, prompt, or session_dir" with no replacement. The guarantee still holds in the actual payload (checked persist_route_decision's dict) — this is a documentation regression, not a behavior one, but it was the one comment warning a future editor not to add such a field at the point that actually emits it.

One claim checked and refuted

The /code-review pass also flagged tui_screens.py's on_key (narrowed to escape-only, dropping "q"/"enter") as letting q fall through to the app-level ("q", "quit", "Quit") binding and quit the whole dashboard while the popup is open. Checked this directly against the installed Textual 8.2.8: App._check_bindings resolves non-priority keys (this app's bindings are all plain tuples, so none are priority=True) via Screen._modal_binding_chain, which explicitly truncates the chain at the first modal screen — it never reaches the App bindings underneath. Verified empirically with run_test(): pressing q with DecisionDetailScreen open leaves the app running and the modal open (it also doesn't close the modal, since q was dropped from on_key — you just can't close it with q anymore, only escape). No quit-while-modal-open bug. This tui_screens.py diff isn't new work from this round anyway — it's the same uncommitted change already assessed twice as fine in the two prior reports.

What's solid

  • Finding #9's fix is correctly and completely closed, with a test that would fail if the deferred-update guard regressed.
  • The intent behind the #5 attempt — a genuine async generator instead of a threadpool-blocking one — is the right direction; the bridging mechanism (asyncio.get_event_loop() called from the publisher's thread) is the one piece that's wrong, not the overall design.
  • _sse_data() reuse and the rest of 943b2d5's prior fixes (#1-4, #7) are untouched by this round and remain correct.

Recommendation

Don't ship this round's #5 attempt as resolving #5 — it currently makes the live feed silently non-functional rather than merely thread-pool-expensive. Fix the loop-capture bug (capture the loop once, from async def events_decisions() or inside subscribe_sse called from that context, not via get_event_loop() in the publisher thread), add a maxsize + eviction/backpressure policy to the SSE queue to match what the old path had, and move the _sse_subscribers iteration in publish_decision inside _subscribers_lock. Once those three land, add a test that exercises the real shape — publish_decision called from a plain thread, delivery checked on the subscribed asyncio.Queue — so this doesn't regress silently again. The dead old subscriber path, the dead sentinel write, and the queue shadowing are all cleanup, not urgent, and can ride along with the same pass. The docstring guarantee is a one-line restore.