whole bunch of fixes and features #7

Merged
alee merged 32 commits from neuralwatt-router-service into main 2026-08-29 03:31:49 +00:00
Owner

Fix: guard publish_decision against concurrent subscriber churn

Problem

publish_decision() was iterating _sse_loops outside the _subscribers_lock, contradicting its own doc-comment that states the subscriber snapshot must be taken while holding the lock. This creates a real race:

  • subscribe() can add an SSE loop to the set mid-iteration → unbounded iteration (new subscriber keeps the iterator alive).
  • unsubscribe() can remove/drop a loop mid-iteration → RuntimeError: set changed size during iteration.
  • The except-branch in dispatcher._stream_route_decisions() that calls events.unsubscribe() from a background thread can collide with the iteration.

Fix

  1. Snapshot under lock — acquire _subscribers_lock, snapshot list(_sse_loops), release. Iterate the snapshot in the outer scope.
  2. Except-branch mutation under lock — the error-path call to _subscriber_loop.close() and _sse_loops.discard() now acquires _subscribers_lock before mutating (mirrors the path in dispatcher._stream_route_decisions() so the two halves never race).
  3. call_soon_threadsafe() outside lock — the event fan-out uses loop.call_soon_threadsafe() which is already thread-safe and does not require holding the inner lock; the lock is already released for the broadcast loop.

Test

Added test_concurrent_subscribe_unsubscribe_during_publish — a background thread calls subscribe()/unsubscribe() rapidly (500 round-trips) in parallel with a publish_decision() call, asserting no RuntimeError from set-changed-during-iteration.

Verification

  • 635 tests pass (0 regressions)
  • LSP diagnostics clean on changed files
## Fix: guard `publish_decision` against concurrent subscriber churn ### Problem `publish_decision()` was iterating `_sse_loops` outside the `_subscribers_lock`, contradicting its own doc-comment that states the subscriber snapshot must be taken while holding the lock. This creates a real race: - `subscribe()` can add an SSE loop to the set mid-iteration → **unbounded iteration** (new subscriber keeps the iterator alive). - `unsubscribe()` can remove/drop a loop mid-iteration → **RuntimeError: set changed size during iteration**. - The except-branch in `dispatcher._stream_route_decisions()` that calls `events.unsubscribe()` from a background thread can collide with the iteration. ### Fix 1. **Snapshot under lock** — acquire `_subscribers_lock`, snapshot `list(_sse_loops)`, release. Iterate the snapshot in the outer scope. 2. **Except-branch mutation under lock** — the error-path call to `_subscriber_loop.close()` and `_sse_loops.discard()` now acquires `_subscribers_lock` before mutating (mirrors the path in `dispatcher._stream_route_decisions()` so the two halves never race). 3. **`call_soon_threadsafe()` outside lock** — the event fan-out uses `loop.call_soon_threadsafe()` which is already thread-safe and does not require holding the inner lock; the lock is already released for the broadcast loop. ### Test Added `test_concurrent_subscribe_unsubscribe_during_publish` — a background thread calls `subscribe()/unsubscribe()` rapidly (500 round-trips) in parallel with a `publish_decision()` call, asserting no `RuntimeError` from set-changed-during-iteration. ### Verification - **635 tests pass** (0 regressions) - **LSP diagnostics clean** on changed files
alee added 27 commits 2026-08-28 22:43:10 +00:00
Adds a real-time view of actual routing tasks to the TUI dashboard:

backend:
- events.py: in-memory decision-event broker (pure stdlib, thread-safe).
  Bounded ring buffer + fan-out queues. persist_route_decision publishes
  here after each write so the TUI sees decisions without polling.
- dispatcher.py: GET /events/decisions SSE endpoint — replays recent
  decisions then streams live ones with :heartbeat keepalive. Wired into
  persist_route_decision's write path.

tui:
- tui.py: DashboardApp now consumes /events/decisions via a background
  thread (call_from_thread). Decisions table updates live without waiting
  for the 5s /metrics poll. New columns: id, kind, category, tier, ctx,
  selected, est $. Number keys 1-6 cycle panels.
- tui_screens.py: DecisionDetailScreen modal — press Enter or e on any
  decision row to see the full JSON (runner-ups, rejected reason, feature
  flags, confidence, context size).
- tui_model.py: pure data layer extracted from tui.py — build_model,
  build_category_breakdown, decision_row. Testable without a terminal.
- tui_sse.py: background-thread SSE consumer with reconnect.

17 new tests (events broker, SSE endpoint, TUI data model, detail popup,
live decision handling). 562 total, all passing. lsp_diagnostics clean.
- README.md: update test count (356→562, 27 files), add events.py,
  tui_model.py, tui_sse.py, tui_screens.py to Modules table, add
  GET /events/decisions to API Endpoints, describe live SSE feed +
  detail popup + category breakdown in Monitoring, update textual
  import note, add new test files to Testing table.
