43 lines
1.7 KiB
Python
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)
|