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
185 lines
11 KiB
Markdown
185 lines
11 KiB
Markdown
# Review: live routing-decisions panel (`events.py`, `tui_sse.py`, SSE endpoint)
|
|
|
|
Status: done -- review of shipped work
|
|
|
|
**Scope.** Commits `f0ebd83` ("feat(tui): live routing-decisions panel with
|
|
SSE, detail popup, breakdown") and `765a4a6` ("docs: sweep for live routing
|
|
panel + AGENTS.md"), i.e. everything since the last-reviewed commit
|
|
`5f7716e`. 13 files, +1188/-130: new `events.py` (in-memory decision
|
|
broker), new `tui_sse.py` (background SSE consumer thread), a new
|
|
`GET /events/decisions` endpoint in `dispatcher.py`, and the corresponding
|
|
`tui.py`/`tui_model.py`/`tui_screens.py` wiring. Full suite: 562/562
|
|
passing.
|
|
|
|
Run via `/code-review high --since 5f7716e` (forked, finder-angle + verify
|
|
phases), then independently re-derived and confirmed against the actual
|
|
files rather than trusted as-is — every finding below was re-read against
|
|
the committed source at the cited file/line before being included.
|
|
|
|
## Findings
|
|
|
|
### 1. `events.py`'s subscriber set is mutated from multiple threads with no lock — `events.py:24-41,61,67`
|
|
|
|
FastAPI runs its sync `def` endpoints (`/dispatch`, `/route`,
|
|
`/v1/chat/completions`, and the new `/events/decisions`) in a thread pool.
|
|
`_subscribers` is a plain `set()`. `publish_decision` iterates it directly
|
|
(`for subscriber in _subscribers:`, line 36) while `subscribe()`/`unsubscribe()`
|
|
mutate it with `.add()`/`.discard()` from whatever thread is serving a
|
|
concurrent dashboard connect/disconnect. A decision recorded at the same
|
|
instant a dashboard connects or drops raises
|
|
`RuntimeError: Set changed size during iteration`. Confirmed by reading —
|
|
this is a plain, unguarded `set`, no `threading.Lock` anywhere in the file.
|
|
|
|
Where it lands matters: raised inside `persist_route_decision` it's caught
|
|
by that function's broad `except Exception` and silently swallowed — the
|
|
live fan-out for that one decision is just dropped, no crash, no log.
|
|
Raised inside `events.subscribe()` at the top of `_decision_event_stream`
|
|
(dispatcher.py:1312), before that generator's own `try/finally`, it's
|
|
unhandled and can break a new dashboard connection outright.
|
|
|
|
Directly contradicts the module's own docstring
|
|
("thread-safe fan-out to SSE subscribers" — `CLAUDE.md`, and `events.py:1-12`
|
|
describes the same intent without ever establishing it).
|
|
|
|
### 2. A subscriber dropped for a full queue is never told, so its SSE stream idles forever — `events.py:39-41` + `dispatcher.py:1317-1323`
|
|
|
|
When a dashboard falls behind and its 100-slot queue fills, `publish_decision`
|
|
silently evicts it from `_subscribers` (no close, no sentinel). The matching
|
|
`_decision_event_stream` generator has no idea — it keeps calling
|
|
`subscriber.get(timeout=SSE_HEARTBEAT_SECONDS)` on the now-orphaned queue,
|
|
which can only ever time out, so it emits `:heartbeat` forever
|
|
(dispatcher.py:1317-1323). The HTTP connection never errors, so
|
|
`tui_sse.DecisionStream`'s reconnect loop never fires. The dashboard looks
|
|
alive — table renders, connection stays open — but silently stops receiving
|
|
any new decision until the process is restarted. Confirmed by reading both
|
|
sides of the queue handoff.
|
|
|
|
### 3. Live decisions are appended with no id-based dedup, so every reconnect (and the very first connect) duplicates rows already shown — `tui.py:234-249`
|
|
|
|
`_decision_event_stream` always replays the ring buffer on connect
|
|
(`events.subscribe(replay=True)`, dispatcher.py:1312), and
|
|
`tui_sse.DecisionStream.run()` reconnects automatically 5s after any
|
|
transient error (`tui_sse.py:45-61`). `_handle_live_decision` (tui.py:238)
|
|
unconditionally `insert(0, ...)`s every decision it receives into
|
|
`recent_decisions` with no check against ids already present. Since the
|
|
initial `/metrics` poll on startup already populates the same recent
|
|
decisions, the very first SSE replay duplicates them immediately; every
|
|
later reconnect duplicates again. `build_category_breakdown` runs over this
|
|
same list (line 246), so the per-category counts and "majority model" in
|
|
the breakdown panel skew from replay noise, not real traffic. Confirmed by
|
|
reading — no `id` set or seen-check anywhere in `_handle_live_decision` or
|
|
`decision_row`.
|
|
|
|
### 4. `tui_sse.py`'s reconnect loop only catches network/parse errors, logs nothing, and can be permanently killed by its own callback — `tui_sse.py:45-61`
|
|
|
|
`self.callback(...)` (line 56) is `App.call_from_thread`, which can raise
|
|
once the Textual app's event loop is gone — e.g. during shutdown, since
|
|
`stop()` (line 63) only sets an `Event` and does not interrupt a blocking
|
|
`iter_lines()` read, so the thread can still be mid-callback for up to
|
|
`STREAM_TIMEOUT` seconds after `on_unmount` calls `stop()`. That's not a
|
|
`requests.RequestException` or `ValueError`, so it isn't caught by the
|
|
`except` on line 57 — it propagates out of `run()` and ends the thread for
|
|
good, no further reconnect attempts, ever.
|
|
|
|
Separately, even the errors that *are* caught are swallowed with a bare
|
|
`pass` (line 60) — no `logs.warning(...)` call, unlike the identical
|
|
"this must never raise" pattern used elsewhere in this codebase (e.g.
|
|
`persist_route_decision`'s except block, which does log). An ordinary
|
|
dropped VPN tunnel or router restart produces zero diagnostic trace here.
|
|
|
|
### 5. `/events/decisions` is a sync endpoint that blocks in the shared threadpool for the life of each SSE connection — `dispatcher.py:1328-1329`
|
|
|
|
Confirmed: `def events_decisions():`, not `async def`. FastAPI runs sync
|
|
endpoints in its bounded default threadpool, the same pool serving
|
|
`/dispatch`, `/route`, and `/v1/chat/completions`. A handful of connected
|
|
dashboards, or `DecisionStream`'s reconnect loop flapping through repeated
|
|
transient failures (finding #4 makes that worse — a dead thread means a
|
|
`textual` restart reconnects from scratch, briefly doubling in-flight
|
|
connections), can hold enough concurrent long-lived streams to exhaust the
|
|
pool and stall real routing/dispatch requests behind idle SSE connections.
|
|
|
|
### 6. `DecisionDetailScreen.on_key` risked a double-dismiss on Enter — `tui_screens.py:65` (as committed in `765a4a6`)
|
|
|
|
As committed: `if event.key in ("escape", "q", "enter"): self.dismiss(None)`
|
|
with no `event.stop()`. When the Close button has focus and Enter is
|
|
pressed, Textual delivers the key to the focused `Button` first; `Button`
|
|
has no key handler of its own, so the event bubbles unstopped to this
|
|
`on_key`, which dismisses — but the event can *also* continue bubbling to
|
|
the App's binding resolution and match `Button.BINDINGS`'s own `enter`
|
|
binding, firing `action_press()` → `Button.Pressed` →
|
|
`on_button_pressed` → a second `self.dismiss(None)` on an already-popped
|
|
modal. Confirmed against the exact committed line via `git show 765a4a6`.
|
|
|
|
**Already independently fixed in the working tree, uncommitted, as of this
|
|
review** — `git diff -- tui_screens.py` shows `on_key` narrowed to
|
|
`if event.key == "escape":` only, dropping `"enter"`/`"q"` handling from
|
|
this method entirely (Enter now only ever reaches `Button`'s own binding,
|
|
`q`'s docstring/label mention removed too). This resolves the race by
|
|
construction rather than by adding `event.stop()`. No action needed here —
|
|
noting it so the fix isn't lost if the working tree changes again before
|
|
it's committed.
|
|
|
|
### 7. Two independent `datetime.now()` calls give the SSE payload a different `observed_at` than the persisted row for the same decision — `dispatcher.py:923` vs `dispatcher.py:954`
|
|
|
|
The `INSERT` is stamped at line 923; `events.publish_decision(...)`'s
|
|
payload is stamped by a second, separate call at line 954, several
|
|
statements later and after `conn.commit()`. A dashboard that correlates the
|
|
live SSE event for decision `id=N` against the same row fetched later via
|
|
`/metrics` sees `observed_at` differ by the insert/commit latency — small in
|
|
practice, but it breaks the assumption (implicit in the code's own comment
|
|
at line 949-950, "the row id becomes the ordering handle") that the SSE
|
|
payload mirrors the persisted row exactly. Trivial fix: reuse one timestamp
|
|
for both.
|
|
|
|
### 8. `_decision_event_stream` hand-rolls SSE frame construction as a third, inconsistent implementation — `dispatcher.py:1303-1323`
|
|
|
|
`dispatcher.py` already builds `text/event-stream` frames twice elsewhere
|
|
(the OpenAI-compatible streaming wrapper and the real streaming proxy),
|
|
both via `yield f"data: {json.dumps(...)}\n\n".encode()`. The new generator
|
|
reimplements the same primitive and is the only one of the three that
|
|
*doesn't* `.encode()` the yielded string (relying on Starlette's
|
|
`StreamingResponse` to encode `str` chunks for it, which does work — this
|
|
isn't a bug — but it's now three copies of one pattern that can silently
|
|
drift apart on the next edit to any one of them). Minor; a shared helper
|
|
would remove the inconsistency rather than fix a defect.
|
|
|
|
### 9. Every live decision triggers a full sort + `Counter` rebuild over the whole recent-decisions list — `tui_model.py` (`build_category_breakdown`) via `tui.py:246-248`
|
|
|
|
`_handle_live_decision` calls `build_category_breakdown` over the entire
|
|
(up to 50-item) list on every single SSE event, doing a full sort and
|
|
`Counter` rebuild on the Textual UI thread each time, when only the one
|
|
`(category, tier)` bucket the new decision falls into actually changed.
|
|
Not a correctness bug — under realistic traffic volumes (a handful of
|
|
decisions/sec at most) this is imperceptible — but under a burst it's
|
|
doing O(n log n) work per event for an O(1) update, on the UI thread.
|
|
Lowest priority of the nine; noted for completeness rather than urgency.
|
|
|
|
## What's solid
|
|
|
|
- The core SSE plumbing works end-to-end: replay-then-live, heartbeats,
|
|
`retry:` hint, and the `unsubscribe`-in-`finally` shutdown path are all
|
|
correctly shaped for the common case (one dashboard, no reconnect storms).
|
|
- `/events/decisions`'s docstring is accurate about what it does and doesn't
|
|
carry (no conversation text, prompt, or `session_dir` — matches the same
|
|
privacy invariant already enforced for `route_decisions` rows).
|
|
- `_drain_queue`'s non-blocking drain-then-block pattern (dispatcher.py:1293-1300,
|
|
1315-1323) is the right shape for "flush anything buffered, then wait" and
|
|
is itself correct in isolation.
|
|
- 562/562 tests pass, including new coverage in `test_events.py`,
|
|
`test_metrics_endpoint.py`, and `test_tui.py` for the parts of this
|
|
feature that are correct.
|
|
|
|
## Recommendation
|
|
|
|
Findings #1 and #2 are the ones worth fixing before this sees real
|
|
multi-dashboard or flaky-network use — #1 is a live `RuntimeError` under
|
|
ordinary concurrent access (not a rare race window; any connect/disconnect
|
|
overlapping a `publish_decision` call triggers it), and #2 silently breaks
|
|
the exact feature this commit exists to ship. #4 (dead reconnect thread) and
|
|
#3 (duplicate rows on every reconnect) compound #1/#2 — a dashboard that hit
|
|
the `RuntimeError` and then can't reconnect because its thread died, showing
|
|
stale-but-plausible data with no error, is a bad failure mode for something
|
|
meant to be watched passively. #5, #7, #8, #9 are all real but low-severity
|
|
and can ride along with the same pass. #6 is already resolved in the
|
|
uncommitted working tree.
|