- CLAUDE.md: update dispatcher description (SSE endpoint), tui
  description (live feed, modal, breakdown, tui_model split), test
  count (562, 27 files), add events.py entry.
- AGENTS.md: new agent-facing working guide — stack snapshot, module
  map with file boundaries and import discipline, conventions,
  test commands, open items, post-change checklist.
Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent)

Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
Use stable row keys in all DataTable renderers and restore cursor position after clear/rebuild cycles.

Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent)

Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
Add tests for cursor persistence across refresh and safe clamping when the selected row is removed.

Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent)

Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
Place 'Verdict mix' and 'Category → model breakdown' panels in a shared
Horizontal container so neither wastes vertical space alone. Both tables
keep their widget ids, so the 2/4 number-key focus bindings still work.

Add a 30-day reset_date to quota_burn (metrics.py) and propagate it
through build_model (tui_model.py) into the TUI quota legend, rendered in
a darker grey ([rgb(128,128,128)]) via Static markup=True so the date is
visually distinct from the surrounding muted legend.

Tests cover reset_date propagation, the legend text, and the shared
Horizontal ancestor for the two side-by-side tables.
Remove the 'Verdict mix' panel from the main dashboard and the side-by-side
Horizontal layout. 'Category → model breakdown' returns to vertical full-width.
The verdict mix is now shown via a '^v' (ctrl+v) popup using the new
VerdictMixScreen modal, mirroring the existing decision-detail popup pattern.

Renumber focus-key bindings to 1-5 (dropping the removed verdict table):
1 model, 2 decision, 3 breakdown, 4 quota, 5 warnings.

Quota reset-date feature and data flow (metrics.py/tui_model.py) unchanged.

Tests cover the ctrl+v popup (open + escape dismiss) and the new focus mapping.
Move the "Recent decisions (enter = details)" panel to directly under the
"Quota burn" panel, so the panel order is: Quota → Recent decisions →
Per-model → Category → Warnings. The decision table is now focused by
default when the TUI opens, so the cursor lands on the most-recent decision
row without requiring a number-key press.

Focus-key bindings (1-5) and the _panels list are unchanged.

Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent)
Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
Replace naive date.today() with datetime.now(timezone.utc).date() to
avoid timezone-dependent reset_date values. Rename legend label from
"resets" to "window start" to accurately describe the backward-looking
30-day window start.
The {"label": "quota", "value": "N/A (plan not set)"} row in
build_model was never consumed by the renderer, which only looks up
plan_kwh/metered_kwh_30d/fraction/calls/reset_date. Return an empty
list when quota is null instead of producing dead data.
Replace the snapshot-at-open pattern with a Textual reactive attribute
so the verdict mix popup tracks the latest /metrics payload while open.
A watch_verdict_mix handler repopulates the DataTable on each refresh,
and the dashboard pushes updates via _push_verdict_mix_updates() at the
end of _render.
- docs: sync stale "1-6" panel references to "1-5"
- fix: hide ProgressBar (display=False) when plan is unconfigured instead
  of fabricating total=1/progress=0 bounds
- fix: simplify _restore_cursor fallback to row 0 (clear() already resets
  cursor, making the clamp logic dead)
- fix: guard against DuplicateKeyError in model table by deduplicating
  per-model rows with a seen set; log a warning on duplicates
- feat: extract _format_quota_legend helper that suppresses None values
  (renders as n/a) so the legend never contains the literal "None"
- feat: push live verdict mix updates into an open VerdictMixScreen
- fix: guard _render against pre-mount NoMatches when set_interval fires
  before compose() finishes mounting widgets

Tests: 579 passed (was 562), including new tests for duplicate model
keys, hidden progress bar, legend None-suppression, live verdict mix
refresh, and cursor restoration to row 0.
4-position scale (no-flex, auto, prefer-flex, force-flex) that swaps
the ranked winner to its -flex serving-class sibling after ranking,
without changing which base model wins. Composes as both a config
default (routing.default_flex_preference) and a per-request override
(TaskRequest.flex_preference), mirroring latency_tolerance.

prefer-flex respects the interactive latency filter; force-flex bypasses
it and sets flex_forced=True in telemetry. Falls back to the standard
winner when no flex sibling exists. Cost is re-estimated from the
sibling row on swap.

