# 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`](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: ```python 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): ```python 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.Thread`s, 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.