Files
6krrt/plans/tui-live-routing-panel-fixes-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

145 lines
11 KiB
Markdown

# 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.