"""Ereignisstrom — ein Kanal sagt der Zentrale, WANN neu laden lohnt, und schickt Metriken. GET /api/stream (v3-Umbau P4, 28.08.2026) — zwei Ereignisarten: · `invalidate` Nur bei Änderung, mit den betroffenen React-Query-Schlüsseln. · `metrik` Jede Sekunde ein Messpunkt (CPU/RAM/GPU/Temp/Token-Zähler). Der alte Kanal /api/events ist mit dem Box-Wart-Umbau (09/2026) entfallen. WARUM METRIKEN JETZT MITKOMMEN: Bis P4 pollte das Frontend `/api/system/status` und `/api/system/token-stats` im 3-Sekunden-Takt — zwei Dauer-Anfragen, unabhängig davon, ob sich etwas geändert hat, plus sechs weitere langsamere Poller auf der Startseite. Der Messpunkt kostet hier 0,2 ms (gemessen); `system_status()` würde 100 ms kosten, weil `psutil.cpu_percent(interval=0.1)` wartet. Deshalb der eigene, leichte `metrik_punkt()`. WAS BEWUSST NICHT DRIN IST: Ein `agent`-Thema für Lucys Denkschritte. MC2 kann Hermes' interne Schritte nicht sehen, ohne dessen Quellcode zu patchen — und das ist per AGENTS.md verboten. Eine leere Leitung zu bauen, wäre eine Zusage, die keiner einlöst. Die Wahrheit bleibt in den bestehenden Endpunkten: `invalidate` ist ein reiner Anstoß-Bus, kein zweites Zustandsmodell. `metrik` ist die einzige Ausnahme — es ist der Wert selbst, weil ein Anstoß für eine Zahl, die sich jede Sekunde ändert, nur Lärm wäre. Frontend-Gegenstück: frontend/src/lib/events.ts """ import asyncio import json import logging from pathlib import Path from config import MODELS_DIR from fastapi import APIRouter, Request from fastapi.responses import StreamingResponse log = logging.getLogger(__name__) router = APIRouter(prefix="/api") METRIK_S = 1.0 # Takt der Messpunkte ABDRUCK_S = 3.0 # Takt der Änderungs-Prüfung (nur Fingerabdrücke, kein Neuberechnen) KEEPALIVE_S = 20.0 # Kommentar-Ping, damit Proxies/Browser die Verbindung halten def _mtime(p: Path) -> float: try: return p.stat().st_mtime except OSError: return 0.0 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: # Wächter-Stand (mc2-steward schreibt ihn jede Minute): Hinweise → Startseite neu laden from services import waechter fp["start"] = _mtime(waechter.STORE_PATH) except Exception: pass try: # Geladene Modelle (Running-Set): Laden/Entladen soll die Modelle-Ansicht anstoßen from services import llamaswap fp["models"] = json.dumps(sorted(str(m) for m in llamaswap.get_running_models())) except Exception: pass try: # Jobs (Downloads, Wartung): Zustand + Fortschritt — spart den 3-s-Poller der Schublade from services import jobengine fp["jobs"] = json.dumps( [(j.get("id"), j.get("state"), j.get("progress")) for j in jobengine.public_jobs()] ) except Exception: pass # Erinnerungen: Datei-mtime (das Auftragsbuch ist seit 04.09.2026 ausgebaut) fp["reminders"] = _mtime(MODELS_DIR / "mc2-reminders.json") return fp async def _strom(request: Request, mit_metrik: bool): """Gemeinsamer Kern beider Endpunkte. Die Basislinie entsteht JE VERBINDUNG (der Client hat beim Verbinden frisch geladen) — ein globaler Snapshot würde bei mehreren Clients Events verschlucken. """ alt = _fingerprints() yield ": verbunden\n\n" seit_abdruck = 0.0 seit_ping = 0.0 takt = METRIK_S if mit_metrik else ABDRUCK_S while True: if await request.is_disconnected(): return await asyncio.sleep(takt) seit_abdruck += takt seit_ping += takt if mit_metrik: try: from services.system import metrik_punkt yield f"event: metrik\ndata: {json.dumps(metrik_punkt())}\n\n" seit_ping = 0.0 except Exception: # Ein kaputter Messpunkt darf den Strom nicht reißen — die Ansicht fällt # dann auf ihre Poller zurück, das ist besser als eine tote Leitung. log.warning("Messpunkt fehlgeschlagen", exc_info=True) if seit_abdruck >= ABDRUCK_S: seit_abdruck = 0.0 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 if seit_ping >= KEEPALIVE_S: yield ": ping\n\n" seit_ping = 0.0 _KOPF = {"Cache-Control": "no-cache", "X-Accel-Buffering": "no"} @router.get("/stream") async def stream(request: Request) -> StreamingResponse: """Anstöße UND Messpunkte.""" return StreamingResponse(_strom(request, mit_metrik=True), media_type="text/event-stream", headers=_KOPF)