syncova-backup/packages/jobs/loop_test.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

654 lines
23 KiB
Go

package jobs
import (
"context"
"errors"
"io"
"log/slog"
"sync"
"testing"
"time"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/syncova/syncova/packages/scheduler"
)
// recordingExecutor merkt sich die ausgeführten Läufe.
//
// Er ersetzt die Backup Engine. Ohne diese Trennung liesse sich das
// Zusammenspiel von Zeitplan, Wiederholung und Abhängigkeit nur mit einem
// echten Repository prüfen — also praktisch gar nicht.
type recordingExecutor struct {
// mutex schützt den Zustand.
mutex sync.Mutex
// executedJobs sind die Namen der ausgeführten Aufträge.
executedJobs []string
// executionCount zählt die Aufrufe.
executionCount int
// resultToReturn ist das gemeldete Ergebnis.
resultToReturn ExecutionResult
// errorToReturn ist der gemeldete Fehler.
errorToReturn error
// blockUntil hält die Ausführung an, bis der Kanal schließt.
blockUntil chan struct{}
// contextObserved hält den Fehler des Kontexts nach der Ausführung.
contextObserved error
}
// Execute erfüllt die Executor-Schnittstelle.
func (executor *recordingExecutor) Execute(executionContext context.Context, executionRequest ExecutionRequest) (ExecutionResult, error) {
executor.mutex.Lock()
executor.executionCount++
executor.executedJobs = append(executor.executedJobs, executionRequest.Job.Name)
blockChannel := executor.blockUntil
resultToReturn := executor.resultToReturn
errorToReturn := executor.errorToReturn
executor.mutex.Unlock()
if blockChannel != nil {
select {
case <-blockChannel:
case <-executionContext.Done():
executor.mutex.Lock()
executor.contextObserved = executionContext.Err()
executor.mutex.Unlock()
return ExecutionResult{}, executionContext.Err()
}
}
return resultToReturn, errorToReturn
}
// count liefert die Zahl der Aufrufe.
func (executor *recordingExecutor) count() int {
executor.mutex.Lock()
defer executor.mutex.Unlock()
return executor.executionCount
}
// newTestLoop baut eine Schleife mit kurzem Takt.
func newTestLoop(testInstance *testing.T, connectionPool *pgxpool.Pool, executor Executor) *Loop {
testInstance.Helper()
testLogger := slog.New(slog.NewJSONHandler(io.Discard, nil))
builtLoop, loopError := NewLoop(NewPostgresStore(connectionPool), executor, LoopOptions{
InstanceName: "test-" + uuid.NewString()[:8],
TickInterval: 50 * time.Millisecond,
HeartbeatInterval: 100 * time.Millisecond,
StaleRunTimeout: time.Second,
ShutdownGracePeriod: 3 * time.Second,
MaximumConcurrentRuns: 4,
}, testLogger)
if loopError != nil {
testInstance.Fatalf("die Schleife ließ sich nicht bauen: %v", loopError)
}
return builtLoop
}
// runLoopUntil führt die Schleife aus, bis die Bedingung erfüllt ist.
func runLoopUntil(testInstance *testing.T, executionLoop *Loop, condition func() bool, timeout time.Duration) {
testInstance.Helper()
loopContext, cancelLoop := context.WithCancel(context.Background())
loopDone := make(chan struct{})
go func() {
defer close(loopDone)
_ = executionLoop.Run(loopContext)
}()
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
if condition() {
break
}
time.Sleep(20 * time.Millisecond)
}
cancelLoop()
select {
case <-loopDone:
case <-time.After(5 * time.Second):
testInstance.Fatal("die Schleife endete nicht innerhalb der Frist")
}
}
// TestLoopExecutesDueJob prüft den Grundfall.
func TestLoopExecutesDueJob(testInstance *testing.T) {
connectionPool := connectTestDatabase(testInstance)
testStore := NewPostgresStore(connectionPool)
repositoryID := createTestRepository(testInstance, connectionPool)
dueTime := time.Now().UTC().Add(-time.Minute)
dueJob := buildStoredJob(repositoryID)
dueJob.Schedule = scheduler.Schedule{ScheduleType: scheduler.ScheduleTypeHourly}
dueJob.NextRunAt = &dueTime
createdJobID, createError := testStore.CreateJob(context.Background(), dueJob)
if createError != nil {
testInstance.Fatalf("der Auftrag ließ sich nicht anlegen: %v", createError)
}
cleanupJob(testInstance, connectionPool, createdJobID)
testExecutor := &recordingExecutor{
resultToReturn: ExecutionResult{BytesProcessed: 1024, FilesProcessed: 7},
}
executionLoop := newTestLoop(testInstance, connectionPool, testExecutor)
runLoopUntil(testInstance, executionLoop, func() bool {
return testExecutor.count() > 0
}, 5*time.Second)
if testExecutor.count() == 0 {
testInstance.Fatal("der fällige Auftrag wurde nicht ausgeführt")
}
finishedRuns, _, listError := testStore.ListRuns(context.Background(), createdJobID, 1, 10)
if listError != nil {
testInstance.Fatalf("die Läufe ließen sich nicht lesen: %v", listError)
}
if len(finishedRuns) == 0 {
testInstance.Fatal("es wurde kein Lauf angelegt")
}
if finishedRuns[0].Status != RunStatusSucceeded {
testInstance.Errorf("der Lauf endete mit %q: %s", finishedRuns[0].Status, finishedRuns[0].ErrorMessage)
}
if finishedRuns[0].BytesProcessed != 1024 || finishedRuns[0].FilesProcessed != 7 {
testInstance.Errorf("die Kennzahlen kamen nicht an: %d Byte, %d Objekte",
finishedRuns[0].BytesProcessed, finishedRuns[0].FilesProcessed)
}
// Der nächste Zeitpunkt muss fortgeschrieben sein, sonst liefe der Auftrag
// im nächsten Durchgang erneut.
updatedJob, readError := testStore.GetJob(context.Background(), createdJobID)
if readError != nil {
testInstance.Fatalf("der Auftrag ließ sich nicht lesen: %v", readError)
}
if updatedJob.NextRunAt == nil || !updatedJob.NextRunAt.After(time.Now().UTC()) {
testInstance.Errorf("der nächste Zeitpunkt wurde nicht fortgeschrieben: %v", updatedJob.NextRunAt)
}
if updatedJob.LastOutcome != scheduler.OutcomeSucceeded {
testInstance.Errorf("der Ausgang wurde am Auftrag als %q vermerkt", updatedJob.LastOutcome)
}
}
// TestLoopTreatsSkippedFilesAsPartialFailure ist der wichtigste Ehrlichkeitstest.
//
// Der Executor meldet **keinen** Fehler — nur übergangene Objekte. Der Lauf darf
// trotzdem nicht als Erfolg dastehen (PROMPT.md §138).
func TestLoopTreatsSkippedFilesAsPartialFailure(testInstance *testing.T) {
connectionPool := connectTestDatabase(testInstance)
testStore := NewPostgresStore(connectionPool)
repositoryID := createTestRepository(testInstance, connectionPool)
dueTime := time.Now().UTC().Add(-time.Minute)
dueJob := buildStoredJob(repositoryID)
dueJob.Schedule = scheduler.Schedule{ScheduleType: scheduler.ScheduleTypeHourly}
dueJob.NextRunAt = &dueTime
createdJobID, _ := testStore.CreateJob(context.Background(), dueJob)
cleanupJob(testInstance, connectionPool, createdJobID)
testExecutor := &recordingExecutor{
resultToReturn: ExecutionResult{
BytesProcessed: 2048,
FilesProcessed: 40,
FilesSkipped: 3,
SkipReasons: []string{"Der Zugriff wurde verweigert."},
},
}
executionLoop := newTestLoop(testInstance, connectionPool, testExecutor)
runLoopUntil(testInstance, executionLoop, func() bool {
return testExecutor.count() > 0
}, 5*time.Second)
finishedRuns, _, _ := testStore.ListRuns(context.Background(), createdJobID, 1, 10)
if len(finishedRuns) == 0 {
testInstance.Fatal("es wurde kein Lauf angelegt")
}
if finishedRuns[0].Status != RunStatusPartialFailure {
testInstance.Fatalf("ein Lauf mit 3 übergangenen Objekten wurde als %q gewertet", finishedRuns[0].Status)
}
// Die Begründung muss mitkommen: „3 Objekte übergangen" ist keine Auskunft.
if finishedRuns[0].ErrorMessage == "" {
testInstance.Error("der Teilfehler wurde nicht begründet")
}
}
// TestPartialFailureBlocksDependentJob prüft die Abhängigkeit.
//
// Wer eine Datenbank sichert und danach das Anwendungsverzeichnis, will nicht
// das Verzeichnis zu einer halben Datenbank.
func TestPartialFailureBlocksDependentJob(testInstance *testing.T) {
connectionPool := connectTestDatabase(testInstance)
testStore := NewPostgresStore(connectionPool)
repositoryID := createTestRepository(testInstance, connectionPool)
// Der vorausgesetzte Auftrag endete mit einem Teilfehler.
predecessorJob := buildStoredJob(repositoryID)
predecessorJob.Schedule = scheduler.Schedule{ScheduleType: scheduler.ScheduleTypeManual}
predecessorID, _ := testStore.CreateJob(context.Background(), predecessorJob)
cleanupJob(testInstance, connectionPool, predecessorID)
if _, execError := connectionPool.Exec(context.Background(),
`UPDATE backup_jobs SET last_outcome = 'partial_failure' WHERE id = $1`, predecessorID); execError != nil {
testInstance.Fatalf("der Ausgang ließ sich nicht setzen: %v", execError)
}
dueTime := time.Now().UTC().Add(-time.Minute)
dependentJob := buildStoredJob(repositoryID)
dependentJob.Schedule = scheduler.Schedule{ScheduleType: scheduler.ScheduleTypeHourly}
dependentJob.NextRunAt = &dueTime
dependentJob.DependsOnJobIDs = []uuid.UUID{predecessorID}
dependentID, createError := testStore.CreateJob(context.Background(), dependentJob)
if createError != nil {
testInstance.Fatalf("der abhängige Auftrag ließ sich nicht anlegen: %v", createError)
}
cleanupJob(testInstance, connectionPool, dependentID)
testExecutor := &recordingExecutor{}
executionLoop := newTestLoop(testInstance, connectionPool, testExecutor)
runLoopUntil(testInstance, executionLoop, func() bool {
runs, _, _ := testStore.ListRuns(context.Background(), dependentID, 1, 10)
return len(runs) > 0 && runs[0].IsFinished()
}, 5*time.Second)
if testExecutor.count() > 0 {
testInstance.Fatal("der abhängige Auftrag lief, obwohl der Vorgänger einen Teilfehler meldete")
}
blockedRuns, _, _ := testStore.ListRuns(context.Background(), dependentID, 1, 10)
if len(blockedRuns) == 0 {
testInstance.Fatal("es wurde kein Lauf angelegt")
}
if blockedRuns[0].ErrorCode != "DEPENDENCY_NOT_SATISFIED" {
testInstance.Errorf("die Blockade wurde als %q vermerkt", blockedRuns[0].ErrorCode)
}
}
// TestLoopSchedulesRetryOnTransientFailure prüft die Wiederholung.
func TestLoopSchedulesRetryOnTransientFailure(testInstance *testing.T) {
connectionPool := connectTestDatabase(testInstance)
testStore := NewPostgresStore(connectionPool)
repositoryID := createTestRepository(testInstance, connectionPool)
dueTime := time.Now().UTC().Add(-time.Minute)
failingJob := buildStoredJob(repositoryID)
failingJob.Schedule = scheduler.Schedule{ScheduleType: scheduler.ScheduleTypeHourly}
failingJob.NextRunAt = &dueTime
failingJob.RetryPolicy = scheduler.RetryPolicy{
MaximumAttempts: 3,
InitialDelay: time.Hour, // weit in der Zukunft, damit der Test nicht wartet
MaximumDelay: 2 * time.Hour,
}
createdJobID, _ := testStore.CreateJob(context.Background(), failingJob)
cleanupJob(testInstance, connectionPool, createdJobID)
testExecutor := &recordingExecutor{
errorToReturn: &ExecutionError{
Code: "REPOSITORY_UNAVAILABLE",
Message: "Das Repository war nicht erreichbar.",
FailureClass: scheduler.FailureTransient,
},
}
executionLoop := newTestLoop(testInstance, connectionPool, testExecutor)
runLoopUntil(testInstance, executionLoop, func() bool {
runs, _, _ := testStore.ListRuns(context.Background(), createdJobID, 1, 10)
return len(runs) >= 2
}, 5*time.Second)
allRuns, _, _ := testStore.ListRuns(context.Background(), createdJobID, 1, 10)
if len(allRuns) < 2 {
testInstance.Fatalf("es wurde keine Wiederholung eingereiht; %d Läufe vorhanden", len(allRuns))
}
// Der neueste Lauf ist der Wiederholungsversuch.
retryRun := allRuns[0]
if retryRun.Trigger != TriggerRetry {
testInstance.Errorf("der neue Lauf trägt den Auslöser %q", retryRun.Trigger)
}
if retryRun.AttemptNumber != 2 {
testInstance.Errorf("der Wiederholungsversuch trägt die Nummer %d", retryRun.AttemptNumber)
}
// Er darf nicht sofort laufen, sondern erst nach der Wartezeit.
if retryRun.ScheduledFor == nil || !retryRun.ScheduledFor.After(time.Now().UTC().Add(30*time.Minute)) {
testInstance.Errorf("die Wartezeit wurde nicht eingehalten: geplant für %v", retryRun.ScheduledFor)
}
}
// TestLoopDoesNotRetryPermanentFailure prüft den dauerhaften Fehler.
//
// Ein Anmeldefehler behebt sich nicht durch Warten.
func TestLoopDoesNotRetryPermanentFailure(testInstance *testing.T) {
connectionPool := connectTestDatabase(testInstance)
testStore := NewPostgresStore(connectionPool)
repositoryID := createTestRepository(testInstance, connectionPool)
dueTime := time.Now().UTC().Add(-time.Minute)
failingJob := buildStoredJob(repositoryID)
failingJob.Schedule = scheduler.Schedule{ScheduleType: scheduler.ScheduleTypeHourly}
failingJob.NextRunAt = &dueTime
createdJobID, _ := testStore.CreateJob(context.Background(), failingJob)
cleanupJob(testInstance, connectionPool, createdJobID)
testExecutor := &recordingExecutor{
errorToReturn: &ExecutionError{
Code: "AUTHENTICATION_FAILED",
Message: "Die Anmeldedaten wurden abgelehnt.",
FailureClass: scheduler.FailureAuthentication,
},
}
executionLoop := newTestLoop(testInstance, connectionPool, testExecutor)
runLoopUntil(testInstance, executionLoop, func() bool {
runs, _, _ := testStore.ListRuns(context.Background(), createdJobID, 1, 10)
return len(runs) > 0 && runs[0].IsFinished()
}, 3*time.Second)
// Kurz warten, damit eine fälschliche Wiederholung Zeit hätte zu erscheinen.
time.Sleep(300 * time.Millisecond)
allRuns, _, _ := testStore.ListRuns(context.Background(), createdJobID, 1, 10)
for _, singleRun := range allRuns {
if singleRun.Trigger == TriggerRetry {
testInstance.Fatal("ein Anmeldefehler wurde wiederholt, obwohl er sich nicht behebt")
}
}
if len(allRuns) == 0 || allRuns[0].Status != RunStatusFailed {
testInstance.Fatalf("der Lauf endete nicht als gescheitert: %+v", allRuns)
}
}
// TestReclaimStaleRunsFreesBlockedJob ist der wichtigste Selbstheilungstest.
//
// Stirbt ein Control-Server mitten im Lauf, bleibt dessen Zeile auf „running"
// stehen. Der Teilindex lässt dann keinen weiteren Lauf zu — die Sicherung
// fiele **dauerhaft** aus, ohne dass jemand einen Fehler sähe.
func TestReclaimStaleRunsFreesBlockedJob(testInstance *testing.T) {
connectionPool := connectTestDatabase(testInstance)
testStore := NewPostgresStore(connectionPool)
repositoryID := createTestRepository(testInstance, connectionPool)
blockedJob := buildStoredJob(repositoryID)
blockedJob.Schedule = scheduler.Schedule{ScheduleType: scheduler.ScheduleTypeHourly}
createdJobID, _ := testStore.CreateJob(context.Background(), blockedJob)
cleanupJob(testInstance, connectionPool, createdJobID)
// Ein Lauf, dessen Server vor zwei Stunden verschwand.
if _, execError := connectionPool.Exec(context.Background(), `
INSERT INTO backup_job_runs (job_id, status, correlation_id, scheduler_instance, started_at, heartbeat_at)
VALUES ($1, 'running', gen_random_uuid(), 'toter-server', now() - interval '2 hours', now() - interval '2 hours')`,
createdJobID); execError != nil {
testInstance.Fatalf("der verwaiste Lauf ließ sich nicht anlegen: %v", execError)
}
// Vor der Freigabe blockiert er jeden neuen Lauf.
if _, blockedError := testStore.CreateManualRun(context.Background(), createdJobID, nil); !errors.Is(blockedError, ErrRunAlreadyActive) {
testInstance.Fatalf("der verwaiste Lauf blockierte nicht wie erwartet: %v", blockedError)
}
reclaimedCount, reclaimError := testStore.ReclaimStaleRuns(context.Background(), time.Minute)
if reclaimError != nil {
testInstance.Fatalf("die Freigabe schlug fehl: %v", reclaimError)
}
if reclaimedCount != 1 {
testInstance.Fatalf("es wurden %d Läufe freigegeben, erwartet war 1", reclaimedCount)
}
// Danach ist der Auftrag wieder ausführbar.
if _, freeError := testStore.CreateManualRun(context.Background(), createdJobID, nil); freeError != nil {
testInstance.Fatalf("der Auftrag blieb nach der Freigabe blockiert: %v", freeError)
}
// Der Abbruch muss als Ereignis erhalten bleiben, nicht gelöscht werden.
allRuns, _, _ := testStore.ListRuns(context.Background(), createdJobID, 1, 10)
var foundReclaimed bool
for _, singleRun := range allRuns {
if singleRun.ErrorCode == "SCHEDULER_LOST" {
foundReclaimed = true
if singleRun.FailureClass != scheduler.FailureTransient {
testInstance.Errorf("der verwaiste Lauf wurde als %q eingeordnet; eine Wiederholung wäre damit ausgeschlossen",
singleRun.FailureClass)
}
}
}
if !foundReclaimed {
testInstance.Error("der freigegebene Lauf wurde nicht als solcher vermerkt")
}
}
// TestReclaimSparesLiveRuns prüft, dass laufende Vorgänge unberührt bleiben.
//
// Eine zu enge Frist gäbe einen noch laufenden Lauf frei — und derselbe Auftrag
// liefe zweimal.
func TestReclaimSparesLiveRuns(testInstance *testing.T) {
connectionPool := connectTestDatabase(testInstance)
testStore := NewPostgresStore(connectionPool)
repositoryID := createTestRepository(testInstance, connectionPool)
liveJob := buildStoredJob(repositoryID)
liveJob.Schedule = scheduler.Schedule{ScheduleType: scheduler.ScheduleTypeManual}
createdJobID, _ := testStore.CreateJob(context.Background(), liveJob)
cleanupJob(testInstance, connectionPool, createdJobID)
// Ein Lauf, der sich gerade eben gemeldet hat.
if _, execError := connectionPool.Exec(context.Background(), `
INSERT INTO backup_job_runs (job_id, status, correlation_id, scheduler_instance, started_at, heartbeat_at)
VALUES ($1, 'running', gen_random_uuid(), 'lebender-server', now(), now())`,
createdJobID); execError != nil {
testInstance.Fatalf("der Lauf ließ sich nicht anlegen: %v", execError)
}
reclaimedCount, reclaimError := testStore.ReclaimStaleRuns(context.Background(), time.Minute)
if reclaimError != nil {
testInstance.Fatalf("die Freigabe schlug fehl: %v", reclaimError)
}
if reclaimedCount != 0 {
testInstance.Fatalf("ein lebender Lauf wurde freigegeben (%d); derselbe Auftrag liefe damit zweimal", reclaimedCount)
}
}
// TestLoopShutdownCancelsRunningExecution prüft das geordnete Beenden.
//
// Ein Prozess, der endet, während noch ein Lauf schreibt, hinterliesse einen
// Lauf auf „running" und blockierte den Auftrag bis zum Ablauf der Frist.
func TestLoopShutdownCancelsRunningExecution(testInstance *testing.T) {
connectionPool := connectTestDatabase(testInstance)
testStore := NewPostgresStore(connectionPool)
repositoryID := createTestRepository(testInstance, connectionPool)
dueTime := time.Now().UTC().Add(-time.Minute)
longJob := buildStoredJob(repositoryID)
longJob.Schedule = scheduler.Schedule{ScheduleType: scheduler.ScheduleTypeHourly}
longJob.NextRunAt = &dueTime
createdJobID, _ := testStore.CreateJob(context.Background(), longJob)
cleanupJob(testInstance, connectionPool, createdJobID)
// Der Executor blockiert, bis der Kontext abgebrochen wird.
testExecutor := &recordingExecutor{blockUntil: make(chan struct{})}
executionLoop := newTestLoop(testInstance, connectionPool, testExecutor)
loopContext, cancelLoop := context.WithCancel(context.Background())
loopDone := make(chan error, 1)
go func() {
loopDone <- executionLoop.Run(loopContext)
}()
// Warten, bis der Lauf begonnen hat.
startDeadline := time.Now().Add(5 * time.Second)
for time.Now().Before(startDeadline) && testExecutor.count() == 0 {
time.Sleep(20 * time.Millisecond)
}
if testExecutor.count() == 0 {
cancelLoop()
testInstance.Fatal("der Lauf begann nicht")
}
cancelLoop()
select {
case shutdownError := <-loopDone:
if shutdownError != nil {
testInstance.Errorf("das Beenden meldete einen Fehler: %v", shutdownError)
}
case <-time.After(6 * time.Second):
testInstance.Fatal("die Schleife endete nicht")
}
// Der entscheidende Punkt: Kein Lauf bleibt auf „running" stehen.
remainingRuns, _, _ := testStore.ListRuns(context.Background(), createdJobID, 1, 10)
for _, singleRun := range remainingRuns {
if !singleRun.IsFinished() {
testInstance.Fatalf("nach dem Beenden steht ein Lauf noch auf %q; der Auftrag bleibt blockiert", singleRun.Status)
}
if singleRun.Status != RunStatusCancelled {
testInstance.Errorf("der abgebrochene Lauf wurde als %q vermerkt", singleRun.Status)
}
}
}
// TestLoopSkipsMissedRunsWithoutBackfill prüft das Nachholverhalten.
//
// 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.
func TestLoopSkipsMissedRunsWithoutBackfill(testInstance *testing.T) {
connectionPool := connectTestDatabase(testInstance)
testStore := NewPostgresStore(connectionPool)
repositoryID := createTestRepository(testInstance, connectionPool)
// Der Zeitpunkt liegt drei Tage zurück; der Plan läuft täglich.
longAgo := time.Now().UTC().Add(-72 * time.Hour)
staleJob := buildStoredJob(repositoryID)
staleJob.Schedule = scheduler.Schedule{ScheduleType: scheduler.ScheduleTypeDaily, Hour: 2}
staleJob.NextRunAt = &longAgo
createdJobID, _ := testStore.CreateJob(context.Background(), staleJob)
cleanupJob(testInstance, connectionPool, createdJobID)
testExecutor := &recordingExecutor{}
executionLoop := newTestLoop(testInstance, connectionPool, testExecutor)
runLoopUntil(testInstance, executionLoop, func() bool {
return testExecutor.count() > 0
}, 5*time.Second)
// Kurz warten, damit weitere Läufe Zeit hätten zu erscheinen.
time.Sleep(300 * time.Millisecond)
if testExecutor.count() != 1 {
testInstance.Fatalf("der Auftrag lief %d mal; versäumte Läufe wurden nachgeholt", testExecutor.count())
}
// Der nächste Zeitpunkt muss in der Zukunft liegen.
updatedJob, _ := testStore.GetJob(context.Background(), createdJobID)
if updatedJob.NextRunAt == nil || !updatedJob.NextRunAt.After(time.Now().UTC()) {
testInstance.Fatalf("der nächste Zeitpunkt liegt nicht in der Zukunft: %v", updatedJob.NextRunAt)
}
}
// TestNotImplementedExecutorFailsLoudly prüft den Standard ohne Ausführung.
//
// Ein Executor, der stillschweigend Erfolg meldete, wäre das gefährlichste
// Fake-Feature der Anlage: Die Oberfläche zeigte grüne Läufe, und im Repository
// läge nichts.
func TestNotImplementedExecutorFailsLoudly(testInstance *testing.T) {
notImplementedExecutor := &NotImplementedExecutor{}
_, executionError := notImplementedExecutor.Execute(context.Background(), ExecutionRequest{})
if executionError == nil {
testInstance.Fatal("ein nicht eingerichteter Executor meldete Erfolg")
}
var typedError *ExecutionError
if !errors.As(executionError, &typedError) {
testInstance.Fatalf("der Fehler war nicht eingeordnet: %v", executionError)
}
// Ein Konfigurationsfehler wird nicht wiederholt — richtig so: Er behebt
// sich nicht durch Warten.
if typedError.FailureClass.IsRetryable() {
testInstance.Error("ein fehlender Executor wurde als wiederholbar eingeordnet")
}
}
// TestLoopOptionsRejectShortStaleTimeout prüft die Einstellungsprüfung.
//
// 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.
func TestLoopOptionsRejectShortStaleTimeout(testInstance *testing.T) {
testLogger := slog.New(slog.NewJSONHandler(io.Discard, nil))
_, loopError := NewLoop(nil, nil, LoopOptions{
HeartbeatInterval: time.Minute,
StaleRunTimeout: 30 * time.Second,
}, testLogger)
if loopError == nil {
testInstance.Fatal("eine zu kurze Frist für verwaiste Läufe wurde angenommen")
}
}