"""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 (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. 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. """ 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 # 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(): # 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"})