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, }) }