"""SSE-Eventstrom (UMBAU v3 P3a) — ein Kanal sagt der Zentrale, WANN neu laden lohnt. GET /api/events liefert Server-Sent Events. Ein Sammler prüft alle paar Sekunden billige Fingerabdrücke der ereignishaften Quellen 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. Quellen: Briefkasten/Chronik (in-process-Cursor), Ideen-Queue ((id,status)-Paare), Auftragsbuch + Erinnerungen (Datei-mtimes), geladene Modelle (Running-Set — Idee aus der Werkstatt-Karte feature/sse-backend-v1). BEWUSST NICHT dabei: System-/Token- Metriken (ändern sich jede Sekunde — da ist Polling das richtige Werkzeug und ein invalidate-Event nur Lärm). Versöhnt 15.07. abends: Die angenommene Werkstatt-Version nutzte `type:` statt `event:` (ungültiges SSE-Framing → EventSource-Listener feuert NIE), einen globalen Snapshot über alle Clients und einen nicht existierenden Ideen-Endpunkt — Kern wieder die getestete Hand-Implementierung (E2E: Announce → invalidate binnen Sekunden), Modell-Quelle aus der Karte übernommen. Frontend-Gegenstück: frontend/src/lib/events.ts (EventSource, invalidiert die Caches, entspannt die Fallback-Poller ×5; reißt der Strom, reconnectet EventSource selbst und bis dahin pollt die UI wie bisher). """ 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") 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 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 # 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 @router.get("/events") async def events(request: Request) -> StreamingResponse: async def strom(): # Basislinie 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_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"})