Новая система уведомлений поверх готовой SSE-шины: события пишутся в БД,
живут 72 часа и чистятся (при чтении списка + фоновой задачей), всплывают
сверху экрана в момент прихода, доступны через колокольчик в правом верхнем
углу (бейдж непрочитанных + панель).
Типы: приглашение в группу (→ /group), старт/финиш партии участникам кроме
инициатора (→ /match/{id}). Титулы — готовый хелпер-задел (не подключён, т.к.
выдача титулов игрокам ещё не реализована).
Бэкенд: модель Notification + миграция 0008 (идемпотентная), notification_service,
notify.notifications_changed, роутер /api/notifications (GET + /read), триггеры
в groups/matches, фоновая чистка в lifespan, защита hub.publish от закрытого loop.
Фронт: useNotifications/useMarkNotificationsRead, NotificationBell/Panel/Toaster,
перекомпоновка top-bar, стили; useServerEvents знает тип notifications.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
60 lines
2.6 KiB
Python
60 lines
2.6 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 or loop.is_closed():
|
||
return
|
||
ids = [uid for uid in user_ids if uid is not None]
|
||
if ids:
|
||
try:
|
||
loop.call_soon_threadsafe(self._fanout, ids, event)
|
||
except RuntimeError:
|
||
pass # loop закрылся между проверкой и вызовом (напр. при shutdown)
|
||
|
||
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()
|