1"""In-process event bus for real-time dashboard updates (SSE)."""23from __future__ import annotations45import asyncio6import json7import time8from collections import deque9from typing import Any101112class EventBus:13 def __init__(self, history: int = 200):14 self._subs: set[asyncio.Queue] = set()15 self._history: deque[dict] = deque(maxlen=history)16 self._seq = 01718 def publish(self, kind: str, data: Any) -> None:19 self._seq += 120 ev = {"seq": self._seq, "ts": time.time(), "type": kind, "data": data}21 self._history.append(ev)22 dead = []23 for q in self._subs:24 try:25 q.put_nowait(ev)26 except asyncio.QueueFull:27 dead.append(q)28 for q in dead:29 self._subs.discard(q)3031 def subscribe(self) -> asyncio.Queue:32 q: asyncio.Queue = asyncio.Queue(maxsize=500)33 self._subs.add(q)34 return q3536 def unsubscribe(self, q: asyncio.Queue) -> None:37 self._subs.discard(q)3839 def recent(self, kinds: set[str] | None = None, limit: int = 50) -> list[dict]:40 evs = [e for e in self._history if not kinds or e["type"] in kinds]41 return evs[-limit:]4243 @staticmethod44 def format_sse(ev: dict) -> str:45 return f"id: {ev['seq']}\nevent: {ev['type']}\ndata: {json.dumps(ev, ensure_ascii=False)}\n\n"464748bus = EventBus()49