57 lines
2.4 KiB
Python
57 lines
2.4 KiB
Python
"""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:
|
||
return
|
||
ids = [uid for uid in user_ids if uid is not None]
|
||
if ids:
|
||
loop.call_soon_threadsafe(self._fanout, ids, event)
|
||
|
||
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()
|