"""Серверные события (SSE): шина публикации (ТЗ 3.14).
События грубозернистые: подписчики получают «kind + данные» и сами решают,
что перечитать (refetch вместо патчей состояния). Издатели — HTTP-хендлеры,
MCP-инструменты (потоки FastMCP) и фоновый detail_task: публикация всегда
идёт через loop.call_soon_threadsafe, поэтому работает из любого потока.
"""
import asyncio
import logging
from dataclasses import dataclass
from typing import Any
logger = logging.getLogger(__name__)
# Размер очереди подписчика: при переполнении события пропускаются —
# coarse-модель самовосстанавливается (клиент перечитывает данные целиком)
QUEUE_MAXSIZE = 64
@dataclass(frozen=True)
class Event:
kind: str
data: dict[str, Any] | None = None
class EventBus:
"""Реестр подписок по user_id (None при publish = broadcast всем)."""
def __init__(self) -> None:
# user_id -> очереди его соединений (несколько вкладок = несколько очередей)
self._subscribers: dict[str, set[asyncio.Queue[Event]]] = {}
# Loop берётся лениво при первой подписке: lifespan в тестах не стартует,
# а после закрытия loop публикация молча сбрасывает его на новый
self._loop: asyncio.AbstractEventLoop | None = None
def subscribe(self, user_id: str) -> asyncio.Queue[Event]:
queue: asyncio.Queue[Event] = asyncio.Queue(maxsize=QUEUE_MAXSIZE)
self._subscribers.setdefault(user_id, set()).add(queue)
try:
self._loop = asyncio.get_running_loop()
except RuntimeError:
pass # вне loop (не бывает для SSE-хендлера) — оставляем прежний
return queue
def unsubscribe(self, user_id: str, queue: asyncio.Queue[Event]) -> None:
queues = self._subscribers.get(user_id)
if queues is not None:
queues.discard(queue)
if not queues:
del self._subscribers[user_id]
def publish(self, user_id: str | None, kind: str, data: dict[str, Any] | None = None) -> None:
"""Положить событие всем подписчикам (None = всем). Потокобезопасно."""
if not self._subscribers or self._loop is None:
return # подписчиков нет — клиент догонит refetch'ем при подключении
event = Event(kind=kind, data=data)
try:
self._loop.call_soon_threadsafe(self._dispatch, user_id, event)
except RuntimeError:
# loop закрыт (перезапуск uvicorn --reload) — сбросим и ждём новую подписку
self._loop = None
def _dispatch(self, user_id: str | None, event: Event) -> None:
if user_id is None:
targets = list(self._subscribers.values())
else:
targets = [self._subscribers.get(user_id, set())]
for queues in targets:
for queue in queues:
try:
queue.put_nowait(event)
except asyncio.QueueFull:
logger.warning("SSE queue overflow for %r: event %r dropped", queue, event.kind)
bus = EventBus()
def publish(user_id: str | None, kind: str, data: dict[str, Any] | None = None) -> None:
"""Точка публикации для хендлеров/агентов (см. ТЗ 3.14)."""
bus.publish(user_id, kind, data)