"""Шина событий для трансляции прогресса в SSE. Простая in-memory реализация на основе asyncio-очередей подписчиков (fan-out). Каждый SSE-клиент подписывается, получает свою очередь и читает из неё события. Завтра при необходимости тут окажется Redis pub/sub — интерфейс не изменится. """ from __future__ import annotations import asyncio import contextlib from collections.abc import AsyncIterator from app.models import DownloadEvent class EventBus: def __init__(self, max_queue: int = 1000) -> None: self._subscribers: set[asyncio.Queue[DownloadEvent]] = set() self._max_queue = max_queue async def publish(self, event: DownloadEvent) -> None: # Рассылаем всем подписчикам. Если чья-то очередь переполнена # (медленный клиент) — дропаем самое старое событие, не блокируясь. for queue in list(self._subscribers): if queue.full(): with contextlib.suppress(asyncio.QueueEmpty): queue.get_nowait() queue.put_nowait(event) @contextlib.asynccontextmanager async def subscribe(self) -> AsyncIterator[asyncio.Queue[DownloadEvent]]: queue: asyncio.Queue[DownloadEvent] = asyncio.Queue(maxsize=self._max_queue) self._subscribers.add(queue) try: yield queue finally: self._subscribers.discard(queue) @property def subscriber_count(self) -> int: return len(self._subscribers)