| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556 |
- import json
- from fastapi import APIRouter, Depends, Request
- from fastapi.responses import StreamingResponse
- from app.api import deps
- from app.core.events import redis_client
- from app.models import OrgRole, User
- router = APIRouter()
- @router.get("/project/{project_id}/stream", summary="Real-Time SSE Stream")
- async def project_stream(
- project_id: str,
- client_id: str,
- request: Request,
- token: str,
- current_user: User = Depends(
- deps.require_org_role_sse(
- OrgRole.ORG_ADMIN, OrgRole.RESPO_BENEVOLE, OrgRole.RESPO_COMMISSION
- )
- ),
- ):
- """
- Connect to this endpoint to receive real-time updates.
- **Required Query Parameters:**
- - `client_id`: A unique UUID generated by the frontend to prevent echoing.
- """
- print(f"create event generator for {client_id}")
- async def event_generator():
- channel = f"project_{project_id}"
- pubsub = redis_client.pubsub()
- await pubsub.subscribe(channel)
- print(f"subscribed to {channel}")
- try:
- async for event in pubsub.listen():
- if event["type"] != "message":
- continue
- if await request.is_disconnected():
- print("disconnected")
- break
- message = json.loads(event["data"])
- if message.get("client_id") == client_id:
- continue
- yield f"data: {json.dumps({'event': message['event'], 'item': message['item']})}\n\n"
- finally:
- await pubsub.unsubscribe(channel)
- await pubsub.aclose()
- print(f"unsubscribed from {channel}")
- print(f"return event generator for {client_id}")
- # Return as text/event-stream so the browser knows it's an ongoing connection
- return StreamingResponse(event_generator(), media_type="text/event-stream")
|