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