"""Live-Board-Hub: Statuswechsel aus der DB-Schicht → alle WebSocket-Clients. Kein Polling; das Frontend hält den Snapshot und wendet Deltas an.""" import asyncio import json import logging log = logging.getLogger("creator2.ws") class Hub: def __init__(self): self._clients: set = set() self._loop: asyncio.AbstractEventLoop | None = None self._queue: asyncio.Queue | None = None def bind_loop(self, loop: asyncio.AbstractEventLoop) -> None: self._loop = loop self._queue = asyncio.Queue() # EIN Sender-Task serialisiert die Broadcasts — ein Task-pro-push konnte # sich bei Backpressure überholen (out-of-order-Deltas, konkurrierende # send_text auf demselben Socket). loop.create_task(self._sende_schleife()) async def _sende_schleife(self) -> None: assert self._queue is not None while True: msg = await self._queue.get() await self._send_all(msg) async def connect(self, websocket) -> None: await websocket.accept() self._clients.add(websocket) def disconnect(self, websocket) -> None: self._clients.discard(websocket) def push(self, tabelle: str, row: dict) -> None: """Thread-sicher (DB-Schicht ruft synchron): in die Queue legen, der Sender-Task verschickt der Reihe nach.""" if not self._clients or self._loop is None or self._queue is None: return msg = json.dumps({"tabelle": tabelle, "row": row}, ensure_ascii=False, default=str) self._loop.call_soon_threadsafe(self._queue.put_nowait, msg) async def _send_all(self, msg: str) -> None: for ws in list(self._clients): try: await ws.send_text(msg) except Exception: self._clients.discard(ws) hub = Hub()