Files
mission-control-v2/backend/services/announce.py
T
Hitonabi a9c0239f8a UMBAU v3 P2: Steward — Waechter-Loops als eigener Prozess mc2-steward
Re-Warm, Health-Sentry und Mem0-Dedupe laufen als eigener Dienst
(backend/steward.py, Restart=always) statt im Steuerpult-Lifespan:
MC2-Neustarts nehmen den Waechtern nicht mehr Timing/Flanken-Gedaechtnis,
und der Sentry ueberwacht erstmals MC2 SELBST + mc2-gateway (Telegram
funktioniert auch bei totem Steuerpult; Briefkasten-Abgabe per HTTP via
MC_ANNOUNCE_HTTP, Store bleibt exklusiv beim MC2-Prozess). Warm-Nudge
nach Config-Aenderung via mtime-Watch (5 s) statt In-Process-Signal.
Reiner Konfig-Split: MC2-Unit setzt die drei ENABLED-Schalter auf 0,
Zeilen entfernen = Rollback. reminders_loop bleibt bewusst in MC2
(teilt Datei+CRUD mit /api/reminders, Zwei-Schreiber-Risiko).

Stellt ausserdem den beim Karten-Neuaufbau (ccc9a25) verlorenen
stack-postcheck-Block fuer mc2-gateway wieder her (+ Steward-Check 4c).

Baut auf feature/von-allein-und-gateway-p1 auf; Annahme schliesst P1 ein.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-15 14:15:36 +02:00

132 lines
5.4 KiB
Python

"""
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 shutil
import subprocess
import threading
import time
from pathlib import Path
import httpx
from config import MODELS_DIR
log = logging.getLogger(__name__)
NOTIFY_SH = str(Path(__file__).resolve().parent.parent.parent / "deploy" / "notify.sh")
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"))
# Steward-Auszug (UMBAU v3 P2): Läuft der Absender AUSSERHALB des MC2-Prozesses
# (mc2-steward), darf er den Store NICHT direkt schreiben — Datei + In-Memory-Cache
# gehören exklusiv dem Steuerpult. Gesetzt (Unit-Env, z. B. http://127.0.0.1:9001)
# liefert add() die Meldung stattdessen per HTTP am /api/voice/announce-Endpunkt ab.
ANNOUNCE_HTTP = os.environ.get("MC_ANNOUNCE_HTTP", "").rstrip("/")
_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.")
if ANNOUNCE_HTTP:
# Fremdprozess-Modus (mc2-steward): per HTTP beim Steuerpult abliefern.
payload = {"text": text[:4000], "subject": (subject or "").strip()[:120],
"source": (source or "").strip()[:60], "priority": priority}
try:
with httpx.Client(timeout=5.0) as c:
r = c.post(f"{ANNOUNCE_HTTP}/api/voice/announce", json=payload)
if r.status_code == 200:
return (r.json() or {}).get("item") or payload
except Exception:
pass
# MC2 down/Endpunkt weg: Alarm ist nicht verloren — der Telegram-Direktweg des
# Aufrufers (notify.sh) läuft separat; hier bleibt nur der Log-Nachweis.
log.warning("announce: HTTP-Abgabe an %s fehlgeschlagen — nur im Log: [%s] %.120s",
ANNOUNCE_HTTP, subject or "-", text)
return payload
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 notify_telegram(subject: str, text: str) -> None:
"""Best-effort auch auf Telegram (User ist evtl. nicht am PC). MC_NOTIFY_NO_ANNOUNCE=1
verhindert, dass notify.sh die Meldung ZURÜCK in den Briefkasten spiegelt — der Absender
(Wächter/Erinnerung) hat sie dort schon selbst abgelegt. Auf Windows (Dev) ein No-op."""
if not (os.name == "posix" and shutil.which("bash")):
return
try:
subprocess.run(["bash", NOTIFY_SH, "-s", subject, text],
timeout=30, capture_output=True, env={**os.environ, "MC_NOTIFY_NO_ANNOUNCE": "1"})
except Exception:
log.warning("notify_telegram: notify.sh fehlgeschlagen", exc_info=True)
def list_recent(limit: int = 150) -> list[dict]:
"""Die jüngsten Einträge (neueste zuerst) — Datenquelle der Chronik: alles, was die Box
dem Commander je aktiv gemeldet hat (Health-Wächter, Updates, Radar, Träume, Alarme)."""
with _lock:
state = _load()
return list(reversed(state["items"][-max(1, min(limit, MAX_ITEMS)):]))
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}