"""Duenne Schicht um `pvesh` - Inventar lesen, Snapshots anlegen/loeschen. Bewusst ueber `pvesh` statt `qm`/`pct`: damit funktioniert alles auch fuer Gaeste, die auf einem anderen Node des Clusters laufen. """ from __future__ import annotations import json import logging import re import shutil import subprocess import time from dataclasses import dataclass, field log = logging.getLogger("pvesnap.proxmox") PVESH = "/usr/bin/pvesh" # UPID:node:PID:PSTART:STARTTIME:TYPE:ID:USER: _UPID_PATTERN = re.compile(r"UPID:[^\s\"']+") class ProxmoxError(Exception): """Ein pvesh-Aufruf ist fehlgeschlagen.""" @dataclass class Guest: vmid: int name: str = "" type: str = "qemu" # qemu | lxc node: str = "" status: str = "unknown" tags: list = field(default_factory=list) pool: str = "" @property def running(self): return self.status == "running" @property def label(self): return "%s %d (%s)" % ("LXC" if self.type == "lxc" else "VM", self.vmid, self.name or "?") @dataclass class Snapshot: name: str description: str = "" snaptime: int = 0 parent: str = "" # "prepare" oder "delete" = halbfertiger Eintrag in der Gast-Konfiguration. # Dahinter steckt kein brauchbarer Snapshot auf dem Storage. snapstate: str = "" @property def complete(self): return not self.snapstate def pvesh_available(): return shutil.which(PVESH) is not None or shutil.which("pvesh") is not None def _pvesh_binary(): return shutil.which(PVESH) or shutil.which("pvesh") or PVESH class Proxmox: """Zugriff auf die Proxmox-API ueber das Kommando `pvesh`.""" def __init__(self, dry_run=False, task_timeout=900, command_timeout=60, retries=2, retry_delay=60): self.dry_run = dry_run self.task_timeout = task_timeout self.command_timeout = command_timeout self.retries = retries self.retry_delay = retry_delay self._inventory_cache = None @classmethod def from_config(cls, config): globals_ = config.globals return cls(dry_run=globals_.dry_run, task_timeout=globals_.task_timeout, retries=globals_.retries, retry_delay=globals_.retry_delay) # -- unterste Ebene --------------------------------------------------- def _run(self, args, timeout=None): command = [_pvesh_binary()] + args log.debug("pvesh: %s", " ".join(command)) try: proc = subprocess.run(command, stdout=subprocess.PIPE, stderr=subprocess.PIPE, timeout=timeout or self.command_timeout) except FileNotFoundError: raise ProxmoxError("'pvesh' nicht gefunden - laeuft pvesnap wirklich " "auf einem Proxmox-VE-Host?") except subprocess.TimeoutExpired: raise ProxmoxError("Zeitueberschreitung bei: %s" % " ".join(command)) stdout = proc.stdout.decode("utf-8", "replace").strip() stderr = proc.stderr.decode("utf-8", "replace").strip() if proc.returncode != 0: raise ProxmoxError((stderr or stdout or "unbekannter Fehler").splitlines()[0]) return stdout, stderr def _json(self, args): stdout, _ = self._run(args + ["--output-format", "json"]) if not stdout: return None try: return json.loads(stdout) except ValueError: raise ProxmoxError("unerwartete Antwort von pvesh: %s" % stdout[:200]) # -- Inventar --------------------------------------------------------- def inventory(self, refresh=False): """Alle VMs und Container des Clusters.""" if self._inventory_cache is not None and not refresh: return self._inventory_cache data = self._json(["get", "/cluster/resources", "--type", "vm"]) or [] guests = [] for entry in data: if entry.get("type") not in ("qemu", "lxc"): continue raw_tags = entry.get("tags") or "" tags = [t.strip().lower() for t in str(raw_tags).replace(",", ";").split(";") if t.strip()] guests.append(Guest( vmid=int(entry.get("vmid")), name=entry.get("name") or "", type=entry.get("type"), node=entry.get("node") or "", status=entry.get("status") or "unknown", tags=tags, pool=entry.get("pool") or "", )) guests.sort(key=lambda g: g.vmid) self._inventory_cache = guests return guests @staticmethod def local_node(): import socket return socket.gethostname().split(".")[0] # -- Snapshots -------------------------------------------------------- def _base_path(self, guest): return "/nodes/%s/%s/%d/snapshot" % (guest.node, guest.type, guest.vmid) def list_snapshots(self, guest): data = self._json(["get", self._base_path(guest)]) or [] result = [] for entry in data: name = entry.get("name") if not name or name == "current": continue # "current" ist kein echter Snapshot result.append(Snapshot( name=name, description=(entry.get("description") or "").strip(), snaptime=int(entry.get("snaptime") or 0), parent=entry.get("parent") or "", snapstate=(entry.get("snapstate") or "").strip(), )) result.sort(key=lambda s: s.snaptime) return result def create_snapshot(self, guest, name, description="", vmstate=False): args = ["create", self._base_path(guest), "--snapname", name] if description: args += ["--description", description] if vmstate and guest.type == "qemu": args += ["--vmstate", "1"] if self.dry_run: log.info("[TESTLAUF] wuerde Snapshot anlegen: %s -> %s", guest.label, name) return None what = "Snapshot %s fuer %s" % (name, guest.label) def exists(): try: return any(snap.name == name for snap in self.list_snapshots(guest)) except ProxmoxError: return False def attempt(): result = self._task(guest, args, what) self._verify_created(guest, name) return result return self._retry(what, attempt, skip_retry=exists) def delete_snapshot(self, guest, name, force=False): if self.dry_run: log.info("[TESTLAUF] wuerde Snapshot loeschen: %s -> %s", guest.label, name) return None args = ["delete", "%s/%s" % (self._base_path(guest), name)] if force: # Nur fuer halbfertige Eintraege: den Eintrag auch dann aus der # Konfiguration nehmen, wenn auf dem Storage nichts (mehr) liegt. args += ["--force", "1"] what = "Loeschen von %s bei %s" % (name, guest.label) return self._retry(what, lambda: self._task(guest, args, what)) # -- Gaeste lesen und steuern ----------------------------------------- def guest_config(self, guest, snapshot=None): """Konfiguration eines Gastes - wahlweise die zu einem Snapshot.""" tail = ("snapshot/%s/config" % snapshot) if snapshot else "config" data = self._json(["get", "/nodes/%s/%s/%d/%s" % (guest.node, guest.type, guest.vmid, tail)]) if not isinstance(data, dict): raise ProxmoxError("Konfiguration von %s nicht lesbar" % guest.label) return data def guest_status(self, guest): try: data = self._json(["get", "/nodes/%s/%s/%d/status/current" % (guest.node, guest.type, guest.vmid)]) except ProxmoxError: return {} return data if isinstance(data, dict) else {} def guest_exists(self, vmid): return any(g.vmid == int(vmid) for g in self.inventory(refresh=True)) def next_vmid(self, wanted=None): """Freie VMID. Mit `wanted` wird geprueft, ob genau die frei ist.""" args = ["get", "/cluster/nextid"] if wanted: args += ["--vmid", str(int(wanted))] return int(self._json(args)) def vmid_free(self, vmid): try: self.next_vmid(vmid) return True except (ProxmoxError, TypeError, ValueError): return False def start_guest(self, guest): return self._guest_task(guest, "create", "status/start", "Starten von %s" % guest.label) def shutdown_guest(self, guest, timeout=180, force=True): """Sauberes Herunterfahren; `force` schaltet nach Ablauf hart ab.""" extra = ["--timeout", str(int(timeout))] if force: extra += ["--forceStop", "1"] return self._guest_task(guest, "create", "status/shutdown", "Herunterfahren von %s" % guest.label, extra, timeout=timeout + 60) def stop_guest(self, guest): return self._guest_task(guest, "create", "status/stop", "Ausschalten von %s" % guest.label) def set_guest_config(self, guest, values): """Konfigurationswerte setzen. Laeuft der Gast, steckt Proxmox hotplug-faehige Geraete dabei gleich an.""" args = ["set", "/nodes/%s/%s/%d/config" % (guest.node, guest.type, guest.vmid)] for key, value in values.items(): args += ["--%s" % key, str(value)] if self.dry_run: log.info("[TESTLAUF] wuerde setzen: %s", " ".join(args)) return None stdout, _ = self._run(args, timeout=self.task_timeout) return self._wait_upid(guest.node, stdout, "Aendern von %s" % guest.label, guest) def resume_guest(self, guest): return self._guest_task(guest, "create", "status/resume", "Fortsetzen von %s" % guest.label) def task_log(self, node, upid, limit=200): """Das Protokoll eines Tasks als Liste von Zeilen.""" try: entries = self._json(["get", "/nodes/%s/tasks/%s/log" % (node, upid), "--limit", str(int(limit))]) or [] except ProxmoxError: return [] return [str(entry.get("t") or "") for entry in entries] def destroy_guest(self, guest, purge=True): extra = ["--purge", "1"] if purge else [] if guest.type == "qemu": extra += ["--destroy-unreferenced-disks", "1"] return self._guest_task(guest, "delete", "", "Entfernen von %s" % guest.label, extra) def _guest_task(self, guest, verb, sub, what, extra=None, timeout=None): path = "/nodes/%s/%s/%d" % (guest.node, guest.type, guest.vmid) if sub: path += "/" + sub args = [verb, path] + (extra or []) if self.dry_run: log.info("[TESTLAUF] wuerde ausfuehren: %s", " ".join(args)) return None stdout, _ = self._run(args, timeout=timeout or self.task_timeout) return self._wait_upid(guest.node, stdout, what, guest) def _task(self, guest, args, what): """Fuehrt eine Snapshot-Aktion aus und wartet, bis sie wirklich fertig ist. Wichtig: hier gilt das lange Task-Zeitlimit, nicht das kurze fuer Lesezugriffe. Proxmox wartet beim Anlegen bis zu 60s auf den Storage-Lock; liefen wir vorher in unser eigenes Zeitlimit, wuerde die naechste VM losgeschickt, waehrend der Task noch laeuft - und genau das laesst dann alle folgenden Snapshots am cfs-Lock scheitern. """ stdout, _ = self._run(args, timeout=self.task_timeout) return self._wait_task(guest, stdout, what) # -- Wiederholung bei belegten Sperren -------------------------------- # Solche Meldungen sind voruebergehend - da lohnt ein zweiter Versuch. _TRANSIENT = ("cfs-lock", "lock request timeout", "trying to acquire", "can't lock file", "got lock timeout", "unable to acquire lock", "resource temporarily unavailable", "storage is locked", "vm is locked", "ct is locked") @classmethod def _is_transient(cls, message): text = str(message).lower() return any(marker in text for marker in cls._TRANSIENT) def _retry(self, what, action, skip_retry=None): attempt = 0 while True: try: return action() except ProxmoxError as exc: if attempt >= self.retries or not self._is_transient(exc): raise if skip_retry is not None and skip_retry(): # Der Snapshot ist trotz Fehler schon da - ein zweiter # Versuch scheiterte nur an "already exists" und wuerde # die eigentliche Ursache verdecken. log.warning("%s: Snapshot ist trotz Fehler vorhanden - " "keine Wiederholung", what) raise attempt += 1 log.warning("%s: %s - Versuch %d von %d in %ds", what, exc, attempt + 1, self.retries + 1, self.retry_delay) time.sleep(self.retry_delay) # -- Task-Verfolgung -------------------------------------------------- def _wait_task(self, guest, output, what): return self._wait_upid(guest.node, output, what, guest) def _wait_upid(self, node, output, what, guest=None): """pvesh liefert bei laengeren Aktionen eine UPID; darauf warten wir.""" upid = self._extract_upid(output) if not upid: # Ohne UPID koennen wir den Task nicht verfolgen - dann warten wir # ersatzweise, bis der Gast nicht mehr gesperrt ist. Sonst wuerde # die naechste Aktion in eine noch laufende hineinlaufen. log.debug("%s: keine UPID in der Antwort von pvesh", what) if guest is not None: self._wait_unlocked(guest) return None deadline = time.time() + self.task_timeout delay = 0.5 while True: status = self._json(["get", "/nodes/%s/tasks/%s/status" % (node, upid)]) or {} if status.get("status") == "stopped": exit_status = status.get("exitstatus") or "unbekannt" if exit_status != "OK": # Nur die Ursache melden - wer und was, ergaenzt der Aufrufer. raise ProxmoxError("%s%s" % (exit_status, self._task_log_tail(node, upid))) return upid if time.time() > deadline: raise ProxmoxError("%s: Zeitueberschreitung nach %ds (Task laeuft weiter: %s)" % (what, self.task_timeout, upid)) time.sleep(delay) delay = min(delay * 1.5, 5.0) def _task_log_tail(self, node, upid, lines=4): """Die letzten Zeilen des Task-Logs - dort steht, woran es lag. Wiederholungen wie "trying to acquire cfs lock ..." werden zusammengefasst, sonst besteht die Meldung nur noch daraus. """ try: entries = self._json(["get", "/nodes/%s/tasks/%s/log" % (node, upid), "--limit", "100"]) or [] except ProxmoxError: return "" texts = [] for entry in entries: text = str(entry.get("t") or "").strip() if not text: continue if texts and texts[-1][0] == text: texts[-1][1] += 1 else: texts.append([text, 1]) tail = ["%s%s" % (text, " (%dx)" % count if count > 1 else "") for text, count in texts[-lines:]] return " [Task-Log: %s]" % " | ".join(tail) if tail else "" def _wait_unlocked(self, guest, timeout=None): """Wartet, bis der Gast keine laufende Sperre mehr hat.""" deadline = time.time() + (timeout or self.task_timeout) path = "/nodes/%s/%s/%d/status/current" % (guest.node, guest.type, guest.vmid) while True: try: status = self._json(["get", path]) or {} except ProxmoxError: return # lieber weitermachen als haengen bleiben if not status.get("lock"): return if time.time() > deadline: raise ProxmoxError("%s ist seit %ds gesperrt (%s)" % (guest.label, timeout or self.task_timeout, status.get("lock"))) time.sleep(2.0) @staticmethod def _extract_upid(output): """UPID aus der Antwort von pvesh fischen. Bei Containern schreibt pvesh den Fortschritt vor die UPID, alles in einer Zeile: Creating snap: 100% complete...done."UPID:node:...:vzsnapshot:100:..." Zeilenweise auf den Anfang zu pruefen findet die UPID dort nicht - und ohne UPID bliebe ein fehlgeschlagener Task unbemerkt. """ match = _UPID_PATTERN.search(output or "") return match.group(0) if match else None def _verify_created(self, guest, name): """Nachsehen, ob der Snapshot wirklich brauchbar entstanden ist. Zweite Sicherung fuer den Fall, dass ein Fehlschlag weder ueber den Rueckgabewert noch ueber den Task sichtbar wird: Proxmox laesst dann einen halbfertigen Eintrag ("snapstate") in der Konfiguration stehen. """ try: snapshots = self.list_snapshots(guest) except ProxmoxError: return # nicht schlimmer machen als noetig found = next((s for s in snapshots if s.name == name), None) if found is None: raise ProxmoxError("Der Snapshot '%s' ist nach dem Anlegen nicht " "vorhanden." % name) if not found.complete: raise ProxmoxError("Der Snapshot '%s' ist unvollstaendig geblieben " "(snapstate: %s) - er enthaelt keine Daten." % (name, found.snapstate))