first commit
This commit is contained in:
@@ -0,0 +1,232 @@
|
||||
"""Der eigentliche Dienst: Zeitplan ueberwachen und Gruppen ausfuehren."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import errno
|
||||
import fcntl
|
||||
import logging
|
||||
import os
|
||||
import signal
|
||||
import threading
|
||||
import time
|
||||
from datetime import datetime
|
||||
|
||||
from .config import ConfigError, load_config
|
||||
from .engine import run_group
|
||||
from .proxmox import Proxmox, ProxmoxError
|
||||
from .schedule import next_due
|
||||
from .state import State
|
||||
|
||||
log = logging.getLogger("pvesnap.daemon")
|
||||
|
||||
|
||||
class SingleInstanceLock:
|
||||
"""Dateisperre.
|
||||
|
||||
Zwei Sperren sind im Spiel:
|
||||
* <lock_file>.daemon haelt der Dienst waehrend seiner gesamten Laufzeit,
|
||||
damit er nicht doppelt startet.
|
||||
* <lock_file> wird nur waehrend der eigentlichen Snapshot-Arbeit
|
||||
gehalten - so kann man jederzeit auch von Hand `pvesnap run` aufrufen,
|
||||
ohne dass sich die Laeufe in die Quere kommen.
|
||||
"""
|
||||
|
||||
def __init__(self, path, wait=0, busy_message=None):
|
||||
self.path = path
|
||||
self.wait = wait
|
||||
self.busy_message = busy_message or "pvesnap laeuft bereits (Sperrdatei %s)" % path
|
||||
self._handle = None
|
||||
|
||||
def __enter__(self):
|
||||
# Die Datei anzulegen muss sofort klappen - fehlende Rechte sind ein
|
||||
# harter Fehler und kein "gerade belegt".
|
||||
try:
|
||||
os.makedirs(os.path.dirname(os.path.abspath(self.path)), exist_ok=True)
|
||||
self._handle = open(self.path, "w")
|
||||
except OSError as exc:
|
||||
self._handle = None
|
||||
raise RuntimeError("Sperrdatei %s nicht nutzbar: %s" % (self.path, exc))
|
||||
|
||||
deadline = time.monotonic() + max(0, self.wait)
|
||||
while True:
|
||||
try:
|
||||
fcntl.flock(self._handle, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||||
self._handle.write("%d\n" % os.getpid())
|
||||
self._handle.flush()
|
||||
return self
|
||||
except OSError as exc:
|
||||
busy = exc.errno in (errno.EAGAIN, errno.EWOULDBLOCK, errno.EACCES)
|
||||
if not busy or time.monotonic() >= deadline:
|
||||
self._handle.close()
|
||||
self._handle = None
|
||||
if busy:
|
||||
raise RuntimeError(self.busy_message)
|
||||
raise RuntimeError("Sperrdatei %s nicht nutzbar: %s" % (self.path, exc))
|
||||
time.sleep(1.0)
|
||||
|
||||
def __exit__(self, *_exc):
|
||||
if self._handle:
|
||||
try:
|
||||
fcntl.flock(self._handle, fcntl.LOCK_UN)
|
||||
except OSError:
|
||||
pass
|
||||
self._handle.close()
|
||||
self._handle = None
|
||||
return False
|
||||
|
||||
|
||||
class Daemon:
|
||||
def __init__(self, config_path, dry_run=False):
|
||||
self.config_path = config_path
|
||||
self.dry_run_override = dry_run
|
||||
self.config = None
|
||||
self.state = None
|
||||
self._stop = threading.Event()
|
||||
self._reload = threading.Event()
|
||||
self._wake = threading.Event() # bricht das Warten bei Signalen ab
|
||||
|
||||
# -- Signale ----------------------------------------------------------
|
||||
|
||||
def install_signal_handlers(self):
|
||||
signal.signal(signal.SIGTERM, self._on_stop)
|
||||
signal.signal(signal.SIGINT, self._on_stop)
|
||||
signal.signal(signal.SIGHUP, self._on_reload)
|
||||
|
||||
def _on_stop(self, signum, _frame):
|
||||
log.info("Signal %s empfangen - beende nach dem aktuellen Durchlauf", signum)
|
||||
self._stop.set()
|
||||
self._wake.set()
|
||||
|
||||
def _on_reload(self, _signum, _frame):
|
||||
log.info("SIGHUP empfangen - Konfiguration wird neu geladen")
|
||||
self._reload.set()
|
||||
self._wake.set()
|
||||
|
||||
# -- Konfiguration ----------------------------------------------------
|
||||
|
||||
def load(self):
|
||||
config = load_config(self.config_path)
|
||||
problems = config.validate()
|
||||
if problems:
|
||||
raise ConfigError("Konfiguration fehlerhaft:\n - %s" % "\n - ".join(problems))
|
||||
if self.dry_run_override:
|
||||
config.globals.dry_run = True
|
||||
self.config = config
|
||||
if self.state is None or self.state.path != config.globals.state_file:
|
||||
self.state = State(config.globals.state_file).load()
|
||||
return config
|
||||
|
||||
def _reload_if_requested(self):
|
||||
if not self._reload.is_set():
|
||||
return
|
||||
self._reload.clear()
|
||||
try:
|
||||
self.load()
|
||||
log.info("Konfiguration neu geladen: %d Gruppe(n)", len(self.config.groups))
|
||||
except ConfigError as exc:
|
||||
log.error("Neuladen fehlgeschlagen, behalte alte Konfiguration: %s", exc)
|
||||
|
||||
# -- Hauptschleife ----------------------------------------------------
|
||||
|
||||
def run(self):
|
||||
self.load()
|
||||
self.install_signal_handlers()
|
||||
|
||||
globals_ = self.config.globals
|
||||
log.info("pvesnap gestartet (%d Gruppe(n), Praefix '%s'%s)",
|
||||
len(self.config.groups), globals_.prefix,
|
||||
", TESTLAUF" if globals_.dry_run else "")
|
||||
|
||||
daemon_lock = SingleInstanceLock(
|
||||
globals_.lock_file + ".daemon",
|
||||
busy_message="Der pvesnap-Dienst laeuft bereits (%s.daemon)" % globals_.lock_file)
|
||||
with daemon_lock:
|
||||
while not self._stop.is_set():
|
||||
self._wake.clear()
|
||||
self._reload_if_requested()
|
||||
try:
|
||||
self.tick()
|
||||
except ProxmoxError as exc:
|
||||
log.error("Proxmox nicht erreichbar: %s", exc)
|
||||
except Exception: # pragma: no cover - Dienst darf nie sterben
|
||||
log.exception("Unerwarteter Fehler im Durchlauf")
|
||||
self._sleep()
|
||||
|
||||
if self.state:
|
||||
self.state.save()
|
||||
log.info("pvesnap beendet")
|
||||
|
||||
def tick(self, now=None):
|
||||
"""Einmal alle Gruppen pruefen und faellige ausfuehren."""
|
||||
now = now or datetime.now()
|
||||
config = self.config
|
||||
state = self.state
|
||||
state.prune_unknown([g.name for g in config.groups])
|
||||
|
||||
due = []
|
||||
for group in config.groups:
|
||||
if not group.enabled:
|
||||
continue
|
||||
last = state.last_run(group.name)
|
||||
if last is None and not config.globals.run_on_start:
|
||||
# Erster Start: nicht sofort feuern, sondern auf den naechsten
|
||||
# regulaeren Termin warten.
|
||||
state.record_run(group.name, now, created=0, deleted=0)
|
||||
log.info("Gruppe '%s': erster Termin am %s", group.name,
|
||||
next_due(group, now, now).strftime("%Y-%m-%d %H:%M:%S"))
|
||||
continue
|
||||
if last is None or next_due(group, last, now) <= now:
|
||||
due.append(group)
|
||||
|
||||
if not due:
|
||||
state.save()
|
||||
return []
|
||||
|
||||
proxmox = Proxmox(dry_run=config.globals.dry_run,
|
||||
task_timeout=config.globals.task_timeout)
|
||||
|
||||
results = []
|
||||
# Die Arbeitssperre wird nur waehrend des Laufs gehalten, damit
|
||||
# `pvesnap run` von Hand weiterhin moeglich bleibt.
|
||||
work_lock = SingleInstanceLock(
|
||||
config.globals.lock_file, wait=300,
|
||||
busy_message="Ein anderer pvesnap-Lauf ist noch aktiv - "
|
||||
"dieser Durchlauf wird uebersprungen")
|
||||
try:
|
||||
with work_lock:
|
||||
guests = proxmox.inventory(refresh=True)
|
||||
for group in due:
|
||||
if self._stop.is_set():
|
||||
break
|
||||
result = run_group(proxmox, config, group, now=datetime.now(),
|
||||
guests=guests)
|
||||
state.record_run(group.name, datetime.now(),
|
||||
created=len(result.created),
|
||||
deleted=len(result.deleted), errors=result.errors)
|
||||
results.append(result)
|
||||
except RuntimeError as exc:
|
||||
log.warning("%s", exc)
|
||||
|
||||
state.save()
|
||||
return results
|
||||
|
||||
def _sleep(self):
|
||||
"""Bis zum naechsten Termin schlafen, hoechstens aber check_interval."""
|
||||
interval = self.config.globals.check_interval
|
||||
now = datetime.now()
|
||||
wait = float(interval)
|
||||
for group in self.config.groups:
|
||||
if not group.enabled:
|
||||
continue
|
||||
last = self.state.last_run(group.name)
|
||||
if last is None:
|
||||
continue
|
||||
try:
|
||||
remaining = (next_due(group, last, now) - now).total_seconds()
|
||||
except ValueError:
|
||||
continue
|
||||
wait = min(wait, max(1.0, remaining))
|
||||
# Auf Signale reagieren wir sofort - das Event bricht das Warten ab.
|
||||
# Geleert wird es am Anfang des Schleifendurchlaufs, damit ein Signal
|
||||
# waehrend tick() nicht verloren geht.
|
||||
self._wake.wait(max(1.0, min(wait, float(interval))))
|
||||
Reference in New Issue
Block a user