177 lines
6.3 KiB
Python
177 lines
6.3 KiB
Python
"""Kanal zum Ausführer auf dem Proxmox-Host (Phase 3/4, 24.09.2026).
|
|
|
|
Der Ausführer (deploy/homelab/ausfuehrer.py) holt sich Aufträge hier ab, liefert die Ergebnisse und
|
|
alle zehn Minuten einen Bericht. Er weist sich mit einem gemeinsamen Geheimnis aus (Kopfzeile
|
|
X-MC2-Ausfuehrer). Das erzeugt dieser Teil beim ersten Bedarf in <Datenordner>/ausfuehrer.token (nur
|
|
für den Dienst lesbar); die Einrichtung auf dem Proxmox-Host kopiert es in /etc/mc2-ausfuehrer.json.
|
|
|
|
Aufträge liegen in <Datenordner>/ausfuehrer-auftraege.json und überleben so einen Neustart dieses Teils.
|
|
"""
|
|
|
|
import hmac
|
|
import json
|
|
import os
|
|
import secrets
|
|
import threading
|
|
import time
|
|
import uuid
|
|
from pathlib import Path
|
|
|
|
from kern.einstellungen import einstellungen
|
|
|
|
AKTIONEN = {"bericht", "snapshot", "update", "os_update", "suchen", "zurueck", "snapshot_loeschen",
|
|
"host_update", "host_neustart"}
|
|
VERLOREN_S = 3 * 3600 # abgeholt, aber nie beantwortet (Ausführer abgestürzt, Host neu gestartet)
|
|
BEHALTEN = 60 # so viele erledigte Aufträge bleiben sichtbar
|
|
|
|
_lock = threading.Lock()
|
|
_kontakt: dict = {"zuletzt": None}
|
|
|
|
|
|
def _dir() -> Path:
|
|
return einstellungen().daten_dir
|
|
|
|
|
|
def _json_schreiben(pfad: Path, daten: object) -> None:
|
|
pfad.parent.mkdir(parents=True, exist_ok=True)
|
|
tmp = pfad.with_suffix(pfad.suffix + ".tmp")
|
|
tmp.write_text(json.dumps(daten, ensure_ascii=False), encoding="utf-8")
|
|
tmp.replace(pfad)
|
|
|
|
|
|
def _json_lesen(pfad: Path, leer: object) -> object:
|
|
try:
|
|
return json.loads(pfad.read_text(encoding="utf-8"))
|
|
except (OSError, ValueError):
|
|
return leer
|
|
|
|
|
|
# --- Gemeinsames Geheimnis ---------------------------------------------------------------
|
|
|
|
def token() -> str:
|
|
"""Das Geheimnis des Kanals; beim ersten Aufruf erzeugt (0600). MC_AUSFUEHRER_TOKEN geht vor."""
|
|
if wert := os.environ.get("MC_AUSFUEHRER_TOKEN", "").strip():
|
|
return wert
|
|
pfad = _dir() / "ausfuehrer.token"
|
|
try:
|
|
return pfad.read_text(encoding="utf-8").strip()
|
|
except OSError:
|
|
pass
|
|
wert = secrets.token_urlsafe(32)
|
|
pfad.parent.mkdir(parents=True, exist_ok=True)
|
|
try:
|
|
fd = os.open(pfad, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600)
|
|
except FileExistsError: # gleichzeitig erzeugt: den anderen nehmen
|
|
return pfad.read_text(encoding="utf-8").strip()
|
|
with os.fdopen(fd, "w", encoding="utf-8") as f:
|
|
f.write(wert + "\n")
|
|
return wert
|
|
|
|
|
|
def token_gueltig(wert: str | None) -> bool:
|
|
return bool(wert) and hmac.compare_digest(str(wert), token())
|
|
|
|
|
|
def kontakt() -> None:
|
|
_kontakt["zuletzt"] = time.time()
|
|
|
|
|
|
def zuletzt() -> float | None:
|
|
"""Letztes Lebenszeichen des Ausführers (nach einem Neustart: Zeit des letzten Berichts)."""
|
|
return _kontakt["zuletzt"] or (bericht() or {}).get("empfangen")
|
|
|
|
|
|
# --- Bericht ------------------------------------------------------------------------------
|
|
|
|
def bericht_speichern(daten: dict) -> None:
|
|
with _lock:
|
|
_json_schreiben(_dir() / "pve-bericht.json", {"empfangen": time.time(), "bericht": daten})
|
|
|
|
|
|
def bericht() -> dict | None:
|
|
"""{"empfangen": Unix-Sekunden, "bericht": {...}} oder None, solange nie einer kam."""
|
|
with _lock:
|
|
daten = _json_lesen(_dir() / "pve-bericht.json", None)
|
|
return daten if isinstance(daten, dict) and isinstance(daten.get("bericht"), dict) else None
|
|
|
|
|
|
# --- Aufträge -------------------------------------------------------------------------------
|
|
|
|
def _pfad() -> Path:
|
|
return _dir() / "ausfuehrer-auftraege.json"
|
|
|
|
|
|
def _alle() -> list[dict]:
|
|
daten = _json_lesen(_pfad(), [])
|
|
return daten if isinstance(daten, list) else []
|
|
|
|
|
|
def _merken(auftraege: list[dict]) -> None:
|
|
offen = [a for a in auftraege if a["status"] in ("wartet", "laeuft")]
|
|
fertig = [a for a in auftraege if a["status"] not in ("wartet", "laeuft")][-BEHALTEN:]
|
|
_json_schreiben(_pfad(), sorted(offen + fertig, key=lambda a: a["erstellt"]))
|
|
|
|
|
|
def anlegen(aktion: str, parameter: dict | None = None) -> str:
|
|
if aktion not in AKTIONEN:
|
|
raise ValueError(f"Unbekannte Aktion: {aktion}")
|
|
auftrag = {"id": uuid.uuid4().hex[:12], "aktion": aktion, "parameter": parameter or {},
|
|
"status": "wartet", "erstellt": time.time(), "abgeholt": None, "fertig": None,
|
|
"code": None, "text": None}
|
|
with _lock:
|
|
auftraege = _alle()
|
|
auftraege.append(auftrag)
|
|
_merken(auftraege)
|
|
return auftrag["id"]
|
|
|
|
|
|
def naechster() -> dict | None:
|
|
"""Für den Ausführer: der älteste wartende Auftrag, ab jetzt „läuft“. Hängengebliebene werden
|
|
vorher als verloren abgeschlossen."""
|
|
jetzt = time.time()
|
|
with _lock:
|
|
auftraege = _alle()
|
|
for a in auftraege:
|
|
if a["status"] == "laeuft" and jetzt - (a["abgeholt"] or jetzt) > VERLOREN_S:
|
|
a.update(status="verloren", fertig=jetzt, text="Der Ausführer hat nie geantwortet.")
|
|
auftrag = next((a for a in auftraege if a["status"] == "wartet"), None)
|
|
if auftrag:
|
|
auftrag.update(status="laeuft", abgeholt=jetzt)
|
|
_merken(auftraege)
|
|
return {k: auftrag[k] for k in ("id", "aktion", "parameter")} if auftrag else None
|
|
|
|
|
|
def ergebnis(auftrag_id: str, code: int, text: str) -> bool:
|
|
with _lock:
|
|
auftraege = _alle()
|
|
auftrag = next((a for a in auftraege if a["id"] == auftrag_id), None)
|
|
if auftrag is None or auftrag["status"] != "laeuft":
|
|
return False
|
|
auftrag.update(status="fertig" if code == 0 else "fehler", fertig=time.time(), code=code,
|
|
text=str(text)[-20000:])
|
|
_merken(auftraege)
|
|
return True
|
|
|
|
|
|
def auftrag(auftrag_id: str) -> dict | None:
|
|
# Lesen unter derselben Sperre wie das Schreiben: Unter Windows scheitert das atomare Ersetzen,
|
|
# solange ein anderer Thread die Datei offen hat.
|
|
with _lock:
|
|
return next((a for a in _alle() if a["id"] == auftrag_id), None)
|
|
|
|
|
|
def warten(auftrag_id: str, zeitlimit_s: float, takt_s: float = 2.0) -> dict | None:
|
|
"""Bis der Auftrag fertig ist (oder das Zeitlimit abläuft). Rückgabe: der Auftrag, None bei Zeitablauf."""
|
|
ende = time.time() + zeitlimit_s
|
|
while time.time() < ende:
|
|
a = auftrag(auftrag_id)
|
|
if a and a["status"] not in ("wartet", "laeuft"):
|
|
return a
|
|
time.sleep(takt_s)
|
|
return None
|
|
|
|
|
|
def liste() -> list[dict]:
|
|
with _lock:
|
|
return list(reversed(_alle()))
|