SSEBroadcasterRoute.py 1.6 KB

1234567891011121314151617181920212223242526272829303132333435363738394041
  1. import json
  2. from fastapi import Request, Response
  3. from fastapi.routing import APIRoute
  4. from app.core.events import redis_client
  5. class SSEBroadcasterRoute(APIRoute):
  6. def get_route_handler(self):
  7. original_route_handler = super().get_route_handler()
  8. async def custom_route_handler(request: Request) -> Response:
  9. # Execute the route and let FastAPI serialize the response
  10. response = await original_route_handler(request)
  11. openapi_extra = self.openapi_extra or {}
  12. event_name = openapi_extra.get("sse_event")
  13. # Check if an event is declared and the request was successful
  14. if event_name and response.status_code in (200, 201):
  15. project_id = request.path_params.get("project_id")
  16. client_id = request.headers.get("X-Client-ID", "unknown")
  17. if request.method == "DELETE":
  18. # For DELETE, the response body is often empty. We broadcast the path parameters.
  19. item_json = json.dumps(request.path_params)
  20. else:
  21. # For POST/PUT, response.body is already serialized as a JSON string by FastAPI
  22. item_json = response.body.decode("utf-8")
  23. # Assemble the raw payload string
  24. payload_str = (
  25. f'{{"client_id": "{client_id}", "event": "{event_name}", "item": {item_json}}}'
  26. )
  27. channel = f"project_{project_id}"
  28. await redis_client.publish(channel, payload_str)
  29. return response
  30. return custom_route_handler