package jobs import ( "context" "errors" "fmt" "log/slog" "sync" "time" "github.com/google/uuid" "github.com/syncova/syncova/packages/platform/logging" "github.com/syncova/syncova/packages/scheduler" ) // LoopOptions steuern die Ausführungsschleife. type LoopOptions struct { // InstanceName benennt diesen Control-Server. // // Er steht an jedem übernommenen Lauf. Ohne ihn liesse sich nach einem // Ausfall nicht feststellen, welcher Server welche Läufe hielt. InstanceName string // TickInterval ist der Abstand zwischen zwei Durchgängen. TickInterval time.Duration // MaximumConcurrentRuns begrenzt die gleichzeitig ausgeführten Läufe. MaximumConcurrentRuns int // HeartbeatInterval ist der Abstand der Lebendmeldungen. HeartbeatInterval time.Duration // StaleRunTimeout ist die Zeit, nach der ein Lauf als verwaist gilt. // // Sie muss deutlich über dem Meldeabstand liegen: Andernfalls gäbe ein // kurzer Aussetzer — eine langsame Datenbank, ein überlasteter Server — // einen noch laufenden Lauf frei, und derselbe Auftrag liefe zweimal. StaleRunTimeout time.Duration // ShutdownGracePeriod ist die Frist für laufende Vorgänge beim Beenden. ShutdownGracePeriod time.Duration // MaintenanceWindows unterdrücken Läufe. MaintenanceWindows *scheduler.WindowSet // ConcurrencyLimits begrenzen die gleichzeitige Ausführung. ConcurrencyLimits scheduler.ConcurrencyLimits } // defaultTickInterval ist der Standardabstand zwischen zwei Durchgängen. // // Zehn Sekunden sind ein Ausgleich: Ein Zeitplan hat Minutenauflösung, also // genügt das für die Pünktlichkeit; häufiger belastete die Datenbank ohne // Erkenntnisgewinn. const defaultTickInterval = 10 * time.Second // defaultHeartbeatInterval ist der Standardabstand der Lebendmeldungen. const defaultHeartbeatInterval = 30 * time.Second // defaultStaleRunTimeout ist die Standardfrist für verwaiste Läufe. // // Fünf Minuten sind das Zehnfache des Meldeabstands. Enger gesetzt gäbe ein // kurzer Aussetzer einen noch laufenden Lauf frei. const defaultStaleRunTimeout = 5 * time.Minute // defaultShutdownGracePeriod ist die Standardfrist beim Beenden. const defaultShutdownGracePeriod = 30 * time.Second // applyDefaults füllt fehlende Werte mit sinnvollen Vorgaben. func (loopOptions *LoopOptions) applyDefaults() { if loopOptions.InstanceName == "" { loopOptions.InstanceName = "control-" + uuid.NewString()[:8] } if loopOptions.TickInterval <= 0 { loopOptions.TickInterval = defaultTickInterval } if loopOptions.MaximumConcurrentRuns <= 0 { loopOptions.MaximumConcurrentRuns = 4 } if loopOptions.HeartbeatInterval <= 0 { loopOptions.HeartbeatInterval = defaultHeartbeatInterval } if loopOptions.StaleRunTimeout <= 0 { loopOptions.StaleRunTimeout = defaultStaleRunTimeout } if loopOptions.ShutdownGracePeriod <= 0 { loopOptions.ShutdownGracePeriod = defaultShutdownGracePeriod } } // Validate prüft die Einstellungen auf Widerspruchsfreiheit. func (loopOptions *LoopOptions) Validate() error { // Eine Frist unterhalb des Meldeabstands gäbe laufende Läufe frei. Der // Fehler wäre im Betrieb kaum zu finden: Ein Auftrag liefe gelegentlich // doppelt, ohne erkennbares Muster. if loopOptions.StaleRunTimeout <= loopOptions.HeartbeatInterval { return fmt.Errorf("die frist für verwaiste läufe (%s) muss über dem meldeabstand (%s) liegen", loopOptions.StaleRunTimeout, loopOptions.HeartbeatInterval) } return nil } // Loop führt fällige Sicherungsaufträge aus. // // Der Ablauf je Durchgang: // // 1. Verwaiste Läufe abgestürzter Server freigeben. // 2. Fällige Aufträge übernehmen (legt Läufe an). // 3. Anstehende Läufe übernehmen und ausführen. // 4. Für jeden abgeschlossenen Lauf: Ergebnis festschreiben, gegebenenfalls // Wiederholung einreihen, nächsten Zeitpunkt fortschreiben. type Loop struct { // store ist die Datenzugriffsschicht. store *PostgresStore // executor führt die eigentliche Sicherung aus. executor Executor // options sind die Einstellungen. options LoopOptions // queue verwaltet Prioritäten, Nebenläufigkeit und Abhängigkeiten. queue *scheduler.Queue // agentTaskReclaimer gibt Aufträge abgestürzter Agenten frei. // // Als schmale Schnittstelle statt des ganzen Auftragsspeichers: Die // Schleife braucht genau diese eine Fähigkeit, und jobs darf nicht von // agenttasks abhängen — sonst zeigten zwei Pakete aufeinander. agentTaskReclaimer AgentTaskReclaimer // logger protokolliert den Verlauf. logger *slog.Logger // runningGroup wartet beim Beenden auf die laufenden Vorgänge. runningGroup sync.WaitGroup // activeRuns hält die Abbruchfunktionen der laufenden Vorgänge. activeRuns sync.Map } // NewLoop erzeugt die Ausführungsschleife. func NewLoop(store *PostgresStore, executor Executor, loopOptions LoopOptions, baseLogger *slog.Logger) (*Loop, error) { loopOptions.applyDefaults() if validationError := loopOptions.Validate(); validationError != nil { return nil, validationError } if executor == nil { // Kein Executor bedeutet keine Ausführung — das wird gemeldet, nicht // durch stillen Erfolg ersetzt. executor = &NotImplementedExecutor{} } return &Loop{ store: store, executor: executor, options: loopOptions, queue: scheduler.NewQueue(loopOptions.ConcurrencyLimits), logger: logging.WithComponent(baseLogger, "job-loop"), }, nil } // Run führt die Schleife aus, bis der Kontext endet. // // Beim Beenden wird auf die laufenden Vorgänge gewartet. Ein abgebrochener Lauf // wird als solcher vermerkt, statt auf „running" stehen zu bleiben: Sonst // blockierte er den Auftrag, bis die Frist für verwaiste Läufe abgelaufen ist. func (loop *Loop) Run(runContext context.Context) error { loop.logger.Info("ausführungsschleife gestartet", slog.String("instanz", loop.options.InstanceName), slog.Duration("takt", loop.options.TickInterval), slog.Int("gleichzeitig", loop.options.MaximumConcurrentRuns)) tickTimer := time.NewTicker(loop.options.TickInterval) defer tickTimer.Stop() // Der erste Durchgang läuft sofort. Ohne ihn bliebe nach einem Neustart bis // zum ersten Takt alles liegen — samt der verwaisten Läufe, die den // Neustart überhaupt nötig gemacht haben könnten. 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("ausführungsschleife wird beendet", slog.Duration("frist", loop.options.ShutdownGracePeriod)) // Die laufenden Vorgänge werden abgebrochen. Sie zu Ende laufen zu lassen // hiesse, das Herunterfahren an das längste Backup zu hängen — bei einem // mehrstündigen Lauf ein Neustart, der nie endet. loop.activeRuns.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: loop.logger.Info("alle läufe beendet") return nil case <-time.After(loop.options.ShutdownGracePeriod): // Die verbliebenen Läufe bleiben auf „running" stehen und werden von // der Freigabe verwaister Läufe eingesammelt. Das ist der ehrliche // Ausgang: Wir wissen nicht, ob sie noch etwas schreiben. loop.logger.Warn("die frist verstrich, während noch läufe aktiv waren; " + "sie werden nach Ablauf der Frist für verwaiste Läufe freigegeben") return errors.New("die ausführungsschleife wurde beendet, während noch läufe aktiv waren") } } // performTick führt einen Durchgang aus. func (loop *Loop) performTick(tickContext context.Context) { if tickContext.Err() != nil { return } loop.reclaimStaleRuns(tickContext) loop.reclaimStaleAgentTasks(tickContext) loop.claimDueJobs(tickContext) loop.dispatchQueuedRuns(tickContext) } // reclaimStaleRuns gibt Läufe abgestürzter Server frei. func (loop *Loop) reclaimStaleRuns(reclaimContext context.Context) { reclaimedCount, reclaimError := loop.store.ReclaimStaleRuns(reclaimContext, loop.options.StaleRunTimeout) if reclaimError != nil { loop.logger.Error("verwaiste läufe konnten nicht freigegeben werden", slog.String("grund", reclaimError.Error())) return } if reclaimedCount > 0 { // Kein stiller Vorgang: Ein freigegebener Lauf bedeutet, dass ein // Control-Server ausgefallen ist. loop.logger.Warn("verwaiste läufe freigegeben", slog.Int("anzahl", reclaimedCount), slog.Duration("frist", loop.options.StaleRunTimeout)) } } // AgentTaskReclaimer gibt Aufträge abgestürzter Agenten frei. type AgentTaskReclaimer interface { // ReclaimStaleTasks vermerkt Aufträge ohne Lebendmeldung als gescheitert. ReclaimStaleTasks(reclaimContext context.Context, staleAfter time.Duration) (int, error) } // SetAgentTaskReclaimer hinterlegt die Freigabe verwaister Agentenaufträge. func (loop *Loop) SetAgentTaskReclaimer(reclaimer AgentTaskReclaimer) { loop.agentTaskReclaimer = reclaimer } // reclaimStaleAgentTasks gibt Aufträge abgestürzter Agenten frei. // // Ohne diesen Schritt bliebe ein Auftrag nach dem Absturz eines Agenten // dauerhaft auf „läuft" stehen — und der Teilindex sperrte den Agenten für // jeden weiteren Auftrag, ohne dass jemand einen Fehler sähe. Dieselbe // Überlegung wie bei den verwaisten Läufen (Phase 8). func (loop *Loop) reclaimStaleAgentTasks(reclaimContext context.Context) { if loop.agentTaskReclaimer == nil { return } reclaimedCount, reclaimError := loop.agentTaskReclaimer.ReclaimStaleTasks(reclaimContext, loop.options.StaleRunTimeout) if reclaimError != nil { loop.logger.Error("verwaiste agentenauftraege konnten nicht freigegeben werden", slog.String("grund", reclaimError.Error())) return } if reclaimedCount > 0 { loop.logger.Warn("verwaiste agentenauftraege freigegeben", slog.Int("anzahl", reclaimedCount), slog.Duration("frist", loop.options.StaleRunTimeout)) } } // claimDueJobs übernimmt fällige Aufträge und legt ihre Läufe an. func (loop *Loop) claimDueJobs(claimContext context.Context) { currentTime := time.Now().UTC() claimedJobIDs, claimError := loop.store.ClaimDueJobs(claimContext, currentTime, loop.options.MaximumConcurrentRuns*2, loop.options.InstanceName) if claimError != nil { loop.logger.Error("fällige aufträge konnten nicht übernommen werden", slog.String("grund", claimError.Error())) return } // Der nächste Zeitpunkt wird sofort fortgeschrieben, nicht erst nach dem // Lauf. Sonst bliebe der Auftrag bei einem langen Backup fällig und würde // im nächsten Durchgang erneut übernommen. for _, claimedJobID := range claimedJobIDs { loop.advanceNextRun(claimContext, claimedJobID, currentTime) } } // advanceNextRun berechnet den nächsten Zeitpunkt eines Auftrags. // // Zwei Fallen werden hier behandelt: // // 1. **Kein Nachholen versäumter Läufe.** War der Server drei Tage aus, liegt // der geplante Zeitpunkt drei Tage zurück. Der Auftrag läuft **einmal**, // nicht dreimal: Drei Sicherungen desselben Bestands hintereinander kosten // Zeit und Platz, ohne einen einzigen zusätzlichen Wiederherstellungspunkt // zu schaffen — die Daten von vorgestern gibt es nicht mehr. // 2. **Kein Drift.** Gerechnet wird vom geplanten Zeitpunkt aus, nicht vom // tatsächlichen Beginn. Sonst schöbe sich ein Auftrag mit jeder Verspätung // weiter nach hinten. func (loop *Loop) advanceNextRun(advanceContext context.Context, jobIdentifier uuid.UUID, currentTime time.Time) { loadedJob, readError := loop.store.GetJob(advanceContext, jobIdentifier) if readError != nil { loop.logger.Error("der auftrag konnte für die zeitplanung nicht gelesen werden", slog.String("job_id", jobIdentifier.String()), slog.String("grund", readError.Error())) return } if loadedJob.Schedule.ScheduleType == scheduler.ScheduleTypeManual { return } // Ausgangspunkt ist der geplante Zeitpunkt, damit kein Drift entsteht. referenceTime := currentTime if loadedJob.NextRunAt != nil && loadedJob.NextRunAt.Before(currentTime) { referenceTime = *loadedJob.NextRunAt } nextRun, nextError := loadedJob.Schedule.NextRun(referenceTime) if nextError != nil { loop.logger.Error("der nächste zeitpunkt war nicht berechenbar", slog.String("job_id", jobIdentifier.String()), slog.String("grund", nextError.Error())) return } // Liegt der berechnete Zeitpunkt weiterhin in der Vergangenheit, wurden // Läufe versäumt. Sie werden übersprungen, bis der nächste in der Zukunft // liegt — und die Zahl wird genannt, statt sie zu verschweigen. var skippedRuns int for nextRun.Before(currentTime) { skippedRuns++ followingRun, followingError := loadedJob.Schedule.NextRun(nextRun) if followingError != nil { break } nextRun = followingRun } if skippedRuns > 0 { loop.logger.Warn("versäumte läufe übersprungen", slog.String("job_id", jobIdentifier.String()), slog.String("auftrag", loadedJob.Name), slog.Int("uebersprungen", skippedRuns), slog.Time("naechster_lauf", nextRun)) } // Ein Wartungsfenster verschiebt den Zeitpunkt, statt ihn ausfallen zu // lassen. if loop.options.MaintenanceWindows != nil { allowedTime, windowError := loop.options.MaintenanceWindows.NextAllowedTime(nextRun, jobIdentifier.String()) if windowError != nil { loop.logger.Error("es war kein zulässiger zeitpunkt zu finden", slog.String("job_id", jobIdentifier.String()), slog.String("grund", windowError.Error())) } else if !allowedTime.Equal(nextRun) { loop.logger.Info("der lauf wurde wegen eines wartungsfensters verschoben", slog.String("job_id", jobIdentifier.String()), slog.Time("geplant", nextRun), slog.Time("verschoben_auf", allowedTime)) nextRun = allowedTime } } if updateError := loop.store.SetNextRun(advanceContext, jobIdentifier, &nextRun); updateError != nil { loop.logger.Error("der nächste zeitpunkt konnte nicht gespeichert werden", slog.String("job_id", jobIdentifier.String()), slog.String("grund", updateError.Error())) } } // dispatchQueuedRuns übernimmt anstehende Läufe und startet sie. func (loop *Loop) dispatchQueuedRuns(dispatchContext context.Context) { availableSlots := loop.options.MaximumConcurrentRuns - loop.queue.RunningCount() if availableSlots <= 0 { return } claimedRuns, claimError := loop.store.ClaimQueuedRuns(dispatchContext, time.Now().UTC(), availableSlots, loop.options.InstanceName) if claimError != nil { loop.logger.Error("anstehende läufe konnten nicht übernommen werden", slog.String("grund", claimError.Error())) return } for runIndex := range claimedRuns { loop.startRun(dispatchContext, claimedRuns[runIndex]) } } // startRun beginnt die Ausführung eines Laufs. func (loop *Loop) startRun(startContext context.Context, queuedRun Run) { loadedJob, readError := loop.store.GetJob(startContext, queuedRun.JobID) if readError != nil { loop.failRun(startContext, queuedRun, "JOB_UNREADABLE", "Der zugehörige Auftrag war nicht lesbar.", scheduler.FailureTransient) return } // Die Abhängigkeiten werden erst hier geprüft, unmittelbar vor der // Ausführung: Zwischen dem Einreihen und dem Start kann der vorausgesetzte // Auftrag gescheitert sein. if blockingReason := loop.checkDependencies(startContext, loadedJob); blockingReason != "" { loop.logger.Info("der lauf wartet auf einen vorausgesetzten auftrag", slog.String("job_id", loadedJob.ID.String()), slog.String("grund", blockingReason)) loop.failRun(startContext, queuedRun, "DEPENDENCY_NOT_SATISFIED", blockingReason, scheduler.FailureTransient) return } if beginError := loop.store.StartRun(startContext, queuedRun.ID, loop.options.InstanceName); beginError != nil { // Ein anderer Server war schneller oder der Lauf wurde abgebrochen. // Beides ist kein Fehler dieses Servers. loop.logger.Debug("der lauf war nicht mehr startbar", slog.String("run_id", queuedRun.ID.String()), slog.String("grund", beginError.Error())) return } loop.queue.MarkRunning(loadedJob.ID.String(), time.Now().UTC()) // Der Lauf bekommt einen eigenen Kontext, damit er beim Herunterfahren // gezielt abgebrochen werden kann. executionContext, cancelExecution := context.WithCancel(context.WithoutCancel(startContext)) loop.activeRuns.Store(queuedRun.ID, cancelExecution) loop.runningGroup.Add(1) go func() { defer loop.runningGroup.Done() defer cancelExecution() defer loop.activeRuns.Delete(queuedRun.ID) loop.executeRun(executionContext, loadedJob, queuedRun) }() } // checkDependencies prüft die Abhängigkeiten eines Auftrags. func (loop *Loop) checkDependencies(checkContext context.Context, loadedJob *Job) string { for _, dependencyID := range loadedJob.DependsOnJobIDs { dependencyJob, readError := loop.store.GetJob(checkContext, dependencyID) if readError != nil { return fmt.Sprintf("Der vorausgesetzte Auftrag %s ist nicht lesbar.", dependencyID) } if dependencyJob.LastOutcome == "" { return fmt.Sprintf("Der vorausgesetzte Auftrag %q ist noch nie gelaufen.", dependencyJob.Name) } // Ein Teilfehler genügt nicht: Wer eine Datenbank sichert und danach // das Anwendungsverzeichnis, will nicht das Verzeichnis zu einer // halben Datenbank. if string(dependencyJob.LastOutcome) != string(RunStatusSucceeded) { return fmt.Sprintf("Der vorausgesetzte Auftrag %q endete zuletzt mit %s.", dependencyJob.Name, dependencyJob.LastOutcome) } } return "" } // executeRun führt einen Lauf aus und schreibt sein Ergebnis fest. func (loop *Loop) executeRun(executionContext context.Context, loadedJob *Job, currentRun Run) { startTime := time.Now() runLogger := loop.logger.With( slog.String("job_id", loadedJob.ID.String()), slog.String("run_id", currentRun.ID.String()), slog.String("auftrag", loadedJob.Name), slog.Int("versuch", currentRun.AttemptNumber)) runLogger.Info("lauf gestartet") // Die Lebendmeldung läuft nebenher. Ohne sie wäre ein stundenlanges Backup // von einem abgestürzten Server nicht zu unterscheiden — und die Freigabe // verwaister Läufe würde es abwürgen. heartbeatStop := loop.startHeartbeat(executionContext, currentRun.ID) defer close(heartbeatStop) executionResult, executionError := loop.executor.Execute(executionContext, ExecutionRequest{ Job: loadedJob, RunID: currentRun.ID, CorrelationID: currentRun.CorrelationID, AttemptNumber: currentRun.AttemptNumber, }) outcome := evaluateExecution(executionResult, executionError, executionContext.Err()) outcome.Duration = time.Since(startTime) loop.finishRun(loadedJob, currentRun, outcome, runLogger) } // startHeartbeat meldet einen laufenden Vorgang regelmäßig als lebendig. func (loop *Loop) startHeartbeat(heartbeatContext context.Context, runIdentifier 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: // Ein eigener Kontext: Der Lauf könnte gerade abgebrochen // werden, und dann käme die Meldung nicht mehr durch. updateContext, cancelUpdate := context.WithTimeout(context.Background(), 5*time.Second) if heartbeatError := loop.store.RecordHeartbeat(updateContext, runIdentifier); heartbeatError != nil { loop.logger.Warn("die lebendmeldung schlug fehl", slog.String("run_id", runIdentifier.String()), slog.String("grund", heartbeatError.Error())) } cancelUpdate() } } }() return stopChannel } // evaluateExecution wertet das Ergebnis eines Laufs aus. // // Hier fällt die wichtigste Entscheidung der ganzen Schleife: **Ein Lauf mit // übergangenen Objekten ist ein Teilfehler, auch wenn der Executor keinen // Fehler meldet** (PROMPT.md §138). Die Reihenfolge der Prüfungen ist deshalb // festgelegt und nicht beliebig. func evaluateExecution(executionResult ExecutionResult, executionError error, contextError error) ExecutionOutcome { outcome := ExecutionOutcome{Result: executionResult} // Ein Abbruch wiegt schwerer als ein Fehler: Wer abbricht, will kein // Ergebnis mehr, und ein „gescheitert" löste eine sinnlose Wiederholung aus. if contextError != nil { outcome.Status = RunStatusCancelled outcome.ErrorCode = "CANCELLED" outcome.ErrorMessage = "Der Lauf wurde abgebrochen." return outcome } if executionError != nil { outcome.Status = RunStatusFailed outcome.ErrorCode = "EXECUTION_FAILED" outcome.ErrorMessage = executionError.Error() outcome.FailureClass = scheduler.FailurePermanent var typedError *ExecutionError if errors.As(executionError, &typedError) { outcome.ErrorCode = typedError.Code outcome.ErrorMessage = typedError.Message if typedError.FailureClass != "" { outcome.FailureClass = typedError.FailureClass } } return outcome } // Kein Fehler, aber übergangene Objekte: Teilfehler. if executionResult.FilesSkipped > 0 { outcome.Status = RunStatusPartialFailure outcome.ErrorCode = "PARTIAL_FAILURE" outcome.ErrorMessage = fmt.Sprintf("%d Objekte wurden übergangen.", executionResult.FilesSkipped) if len(executionResult.SkipReasons) > 0 { outcome.ErrorMessage = fmt.Sprintf("%d Objekte wurden übergangen: %s", executionResult.FilesSkipped, executionResult.SkipReasons[0]) } return outcome } outcome.Status = RunStatusSucceeded return outcome } // finishRun schreibt das Ergebnis fest und plant gegebenenfalls eine Wiederholung. func (loop *Loop) finishRun(loadedJob *Job, currentRun Run, outcome ExecutionOutcome, runLogger *slog.Logger) { // Ein eigener Kontext: Das Ergebnis muss auch dann geschrieben werden, wenn // der Lauf abgebrochen wurde. Sonst bliebe die Zeile auf „running" stehen // und blockierte den Auftrag bis zum Ablauf der Frist. finishContext, cancelFinish := context.WithTimeout(context.Background(), 30*time.Second) defer cancelFinish() if writeError := loop.store.FinishRun(finishContext, currentRun.ID, RunOutcome{ Status: outcome.Status, BytesProcessed: outcome.Result.BytesProcessed, BytesWritten: outcome.Result.BytesWritten, FilesProcessed: outcome.Result.FilesProcessed, FilesSkipped: outcome.Result.FilesSkipped, ErrorCode: outcome.ErrorCode, ErrorMessage: outcome.ErrorMessage, FailureClass: outcome.FailureClass, }); writeError != nil { runLogger.Error("das ergebnis konnte nicht festgeschrieben werden", slog.String("grund", writeError.Error())) } loop.queue.MarkFinished(loadedJob.ID.String(), scheduler.JobOutcome(outcome.Status)) switch outcome.Status { case RunStatusSucceeded: runLogger.Info("lauf erfolgreich", slog.Int64("bytes", outcome.Result.BytesProcessed), slog.Int64("objekte", outcome.Result.FilesProcessed), slog.String("dauer", outcome.Duration.Round(time.Millisecond).String())) case RunStatusPartialFailure: // Ein Teilfehler wird nie als Erfolg protokolliert. runLogger.Warn("lauf TEILWEISE FEHLGESCHLAGEN", slog.Int64("gesichert", outcome.Result.FilesProcessed), slog.Int64("uebergangen", outcome.Result.FilesSkipped), slog.String("meldung", outcome.ErrorMessage)) case RunStatusCancelled: runLogger.Info("lauf abgebrochen") default: runLogger.Error("lauf gescheitert", slog.String("fehlercode", outcome.ErrorCode), slog.String("meldung", outcome.ErrorMessage), slog.String("fehlerart", string(outcome.FailureClass))) } loop.scheduleRetryIfNeeded(finishContext, loadedJob, currentRun, outcome, runLogger) } // scheduleRetryIfNeeded reiht bei Bedarf einen Wiederholungsversuch ein. func (loop *Loop) scheduleRetryIfNeeded(retryContext context.Context, loadedJob *Job, currentRun Run, outcome ExecutionOutcome, runLogger *slog.Logger) { // Wiederholt wird nur nach einem Fehlschlag. Ein Teilfehler wird **nicht** // wiederholt: Die übergangenen Objekte wären beim nächsten Versuch mit // hoher Wahrscheinlichkeit dieselben, und der Lauf kostete die volle Zeit // für dasselbe Ergebnis. Er verlangt einen Blick, keine Wiederholung. if outcome.Status != RunStatusFailed { return } retryDecision := loadedJob.RetryPolicy.Decide(currentRun.AttemptNumber, outcome.FailureClass) if !retryDecision.ShouldRetry { runLogger.Warn("keine weitere wiederholung", slog.String("begruendung", retryDecision.Reason)) return } retryTime := time.Now().UTC().Add(retryDecision.Delay) if _, retryError := loop.store.ScheduleRetry(retryContext, loadedJob.ID, currentRun.AttemptNumber+1, retryTime); retryError != nil { runLogger.Error("der wiederholungsversuch konnte nicht eingereiht werden", slog.String("grund", retryError.Error())) return } runLogger.Info("wiederholung eingereiht", slog.Int("naechster_versuch", currentRun.AttemptNumber+1), slog.String("wartezeit", retryDecision.Delay.Round(time.Second).String()), slog.Time("geplant_fuer", retryTime)) } // failRun schreibt einen Lauf als gescheitert fest, ohne ihn auszuführen. func (loop *Loop) failRun(failContext context.Context, currentRun Run, errorCode string, errorMessage string, failureClass scheduler.FailureClass) { if writeError := loop.store.FinishRun(failContext, currentRun.ID, RunOutcome{ Status: RunStatusFailed, ErrorCode: errorCode, ErrorMessage: errorMessage, FailureClass: failureClass, }); writeError != nil { loop.logger.Error("der gescheiterte lauf konnte nicht festgeschrieben werden", slog.String("run_id", currentRun.ID.String()), slog.String("grund", writeError.Error())) } } // RunningCount liefert die Zahl der gerade ausgeführten Läufe. func (loop *Loop) RunningCount() int { return loop.queue.RunningCount() } // InstanceName liefert den Namen dieses Control-Servers. func (loop *Loop) InstanceName() string { return loop.options.InstanceName }