"""In-process revision notifications for SSE clients.""" from __future__ import annotations import asyncio import contextlib from collections import defaultdict class EventBroker: """Fan out room revisions; clients always reload the latest snapshot.""" def __init__(self) -> None: """Initialize an empty subscription registry.""" self._queues: dict[str, set[asyncio.Queue[int]]] = defaultdict(set) def subscribe(self, room_code: str) -> asyncio.Queue[int]: """Subscribe a bounded notification queue to a room.""" queue: asyncio.Queue[int] = asyncio.Queue(maxsize=1) self._queues[room_code].add(queue) return queue def unsubscribe(self, room_code: str, queue: asyncio.Queue[int]) -> None: """Remove a room notification queue from the registry.""" self._queues[room_code].discard(queue) if not self._queues[room_code]: self._queues.pop(room_code, None) def publish(self, room_code: str, revision: int) -> None: """Publish the newest room revision to every subscriber.""" for queue in tuple(self._queues.get(room_code, ())): if queue.full(): with contextlib.suppress(asyncio.QueueEmpty): queue.get_nowait() queue.put_nowait(revision)