Lucy-Proaktivität (A3): Melde-Briefkasten + Health-Wächter
- 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 <noreply@anthropic.com>
This commit is contained in:
+9
-5
@@ -19,7 +19,7 @@ from starlette.requests import Request
|
|||||||
|
|
||||||
from config import FRONTEND_DIST, VERSION
|
from config import FRONTEND_DIST, VERSION
|
||||||
from routers import agent, connect, gateway_proxy, health, maintenance, memory, models, routing, system, voice
|
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
|
# Zentrales Logging — Level via MC_LOG_LEVEL (INFO default). Eine Konfiguration
|
||||||
# für alle Module (logging.getLogger(__name__)).
|
# für alle Module (logging.getLogger(__name__)).
|
||||||
@@ -31,14 +31,18 @@ log = logging.getLogger(__name__)
|
|||||||
|
|
||||||
@asynccontextmanager
|
@asynccontextmanager
|
||||||
async def lifespan(app: FastAPI):
|
async def lifespan(app: FastAPI):
|
||||||
"""Hintergrund-Tasks an den App-Lebenszyklus binden: Re-Warm-Wächter fürs Agent-Hirn."""
|
"""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
|
+ Health-Wächter (meldet Ausfälle/Erholung in den Lucy-Briefkasten und auf Telegram)."""
|
||||||
if task:
|
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)
|
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:
|
try:
|
||||||
yield
|
yield
|
||||||
finally:
|
finally:
|
||||||
if task:
|
for task in tasks:
|
||||||
task.cancel()
|
task.cancel()
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -19,6 +19,7 @@ from fastapi.responses import Response, StreamingResponse
|
|||||||
from pydantic import BaseModel
|
from pydantic import BaseModel
|
||||||
|
|
||||||
from config import HERMES_API_KEY, HERMES_API_MODEL, HERMES_API_URL, LLAMA_SWAP_URL, VOICE_SERVICE_URL
|
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)
|
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).
|
# 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")
|
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")
|
@router.get("/voice/metrics")
|
||||||
def voice_metrics() -> dict:
|
def voice_metrics() -> dict:
|
||||||
"""Per-Stage-Latenz (STT/Vision/Chat-TTFB/TTS) — rollende Statistik, macht die Voice-Pipeline
|
"""Per-Stage-Latenz (STT/Vision/Chat-TTFB/TTS) — rollende Statistik, macht die Voice-Pipeline
|
||||||
|
|||||||
@@ -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}
|
||||||
@@ -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)
|
||||||
@@ -28,6 +28,19 @@ fi
|
|||||||
|
|
||||||
TS="$(date '+%Y-%m-%d %H:%M:%S')"
|
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).
|
# Primärweg: Telegram via hermes send (Login-Shell-PATH, falls aus Timer/cron aufgerufen).
|
||||||
if [ -n "$SUBJECT" ]; then
|
if [ -n "$SUBJECT" ]; then
|
||||||
SEND_OUT="$(bash -lc 'hermes send --to telegram --subject "$1" -- "$2"' _ "$SUBJECT" "$MSG" 2>&1)"
|
SEND_OUT="$(bash -lc 'hermes send --to telegram --subject "$1" -- "$2"' _ "$SUBJECT" "$MSG" 2>&1)"
|
||||||
|
|||||||
Reference in New Issue
Block a user