47f7a85510
Die Ampel-Nachruestung deckte 317/337 vorbestehende ruff-Verstoesse im ganzen Repo auf. Aufgeraeumt: - ruff.toml: intentionale Muster als Projekt-Politik ausgenommen (BLE001 blind-except, S110/S112 try-except-pass/continue, PLW1510 subprocess-best-effort, B008 FastAPI- Depends/File-Idiom, EXE001 Shebang, + wenige Stil-Regeln). __init__.py-Re-Exports geschuetzt (F401). - ruff --fix: 128 mechanische (Import-Sortierung, PEP585/604-Annotationen, tote Imports, ueberfluessige noqa) auto-behoben. - 12 echte Reste von Hand: PERF402/102, PLC3002 (Lambda->walrus), ISC004 (String-Concat geklammert), F841/RUF059 (ungenutzte Vars), PIE810 (startswith-Tuple), UP031 (f-string), UP035 (veraltete typing-Imports). Ergebnis: 'ruff check .' = 0, 'compileall' grün. Kein Verhaltenswechsel (nur Stil/Modernisierung). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
101 lines
4.2 KiB
Python
101 lines
4.2 KiB
Python
"""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"})
|