syncova-backup/packages/backupexecutor/executor.go
Jerrit Fritzsche 8e98cc7510
Some checks failed
CI / Backend (Go) (push) Failing after 31s
CI / Frontend (React/TypeScript) (push) Successful in 46s
CI / Sicherheitsprüfungen (push) Successful in 28s
Sicherungsart je Auftrag, Agenten-Token und -Anleitung, update.sh
**Sicherungsart.** Bisher entschied der Executor allein: Liegt ein Elternbackup
vor, wird inkrementell gesichert. Jetzt waehlbar je Auftrag —

- `incremental` (Standard, bisheriges Verhalten),
- `always_full`, oder
- inkrementell **mit einem festen Volltag** ("immer freitags").

Migration 000014 mit drei CHECKs. Der dritte lehnt "immer voll" zusammen mit
einem Wochentag ab: Dann ist ohnehin jeder Lauf voll, und die Regel gehoert in
die Datenbank, weil im Code jede Stelle sie einhalten muesste — eine vergisst
es. Real geprueft: der Widerspruch wird abgewiesen.

Der Wochentag wird in der **Zeitzone des Zeitplans** bestimmt. Rechnete der
Server in UTC, bekaeme ein Betreiber in Berlin seine Vollsicherung am
Donnerstagabend und wunderte sich, warum sie freitags fehlt. Vier Tests, der
entscheidende durch Mutation als fangend bestaetigt.

Zur Einordnung, weil es leicht verwechselt wird: Der Platzbedarf steigt bei
"immer voll" **nicht** nennenswert — unveraenderte Bloecke werden dedupliziert
und liegen weiterhin nur einmal im Repository. Was steigt, ist die Laufzeit.
Steht so in der Maske.

**Aufnahme-Token zeigte "undefined".** Das Feld heisst `token`, nicht
`enrollment_token` — Letzteres ist der Name im *Anfrage*koerper der
Registrierung. Der dritte Formfehler dieser Art; alle konsumierten Endpunkte
sind jetzt gegen den laufenden Dienst abgeglichen.

**Der Aufnahmedialog** hat jetzt eine vollstaendige Anleitung fuer Linux und
Windows mit fertig ausgefuellten Befehlen — Serveradresse und Token eingesetzt,
je Schritt einzeln kopierbar. Eine Anleitung mit Platzhaltern fuehrt
zuverlaessig dazu, dass jemand `<token>` woertlich einsetzt und dann eine
Fehlermeldung sucht, die nichts mit seinem Problem zu tun hat. Dazu die beiden
Stolperstellen: `--state` will eine Datei, und der Agent braucht Schreibzugriff
aufs Repository. Beim Windows-Weg steht dabei, dass der Dienst nie auf echter
Hardware lief.

**update.sh ruestet die Wiederherstellungsflaeche nach** — anlegen und in
ReadWritePaths eintragen. Ein Schritt, den man von Hand ausfuehren muss, wird
uebersehen und faellt erst im Ernstfall auf.

84 Tests im Frontend, alle Go-Tests gruen, shellcheck sauber.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-18 17:13:22 +02:00

739 lines
29 KiB
Go

