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>
This commit is contained in:
@@ -136,7 +136,6 @@ app.include_router(maintenance.router)
|
|||||||
app.include_router(auftragsbuch.router) # Vorschlags-Inbox (Mensch-Gate als Klick)
|
app.include_router(auftragsbuch.router) # Vorschlags-Inbox (Mensch-Gate als Klick)
|
||||||
app.include_router(ideen.router) # Ideen-Queue (natives Hermes-Kanban) — Tür der Zentrale
|
app.include_router(ideen.router) # Ideen-Queue (natives Hermes-Kanban) — Tür der Zentrale
|
||||||
app.include_router(chronik.router) # Timeline der autonomen Taten (Announce-Store)
|
app.include_router(chronik.router) # Timeline der autonomen Taten (Announce-Store)
|
||||||
app.include_router(events.router) # SSE-Eventstrom /api/events (UMBAU v3 P3a)
|
|
||||||
app.include_router(eigenleben.router) # „Von allein": Skills + Vorschlags-Bilanz der Box
|
app.include_router(eigenleben.router) # „Von allein": Skills + Vorschlags-Bilanz der Box
|
||||||
app.include_router(events.router) # SSE-Eventstrom /api/events (P3a) — Invalidation-Bus
|
app.include_router(events.router) # SSE-Eventstrom /api/events (P3a) — Invalidation-Bus
|
||||||
app.include_router(wissen.router) # Wissens-Vault (Traum-Notizen) read-only
|
app.include_router(wissen.router) # Wissens-Vault (Traum-Notizen) read-only
|
||||||
|
|||||||
+74
-135
@@ -1,162 +1,101 @@
|
|||||||
"""SSE-Eventstrom /api/events — Server-Sent Events für Frontend-Updates.
|
"""SSE-Eventstrom (UMBAU v3 P3a) — ein Kanal sagt der Zentrale, WANN neu laden lohnt.
|
||||||
|
|
||||||
EIN EventSource-Endpunkt, der Changes der heute gepollten Zustände pusht:
|
GET /api/events liefert Server-Sent Events. Ein Sammler prüft alle paar Sekunden
|
||||||
- System-Status (SystemStatus)
|
billige Fingerabdrücke der ereignishaften Quellen und schickt NUR bei Änderung ein
|
||||||
- Modelle/running (ModelsResp)
|
`invalidate`-Event mit den React-Query-Keys. Die Wahrheit bleibt in den bestehenden
|
||||||
- Agent-Status (AgentStatus)
|
Endpunkten — der Strom ist ein reiner Invalidation-Bus, kein zweites Zustandsmodell.
|
||||||
- Auftragsbuch/Ideen-Zaehler (IdeenResp)
|
|
||||||
- Token-Stats (TokenStats)
|
|
||||||
|
|
||||||
Event-Schema:
|
Quellen: Briefkasten/Chronik (in-process-Cursor), Ideen-Queue ((id,status)-Paare),
|
||||||
type: "invalidate"
|
Auftragsbuch + Erinnerungen (Datei-mtimes), geladene Modelle (Running-Set — Idee aus
|
||||||
data: {"keys": ["system-status", "models", "agent-status", "ideen", "token-stats"]}
|
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 asyncio
|
||||||
import json
|
import json
|
||||||
import logging
|
import logging
|
||||||
import time
|
from pathlib import Path
|
||||||
from typing import Any, AsyncGenerator
|
|
||||||
|
|
||||||
from fastapi import APIRouter, Request
|
from fastapi import APIRouter, Request
|
||||||
from fastapi.responses import StreamingResponse
|
from fastapi.responses import StreamingResponse
|
||||||
|
|
||||||
from config import HERMES_API_URL, LLAMA_SWAP_URL, GATEWAY_URL
|
from config import MODELS_DIR
|
||||||
import httpx
|
|
||||||
|
|
||||||
from services import agent, llamaswap
|
|
||||||
from services.system import system_status
|
|
||||||
from services.token_stats import get_stats
|
|
||||||
|
|
||||||
log = logging.getLogger(__name__)
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
router = APIRouter(prefix="/api")
|
router = APIRouter(prefix="/api")
|
||||||
|
|
||||||
# Aktueller Snapshot für Diff-Prüfung (kein Lock nötig: Python GIL, single writer)
|
TICK_S = 3.0 # Prüf-Takt des Sammlers (nur Fingerabdrücke, kein Neuberechnen)
|
||||||
_last_snapshot: dict[str, dict[str, Any]] = {
|
KEEPALIVE_S = 20.0 # Kommentar-Ping, damit Proxies/Browser die Verbindung halten
|
||||||
"system-status": {},
|
|
||||||
"models": {},
|
|
||||||
"agent-status": {},
|
|
||||||
"ideen": {},
|
|
||||||
"token-stats": {},
|
|
||||||
}
|
|
||||||
|
|
||||||
# Heartbeat-Intervall (Sekunden) — EventSource reconnectet nach ~3s ohne Daten
|
|
||||||
HEARTBEAT_INTERVAL = 25
|
|
||||||
|
|
||||||
|
|
||||||
async def _get_system_status() -> dict:
|
def _mtime(p: Path) -> float:
|
||||||
try:
|
try:
|
||||||
return system_status()
|
return p.stat().st_mtime
|
||||||
except Exception:
|
except OSError:
|
||||||
return {}
|
return 0.0
|
||||||
|
|
||||||
|
|
||||||
async def _get_models() -> dict:
|
def _fingerprints() -> dict[str, object]:
|
||||||
try:
|
"""Billige Änderungs-Signale je Quelle → React-Query-Key. Fehler einer Quelle
|
||||||
items = llamaswap.list_models()
|
dürfen den Strom nie reißen (dann bleibt ihr Abdruck einfach stehen)."""
|
||||||
return {"models": items, "count": len(items), "running": llamaswap.get_running_models()}
|
fp: dict[str, object] = {}
|
||||||
except Exception:
|
try: # Briefkasten trägt Chronik UND Kontext-Limit-Warnungen — in-process, spottbillig
|
||||||
return {"models": [], "count": 0, "running": []}
|
from services import announce
|
||||||
|
fp["chronik"] = announce.list_after(None)["latest"]
|
||||||
|
|
||||||
async def _get_agent_status() -> dict:
|
|
||||||
try:
|
|
||||||
return agent.agent_status()
|
|
||||||
except Exception:
|
|
||||||
return {}
|
|
||||||
|
|
||||||
|
|
||||||
async def _get_ideen_count() -> dict:
|
|
||||||
"""Ideen-Zaehler (nur count, nicht volle Liste — zu teuer)."""
|
|
||||||
try:
|
|
||||||
async with httpx.AsyncClient(timeout=5.0) as client:
|
|
||||||
r = await client.get(f"{GATEWAY_URL}/v1/ideen", headers={"Authorization": "Bearer placeholder"})
|
|
||||||
if r.status_code == 200:
|
|
||||||
data = r.json()
|
|
||||||
return {"count": len(data.get("items", []))}
|
|
||||||
except Exception:
|
except Exception:
|
||||||
pass
|
pass
|
||||||
return {"count": 0}
|
try: # Queue: (id,status)-Paare; list_queue cached selbst ~15 s, der Tick kostet nichts
|
||||||
|
from services import ideen
|
||||||
|
d = ideen.list_queue()
|
||||||
async def _get_token_stats() -> dict:
|
fp["ideen"] = json.dumps([(i.get("id"), i.get("status")) for i in d.get("items") or []])
|
||||||
try:
|
|
||||||
return get_stats()
|
|
||||||
except Exception:
|
except Exception:
|
||||||
return {}
|
pass
|
||||||
|
try: # Geladene Modelle (Running-Set): Laden/Entladen soll die Modelle-Ansicht anstoßen
|
||||||
|
from services import llamaswap
|
||||||
def _snapshot_diff(key: str, new: dict) -> bool:
|
fp["models"] = json.dumps(sorted(str(m) for m in llamaswap.get_running_models()))
|
||||||
"""Gibt True zurück, wenn sich das Snapshot geändert hat."""
|
except Exception:
|
||||||
old = _last_snapshot.get(key, {})
|
pass
|
||||||
changed = old != new
|
# Auftragsbuch (Annahme-Status + Karten-Meldungen) & Erinnerungen: Datei-mtimes
|
||||||
if changed:
|
fp["auftragsbuch"] = (_mtime(MODELS_DIR / "mc2-auftragsbuch.json"),
|
||||||
_last_snapshot[key] = new
|
_mtime(MODELS_DIR / "mc2-announce-branches.json"))
|
||||||
return changed
|
fp["reminders"] = _mtime(MODELS_DIR / "mc2-reminders.json")
|
||||||
|
return fp
|
||||||
|
|
||||||
async def event_stream() -> AsyncGenerator[str, None]:
|
|
||||||
"""Generiert SSE-Events."""
|
|
||||||
# Initialer Heartbeat sofort
|
|
||||||
yield f": heartbeat\n\n"
|
|
||||||
|
|
||||||
while True:
|
|
||||||
try:
|
|
||||||
# Alle Status abrufen
|
|
||||||
tasks = {
|
|
||||||
"system-status": asyncio.create_task(_get_system_status()),
|
|
||||||
"models": asyncio.create_task(_get_models()),
|
|
||||||
"agent-status": asyncio.create_task(_get_agent_status()),
|
|
||||||
"ideen": asyncio.create_task(_get_ideen_count()),
|
|
||||||
"token-stats": asyncio.create_task(_get_token_stats()),
|
|
||||||
}
|
|
||||||
|
|
||||||
# Warten auf alle (mit Timeout)
|
|
||||||
await asyncio.wait(tasks.values(), timeout=10)
|
|
||||||
|
|
||||||
# Prüfen, ob sich etwas geändert hat
|
|
||||||
keys_to_invalidate: list[str] = []
|
|
||||||
|
|
||||||
for key, task in tasks.items():
|
|
||||||
try:
|
|
||||||
new = task.result()
|
|
||||||
if _snapshot_diff(key, new):
|
|
||||||
keys_to_invalidate.append(key)
|
|
||||||
except Exception as exc:
|
|
||||||
log.warning(f"Event {key} failed: {exc}")
|
|
||||||
|
|
||||||
if keys_to_invalidate:
|
|
||||||
event = json.dumps({"keys": keys_to_invalidate})
|
|
||||||
yield f"type: invalidate\n"
|
|
||||||
yield f"data: {event}\n\n"
|
|
||||||
log.info(f"SSE invalidate: {keys_to_invalidate}")
|
|
||||||
|
|
||||||
except asyncio.CancelledError:
|
|
||||||
log.info("SSE stream cancelled")
|
|
||||||
break
|
|
||||||
except Exception as exc:
|
|
||||||
log.error(f"SSE stream error: {exc}")
|
|
||||||
break
|
|
||||||
|
|
||||||
# Heartbeat alle HEARTBEAT_INTERVAL Sekunden
|
|
||||||
try:
|
|
||||||
for _ in range(HEARTBEAT_INTERVAL * 10): # 100ms Intervall
|
|
||||||
await asyncio.sleep(0.1)
|
|
||||||
yield f": heartbeat\n\n"
|
|
||||||
except asyncio.CancelledError:
|
|
||||||
break
|
|
||||||
|
|
||||||
|
|
||||||
@router.get("/events")
|
@router.get("/events")
|
||||||
async def events_endpoint(request: Request):
|
async def events(request: Request) -> StreamingResponse:
|
||||||
"""SSE-Endpunkt für EventSource."""
|
async def strom():
|
||||||
return StreamingResponse(
|
# Basislinie JE VERBINDUNG (der Client hat beim Verbinden frisch geladen) —
|
||||||
event_stream(),
|
# ein globaler Snapshot würde bei mehreren Clients Events verschlucken.
|
||||||
media_type="text/event-stream",
|
alt = _fingerprints()
|
||||||
headers={
|
yield ": verbunden\n\n"
|
||||||
"Cache-Control": "no-cache",
|
seit_ping = 0.0
|
||||||
"Connection": "keep-alive",
|
while True:
|
||||||
"X-Accel-Buffering": "no", # Nginx-Buffering deaktivieren
|
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"})
|
||||||
|
|||||||
Reference in New Issue
Block a user