package recovery import ( "context" "errors" "fmt" "log/slog" "sync" "time" "github.com/google/uuid" "github.com/syncova/syncova/packages/platform/logging" ) // LoopOptions steuern die Ausfuehrung von Wiederherstellungen. type LoopOptions struct { // InstanceName benennt diesen Control-Server. InstanceName string // TickInterval ist der Abstand zwischen zwei Durchgaengen. TickInterval time.Duration // MaximumConcurrentRestores begrenzt die gleichzeitigen Wiederherstellungen. // // Standard ist 1. Eine Wiederherstellung laeuft im Ernstfall, und dann // zaehlt die Geschwindigkeit **einer** Wiederherstellung, nicht der // Durchsatz mehrerer: Wer zwei gleichzeitig laufen laesst, halbiert die // Geschwindigkeit derjenigen, auf die alle warten. MaximumConcurrentRestores int // HeartbeatInterval ist der Abstand der Lebendmeldungen. HeartbeatInterval time.Duration // StaleRestoreTimeout ist die Zeit, nach der ein Auftrag als verwaist gilt. StaleRestoreTimeout time.Duration // ShutdownGracePeriod ist die Frist fuer laufende Vorgaenge beim Beenden. ShutdownGracePeriod time.Duration // CheckpointInterval ist der Abstand zwischen zwei Pruefpunkten. // // Nicht jede Datei: Bei einer Million kleiner Dateien waere das eine // Million Schreibvorgaenge in die Datenbank. Der Preis eines groesseren // Abstands ist, dass eine Fortsetzung etwas mehr wiederholt. CheckpointInterval time.Duration } // defaultRestoreTickInterval ist der Standardabstand zwischen zwei Durchgaengen. const defaultRestoreTickInterval = 5 * time.Second // defaultRestoreHeartbeatInterval ist der Standardabstand der Lebendmeldungen. const defaultRestoreHeartbeatInterval = 30 * time.Second // defaultStaleRestoreTimeout ist die Standardfrist fuer verwaiste Auftraege. const defaultStaleRestoreTimeout = 5 * time.Minute // defaultRestoreShutdownGrace ist die Standardfrist beim Beenden. const defaultRestoreShutdownGrace = 30 * time.Second // defaultCheckpointInterval ist der Standardabstand der Pruefpunkte. const defaultCheckpointInterval = 5 * time.Second // applyDefaults fuellt fehlende Werte mit sinnvollen Vorgaben. func (loopOptions *LoopOptions) applyDefaults() { if loopOptions.InstanceName == "" { loopOptions.InstanceName = "recovery-" + uuid.NewString()[:8] } if loopOptions.TickInterval <= 0 { loopOptions.TickInterval = defaultRestoreTickInterval } if loopOptions.MaximumConcurrentRestores <= 0 { loopOptions.MaximumConcurrentRestores = 1 } if loopOptions.HeartbeatInterval <= 0 { loopOptions.HeartbeatInterval = defaultRestoreHeartbeatInterval } if loopOptions.StaleRestoreTimeout <= 0 { loopOptions.StaleRestoreTimeout = defaultStaleRestoreTimeout } if loopOptions.ShutdownGracePeriod <= 0 { loopOptions.ShutdownGracePeriod = defaultRestoreShutdownGrace } if loopOptions.CheckpointInterval <= 0 { loopOptions.CheckpointInterval = defaultCheckpointInterval } } // Validate prueft die Einstellungen auf Widerspruchsfreiheit. func (loopOptions *LoopOptions) Validate() error { if loopOptions.StaleRestoreTimeout <= loopOptions.HeartbeatInterval { return fmt.Errorf("die frist fuer verwaiste auftraege (%s) muss ueber dem meldeabstand (%s) liegen", loopOptions.StaleRestoreTimeout, loopOptions.HeartbeatInterval) } return nil } // ExecutionRequest beschreibt eine auszufuehrende Wiederherstellung. type ExecutionRequest struct { // Job ist der auszufuehrende Auftrag. Job *RestoreJob // ResumeAfterPath setzt eine abgebrochene Wiederherstellung fort. ResumeAfterPath string // EntryRestoredCallback meldet jedes fertige Objekt fuer den Pruefpunkt. EntryRestoredCallback func(entryPath string, entryBytes int64) } // ExecutionResult ist das Ergebnis einer Wiederherstellung. type ExecutionResult struct { // BytesRestored ist die zurueckgeschriebene Datenmenge. BytesRestored int64 // FilesRestored ist die Zahl zurueckgeschriebener Objekte. FilesRestored int64 // FilesSkipped ist die Zahl uebergangener Objekte. FilesSkipped int64 // SkipReasons erklaeren die uebergangenen Objekte. SkipReasons []string } // Executor fuehrt eine Wiederherstellung aus. // // Wie bei den Sicherungen haelt die Schnittstelle die Schleife frei von // Repository und Engine — und macht das Zusammenspiel von Sitzung, Pruefpunkt // und Abbruch ohne echtes Repository pruefbar. type Executor interface { // Execute fuehrt die Wiederherstellung aus. Execute(executionContext context.Context, executionRequest ExecutionRequest) (ExecutionResult, error) } // ExecutionError ist ein eingeordneter Fehler einer Wiederherstellung. type ExecutionError struct { // Code ist die Fehlerkennung in SCREAMING_SNAKE_CASE. Code string // Message ist die verstaendliche Meldung. Message string // Cause ist der zugrunde liegende Fehler. Cause error } // Error erfuellt die Fehlerschnittstelle. func (executionError *ExecutionError) Error() string { return executionError.Message } // Unwrap gibt die Ursache frei. func (executionError *ExecutionError) Unwrap() error { return executionError.Cause } // NotImplementedExecutor meldet, dass keine Ausfuehrung eingerichtet ist. type NotImplementedExecutor struct{} // Execute meldet die fehlende Ausfuehrung als Fehler. func (executor *NotImplementedExecutor) Execute(_ context.Context, _ ExecutionRequest) (ExecutionResult, error) { return ExecutionResult{}, &ExecutionError{ Code: "EXECUTOR_NOT_CONFIGURED", Message: "Fuer diesen Dienst ist keine Wiederherstellung eingerichtet.", } } // Loop fuehrt Wiederherstellungsauftraege aus. type Loop struct { // store ist die Datenzugriffsschicht. store *Store // executor fuehrt die eigentliche Wiederherstellung aus. executor Executor // options sind die Einstellungen. options LoopOptions // logger protokolliert den Verlauf. logger *slog.Logger // runningGroup wartet beim Beenden auf die laufenden Vorgaenge. runningGroup sync.WaitGroup // activeRestores haelt die Abbruchfunktionen der laufenden Vorgaenge. activeRestores sync.Map // runningCount zaehlt die laufenden Vorgaenge. runningCount atomicCounter } // atomicCounter ist ein einfacher, nebenlaeufigkeitssicherer Zaehler. type atomicCounter struct { // mutex schuetzt den Wert. mutex sync.Mutex // value ist der aktuelle Stand. value int } // add veraendert den Zaehler und liefert den neuen Stand. func (counter *atomicCounter) add(delta int) int { counter.mutex.Lock() defer counter.mutex.Unlock() counter.value += delta return counter.value } // current liefert den aktuellen Stand. func (counter *atomicCounter) current() int { counter.mutex.Lock() defer counter.mutex.Unlock() return counter.value } // NewLoop erzeugt die Ausfuehrungsschleife fuer Wiederherstellungen. func NewLoop(store *Store, executor Executor, loopOptions LoopOptions, baseLogger *slog.Logger) (*Loop, error) { loopOptions.applyDefaults() if validationError := loopOptions.Validate(); validationError != nil { return nil, validationError } if executor == nil { executor = &NotImplementedExecutor{} } return &Loop{ store: store, executor: executor, options: loopOptions, logger: logging.WithComponent(baseLogger, "recovery-loop"), }, nil } // Run fuehrt die Schleife aus, bis der Kontext endet. func (loop *Loop) Run(runContext context.Context) error { loop.logger.Info("wiederherstellungsschleife gestartet", slog.String("instanz", loop.options.InstanceName), slog.Int("gleichzeitig", loop.options.MaximumConcurrentRestores)) tickTimer := time.NewTicker(loop.options.TickInterval) defer tickTimer.Stop() loop.performTick(runContext) for { select { case <-runContext.Done(): return loop.shutdown() case <-tickTimer.C: loop.performTick(runContext) } } } // shutdown beendet die Schleife geordnet. func (loop *Loop) shutdown() error { loop.logger.Info("wiederherstellungsschleife wird beendet") loop.activeRestores.Range(func(_ any, cancelValue any) bool { if cancelFunction, isFunction := cancelValue.(context.CancelFunc); isFunction { cancelFunction() } return true }) waitDone := make(chan struct{}) go func() { loop.runningGroup.Wait() close(waitDone) }() select { case <-waitDone: return nil case <-time.After(loop.options.ShutdownGracePeriod): return errors.New("die wiederherstellungsschleife wurde beendet, waehrend noch auftraege liefen") } } // performTick fuehrt einen Durchgang aus. func (loop *Loop) performTick(tickContext context.Context) { if tickContext.Err() != nil { return } reclaimedCount, reclaimError := loop.store.ReclaimStaleRestores(tickContext, loop.options.StaleRestoreTimeout) if reclaimError != nil { loop.logger.Error("verwaiste wiederherstellungen konnten nicht freigegeben werden", slog.String("grund", reclaimError.Error())) } else if reclaimedCount > 0 { loop.logger.Warn("verwaiste wiederherstellungen freigegeben", slog.Int("anzahl", reclaimedCount)) } availableSlots := loop.options.MaximumConcurrentRestores - loop.runningCount.current() if availableSlots <= 0 { return } claimedJobs, claimError := loop.store.ClaimQueuedRestores(tickContext, availableSlots, loop.options.InstanceName) if claimError != nil { loop.logger.Error("anstehende wiederherstellungen konnten nicht uebernommen werden", slog.String("grund", claimError.Error())) return } for jobIndex := range claimedJobs { loop.startRestore(tickContext, claimedJobs[jobIndex]) } } // startRestore beginnt die Ausfuehrung eines Auftrags. func (loop *Loop) startRestore(startContext context.Context, restoreJob RestoreJob) { activeSession, sessionError := loop.store.GetActiveSession(startContext, restoreJob.ID) if sessionError != nil { loop.logger.Error("die sitzung war nicht lesbar", slog.String("restore_id", restoreJob.ID.String()), slog.String("grund", sessionError.Error())) return } if beginError := loop.store.StartRestore(startContext, restoreJob.ID, loop.options.InstanceName); beginError != nil { loop.logger.Debug("der auftrag war nicht mehr startbar", slog.String("restore_id", restoreJob.ID.String()), slog.String("grund", beginError.Error())) return } loop.runningCount.add(1) executionContext, cancelExecution := context.WithCancel(context.WithoutCancel(startContext)) loop.activeRestores.Store(restoreJob.ID, cancelExecution) loop.runningGroup.Add(1) go func() { defer loop.runningGroup.Done() defer cancelExecution() defer loop.activeRestores.Delete(restoreJob.ID) defer loop.runningCount.add(-1) loop.executeRestore(executionContext, restoreJob, activeSession) }() } // executeRestore fuehrt einen Auftrag aus und schreibt sein Ergebnis fest. func (loop *Loop) executeRestore(executionContext context.Context, restoreJob RestoreJob, activeSession *RestoreSession) { startTime := time.Now() restoreLogger := loop.logger.With( slog.String("restore_id", restoreJob.ID.String()), slog.String("ziel", restoreJob.TargetRef)) resumeAfterPath := "" if activeSession != nil && activeSession.Checkpoint.LastCompletedPath != "" { resumeAfterPath = activeSession.Checkpoint.LastCompletedPath restoreLogger.Info("wiederherstellung wird fortgesetzt", slog.String("ab_pfad", resumeAfterPath), slog.Int64("bereits_fertig", activeSession.Checkpoint.FilesCompleted)) } else { restoreLogger.Info("wiederherstellung gestartet", slog.Bool("ueberschreibt", restoreJob.OverwriteExisting)) } heartbeatStop := loop.startHeartbeat(executionContext, restoreJob.ID) defer close(heartbeatStop) checkpointTracker := newCheckpointTracker(loop.store, activeSession, loop.options.CheckpointInterval) executionResult, executionError := loop.executor.Execute(executionContext, ExecutionRequest{ Job: &restoreJob, ResumeAfterPath: resumeAfterPath, EntryRestoredCallback: checkpointTracker.recordEntry, }) // Der letzte Stand wird immer geschrieben — gerade nach einem Abbruch ist // er das Wertvollste, was der Lauf hinterlaesst. checkpointTracker.flush() outcome := evaluateRestore(executionResult, executionError, executionContext.Err()) outcome.BytesRestored += checkpointTracker.baseBytes outcome.FilesRestored += checkpointTracker.baseFiles loop.finishRestore(restoreJob, activeSession, outcome, restoreLogger, time.Since(startTime)) } // startHeartbeat meldet einen laufenden Auftrag regelmaessig als lebendig. func (loop *Loop) startHeartbeat(heartbeatContext context.Context, restoreIdentifier uuid.UUID) chan struct{} { stopChannel := make(chan struct{}) go func() { heartbeatTimer := time.NewTicker(loop.options.HeartbeatInterval) defer heartbeatTimer.Stop() for { select { case <-stopChannel: return case <-heartbeatContext.Done(): return case <-heartbeatTimer.C: updateContext, cancelUpdate := context.WithTimeout(context.Background(), 5*time.Second) if heartbeatError := loop.store.RecordHeartbeat(updateContext, restoreIdentifier); heartbeatError != nil { loop.logger.Warn("die lebendmeldung schlug fehl", slog.String("restore_id", restoreIdentifier.String()), slog.String("grund", heartbeatError.Error())) } cancelUpdate() } } }() return stopChannel } // evaluateRestore wertet das Ergebnis einer Wiederherstellung aus. // // Wie bei den Sicherungen gilt: Uebergangene Objekte machen den Lauf zum // Teilfehler, auch ohne gemeldeten Fehler. Eine Wiederherstellung, die zwei // Dateien nicht zurueckschreiben konnte, ist unvollstaendig — und der Anwender // glaubte sonst, alles sei wieder da. func evaluateRestore(executionResult ExecutionResult, executionError error, contextError error) RestoreOutcome { outcome := RestoreOutcome{ BytesRestored: executionResult.BytesRestored, FilesRestored: executionResult.FilesRestored, FilesSkipped: executionResult.FilesSkipped, } if contextError != nil { outcome.Status = RestoreStatusCancelled outcome.ErrorCode = "CANCELLED" outcome.ErrorMessage = "Die Wiederherstellung wurde abgebrochen. Der Pruefpunkt der Sitzung " + "erlaubt eine Fortsetzung." return outcome } if executionError != nil { outcome.Status = RestoreStatusFailed outcome.ErrorCode = "RESTORE_FAILED" outcome.ErrorMessage = executionError.Error() var typedError *ExecutionError if errors.As(executionError, &typedError) { outcome.ErrorCode = typedError.Code outcome.ErrorMessage = typedError.Message } return outcome } if executionResult.FilesSkipped > 0 { outcome.Status = RestoreStatusPartialFailure outcome.ErrorCode = "PARTIAL_RESTORE" outcome.ErrorMessage = fmt.Sprintf("%d Objekte wurden nicht zurueckgeschrieben.", executionResult.FilesSkipped) if len(executionResult.SkipReasons) > 0 { outcome.ErrorMessage += " " + executionResult.SkipReasons[0] } return outcome } outcome.Status = RestoreStatusSucceeded return outcome } // finishRestore schreibt das Ergebnis fest und schliesst die Sitzung. func (loop *Loop) finishRestore(restoreJob RestoreJob, activeSession *RestoreSession, outcome RestoreOutcome, restoreLogger *slog.Logger, elapsedDuration time.Duration) { finishContext, cancelFinish := context.WithTimeout(context.Background(), 30*time.Second) defer cancelFinish() if writeError := loop.store.FinishRestore(finishContext, restoreJob.ID, outcome); writeError != nil { restoreLogger.Error("das ergebnis konnte nicht festgeschrieben werden", slog.String("grund", writeError.Error())) } // Die Sitzung bleibt nach einem Abbruch **offen**: Ihr Pruefpunkt ist die // Grundlage jeder Fortsetzung. Nur ein abgeschlossener Lauf schliesst sie. if activeSession != nil && outcome.Status == RestoreStatusSucceeded { if closeError := loop.store.CloseSession(finishContext, activeSession.ID, "completed"); closeError != nil { restoreLogger.Warn("die sitzung konnte nicht geschlossen werden", slog.String("grund", closeError.Error())) } } switch outcome.Status { case RestoreStatusSucceeded: restoreLogger.Info("wiederherstellung erfolgreich", slog.Int64("objekte", outcome.FilesRestored), slog.Int64("bytes", outcome.BytesRestored), slog.String("dauer", elapsedDuration.Round(time.Millisecond).String())) case RestoreStatusPartialFailure: restoreLogger.Warn("wiederherstellung TEILWEISE FEHLGESCHLAGEN", slog.Int64("zurueckgeschrieben", outcome.FilesRestored), slog.Int64("uebergangen", outcome.FilesSkipped), slog.String("meldung", outcome.ErrorMessage)) case RestoreStatusCancelled: restoreLogger.Info("wiederherstellung abgebrochen", slog.Int64("bereits_zurueckgeschrieben", outcome.FilesRestored)) default: restoreLogger.Error("wiederherstellung gescheitert", slog.String("fehlercode", outcome.ErrorCode), slog.String("meldung", outcome.ErrorMessage)) } } // RunningCount liefert die Zahl der gerade laufenden Wiederherstellungen. func (loop *Loop) RunningCount() int { return loop.runningCount.current() } // checkpointTracker fuehrt den Pruefpunkt einer Sitzung fort. // // Nicht jede Datei wird geschrieben: Bei einer Million kleiner Dateien waere // das eine Million Schreibvorgaenge in die Datenbank. Der Preis eines // groesseren Abstands ist, dass eine Fortsetzung etwas mehr wiederholt — bei // wenigen Sekunden Abstand also hoechstens wenige Sekunden Arbeit. type checkpointTracker struct { // store schreibt den Pruefpunkt. store *Store // session ist die zugehoerige Sitzung; nil schaltet den Pruefpunkt ab. session *RestoreSession // interval ist der Mindestabstand zwischen zwei Schreibvorgaengen. interval time.Duration // mutex schuetzt den Zustand. mutex sync.Mutex // current ist der aktuelle Stand. current Checkpoint // lastWriteTime ist der Zeitpunkt des letzten Schreibvorgangs. lastWriteTime time.Time // baseFiles ist die Zahl bereits vor diesem Lauf fertiger Objekte. baseFiles int64 // baseBytes ist die Menge bereits vor diesem Lauf fertiger Daten. baseBytes int64 } // newCheckpointTracker erzeugt die Fortschreibung. func newCheckpointTracker(store *Store, session *RestoreSession, interval time.Duration) *checkpointTracker { tracker := &checkpointTracker{ store: store, session: session, interval: interval, lastWriteTime: time.Now(), } if session != nil { // Bei einer Fortsetzung zaehlt der bisherige Stand weiter, damit die // Endzahlen den gesamten Vorgang beschreiben und nicht nur den letzten // Versuch. tracker.current = session.Checkpoint tracker.baseFiles = session.Checkpoint.FilesCompleted tracker.baseBytes = session.Checkpoint.BytesCompleted } return tracker } // recordEntry vermerkt ein fertig zurueckgeschriebenes Objekt. func (tracker *checkpointTracker) recordEntry(entryPath string, entryBytes int64) { if tracker.session == nil { return } tracker.mutex.Lock() tracker.current.LastCompletedPath = entryPath tracker.current.FilesCompleted++ tracker.current.BytesCompleted += entryBytes shouldWrite := time.Since(tracker.lastWriteTime) >= tracker.interval if shouldWrite { tracker.lastWriteTime = time.Now() } checkpointCopy := tracker.current tracker.mutex.Unlock() if shouldWrite { tracker.writeCheckpoint(checkpointCopy) } } // flush schreibt den letzten Stand. func (tracker *checkpointTracker) flush() { if tracker.session == nil { return } tracker.mutex.Lock() checkpointCopy := tracker.current tracker.mutex.Unlock() tracker.writeCheckpoint(checkpointCopy) } // writeCheckpoint schreibt einen Pruefpunkt in die Datenbank. func (tracker *checkpointTracker) writeCheckpoint(checkpoint Checkpoint) { // Ein eigener Kontext: Der Pruefpunkt muss auch dann geschrieben werden, // wenn der Lauf gerade abgebrochen wird — dann ist er am wertvollsten. writeContext, cancelWrite := context.WithTimeout(context.Background(), 5*time.Second) defer cancelWrite() _ = tracker.store.SaveCheckpoint(writeContext, tracker.session.ID, checkpoint) }