// Package backupexecutor verbindet die Ausführungsschleife mit der Backup Engine.
//
// Es ist die Stelle, an der aus einem geplanten Auftrag tatsächlich ein Backup
// wird. Bewusst ein eigenes Paket: Das Paket jobs bleibt damit frei von
// Repository, Engine und Providern, und das Zusammenspiel von Zeitplan,
// Wartungsfenster und Wiederholung ist ohne echtes Repository prüfbar.
package backupexecutor
import (
"context"
"errors"
"fmt"
"log/slog"
"os"
"strings"
"syscall"
"time"
"github.com/google/uuid"
"github.com/syncova/syncova/packages/agent"
"github.com/syncova/syncova/packages/agenttasks"
"github.com/syncova/syncova/packages/backupengine"
"github.com/syncova/syncova/packages/hypervisor"
"github.com/syncova/syncova/packages/jobs"
"github.com/syncova/syncova/packages/platform/crypto"
"github.com/syncova/syncova/packages/platform/logging"
"github.com/syncova/syncova/packages/platform/ratelimit"
"github.com/syncova/syncova/packages/repository"
"github.com/syncova/syncova/packages/scheduler"
)
// Options steuern den Executor.
type Options struct {
// SecretStore verschlüsselt die Datenschlüssel der Repositories.
//
// Fehlt er, sind ausschließlich unverschlüsselte Backups möglich — und die
// verweigert der Executor. Verschlüsselung ist der Standard (PROMPT.md §119).
SecretStore crypto.SecretStore
// HypervisorStore öffnet die Virtualisierungsumgebungen (Phase 7).
//
// Er darf nil sein; dann meldet eine Proxmox-Quelle einen
// Konfigurationsfehler, statt stillschweigend übergangen zu werden.
HypervisorStore *hypervisor.Store
// AgentTaskStore vermittelt Sicherungsaufträge an Agenten (Phase 5).
//
// Er darf nil sein; dann meldet eine Quelle mit zugewiesenem Agenten einen
// Konfigurationsfehler, statt stillschweigend übergangen zu werden.
AgentTaskStore *agenttasks.Store
// AllowUnencrypted lässt Sicherungen ohne Verschlüsselung zu.
//
// Ausschließlich für Testumgebungen. Der Executor protokolliert jeden
// solchen Lauf als Warnung.
AllowUnencrypted bool
// CompressionLevel ist die gewünschte Kompressionsstufe.
CompressionLevel backupengine.CompressionLevel
// CreatedByVersion ist die erzeugende Programmversion.
CreatedByVersion string
}
// Executor führt Sicherungsläufe über die Backup Engine aus.
type Executor struct {
// store ist die Datenzugriffsschicht der Control Plane.
store *jobs.PostgresStore
// options sind die Einstellungen.
options Options
// logger protokolliert den Verlauf.
logger *slog.Logger
}
// Sicherstellen, dass die Schnittstelle erfüllt wird.
var _ jobs.Executor = (*Executor)(nil)
// ErrEncryptionRequired meldet ein fehlendes Schlüsselmaterial.
var ErrEncryptionRequired = errors.New("es ist kein schlüsselmaterial eingerichtet")
// New erzeugt den Executor.
func New(store *jobs.PostgresStore, executorOptions Options, baseLogger *slog.Logger) (*Executor, error) {
if store == nil {
return nil, errors.New("der executor braucht eine datenzugriffsschicht")
}
// Ein Executor ohne Schlüssel, der trotzdem sichert, legte unverschlüsselte
// Backups an, ohne dass es jemandem auffiele. Deshalb wird die Kombination
// beim Einrichten abgelehnt, nicht erst beim ersten Lauf um zwei Uhr nachts.
if executorOptions.SecretStore == nil && !executorOptions.AllowUnencrypted {
return nil, fmt.Errorf("%w: setzen Sie SYNCOVA_ENCRYPTION_KEYS oder lassen Sie unverschlüsselte Sicherungen ausdrücklich zu",
ErrEncryptionRequired)
}
if executorOptions.CompressionLevel == "" {
executorOptions.CompressionLevel = backupengine.CompressionBalanced
}
executorLogger := logging.WithComponent(baseLogger, "backup-executor")
if executorOptions.SecretStore == nil {
executorLogger.Warn("es ist kein schlüsselmaterial eingerichtet; alle sicherungen werden UNVERSCHLÜSSELT abgelegt")
}
return &Executor{
store: store,
options: executorOptions,
logger: executorLogger,
}, nil
}
// Execute führt einen Sicherungslauf aus.
//
// Der Ablauf je Lauf: Repository auflösen und öffnen → jede Quelle sichern →
// Backup in der Control Plane vermerken → Kennzahlen zurückgeben.
//
// Das Repository wird **einmal** für den ganzen Lauf geöffnet und hält dabei
// die Schreibsperre. Es je Quelle zu öffnen und zu schließen liefe auf ein
// Wechselspiel um dieselbe Sperre hinaus.
func (executor *Executor) Execute(executionContext context.Context, executionRequest jobs.ExecutionRequest) (jobs.ExecutionResult, error) {
executionJob := executionRequest.Job
runLogger := executor.logger.With(
slog.String("job_id", executionJob.ID.String()),
slog.String("run_id", executionRequest.RunID.String()),
slog.String("auftrag", executionJob.Name))
// Das Repository wird nur geöffnet, wenn der Server selbst sichert.
//
// **Dies ist die entscheidende Stelle für Agentensicherungen.** Ein
// geöffnetes Repository hält die Schreibsperre — und der Agent, der auf
// dasselbe Ziel schreiben soll, käme dann nicht hinein. Im Nachweis der
// Phase 5 schlug genau das zu: Der Agent meldete „Repository nicht
// erreichbar", während es einwandfrei dalag; gesperrt hatte es der Server,
// der auf den Agenten wartete.
//
// Hat der Lauf ausschließlich Agentenquellen, bleibt das Repository
// deshalb ungeöffnet. Der Datensatz wird trotzdem gebraucht: Der Agent
// bekommt seinen Pfad im Auftrag.
repositoryRecord, resolveError := executor.resolveRepository(executionContext, executionJob)
if resolveError != nil {
return jobs.ExecutionResult{}, resolveError
}
var (
openedRepository *repository.LocalRepository
backupRunner *agent.BackupRunner
backupEngine *backupengine.Engine
)
if hasServerSideSource(executionJob) {
var openError error
openedRepository, openError = executor.openRepositoryRecord(executionContext,
repositoryRecord, runLogger)
if openError != nil {
return jobs.ExecutionResult{}, openError
}
defer func() {
if closeError := openedRepository.Close(); closeError != nil {
runLogger.Error("das repository konnte nicht geschlossen werden",
slog.String("grund", closeError.Error()))
}
}()
backupEngine = backupengine.NewEngine(openedRepository, executor.options.SecretStore, runLogger)
backupRunner = agent.NewBackupRunner(backupEngine, runLogger)
} else {
runLogger.Info("alle quellen dieses laufs werden von agenten gesichert; " +
"das repository bleibt fuer sie ungeoeffnet")
}
// Der Begrenzer wird **einmal je Lauf** gebaut und über alle Quellen
// geteilt. Je Quelle einen neuen zu bauen ergäbe bei einem Auftrag mit drei
// Verzeichnissen das Dreifache der vereinbarten Rate — und der Eimer
// begänne jedes Mal wieder voll.
bandwidthLimiter, limiterError := buildBandwidthLimiter(executionJob.BandwidthLimitBytesPerSecond)
if limiterError != nil {
return jobs.ExecutionResult{}, &jobs.ExecutionError{
Code: "BANDWIDTH_LIMIT_INVALID",
Message: "Die Bandbreitengrenze des Auftrags ist unbrauchbar: " + limiterError.Error(),
FailureClass: scheduler.FailureConfiguration,
Cause: limiterError,
}
}
if !bandwidthLimiter.IsUnlimited() {
runLogger.Info("der lauf ist in der bandbreite begrenzt",
slog.String("grenze", ratelimit.FormatBandwidthLimit(bandwidthLimiter.BytesPerSecond())))
}
aggregatedResult := jobs.ExecutionResult{}
var failedSources []string
for sourceIndex := range executionJob.Sources {
currentSource := executionJob.Sources[sourceIndex]
if executionContext.Err() != nil {
return aggregatedResult, executionContext.Err()
}
sourceResult, sourceError := executor.backupSingleSource(executionContext, sourceBackupRequest{
Job: executionJob,
Source: currentSource,
RunID: executionRequest.RunID,
RepositoryRecord: repositoryRecord,
Runner: backupRunner,
Engine: backupEngine,
Limiter: bandwidthLimiter,
Logger: runLogger,
})
aggregatedResult.BytesProcessed += sourceResult.BytesProcessed
aggregatedResult.BytesWritten += sourceResult.BytesWritten
aggregatedResult.FilesProcessed += sourceResult.FilesProcessed
aggregatedResult.FilesSkipped += sourceResult.FilesSkipped
aggregatedResult.SkipReasons = append(aggregatedResult.SkipReasons, sourceResult.SkipReasons...)
if sourceError != nil {
// Ein volles Repository beendet den Lauf sofort.
//
// Es ist kein Fehler der Quelle: Die naechste Quelle traefe auf
// dieselbe volle Platte, und jeder Versuch legte weitere Bloecke
// ab — die Wiederholung machte die Lage schlimmer. Der Fall
// verlangt einen Eingriff, keinen neuen Versuch.
if isOutOfSpace(sourceError) {
return aggregatedResult, &jobs.ExecutionError{
Code: "REPOSITORY_FULL",
Message: fmt.Sprintf("Im Repository ist kein Platz mehr (Quelle %q). "+
"Der Lauf wird nicht wiederholt: Ein neuer Versuch legte weitere "+
"Blöcke ab und verschärfte die Lage. Schaffen Sie Platz — über eine "+
"Aufbewahrungsregel oder mehr Speicher.", currentSource.SourceID),
FailureClass: scheduler.FailureConfiguration,
Cause: sourceError,
}
}
// Eine gescheiterte Quelle beendet den Lauf nicht: Die übrigen
// sollen gesichert werden. Verschwiegen wird sie trotzdem nicht.
runLogger.Error("eine quelle konnte nicht gesichert werden",
slog.String("quelle", currentSource.SourceID),
slog.String("grund", sourceError.Error()))
failedSources = append(failedSources,
fmt.Sprintf("%s: %s", currentSource.SourceID, sourceError.Error()))
}
}
// Sind alle Quellen gescheitert, ist der Lauf gescheitert. Bleibt ein Teil
// übrig, ist er ein Teilfehler — und das entscheidet die Schleife anhand
// von FilesSkipped.
if len(failedSources) == len(executionJob.Sources) {
return aggregatedResult, &jobs.ExecutionError{
Code: "ALL_SOURCES_FAILED",
Message: "Keine der Quellen konnte gesichert werden: " + strings.Join(failedSources, "; "),
FailureClass: scheduler.FailureSource,
}
}
if len(failedSources) > 0 {
// Eine gescheiterte Quelle zählt als übergangenes Objekt: Der Lauf wird
// damit zum Teilfehler und niemals zum Erfolg (PROMPT.md §138).
aggregatedResult.FilesSkipped += int64(len(failedSources))
aggregatedResult.SkipReasons = append(aggregatedResult.SkipReasons, failedSources...)
}
return aggregatedResult, nil
}
// openRepository löst das Ziel-Repository auf und öffnet es.
func (executor *Executor) resolveRepository(resolveContext context.Context, executionJob *jobs.Job) (*jobs.Repository, error) {
repositoryRecord, readError := executor.store.GetRepository(resolveContext, executionJob.RepositoryID)
if readError != nil {
return nil, &jobs.ExecutionError{
Code: "REPOSITORY_UNKNOWN",
Message: "Das Ziel-Repository des Auftrags ist nicht bekannt.",
FailureClass: scheduler.FailureConfiguration,
Cause: readError,
}
}
// Ein Repository im Wartungs- oder Nur-Lese-Zustand nimmt keine Sicherung
// an. Das jetzt zu erkennen erspart einen Abbruch mitten im Lauf.
if !repositoryRecord.Status.AcceptsWrites() {
return nil, &jobs.ExecutionError{
Code: "REPOSITORY_NOT_WRITABLE",
Message: fmt.Sprintf("Das Repository %q steht auf %q und nimmt keine Sicherungen an.",
repositoryRecord.Name, repositoryRecord.Status),
FailureClass: scheduler.FailureRepository,
}
}
if repositoryRecord.RepositoryType != "local" {
return nil, &jobs.ExecutionError{
Code: "REPOSITORY_TYPE_UNSUPPORTED",
Message: fmt.Sprintf("Repositories der Art %q werden noch nicht unterstützt.",
repositoryRecord.RepositoryType),
FailureClass: scheduler.FailureConfiguration,
}
}
return repositoryRecord, nil
}
// openRepositoryRecord öffnet ein aufgelöstes Repository.
//
// Getrennt vom Auflösen, weil ein Lauf mit ausschließlich Agentenquellen den
// Datensatz braucht, das geöffnete Repository aber nicht — im Gegenteil: Ein
// geöffnetes Repository hielte die Schreibsperre, die der Agent selbst braucht.
func (executor *Executor) openRepositoryRecord(openContext context.Context,
repositoryRecord *jobs.Repository, runLogger *slog.Logger) (*repository.LocalRepository, error) {
openedRepository, openError := repository.Open(openContext, repositoryRecord.Location,
repository.OpenOptions{}, runLogger)
if openError != nil {
// Ein belegtes oder nicht erreichbares Repository ist ein
// vorübergehender Fehler: Beim nächsten Versuch kann es frei sein.
return nil, &jobs.ExecutionError{
Code: "REPOSITORY_UNAVAILABLE",
Message: fmt.Sprintf("Das Repository %q unter %s ließ sich nicht öffnen.",
repositoryRecord.Name, repositoryRecord.Location),
FailureClass: scheduler.FailureRepository,
Cause: openError,
}
}
executor.verifyRepositoryIdentity(openContext, openedRepository, repositoryRecord, runLogger)
return openedRepository, nil
}
// verifyRepositoryIdentity vergleicht die hinterlegte mit der vorgefundenen Kennung.
//
// Weicht sie ab, wurde das Verzeichnis ausgetauscht. Der Lauf wird deswegen
// nicht abgebrochen — das Repository ist in sich stimmig und nimmt die
// Sicherung an —, aber die Angaben der Control Plane zu Backups und Belegung
// beziehen sich dann auf einen fremden Bestand. Das gehört ins Protokoll.
func (executor *Executor) verifyRepositoryIdentity(verifyContext context.Context, openedRepository *repository.LocalRepository, repositoryRecord *jobs.Repository, runLogger *slog.Logger) {
foundIdentifier := openedRepository.Descriptor().RepositoryID
if repositoryRecord.RepositoryUUID == "" {
if recordError := executor.store.RecordRepositoryIdentity(verifyContext,
repositoryRecord.ID, foundIdentifier); recordError != nil {
runLogger.Warn("die repository-kennung konnte nicht vermerkt werden",
slog.String("grund", recordError.Error()))
}
return
}
if repositoryRecord.RepositoryUUID != foundIdentifier {
runLogger.Error("das repository unter diesem pfad ist nicht mehr dasselbe",
slog.String("erwartet", repositoryRecord.RepositoryUUID),
slog.String("vorgefunden", foundIdentifier),
slog.String("pfad", repositoryRecord.Location))
}
}
// sourceBackupRequest bündelt die Angaben zur Sicherung einer einzelnen Quelle.
type sourceBackupRequest struct {
// Job ist der ausgeführte Auftrag.
Job *jobs.Job
// Source ist die zu sichernde Quelle.
Source jobs.JobSource
// RunID ist der Lauf.
RunID uuid.UUID
// RepositoryRecord ist das Ziel-Repository.
RepositoryRecord *jobs.Repository
// Runner führt die Sicherung einer Dateisystemquelle aus.
Runner *agent.BackupRunner
// Engine ist die geöffnete Backup Engine.
//
// Eine Proxmox-Quelle geht nicht über den BackupRunner: Der durchläuft ein
// Dateisystem, und bei einem Gast gibt es keines zu durchlaufen — es gibt
// einen Datenstrom.
Engine *backupengine.Engine
// Limiter begrenzt den Lesedurchsatz; er gilt für den ganzen Lauf.
Limiter *ratelimit.Limiter
// Logger protokolliert den Verlauf.
Logger *slog.Logger
}
// backupSingleSource sichert eine einzelne Quelle.
func (executor *Executor) backupSingleSource(backupContext context.Context, backupRequest sourceBackupRequest) (jobs.ExecutionResult, error) {
// Nur Dateisystemquellen sind umgesetzt. Alles andere wird benannt, statt
// stillschweigend übergangen zu werden — eine übergangene Quelle ergäbe ein
// Backup, das vollständig aussieht und es nicht ist.
// Ist der Quelle ein Agent zugewiesen, führt dieser die Sicherung aus.
//
// Der Control-Server kann das Dateisystem eines fremden Rechners nicht
// lesen. Genau dafür gibt es den Agenten — und bis zu dieser Phase wurde
// eine solche Quelle abgewiesen.
if backupRequest.Source.AgentID != nil {
return executor.delegateToAgent(backupContext, backupRequest)
}
// Ein Gast einer Virtualisierungsumgebung wird über den Provider gesichert.
//
// Der Weg ist ein anderer: kein Dateisystem zum Durchlaufen, sondern ein
// einzelner Datenstrom aus vzdump. Er durchläuft dieselbe Pipeline.
if backupRequest.Source.SourceType == jobs.SourceTypeProxmoxVM ||
backupRequest.Source.SourceType == jobs.SourceTypeProxmoxContainer {
return executor.backupProxmoxGuest(backupContext, backupRequest)
}
if backupRequest.Source.SourceType != jobs.SourceTypeFilesystem {
return jobs.ExecutionResult{}, fmt.Errorf(
"quellen der art %q können noch nicht gesichert werden", backupRequest.Source.SourceType)
}
sourcePath := backupRequest.Source.SourceID
if _, statError := os.Stat(sourcePath); statError != nil {
return jobs.ExecutionResult{}, fmt.Errorf("die quelle %s ist nicht erreichbar: %w", sourcePath, statError)
}
chainIdentifier, chainError := executor.store.EnsureChain(backupContext,
backupRequest.Source.SourceReference(), backupRequest.RepositoryRecord.ID)
if chainError != nil {
return jobs.ExecutionResult{}, chainError
}
// Liegt ein Elternbackup vor, wird inkrementell gesichert. Das ist der
// richtige Standard: Der Gewinn ist Lesezeit, und die ist bei jedem Lauf
// nach dem ersten der begrenzende Faktor.
parentBackupID, parentBackupInRepository, parentError := executor.store.FindLatestBackup(backupContext, chainIdentifier)
if parentError != nil {
return jobs.ExecutionResult{}, parentError
}
// Der Auftrag kann davon abweichen — dauerhaft oder an einem Wochentag.
//
// Der Platzbedarf steigt dadurch nicht nennenswert: Unveraenderte Bloecke
// werden dedupliziert und liegen weiterhin nur einmal im Repository. Was
// steigt, ist die Laufzeit.
forceFullBackup := shouldForceFullBackup(backupRequest.Job, time.Now())
backupIdentifier := buildBackupIdentifier(backupRequest.RunID, backupRequest.Source)
startTime := time.Now().UTC()
runOptions := agent.BackupRunOptions{
BackupID: backupIdentifier,
SourcePath: sourcePath,
SourceName: backupRequest.Source.SourceName,
DiscoveryOptions: agent.DiscoveryOptions{
IncludePatterns: backupRequest.Source.IncludePatterns,
ExcludePatterns: backupRequest.Source.ExcludePatterns,
},
CompressionLevel: executor.options.CompressionLevel,
EncryptionEnabled: executor.options.SecretStore != nil,
ChainID: chainIdentifier.String(),
CreatedByVersion: executor.options.CreatedByVersion,
Incremental: parentBackupInRepository != "" && !forceFullBackup,
ParentBackupID: parentBackupInRepository,
BandwidthLimiter: backupRequest.Limiter,
}
backupRequest.Logger.Info("quelle wird gesichert",
slog.String("quelle", sourcePath),
slog.String("backup_id", backupIdentifier),
slog.Bool("inkrementell", runOptions.Incremental),
slog.Bool("verschluesselt", runOptions.EncryptionEnabled))
runResult, backupError := backupRequest.Runner.RunBackup(backupContext, runOptions)
if backupError != nil {
return jobs.ExecutionResult{}, backupError
}
executionResult := jobs.ExecutionResult{
BytesProcessed: runResult.Progress.BytesProcessed,
BytesWritten: runResult.Progress.BytesWritten,
FilesProcessed: int64(runResult.FilesBackedUp + runResult.DirectoriesRecorded + runResult.SymlinksRecorded),
FilesSkipped: int64(len(runResult.Problems)),
}
for _, discoveryProblem := range runResult.Problems {
executionResult.SkipReasons = append(executionResult.SkipReasons,
discoveryProblem.Path+": "+discoveryProblem.Reason)
}
executor.recordBackup(backupContext, recordRequest{
BackupRequest: backupRequest,
BackupIdentifier: backupIdentifier,
ChainIdentifier: chainIdentifier,
ParentBackupID: parentBackupID,
IsIncremental: runOptions.Incremental,
Result: executionResult,
StartedAt: startTime,
Signals: collectRansomwareSignals(runResult),
})
backupRequest.Logger.Info("quelle gesichert",
slog.String("quelle", sourcePath),
slog.Int64("bytes_gelesen", executionResult.BytesProcessed),
slog.Int64("bytes_abgelegt", executionResult.BytesWritten),
slog.Int64("objekte", executionResult.FilesProcessed),
slog.Int64("uebergangen", executionResult.FilesSkipped))
return executionResult, nil
}
// recordRequest bündelt die Angaben für den Datenbankeintrag.
type recordRequest struct {
// BackupRequest sind die Angaben zur Quelle.
BackupRequest sourceBackupRequest
// BackupIdentifier ist die Kennung im Repository.
BackupIdentifier string
// ChainIdentifier ist die Sicherungskette.
ChainIdentifier uuid.UUID
// ParentBackupID ist das Elternbackup.
ParentBackupID uuid.UUID
// IsIncremental meldet eine Zusatzsicherung.
IsIncremental bool
// Result sind die Kennzahlen.
Result jobs.ExecutionResult
// Signals sind die Kennzahlen der Ransomware-Erkennung.
Signals ransomwareSignals
// StartedAt ist der Beginn in UTC.
StartedAt time.Time
}
// recordBackup vermerkt das Backup in der Control Plane.
//
// Ein Fehler hier macht das Backup nicht ungültig — es liegt bereits vollständig
// im Repository und ist von dort auch ohne Datenbank wiederherstellbar. Deshalb
// wird er protokolliert und nicht geworfen: Den Lauf als gescheitert zu melden,
// obwohl die Daten sicher sind, führte zu einer sinnlosen Wiederholung.
func (executor *Executor) recordBackup(recordContext context.Context, request recordRequest) {
backupType := "full"
var parentReference *uuid.UUID
if request.IsIncremental {
backupType = "incremental"
parentReference = &request.ParentBackupID
}
if _, recordError := executor.store.RecordBackup(recordContext, jobs.BackupRecord{
JobRunID: request.BackupRequest.RunID,
RepositoryID: request.BackupRequest.RepositoryRecord.ID,
ChainID: request.ChainIdentifier,
ParentBackupID: parentReference,
BackupIDInRepository: request.BackupIdentifier,
BackupType: backupType,
// Absturzkonsistent — bei einer Dateisicherung im laufenden Betrieb
// ebenso wie bei einem vzdump im Modus "snapshot". Mehr zu behaupten
// wäre falsch: Wir halten keine Anwendung an, und ob Proxmox das
// Dateisystem des Gasts über den Gastdienst eingefroren hat, lässt
// sich hinterher nicht nachweisen. Eine Konsistenzstufe, die man nicht
// belegen kann, wird nicht vergeben (PROMPT.md §138).
ConsistencyLevel: "crash_consistent",
ManifestRef: "manifests/" + request.BackupIdentifier + ".json",
LogicalBytes: request.Result.BytesProcessed,
UniqueBytes: request.Result.BytesWritten,
EncryptedBytes: request.Result.BytesWritten,
StartedAt: request.StartedAt,
CompletedAt: time.Now().UTC(),
// Die Signale der Ransomware-Erkennung (Phase 16). Sie entstehen beim
// Sichern und wären ohne diesen Schritt verloren.
IncompressibleChunks: request.Signals.IncompressibleChunks,
NewChunkCount: request.Signals.NewChunkCount,
ChangedFileCount: request.Signals.ChangedFileCount,
DeletedFileCount: request.Signals.DeletedFileCount,
ExtensionDistribution: request.Signals.ExtensionDistribution,
}); recordError != nil {
request.BackupRequest.Logger.Error("das backup liegt im repository, konnte aber nicht in der datenbank vermerkt werden",
slog.String("backup_id", request.BackupIdentifier),
slog.String("grund", recordError.Error()))
}
}
// buildBandwidthLimiter baut den Begrenzer eines Laufs.
//
// Eine Grenze von 0 bedeutet unbegrenzt und ergibt einen Begrenzer, der nichts
// begrenzt. Das ist bequemer als überall zu prüfen, ob überhaupt einer
// vorhanden ist — und verhindert den Fehler, ihn versehentlich zu übergehen.
func buildBandwidthLimiter(bytesPerSecond int64) (*ratelimit.Limiter, error) {
return ratelimit.NewLimiter(bytesPerSecond)
}
// buildBackupIdentifier bildet die Kennung eines Backups im Repository.
//
// Sie enthält die Laufkennung: Damit lässt sich von einem Backup im Repository
// aus zurückverfolgen, welcher Lauf es erzeugt hat — auch dann, wenn die
// Datenbank verloren ging. Die Quellkennung kommt hinzu, weil ein Auftrag
// mehrere Quellen haben kann und jede ihr eigenes Backup bekommt.
func buildBackupIdentifier(runIdentifier uuid.UUID, jobSource jobs.JobSource) string {
sanitizedSource := sanitizeIdentifierPart(jobSource.SourceID)
return backupIdentifierPrefix + runIdentifier.String() + "-" + sanitizedSource
}
// maximumBackupIdentifierLength ist die Obergrenze des Repositorys.
//
// Sie steht in packages/repository/layout.go: Die Kennung wird zu einem
// Dateinamen, und 64 Zeichen sind dort die Grenze.
const maximumBackupIdentifierLength = 64
// backupIdentifierPrefix leitet jede Kennung ein.
const backupIdentifierPrefix = "run-"
// maximumSourcePartLength begrenzt den Quellanteil der Kennung.
//
// Errechnet aus der Obergrenze abzüglich Präfix, UUID und Trennzeichen. Die
// Laufkennung hat Vorrang: Aus ihr lässt sich das Backup einem Lauf zuordnen,
// auch wenn die Datenbank verloren ging. Der Quellanteil ist nur eine
// Lesehilfe und darf gekürzt werden.
const maximumSourcePartLength = maximumBackupIdentifierLength - len(backupIdentifierPrefix) - 36 - 1
// sanitizeIdentifierPart macht eine Quellkennung für einen Dateinamen tauglich.
//
// Eine Backup-Kennung wird zum Dateinamen des Manifests. Ein Pfad mit
// Schrägstrichen erzeugte dort Unterverzeichnisse, ein Punkt am Anfang eine
// versteckte Datei.
func sanitizeIdentifierPart(rawIdentifier string) string {
var builder strings.Builder
for _, currentRune := range rawIdentifier {
switch {
case currentRune >= 'a' && currentRune <= 'z',
currentRune >= 'A' && currentRune <= 'Z',
currentRune >= '0' && currentRune <= '9':
builder.WriteRune(currentRune)
default:
builder.WriteByte('-')
}
}
sanitized := strings.Trim(builder.String(), "-")
if sanitized == "" {
sanitized = "quelle"
}
if len(sanitized) > maximumSourcePartLength {
sanitized = sanitized[:maximumSourcePartLength]
}
return sanitized
}
// ransomwareSignals bündelt die Kennzahlen der Erkennung.
type ransomwareSignals struct {
// IncompressibleChunks ist die Zahl nicht verkleinerbarer neuer Blöcke.
IncompressibleChunks int64
// NewChunkCount ist die Zahl neu geschriebener Blöcke.
NewChunkCount int64
// ChangedFileCount ist die Zahl geänderter oder neuer Objekte.
ChangedFileCount int64
// DeletedFileCount ist die Zahl verschwundener Objekte.
DeletedFileCount int64
// ExtensionDistribution zählt die Dateiendungen.
ExtensionDistribution map[string]int64
}
// collectRansomwareSignals liest die Signale aus dem Laufergebnis.
//
// Alle vier Größen entstehen beim Sichern: Der Entropie-Indikator in der
// Pipeline, die Änderungs- und Löschzahlen bei der Änderungserkennung, die
// Endungsverteilung bei der Erfassung. Sie hier abzuholen kostet nichts; sie
// nachträglich zu ermitteln hiesse, Millionen Blöcke erneut zu lesen.
func collectRansomwareSignals(runResult *agent.BackupRunResult) ransomwareSignals {
collectedSignals := ransomwareSignals{
IncompressibleChunks: runResult.Progress.ChunksIncompressible,
NewChunkCount: runResult.Progress.ChunksWritten,
ExtensionDistribution: runResult.ExtensionDistribution,
}
if runResult.Changes != nil {
collectedSignals.ChangedFileCount = int64(runResult.Changes.ChangedFileCount())
collectedSignals.DeletedFileCount = int64(len(runResult.Changes.DeletedPaths))
return collectedSignals
}
// Bei einer Vollsicherung gibt es kein Elternbackup und damit keine
// Änderungserkennung: Alles ist neu, nichts ist gelöscht. Das als „viele
// Änderungen" zu werten wäre ein Fehlalarm mit Ansage.
collectedSignals.ChangedFileCount = int64(runResult.FilesBackedUp)
return collectedSignals
}
// isOutOfSpace erkennt ein volles Dateisystem.
//
// Die Pruefung geht ueber den Fehlerwert des Betriebssystems, nicht ueber den
// Meldungstext: Ein Textvergleich braeche bei der ersten uebersetzten
// Fehlermeldung, und zwar unbemerkt — der Lauf wuerde dann wieder als
// Quellfehler gefuehrt und wiederholt.
func isOutOfSpace(occurredError error) bool {
return errors.Is(occurredError, syscall.ENOSPC)
}
// hasServerSideSource meldet mindestens eine Quelle, die der Server selbst sichert.
//
// Nur dann braucht er das Repository — und nur dann darf er dessen Schreibsperre
// halten.
func hasServerSideSource(executionJob *jobs.Job) bool {
for sourceIndex := range executionJob.Sources {
if executionJob.Sources[sourceIndex].AgentID == nil {
return true
}
}
return false
}
// shouldForceFullBackup meldet, ob dieser Lauf die Quelle vollstaendig lesen soll.
//
// Zwei Gruende, beide betrieblich:
//
// - **Immer voll.** Wer sein Backup ausser Haus gibt oder auf einen
// Datentraeger schreibt, der einzeln weggetragen wird, will nicht, dass ein
// Wiederherstellungspunkt an einem frueheren haengt.
// - **Woechentlich voll.** Der uebliche Kompromiss: unter der Woche schnell,
// an einem festen Tag einmal vollstaendig.
//
// Der Wochentag wird in der **Zeitzone des Zeitplans** bestimmt. Ohne diese
// Umrechnung liefe derselbe Auftrag auf zwei Servern an verschiedenen Tagen
// voll — und ein Betreiber in Berlin bekaeme seine Vollsicherung am
// Donnerstagabend, weil der Server in UTC rechnet.
func shouldForceFullBackup(executionJob *jobs.Job, currentTime time.Time) bool {
if executionJob == nil {
return false
}
if executionJob.BackupMode == jobs.BackupModeAlwaysFull {
return true
}
if executionJob.FullBackupWeekday == nil {
return false
}
scheduleLocation := time.UTC
if executionJob.Schedule.TimeZone != "" {
if loadedLocation, loadError := time.LoadLocation(executionJob.Schedule.TimeZone); loadError == nil {
scheduleLocation = loadedLocation
}
}
return currentTime.In(scheduleLocation).Weekday() == *executionJob.FullBackupWeekday
}