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
145 lines
11 KiB
Markdown
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.
|