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()