"""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 shutil import subprocess import time from dataclasses import dataclass, field log = logging.getLogger("pvesnap.proxmox") PVESH = "/usr/bin/pvesh" 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 = "" 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 "", )) 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) return self._retry(what, lambda: self._task(guest, args, what)) def delete_snapshot(self, guest, name): 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)] what = "Loeschen von %s bei %s" % (name, guest.label) return self._retry(what, lambda: self._task(guest, args, what)) 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): attempt = 0 while True: try: return action() except ProxmoxError as exc: if attempt >= self.retries or not self._is_transient(exc): 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): """pvesh liefert bei Snapshot-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) 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" % (guest.node, upid)]) or {} if status.get("status") == "stopped": exit_status = status.get("exitstatus") or "unbekannt" if exit_status != "OK": raise ProxmoxError("%s fehlgeschlagen: %s" % (what, exit_status)) 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 _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): for line in (output or "").splitlines(): candidate = line.strip().strip('"') if candidate.startswith("UPID:"): return candidate return None