syncova-backup/packages/jobs/loop.go
Jerrit Fritzsche 610719c316
Some checks failed
CI / Backend (Go) (push) Failing after 3m7s
CI / Frontend (React/TypeScript) (push) Successful in 37s
CI / Sicherheitsprüfungen (push) Successful in 44s
Syncova Backups V1
Enterprise-Backup-, Recovery-, Verification-, Security- und
Monitoring-Plattform fuer Proxmox VE, Windows, Linux und Dateisysteme.

Der Leitsatz, der fast jede Entscheidung erklaert: Ein Backup gilt erst als
vertrauenswuerdig, wenn Integritaet geprueft und Wiederherstellbarkeit
nachgewiesen wurde. Deshalb steigt ein Wiederherstellungspunkt erst nach einem
tatsaechlich durchgefuehrten Restore-Test auf "recoverable", und Unbekanntes
geht in keine Bewertung als "gut" ein.

Umfang (Phasen 0-23):

- Repository Engine: inhaltsadressierte Bloecke, atomares Commit-Protokoll,
  Katalogaufbau allein aus den Manifesten — ohne Datenbank
- Backup Engine: inhaltsabhaengiges Chunking, Deduplizierung trotz
  Verschluesselung, zstd, AES-256-GCM, Streaming mit Gegendruck
- Agenten fuer Windows und Linux mit Auftragsabholung (Pull-Modell)
- Proxmox-Provider mit beiden Zugriffswegen auf die Sicherungsarchive
- Scheduler, Recovery Engine mit Pruefpunkt, Verification, Unveraenderlichkeit
- Weboberflaeche, Kennzahlen, Meldungen, Berichte, Security Center,
  Ransomware-Heuristik (meldet, handelt nie)
- Disaster Recovery, Haertung, Leistungsmessung, Chaos Testing
- Eingefrorene Vertraege fuer API, Migrationen, Backup-Format und Repository
- Auslieferungspaket fuer linux/amd64, linux/arm64 und windows/amd64

Nicht enthalten und als solches gekennzeichnet: Kapazitaetsprognose, Backup
Copy, Changed Block Tracking bei Proxmox, erweiterte Attribute und ACLs.

Gebaut, aber nie auf echter Hardware gefahren: der Windows-Dienst, die
systemd-Einheit und der verpflichtende Proxmox-Meilenstein — ob eine
wiederhergestellte VM startet, ist ungeprueft. Einzelheiten in CHANGELOG.md
und docs/release-candidate.md.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-17 09:10:54 +02:00

734 lines
26 KiB
Go

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
}