55 lines
1.7 KiB
Python
55 lines
1.7 KiB
Python
"""Server-Sent Events manager (ROADMAP #85, tranche 4).
|
|
|
|
Singleton extrait de :mod:`backend.main` sans changement de comportement :
|
|
les routers montés par ``main`` partagent la même instance (les clients SSE
|
|
connectés sur ``/api/events`` reçoivent les broadcasts émis depuis
|
|
n'importe quel router).
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import json as _json
|
|
import logging
|
|
|
|
logger = logging.getLogger("obsigate")
|
|
|
|
|
|
class SSEManager:
|
|
"""Manages SSE client connections and broadcasts events."""
|
|
|
|
def __init__(self):
|
|
self._clients: list[asyncio.Queue] = []
|
|
|
|
async def connect(self) -> asyncio.Queue:
|
|
"""Register a new SSE client and return its message queue."""
|
|
queue: asyncio.Queue = asyncio.Queue()
|
|
self._clients.append(queue)
|
|
logger.debug(f"SSE client connected (total: {len(self._clients)})")
|
|
return queue
|
|
|
|
def disconnect(self, queue: asyncio.Queue):
|
|
"""Remove a disconnected SSE client."""
|
|
if queue in self._clients:
|
|
self._clients.remove(queue)
|
|
logger.debug(f"SSE client disconnected (total: {len(self._clients)})")
|
|
|
|
async def broadcast(self, event_type: str, data: dict):
|
|
"""Send an event to all connected SSE clients."""
|
|
message = _json.dumps(data, ensure_ascii=False)
|
|
dead: list[asyncio.Queue] = []
|
|
for q in self._clients:
|
|
try:
|
|
q.put_nowait({"event": event_type, "data": message})
|
|
except asyncio.QueueFull:
|
|
dead.append(q)
|
|
for q in dead:
|
|
self.disconnect(q)
|
|
|
|
@property
|
|
def client_count(self) -> int:
|
|
return len(self._clients)
|
|
|
|
|
|
sse_manager = SSEManager()
|