project_stream.py 1.8 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556
  1. import json
  2. from fastapi import APIRouter, Depends, Request
  3. from fastapi.responses import StreamingResponse
  4. from app.api import deps
  5. from app.core.events import redis_client
  6. from app.models import OrgRole, User
  7. router = APIRouter()
  8. @router.get("/project/{project_id}/stream", summary="Real-Time SSE Stream")
  9. async def project_stream(
  10. project_id: str,
  11. client_id: str,
  12. request: Request,
  13. token: str,
  14. current_user: User = Depends(
  15. deps.require_org_role_sse(
  16. OrgRole.ORG_ADMIN, OrgRole.RESPO_BENEVOLE, OrgRole.RESPO_COMMISSION
  17. )
  18. ),
  19. ):
  20. """
  21. Connect to this endpoint to receive real-time updates.
  22. **Required Query Parameters:**
  23. - `client_id`: A unique UUID generated by the frontend to prevent echoing.
  24. """
  25. print(f"create event generator for {client_id}")
  26. async def event_generator():
  27. channel = f"project_{project_id}"
  28. pubsub = redis_client.pubsub()
  29. await pubsub.subscribe(channel)
  30. print(f"subscribed to {channel}")
  31. try:
  32. async for event in pubsub.listen():
  33. if event["type"] != "message":
  34. continue
  35. if await request.is_disconnected():
  36. print("disconnected")
  37. break
  38. message = json.loads(event["data"])
  39. if message.get("client_id") == client_id:
  40. continue
  41. yield f"data: {json.dumps({'event': message['event'], 'item': message['item']})}\n\n"
  42. finally:
  43. await pubsub.unsubscribe(channel)
  44. await pubsub.aclose()
  45. print(f"unsubscribed from {channel}")
  46. print(f"return event generator for {client_id}")
  47. # Return as text/event-stream so the browser knows it's an ongoing connection
  48. return StreamingResponse(event_generator(), media_type="text/event-stream")