syncova-backup/packages/scheduler/queue.go
Jerrit Fritzsche 610719c316
Some checks failed
CI / Backend (Go) (push) Failing after 3m7s
CI / Frontend (React/TypeScript) (push) Successful in 37s
CI / Sicherheitsprüfungen (push) Successful in 44s
Syncova Backups V1
Enterprise-Backup-, Recovery-, Verification-, Security- und
Monitoring-Plattform fuer Proxmox VE, Windows, Linux und Dateisysteme.

Der Leitsatz, der fast jede Entscheidung erklaert: Ein Backup gilt erst als
vertrauenswuerdig, wenn Integritaet geprueft und Wiederherstellbarkeit
nachgewiesen wurde. Deshalb steigt ein Wiederherstellungspunkt erst nach einem
tatsaechlich durchgefuehrten Restore-Test auf "recoverable", und Unbekanntes
geht in keine Bewertung als "gut" ein.

Umfang (Phasen 0-23):

- Repository Engine: inhaltsadressierte Bloecke, atomares Commit-Protokoll,
  Katalogaufbau allein aus den Manifesten — ohne Datenbank
- Backup Engine: inhaltsabhaengiges Chunking, Deduplizierung trotz
  Verschluesselung, zstd, AES-256-GCM, Streaming mit Gegendruck
- Agenten fuer Windows und Linux mit Auftragsabholung (Pull-Modell)
- Proxmox-Provider mit beiden Zugriffswegen auf die Sicherungsarchive
- Scheduler, Recovery Engine mit Pruefpunkt, Verification, Unveraenderlichkeit
- Weboberflaeche, Kennzahlen, Meldungen, Berichte, Security Center,
  Ransomware-Heuristik (meldet, handelt nie)
- Disaster Recovery, Haertung, Leistungsmessung, Chaos Testing
- Eingefrorene Vertraege fuer API, Migrationen, Backup-Format und Repository
- Auslieferungspaket fuer linux/amd64, linux/arm64 und windows/amd64

Nicht enthalten und als solches gekennzeichnet: Kapazitaetsprognose, Backup
Copy, Changed Block Tracking bei Proxmox, erweiterte Attribute und ACLs.

Gebaut, aber nie auf echter Hardware gefahren: der Windows-Dienst, die
systemd-Einheit und der verpflichtende Proxmox-Meilenstein — ob eine
wiederhergestellte VM startet, ist ungeprueft. Einzelheiten in CHANGELOG.md
und docs/release-candidate.md.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-17 09:10:54 +02:00

514 lines
16 KiB
Go

