Spaces:
Sleeping
Sleeping
| import asyncio | |
| import json | |
| from typing import Dict, AsyncGenerator | |
| class SSEStreamManager: | |
| def __init__(self): | |
| self._queues: Dict[str, asyncio.Queue] = {} | |
| def get_queue(self, session_id: str) -> asyncio.Queue: | |
| if session_id not in self._queues: | |
| self._queues[session_id] = asyncio.Queue() | |
| return self._queues[session_id] | |
| async def publish(self, session_id: str, event_type: str, data: dict): | |
| queue = self.get_queue(session_id) | |
| payload = f"event: {event_type}\ndata: {json.dumps(data)}\n\n" | |
| await queue.put(payload) | |
| async def event_generator(self, session_id: str) -> AsyncGenerator[str, None]: | |
| queue = self.get_queue(session_id) | |
| try: | |
| while True: | |
| data = await queue.get() | |
| yield data | |
| queue.task_done() | |
| except asyncio.CancelledError: | |
| pass | |
| sse_manager = SSEStreamManager() | |