Bei Containern schreibt pvesh den Fortschritt vor die UPID, alles in eine
Zeile:
Creating snap: 100% complete...done."UPID:node:...:vzsnapshot:100:..."
Die UPID wurde zeilenweise mit startswith() gesucht und deshalb nicht
gefunden. Ohne UPID liess sich der Task nicht verfolgen - ein Fehlschlag
blieb damit unbemerkt und pvesnap protokollierte "Snapshot angelegt",
obwohl Proxmox nur einen halbfertigen Eintrag (snapstate "prepare") in der
Konfiguration hinterlassen hatte. Genau so sind die beiden Leichen bei
LXC 100 entstanden.
* Die UPID wird jetzt per Suchmuster aus der gesamten Ausgabe geholt,
egal woran sie haengt.
* Zusaetzlich wird nach dem Anlegen nachgesehen, ob der Snapshot wirklich
da und vollstaendig ist. Damit faellt ein stiller Fehlschlag auch dann
auf, wenn weder Rueckgabewert noch Task ihn melden.
Auf pvetest01 mit Container 100 geprueft: anlegen (Task wird verfolgt,
Eintrag vollstaendig, rbd-Snapshot vorhanden), einbinden (rootfs als ext4,
/etc/hostname liefert "unifi"), Explorer und Web-Oberflaeche inklusive
ZIP, sowie loeschen ueber die Vorhaltezahl.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
363 lines
14 KiB
Python
363 lines
14 KiB
Python
"""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))
|
|
|
|
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):
|
|
"""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":
|
|
# Nur die Ursache melden - wer und was, ergaenzt der Aufrufer.
|
|
raise ProxmoxError("%s%s" % (exit_status,
|
|
self._task_log_tail(guest.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))
|