46 lines
1.6 KiB
Python
46 lines
1.6 KiB
Python
"""SSE: поток событий «данные изменились» для текущего пользователя.
|
|
|
|
Клиент (EventSource) держит одно соединение и по событию точечно перезапрашивает данные
|
|
(push-to-invalidate). Мутации остаются обычным REST.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import json
|
|
|
|
from fastapi import APIRouter, Depends, Request
|
|
from fastapi.responses import StreamingResponse
|
|
|
|
from app.auth.deps import get_current_user
|
|
from app.core.events import hub
|
|
from app.models import User
|
|
|
|
router = APIRouter(tags=["events"])
|
|
|
|
_HEARTBEAT_SECONDS = 25
|
|
|
|
|
|
@router.get("/events")
|
|
async def events(request: Request, user: User = Depends(get_current_user)) -> StreamingResponse:
|
|
queue = hub.subscribe(user.id) # type: ignore[arg-type]
|
|
|
|
async def stream():
|
|
try:
|
|
yield ": connected\n\n"
|
|
while True:
|
|
if await request.is_disconnected():
|
|
break
|
|
try:
|
|
event = await asyncio.wait_for(queue.get(), timeout=_HEARTBEAT_SECONDS)
|
|
yield f"data: {json.dumps(event, ensure_ascii=False)}\n\n"
|
|
except asyncio.TimeoutError:
|
|
yield ": ping\n\n" # heartbeat против idle-таймаутов прокси/туннеля
|
|
finally:
|
|
hub.unsubscribe(user.id, queue) # type: ignore[arg-type]
|
|
|
|
return StreamingResponse(
|
|
stream(),
|
|
media_type="text/event-stream",
|
|
headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"},
|
|
)
|