"""SSE-эндпоинт (ТЗ 3.14): поток событий для SPA.
EventSource фронтенда держит соединение открытым; сервер шлёт события шины
и ping каждые PING_SECONDS секунд (keep-alive, чтобы прокси не рвал поток).
"""
import asyncio
import json
from collections.abc import AsyncIterator
from typing import Any
from fastapi import APIRouter
from starlette.requests import Request
from starlette.responses import StreamingResponse
from app.dependencies import UserDep
from app.realtime import Event, bus
router = APIRouter(prefix="/api/events", tags=["events"])
# keep-alive: прокси (Vite dev, nginx) закрывают «тихие» соединения
PING_SECONDS = 20.0
def _sse_json(kind: str, data: dict[str, Any] | None = None) -> str:
payload: dict[str, Any] = {"kind": kind}
if data:
payload["data"] = data
return json.dumps(payload, ensure_ascii=False)
async def _event_source(request: Request, user_id: str) -> AsyncIterator[str]:
queue = bus.subscribe(user_id)
try:
yield f"data: {_sse_json('ready')}\n\n"
while True:
if await request.is_disconnected():
return
try:
event: Event = await asyncio.wait_for(queue.get(), timeout=PING_SECONDS)
payload = _sse_json(event.kind, event.data)
except TimeoutError:
payload = _sse_json("ping")
yield f"data: {payload}\n\n"
finally:
bus.unsubscribe(user_id, queue)
@router.get("")
async def events(request: Request, user: UserDep) -> StreamingResponse:
"""SSE-поток: события изменений + ping (keep-alive)."""
return StreamingResponse(
_event_source(request, str(user["user_id"])),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
# nginx в проде не буферизует стрим
"X-Accel-Buffering": "no",
},
)