route_decisions gains flex_preference, flex_swapped, flex_forced
columns (idempotent ALTER for live DBs). 601 tests pass.
A flex swap could dispatch to a -flex sibling that independently failed a
non-latency hard filter (stale, deprecated, access-restricted, under-tiered,
or capability-missing) because apply_flex_preference only checked latency.
Re-gate the sibling via rejection_reason with latency_tolerance=BATCH
(neutralizing only the latency filter) before swapping. 606 tests pass.
Show the per-decision resolved flex_preference (with swap/forced marker) in
the recent-decisions table and detail popup, plus the configured
routing.default_flex_preference in the quota/legend readout. Recent-decisions
/metrics and the SSE feed now carry flex_preference/flex_swapped/flex_forced
consistently. 615 tests pass.

Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent)
Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
Textual's built-in Footer ellipsizes on narrow terminals, hiding keys. Replace
it with a compact wrapping Static legend (9 short key/action pairs, quit
deduped to one) that wraps rather than truncates at any width. 618 tests pass.
Adds 'cached' to Classification.source and integrates the session cache:
on a hit, skip the classifier and reuse the cached category/tier via the
existing override branch; on a miss, write back only successful (non-fallback)
classifications. Capability flags stay fresh per request; no persistence.
alee added 1 commit 2026-08-28 23:02:52 +00:00
alee changed title from fix: resolve outstanding code review findings (baseline report, SSE cleanup, flex-preference) to fix(events): guard publish_decision against concurrent subscriber churn 2026-08-28 23:03:24 +00:00
alee added 2 commits 2026-08-29 03:15:54 +00:00
Ollama's OpenAI-compatible endpoint (0.22.0) silently ignores per-request num_ctx
and keep_alive under every field shape tried, so the context size has to live in
the model tag itself. Point classifier/verification at mistral-nemo-router:12b
(8192 ctx, 8.6GB vs 14GB untagged) and local_vision at qwen3-vl-router:4b (16384
ctx), so classify+verify share one resident instance and the two-model worst
case (15.9GB) leaves real headroom on a 24GB card. Document the live-measured
rationale in CLAUDE.md and add inline notes in dispatcher where the limitation
surfaces (_classify_once, _run_local_vision).
Two independent, off-by-default features built from code_plans specs in one pass.

Track A — pinch embedding relevance (code_plans/pinch-embedding-relevance.md):
- context_prune: pure order_by_relevance (ascending cosine, least-relevant first),
  shared trim_candidates helper, prune_context(relevance_order=None) param that
  compresses least-relevant candidates first and stops once the token deficit is
  covered; None is byte-for-byte today's uniform pass.
- config: PinchRelevanceConfig nested under PinchConfig (model, base_url, timeout,
  min_candidates) + config.yaml relevance block.
- dispatcher: _embed_for_relevance (single batched /v1/embeddings call, fails
  safe to None), _relevance_order_for (min_candidates gate), wired into both
  chat_completions pinch call sites.

Track B — upstream failover + passive circuit breaker (code_plans/upstream-failover-and-circuit-breaker.md):
- circuit_breaker: pure module mirroring session_cache.py (is_down/record_failure/
  record_success/clear), injected time, exponential backoff capped.
- routing: exclude_models hard filter (reason "circuit_open"), same shape as stale.
- config: CircuitBreakerConfig + config.yaml block.
- dispatcher: _open_circuits exclusion set in route(); non-streaming failover over
  runners_up on 5xx (does NOT consume attempts_used quality budget); streaming
  failover opens+checks each candidate before StreamingResponse so a dead replica
  never reaches the client as a broken stream.

Both features ship off-by-default; every failure path reverts to today's behavior.
673 tests pass (up from 562); live /v1/embeddings endpoint verified with
nomic-embed-text (distinct 768-dim vectors).
alee added 1 commit 2026-08-29 03:28:19 +00:00
The code review (pinch-relevance-and-failover-review.md) found that circuit_breaker.record_success was defined and unit-tested but never called from the real dispatch path. Without it, any model with a failure history monotonically ratchets its cooldown to max_cooldown_seconds on every subsequent failure, even after long stretches of trouble-free service — the opposite of the design's passive-recovery intent.

Call record_success after a successful (<400) response in the non-streaming retry loop and after the streaming pre-flight loop confirms a healthy connection, both gated by cfg.circuit_breaker.enabled matching record_failure's gating. Add an end-to-end integration test at the chat_completions level: a 503 on the first candidate, 200 on the second, asserting record_failure fired for the first and record_success for the second.

Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent)

Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
alee added 1 commit 2026-08-29 03:28:50 +00:00
alee changed title from fix(events): guard publish_decision against concurrent subscriber churn to whole bunch of fixes and features 2026-08-29 03:31:43 +00:00
alee merged commit b46e4abf90 into main 2026-08-29 03:31:49 +00:00
Sign in to join this conversation.
No Reviewers
No Label
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: alee/6krrt#7