533 lines
21 KiB
Python
533 lines
21 KiB
Python
"""
|
|
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-<id>“ (systemd-run), nicht mehr als
|
|
Kindprozess von MC2. Ein Deploy oder Absturz von MC2 würgt ihn nicht mehr ab.
|
|
- Akte (<id>.json), Protokoll (<id>.log) und Exit-Code (<id>.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 </dev/null; '
|
|
'"$@"; code=$?; echo "$code" > "$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 and not k.startswith("_")}
|
|
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). Bei Erfolg läuft die Nacharbeit VOR dem Endzustand:
|
|
Wer „done“ sieht, sieht auch schon ihre Wirkung (gesetzte Rolle, geleerte Zwischenspeicher)."""
|
|
with _LOCK:
|
|
if job["state"] in _ENDE or job.get("_schliesst"):
|
|
return
|
|
job["_schliesst"] = True
|
|
felder: dict = {"returncode": code}
|
|
if job.get("canceled"):
|
|
felder["state"] = "canceled"
|
|
elif job.get("zeitlimit") or (code is None and time.time() - job["started_at"] >= job["zeitlimit_s"] - 5):
|
|
felder.update(state="failed", zeitlimit=True, error="Zeitlimit überschritten")
|
|
elif code is None:
|
|
felder.update(state="failed", error=job.get("error") or "Der Auftrag endete ohne Rückmeldung.")
|
|
else:
|
|
felder["state"] = "done" if code == 0 else "failed"
|
|
if felder["state"] == "done":
|
|
_nacharbeit_ausfuehren(job)
|
|
with _LOCK:
|
|
job.update(felder, finished_at=time.time())
|
|
job.pop("_schliesst", None)
|
|
_PROCS.pop(job["id"], None)
|
|
_speichern(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 and not k.startswith("_")} 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)
|