191 lines
6.1 KiB
Python
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 == {}
|