From 019bd6d18b8e270a2667287cc565232b86680ead Mon Sep 17 00:00:00 2001 From: Hitonabi Date: Fri, 3 Jul 2026 13:45:07 +0200 Subject: [PATCH] =?UTF-8?q?Lucy-Proaktivit=C3=A4t=20(A3):=20Melde-Briefkas?= =?UTF-8?q?ten=20+=20Health-W=C3=A4chter?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - services/announce.py: persistenter Briefkasten (/srv/models/mc2-announce.json), POST /api/voice/announce + GET /api/voice/announcements (Cursor-Polling) - services/sentry.py: Health-Wächter (Engine/Hirn/Hermes/Mem0/Voice/Platte), flankenerkannt (Alarm nach 3 Fehl-Ticks, Entwarnung, 6h-Erinnerung), meldet in Briefkasten + Telegram; Hirn-Verdrängung durch IDE-Last = kein Alarm - notify.sh spiegelt jede Telegram-Meldung in den Briefkasten (Updates/Radar erreichen damit auch die Desktop-Lucy) Co-Authored-By: Claude Fable 5 --- backend/app.py | 14 ++-- backend/routers/voice.py | 25 ++++++ backend/services/announce.py | 82 +++++++++++++++++++ backend/services/sentry.py | 152 +++++++++++++++++++++++++++++++++++ deploy/notify.sh | 13 +++ 5 files changed, 281 insertions(+), 5 deletions(-) create mode 100644 backend/services/announce.py create mode 100644 backend/services/sentry.py diff --git a/backend/app.py b/backend/app.py index d112fed..fb6337f 100644 --- a/backend/app.py +++ b/backend/app.py @@ -19,7 +19,7 @@ from starlette.requests import Request from config import FRONTEND_DIST, VERSION from routers import agent, connect, gateway_proxy, health, maintenance, memory, models, routing, system, voice -from services import warmer +from services import sentry, warmer # Zentrales Logging — Level via MC_LOG_LEVEL (INFO default). Eine Konfiguration # für alle Module (logging.getLogger(__name__)). @@ -31,14 +31,18 @@ log = logging.getLogger(__name__) @asynccontextmanager async def lifespan(app: FastAPI): - """Hintergrund-Tasks an den App-Lebenszyklus binden: Re-Warm-Wächter fürs Agent-Hirn.""" - task = asyncio.create_task(warmer.rewarm_loop()) if warmer.ENABLED else None - if task: + """Hintergrund-Tasks an den App-Lebenszyklus binden: Re-Warm-Wächter fürs Agent-Hirn + + Health-Wächter (meldet Ausfälle/Erholung in den Lucy-Briefkasten und auf Telegram).""" + tasks = [] + if warmer.ENABLED: + tasks.append(asyncio.create_task(warmer.rewarm_loop())) log.info("Hirn-Re-Warm-Wächter aktiv (Intervall %ss, Hirn dynamisch aus Hermes-Config)", warmer.INTERVAL) + if sentry.ENABLED: + tasks.append(asyncio.create_task(sentry.sentry_loop())) try: yield finally: - if task: + for task in tasks: task.cancel() diff --git a/backend/routers/voice.py b/backend/routers/voice.py index 9cbb602..a982585 100644 --- a/backend/routers/voice.py +++ b/backend/routers/voice.py @@ -19,6 +19,7 @@ from fastapi.responses import Response, StreamingResponse from pydantic import BaseModel from config import HERMES_API_KEY, HERMES_API_MODEL, HERMES_API_URL, LLAMA_SWAP_URL, VOICE_SERVICE_URL +from services import announce from services.voice_metrics import Timer, get_metrics, record_stage # Per-Stage-Latenz (C2) # Injection-Schutz (Stufe 0): guard.py liegt im mcp/-Verzeichnis. Per Pfad laden (eigene MC2-Venv). @@ -90,6 +91,30 @@ class ChatIn(BaseModel): images: list[str] = [] # optionale Bildschirm-Sicht: ein data:-URL je Monitor (Lucys „Augen") +class AnnounceIn(BaseModel): + text: str # die Meldung (wird von Lucy gesprochen) + subject: str = "" # kurze Betreffzeile (z.B. "[Update]") + source: str = "" # Absender (sentry/notify/cron …) — nur fürs Log/Panel + priority: str = "normal" # 'silent' = nur im Verlauf zeigen, nicht sprechen + + +@router.post("/voice/announce") +def voice_announce(body: AnnounceIn) -> dict: + """Meldung in den Briefkasten legen (Lucy-Proaktivität). Absender: Health-Wächter, + notify.sh (Updates/Radar/Telegram-Spiegel), Hermes-cron. LAN-only wie alle MC2-Endpoints.""" + try: + return {"ok": True, "item": announce.add(body.text, body.subject, body.source, body.priority)} + except ValueError as exc: + raise HTTPException(400, str(exc)) + + +@router.get("/voice/announcements") +def voice_announcements(after: int | None = None, limit: int = 20) -> dict: + """Neue Meldungen nach Cursor `after` abholen (Lucy pollt). Ohne `after` nur den + aktuellen Cursor-Stand (latest) — Erststart plappert so keine alten Meldungen nach.""" + return announce.list_after(after, limit) + + @router.get("/voice/metrics") def voice_metrics() -> dict: """Per-Stage-Latenz (STT/Vision/Chat-TTFB/TTS) — rollende Statistik, macht die Voice-Pipeline diff --git a/backend/services/announce.py b/backend/services/announce.py new file mode 100644 index 0000000..d871122 --- /dev/null +++ b/backend/services/announce.py @@ -0,0 +1,82 @@ +""" +Melde-Briefkasten der Box (Lucy-Proaktivität, Faden A3). + +Alles, was die Box dem Commander aktiv sagen will (Health-Wächter, Auto-Updates, +Radar, Hermes-cron via notify.sh), landet als Eintrag hier. Die Lucy-Desktop-App +pollt `/api/voice/announcements` und SPRICHT neue Einträge von sich aus. + +Persistenz als JSON neben den Modellen (übersteht Deploys/Neustarts, wie der +Discover-Cache). Bewusst klein: fortlaufende IDs als Cursor, Ring der letzten +MAX_ITEMS Einträge, ein Lock für die FastAPI-Threadpool-Worker. +""" + +import json +import logging +import os +import threading +import time +from pathlib import Path + +from config import MODELS_DIR + +log = logging.getLogger(__name__) + +STORE_PATH = Path(os.environ.get("MC_ANNOUNCE_STORE", str(MODELS_DIR / "mc2-announce.json"))) +MAX_ITEMS = int(os.environ.get("MC_ANNOUNCE_MAX", "200")) + +_lock = threading.Lock() +_state: dict | None = None # {"next_id": int, "items": [...]} + + +def _load() -> dict: + global _state + if _state is None: + try: + _state = json.loads(STORE_PATH.read_text(encoding="utf-8")) + assert isinstance(_state.get("next_id"), int) and isinstance(_state.get("items"), list) + except Exception: + _state = {"next_id": 1, "items": []} + return _state + + +def _save(state: dict) -> None: + try: + tmp = STORE_PATH.with_suffix(".tmp") + tmp.write_text(json.dumps(state, ensure_ascii=False), encoding="utf-8") + tmp.replace(STORE_PATH) + except OSError: + # Briefkasten darf den Absender nie blockieren — dann eben nur in-memory. + log.warning("announce: Store %s nicht schreibbar", STORE_PATH, exc_info=True) + + +def add(text: str, subject: str = "", source: str = "", priority: str = "normal") -> dict: + """Eintrag anhängen. priority: 'normal' (sprechen) | 'silent' (nur Verlauf/Panel).""" + text = (text or "").strip() + if not text: + raise ValueError("Leere Meldung.") + with _lock: + state = _load() + item = { + "id": state["next_id"], + "ts": time.time(), + "subject": (subject or "").strip()[:120], + "text": text[:4000], + "source": (source or "").strip()[:60], + "priority": priority if priority in ("normal", "silent") else "normal", + } + state["next_id"] += 1 + state["items"].append(item) + del state["items"][:-MAX_ITEMS] + _save(state) + log.info("announce #%s [%s] %s: %.80s", item["id"], item["source"] or "-", item["subject"] or "-", text) + return item + + +def list_after(after: int | None, limit: int = 20) -> dict: + """Einträge NACH Cursor `after` (aufsteigend). Ohne Cursor nur den aktuellen + Stand liefern (latest) — so initialisiert Lucy ihren Cursor, ohne Altes nachzuplappern.""" + with _lock: + state = _load() + latest = state["next_id"] - 1 + items = [] if after is None else [i for i in state["items"] if i["id"] > after][: max(1, min(limit, 100))] + return {"latest": latest, "items": items} diff --git a/backend/services/sentry.py b/backend/services/sentry.py new file mode 100644 index 0000000..256730b --- /dev/null +++ b/backend/services/sentry.py @@ -0,0 +1,152 @@ +""" +Health-Wächter der Box (Lucy-Proaktivität, Faden A3). + +Prüft periodisch die Kern-Dienste (Engine, Agent-Hirn, Hermes-Gateway, Mem0, +Voice-Sidecar, Platte) und meldet ZUSTANDSWECHSEL in den Melde-Briefkasten +(services/announce.py → Lucy spricht es) und via notify.sh (Telegram). + +Flankenerkennung statt Dauerfeuer: Alarm erst nach FAIL_AFTER Fehl-Ticks in +Folge (überlebt Neustarts/Update-Fenster), Entwarnung beim ersten grünen Tick +nach einem Alarm. Hält ein Problem an, wird frühestens nach REMIND_S erinnert. + +Braucht KEIN sudo, keine neuen Dienste — läuft als asyncio-Task im MC2-Backend +(wie der Re-Warm-Wächter). Abschaltbar via MC_SENTRY_ENABLED=0. +""" + +import asyncio +import logging +import os +import shutil +import subprocess +import time + +import httpx +import psutil + +from config import HERMES_API_URL, LLAMA_SWAP_URL, MEM0_SERVICE_URL, MODELS_DIR, VOICE_SERVICE_URL +from services import announce, llamaswap + +log = logging.getLogger(__name__) + +ENABLED = os.environ.get("MC_SENTRY_ENABLED", "1") != "0" +INTERVAL = int(os.environ.get("MC_SENTRY_INTERVAL", "120")) # Sekunden zwischen Ticks +START_DELAY = int(os.environ.get("MC_SENTRY_START_DELAY", "90")) # Dienste nach Boot setzen lassen +FAIL_AFTER = int(os.environ.get("MC_SENTRY_FAIL_AFTER", "3")) # Fehl-Ticks bis Alarm (3×120s = 6 min) +REMIND_S = int(os.environ.get("MC_SENTRY_REMIND_S", "21600")) # Erinnerung bei Dauerproblem: 6 h +DISK_ALARM_PCT = float(os.environ.get("MC_SENTRY_DISK_PCT", "90")) +NOTIFY_SH = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), + "deploy", "notify.sh") + + +def _reach(url: str, path: str = "/health") -> bool: + try: + with httpx.Client(timeout=5.0) as c: + return c.get(f"{url}{path}").status_code < 500 + except Exception: + return False + + +def _check_engine() -> bool: + return llamaswap.engine_reachable() + + +def _check_brain() -> bool: + """Hirn tot? Verdrängung durch ein ANDERES laufendes Modell (IDE-Last) ist NORMAL — + Alarm nur, wenn gar nichts läuft und das Hirn trotz Re-Warm-Wächter kalt bleibt.""" + st = llamaswap.brain_status() + if st.get("ready"): + return True + return bool(llamaswap.get_running_models()) # anderes Modell aktiv → Verdrängung, kein Defekt + + +def _check_disk() -> bool: + try: + return psutil.disk_usage(str(MODELS_DIR) if MODELS_DIR.exists() else os.getcwd()).percent < DISK_ALARM_PCT + except Exception: + return True # kein Messwert ≠ Alarm + + +# name → (Checker, Alarm-Text, Entwarnungs-Text) — Texte sind Lucy-sprechbar (kurz, Alltagssprache). +CHECKS: dict[str, tuple] = { + "engine": (_check_engine, + "Die Modell-Engine antwortet nicht mehr. Ohne sie laufen keine KI-Modelle.", + "Die Modell-Engine ist wieder da."), + "brain": (_check_brain, + "Mein Gehirn lädt nicht — ich kann gerade nicht richtig denken. Ein Neustart der Engine könnte helfen.", + "Mein Gehirn ist wieder geladen. Alles klar bei mir."), + "hermes": (lambda: _reach(HERMES_API_URL), + "Der Agent-Dienst ist ausgefallen — Telegram und meine Tools gehen gerade nicht.", + "Der Agent-Dienst läuft wieder."), + "mem0": (lambda: _reach(MEM0_SERVICE_URL), + "Mein Gedächtnis-Dienst ist ausgefallen — ich merke mir vorübergehend nichts Neues.", + "Mein Gedächtnis ist wieder online."), + "voice": (lambda: _reach(VOICE_SERVICE_URL), + "Der Hör-Dienst auf der Box ist ausgefallen — Spracheingabe könnte haken.", + "Der Hör-Dienst läuft wieder."), + "disk": (_check_disk, + f"Die Platte der Box ist zu über {DISK_ALARM_PCT:.0f} Prozent voll. Es wird eng für Modelle und Backups.", + "Die Platte hat wieder genug Luft."), +} + + +class _Watch: + __slots__ = ("fails", "alerted", "alert_ts") + + def __init__(self) -> None: + self.fails = 0 # Fehl-Ticks in Folge + self.alerted = False # Alarm ist raus, Entwarnung steht aus + self.alert_ts = 0.0 # Zeitpunkt des letzten Alarms (für REMIND_S) + + +_watches: dict[str, _Watch] = {name: _Watch() for name in CHECKS} + + +def _notify_telegram(subject: str, text: str) -> None: + """Best-effort auch auf Telegram (User ist evtl. nicht am PC). notify.sh spiegelt + seinerseits in den Briefkasten — als Quelle 'sentry' markierte Einträge legt der + Wächter aber schon selbst ab, darum hier der Direktweg NUR für Telegram.""" + if not (os.name == "posix" and shutil.which("bash")): + return # Dev auf Windows: kein hermes/notify + try: + subprocess.run(["bash", NOTIFY_SH, "-s", subject, text + " (Diese Meldung kam auch an Lucy.)"], + timeout=30, capture_output=True, env={**os.environ, "MC_NOTIFY_NO_ANNOUNCE": "1"}) + except Exception: + log.warning("sentry: notify.sh fehlgeschlagen", exc_info=True) + + +def _tick() -> None: + now = time.time() + for name, (check, fail_msg, ok_msg) in CHECKS.items(): + w = _watches[name] + try: + ok = bool(check()) + except Exception: + ok = False + if ok: + w.fails = 0 + if w.alerted: + w.alerted = False + announce.add(ok_msg, subject="[Box wieder ok]", source="sentry") + _notify_telegram("[Box wieder ok]", ok_msg) + continue + w.fails += 1 + due = (not w.alerted and w.fails >= FAIL_AFTER) or (w.alerted and now - w.alert_ts >= REMIND_S) + if due: + prefix = "" if not w.alerted else "Immer noch: " + w.alerted = True + w.alert_ts = now + announce.add(prefix + fail_msg, subject="[Box-Problem]", source="sentry") + _notify_telegram("[Box-Problem]", prefix + fail_msg) + log.warning("sentry: %s ALARM (%s Fehl-Ticks)", name, w.fails) + + +async def sentry_loop() -> None: + """Endlos-Schleife (Hintergrund-Task im MC2-Lifespan).""" + await asyncio.sleep(START_DELAY) + log.info("sentry: Health-Wächter aktiv (Intervall %ss, Alarm nach %s Fehl-Ticks)", INTERVAL, FAIL_AFTER) + while True: + try: + await asyncio.to_thread(_tick) + except Exception: + log.debug("sentry: Tick fehlgeschlagen", exc_info=True) + await asyncio.sleep(INTERVAL) diff --git a/deploy/notify.sh b/deploy/notify.sh index 1ec8aa9..c2d4e39 100644 --- a/deploy/notify.sh +++ b/deploy/notify.sh @@ -28,6 +28,19 @@ fi TS="$(date '+%Y-%m-%d %H:%M:%S')" +# Lucy-Briefkasten (Proaktivität A3): Meldung zusätzlich an MC2 spiegeln — die Desktop-Lucy +# spricht sie dann von sich aus. Best-effort (darf den Telegram-Weg nie aufhalten). +# MC_NOTIFY_NO_ANNOUNCE=1 unterdrückt das (der Health-Wächter legt seine Einträge selbst ab). +if [ "${MC_NOTIFY_NO_ANNOUNCE:-0}" != "1" ]; then + curl -sf -m 3 -X POST "${MC_ANNOUNCE_URL:-http://127.0.0.1:9001/api/voice/announce}" \ + -H 'Content-Type: application/json' \ + --data "$(python3 - "$SUBJECT" "$MSG" <<'PY' +import json, sys +print(json.dumps({"text": sys.argv[2], "subject": sys.argv[1], "source": "notify"})) +PY +)" >/dev/null 2>&1 || true +fi + # Primärweg: Telegram via hermes send (Login-Shell-PATH, falls aus Timer/cron aufgerufen). if [ -n "$SUBJECT" ]; then SEND_OUT="$(bash -lc 'hermes send --to telegram --subject "$1" -- "$2"' _ "$SUBJECT" "$MSG" 2>&1)"