package verification import ( "context" "errors" "fmt" "log/slog" "sync" "time" "github.com/google/uuid" "github.com/syncova/syncova/packages/platform/logging" ) // Runner fuehrt eine einzelne Pruefung aus. // // Die Schnittstelle haelt die Schleife frei von Repository und Engine — und // macht Uebernahme, Lebendmeldung und Abbruch ohne echtes Repository pruefbar. type Runner interface { // Run fuehrt die Pruefung aus. // // Der Wiederherstellungsbericht ist nur bei TypeRestoreTest gefuellt. Run(runContext context.Context, verificationJob Job) (*Report, *RestoreTestReport, error) } // LoopOptions steuern die Pruefschleife. type LoopOptions struct { // InstanceName benennt diesen Control-Server. InstanceName string // TickInterval ist der Abstand zwischen zwei Durchgaengen. TickInterval time.Duration // MaximumConcurrentVerifications begrenzt die gleichzeitigen Pruefungen. // // Standard ist 1. Eine Pruefung liest das gesamte Backup; mehrere // gleichzeitig teilen sich denselben Datentraeger und werden nicht // schneller, verdraengen aber die laufenden Sicherungen. MaximumConcurrentVerifications int // HeartbeatInterval ist der Abstand der Lebendmeldungen. HeartbeatInterval time.Duration // StaleJobTimeout ist die Zeit, nach der eine Pruefung als verwaist gilt. StaleJobTimeout time.Duration // ShutdownGracePeriod ist die Frist fuer laufende Pruefungen beim Beenden. ShutdownGracePeriod time.Duration } // Standardwerte der Schleife. const ( // defaultTickInterval ist der Standardabstand zwischen zwei Durchgaengen. defaultTickInterval = 5 * time.Second // defaultHeartbeatInterval ist der Standardabstand der Lebendmeldungen. defaultHeartbeatInterval = 30 * time.Second // defaultStaleJobTimeout ist die Standardfrist fuer verwaiste Pruefungen. defaultStaleJobTimeout = 30 * time.Minute // defaultShutdownGracePeriod ist die Standardfrist beim Beenden. defaultShutdownGracePeriod = 30 * time.Second ) // applyDefaults fuellt fehlende Werte mit sinnvollen Vorgaben. func (loopOptions *LoopOptions) applyDefaults() { if loopOptions.InstanceName == "" { loopOptions.InstanceName = "verification-" + uuid.NewString()[:8] } if loopOptions.TickInterval <= 0 { loopOptions.TickInterval = defaultTickInterval } if loopOptions.MaximumConcurrentVerifications <= 0 { loopOptions.MaximumConcurrentVerifications = 1 } if loopOptions.HeartbeatInterval <= 0 { loopOptions.HeartbeatInterval = defaultHeartbeatInterval } if loopOptions.StaleJobTimeout <= 0 { loopOptions.StaleJobTimeout = defaultStaleJobTimeout } if loopOptions.ShutdownGracePeriod <= 0 { loopOptions.ShutdownGracePeriod = defaultShutdownGracePeriod } } // Validate prueft die Einstellungen auf Widerspruchsfreiheit. func (loopOptions *LoopOptions) Validate() error { if loopOptions.StaleJobTimeout <= loopOptions.HeartbeatInterval { return fmt.Errorf("die frist fuer verwaiste pruefungen (%s) muss ueber dem meldeabstand (%s) liegen", loopOptions.StaleJobTimeout, loopOptions.HeartbeatInterval) } return nil } // Loop fuehrt eingereihte Pruefauftraege aus. type Loop struct { // store ist die Datenzugriffsschicht der Pruefauftraege. store *Store // runner fuehrt die einzelne Pruefung aus. runner Runner // options sind die Einstellungen. options LoopOptions // logger protokolliert den Verlauf. logger *slog.Logger // activeJobs wartet beim Beenden auf die laufenden Pruefungen. activeJobs sync.WaitGroup // runningMutex schuetzt den Zaehler der laufenden Pruefungen. runningMutex sync.Mutex // runningJobCount ist die Zahl laufender Pruefungen. // // Ein eigener Zaehler neben der WaitGroup: Deren Stand laesst sich nicht // auslesen, und ohne ihn wuesste die Schleife nicht, wie viele Plaetze frei // sind — sie uebernaehme in jedem Durchgang erneut die volle Zahl. runningJobCount int } // NewLoop erzeugt die Pruefschleife. func NewLoop(store *Store, runner Runner, loopOptions LoopOptions, baseLogger *slog.Logger) (*Loop, error) { if store == nil { return nil, errors.New("die pruefschleife braucht eine datenzugriffsschicht") } if runner == nil { return nil, errors.New("die pruefschleife braucht einen ausfuehrenden") } loopOptions.applyDefaults() if validationError := loopOptions.Validate(); validationError != nil { return nil, validationError } return &Loop{ store: store, runner: runner, options: loopOptions, logger: logging.WithComponent(baseLogger, "verification-loop"), }, nil } // Run laeuft bis zum Abbruch des Kontexts. func (loop *Loop) Run(runContext context.Context) error { loop.logger.Info("die pruefschleife beginnt", slog.String("instanz", loop.options.InstanceName), slog.Int("gleichzeitig", loop.options.MaximumConcurrentVerifications)) tickTimer := time.NewTicker(loop.options.TickInterval) defer tickTimer.Stop() // Beim Start werden Pruefungen abgestuerzter Control-Server freigegeben. // Ohne diesen Schritt blockierte der Eindeutigkeitsindex das Backup dauerhaft // gegen jede weitere Pruefung. loop.reclaimStaleJobs(runContext) for { select { case <-runContext.Done(): return loop.shutdown() case <-tickTimer.C: loop.processTick(runContext) } } } // processTick uebernimmt und startet anstehende Pruefungen. func (loop *Loop) processTick(tickContext context.Context) { availableSlots := loop.options.MaximumConcurrentVerifications - loop.runningCount() if availableSlots <= 0 { return } claimedJobs, claimError := loop.store.ClaimQueuedJobs(tickContext, availableSlots, loop.options.InstanceName) if claimError != nil { loop.logger.Error("anstehende pruefungen konnten nicht ermittelt werden", slog.String("grund", claimError.Error())) return } for _, claimedJob := range claimedJobs { loop.startJob(tickContext, claimedJob) } } // runningCount liefert die Zahl laufender Pruefungen dieser Schleife. func (loop *Loop) runningCount() int { loop.runningMutex.Lock() defer loop.runningMutex.Unlock() return loop.runningJobCount } // adjustRunningCount veraendert den Zaehler laufender Pruefungen. func (loop *Loop) adjustRunningCount(countDelta int) { loop.runningMutex.Lock() defer loop.runningMutex.Unlock() loop.runningJobCount += countDelta } // startJob fuehrt eine uebernommene Pruefung nebenlaeufig aus. func (loop *Loop) startJob(startContext context.Context, verificationJob Job) { if startError := loop.store.StartJob(startContext, verificationJob.ID, loop.options.InstanceName); startError != nil { loop.logger.Warn("eine uebernommene pruefung liess sich nicht starten", slog.String("pruefung", verificationJob.ID.String()), slog.String("grund", startError.Error())) return } loop.activeJobs.Add(1) loop.adjustRunningCount(1) go func() { defer loop.activeJobs.Done() defer loop.adjustRunningCount(-1) loop.executeJob(startContext, verificationJob) }() } // executeJob fuehrt eine Pruefung aus und schreibt ihr Ergebnis fest. func (loop *Loop) executeJob(executionContext context.Context, verificationJob Job) { jobLogger := loop.logger.With( slog.String("pruefung", verificationJob.ID.String()), slog.String("backup", verificationJob.BackupID.String()), slog.String("art", string(verificationJob.VerificationType)), slog.String("correlation_id", verificationJob.CorrelationID.String())) jobLogger.Info("die pruefung beginnt") heartbeatContext, stopHeartbeat := context.WithCancel(context.Background()) defer stopHeartbeat() go loop.sendHeartbeats(heartbeatContext, verificationJob.ID) startTime := time.Now() report, restoreReport, runError := loop.runner.Run(executionContext, verificationJob) stopHeartbeat() // Ein Ergebnis wird auch dann festgeschrieben, wenn der Kontext abgebrochen // wurde: Sonst bliebe die Pruefung auf „laeuft" stehen und der // Eindeutigkeitsindex sperrte das Backup gegen jede weitere Pruefung. finishContext, cancelFinish := context.WithTimeout(context.Background(), loop.options.ShutdownGracePeriod) defer cancelFinish() if runError != nil { // Ein Fehler der Pruefung selbst ist **kein** Befund am Backup. Ihn als // Beschaedigung zu werten waere ein Fehlalarm — und ein Pruefwerkzeug, // das grundlos Alarm schlaegt, wird bald nicht mehr ernst genommen. jobLogger.Error("die pruefung konnte nicht durchgefuehrt werden", slog.String("grund", runError.Error())) if failError := loop.store.FailJob(finishContext, verificationJob.ID, runError.Error()); failError != nil { jobLogger.Error("der fehlschlag konnte nicht vermerkt werden", slog.String("grund", failError.Error())) } return } if finishError := loop.store.FinishJob(finishContext, verificationJob.ID, report, restoreReport); finishError != nil { jobLogger.Error("das ergebnis konnte nicht festgeschrieben werden", slog.String("grund", finishError.Error())) return } logLevel := slog.LevelInfo if !report.IsClean() { // Ein Befund gehoert nicht in eine Zeile, die im Betrieb niemand liest. logLevel = slog.LevelError } jobLogger.Log(finishContext, logLevel, "die pruefung ist abgeschlossen", slog.String("ergebnis", report.Summary()), slog.Int("geprueft", report.ChunksChecked), slog.Int("fehlend", report.ChunksMissing), slog.Int("beschaedigt", report.ChunksCorrupted), slog.Duration("dauer", time.Since(startTime))) } // sendHeartbeats meldet eine laufende Pruefung regelmaessig als lebendig. func (loop *Loop) sendHeartbeats(heartbeatContext context.Context, jobIdentifier uuid.UUID) { heartbeatTicker := time.NewTicker(loop.options.HeartbeatInterval) defer heartbeatTicker.Stop() for { select { case <-heartbeatContext.Done(): return case <-heartbeatTicker.C: if heartbeatError := loop.store.RecordHeartbeat(heartbeatContext, jobIdentifier); heartbeatError != nil { loop.logger.Warn("die lebendmeldung schlug fehl", slog.String("pruefung", jobIdentifier.String()), slog.String("grund", heartbeatError.Error())) } } } } // reclaimStaleJobs gibt Pruefungen abgestuerzter Control-Server frei. func (loop *Loop) reclaimStaleJobs(reclaimContext context.Context) { reclaimedCount, reclaimError := loop.store.ReclaimStaleJobs(reclaimContext, loop.options.StaleJobTimeout) if reclaimError != nil { loop.logger.Error("verwaiste pruefungen konnten nicht freigegeben werden", slog.String("grund", reclaimError.Error())) return } if reclaimedCount > 0 { loop.logger.Warn("verwaiste pruefungen wurden freigegeben", slog.Int("anzahl", reclaimedCount)) } } // shutdown wartet auf die laufenden Pruefungen. func (loop *Loop) shutdown() error { loop.logger.Info("die pruefschleife wird beendet") finishedChannel := make(chan struct{}) go func() { loop.activeJobs.Wait() close(finishedChannel) }() select { case <-finishedChannel: loop.logger.Info("alle pruefungen wurden abgeschlossen") return nil case <-time.After(loop.options.ShutdownGracePeriod): // Eine abgebrochene Pruefung ist unangenehm, aber harmlos: Sie hat nichts // geschrieben. Der naechste Start gibt sie ueber die Frist wieder frei. loop.logger.Warn("pruefungen liefen beim beenden noch", slog.Duration("frist", loop.options.ShutdownGracePeriod)) return nil } }