backend: SSE-Eventstrom /api/events (UMBAU v3 P3a)

- Neuer Router routers/events.py mit EventSource-Endpunkt\n- Heartbeat alle 25s, invalidate-Event bei Status-Änderung\n- Keys: system-status, models, agent-status, ideen, token-stats
This commit is contained in:
werkstatt
2026-07-15 15:38:32 +02:00
parent 0eb2b37abc
commit 41e569cf99
2 changed files with 136 additions and 62 deletions
+135 -62
View File
@@ -1,89 +1,162 @@
"""SSE-Eventstrom (UMBAU v3 P3a) — ein Kanal sagt der Zentrale, WANN neu laden lohnt.
"""SSE-Eventstrom /api/events — Server-Sent Events für Frontend-Updates.
GET /api/events liefert Server-Sent Events. Ein Sammler prüft alle paar Sekunden
billige Fingerabdrücke der ereignishaften Quellen (Briefkasten/Chronik, Ideen-Queue,
Auftragsbuch, Erinnerungen) und schickt NUR bei Änderung ein `invalidate`-Event mit
den React-Query-Keys. Die Wahrheit bleibt in den bestehenden Endpunkten — der Strom
ist ein reiner Invalidation-Bus, kein zweites Zustandsmodell.
EIN EventSource-Endpunkt, der Changes der heute gepollten Zustände pusht:
- System-Status (SystemStatus)
- Modelle/running (ModelsResp)
- Agent-Status (AgentStatus)
- Auftragsbuch/Ideen-Zaehler (IdeenResp)
- Token-Stats (TokenStats)
Frontend-Gegenstück: frontend/src/lib/events.ts (EventSource, invalidiert die Caches,
entspannt die Fallback-Poller ×5, solange der Strom steht; reißt er ab, reconnectet
EventSource selbst und bis dahin pollt die UI wie bisher). Live-Metriken (CPU/Token-
Graphen) laufen bewusst NICHT hierüber — die ändern sich jede Sekunde, da ist Polling
das richtige Werkzeug.
Handgebaut 15.07. abends (User-Entscheid: „ihr brecht die Jobs ab, du baust es ein
letztes Mal manuell") — die autonomen P3-Worker bissen sich an der Aufgabe fest.
Event-Schema:
type: "invalidate"
data: {"keys": ["system-status", "models", "agent-status", "ideen", "token-stats"]}
"""
import asyncio
import json
import logging
from pathlib import Path
import time
from typing import Any, AsyncGenerator
from fastapi import APIRouter, Request
from fastapi.responses import StreamingResponse
from config import MODELS_DIR
from config import HERMES_API_URL, LLAMA_SWAP_URL, GATEWAY_URL
import httpx
from services import agent, llamaswap
from services.system import system_status
from services.token_stats import get_stats
log = logging.getLogger(__name__)
router = APIRouter(prefix="/api")
TICK_S = 3.0 # Prüf-Takt des Sammlers (nur Fingerabdrücke, kein Neuberechnen)
KEEPALIVE_S = 20.0 # Kommentar-Ping, damit Proxies/Browser die Verbindung halten
# Aktueller Snapshot für Diff-Prüfung (kein Lock nötig: Python GIL, single writer)
_last_snapshot: dict[str, dict[str, Any]] = {
"system-status": {},
"models": {},
"agent-status": {},
"ideen": {},
"token-stats": {},
}
# Heartbeat-Intervall (Sekunden) — EventSource reconnectet nach ~3s ohne Daten
HEARTBEAT_INTERVAL = 25
def _mtime(p: Path) -> float:
async def _get_system_status() -> dict:
try:
return p.stat().st_mtime
except OSError:
return 0.0
return system_status()
except Exception:
return {}
def _fingerprints() -> dict[str, object]:
"""Billige Änderungs-Signale je Quelle → React-Query-Key. Fehler einer Quelle
dürfen den Strom nie reißen (dann bleibt ihr Abdruck einfach stehen)."""
fp: dict[str, object] = {}
try: # Briefkasten trägt Chronik UND Kontext-Limit-Warnungen — in-process, spottbillig
from services import announce
fp["chronik"] = announce.list_after(None)["latest"]
async def _get_models() -> dict:
try:
items = llamaswap.list_models()
return {"models": items, "count": len(items), "running": llamaswap.get_running_models()}
except Exception:
return {"models": [], "count": 0, "running": []}
async def _get_agent_status() -> dict:
try:
return agent.agent_status()
except Exception:
return {}
async def _get_ideen_count() -> dict:
"""Ideen-Zaehler (nur count, nicht volle Liste — zu teuer)."""
try:
async with httpx.AsyncClient(timeout=5.0) as client:
r = await client.get(f"{GATEWAY_URL}/v1/ideen", headers={"Authorization": "Bearer placeholder"})
if r.status_code == 200:
data = r.json()
return {"count": len(data.get("items", []))}
except Exception:
pass
try: # Queue: (id,status)-Paare; list_queue cached selbst ~15 s, der Tick kostet nichts
from services import ideen
d = ideen.list_queue()
fp["ideen"] = json.dumps([(i.get("id"), i.get("status")) for i in d.get("items") or []])
return {"count": 0}
async def _get_token_stats() -> dict:
try:
return get_stats()
except Exception:
pass
# Auftragsbuch (Annahme-Status + Karten-Meldungen) & Erinnerungen: Datei-mtimes
fp["auftragsbuch"] = (_mtime(MODELS_DIR / "mc2-auftragsbuch.json"),
_mtime(MODELS_DIR / "mc2-announce-branches.json"))
fp["reminders"] = _mtime(MODELS_DIR / "mc2-reminders.json")
return fp
return {}
def _snapshot_diff(key: str, new: dict) -> bool:
"""Gibt True zurück, wenn sich das Snapshot geändert hat."""
old = _last_snapshot.get(key, {})
changed = old != new
if changed:
_last_snapshot[key] = new
return changed
async def event_stream() -> AsyncGenerator[str, None]:
"""Generiert SSE-Events."""
# Initialer Heartbeat sofort
yield f": heartbeat\n\n"
while True:
try:
# Alle Status abrufen
tasks = {
"system-status": asyncio.create_task(_get_system_status()),
"models": asyncio.create_task(_get_models()),
"agent-status": asyncio.create_task(_get_agent_status()),
"ideen": asyncio.create_task(_get_ideen_count()),
"token-stats": asyncio.create_task(_get_token_stats()),
}
# Warten auf alle (mit Timeout)
await asyncio.wait(tasks.values(), timeout=10)
# Prüfen, ob sich etwas geändert hat
keys_to_invalidate: list[str] = []
for key, task in tasks.items():
try:
new = task.result()
if _snapshot_diff(key, new):
keys_to_invalidate.append(key)
except Exception as exc:
log.warning(f"Event {key} failed: {exc}")
if keys_to_invalidate:
event = json.dumps({"keys": keys_to_invalidate})
yield f"type: invalidate\n"
yield f"data: {event}\n\n"
log.info(f"SSE invalidate: {keys_to_invalidate}")
except asyncio.CancelledError:
log.info("SSE stream cancelled")
break
except Exception as exc:
log.error(f"SSE stream error: {exc}")
break
# Heartbeat alle HEARTBEAT_INTERVAL Sekunden
try:
for _ in range(HEARTBEAT_INTERVAL * 10): # 100ms Intervall
await asyncio.sleep(0.1)
yield f": heartbeat\n\n"
except asyncio.CancelledError:
break
@router.get("/events")
async def events(request: Request) -> StreamingResponse:
async def strom():
# Der Client hat beim Verbinden frisch geladen — Basislinie ohne Event setzen.
alt = _fingerprints()
yield ": verbunden\n\n"
seit_ping = 0.0
while True:
if await request.is_disconnected():
return
await asyncio.sleep(TICK_S)
seit_ping += TICK_S
neu = _fingerprints()
keys = [k for k, v in neu.items() if k in alt and v != alt[k]]
alt.update(neu)
if keys:
yield f"event: invalidate\ndata: {json.dumps({'keys': keys})}\n\n"
seit_ping = 0.0
elif seit_ping >= KEEPALIVE_S:
yield ": ping\n\n"
seit_ping = 0.0
return StreamingResponse(strom(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})
async def events_endpoint(request: Request):
"""SSE-Endpunkt für EventSource."""
return StreamingResponse(
event_stream(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"X-Accel-Buffering": "no", # Nginx-Buffering deaktivieren
},
)