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")