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>
This commit is contained in:
Hitonabi
2026-07-15 14:15:36 +02:00
parent ccc9a25a25
commit a9c0239f8a
8 changed files with 204 additions and 5 deletions
+24
View File
@@ -19,6 +19,8 @@ import threading
import time
from pathlib import Path
import httpx
from config import MODELS_DIR
log = logging.getLogger(__name__)
@@ -28,6 +30,12 @@ NOTIFY_SH = str(Path(__file__).resolve().parent.parent.parent / "deploy" / "noti
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": [...]}
@@ -58,6 +66,22 @@ def add(text: str, subject: str = "", source: str = "", priority: str = "normal"
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 = {
+14
View File
@@ -85,6 +85,20 @@ CHECKS: dict[str, tuple] = {
}
# Steward-Modus (UMBAU v3 P2): Läuft der Wächter als EIGENER Prozess (mc2-steward),
# beobachtet er zusätzlich das Steuerpult selbst und den mc2-gateway — genau die zwei
# Ausfälle, die der alte In-Process-Wächter prinzipbedingt nie melden konnte (er starb mit).
if os.environ.get("MC_SENTRY_WATCH_MC2", "") == "1":
_MC2_URL = os.environ.get("MC_SENTRY_MC2_URL", "http://127.0.0.1:9001")
_GW_URL = os.environ.get("MC_SENTRY_GATEWAY_URL", "http://127.0.0.1:9010")
CHECKS["mc2"] = (lambda: _reach(_MC2_URL, "/api/health"),
"Das Steuerpult ist ausgefallen — Dashboard, Briefkasten und Auftragsbuch gehen gerade nicht.",
"Das Steuerpult ist wieder da.")
CHECKS["gateway"] = (lambda: _reach(_GW_URL, "/gw/health"),
"Der Modell-Gateway ist ausgefallen — meine Denk-Anfragen und die der Worker hängen gerade.",
"Der Modell-Gateway läuft wieder.")
class _Watch:
__slots__ = ("fails", "alerted", "alert_ts")
+80
View File
@@ -0,0 +1,80 @@
"""
MC2-Steward — die Wächter-Loops als EIGENER Prozess (UMBAU v3, P2).
Bisher hingen Re-Warm-Wächter, Health-Wächter (sentry) und Mem0-Auto-Dedupe am
Lebenszyklus des Steuerpult-Webservers (app.py-Lifespan): jeder MC2-Neustart riss
den Wächtern Timing und Flanken-Gedächtnis weg — und ein TOTES Steuerpult konnte
sich prinzipbedingt nicht selbst melden. Hier laufen DIESELBEN Loops (unveränderte
Module) als eigener Mini-Dienst (mc2-steward.service, Restart=always):
• warmer.rewarm_loop — Warm-Set nachladen (+ Config-Watch ersetzt den
In-Process-Nudge aus llamaswap.write_config)
• sentry.sentry_loop — Health-Flanken → Briefkasten (HTTP an MC2) + Telegram;
beobachtet im Steward-Modus AUCH MC2 selbst + mc2-gateway
• memory.auto_dedupe_loop — Gedächtnis-Dubletten (HTTP an den Mem0-Sidecar)
BEWUSST NICHT hier: reminders_loop — der teilt sich Datei UND CRUD-Pfade mit dem
/api/reminders-Router (Zwei-Schreiber-Risiko auf mc2-reminders.json); er bleibt im
Steuerpult. Der Briefkasten-Store gehört weiter EXKLUSIV dem MC2-Prozess — dieser
Prozess liefert Meldungen per HTTP ab (announce.py, MC_ANNOUNCE_HTTP).
Die Loop-Schalter (MC_REWARM_ENABLED / MC_SENTRY_ENABLED / MC_MEM_DEDUPE_ENABLED)
stehen in der MC2-Unit auf 0 und hier auf Default 1 — reiner Konfig-Split, kein
Verhaltens-Code im Steuerpult angefasst. Zeilen dort entfernen = Rollback.
"""
import asyncio
import logging
import os
from config import CONFIG_PATH
from services import memory as memory_svc
from services import sentry, warmer
logging.basicConfig(
level=os.environ.get("MC_LOG_LEVEL", "INFO").upper(),
format="%(asctime)s %(levelname)-7s %(name)s: %(message)s",
)
log = logging.getLogger("steward")
# Wie oft nach einer llama-swap-Config-Änderung geschaut wird (mtime-Watch). Ersetzt
# warmer.nudge() aus dem MC2-Prozess: ein -watch-config-Reload verwirft das Warm-Set,
# wir sehen die Änderung binnen Sekunden statt erst beim nächsten 90s-Tick.
CONFIG_WATCH_S = int(os.environ.get("MC_STEWARD_CONFIG_WATCH_S", "5"))
async def config_watch() -> None:
last: float | None = None
while True:
try:
mtime = CONFIG_PATH.stat().st_mtime
except OSError:
mtime = None
if last is not None and mtime is not None and mtime != last:
log.info("steward: %s geändert → Warm-Set-Nudge", CONFIG_PATH)
warmer.nudge()
last = mtime
await asyncio.sleep(CONFIG_WATCH_S)
async def main() -> None:
tasks: list[asyncio.Task] = []
if warmer.ENABLED:
tasks.append(asyncio.create_task(warmer.rewarm_loop()))
tasks.append(asyncio.create_task(config_watch()))
log.info("Re-Warm-Wächter aktiv (Intervall %ss, Config-Watch %ss)",
warmer.INTERVAL, CONFIG_WATCH_S)
if sentry.ENABLED:
tasks.append(asyncio.create_task(sentry.sentry_loop()))
if memory_svc.AUTO_DEDUPE_ENABLED:
tasks.append(asyncio.create_task(memory_svc.auto_dedupe_loop()))
log.info("Mem0-Auto-Dedupe aktiv (alle %ss, Schwelle %s)",
memory_svc.AUTO_DEDUPE_INTERVAL, memory_svc.AUTO_DEDUPE_THRESHOLD)
if not tasks:
log.warning("steward: alle Loops per Env deaktiviert — nichts zu tun, Ende.")
return
await asyncio.gather(*tasks)
if __name__ == "__main__":
asyncio.run(main())