"""Ereignisstrom — ein Kanal sagt der Zentrale, WANN neu laden lohnt, und schickt Metriken. Zwei Endpunkte, ein Sammler: GET /api/stream (v3-Umbau P4, 28.08.2026) — der aktuelle Kanal. 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). GET /api/events — der alte Kanal, nur `invalidate`. Bleibt EINE Fassung lang stehen, weil ein Browser-Tab nach einem Deploy noch das vorige Bündel halten kann und dieses nur `/api/events` kennt. Danach entfernen. 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: # Briefkasten trägt Chronik UND Kontext-Limit-Warnungen — in-process, spottbillig from services import announce fp["chronik"] = announce.list_after(None)["latest"] 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 []]) 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: """Der aktuelle Kanal: Anstöße UND Messpunkte.""" return StreamingResponse(_strom(request, mit_metrik=True), media_type="text/event-stream", headers=_KOPF) @router.get("/events") async def events(request: Request) -> StreamingResponse: """Alt-Kanal ohne Messpunkte. Nur für Browser-Tabs, die noch ein Bündel von vor dem 28.08.2026 halten. Mit der übernächsten Fassung entfernen.""" return StreamingResponse(_strom(request, mit_metrik=False), media_type="text/event-stream", headers=_KOPF)