SSEBroadcasterRoute.py 1.6 KB

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