""" Job-Engine: Download eines macOS-Installers + Beschreiben eines USB-Sticks. Zwei Pfade, unterschiedlich robust: legacy_dmg (macOS bis Catalina/10.15): Apple liefert ein fertiges, bootfaehiges BaseSystem.dmg/InstallESD.dmg. Wir wandeln es mit dmg2img in ein rohes Image und schreiben es 1:1 mit dd auf den Stick. Etablierter, zuverlaessiger Weg. recovery (macOS ab Big Sur/11): Apple liefert nur noch ein InstallAssistant.pkg. Das eigentliche createinstallmedia laeuft nur unter echtem macOS. Wir extrahieren das .pkg (xar -> pbzx -> cpio) und suchen darin ein SharedSupport.dmg/BaseSystem.dmg, das wir dann genauso per dd schreiben. EXPERIMENTELL: Apple hat den internen Aufbau von InstallAssistant.pkg mehrfach leicht geaendert, das kann brechen. Klar als "experimentell" markiert. Sicherheits-Leitplanken: - Ziel-Device wird gegen die Root-Disk des Hosts geprueft (Refuse). - Ziel-Device muss "removable" oder USB-Transport sein (Refuse sonst). - Aufrufer muss den Device-Pfad im Request nochmal exakt bestaetigen. """ import logging import os import shutil import subprocess import tempfile import threading import time import uuid from queue import Queue, Empty import requests from devices import list_usb_sticks log = logging.getLogger("installer") DOWNLOAD_DIR = "/data/downloads" JOBS: dict[str, "Job"] = {} class JobAborted(Exception): pass class Job: def __init__(self, product: dict, device_path: str): self.id = str(uuid.uuid4()) self.product = product self.device_path = device_path self.status = "queued" # queued -> running -> done | failed | aborted self.phase = "warten" self.progress = 0 # 0-100 im aktuellen Phase self.error = None self._q: Queue = Queue() self._lock = threading.Lock() self.history = [] def emit(self, phase=None, progress=None, msg=None, status=None): with self._lock: if phase is not None: self.phase = phase if progress is not None: self.progress = progress if status is not None: self.status = status event = { "phase": self.phase, "progress": self.progress, "status": self.status, "msg": msg, "ts": time.time(), } self.history.append(event) self._q.put(event) def stream(self): # erst die bisherige Historie nachliefern, dann live weiter for event in list(self.history): yield event while True: try: yield self._q.get(timeout=30) except Empty: if self.status in ("done", "failed", "aborted"): return yield {"phase": self.phase, "progress": self.progress, "status": self.status, "msg": None, "ts": time.time()} def snapshot(self): with self._lock: return { "id": self.id, "status": self.status, "phase": self.phase, "progress": self.progress, "error": self.error, } def create_job(product: dict, device_path: str, confirm_path: str) -> Job: if confirm_path != device_path: raise ValueError("Bestaetigungs-Pfad stimmt nicht mit dem Ziel-Geraet ueberein.") sticks = list_usb_sticks() match = next((s for s in sticks if s["path"] == device_path), None) if not match: raise ValueError(f"Geraet {device_path} nicht in der aktuellen USB-Liste gefunden.") if match["is_root_disk"]: raise ValueError("Verweigert: das ist die System-Platte des Hosts.") if match["mounted"]: raise ValueError("Geraet ist aktuell gemountet -- bitte vorher aushaengen.") if not (match["transport"] == "usb"): raise ValueError("Geraet ist nicht als USB-Transport erkannt -- Sicherheitsstop.") job = Job(product, device_path) JOBS[job.id] = job t = threading.Thread(target=_run_job, args=(job,), daemon=True) t.start() return job def _run(job: Job, cmd: list[str], phase: str, progress_from_stderr=False): log.info("job %s: %s", job.id, " ".join(cmd)) proc = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, bufsize=1) for line in proc.stdout: line = line.strip() if line: job.emit(phase=phase, msg=line) ret = proc.wait() if ret != 0: raise RuntimeError(f"Befehl fehlgeschlagen ({ret}): {' '.join(cmd)}") def _download(job: Job, url: str, dest_path: str): job.emit(phase="download", progress=0, msg=f"Lade {url}") with requests.get(url, stream=True, timeout=30) as r: r.raise_for_status() total = int(r.headers.get("Content-Length", 0)) or None done = 0 last_pct = -1 with open(dest_path, "wb") as f: for chunk in r.iter_content(chunk_size=1024 * 1024 * 4): if not chunk: continue f.write(chunk) done += len(chunk) if total: pct = int(done * 100 / total) if pct != last_pct: job.emit(phase="download", progress=pct) last_pct = pct job.emit(phase="download", progress=100, msg="Download fertig") def _dmg_to_raw(job: Job, dmg_path: str) -> str: raw_path = dmg_path + ".img" job.emit(phase="konvertieren", progress=0, msg="dmg2img: DMG -> rohes Image") _run(job, ["dmg2img", "-i", dmg_path, "-o", raw_path], phase="konvertieren") return raw_path def _write_raw_to_device(job: Job, raw_path: str, device_path: str): job.emit(phase="schreiben", progress=0, msg=f"Schreibe {raw_path} -> {device_path} (dd)") _run(job, ["wipefs", "-a", device_path], phase="schreiben") size = os.path.getsize(raw_path) cmd = ["dd", f"if={raw_path}", f"of={device_path}", "bs=4M", "status=progress", "conv=fsync"] proc = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, bufsize=1) for line in proc.stdout: line = line.strip() if not line: continue job.emit(phase="schreiben", msg=line) # dd status=progress schreibt Zeilen wie "1234567890 bytes (1.2 GB, 1.1 GiB) copied, ..." try: bytes_done = int(line.split(" ")[0]) pct = min(100, int(bytes_done * 100 / size)) job.emit(phase="schreiben", progress=pct) except (ValueError, IndexError, ZeroDivisionError): pass ret = proc.wait() if ret != 0: raise RuntimeError("dd fehlgeschlagen") job.emit(phase="schreiben", progress=100, msg="dd fertig, sync...") subprocess.run(["sync"]) def _extract_recovery_dmg(job: Job, pkg_path: str, work_dir: str) -> str: """Experimentell: InstallAssistant.pkg -> xar -> pbzx -> cpio, sucht SharedSupport/BaseSystem.dmg.""" job.emit(phase="extrahieren", progress=0, msg="xar: InstallAssistant.pkg entpacken (experimentell)") extract_dir = os.path.join(work_dir, "xar") os.makedirs(extract_dir, exist_ok=True) _run(job, ["xar", "-xf", pkg_path, "-C", extract_dir], phase="extrahieren") payloads = [] for root, _dirs, files in os.walk(extract_dir): for fn in files: if fn == "Payload": payloads.append(os.path.join(root, fn)) if not payloads: raise RuntimeError("Kein Payload in InstallAssistant.pkg gefunden -- Apple-Struktur hat sich vermutlich geaendert.") cpio_dir = os.path.join(work_dir, "cpio") os.makedirs(cpio_dir, exist_ok=True) for payload in payloads: job.emit(phase="extrahieren", msg=f"pbzx+cpio: {payload}") # pbzx dekomprimiert das payload, cpio packt den entstandenen Stream aus pbzx_proc = subprocess.Popen(["pbzx", "-n", payload], stdout=subprocess.PIPE) cpio_proc = subprocess.Popen( ["cpio", "-idm"], cwd=cpio_dir, stdin=pbzx_proc.stdout, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, ) pbzx_proc.stdout.close() for line in cpio_proc.stdout: line = line.strip() if line: job.emit(phase="extrahieren", msg=line) cpio_proc.wait() pbzx_proc.wait() dmg_candidates = [] for root, _dirs, files in os.walk(cpio_dir): for fn in files: if fn in ("SharedSupport.dmg", "BaseSystem.dmg"): dmg_candidates.append(os.path.join(root, fn)) if not dmg_candidates: raise RuntimeError( "Kein SharedSupport.dmg/BaseSystem.dmg im entpackten Installer gefunden. " "Das ist der bekannte fragile Punkt bei macOS Big Sur+ ohne echten Mac." ) # SharedSupport.dmg bevorzugen (enthaelt das vollstaendige Recovery-System) dmg_candidates.sort(key=lambda p: 0 if "SharedSupport" in p else 1) return dmg_candidates[0] def _run_job(job: Job): os.makedirs(DOWNLOAD_DIR, exist_ok=True) work_dir = tempfile.mkdtemp(dir=DOWNLOAD_DIR) try: job.emit(status="running", phase="start", progress=0, msg="Job gestartet") product = job.product mode = product["mode"] if mode == "legacy_dmg": url = product["base_system_url"] dmg_path = os.path.join(work_dir, "BaseSystem.dmg") _download(job, url, dmg_path) raw_path = _dmg_to_raw(job, dmg_path) _write_raw_to_device(job, raw_path, job.device_path) else: job.emit(phase="hinweis", msg="Big Sur+ Pfad ist experimentell (siehe README).") url = product["install_assistant_url"] pkg_path = os.path.join(work_dir, "InstallAssistant.pkg") _download(job, url, pkg_path) dmg_path = _extract_recovery_dmg(job, pkg_path, work_dir) raw_path = _dmg_to_raw(job, dmg_path) _write_raw_to_device(job, raw_path, job.device_path) job.emit(status="done", phase="fertig", progress=100, msg="Stick ist fertig.") except Exception as exc: # noqa: BLE001 log.exception("job %s failed", job.id) job.error = str(exc) job.emit(status="failed", msg=f"FEHLER: {exc}") finally: shutil.rmtree(work_dir, ignore_errors=True)