Drei Seiten statt zehn: Start, Updates, Modelle — am PC mit Kopfnavigation, am Handy mit Leiste unten. Umsetzung von Mockup A („Cockpit“), vom User am 23.09. gewaehlt. - Start: Hauptwarnleuchte + 8 Warnlampen, Rundinstrumente mit Live-Werten aus dem Strom (Speicher, Temperatur, Platte) und Laufzeit-Zaehlwerk, Checkliste der Waechter-Hinweise mit ihren Knoepfen, Flugplan (heute gelaufen / geplant), Radar-Kasten. - Updates: Bausteine mit „Laeuft → Neu“ und Zusammenfassung, laufende Auftraege, Verlauf der Sonntagslaeufe (neu: GET /api/updates/verlauf), Sicherungen samt Zurueckspielen mit Rueckfrage. - Modelle: Speicherbalken, Rollen Hirn/Coder/Dritte Rolle, wer die Modelle nutzt (7 Tage + 24 h je Stunde), Modell-Radar, weitere Eintraege, Modelle selbst suchen. - Schubladen: Dienste mit Protokoll und Neustart, Einstellungen (HF-Zugang), Hermes-Link. - Werkzeuge: Vite 8, React 19.3 mit React Compiler 1.0 (Babel), vitest 5, Tailwind 4.3, shadcn 4 (Radix) fuer Dialog/Schublade/Knopf, Schriften Barlow/Barlow Condensed/ JetBrains Mono. Entfernt: recharts, cmdk, Kraftgraph, zustand, Inter, Space Grotesk. - Startbuendel 115 KB gzip (Budget 140); Updates/Modelle/Schubladen laden bei Bedarf. - 16 Oberflaechen-Tests (Instrument-Bogen, Hauptleuchte, Zeitformate, Versionen). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
558 lines
24 KiB
Python
558 lines
24 KiB
Python
"""
|
|
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 update_laeufe(anzahl: int = 8) -> list[dict]:
|
|
"""Letzte Läufe aller Update-Jobs (Name enthält „update“), neueste zuerst."""
|
|
ids = {str(j.get("id")): str(j.get("name", "")) for j in _hermes_jobs()
|
|
if "update" in str(j.get("name", "")).lower()}
|
|
if not ids or not HERMES_EXEC_DB.exists():
|
|
return []
|
|
platz = ",".join("?" * len(ids))
|
|
try:
|
|
con = sqlite3.connect(f"file:{HERMES_EXEC_DB}?mode=ro", uri=True, timeout=2)
|
|
try:
|
|
rows = con.execute(
|
|
f"select job_id, status, started_at, finished_at, error from executions "
|
|
f"where job_id in ({platz}) order by started_at desc limit ?",
|
|
(*ids, anzahl)).fetchall()
|
|
finally:
|
|
con.close()
|
|
except sqlite3.Error:
|
|
return []
|
|
return [{"job": ids.get(r[0], r[0]), "status": r[1], "start": r[2], "ende": r[3],
|
|
"fehler": (r[4] or None)} for r in rows]
|
|
|
|
|
|
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}."}
|