Align instance dashboard with snapshot+SSE and options net PnL.

Only dashboard path changes: background snapshot, SSE refresh, smaller header type, options net PnL.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
dekun
2026-07-17 22:10:58 +08:00
parent 4fe0df6c6d
commit 1d2bdc3aa1
14 changed files with 510 additions and 99 deletions
+175
View File
@@ -0,0 +1,175 @@
"""实例数据看板:后台定时聚合,内存快照,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()