Files
mission-control-v2/backend/routers/events.py
T
Hitonabi e91f6b7ba8 SSE-Eventstrom versoehnt: getesteter Kern + Modell-Quelle aus der Werkstatt-Karte
Die angenommene Karte feature/sse-backend-v1 ersetzte die bereits E2E-
getestete Hand-Implementierung mit drei harten Fehlern: 'type:' statt
'event:' (ungueltiges SSE-Framing - EventSource-Listener feuern NIE),
globaler Snapshot ueber alle Clients (verschluckt Events), Ideen-Zaehler
gegen nicht existierenden Endpunkt :9010/v1/ideen. Dazu doppelter
Router-Mount in app.py (Merge-Folge).

Zurueck auf den bewiesenen Kern (3s-Fingerprints, Basislinie je
Verbindung, is_disconnected, korrektes Framing) + die GUTE Idee der
Karte uebernommen: Running-Set der Modelle als Quelle (Laden/Entladen
invalidiert die Modelle-Ansicht). System-/Token-Metriken bleiben bewusst
beim Polling (aendern sich jede Sekunde = invalidate-Laerm).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-15 19:30:27 +02:00

102 lines
4.2 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""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"})