""" Auftrags-System: Hintergrund-Aufträge (Updates, Modell-Downloads) mit Protokoll und Fortschritt. Ursprünglich aus Mission Control v1 portiert. Seit Phase 2 (24.09.2026) überleben Aufträge einen Neustart von MC2: - Jeder Auftrag läuft als eigene systemd-Einheit „mc2-job-“ (systemd-run), nicht mehr als Kindprozess von MC2. Ein Deploy oder Absturz von MC2 würgt ihn nicht mehr ab. - Akte (.json), Protokoll (.log) und Exit-Code (.exit) liegen unter JOBS_DIR. MC2 liest die Akten beim Start wieder ein (wiederaufnehmen()) und beobachtet weiter. - Nacharbeiten („Rolle setzen, wenn der Download fertig ist“) sind benannt (@nacharbeit) und stehen mit ihren Daten in der Akte — sie laufen auch, wenn MC2 zwischendurch neu gestartet hat. - Geheimnisse (z. B. HF_TOKEN) gehen über eine Umgebungsdatei (nur für den Nutzer lesbar), die der Auftrag beim Start liest und löscht — nicht über die Befehlszeile oder die Einheit. - Ohne systemd (Entwicklungsrechner, Tests: MC_JOBS_ART=prozess) läuft der Auftrag als Kindprozess; dann überlebt er einen Neustart nicht. Weiterhin: „Abbrechen“ beendet den ganzen Auftrag samt Kindern, jede Auftragsart hat ein Zeitlimit, start_job_exklusiv prüft und trägt unter EINER Sperre ein, beendete Aufträge verschwinden nach einem Tag bzw. ab 40 Stück. Nur MC2 selbst beobachtet Aufträge (wiederaufnehmen() im Lebenszyklus der App); andere Prozesse starten keine. """ import glob import json import logging import os import re import shlex import shutil import signal import subprocess import threading import time import uuid from collections.abc import Callable from pathlib import Path from kern.einstellungen import einstellungen log = logging.getLogger(__name__) JOBS: dict[str, dict] = {} _PROCS: dict[str, subprocess.Popen] = {} _LOCK = threading.Lock() _LOG_CAP = 400 _LOG_LESEN_MAX = 256_000 # so viel vom Protokollende wird für die Anzeige gelesen BEHALTEN_MAX = 40 # so viele beendete Aufträge bleiben sichtbar BEHALTEN_S = 24 * 3600 # beendete Aufträge verschwinden nach einem Tag STANDARD_ZEITLIMIT_S = 3 * 3600 _ENDE = ("done", "failed", "canceled") _TAKT_S = 1.0 # wie oft ein Beobachter nach dem Exit-Code schaut _EINHEIT_PRUEFEN_S = 5.0 # wie oft er zusätzlich fragt, ob die Einheit noch lebt # Nur im Speicher (nicht in der Akte): Fortschritt eines Downloads. _LAUFZEIT = {"progress", "done_bytes", "rate_bps", "eta_s"} # Nicht an die Oberfläche: Verwaltung des Auftrags. _INTERN = {"art", "nacharbeit", "nacharbeit_daten", "nacharbeit_erledigt", "fortschritt", "zeitlimit_s"} # Nicht in die Umgebung des Auftrags: gehört zur Einheit von MC2 oder ist in bash schreibgeschützt. _UMGEBUNG_OHNE = {"INVOCATION_ID", "JOURNAL_STREAM", "SYSTEMD_EXEC_PID", "MANAGERPID", "NOTIFY_SOCKET", "LISTEN_PID", "LISTEN_FDS", "LISTEN_FDNAMES", "WATCHDOG_PID", "WATCHDOG_USEC", "MAINPID", "_", "SHLVL", "PWD", "OLDPWD", "SHELLOPTS", "BASHOPTS", "BASH_VERSINFO", "EUID", "UID", "PPID", "GROUPS"} _NAME = re.compile(r"^[A-Za-z_][A-Za-z0-9_]*$") # Die Hülle um den eigentlichen Befehl: Umgebung lesen (und die Datei löschen), Ausgabe ins # Protokoll, Exit-Code atomar ablegen. $1 = Pfadpräfix der Akte, danach der Befehl. _HUELLE = ('akte="$1"; shift; ' 'if [ -f "$akte.env" ]; then . "$akte.env"; rm -f "$akte.env"; fi; ' 'exec >>"$akte.log" 2>&1 "$akte.exit.tmp" && mv -f "$akte.exit.tmp" "$akte.exit"') _NACHARBEITEN: dict[str, Callable[..., None]] = {} _geladen = False def nacharbeit(name: str) -> Callable[[Callable[..., None]], Callable[..., None]]: """Eine Nacharbeit unter festem Namen anmelden. Sie bekommt die nacharbeit_daten des Auftrags als Schlüsselwort-Argumente und läuft nur, wenn der Auftrag erfolgreich war.""" def anmelden(fn: Callable[..., None]) -> Callable[..., None]: _NACHARBEITEN[name] = fn return fn return anmelden # --- Ablage ------------------------------------------------------------------------ def _jobs_dir() -> Path: return Path(os.environ.get("MC_JOBS_DIR", str(einstellungen().daten_dir / "mc2-jobs"))) def _pfad(job_id: str, endung: str) -> Path: return _jobs_dir() / f"{job_id}{endung}" def _speichern(job: dict) -> None: """Akte ohne Laufzeitfelder atomar schreiben.""" daten = {k: v for k, v in job.items() if k not in _LAUFZEIT} try: _jobs_dir().mkdir(parents=True, exist_ok=True) tmp = _pfad(job["id"], ".json.tmp") tmp.write_text(json.dumps(daten, ensure_ascii=False), encoding="utf-8") tmp.replace(_pfad(job["id"], ".json")) except OSError: log.warning("jobengine: Akte %s nicht schreibbar", job.get("id"), exc_info=True) def _protokoll(job_id: str, zeile: str) -> None: """Eine Zeile von MC2 ins Protokoll des Auftrags schreiben.""" try: _jobs_dir().mkdir(parents=True, exist_ok=True) with _pfad(job_id, ".log").open("a", encoding="utf-8", newline="\n") as f: f.write(zeile.rstrip("\n") + "\n") except OSError: pass def zeilen_aus(text: str) -> list[str]: """Protokolltext in Anzeigezeilen: `\\r` (Fortschrittsbalken von hf/tqdm/apt) überschreibt die Zeile.""" zeilen = [] for roh in text.split("\n"): if "\r" in roh: teile = [t for t in roh.split("\r") if t] roh = teile[-1] if teile else "" zeilen.append(roh) if zeilen and zeilen[-1] == "": zeilen.pop() return zeilen[-_LOG_CAP:] def _protokoll_lesen(job_id: str) -> list[str]: try: with _pfad(job_id, ".log").open("rb") as f: f.seek(0, os.SEEK_END) groesse = f.tell() f.seek(max(0, groesse - _LOG_LESEN_MAX)) roh = f.read() except OSError: return [] text = roh.decode("utf-8", "replace") if groesse > _LOG_LESEN_MAX: text = text.split("\n", 1)[-1] # angeschnittene erste Zeile verwerfen return zeilen_aus(text) def _exit_code(job_id: str) -> int | None: try: return int(_pfad(job_id, ".exit").read_text(encoding="utf-8").strip()) except (OSError, ValueError): return None def _exit_ablegen(job_id: str, code: int) -> None: try: tmp = _pfad(job_id, ".exit.tmp") tmp.write_text(str(code), encoding="utf-8") tmp.replace(_pfad(job_id, ".exit")) except OSError: pass def _umgebungsdatei(job_id: str, env: dict | None) -> None: """Umgebung des Auftrags (MC2-Umgebung + env) als nur für den Nutzer lesbare Datei.""" werte = {**os.environ, **(env or {})} zeilen = [f"export {k}={shlex.quote(str(v))}" for k, v in werte.items() if _NAME.match(k) and k not in _UMGEBUNG_OHNE] fd = os.open(_pfad(job_id, ".env"), os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600) with os.fdopen(fd, "w", encoding="utf-8", newline="\n") as f: f.write("\n".join(zeilen) + "\n") # --- Ausführung ---------------------------------------------------------------------- def _art() -> str: """„systemd“, wenn Aufträge als eigene Einheiten laufen können, sonst „prozess“.""" gewuenscht = os.environ.get("MC_JOBS_ART", "").strip().lower() if gewuenscht in ("systemd", "prozess"): return gewuenscht if os.name != "posix" or not shutil.which("systemd-run"): return "prozess" if os.geteuid() != 0 and not os.environ.get("XDG_RUNTIME_DIR"): return "prozess" return "systemd" def _systemd(programm: str, *args: str) -> list[str]: """systemctl/systemd-run für den richtigen Manager: als root der System-, sonst der Nutzer-Manager.""" als_root = os.name == "posix" and os.geteuid() == 0 return [programm, *(() if als_root else ("--user",)), *args] def _einheit(job_id: str) -> str: return f"mc2-job-{job_id}" def _einheit_lebt(job_id: str) -> bool: try: r = subprocess.run(_systemd("systemctl", "is-active", _einheit(job_id)), capture_output=True, text=True, timeout=10) except Exception: return True # Im Zweifel weiter beobachten statt einen laufenden Auftrag totzusagen. return r.stdout.strip() in ("active", "activating", "deactivating", "reloading") def _starte_systemd(job: dict, args: list[str], env: dict | None) -> None: job_id = job["id"] _umgebungsdatei(job_id, env) befehl = _systemd("systemd-run", "--unit", _einheit(job_id), "--collect", "--quiet", "--working-directory", os.getcwd(), "--property", f"RuntimeMaxSec={int(job['zeitlimit_s'])}", "--", "/bin/bash", "-c", _HUELLE, "mc2-job", str(_jobs_dir() / job_id), *args) r = subprocess.run(befehl, capture_output=True, text=True, timeout=30) if r.returncode != 0: _pfad(job_id, ".env").unlink(missing_ok=True) raise RuntimeError((r.stderr or r.stdout or "systemd-run scheiterte").strip()[:300]) threading.Thread(target=_beobachten, args=(job_id,), daemon=True).start() def _starte_prozess(job: dict, args: list[str], env: dict | None) -> None: """Rückfall ohne systemd: Kindprozess, Ausgabe ins Protokoll; ein Warte-Thread schließt ab.""" job_id = job["id"] with _pfad(job_id, ".log").open("ab") as protokoll: proc = subprocess.Popen(list(args), stdout=protokoll, stderr=subprocess.STDOUT, stdin=subprocess.DEVNULL, env={**os.environ, **(env or {})}, start_new_session=(os.name == "posix")) _PROCS[job_id] = proc grenze = int(job["zeitlimit_s"]) def zu_lang() -> None: job["zeitlimit"] = True _protokoll(job_id, f"[mc] Zeitlimit ({grenze // 60} min) überschritten, Auftrag wird beendet.") _beende_gruppe(proc) wecker = threading.Timer(grenze, zu_lang) wecker.daemon = True wecker.start() def warten() -> None: code = proc.wait() wecker.cancel() _exit_ablegen(job_id, code) _abschliessen(job, code) threading.Thread(target=warten, daemon=True).start() def _beende_gruppe(proc: subprocess.Popen) -> None: """Den Auftragsprozess samt Kindern beenden (posix: ganze Prozessgruppe, sonst nur den Prozess).""" try: if os.name == "posix": os.killpg(os.getpgid(proc.pid), signal.SIGTERM) else: proc.terminate() except (ProcessLookupError, PermissionError, OSError): pass def _abschliessen(job: dict, code: int | None) -> None: """Endzustand festhalten (genau einmal) und, wenn erfolgreich, die Nacharbeit ausführen.""" with _LOCK: if job["state"] in _ENDE: return job["returncode"] = code job["finished_at"] = time.time() if job.get("canceled"): job["state"] = "canceled" elif job.get("zeitlimit") or (code is None and time.time() - job["started_at"] >= job["zeitlimit_s"] - 5): job.update(state="failed", zeitlimit=True, error="Zeitlimit überschritten") elif code is None: job["state"] = "failed" job.setdefault("error", "Der Auftrag endete ohne Rückmeldung.") else: job["state"] = "done" if code == 0 else "failed" _PROCS.pop(job["id"], None) _speichern(job) if job["state"] == "done": _nacharbeit_ausfuehren(job) def _nacharbeit_ausfuehren(job: dict) -> None: name = job.get("nacharbeit") if not name or job.get("nacharbeit_erledigt"): return job["nacharbeit_erledigt"] = True _speichern(job) try: fn = _NACHARBEITEN.get(name) if fn is None: raise RuntimeError(f"Nacharbeit „{name}“ ist in dieser Instanz unbekannt") fn(**(job.get("nacharbeit_daten") or {})) except Exception as exc: _protokoll(job["id"], f"[mc] Nachbearbeitung-Fehler: {exc}") log.warning("jobengine: Nacharbeit %s von %s gescheitert", name, job["id"], exc_info=True) def _beobachten(job_id: str) -> None: """Wartet auf den Exit-Code einer systemd-Einheit. Endet sie ohne einen (Abbruch, Zeitlimit, Absturz), schließt er den Auftrag trotzdem ab.""" naechste_pruefung = time.time() + _EINHEIT_PRUEFEN_S while True: job = JOBS.get(job_id) if job is None or job["state"] in _ENDE: return if (code := _exit_code(job_id)) is not None: _abschliessen(job, code) return if time.time() >= naechste_pruefung: if not _einheit_lebt(job_id): time.sleep(0.5) # der Exit-Code kann gerade erst geschrieben werden _abschliessen(job, _exit_code(job_id)) return naechste_pruefung = time.time() + _EINHEIT_PRUEFEN_S time.sleep(_TAKT_S) def _starten(job: dict, args: list[str], env: dict | None) -> None: job["state"] = "running" _speichern(job) try: if job["art"] == "systemd": _starte_systemd(job, args, env) else: _starte_prozess(job, args, env) except Exception as exc: _protokoll(job["id"], f"[mc] Fehler: {exc}") job["error"] = str(exc) _abschliessen(job, -1) # --- Fortschritt (Downloads) ------------------------------------------------------------ def attach_download_progress(job_id: str, local_dir: str, total_bytes: int) -> None: """Fortschritt in % aus wachsenden *.incomplete-Dateien (hf schreibt sie).""" job = JOBS.get(job_id) if job is None or not total_bytes or total_bytes <= 0: return job["fortschritt"] = {"ordner": local_dir, "gesamt": total_bytes} job["total_bytes"] = total_bytes job["progress"] = 0 _speichern(job) _fortschritt_beobachten(job_id) def _fortschritt_beobachten(job_id: str) -> None: ziel = (JOBS.get(job_id) or {}).get("fortschritt") or {} local_dir, total_bytes = ziel.get("ordner"), ziel.get("gesamt") if not local_dir or not total_bytes: return def _watch(): pat = os.path.join(local_dir, ".cache", "huggingface", "download", "**", "*.incomplete") prev_t = prev_b = None rate = 0.0 while True: j = JOBS.get(job_id) if not j or j["state"] in _ENDE: break try: inc = glob.glob(pat, recursive=True) cur = sum(os.path.getsize(f) for f in inc) if inc else 0 if cur: j["progress"] = min(99, int(cur * 100 / total_bytes)) j["done_bytes"] = cur now = time.time() if prev_t is not None and now > prev_t and cur >= prev_b: inst = (cur - prev_b) / (now - prev_t) rate = inst if rate == 0 else 0.3 * inst + 0.7 * rate if rate > 0: j["rate_bps"] = rate j["eta_s"] = int((total_bytes - cur) / rate) prev_t, prev_b = now, cur except Exception: pass time.sleep(1.0) j = JOBS.get(job_id) if j and j["state"] == "done": j["progress"] = 100 j.pop("eta_s", None) threading.Thread(target=_watch, daemon=True).start() # --- Verwaltung ------------------------------------------------------------------- def _laden() -> None: """Akten von der Platte einlesen (einmal je Prozess; unter _LOCK aufrufen).""" global _geladen if _geladen: return _geladen = True for datei in sorted(_jobs_dir().glob("*.json")): try: akte = json.loads(datei.read_text(encoding="utf-8")) except (OSError, ValueError): continue if isinstance(akte, dict) and akte.get("id") and akte["id"] not in JOBS: JOBS[akte["id"]] = akte def _aufraeumen() -> None: """Alte beendete Aufträge samt Dateien entfernen (unter _LOCK aufrufen).""" jetzt = time.time() fertig = sorted((j for j in JOBS.values() if j["state"] in _ENDE), key=lambda j: j.get("finished_at") or 0, reverse=True) for i, j in enumerate(fertig): if i >= BEHALTEN_MAX or jetzt - (j.get("finished_at") or jetzt) > BEHALTEN_S: JOBS.pop(j["id"], None) for endung in (".json", ".log", ".exit", ".env", ".exit.tmp", ".json.tmp"): _pfad(j["id"], endung).unlink(missing_ok=True) def _eintragen(args: list[str], label: str, group: str | None, zeitlimit_s: int | None, nacharbeit_name: str | None, nacharbeit_daten: dict | None) -> dict: """Neue Akte anlegen (unter _LOCK aufrufen).""" if nacharbeit_name and nacharbeit_name not in _NACHARBEITEN: raise ValueError(f"Unbekannte Nacharbeit: {nacharbeit_name}") job_id = uuid.uuid4().hex[:12] job = { "id": job_id, "label": label, "state": "queued", "group": group, "art": _art(), "returncode": None, "started_at": time.time(), "finished_at": None, "zeitlimit_s": int(zeitlimit_s or STANDARD_ZEITLIMIT_S), "nacharbeit": nacharbeit_name, "nacharbeit_daten": nacharbeit_daten or {}, } JOBS[job_id] = job _protokoll(job_id, "$ " + " ".join(shlex.quote(a) for a in args)) _speichern(job) return job def _aktiv_in(group: str) -> dict | None: for j in list(JOBS.values()): if j.get("group") == group and j.get("state") in ("running", "queued"): return j return None def start_job(args: list[str], label: str, env: dict | None = None, group: str | None = None, zeitlimit_s: int | None = None, nacharbeit: str | None = None, nacharbeit_daten: dict | None = None) -> str: with _LOCK: _laden() _aufraeumen() job = _eintragen(list(args), label, group, zeitlimit_s, nacharbeit, nacharbeit_daten) _starten(job, list(args), env) return job["id"] def start_job_exklusiv(group: str, args: list[str], label: str, env: dict | None = None, zeitlimit_s: int | None = None, nacharbeit: str | None = None, nacharbeit_daten: dict | None = None) -> tuple[str | None, dict | None]: """Wie start_job, aber nur, wenn in der Gruppe nichts läuft — Prüfen und Eintragen unter einer Sperre. Rückgabe: (neue Auftrags-ID, None) oder (None, der schon laufende Auftrag).""" with _LOCK: _laden() if (laeuft := _aktiv_in(group)) is not None: return None, laeuft _aufraeumen() job = _eintragen(list(args), label, group, zeitlimit_s, nacharbeit, nacharbeit_daten) _starten(job, list(args), env) return job["id"], None def active_in_group(group: str) -> dict | None: """Erster laufender/wartender Auftrag einer Gruppe (z. B. 'maintenance'), sonst None. Basis für den Wartungs-Riegel: nur EIN System-Update gleichzeitig.""" with _LOCK: _laden() return _aktiv_in(group) def wiederaufnehmen() -> int: """Beim Start von MC2: Akten einlesen und laufende Aufträge weiter beobachten. Rückgabe: Zahl der wieder aufgenommenen Aufträge.""" with _LOCK: _laden() offen = [j for j in JOBS.values() if j["state"] in ("queued", "running")] nacharbeit_offen = [j for j in JOBS.values() if j["state"] == "done" and j.get("nacharbeit") and not j.get("nacharbeit_erledigt")] for job in offen: if job.get("art") == "systemd": threading.Thread(target=_beobachten, args=(job["id"],), daemon=True).start() _fortschritt_beobachten(job["id"]) elif (code := _exit_code(job["id"])) is not None: _abschliessen(job, code) else: # Kindprozess des vorigen MC2: sein Warte-Thread ist mit ihm gegangen. _protokoll(job["id"], "[mc] MC2 wurde neu gestartet; der Auftrag lief als Kindprozess mit.") job["error"] = "MC2 wurde während des Auftrags neu gestartet." _abschliessen(job, None) for job in nacharbeit_offen: _nacharbeit_ausfuehren(job) if offen: log.info("jobengine: %d laufende Aufträge wieder aufgenommen", len(offen)) return len(offen) def cancel_job(job_id: str) -> bool: with _LOCK: _laden() job = JOBS.get(job_id) if not job or job["state"] in _ENDE: return False job["canceled"] = True _speichern(job) _protokoll(job_id, "[mc] Abbruch angefordert…") if job["state"] == "queued": _abschliessen(job, None) elif job.get("art") == "systemd": try: subprocess.run(_systemd("systemctl", "stop", "--no-block", _einheit(job_id)), capture_output=True, timeout=15) except Exception: log.warning("jobengine: Abbruch von %s gescheitert", job_id, exc_info=True) elif (proc := _PROCS.get(job_id)) is not None: _beende_gruppe(proc) else: _abschliessen(job, None) return True def public_jobs(mit_log: bool = True) -> list[dict]: """Aufträge ohne interne Felder für die API; das Protokoll kommt aus der Datei (Ende).""" with _LOCK: _laden() _aufraeumen() jobs = [{k: v for k, v in j.items() if k not in _INTERN} for j in JOBS.values()] for j in jobs: j["log"] = _protokoll_lesen(j["id"]) if mit_log else [] return sorted(jobs, key=lambda j: j.get("started_at") or 0)