From e91f6b7ba8db7a977184d0b1e4dd561493ef97bd Mon Sep 17 00:00:00 2001 From: Hitonabi Date: Wed, 15 Jul 2026 19:30:27 +0200 Subject: [PATCH] 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 --- backend/app.py | 1 - backend/routers/events.py | 209 ++++++++++++++------------------------ 2 files changed, 74 insertions(+), 136 deletions(-) diff --git a/backend/app.py b/backend/app.py index c6fbedf..553c658 100644 --- a/backend/app.py +++ b/backend/app.py @@ -136,7 +136,6 @@ app.include_router(maintenance.router) 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(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(events.router) # SSE-Eventstrom /api/events (P3a) — Invalidation-Bus app.include_router(wissen.router) # Wissens-Vault (Traum-Notizen) read-only diff --git a/backend/routers/events.py b/backend/routers/events.py index e72b26a..25f43b8 100644 --- a/backend/routers/events.py +++ b/backend/routers/events.py @@ -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: -- System-Status (SystemStatus) -- Modelle/running (ModelsResp) -- Agent-Status (AgentStatus) -- Auftragsbuch/Ideen-Zaehler (IdeenResp) -- Token-Stats (TokenStats) +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. -Event-Schema: -type: "invalidate" -data: {"keys": ["system-status", "models", "agent-status", "ideen", "token-stats"]} +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 -import time -from typing import Any, AsyncGenerator +from pathlib import Path from fastapi import APIRouter, Request from fastapi.responses import StreamingResponse -from config import HERMES_API_URL, LLAMA_SWAP_URL, GATEWAY_URL -import httpx - -from services import agent, llamaswap -from services.system import system_status -from services.token_stats import get_stats +from config import MODELS_DIR log = logging.getLogger(__name__) router = APIRouter(prefix="/api") -# Aktueller Snapshot für Diff-Prüfung (kein Lock nötig: Python GIL, single writer) -_last_snapshot: dict[str, dict[str, Any]] = { - "system-status": {}, - "models": {}, - "agent-status": {}, - "ideen": {}, - "token-stats": {}, -} - -# Heartbeat-Intervall (Sekunden) — EventSource reconnectet nach ~3s ohne Daten -HEARTBEAT_INTERVAL = 25 +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 -async def _get_system_status() -> dict: +def _mtime(p: Path) -> float: try: - return system_status() - except Exception: - return {} + return p.stat().st_mtime + except OSError: + return 0.0 -async def _get_models() -> dict: - try: - items = llamaswap.list_models() - return {"models": items, "count": len(items), "running": llamaswap.get_running_models()} - except Exception: - return {"models": [], "count": 0, "running": []} - - -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", []))} +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 - return {"count": 0} - - -async def _get_token_stats() -> dict: - try: - return get_stats() + 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: - return {} - - -def _snapshot_diff(key: str, new: dict) -> bool: - """Gibt True zurück, wenn sich das Snapshot geändert hat.""" - old = _last_snapshot.get(key, {}) - changed = old != new - if changed: - _last_snapshot[key] = new - return changed - - -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 + 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_endpoint(request: Request): - """SSE-Endpunkt für EventSource.""" - return StreamingResponse( - event_stream(), - media_type="text/event-stream", - headers={ - "Cache-Control": "no-cache", - "Connection": "keep-alive", - "X-Accel-Buffering": "no", # Nginx-Buffering deaktivieren - }, - ) +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"})