diff --git a/backend/app.py b/backend/app.py index 553c658..c6fbedf 100644 --- a/backend/app.py +++ b/backend/app.py @@ -136,6 +136,7 @@ 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 eaf1b67..e72b26a 100644 --- a/backend/routers/events.py +++ b/backend/routers/events.py @@ -1,89 +1,162 @@ -"""SSE-Eventstrom (UMBAU v3 P3a) — ein Kanal sagt der Zentrale, WANN neu laden lohnt. +"""SSE-Eventstrom /api/events — Server-Sent Events für Frontend-Updates. -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. +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) -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. +Event-Schema: +type: "invalidate" +data: {"keys": ["system-status", "models", "agent-status", "ideen", "token-stats"]} """ import asyncio import json import logging -from pathlib import Path +import time +from typing import Any, AsyncGenerator from fastapi import APIRouter, Request from fastapi.responses import StreamingResponse -from config import MODELS_DIR +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 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 +# 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 -def _mtime(p: Path) -> float: +async def _get_system_status() -> dict: try: - return p.stat().st_mtime - except OSError: - return 0.0 + return system_status() + except Exception: + return {} -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"] +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", []))} 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 []]) + return {"count": 0} + + +async def _get_token_stats() -> dict: + try: + return get_stats() 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 + 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 @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"}) +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 + }, + )