2026-08-13 16:23:18 +04:00

108 lines
3.7 KiB
Python

"""
Fan-out hub between the robot bridge and every connected browser.
One bridge produces state; N browsers consume it. Each client gets its own
bounded queue so a slow tab (a backgrounded phone, say) can never stall the
telemetry loop for everyone else - it just drops frames.
"""
from __future__ import annotations
import asyncio
import json
import time
from collections import deque
from typing import Any
class Hub:
def __init__(self, history_seconds: int = 120, sample_hz: int = 10):
self._clients: set[asyncio.Queue] = set()
self._lock = asyncio.Lock()
self._series: dict[str, deque] = {}
self._history_points = max(60, int(history_seconds * sample_hz))
self._events: deque = deque(maxlen=400)
# -- client lifecycle ---------------------------------------------------
async def register(self) -> asyncio.Queue:
queue: asyncio.Queue = asyncio.Queue(maxsize=32)
async with self._lock:
self._clients.add(queue)
return queue
async def unregister(self, queue: asyncio.Queue) -> None:
async with self._lock:
self._clients.discard(queue)
@property
def client_count(self) -> int:
return len(self._clients)
# -- broadcast ----------------------------------------------------------
async def broadcast(self, kind: str, payload: Any) -> None:
message = json.dumps({"type": kind, "ts": time.time(), "data": payload},
separators=(",", ":"), default=str)
dead = []
for queue in list(self._clients):
try:
queue.put_nowait(message)
except asyncio.QueueFull:
# Drop the oldest frame rather than the newest - stale telemetry
# is worth less than current telemetry.
try:
queue.get_nowait()
queue.put_nowait(message)
except (asyncio.QueueEmpty, asyncio.QueueFull):
dead.append(queue)
for queue in dead:
await self.unregister(queue)
# -- time series --------------------------------------------------------
def record(self, key: str, value: float, ts: float | None = None) -> None:
"""Append one sample to a named series for the dashboard's charts."""
if value is None:
return
try:
value = float(value)
except (TypeError, ValueError):
return
series = self._series.get(key)
if series is None:
series = deque(maxlen=self._history_points)
self._series[key] = series
series.append((ts if ts is not None else time.time(), round(value, 4)))
def series(self, key: str, limit: int | None = None) -> list[list[float]]:
data = list(self._series.get(key, ()))
if limit:
data = data[-limit:]
return [[t, v] for t, v in data]
def all_series(self, limit: int | None = None) -> dict[str, list]:
return {k: self.series(k, limit) for k in self._series}
def series_keys(self) -> list[str]:
return sorted(self._series)
# -- event log ----------------------------------------------------------
def log(self, level: str, source: str, message: str, detail: Any = None) -> dict:
entry = {
"ts": time.time(),
"level": level,
"source": source,
"message": message,
"detail": detail,
}
self._events.append(entry)
return entry
def events(self, limit: int = 200) -> list[dict]:
return list(self._events)[-limit:]
async def emit(self, level: str, source: str, message: str, detail: Any = None) -> None:
await self.broadcast("event", self.log(level, source, message, detail))