Spaces:
Running
Running
| import asyncio | |
| import logging | |
| from typing import List | |
| from fastapi import WebSocket | |
| logger = logging.getLogger(__name__) | |
| class ConnectionManager: | |
| def __init__(self): | |
| self.active_connections: List[WebSocket] = [] | |
| self.loop = None # Se registra al iniciar la app | |
| async def connect(self, websocket: WebSocket): | |
| await websocket.accept() | |
| self.active_connections.append(websocket) | |
| # Capturar el event loop principal en la primera conexión | |
| if self.loop is None: | |
| self.loop = asyncio.get_running_loop() | |
| logger.info(f"[WS Admin] Nueva conexión. Total activas: {len(self.active_connections)}") | |
| def disconnect(self, websocket: WebSocket): | |
| if websocket in self.active_connections: | |
| self.active_connections.remove(websocket) | |
| logger.info(f"[WS Admin] Conexión cerrada. Total activas: {len(self.active_connections)}") | |
| async def broadcast(self, message: dict): | |
| """Envía un mensaje JSON a todas las conexiones activas.""" | |
| if not self.active_connections: | |
| return | |
| disconnected = [] | |
| for connection in self.active_connections: | |
| try: | |
| await connection.send_json(message) | |
| except Exception as e: | |
| logger.warning(f"[WS Admin] Error al enviar mensaje, desconectando: {e}") | |
| disconnected.append(connection) | |
| for ws in disconnected: | |
| self.disconnect(ws) | |
| def broadcast_sync(self, message: dict): | |
| """Helper para broadcasts desde endpoints síncronos (def) de FastAPI. | |
| Los endpoints sync corren en un threadpool separado del event loop | |
| principal de uvicorn, por lo que usamos run_coroutine_threadsafe | |
| con la referencia guardada al loop principal. | |
| """ | |
| if not self.active_connections: | |
| return | |
| if self.loop is not None and self.loop.is_running(): | |
| asyncio.run_coroutine_threadsafe(self.broadcast(message), self.loop) | |
| else: | |
| # Fallback: intentar obtener el loop del hilo actual | |
| try: | |
| loop = asyncio.get_event_loop() | |
| if loop.is_running(): | |
| loop.create_task(self.broadcast(message)) | |
| else: | |
| loop.run_until_complete(self.broadcast(message)) | |
| except RuntimeError: | |
| logger.warning("[WS Admin] No se pudo enviar broadcast: sin event loop disponible") | |
| ws_manager = ConnectionManager() | |