SPB Git forge

spb/llm-api

Public
0commits 0branches 0releases
0 Bsize
maindefault branch
—last push
1.4 KB · 49 lines python
Raw Blame History
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