"""Tests for the in-memory decision-event broker (events.py). Offline, no network. Drive the broker directly: publish, replay into a new SSE subscriber, and the interaction where a decision is published with no subscriber (must not error and must stay in the ring). """ from __future__ import annotations import threading import pytest import events @pytest.fixture(autouse=True) def _clean_broker(): events.clear() yield events.clear() def _decision(decision_id: int, model: str = "m") -> dict: return {"id": decision_id, "selected_model": model} def test_publish_adds_to_ring_and_stays_for_replay(): events.publish_decision(_decision(1)) events.publish_decision(_decision(2)) assert [d["id"] for d in events.recent_decisions()] == [1, 2] def test_publish_with_no_subscriber_keeps_ring_and_does_not_raise(): events.publish_decision(_decision(1)) assert [d["id"] for d in events.recent_decisions()] == [1] def test_e2e_sse_publish_from_thread(): """Verifies that publish_decision bridges to a subscribed asyncio.Queue when called from a plain threading.Thread — the exact production shape (persist_route_decision runs on a Starlette/anyio worker thread, not the event loop thread). This is the regression test for finding #5 in tui-live-routing-panel-sse-fix-review.md: asyncio.get_event_loop() from a non-loop thread raised RuntimeError on Python 3.10+ and the bare except swallowed it, so decisions silently never reached SSE subscribers. """ import asyncio as _asyncio loop = _asyncio.new_event_loop() _asyncio.set_event_loop(loop) async def _run(): sse_queue = _asyncio.Queue() events.subscribe_sse(sse_queue, replay=False) # Publish from a separate thread, exactly like # persist_route_decision does in the dispatcher. arrived = _asyncio.Event() def _publisher(): events.publish_decision({"id": 911, "model": "x"}) events.publish_decision({"id": 912, "model": "y"}) th = threading.Thread(target=_publisher) th.start() # The two decisions must arrive on the asyncio.Queue. d1 = await _asyncio.wait_for(sse_queue.get(), timeout=2.0) assert d1 == {"id": 911, "model": "x"} d2 = await _asyncio.wait_for(sse_queue.get(), timeout=2.0) assert d2 == {"id": 912, "model": "y"} arrived.set() th.join(timeout=5) assert not th.is_alive() events.unsubscribe_sse(sse_queue) loop.run_until_complete(_run()) loop.close() @pytest.mark.anyio async def test_sse_bridge_from_plain_thread_delivers_to_queue(): """A decision published from a plain threading.Thread must arrive on an asyncio.Queue registered via subscribe_sse. This is the production shape — persist_route_decision runs on an anyio/Starlette worker thread, not the /events/decisions loop thread, and calls publish_decision, which must bridge onto the subscriber's loop via call_soon_threadsafe. Regression for finding #5 in tui-live-routing-panel-sse-fix-review.md. """ import asyncio events.clear() sse_queue: asyncio.Queue[dict] = asyncio.Queue() events.subscribe_sse(sse_queue, replay=False) th = threading.Thread( target=events.publish_decision, args=({"id": 911, "selected_model": "x"},), ) th.start() th.join(timeout=5) assert not th.is_alive(), "publisher thread did not finish" decision = await asyncio.wait_for(sse_queue.get(), timeout=2.0) assert decision == {"id": 911, "selected_model": "x"} events.unsubscribe_sse(sse_queue) @pytest.mark.anyio async def test_sse_push_drops_oldest_when_queue_full(): """_sse_push must bound memory by dropping the oldest item when a subscriber's queue is full, never the newest decision. Pushing 101 items into a 100-slot queue must evict the first (oldest) and keep the rest.""" import asyncio events.clear() sse_queue: asyncio.Queue[dict] = asyncio.Queue(maxsize=events.DEFAULT_QUEUE_SIZE) for i in range(events.DEFAULT_QUEUE_SIZE + 1): events._sse_push(sse_queue, _decision(i)) assert sse_queue.qsize() == events.DEFAULT_QUEUE_SIZE ids = [d["id"] for d in _drain_sse(sse_queue)] assert ids == list(range(1, events.DEFAULT_QUEUE_SIZE + 1)) def _drain_sse(q) -> list: import asyncio out = [] while True: try: out.append(q.get_nowait()) except asyncio.QueueEmpty: return out @pytest.mark.anyio async def test_concurrent_subscribe_unsubscribe_publish_no_error(): """Repeated subscribe/unsubscribe concurrent with publish must not race. Regression pin for the lock gap in `publish_decision` (review: resolve-review-findings-sse-lock-gap-review.md): the snapshot of `_sse_loops` and the `except RuntimeError` mutation must both hold `_subscribers_lock` like the other three accessors do. Honest caveat: this is a regression pin, not a guaranteed reproducer. Under GIL CPython `list(dict.items())` on identity-hashed keys is a single C-level op that does not yield the GIL mid-iteration, so the race is unlikely to fire here even with the bug present. It exists to exercise the concurrent churn path under the version matrix (including a future free-threaded build) where the race is genuine. """ import asyncio events.clear() stop = threading.Event() errors: list[BaseException] = [] def _publisher(): try: for _ in range(2000): if stop.is_set(): return events.publish_decision(_decision(1)) except BaseException as exc: # noqa: BLE001 errors.append(exc) th = threading.Thread(target=_publisher) th.start() for _ in range(500): q: asyncio.Queue[dict] = asyncio.Queue() events.subscribe_sse(q, replay=False) events.unsubscribe_sse(q) stop.set() th.join(timeout=5) assert not th.is_alive(), "publisher thread did not finish" assert errors == [], f"publisher raised: {errors}" assert events._sse_loops == {}