# Review: fixes for `tui-live-routing-panel-review.md` Status: done -- review of shipped work **What it was reviewing:** commit `943b2d5` ("fix: resolve TUI live-routing panel review findings"), fixing findings #1-5 and #7-9 from [`tui-live-routing-panel-review.md`](tui-live-routing-panel-review.md) (#6 was already resolved separately, uncommitted, at review time — still uncommitted now, in `tui_screens.py`). 7 files, +354/-34. Reviewed by reading the commit diff against each finding, re-deriving the fix logic by hand, and — for the one finding below that isn't fully resolved — tracing it against the installed `starlette==1.6.0` source rather than assuming Starlette's behavior. Full suite: 568/568 passing (up from 562; 6 new tests). ## Verdict: 7 of 9 fixed and verified; 1 not actually fixed by the change made; 1 test-coverage gap opened by the fix itself ### Fixed and verified | # | Original finding | Fix | Verified | |---|---|---|---| | 1 | `events._subscribers` (a plain `set`) mutated from multiple threads with no lock — concurrent iterate/add/discard raises `RuntimeError` | `threading.Lock` now guards every mutation and the iteration in `publish_decision`, `subscribe`, `unsubscribe`, `clear` | Read the full diff: iteration and mutation are both inside the same `with _subscribers_lock:` block in every function, so there's no window where one thread iterates while another mutates. New `test_concurrent_publish_subscribe_unsubscribe_raises_no_error` stresses it with 4 publisher + 4 subscribe/unsubscribe threads; passes | | 2 | A subscriber evicted for a full queue was silently dropped — its SSE generator kept blocking on the orphaned queue forever | New `_EVICTED` sentinel: on eviction, `publish_decision` writes the sentinel into the subscriber's queue (displacing the oldest item if still full); `_decision_event_stream` and `_drain_queue` both check for it and terminate/stop-draining on sight, letting `finally: events.unsubscribe(...)` close the stream | Traced the full path by hand: eviction → sentinel write → generator's `subscriber.get()` returns the sentinel → `break` → `finally` unsubscribes → SSE connection closes → client's `DecisionStream` reconnect loop (5s) picks it back up. New `test_full_subscriber_eviction_terminates_decision_event_stream` reproduces exactly this against the real `dispatcher._decision_event_stream()` generator (not a mock) and asserts `StopIteration`; passes | | 3 | Live decisions inserted with no id-based dedup — every SSE replay/reconnect duplicated rows and skewed the breakdown | New `self._seen_ids` set; `_handle_live_decision` skips ids already seen, `_render` (the periodic `/metrics` poll) reseeds `_seen_ids` from the freshly fetched model so a full refresh can't wedge a stale id into permanent exclusion | Read the logic: dedup is keyed correctly and reset on every authoritative `/metrics` fetch, so a decision that later legitimately reappears (e.g. after a `/metrics`-driven eviction from the 50-row cap) isn't permanently blocked. `test_live_decision_dedup_skips_duplicate_id` reproduces a duplicate id and confirms the row count doesn't grow and the front-of-list row is the new one | | 4 | `tui_sse.py`'s reconnect loop only caught `RequestException`/`ValueError`, logged nothing, and a callback exception (e.g. `call_from_thread` post-shutdown) killed the thread permanently | Added a separate `except RuntimeError` branch for the callback-failure case; both branches now call `logs.warning(...)`; a `self._stopped.is_set()` check runs right after error handling (before the reconnect delay), and the delay itself changed from `time.sleep` to the interruptible `self._stopped.wait(...)` | Confirmed the exception no longer escapes `run()` — it's now caught, logged, and the loop continues. The `Event.wait()` swap is a genuine (if unrequested) improvement: shutdown no longer waits out a stale `RECONNECT_SECONDS` sleep. `logs.warning(name, **fields)`'s signature matches the call sites, so the fix's own error-logging code can't itself raise | | 7 | Two independent `datetime.now()` calls gave the DB row and the SSE payload different `observed_at` for the same decision | `observed_at` computed once, reused for both the `INSERT` and the `publish_decision` payload | Trivial, confirmed by reading — one call, two uses | | 8 (partial — see below) | `_decision_event_stream` hand-rolled SSE frame encoding, the only one of several copies missing `.encode()` | New `_sse_data()` helper, used by both yield points in `_decision_event_stream` | The missing-`.encode()` inconsistency this finding actually described as a defect risk is gone. The broader "three-plus copies of this pattern" observation is not — see below | | 9 | Every live decision triggered a full sort + `Counter` rebuild over the whole recent-decisions list | New `_update_category_breakdown`: only calls `build_category_breakdown` (full rebuild) when the new decision's `(category, tier)` bucket doesn't already exist; an existing bucket's count is left alone and catches up on the next `/metrics` poll | Confirmed the O(n log n) rebuild-per-event is gone for the common case (an existing bucket). This is a real design trade — an existing bucket's displayed count can lag by up to one poll interval (`refresh_seconds`, default part of the periodic timer) rather than updating immediately — but it's exactly what the docstring says it does, and it does fix the stated complaint (needless full rebuild on every event). See the test-coverage gap this opened, below | ### #6, unrelated to this commit Still resolved as reported previously: `tui_screens.py`'s `on_key` narrowed to `escape`-only remains uncommitted in the working tree, unaffected by `943b2d5`. No action needed here. ### Not actually fixed: #5, the async conversion doesn't move `_decision_event_stream` off the threadpool `events_decisions()` is now `async def`, but `_decision_event_stream()` — the generator actually doing the work, including the blocking `subscriber.get(timeout=SSE_HEARTBEAT_SECONDS)` call — is still a plain **sync** generator (`def`, no `await` anywhere in it). Making the *endpoint function* async doesn't change what Starlette does with the object it returns. Traced directly against the installed `starlette==1.6.0` source: ```python # starlette/responses.py, StreamingResponse.__init__ if isinstance(content, AsyncIterable): self.body_iterator = content else: self.body_iterator = iterate_in_threadpool(content) # starlette/concurrency.py async def iterate_in_threadpool(iterator): as_iterator = iter(iterator) while True: try: yield await anyio.to_thread.run_sync(_next, as_iterator) except _StopIteration: break ``` `_decision_event_stream()` is a plain generator, not an `AsyncIterable`, so Starlette wraps it in `iterate_in_threadpool` exactly as it did before this commit. Every `next()` call — including the one that blocks for up to `SSE_HEARTBEAT_SECONDS` (15s) waiting on `subscriber.get(...)` — is dispatched via `anyio.to_thread.run_sync`, which borrows a worker from anyio's thread-pool `CapacityLimiter` for the duration of that call and returns it only between iterations (i.e. for the brief moment between one `next()` returning and the response layer requesting the next one — not a meaningful release under a 15-second block repeated indefinitely for the life of the connection). The finding's actual concern — several long-lived `/events/decisions` connections competing with `/dispatch`/`/route`/ `/v1/chat/completions` for a bounded pool of worker threads — is unchanged by this commit. The `async def` on `events_decisions()` is not wrong, just inert here: the function body is `return StreamingResponse(...)`, which doesn't block regardless. An actual fix needs `_decision_event_stream` itself to stop blocking a thread: either make it a genuine `async def` generator against an `asyncio.Queue` (bridged from the sync `publish_decision` call sites via `loop.call_soon_threadsafe`, since those run in FastAPI's sync-endpoint threadpool), or accept the current design and size anyio's thread pool (`anyio.to_thread.current_default_thread_limiter().total_tokens`) for the expected number of concurrent dashboard connections instead. ### Test-coverage gap opened by the #9 fix `test_live_decision_inserts_row_at_front_and_rerenders`'s prior assertion — that an existing category's count rises after a live decision lands in it — was **deleted** by this commit's diff to `tests/test_tui.py`, not updated: ```diff - # Category breakdown reflects the new decision. - breakdown = app._last_model["category_breakdown"] - assert any( - r["category"] == "coding_general" and r["count"] >= 2 - for r in breakdown - ) ``` That's the correct call given the new deferred-update design (the old assertion would now fail, since an existing bucket's count intentionally doesn't move until the next `/metrics` poll) — but nothing replaced it. `test_live_decision_new_bucket_rebuilds_breakdown` only exercises the brand-new-bucket path. There is currently no test asserting the existing-bucket path's actual documented behavior — that the count is deliberately left unchanged rather than incremented — so a future change that accidentally reintroduces a per-event increment (partial, wrong, or otherwise) for an existing bucket would pass the suite silently. Worth one small test: post a live decision into a bucket that already has one, assert the count is unchanged immediately after, then assert it's correct after the next full `_render`. ## What's solid - The lock (#1) and sentinel (#2) fixes compose correctly together — traced the full concurrent path by hand and it holds: eviction always happens while holding the lock (so `_subscribers` state is never inconsistent), and the sentinel write to a genuinely-full queue's fallback path (discard-oldest, then put) is a narrow theoretical race against a simultaneously-draining consumer, but at `DEFAULT_QUEUE_SIZE = 100` it would need the consumer to drain the entire queue to empty inside that handful of Python bytecodes to matter — not reachable at the shipped default, not worth a finding. - `logs.warning` calls added in `tui_sse.py` use the module's real `(name: str, **fields)` signature, so the new error-handling code can't itself throw and reintroduce the bug it's fixing. - Reusing `observed_at` (#7) and the `_sse_data()` helper (#8) are both clean, minimal, correctly-scoped diffs. ## Recommendation #1, #2, #3, #4, #7, #9 (with the one small test gap noted) are solid — no further action needed on those. #8 is fine as far as it goes; the remaining duplicate `data: {json.dumps(...)}\n\n".encode()` call sites (dispatcher.py:1873-1874, 2445, 2486) were never in scope for a defect fix, just a style observation, so leaving them is a reasonable call. #5 is the one to send back: the fix as shipped doesn't change the threadpool-pinning behavior the finding described, and the commit message ("no longer pin FastAPI worker threads") states something that isn't what the code now does. Recommend either the `asyncio.Queue`-bridge approach above, or explicitly deciding the thread-pool-per-connection cost is acceptable at expected dashboard-connection counts and documenting that instead of carrying an `async def` that implies a fix that isn't there.