Files
proxmox-snapshot-service/pvesnap/proxmox.py
T
duffyduckandClaude Opus 5 e26d966a8b Bei fehlgeschlagenen Tasks das Proxmox-Task-Log mitmelden
Bisher stand im Fehlerfall nur der Exit-Status da ("cfs-lock ... error: got
lock request timeout"), nicht aber, was Proxmox davor protokolliert hat.
Fuer die Fehlersuche fehlte damit genau der interessante Teil.

* Schlaegt ein Task fehl, werden die letzten Zeilen des Task-Logs an die
  Meldung angehaengt; identische Wiederholungen werden zusammengefasst
  ("trying to acquire cfs lock 'storage-data' ... (9x)").
* Die Meldung wiederholt nicht mehr VM und Snapshot-Namen, die der
  Aufrufer ohnehin voranstellt.

Ausserdem ein Fehler in der Wiederholung: schlug der Task fehl, blieb der
Snapshot-Eintrag teils in der VM-Konfiguration stehen. Der zweite Versuch
lief dann in "snapshot already exists" - und diese Meldung verdeckte die
eigentliche Ursache. Ist der Snapshot bereits vorhanden, wird jetzt nicht
wiederholt und der urspruengliche Fehler gemeldet.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-31 02:07:43 +02:00

315 lines
12 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 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)
def exists():
try:
return any(snap.name == name for snap in self.list_snapshots(guest))
except ProxmoxError:
return False
return self._retry(what, lambda: self._task(guest, args, what), skip_retry=exists)
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, 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):
for line in (output or "").splitlines():
candidate = line.strip().strip('"')
if candidate.startswith("UPID:"):
return candidate
return None