108 lines
3.7 KiB
Python
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))
|