syncova-backup/packages/backupexecutor/agent_delegation.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

243 lines
9.6 KiB
Go

package backupexecutor
import (
"context"
"errors"
"fmt"
"log/slog"
"time"
"github.com/google/uuid"
"github.com/syncova/syncova/packages/agenttasks"
"github.com/syncova/syncova/packages/jobs"
"github.com/syncova/syncova/packages/scheduler"
)
// ErrAgentTasksUnavailable meldet eine fehlende Auftragsvermittlung.
var ErrAgentTasksUnavailable = errors.New("die auftragsvermittlung an agenten ist nicht eingerichtet")
// agentPollInterval ist der Abstand zweier Ergebnisabfragen.
//
// Zwei Sekunden: Der Agent meldet Fortschritt, während er arbeitet; häufiger zu
// fragen erzeugte nur Last auf der Datenbank, ohne dass etwas früher fertig
// waere.
const agentPollInterval = 2 * time.Second
// delegateToAgent stellt einen Auftrag ein und wartet auf sein Ergebnis.
//
// **Der Lauf wartet synchron.** Das ist eine bewusste Entscheidung: Ein
// Sicherungslauf gilt erst als abgeschlossen, wenn alle seine Quellen gesichert
// sind. Würde der Server den Auftrag nur einstellen und den Lauf als erfolgreich
// vermerken, stünde in der Übersicht ein grüner Lauf, während der Agent noch
// arbeitet — oder bereits gescheitert ist. Das wäre der gefährlichste aller
// Fake-Erfolge.
//
// Das Warten hat eine Zeitgrenze, die aus dem Kontext des Laufs kommt. Läuft sie
// ab, bleibt der Auftrag beim Agenten stehen; der Lauf meldet einen Zeitfehler,
// und die Freigabe verwaister Aufträge räumt später auf.
func (executor *Executor) delegateToAgent(delegateContext context.Context,
backupRequest sourceBackupRequest) (jobs.ExecutionResult, error) {
if executor.options.AgentTaskStore == nil {
return jobs.ExecutionResult{}, &jobs.ExecutionError{
Code: "AGENT_TASKS_UNAVAILABLE",
Message: "Die Quelle ist einem Agenten zugewiesen, aber die Auftragsvermittlung " +
"ist nicht eingerichtet.",
FailureClass: scheduler.FailureConfiguration,
Cause: ErrAgentTasksUnavailable,
}
}
agentIdentifier := *backupRequest.Source.AgentID
delegationStart := time.Now().UTC()
chainIdentifier, chainError := executor.store.EnsureChain(delegateContext,
backupRequest.Source.SourceReference(), backupRequest.RepositoryRecord.ID)
if chainError != nil {
return jobs.ExecutionResult{}, chainError
}
parentBackupID, parentBackupInRepository, parentError := executor.store.FindLatestBackup(
delegateContext, chainIdentifier)
if parentError != nil {
backupRequest.Logger.Warn("das elternbackup liess sich nicht ermitteln; es wird voll gesichert",
slog.String("grund", parentError.Error()))
}
taskPayload := agenttasks.BackupPayload{
BackupID: buildBackupIdentifier(backupRequest.RunID, backupRequest.Source),
SourcePath: backupRequest.Source.SourceID,
SourceName: backupRequest.Source.SourceReference(),
RepositoryPath: backupRequest.RepositoryRecord.Location,
RepositoryID: backupRequest.RepositoryRecord.ID.String(),
ChainID: chainIdentifier.String(),
ParentBackupID: parentBackupInRepository,
Incremental: parentBackupInRepository != "",
CompressionLevel: "balanced",
EncryptionEnabled: executor.options.SecretStore != nil,
IncludePatterns: backupRequest.Source.IncludePatterns,
ExcludePatterns: backupRequest.Source.ExcludePatterns,
BandwidthLimitBytesPerSecond: backupRequest.Limiter.BytesPerSecond(),
}
createdTask, createError := executor.options.AgentTaskStore.CreateBackupTask(delegateContext,
agentIdentifier, backupRequest.RunID, backupRequest.Source.ID, taskPayload)
if createError != nil {
return jobs.ExecutionResult{}, &jobs.ExecutionError{
Code: "AGENT_TASK_NOT_CREATED",
Message: "Der Auftrag ließ sich nicht an den Agenten übergeben.",
FailureClass: scheduler.FailureTransient,
Cause: createError,
}
}
backupRequest.Logger.Info("ein auftrag wurde an einen agenten uebergeben",
slog.String("agent", agentIdentifier.String()),
slog.String("auftrag", createdTask.ID.String()),
slog.String("quelle", backupRequest.Source.SourceID),
slog.Bool("inkrementell", taskPayload.Incremental))
return executor.awaitAgentResult(delegateContext, createdTask.ID, backupRequest,
parentBackupID, delegationStart)
}
// awaitAgentResult wartet, bis der Agent den Auftrag abgeschlossen hat.
func (executor *Executor) awaitAgentResult(waitContext context.Context, taskIdentifier uuid.UUID,
backupRequest sourceBackupRequest, parentBackupID uuid.UUID, startedAt time.Time) (jobs.ExecutionResult, error) {
pollTicker := time.NewTicker(agentPollInterval)
defer pollTicker.Stop()
for {
select {
case <-waitContext.Done():
// Der Auftrag bleibt beim Agenten stehen; die Freigabe verwaister
// Auftraege raeumt ihn spaeter auf. Ihn hier zu loeschen waere
// falsch: Der Agent koennte gerade schreiben.
return jobs.ExecutionResult{}, &jobs.ExecutionError{
Code: "AGENT_TASK_TIMEOUT",
Message: fmt.Sprintf("Der Agent hat den Auftrag %s nicht abgeschlossen, "+
"bevor der Lauf endete.", taskIdentifier),
FailureClass: scheduler.FailureTransient,
Cause: waitContext.Err(),
}
case <-pollTicker.C:
currentTask, loadError := executor.options.AgentTaskStore.LoadTask(waitContext, taskIdentifier)
if loadError != nil {
return jobs.ExecutionResult{}, &jobs.ExecutionError{
Code: "AGENT_TASK_UNREADABLE",
Message: "Der Zustand des Agentenauftrags ließ sich nicht lesen.",
FailureClass: scheduler.FailureTransient,
Cause: loadError,
}
}
if !currentTask.Status.IsTerminal() {
continue
}
return executor.convertAgentResult(waitContext, currentTask, backupRequest,
parentBackupID, startedAt)
}
}
}
// convertAgentResult übernimmt das Ergebnis des Agenten unverändert.
//
// **Unverändert** ist hier das Wesentliche: Ein Teilfehler des Agenten bleibt
// ein Teilfehler des Laufs. Ihn zu einem Erfolg zu glätten, weil die Quelle ja
// „im Großen und Ganzen" gesichert wurde, wäre genau der vorgetäuschte Erfolg,
// den die verbindliche Regel 1 verbietet.
func (executor *Executor) convertAgentResult(convertContext context.Context, completedTask *agenttasks.Task,
backupRequest sourceBackupRequest, parentBackupID uuid.UUID, startedAt time.Time) (jobs.ExecutionResult, error) {
executionResult := jobs.ExecutionResult{
BytesProcessed: completedTask.BytesProcessed,
BytesWritten: completedTask.BytesWritten,
FilesProcessed: completedTask.FilesProcessed,
FilesSkipped: completedTask.FilesSkipped,
}
switch completedTask.Status {
case agenttasks.StatusSucceeded, agenttasks.StatusPartialFailure:
// Das Backup liegt im Repository. Der Verweis in der Datenbank ist ein
// Nachtrag: Scheitert er, ist das Backup trotzdem sicher — deshalb wird
// er protokolliert und nicht geworfen (dieselbe Ueberlegung wie bei
// RecordBackup in Phase 8).
executor.recordAgentBackup(convertContext, completedTask, backupRequest, parentBackupID, startedAt)
if completedTask.Status == agenttasks.StatusPartialFailure {
executionResult.SkipReasons = append(executionResult.SkipReasons,
fmt.Sprintf("%s: der Agent hat %d Objekte übergangen",
backupRequest.Source.SourceID, completedTask.FilesSkipped))
}
return executionResult, nil
case agenttasks.StatusCancelled:
return executionResult, &jobs.ExecutionError{
Code: "AGENT_TASK_CANCELLED",
Message: "Der Agent hat den Auftrag abgebrochen.",
FailureClass: scheduler.FailureTransient,
}
default:
failureClass := scheduler.FailureClass(completedTask.FailureClass)
if failureClass == "" {
// Ein unbekannter Fehler gilt als dauerhaft. Der umgekehrte
// Standard verdeckte die Ursache und wiederholte endlos.
failureClass = scheduler.FailurePermanent
}
errorCode := completedTask.ErrorCode
if errorCode == "" {
errorCode = "AGENT_BACKUP_FAILED"
}
return executionResult, &jobs.ExecutionError{
Code: errorCode,
Message: fmt.Sprintf("Der Agent konnte die Quelle %q nicht sichern: %s",
backupRequest.Source.SourceID, completedTask.ErrorMessage),
FailureClass: failureClass,
}
}
}
// recordAgentBackup trägt das entstandene Backup in die Control Plane ein.
//
// Er nutzt denselben Weg wie eine serverseitige Sicherung. Ein eigener
// Einfuegepfad fuer Agentensicherungen waere ein zweiter Ort, an dem ein Feld
// fehlen kann — und der Unterschied faellt erst auf, wenn jemand einen
// Wiederherstellungspunkt sucht, den es laut Datenbank nicht gibt.
func (executor *Executor) recordAgentBackup(recordContext context.Context, completedTask *agenttasks.Task,
backupRequest sourceBackupRequest, parentBackupID uuid.UUID, startedAt time.Time) {
if completedTask.BackupIDInRepository == "" {
backupRequest.Logger.Warn("der agent nannte keine backupkennung; der wiederherstellungspunkt "+
"wird nicht eingetragen",
slog.String("auftrag", completedTask.ID.String()))
return
}
chainIdentifier, chainError := uuid.Parse(completedTask.Backup.ChainID)
if chainError != nil {
backupRequest.Logger.Warn("die kettenkennung des auftrags ist unlesbar",
slog.String("grund", chainError.Error()))
return
}
executor.recordBackup(recordContext, recordRequest{
BackupRequest: backupRequest,
BackupIdentifier: completedTask.BackupIDInRepository,
ChainIdentifier: chainIdentifier,
ParentBackupID: parentBackupID,
IsIncremental: completedTask.Backup.Incremental,
Result: jobs.ExecutionResult{
BytesProcessed: completedTask.BytesProcessed,
BytesWritten: completedTask.BytesWritten,
FilesProcessed: completedTask.FilesProcessed,
FilesSkipped: completedTask.FilesSkipped,
},
StartedAt: startedAt,
})
}