"""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 fastapi import APIRouter, Request from fastapi.responses import StreamingResponse from config import MODELS_DIR 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"})