"""In-memory SSE-шина (pub/sub) для пушей «данные изменились». Рассчитана на один процесс (uvicorn `--workers 1`). При переходе на несколько воркеров шину нужно вынести во внешний брокер (например, Redis pub/sub). Публикация вызывается из СИНХРОННЫХ роутеров (FastAPI выполняет их в threadpool), а очереди подписчиков живут в event-loop'е — поэтому публикация перекидывается в loop через `call_soon_threadsafe`, а весь доступ к подпискам происходит на потоке loop'а. """ from __future__ import annotations import asyncio from collections import defaultdict from typing import Any, Iterable class EventHub: def __init__(self) -> None: self._subs: dict[int, set[asyncio.Queue]] = defaultdict(set) self._loop: asyncio.AbstractEventLoop | None = None def bind_loop(self, loop: asyncio.AbstractEventLoop) -> None: self._loop = loop def subscribe(self, user_id: int) -> asyncio.Queue: queue: asyncio.Queue = asyncio.Queue(maxsize=100) self._subs[user_id].add(queue) return queue def unsubscribe(self, user_id: int, queue: asyncio.Queue) -> None: subs = self._subs.get(user_id) if subs is not None: subs.discard(queue) if not subs: self._subs.pop(user_id, None) def publish(self, user_ids: Iterable[int], event: dict[str, Any]) -> None: """Доставить событие подписчикам с указанными user_id (потокобезопасно).""" loop = self._loop if loop is None or loop.is_closed(): return ids = [uid for uid in user_ids if uid is not None] if ids: try: loop.call_soon_threadsafe(self._fanout, ids, event) except RuntimeError: pass # loop закрылся между проверкой и вызовом (напр. при shutdown) def _fanout(self, ids: list[int], event: dict[str, Any]) -> None: # Выполняется на потоке loop'а — доступ к _subs безопасен. for uid in ids: for queue in list(self._subs.get(uid, ())): try: queue.put_nowait(event) except asyncio.QueueFull: pass # медленный клиент — пропускаем (догонит при reconnect/focus) hub = EventHub()