package scheduler
import (
"errors"
"fmt"
"sort"
"sync"
"time"
)
// Priority ist die Dringlichkeit eines Auftrags.
type Priority string
const (
// PriorityCritical ist die höchste Stufe.
PriorityCritical Priority = "critical"
// PriorityHigh ist hohe Dringlichkeit.
PriorityHigh Priority = "high"
// PriorityNormal ist der Standard.
PriorityNormal Priority = "normal"
// PriorityLow ist niedrige Dringlichkeit.
PriorityLow Priority = "low"
)
// priorityRanks bilden die Stufen auf vergleichbare Zahlen ab.
var priorityRanks = map[Priority]int{
PriorityCritical: 400,
PriorityHigh: 300,
PriorityNormal: 200,
PriorityLow: 100,
}
// Rank liefert den Zahlenwert einer Stufe.
func (priority Priority) Rank() int {
if rankValue, isKnown := priorityRanks[priority]; isKnown {
return rankValue
}
// Eine unbekannte Stufe gilt als normal statt als höchste. Ein Tippfehler
// soll einen Auftrag nicht an die Spitze der Warteschlange befördern.
return priorityRanks[PriorityNormal]
}
// IsValid meldet eine bekannte Stufe.
func (priority Priority) IsValid() bool {
_, isKnown := priorityRanks[priority]
return isKnown
}
// agingBonusPerHour ist der Prioritätsgewinn je Wartestunde.
//
// Ohne ihn verhungerte ein Auftrag niedriger Priorität in einer Anlage mit
// genug dringenden Aufträgen für immer — und niemand bemerkte es, weil er
// technisch „wartet" statt zu scheitern. Mit 50 Punkten je Stunde überholt ein
// wartender Auftrag nach zwei Stunden die nächsthöhere Stufe.
const agingBonusPerHour = 50
// maximumAgingBonus begrenzt den Alterungsgewinn.
//
// Ohne Deckel überholte ein tagelang wartender Auftrag jede kritische
// Sicherung. Der Deckel entspricht dem Abstand zweier Stufen: ein wartender
// Auftrag steigt höchstens um eine Stufe auf.
const maximumAgingBonus = 100
// QueuedJob ist ein zur Ausführung anstehender Auftrag.
type QueuedJob struct {
// JobID ist die Kennung des Auftrags.
JobID string
// Priority ist die Dringlichkeit.
Priority Priority
// ScheduledFor ist der geplante Zeitpunkt.
ScheduledFor time.Time
// EnqueuedAt ist der Zeitpunkt der Einreihung.
EnqueuedAt time.Time
// DependsOnJobIDs sind Aufträge, die zuerst gelingen müssen.
DependsOnJobIDs []string
// AttemptNumber ist die Nummer des Versuchs, beginnend bei 1.
AttemptNumber int
}
// effectivePriority berechnet die Dringlichkeit einschließlich Alterung.
func (queuedJob *QueuedJob) effectivePriority(currentTime time.Time) int {
waitedHours := currentTime.Sub(queuedJob.EnqueuedAt).Hours()
if waitedHours < 0 {
waitedHours = 0
}
agingBonus := int(waitedHours * agingBonusPerHour)
if agingBonus > maximumAgingBonus {
agingBonus = maximumAgingBonus
}
return queuedJob.Priority.Rank() + agingBonus
}
// ConcurrencyLimits begrenzen die gleichzeitige Ausführung.
type ConcurrencyLimits struct {
// GlobalLimit ist die Gesamtzahl gleichzeitiger Läufe; 0 bedeutet unbegrenzt.
GlobalLimit int
// PerRepositoryLimit begrenzt Läufe je Repository.
//
// Nötig, weil das Repository beim Schreiben eine Sperre hält: mehrere
// gleichzeitige Läufe auf dasselbe Repository warten ohnehin aufeinander
// und belegen dabei nur Arbeitsspeicher.
PerRepositoryLimit int
// PerSourceLimit begrenzt Läufe je Quelle.
//
// Zwei gleichzeitige Sicherungen derselben VM lasten den Wirt aus, ohne
// zusätzlichen Nutzen.
PerSourceLimit int
}
// RunningJob beschreibt einen laufenden Auftrag.
type RunningJob struct {
// JobID ist die Kennung des Auftrags.
JobID string
// RepositoryID ist das Ziel-Repository.
RepositoryID string
// SourceID ist die gesicherte Quelle.
SourceID string
// StartedAt ist der Beginn.
StartedAt time.Time
}
// JobOutcome ist der Ausgang eines abgeschlossenen Laufs.
type JobOutcome string
const (
// OutcomeSucceeded ist ein vollständig gelungener Lauf.
OutcomeSucceeded JobOutcome = "succeeded"
// OutcomePartialFailure ist ein Lauf mit übergangenen Objekten.
//
// Er gilt ausdrücklich **nicht** als Erfolg (PROMPT.md §138). Abhängige
// Aufträge laufen nach einem Teilfehler nicht an.
OutcomePartialFailure JobOutcome = "partial_failure"
// OutcomeFailed ist ein gescheiterter Lauf.
OutcomeFailed JobOutcome = "failed"
// OutcomeCancelled ist ein abgebrochener Lauf.
OutcomeCancelled JobOutcome = "cancelled"
)
// IsSuccess meldet einen vollständig gelungenen Ausgang.
func (outcome JobOutcome) IsSuccess() bool {
return outcome == OutcomeSucceeded
}
// ErrDependencyCycle meldet einen Ringschluss zwischen Aufträgen.
var ErrDependencyCycle = errors.New("die abhängigkeiten der aufträge bilden einen ring")
// ErrJobAlreadyQueued meldet einen bereits eingereihten Auftrag.
var ErrJobAlreadyQueued = errors.New("der auftrag steht bereits in der warteschlange")
// Queue verwaltet anstehende und laufende Aufträge.
//
// Sie ist bewusst zustandsbehaftet und nebenläufigkeitssicher, aber ohne
// Datenbank: Was ausgeführt werden darf, ist eine Rechenfrage. Die Beständigkeit
// über Neustarts hinweg liegt eine Schicht darüber.
type Queue struct {
// mutex schützt den Zustand.
mutex sync.Mutex
// pendingJobs sind die wartenden Aufträge.
pendingJobs []QueuedJob
// runningJobs sind die laufenden Aufträge, nach Auftragskennung.
runningJobs map[string]RunningJob
// lastOutcomes hält den letzten Ausgang je Auftrag für die Abhängigkeiten.
lastOutcomes map[string]JobOutcome
// limits sind die Nebenläufigkeitsgrenzen.
limits ConcurrencyLimits
// jobMetadata hält Repository und Quelle je Auftrag.
jobMetadata map[string]RunningJob
}
// NewQueue erzeugt eine Warteschlange mit den angegebenen Grenzen.
func NewQueue(limits ConcurrencyLimits) *Queue {
return &Queue{
pendingJobs: make([]QueuedJob, 0, 16),
runningJobs: make(map[string]RunningJob),
lastOutcomes: make(map[string]JobOutcome),
jobMetadata: make(map[string]RunningJob),
limits: limits,
}
}
// RegisterJobMetadata hinterlegt Repository und Quelle eines Auftrags.
//
// Ohne diese Angaben liessen sich die auftragsbezogenen Grenzen nicht prüfen.
func (queue *Queue) RegisterJobMetadata(jobIdentifier string, repositoryIdentifier string, sourceIdentifier string) {
queue.mutex.Lock()
defer queue.mutex.Unlock()
queue.jobMetadata[jobIdentifier] = RunningJob{
JobID: jobIdentifier,
RepositoryID: repositoryIdentifier,
SourceID: sourceIdentifier,
}
}
// Enqueue reiht einen Auftrag ein.
//
// Ein bereits wartender oder laufender Auftrag wird abgelehnt. Andernfalls
// entstünden bei einem überlasteten System immer mehr Einträge desselben
// Auftrags, bis die Warteschlange den Speicher füllt.
func (queue *Queue) Enqueue(queuedJob QueuedJob) error {
queue.mutex.Lock()
defer queue.mutex.Unlock()
if _, isRunning := queue.runningJobs[queuedJob.JobID]; isRunning {
return fmt.Errorf("%w: er läuft bereits", ErrJobAlreadyQueued)
}
for _, pendingJob := range queue.pendingJobs {
if pendingJob.JobID == queuedJob.JobID {
return ErrJobAlreadyQueued
}
}
if queuedJob.EnqueuedAt.IsZero() {
queuedJob.EnqueuedAt = time.Now().UTC()
}
if queuedJob.AttemptNumber <= 0 {
queuedJob.AttemptNumber = 1
}
queue.pendingJobs = append(queue.pendingJobs, queuedJob)
return nil
}
// DispatchDecision erklärt, warum ein Auftrag läuft oder wartet.
type DispatchDecision struct {
// Job ist der betroffene Auftrag.
Job QueuedJob
// CanRun meldet die Ausführbarkeit.
CanRun bool
// Reason erklärt eine Ablehnung verständlich.
//
// Ein wartender Auftrag ohne Begründung ist für den Betrieb wertlos: Man
// sieht, dass nichts geschieht, aber nicht warum.
Reason string
}
// Dispatch wählt die als Nächstes auszuführenden Aufträge aus.
//
// Zurückgegeben werden alle Aufträge, die zum angegebenen Zeitpunkt starten
// dürfen — in der Reihenfolge ihrer Dringlichkeit. Der Aufrufer startet sie und
// meldet den Beginn über MarkRunning.
func (queue *Queue) Dispatch(currentTime time.Time) []DispatchDecision {
queue.mutex.Lock()
defer queue.mutex.Unlock()
dueJobs := make([]QueuedJob, 0, len(queue.pendingJobs))
for _, pendingJob := range queue.pendingJobs {
if !pendingJob.ScheduledFor.After(currentTime) {
dueJobs = append(dueJobs, pendingJob)
}
}
queue.sortByEffectivePriority(dueJobs, currentTime)
// Die Zählungen werden mitgeführt, damit mehrere in einem Durchgang
// freigegebene Aufträge die Grenzen nicht gemeinsam überschreiten.
runningPerRepository := make(map[string]int)
runningPerSource := make(map[string]int)
runningTotal := len(queue.runningJobs)
for _, runningJob := range queue.runningJobs {
runningPerRepository[runningJob.RepositoryID]++
runningPerSource[runningJob.SourceID]++
}
decisions := make([]DispatchDecision, 0, len(dueJobs))
for _, dueJob := range dueJobs {
jobMetadata := queue.jobMetadata[dueJob.JobID]
if blockingReason := queue.checkDependencies(dueJob); blockingReason != "" {
decisions = append(decisions, DispatchDecision{Job: dueJob, Reason: blockingReason})
continue
}
if queue.limits.GlobalLimit > 0 && runningTotal >= queue.limits.GlobalLimit {
decisions = append(decisions, DispatchDecision{
Job: dueJob,
Reason: fmt.Sprintf("Es laufen bereits %d Aufträge; mehr als %d sind nicht zugelassen.", runningTotal, queue.limits.GlobalLimit),
})
continue
}
if queue.limits.PerRepositoryLimit > 0 && jobMetadata.RepositoryID != "" &&
runningPerRepository[jobMetadata.RepositoryID] >= queue.limits.PerRepositoryLimit {
decisions = append(decisions, DispatchDecision{
Job: dueJob,
Reason: fmt.Sprintf("Auf dem Repository %s laufen bereits %d Aufträge.", jobMetadata.RepositoryID, queue.limits.PerRepositoryLimit),
})
continue
}
if queue.limits.PerSourceLimit > 0 && jobMetadata.SourceID != "" &&
runningPerSource[jobMetadata.SourceID] >= queue.limits.PerSourceLimit {
decisions = append(decisions, DispatchDecision{
Job: dueJob,
Reason: fmt.Sprintf("Die Quelle %s wird bereits gesichert.", jobMetadata.SourceID),
})
continue
}
decisions = append(decisions, DispatchDecision{Job: dueJob, CanRun: true})
runningTotal++
runningPerRepository[jobMetadata.RepositoryID]++
runningPerSource[jobMetadata.SourceID]++
}
return decisions
}
// sortByEffectivePriority ordnet Aufträge nach Dringlichkeit.
//
// Bei gleicher Dringlichkeit entscheidet der geplante Zeitpunkt, dann die
// Kennung. Die letzte Stufe ist wichtig: Ohne sie wäre die Reihenfolge bei
// gleichem Zeitpunkt zufällig und zwei Läufe derselben Anlage nicht
// vergleichbar.
func (queue *Queue) sortByEffectivePriority(jobsToSort []QueuedJob, currentTime time.Time) {
sort.SliceStable(jobsToSort, func(firstIndex int, secondIndex int) bool {
firstPriority := jobsToSort[firstIndex].effectivePriority(currentTime)
secondPriority := jobsToSort[secondIndex].effectivePriority(currentTime)
if firstPriority != secondPriority {
return firstPriority > secondPriority
}
if !jobsToSort[firstIndex].ScheduledFor.Equal(jobsToSort[secondIndex].ScheduledFor) {
return jobsToSort[firstIndex].ScheduledFor.Before(jobsToSort[secondIndex].ScheduledFor)
}
return jobsToSort[firstIndex].JobID < jobsToSort[secondIndex].JobID
})
}
// checkDependencies prüft die Abhängigkeiten eines Auftrags.
//
// Eine Abhängigkeit ist erfüllt, wenn der Vorgänger zuletzt **vollständig**
// gelungen ist. Ein Teilfehler genügt nicht: Wer eine Datenbank sichert und
// danach das Anwendungsverzeichnis, will nicht das Verzeichnis zu einer
// halben Datenbank.
func (queue *Queue) checkDependencies(queuedJob QueuedJob) string {
for _, dependencyJobID := range queuedJob.DependsOnJobIDs {
if _, isRunning := queue.runningJobs[dependencyJobID]; isRunning {
return fmt.Sprintf("Der vorausgesetzte Auftrag %s läuft noch.", dependencyJobID)
}
lastOutcome, hasOutcome := queue.lastOutcomes[dependencyJobID]
if !hasOutcome {
return fmt.Sprintf("Der vorausgesetzte Auftrag %s ist noch nie gelaufen.", dependencyJobID)
}
if !lastOutcome.IsSuccess() {
return fmt.Sprintf("Der vorausgesetzte Auftrag %s endete zuletzt mit %s.", dependencyJobID, lastOutcome)
}
}
return ""
}
// MarkRunning vermerkt den Beginn eines Laufs.
func (queue *Queue) MarkRunning(jobIdentifier string, startTime time.Time) {
queue.mutex.Lock()
defer queue.mutex.Unlock()
remainingJobs := make([]QueuedJob, 0, len(queue.pendingJobs))
for _, pendingJob := range queue.pendingJobs {
if pendingJob.JobID != jobIdentifier {
remainingJobs = append(remainingJobs, pendingJob)
}
}
queue.pendingJobs = remainingJobs
jobMetadata := queue.jobMetadata[jobIdentifier]
queue.runningJobs[jobIdentifier] = RunningJob{
JobID: jobIdentifier,
RepositoryID: jobMetadata.RepositoryID,
SourceID: jobMetadata.SourceID,
StartedAt: startTime,
}
}
// MarkFinished vermerkt das Ende eines Laufs.
func (queue *Queue) MarkFinished(jobIdentifier string, outcome JobOutcome) {
queue.mutex.Lock()
defer queue.mutex.Unlock()
delete(queue.runningJobs, jobIdentifier)
queue.lastOutcomes[jobIdentifier] = outcome
}
// RunningCount liefert die Zahl laufender Aufträge.
func (queue *Queue) RunningCount() int {
queue.mutex.Lock()
defer queue.mutex.Unlock()
return len(queue.runningJobs)
}
// PendingCount liefert die Zahl wartender Aufträge.
func (queue *Queue) PendingCount() int {
queue.mutex.Lock()
defer queue.mutex.Unlock()
return len(queue.pendingJobs)
}
// Remove nimmt einen Auftrag aus der Warteschlange.
func (queue *Queue) Remove(jobIdentifier string) bool {
queue.mutex.Lock()
defer queue.mutex.Unlock()
remainingJobs := make([]QueuedJob, 0, len(queue.pendingJobs))
var wasRemoved bool
for _, pendingJob := range queue.pendingJobs {
if pendingJob.JobID == jobIdentifier {
wasRemoved = true
continue
}
remainingJobs = append(remainingJobs, pendingJob)
}
queue.pendingJobs = remainingJobs
return wasRemoved
}
// ValidateDependencyGraph prüft eine Abhängigkeitsmenge auf Ringschlüsse.
//
// Ein Ring bliebe sonst unbemerkt: Jeder beteiligte Auftrag wartete auf einen
// anderen, keiner liefe je an, und die Oberfläche zeigte lauter „wartende"
// Aufträge ohne erkennbaren Grund. Deshalb wird beim Anlegen geprüft, nicht
// beim Ausführen.
func ValidateDependencyGraph(dependenciesByJob map[string][]string) error {
// Tiefensuche mit drei Zuständen: unbesucht, in Bearbeitung, fertig.
const (
stateUnvisited = 0
stateInProgress = 1
stateCompleted = 2
)
visitState := make(map[string]int, len(dependenciesByJob))
// Die feste Reihenfolge macht die Fehlermeldung bei mehreren Ringen
// reproduzierbar.
jobIdentifiers := make([]string, 0, len(dependenciesByJob))
for jobIdentifier := range dependenciesByJob {
jobIdentifiers = append(jobIdentifiers, jobIdentifier)
}
sort.Strings(jobIdentifiers)
var visitJob func(jobIdentifier string, visitPath []string) error
visitJob = func(jobIdentifier string, visitPath []string) error {
switch visitState[jobIdentifier] {
case stateInProgress:
return fmt.Errorf("%w: %s", ErrDependencyCycle,
formatCyclePath(append(visitPath, jobIdentifier)))
case stateCompleted:
return nil
}
visitState[jobIdentifier] = stateInProgress
for _, dependencyJobID := range dependenciesByJob[jobIdentifier] {
if dependencyJobID == jobIdentifier {
return fmt.Errorf("%w: der auftrag %s hängt von sich selbst ab", ErrDependencyCycle, jobIdentifier)
}
if visitError := visitJob(dependencyJobID, append(visitPath, jobIdentifier)); visitError != nil {
return visitError
}
}
visitState[jobIdentifier] = stateCompleted
return nil
}
for _, jobIdentifier := range jobIdentifiers {
if visitError := visitJob(jobIdentifier, nil); visitError != nil {
return visitError
}
}
return nil
}
// formatCyclePath beschreibt einen Ringschluss lesbar.
func formatCyclePath(visitPath []string) string {
pathText := ""
for pathIndex, jobIdentifier := range visitPath {
if pathIndex > 0 {
pathText += " → "
}
pathText += jobIdentifier
}
return pathText
}