spark_colony / telemetry /websocket.py
diwash-barla1's picture
refactor: decompose app into modular domain packages for v2.5
0f336cf
Raw
History Blame Contribute Delete
1.09 kB
from typing import Any, Dict, Set
from fastapi import WebSocket
from core.logging import logger
class WebSocketManager:
"""Manages active WebSocket client connections for real-time live events."""
def __init__(self):
self.active_connections: Set[WebSocket] = set()
async def connect(self, websocket: WebSocket):
await websocket.accept()
self.active_connections.add(websocket)
logger.info(f"WebSocket client connected. Total clients: {len(self.active_connections)}")
def disconnect(self, websocket: WebSocket):
self.active_connections.discard(websocket)
logger.info(f"WebSocket client disconnected. Total clients: {len(self.active_connections)}")
async def broadcast(self, payload: Dict[str, Any]):
if not self.active_connections:
return
stale = set()
for conn in list(self.active_connections):
try:
await conn.send_json(payload)
except Exception:
stale.add(conn)
for conn in stale:
self.disconnect(conn)