Files

43 lines
1.7 KiB
Python

"""Шина событий для трансляции прогресса в 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)