Files
6krrt/tests/test_events.py

191 lines
6.1 KiB
Python

"""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 == {}