umbau(boxwart): Backend auf Box-Wart umgestellt – Waechter, Modell-Nutzung, Zeitplan
MC2 wird Updater, Waechter und Modell-Radar (Konzept „MC2 als Box-Wart“, 23.09.2026). - Neuer Waechter (services/waechter.py) loest sentry.py ab: Dienste, Timer-Laeufe, Hermes-Jobs samt Werkzeugfehlern, Kern-HTTP-Proben, Platte. Abgestuerzte Dienste startet er selbst neu (max. 2/h), rote Hinweise gehen an Telegram und Lucy. Laeuft im mc2-steward; waehrend eines Updates haelt er still. - Neue Schnittstellen (routers/boxwart.py): /api/start, /api/hinweise (+ Aktionen), /api/modelle/nutzung, /api/zeitplan. - Modell-Nutzung aus dem llama-swap-Journal (wer fragt wie oft, 24 h je Stunde). - Entfernt: Ideen, Wissen, Chronik, Skills, Verbinden, Konsolen-Proxy, /api/events; Lucys Werkzeug idee_notieren; box_status nennt jetzt die offenen Hinweise. - Behoben: projekte-sync ueberspringt leere Gitea-Repos (lief seit 07.09. stuendlich rot); Motor-Version kam aus dem verwaisten /opt/llamacpp statt /opt/llamacpp-vulkan. - Erste Backend-Tests (8) fuer Waechter und Modell-Nutzung. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
3669a6e7dc
commit
12dadfe6ef
@@ -0,0 +1,535 @@
|
||||
"""
|
||||
Wächter der Box (Box-Wart, Umbau 09/2026) — löst den alten Health-Wächter (sentry.py) ab.
|
||||
|
||||
Der alte Wächter kannte nur sieben Dienste per HTTP. Dass projekte-sync seit dem 07.09.
|
||||
stündlich scheiterte und der Nachrichten-Job jeden Morgen Werkzeugfehler warf, hat er nie
|
||||
gesehen. Dieser Wächter schaut deshalb breiter hin:
|
||||
|
||||
• Dienste systemd-Zustand der Kern-Dienste (System + User)
|
||||
• Timer-Läufe Einmal-Dienste hinter Timern (projekte-sync, Sicherung)
|
||||
• Hermes-Jobs Status der Cron-Jobs + Werkzeugfehler im letzten Lauf (errors.log)
|
||||
• Kern Engine, Hirn, Hermes, Hör-Dienst, MC2, Gateway antworten per HTTP
|
||||
• Platte Füllstand des Modell-Laufwerks
|
||||
|
||||
Jeder Befund wird zum Hinweis mit Stufe (rot = jetzt, gelb = bei Gelegenheit), Beginn und
|
||||
Knöpfen. Einfaches behebt er selbst (User-Entscheid 23.09.): Ein ABGESTÜRZTER Dienst
|
||||
(ActiveState=failed) wird neu gestartet, höchstens AUTO_MAX_PRO_STUNDE-mal pro Stunde. Ein
|
||||
bewusst gestoppter Dienst (inactive) bleibt aus und wird nur gemeldet. Rote Hinweise gehen
|
||||
zusätzlich an Telegram und in Lucys Briefkasten.
|
||||
|
||||
Läuft im mc2-steward (eigener Prozess, übersteht MC2-Neustarts) und ist dort der EINZIGE
|
||||
Schreiber von STORE_PATH. MC2 liest den Stand nur (lese_stand()) und führt Knopf-Aktionen
|
||||
selbst aus (fuehre_aktion_aus()); der nächste Takt sieht das Ergebnis.
|
||||
|
||||
Während eines Updates (autoupdate.sh, update-engine.sh, update-swap.sh, hermes update)
|
||||
hält er still: keine neuen Dienst-Hinweise, keine Selbstreparatur — Neustarts gehören dort
|
||||
zum Ablauf.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
import sqlite3
|
||||
import subprocess
|
||||
import threading
|
||||
import time
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
|
||||
import httpx
|
||||
import psutil
|
||||
from config import HERMES_API_KEY, HERMES_API_URL, HERMES_HOME, MODELS_DIR, VOICE_SERVICE_URL
|
||||
|
||||
from services import announce, llamaswap, maintenance
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
# MC_SENTRY_ENABLED bleibt als Rückfall-Schalter gültig: Die bestehenden Units setzen ihn.
|
||||
ENABLED = os.environ.get("MC_WAECHTER_ENABLED", os.environ.get("MC_SENTRY_ENABLED", "1")) != "0"
|
||||
INTERVAL = int(os.environ.get("MC_WAECHTER_INTERVAL", "60")) # Sekunden zwischen Takten
|
||||
START_DELAY = int(os.environ.get("MC_WAECHTER_START_DELAY", "60")) # Dienste nach Boot setzen lassen
|
||||
FAIL_AFTER = int(os.environ.get("MC_WAECHTER_FAIL_AFTER", "2")) # Takte bis zum Hinweis
|
||||
REMIND_S = int(os.environ.get("MC_WAECHTER_REMIND_S", "21600")) # Telegram-Erinnerung: 6 h
|
||||
AUTO_MAX_PRO_STUNDE = int(os.environ.get("MC_WAECHTER_AUTO_MAX", "2"))
|
||||
DISK_ROT_PCT = float(os.environ.get("MC_WAECHTER_DISK_ROT", "90"))
|
||||
DISK_GELB_PCT = float(os.environ.get("MC_WAECHTER_DISK_GELB", "80"))
|
||||
STORE_PATH = Path(os.environ.get("MC_WAECHTER_STORE", str(MODELS_DIR / "mc2-waechter.json")))
|
||||
VERLAUF_MAX = 200
|
||||
|
||||
# Im Steward-Modus beobachtet er auch MC2 selbst und den Gateway (wie der alte Wächter).
|
||||
WATCH_MC2 = os.environ.get("MC_SENTRY_WATCH_MC2", os.environ.get("MC_WAECHTER_WATCH_MC2", "")) == "1"
|
||||
MC2_URL = os.environ.get("MC_SENTRY_MC2_URL", "http://127.0.0.1:9001")
|
||||
GATEWAY_URL = os.environ.get("MC_SENTRY_GATEWAY_URL", "http://127.0.0.1:9010")
|
||||
|
||||
# Dienst → (System-Dienst?, wichtig?, Anzeigename). Wichtig = rot, sonst gelb.
|
||||
DIENSTE: dict[str, tuple[bool, bool, str]] = {
|
||||
"llama-swap": (True, True, "Motor (llama-swap)"),
|
||||
"hermes-gateway": (False, True, "Hermes"),
|
||||
"mc2-gateway": (False, True, "Modell-Gateway"),
|
||||
"mission-control-2": (False, True, "MC2"),
|
||||
"voice-service": (False, False, "Hör-Dienst"),
|
||||
"lucy-stimme": (False, False, "Lucys Stimme"),
|
||||
"hermes-builtin-ui": (False, False, "Hermes-Dashboard"),
|
||||
}
|
||||
|
||||
# Einmal-Dienste hinter Timern: Name → (Anzeigename, wichtig?)
|
||||
TIMER_DIENSTE: dict[str, tuple[str, bool]] = {
|
||||
"projekte-sync": ("Projekte-Abgleich", False),
|
||||
"mc2-backup": ("Sicherung", True),
|
||||
}
|
||||
|
||||
HERMES_JOBS_PATH = HERMES_HOME / "cron" / "jobs.json"
|
||||
HERMES_EXEC_DB = HERMES_HOME / "cron" / "executions.db"
|
||||
HERMES_ERRORS_LOG = HERMES_HOME / "logs" / "errors.log"
|
||||
|
||||
# Prozesse, an denen ein laufendes Update zu erkennen ist.
|
||||
_UPDATE_PROZESSE = re.compile(r"autoupdate\.sh|update-engine\.sh|update-swap\.sh|hermes(\s+\S+)*\s+update\b")
|
||||
|
||||
|
||||
@dataclass
|
||||
class Befund:
|
||||
"""Ein aktuell festgestelltes Problem. Wird nach FAIL_AFTER Takten zum Hinweis."""
|
||||
id: str
|
||||
stufe: str # "rot" | "gelb"
|
||||
titel: str
|
||||
text: str
|
||||
quelle: str
|
||||
aktionen: list[dict] = field(default_factory=list)
|
||||
auto: str | None = None # Name der Selbstreparatur, falls erlaubt
|
||||
sofort: bool = False # ohne FAIL_AFTER-Wartezeit (dauerhafte Befunde)
|
||||
|
||||
|
||||
def _aktion(aid: str, label: str, **extra: str) -> dict:
|
||||
return {"id": aid, "label": label, **extra}
|
||||
|
||||
|
||||
# --- Rohdaten -------------------------------------------------------------------
|
||||
|
||||
def _systemctl_show(name: str, system: bool) -> dict[str, str]:
|
||||
"""Zustand einer Unit. Auf Windows (Dev) oder ohne systemd: leeres Dict (harmlos)."""
|
||||
cmd = ["systemctl"] if system else ["systemctl", "--user"]
|
||||
cmd += ["show", name, "-p",
|
||||
"LoadState,ActiveState,SubState,Result,ExecMainStatus,InactiveExitTimestamp"]
|
||||
try:
|
||||
out = subprocess.run(cmd, capture_output=True, text=True, timeout=10).stdout
|
||||
except Exception:
|
||||
return {}
|
||||
return dict(line.split("=", 1) for line in out.splitlines() if "=" in line)
|
||||
|
||||
|
||||
def _journal_fehlerzeile(name: str, system: bool = False) -> str:
|
||||
"""Die aussagekräftigste Fehlerzeile aus den letzten Journal-Zeilen einer Unit."""
|
||||
cmd = ["journalctl"] + ([] if system else ["--user"]) + ["-u", name, "-n", "40", "-o", "cat", "--no-pager"]
|
||||
try:
|
||||
zeilen = subprocess.run(cmd, capture_output=True, text=True, timeout=10).stdout.splitlines()
|
||||
except Exception:
|
||||
return ""
|
||||
return waehle_fehlerzeile(zeilen)
|
||||
|
||||
|
||||
def waehle_fehlerzeile(zeilen: list[str]) -> str:
|
||||
"""Letzte Zeile, die nach Fehler aussieht (projekte-sync markiert Fehler mit „!“)."""
|
||||
muster = re.compile(r"(^\s*!|fehlgeschlagen|ABBRUCH|error|failed|fatal)", re.IGNORECASE)
|
||||
for zeile in reversed(zeilen):
|
||||
if muster.search(zeile):
|
||||
# Zeitstempel "2026-09-23 21:02:46 " des Skript-Logs abschneiden
|
||||
return re.sub(r"^\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}\s+", "", zeile).strip()[:300]
|
||||
return ""
|
||||
|
||||
|
||||
def _reach(url: str, path: str, headers: dict[str, str] | None = None) -> bool:
|
||||
"""Antwortet der Dienst? < 500 genügt — ein 401 beweist, dass jemand zuhört."""
|
||||
try:
|
||||
with httpx.Client(timeout=5.0) as c:
|
||||
return c.get(f"{url}{path}", headers=headers or {}).status_code < 500
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
def _update_laeuft() -> bool:
|
||||
"""Läuft gerade ein Update? Dann gehören Neustarts zum Ablauf und sind kein Befund."""
|
||||
try:
|
||||
for p in psutil.process_iter(["cmdline"]):
|
||||
cmd = " ".join(p.info.get("cmdline") or [])
|
||||
if cmd and _UPDATE_PROZESSE.search(cmd):
|
||||
return True
|
||||
except Exception:
|
||||
return False
|
||||
return False
|
||||
|
||||
|
||||
# --- Prüfungen ------------------------------------------------------------------
|
||||
|
||||
def pruefe_dienste() -> list[Befund]:
|
||||
befunde: list[Befund] = []
|
||||
for name, (system, wichtig, anzeige) in DIENSTE.items():
|
||||
z = _systemctl_show(name, system)
|
||||
if not z or z.get("LoadState") in ("not-found", ""):
|
||||
continue # nicht installiert oder kein systemd (Dev)
|
||||
zustand = z.get("ActiveState", "")
|
||||
if zustand in ("active", "activating", "reloading", "deactivating"):
|
||||
continue
|
||||
stufe = "rot" if wichtig else "gelb"
|
||||
aktionen = [_aktion("neustart", "Neu starten", dienst=name), _aktion("protokoll", "Protokoll", dienst=name)]
|
||||
if zustand == "failed":
|
||||
befunde.append(Befund(
|
||||
id=f"dienst:{name}", stufe=stufe, titel=f"{anzeige} ist abgestürzt",
|
||||
text=_journal_fehlerzeile(name, system) or f"Die Unit {name} steht auf „failed“.",
|
||||
quelle=name, aktionen=aktionen, auto="neustart"))
|
||||
else: # inactive: vermutlich bewusst gestoppt → nur melden, nicht eigenmächtig starten
|
||||
befunde.append(Befund(
|
||||
id=f"dienst:{name}", stufe=stufe, titel=f"{anzeige} ist gestoppt",
|
||||
text=f"Die Unit {name} läuft nicht. Sie wurde vermutlich angehalten.",
|
||||
quelle=name, aktionen=[_aktion("neustart", "Starten", dienst=name), aktionen[1]]))
|
||||
return befunde
|
||||
|
||||
|
||||
def pruefe_timer_dienste() -> list[Befund]:
|
||||
befunde: list[Befund] = []
|
||||
for name, (anzeige, wichtig) in TIMER_DIENSTE.items():
|
||||
z = _systemctl_show(name, False)
|
||||
if not z or z.get("LoadState") in ("not-found", ""):
|
||||
continue
|
||||
if z.get("ActiveState") != "failed" and z.get("Result", "success") == "success":
|
||||
continue
|
||||
grund = _journal_fehlerzeile(name) or f"Letzter Lauf endete mit Code {z.get('ExecMainStatus', '?')}."
|
||||
befunde.append(Befund(
|
||||
id=f"timer:{name}", stufe="rot" if wichtig else "gelb",
|
||||
titel=f"{anzeige} scheitert beim letzten Lauf", text=grund, quelle=name,
|
||||
aktionen=[_aktion("protokoll", "Protokoll", dienst=name),
|
||||
_aktion("neustart", "Jetzt erneut laufen lassen", dienst=name)],
|
||||
sofort=True))
|
||||
return befunde
|
||||
|
||||
|
||||
def _hermes_jobs() -> list[dict]:
|
||||
try:
|
||||
daten = json.loads(HERMES_JOBS_PATH.read_text(encoding="utf-8"))
|
||||
except (OSError, ValueError):
|
||||
return []
|
||||
jobs = daten.get("jobs", daten) if isinstance(daten, dict) else daten
|
||||
if isinstance(jobs, dict):
|
||||
jobs = list(jobs.values())
|
||||
return [j for j in jobs if isinstance(j, dict)]
|
||||
|
||||
|
||||
def _letzter_lauf(job_id: str) -> dict | None:
|
||||
"""Letzter Lauf eines Jobs aus executions.db (nur lesend geöffnet)."""
|
||||
if not HERMES_EXEC_DB.exists():
|
||||
return None
|
||||
try:
|
||||
con = sqlite3.connect(f"file:{HERMES_EXEC_DB}?mode=ro", uri=True, timeout=2)
|
||||
try:
|
||||
row = con.execute(
|
||||
"select status, started_at, finished_at, error from executions "
|
||||
"where job_id=? order by started_at desc limit 1", (job_id,)).fetchone()
|
||||
finally:
|
||||
con.close()
|
||||
except sqlite3.Error:
|
||||
return None
|
||||
if not row:
|
||||
return None
|
||||
return {"status": row[0], "started_at": row[1], "finished_at": row[2], "error": row[3]}
|
||||
|
||||
|
||||
def werkzeugfehler_im_lauf(zeilen: list[str], job_id: str, start_iso: str) -> tuple[int, str]:
|
||||
"""Zählt 'Tool X returned error'-Zeilen eines Laufs in errors.log. Hermes markiert jede
|
||||
Zeile mit [cron_<job>_<JJJJMMTT>_<HHMMSS>]; der Tag des Laufstarts reicht zur Zuordnung."""
|
||||
try:
|
||||
tag = datetime.fromisoformat(start_iso).strftime("%Y%m%d")
|
||||
except (TypeError, ValueError):
|
||||
return 0, ""
|
||||
marke = f"[cron_{job_id}_{tag}_"
|
||||
treffer = [z for z in zeilen if marke in z and "returned error" in z]
|
||||
if not treffer:
|
||||
return 0, ""
|
||||
m = re.search(r"Tool (\S+) returned error[^:]*:\s*(.*)$", treffer[0])
|
||||
beispiel = ""
|
||||
if m:
|
||||
beispiel = f"{m.group(1)}: {m.group(2)}"
|
||||
if (fm := re.search(r'"error":\s*"([^"]+)', m.group(2))):
|
||||
beispiel = f"{m.group(1)}: {fm.group(1)}"
|
||||
return len(treffer), beispiel[:220]
|
||||
|
||||
|
||||
def _errors_log_ende(max_bytes: int = 400_000) -> list[str]:
|
||||
try:
|
||||
with HERMES_ERRORS_LOG.open("rb") as f:
|
||||
f.seek(0, os.SEEK_END)
|
||||
f.seek(max(0, f.tell() - max_bytes))
|
||||
return f.read().decode("utf-8", "replace").splitlines()
|
||||
except OSError:
|
||||
return []
|
||||
|
||||
|
||||
def pruefe_hermes_jobs() -> list[Befund]:
|
||||
befunde: list[Befund] = []
|
||||
jobs = _hermes_jobs()
|
||||
if not jobs:
|
||||
return befunde
|
||||
log_zeilen = _errors_log_ende()
|
||||
jetzt = time.time()
|
||||
for job in jobs:
|
||||
if not job.get("enabled", True):
|
||||
continue
|
||||
jid, name = str(job.get("id", "")), str(job.get("name", "Job"))
|
||||
wichtig = "update" in name.lower()
|
||||
status = job.get("last_status")
|
||||
if status not in (None, "ok", "success") or int(job.get("failure_streak") or 0) > 0:
|
||||
befunde.append(Befund(
|
||||
id=f"job:{jid}", stufe="rot" if wichtig else "gelb",
|
||||
titel=f"Job „{name}“ ist fehlgeschlagen",
|
||||
text=str(job.get("last_error") or "Der letzte Lauf endete mit einem Fehler.")[:300],
|
||||
quelle=f"hermes-cron:{jid}",
|
||||
aktionen=[_aktion("job-wiederholen", "Erneut ausführen", job=jid),
|
||||
_aktion("protokoll", "Protokoll", job=jid)],
|
||||
sofort=True))
|
||||
continue
|
||||
if job.get("no_agent"):
|
||||
continue # Skript-Jobs haben keine Werkzeuge
|
||||
lauf = _letzter_lauf(jid)
|
||||
if not lauf or not lauf.get("started_at"):
|
||||
continue
|
||||
try:
|
||||
alter = jetzt - datetime.fromisoformat(lauf["started_at"]).timestamp()
|
||||
except ValueError:
|
||||
continue
|
||||
if alter > 26 * 3600:
|
||||
continue # nur der aktuelle Lauf zählt
|
||||
anzahl, beispiel = werkzeugfehler_im_lauf(log_zeilen, jid, lauf["started_at"])
|
||||
if anzahl:
|
||||
befunde.append(Befund(
|
||||
id=f"job:{jid}:werkzeug", stufe="gelb",
|
||||
titel=f"Job „{name}“: {anzahl} Werkzeugfehler im letzten Lauf",
|
||||
text=beispiel or "Einzelne Werkzeuge meldeten Fehler, der Lauf selbst kam durch.",
|
||||
quelle=f"hermes-cron:{jid}",
|
||||
aktionen=[_aktion("protokoll", "Protokoll", job=jid)],
|
||||
sofort=True))
|
||||
return befunde
|
||||
|
||||
|
||||
def _hermes_kopf() -> dict[str, str]:
|
||||
return {"Authorization": f"Bearer {HERMES_API_KEY}"} if HERMES_API_KEY else {}
|
||||
|
||||
|
||||
def pruefe_kern() -> list[Befund]:
|
||||
"""HTTP-Proben: läuft der Dienst UND antwortet er? (Der Dienst-Check sieht nur systemd.)"""
|
||||
proben: list[tuple[str, str, str, bool]] = [] # (id, Titel, Text, ok)
|
||||
engine_ok = llamaswap.engine_reachable()
|
||||
proben.append(("kern:engine", "Der Motor antwortet nicht",
|
||||
"llama-swap reagiert nicht. Ohne ihn laufen keine Modelle.", engine_ok))
|
||||
if engine_ok:
|
||||
st = llamaswap.brain_status()
|
||||
hirn_ok = bool(st.get("ready")) or bool(llamaswap.get_running_models())
|
||||
proben.append(("kern:hirn", "Lucys Hirn ist nicht geladen",
|
||||
"Das Hirn-Modell lädt nicht. Der Re-Warm-Wächter versucht es weiter.", hirn_ok))
|
||||
proben.append(("kern:hermes", "Hermes antwortet nicht",
|
||||
"Der Agent-Dienst reagiert nicht. Telegram und Lucys Werkzeuge gehen gerade nicht.",
|
||||
_reach(HERMES_API_URL, "/v1/models", _hermes_kopf())))
|
||||
proben.append(("kern:hoeren", "Der Hör-Dienst antwortet nicht",
|
||||
"Spracheingabe über Lucy-Desktop kann hängen.", _reach(VOICE_SERVICE_URL, "/health")))
|
||||
if WATCH_MC2:
|
||||
proben.append(("kern:mc2", "MC2 antwortet nicht",
|
||||
"Die Oberfläche und Lucys Sprach-Schnittstellen hängen.", _reach(MC2_URL, "/api/health")))
|
||||
proben.append(("kern:gateway", "Der Modell-Gateway antwortet nicht",
|
||||
"Anfragen von Lucy und OpenChamber an die Modelle hängen.", _reach(GATEWAY_URL, "/gw/health")))
|
||||
return [Befund(id=i, stufe="rot", titel=t, text=x, quelle=i.split(":", 1)[1])
|
||||
for (i, t, x, ok) in proben if not ok]
|
||||
|
||||
|
||||
def pruefe_platte() -> list[Befund]:
|
||||
try:
|
||||
pct = psutil.disk_usage(str(MODELS_DIR) if MODELS_DIR.exists() else os.getcwd()).percent
|
||||
except Exception:
|
||||
return []
|
||||
if pct >= DISK_ROT_PCT:
|
||||
stufe = "rot"
|
||||
elif pct >= DISK_GELB_PCT:
|
||||
stufe = "gelb"
|
||||
else:
|
||||
return []
|
||||
return [Befund(id="platte", stufe=stufe, titel=f"Die Platte ist zu {pct:.0f} % voll",
|
||||
text="Es wird eng für Modelle und Sicherungen. Alte Modelldateien löschen hilft am meisten.",
|
||||
quelle="platte", sofort=True)]
|
||||
|
||||
|
||||
PRUEFUNGEN = (pruefe_dienste, pruefe_timer_dienste, pruefe_hermes_jobs, pruefe_kern, pruefe_platte)
|
||||
|
||||
|
||||
# --- Zustand, Takt, Selbstreparatur -----------------------------------------------
|
||||
|
||||
_lock = threading.Lock()
|
||||
_stand: dict = {"hinweise": {}, "kandidaten": {}, "verlauf": [], "auto": {}, "stand": 0.0}
|
||||
|
||||
|
||||
def _lade() -> None:
|
||||
global _stand
|
||||
try:
|
||||
daten = json.loads(STORE_PATH.read_text(encoding="utf-8"))
|
||||
if isinstance(daten.get("hinweise"), dict):
|
||||
_stand = {**_stand, **daten}
|
||||
except (OSError, ValueError):
|
||||
pass
|
||||
|
||||
|
||||
def _speichere() -> None:
|
||||
try:
|
||||
tmp = STORE_PATH.with_suffix(".tmp")
|
||||
tmp.write_text(json.dumps(_stand, ensure_ascii=False), encoding="utf-8")
|
||||
tmp.replace(STORE_PATH)
|
||||
except OSError:
|
||||
log.warning("waechter: Stand %s nicht schreibbar", STORE_PATH, exc_info=True)
|
||||
|
||||
|
||||
def _verlauf(art: str, hinweis_id: str, text: str) -> None:
|
||||
_stand["verlauf"].append({"ts": time.time(), "art": art, "id": hinweis_id, "text": text})
|
||||
del _stand["verlauf"][:-VERLAUF_MAX]
|
||||
|
||||
|
||||
def _telegram(betreff: str, text: str) -> None:
|
||||
announce.add(text, subject=betreff, source="waechter")
|
||||
announce.notify_telegram(betreff, text + " (Diese Meldung kam auch an Lucy.)")
|
||||
|
||||
|
||||
def _auto_erlaubt(schluessel: str, jetzt: float) -> bool:
|
||||
versuche = [t for t in _stand["auto"].get(schluessel, []) if jetzt - t < 3600]
|
||||
_stand["auto"][schluessel] = versuche
|
||||
return len(versuche) < AUTO_MAX_PRO_STUNDE
|
||||
|
||||
|
||||
def _selbst_beheben(b: Befund, jetzt: float) -> None:
|
||||
if b.auto != "neustart" or not _auto_erlaubt(b.id, jetzt):
|
||||
return
|
||||
ergebnis = maintenance.restart_service(b.quelle)
|
||||
_stand["auto"][b.id].append(jetzt)
|
||||
ok = bool(ergebnis.get("ok", ergebnis.get("returncode", 1) == 0))
|
||||
_verlauf("auto", b.id, f"{b.titel}: automatisch neu gestartet" + ("" if ok else " (ohne Erfolg)"))
|
||||
log.warning("waechter: %s → Neustart von %s (%s)", b.id, b.quelle, "ok" if ok else "fehlgeschlagen")
|
||||
|
||||
|
||||
def takt() -> None:
|
||||
"""Ein Prüfdurchlauf: Befunde sammeln, Hinweise führen, Einfaches selbst beheben."""
|
||||
jetzt = time.time()
|
||||
update = _update_laeuft()
|
||||
befunde: list[Befund] = []
|
||||
for pruefung in PRUEFUNGEN:
|
||||
if update and pruefung in (pruefe_dienste, pruefe_kern):
|
||||
continue # während eines Updates sind Neustarts normal
|
||||
try:
|
||||
befunde += pruefung()
|
||||
except Exception:
|
||||
log.debug("waechter: %s fehlgeschlagen", pruefung.__name__, exc_info=True)
|
||||
|
||||
with _lock:
|
||||
hinweise: dict = _stand["hinweise"]
|
||||
kandidaten: dict = _stand["kandidaten"]
|
||||
aktuell = {b.id for b in befunde}
|
||||
|
||||
for b in befunde:
|
||||
k = kandidaten.setdefault(b.id, {"takte": 0, "erstmals": jetzt})
|
||||
k["takte"] += 1
|
||||
if not update and b.auto:
|
||||
_selbst_beheben(b, jetzt)
|
||||
if not (b.sofort or k["takte"] >= FAIL_AFTER):
|
||||
continue
|
||||
h = hinweise.get(b.id)
|
||||
neu = h is None
|
||||
h = h or {"id": b.id, "seit": k["erstmals"], "gemeldet": 0.0}
|
||||
h.update(stufe=b.stufe, titel=b.titel, text=b.text, quelle=b.quelle,
|
||||
aktionen=b.aktionen, zuletzt=jetzt, takte=k["takte"])
|
||||
hinweise[b.id] = h
|
||||
if neu:
|
||||
_verlauf("neu", b.id, b.titel)
|
||||
if b.stufe == "rot" and (neu or jetzt - h.get("gemeldet", 0.0) >= REMIND_S):
|
||||
vorsatz = "" if neu else "Immer noch: "
|
||||
_telegram("[Box-Problem]", f"{vorsatz}{b.titel}. {b.text}")
|
||||
h["gemeldet"] = jetzt
|
||||
|
||||
# Nicht mehr festgestellt → erledigt. Während eines Updates bleiben Dienst- und
|
||||
# Kern-Hinweise stehen (sie wurden in diesem Takt gar nicht geprüft).
|
||||
for hid in list(hinweise):
|
||||
if hid in aktuell:
|
||||
continue
|
||||
if update and hid.startswith(("dienst:", "kern:")):
|
||||
continue
|
||||
h = hinweise.pop(hid)
|
||||
_verlauf("erledigt", hid, h.get("titel", hid))
|
||||
if h.get("stufe") == "rot" and h.get("gemeldet"):
|
||||
_telegram("[Box wieder ok]", f"Erledigt: {h.get('titel', hid)}.")
|
||||
for kid in list(kandidaten):
|
||||
if kid not in aktuell and not (update and kid.startswith(("dienst:", "kern:"))):
|
||||
kandidaten.pop(kid)
|
||||
|
||||
_stand["stand"] = jetzt
|
||||
_stand["update_laeuft"] = update
|
||||
_speichere()
|
||||
|
||||
|
||||
async def waechter_loop() -> None:
|
||||
"""Endlos-Schleife im mc2-steward."""
|
||||
_lade()
|
||||
await asyncio.sleep(START_DELAY)
|
||||
log.info("waechter: aktiv (Takt %ss, Hinweis nach %s Takten, Selbstreparatur max. %s/h)",
|
||||
INTERVAL, FAIL_AFTER, AUTO_MAX_PRO_STUNDE)
|
||||
while True:
|
||||
try:
|
||||
await asyncio.to_thread(takt)
|
||||
except Exception:
|
||||
log.warning("waechter: Takt fehlgeschlagen", exc_info=True)
|
||||
await asyncio.sleep(INTERVAL)
|
||||
|
||||
|
||||
# --- Lesen und Knöpfe (MC2-Prozess) ------------------------------------------------
|
||||
|
||||
def lese_stand() -> dict:
|
||||
"""Stand für die Oberfläche: Hinweise (rot zuerst, dann nach Beginn) + Verlauf."""
|
||||
try:
|
||||
daten = json.loads(STORE_PATH.read_text(encoding="utf-8"))
|
||||
except (OSError, ValueError):
|
||||
daten = {}
|
||||
hinweise = sorted((daten.get("hinweise") or {}).values(),
|
||||
key=lambda h: (h.get("stufe") != "rot", h.get("seit", 0)))
|
||||
for h in hinweise:
|
||||
h.pop("gemeldet", None)
|
||||
return {
|
||||
"hinweise": hinweise,
|
||||
"verlauf": list(reversed((daten.get("verlauf") or [])[-50:])),
|
||||
"stand": daten.get("stand"),
|
||||
"update_laeuft": bool(daten.get("update_laeuft")),
|
||||
"aktiv": bool(daten),
|
||||
}
|
||||
|
||||
|
||||
def _job_protokoll(job_id: str) -> dict:
|
||||
marke = f"[cron_{job_id}_"
|
||||
zeilen = [z for z in _errors_log_ende() if marke in z][-60:]
|
||||
return {"ok": True, "text": "\n".join(zeilen) or "Keine Fehlerzeilen zu diesem Job gefunden."}
|
||||
|
||||
|
||||
def fuehre_aktion_aus(hinweis_id: str, aktion_id: str) -> dict:
|
||||
"""Knopf eines Hinweises ausführen. Nur Aktionen, die der Hinweis selbst anbietet."""
|
||||
stand = lese_stand()
|
||||
hinweis = next((h for h in stand["hinweise"] if h.get("id") == hinweis_id), None)
|
||||
if hinweis is None:
|
||||
return {"ok": False, "detail": "Diesen Hinweis gibt es nicht mehr."}
|
||||
aktion = next((a for a in hinweis.get("aktionen", []) if a.get("id") == aktion_id), None)
|
||||
if aktion is None:
|
||||
return {"ok": False, "detail": "Diese Aktion gehört nicht zu dem Hinweis."}
|
||||
if aktion_id == "neustart":
|
||||
return maintenance.restart_service(aktion["dienst"])
|
||||
if aktion_id == "protokoll":
|
||||
if aktion.get("job"):
|
||||
return _job_protokoll(aktion["job"])
|
||||
return maintenance.logs(aktion["dienst"], lines=120)
|
||||
if aktion_id == "job-wiederholen":
|
||||
try:
|
||||
r = subprocess.run(["hermes", "cron", "run", aktion["job"]],
|
||||
capture_output=True, text=True, timeout=30)
|
||||
return {"ok": r.returncode == 0,
|
||||
"text": (r.stdout or r.stderr).strip()[:500] or "Job läuft beim nächsten Takt."}
|
||||
except Exception as e:
|
||||
return {"ok": False, "detail": f"hermes cron run ging nicht: {e}"}
|
||||
return {"ok": False, "detail": f"Unbekannte Aktion {aktion_id}."}
|
||||
Reference in New Issue
Block a user