"""实例数据看板:后台定时聚合,内存快照,SSE 版本通知(对齐中控 dashboard_store).""" from __future__ import annotations import json import os import queue import threading from collections.abc import Callable, Iterator from typing import Any INSTANCE_DASHBOARD_POLL_SEC = float(os.getenv("INSTANCE_DASHBOARD_POLL_SEC", "5")) INSTANCE_DASHBOARD_SSE_HEARTBEAT_SEC = float(os.getenv("INSTANCE_DASHBOARD_SSE_HEARTBEAT_SEC", "25")) BuildFn = Callable[[], dict[str, Any]] class InstanceDashboardStore: def __init__(self) -> None: self._lock = threading.RLock() self.version = 0 self.payload: dict[str, Any] | None = None self.aggregating = False self.last_error: str | None = None self._subscribers: list[queue.Queue[str | None]] = [] self._stop = threading.Event() self._refresh = threading.Event() self._thread: threading.Thread | None = None self._build_fn: BuildFn | None = None def start(self, build_fn: BuildFn) -> None: self._build_fn = build_fn if self._thread and self._thread.is_alive(): return self._stop.clear() self._thread = threading.Thread( target=self._loop, daemon=True, name="instance-dashboard-poll", ) self._thread.start() def stop(self) -> None: self._stop.set() self._refresh.set() self._broadcast(close=True) def request_refresh(self) -> None: self._refresh.set() def snapshot_dict(self) -> dict[str, Any]: with self._lock: p = dict(self.payload or {}) ver = self.version aggregating = self.aggregating err = self.last_error if not p: return { "ok": False, "dashboard_version": ver, "aggregating": aggregating, "error": err, "msg": err or "看板快照尚未就绪", "poll_interval_sec": INSTANCE_DASHBOARD_POLL_SEC, } return { **p, "dashboard_version": ver, "aggregating": aggregating, "error": err or p.get("error"), "poll_interval_sec": INSTANCE_DASHBOARD_POLL_SEC, } def event_dict(self) -> dict[str, Any]: with self._lock: p = self.payload or {} return { "dashboard_version": self.version, "updated_at": p.get("updated_at"), "aggregating": self.aggregating, "ok": p.get("ok", True) if self.payload else False, "error": self.last_error or p.get("error"), } def _loop(self) -> None: assert self._build_fn is not None while not self._stop.is_set(): self._aggregate_once(self._build_fn) if self._stop.is_set(): break self._refresh.clear() # 周期等待,可被 request_refresh 提前唤醒 self._refresh.wait(timeout=INSTANCE_DASHBOARD_POLL_SEC) def _aggregate_once(self, build_fn: BuildFn) -> None: with self._lock: self.aggregating = True self._broadcast() try: result = build_fn() if not isinstance(result, dict): result = {"ok": False, "msg": "聚合返回无效"} except Exception as e: result = {"ok": False, "msg": str(e), "error": "aggregate_failed"} with self._lock: self.version += 1 prev = self.payload if isinstance(self.payload, dict) else None if result.get("ok") is False and prev and prev.get("ok"): self.payload = prev self.last_error = str(result.get("msg") or result.get("error") or "aggregate_failed") else: self.payload = result self.last_error = ( None if result.get("ok") is not False else str(result.get("msg") or result.get("error") or "aggregate_failed") ) self.aggregating = False self._broadcast() def _broadcast(self, *, close: bool = False) -> None: with self._lock: subs = list(self._subscribers) event = None if close else json.dumps(self.event_dict(), ensure_ascii=False) dead: list[queue.Queue[str | None]] = [] for q in subs: try: q.put_nowait(None if close else event) except queue.Full: try: q.get_nowait() except queue.Empty: pass try: q.put_nowait(event) except queue.Full: dead.append(q) except Exception: dead.append(q) if dead: with self._lock: for q in dead: if q in self._subscribers: self._subscribers.remove(q) def iter_sse(self) -> Iterator[str]: q: queue.Queue[str | None] = queue.Queue(maxsize=32) with self._lock: self._subscribers.append(q) try: yield _sse_frame(self.event_dict()) while True: try: raw = q.get(timeout=INSTANCE_DASHBOARD_SSE_HEARTBEAT_SEC) except queue.Empty: yield ": heartbeat\n\n" continue if raw is None: break try: data = json.loads(raw) except Exception: data = self.event_dict() yield _sse_frame(data) finally: with self._lock: if q in self._subscribers: self._subscribers.remove(q) def _sse_frame(data: dict[str, Any]) -> str: body = json.dumps(data, ensure_ascii=False) return f"event: dashboard\ndata: {body}\n\n" instance_dashboard_store = InstanceDashboardStore()