"""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): self.dry_run = dry_run self.task_timeout = task_timeout self.command_timeout = command_timeout self._inventory_cache = None # -- 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 stdout, _ = self._run(args) return self._wait_task(guest.node, stdout, "Snapshot %s fuer %s" % (name, guest.label)) def delete_snapshot(self, guest, name): if self.dry_run: log.info("[TESTLAUF] wuerde Snapshot loeschen: %s -> %s", guest.label, name) return None stdout, _ = self._run(["delete", "%s/%s" % (self._base_path(guest), name)]) return self._wait_task(guest.node, stdout, "Loeschen von %s bei %s" % (name, guest.label)) # -- Task-Verfolgung -------------------------------------------------- def _wait_task(self, node, output, what): """pvesh liefert bei Snapshot-Aktionen eine UPID; darauf warten wir.""" upid = self._extract_upid(output) if not upid: 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": 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) @staticmethod def _extract_upid(output): for line in (output or "").splitlines(): candidate = line.strip().strip('"') if candidate.startswith("UPID:"): return candidate return None