package recovery import ( "context" "errors" "os" "path/filepath" "sync" "testing" "time" "github.com/google/uuid" "github.com/jackc/pgx/v5/pgxpool" ) // connectTestDatabase oeffnet die Testdatenbank. func connectTestDatabase(testInstance *testing.T) *pgxpool.Pool { testInstance.Helper() connectionString := os.Getenv("SYNCOVA_TEST_DATABASE_URL") if connectionString == "" { testInstance.Skip("SYNCOVA_TEST_DATABASE_URL ist nicht gesetzt; die Datenbanktests werden uebersprungen") } connectContext, cancelConnect := context.WithTimeout(context.Background(), 5*time.Second) defer cancelConnect() connectionPool, poolError := pgxpool.New(connectContext, connectionString) if poolError != nil { testInstance.Skipf("die Testdatenbank war nicht erreichbar: %v", poolError) } if pingError := connectionPool.Ping(connectContext); pingError != nil { connectionPool.Close() testInstance.Skipf("die Testdatenbank antwortete nicht: %v", pingError) } testInstance.Cleanup(connectionPool.Close) return connectionPool } // seedBackupRow legt Repository, Auftrag, Lauf und Backup in der Datenbank an. // // Der Wiederherstellungsauftrag verweist auf ein Backup; ohne diese Kette // haelt der Fremdschluessel nicht. func seedBackupRow(testInstance *testing.T, connectionPool *pgxpool.Pool, repositoryPath string) uuid.UUID { testInstance.Helper() backgroundContext := context.Background() uniqueSuffix := uuid.NewString()[:8] var repositoryID uuid.UUID if scanError := connectionPool.QueryRow(backgroundContext, `INSERT INTO repositories (name, location) VALUES ($1, $2) RETURNING id`, "rec-"+uniqueSuffix, repositoryPath).Scan(&repositoryID); scanError != nil { testInstance.Fatalf("das Repository liess sich nicht eintragen: %v", scanError) } var jobID uuid.UUID if scanError := connectionPool.QueryRow(backgroundContext, `INSERT INTO backup_jobs (name, schedule_type, repository_id) VALUES ($1,'manual',$2) RETURNING id`, "rec-job-"+uniqueSuffix, repositoryID).Scan(&jobID); scanError != nil { testInstance.Fatalf("der Auftrag liess sich nicht eintragen: %v", scanError) } var runID uuid.UUID if scanError := connectionPool.QueryRow(backgroundContext, `INSERT INTO backup_job_runs (job_id, status, completed_at, correlation_id) VALUES ($1,'succeeded',now(),gen_random_uuid()) RETURNING id`, jobID).Scan(&runID); scanError != nil { testInstance.Fatalf("der Lauf liess sich nicht eintragen: %v", scanError) } var backupID uuid.UUID if scanError := connectionPool.QueryRow(backgroundContext, `INSERT INTO backups (job_run_id, repository_id, backup_id_in_repository, backup_type, status, manifest_ref, completed_at) VALUES ($1,$2,'pruef-backup','full','complete','manifests/pruef-backup.json',now()) RETURNING id`, runID, repositoryID).Scan(&backupID); scanError != nil { testInstance.Fatalf("das Backup liess sich nicht eintragen: %v", scanError) } testInstance.Cleanup(func() { cleanupContext := context.Background() _, _ = connectionPool.Exec(cleanupContext, `DELETE FROM restore_sessions WHERE restore_job_id IN (SELECT id FROM restore_jobs WHERE backup_id = $1)`, backupID) _, _ = connectionPool.Exec(cleanupContext, `DELETE FROM restore_jobs WHERE backup_id = $1`, backupID) _, _ = connectionPool.Exec(cleanupContext, `DELETE FROM backups WHERE id = $1`, backupID) _, _ = connectionPool.Exec(cleanupContext, `DELETE FROM backup_job_runs WHERE id = $1`, runID) _, _ = connectionPool.Exec(cleanupContext, `DELETE FROM backup_jobs WHERE id = $1`, jobID) _, _ = connectionPool.Exec(cleanupContext, `DELETE FROM repositories WHERE id = $1`, repositoryID) }) return backupID } // recordingRestoreExecutor merkt sich die ausgefuehrten Wiederherstellungen. type recordingRestoreExecutor struct { // mutex schuetzt den Zustand. mutex sync.Mutex // executionCount zaehlt die Aufrufe. executionCount int // observedResumePaths haelt die gemeldeten Fortsetzungsmarken. observedResumePaths []string // entriesToReport sind die Objekte, die als fertig gemeldet werden. entriesToReport []string // resultToReturn ist das gemeldete Ergebnis. resultToReturn ExecutionResult // errorToReturn ist der gemeldete Fehler. errorToReturn error // blockUntil haelt die Ausfuehrung an, bis der Kanal schliesst. blockUntil chan struct{} } // Execute erfuellt die Executor-Schnittstelle. func (executor *recordingRestoreExecutor) Execute(executionContext context.Context, executionRequest ExecutionRequest) (ExecutionResult, error) { executor.mutex.Lock() executor.executionCount++ executor.observedResumePaths = append(executor.observedResumePaths, executionRequest.ResumeAfterPath) entriesToReport := executor.entriesToReport blockChannel := executor.blockUntil resultToReturn := executor.resultToReturn errorToReturn := executor.errorToReturn executor.mutex.Unlock() for _, entryPath := range entriesToReport { if executionRequest.EntryRestoredCallback != nil { executionRequest.EntryRestoredCallback(entryPath, 1024) } } if blockChannel != nil { select { case <-blockChannel: case <-executionContext.Done(): return ExecutionResult{}, executionContext.Err() } } return resultToReturn, errorToReturn } // count liefert die Zahl der Aufrufe. func (executor *recordingRestoreExecutor) count() int { executor.mutex.Lock() defer executor.mutex.Unlock() return executor.executionCount } // lastResumePath liefert die zuletzt gemeldete Fortsetzungsmarke. func (executor *recordingRestoreExecutor) lastResumePath() string { executor.mutex.Lock() defer executor.mutex.Unlock() if len(executor.observedResumePaths) == 0 { return "" } return executor.observedResumePaths[len(executor.observedResumePaths)-1] } // newTestLoop baut eine Schleife mit kurzem Takt. func newTestLoop(testInstance *testing.T, connectionPool *pgxpool.Pool, executor Executor) *Loop { testInstance.Helper() builtLoop, loopError := NewLoop(NewStore(connectionPool), executor, LoopOptions{ InstanceName: "test-" + uuid.NewString()[:8], TickInterval: 50 * time.Millisecond, HeartbeatInterval: 100 * time.Millisecond, StaleRestoreTimeout: time.Second, ShutdownGracePeriod: 3 * time.Second, MaximumConcurrentRestores: 2, CheckpointInterval: time.Millisecond, }, discardLogger()) if loopError != nil { testInstance.Fatalf("die Schleife liess sich nicht bauen: %v", loopError) } return builtLoop } // runLoopUntil fuehrt die Schleife aus, bis die Bedingung erfuellt 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") } } // TestRestoreLoopExecutesQueuedJob prueft den Grundfall. func TestRestoreLoopExecutesQueuedJob(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewStore(connectionPool) backupID := seedBackupRow(testInstance, connectionPool, testInstance.TempDir()) createdJob, createError := testStore.CreateRestore(context.Background(), CreateRequest{ BackupID: backupID, SourceType: "filesystem", TargetType: TargetFilesystem, TargetRef: filepath.Join(testInstance.TempDir(), "ziel-"+uuid.NewString()[:8]), }) if createError != nil { testInstance.Fatalf("der Auftrag liess sich nicht anlegen: %v", createError) } testExecutor := &recordingRestoreExecutor{ resultToReturn: ExecutionResult{BytesRestored: 4096, FilesRestored: 12}, } executionLoop := newTestLoop(testInstance, connectionPool, testExecutor) runLoopUntil(testInstance, executionLoop, func() bool { loadedJob, _ := testStore.GetRestore(context.Background(), createdJob.ID) return loadedJob != nil && loadedJob.Status.IsFinished() }, 5*time.Second) finishedJob, readError := testStore.GetRestore(context.Background(), createdJob.ID) if readError != nil { testInstance.Fatalf("der Auftrag liess sich nicht lesen: %v", readError) } if finishedJob.Status != RestoreStatusSucceeded { testInstance.Fatalf("der Auftrag endete mit %q: %s", finishedJob.Status, finishedJob.ErrorMessage) } if finishedJob.BytesRestored != 4096 || finishedJob.FilesRestored != 12 { testInstance.Errorf("die Kennzahlen kamen nicht an: %d Byte, %d Objekte", finishedJob.BytesRestored, finishedJob.FilesRestored) } // Eine abgeschlossene Wiederherstellung schliesst ihre Sitzung. activeSession, _ := testStore.GetActiveSession(context.Background(), createdJob.ID) if activeSession != nil { testInstance.Error("die Sitzung blieb nach erfolgreicher Wiederherstellung offen") } } // TestRestoreLoopTreatsSkippedAsPartialFailure ist der Ehrlichkeitstest. // // Eine Wiederherstellung, die zwei Dateien nicht zurueckschreiben konnte, ist // unvollstaendig — der Anwender glaubte sonst, alles sei wieder da. func TestRestoreLoopTreatsSkippedAsPartialFailure(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewStore(connectionPool) backupID := seedBackupRow(testInstance, connectionPool, testInstance.TempDir()) createdJob, _ := testStore.CreateRestore(context.Background(), CreateRequest{ BackupID: backupID, SourceType: "filesystem", TargetType: TargetFilesystem, TargetRef: filepath.Join(testInstance.TempDir(), "ziel-"+uuid.NewString()[:8]), }) testExecutor := &recordingRestoreExecutor{ resultToReturn: ExecutionResult{ BytesRestored: 2048, FilesRestored: 8, FilesSkipped: 2, SkipReasons: []string{"2 vorhandene Objekte wurden nicht ueberschrieben."}, }, } executionLoop := newTestLoop(testInstance, connectionPool, testExecutor) runLoopUntil(testInstance, executionLoop, func() bool { loadedJob, _ := testStore.GetRestore(context.Background(), createdJob.ID) return loadedJob != nil && loadedJob.Status.IsFinished() }, 5*time.Second) finishedJob, _ := testStore.GetRestore(context.Background(), createdJob.ID) if finishedJob.Status != RestoreStatusPartialFailure { testInstance.Fatalf("eine Wiederherstellung mit 2 uebergangenen Objekten wurde als %q gewertet", finishedJob.Status) } if finishedJob.ErrorMessage == "" { testInstance.Error("der Teilfehler wurde nicht begruendet") } } // TestCheckpointEnablesResume ist der wichtigste Test dieser Datei. // // Ein Abbruch bei 90 Prozent darf nicht bedeuten, dass alles von vorn beginnt. func TestCheckpointEnablesResume(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewStore(connectionPool) backupID := seedBackupRow(testInstance, connectionPool, testInstance.TempDir()) createdJob, _ := testStore.CreateRestore(context.Background(), CreateRequest{ BackupID: backupID, SourceType: "filesystem", TargetType: TargetFilesystem, TargetRef: filepath.Join(testInstance.TempDir(), "ziel-"+uuid.NewString()[:8]), }) // Der erste Lauf meldet drei fertige Dateien und scheitert dann. firstExecutor := &recordingRestoreExecutor{ entriesToReport: []string{"a/eins.txt", "a/zwei.txt", "b/drei.txt"}, errorToReturn: &ExecutionError{Code: "IO_ERROR", Message: "Der Datentraeger meldete einen Fehler."}, } firstLoop := newTestLoop(testInstance, connectionPool, firstExecutor) runLoopUntil(testInstance, firstLoop, func() bool { loadedJob, _ := testStore.GetRestore(context.Background(), createdJob.ID) return loadedJob != nil && loadedJob.Status.IsFinished() }, 5*time.Second) // Die Sitzung bleibt offen — ihr Pruefpunkt ist die Grundlage der Fortsetzung. activeSession, sessionError := testStore.GetActiveSession(context.Background(), createdJob.ID) if sessionError != nil { testInstance.Fatalf("die Sitzung liess sich nicht lesen: %v", sessionError) } if activeSession == nil { testInstance.Fatal("die Sitzung wurde nach dem Fehlschlag geschlossen; eine Fortsetzung ist unmoeglich") } if activeSession.Checkpoint.LastCompletedPath != "b/drei.txt" { testInstance.Fatalf("der Pruefpunkt steht auf %q, erwartet war b/drei.txt", activeSession.Checkpoint.LastCompletedPath) } if activeSession.Checkpoint.FilesCompleted != 3 { testInstance.Errorf("der Pruefpunkt zaehlt %d fertige Objekte", activeSession.Checkpoint.FilesCompleted) } // Der Auftrag wird erneut eingereiht — die Fortsetzung. if _, execError := connectionPool.Exec(context.Background(), `UPDATE restore_jobs SET status = 'queued', completed_at = NULL WHERE id = $1`, createdJob.ID); execError != nil { testInstance.Fatalf("der Auftrag liess sich nicht erneut einreihen: %v", execError) } secondExecutor := &recordingRestoreExecutor{ resultToReturn: ExecutionResult{BytesRestored: 1024, FilesRestored: 2}, } secondLoop := newTestLoop(testInstance, connectionPool, secondExecutor) runLoopUntil(testInstance, secondLoop, func() bool { loadedJob, _ := testStore.GetRestore(context.Background(), createdJob.ID) return loadedJob != nil && loadedJob.Status.IsFinished() }, 5*time.Second) // Der Kern: Der zweite Lauf bekommt die Marke und beginnt nicht von vorn. if secondExecutor.lastResumePath() != "b/drei.txt" { testInstance.Fatalf("die Fortsetzung erhielt die Marke %q; sie begaenne von vorn", secondExecutor.lastResumePath()) } finishedJob, _ := testStore.GetRestore(context.Background(), createdJob.ID) // Die Endzahlen beschreiben den gesamten Vorgang, nicht nur den letzten // Versuch: drei Objekte aus dem ersten Lauf plus zwei aus dem zweiten. if finishedJob.FilesRestored != 5 { testInstance.Errorf("es wurden %d Objekte gezaehlt, erwartet waren 5 (3 + 2)", finishedJob.FilesRestored) } } // TestRestoreLoopDoesNotRetryAutomatically prueft den bewussten Verzicht. // // Eine halb geschriebene Wiederherstellung erneut zu starten kann Daten // beschaedigen, die der erste Versuch bereits am Platz hatte. func TestRestoreLoopDoesNotRetryAutomatically(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewStore(connectionPool) backupID := seedBackupRow(testInstance, connectionPool, testInstance.TempDir()) createdJob, _ := testStore.CreateRestore(context.Background(), CreateRequest{ BackupID: backupID, SourceType: "filesystem", TargetType: TargetFilesystem, TargetRef: filepath.Join(testInstance.TempDir(), "ziel-"+uuid.NewString()[:8]), }) testExecutor := &recordingRestoreExecutor{ errorToReturn: &ExecutionError{Code: "IO_ERROR", Message: "Der Datentraeger meldete einen Fehler."}, } executionLoop := newTestLoop(testInstance, connectionPool, testExecutor) runLoopUntil(testInstance, executionLoop, func() bool { loadedJob, _ := testStore.GetRestore(context.Background(), createdJob.ID) return loadedJob != nil && loadedJob.Status.IsFinished() }, 3*time.Second) // Kurz warten, damit eine faelschliche Wiederholung Zeit haette zu erscheinen. time.Sleep(300 * time.Millisecond) if testExecutor.count() != 1 { testInstance.Fatalf("die gescheiterte Wiederherstellung wurde %d mal ausgefuehrt", testExecutor.count()) } } // TestRestoreLoopShutdownMarksCancelled prueft das geordnete Beenden. func TestRestoreLoopShutdownMarksCancelled(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewStore(connectionPool) backupID := seedBackupRow(testInstance, connectionPool, testInstance.TempDir()) createdJob, _ := testStore.CreateRestore(context.Background(), CreateRequest{ BackupID: backupID, SourceType: "filesystem", TargetType: TargetFilesystem, TargetRef: filepath.Join(testInstance.TempDir(), "ziel-"+uuid.NewString()[:8]), }) testExecutor := &recordingRestoreExecutor{ entriesToReport: []string{"a/eins.txt"}, 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) }() 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("die Wiederherstellung begann nicht") } cancelLoop() select { case <-loopDone: case <-time.After(6 * time.Second): testInstance.Fatal("die Schleife endete nicht") } finishedJob, _ := testStore.GetRestore(context.Background(), createdJob.ID) if !finishedJob.Status.IsFinished() { testInstance.Fatalf("nach dem Beenden steht der Auftrag auf %q", finishedJob.Status) } if finishedJob.Status != RestoreStatusCancelled { testInstance.Errorf("der abgebrochene Auftrag wurde als %q vermerkt", finishedJob.Status) } // Der Pruefpunkt muss den Abbruch ueberleben — nach einem Abbruch ist er am // wertvollsten. activeSession, _ := testStore.GetActiveSession(context.Background(), createdJob.ID) if activeSession == nil || activeSession.Checkpoint.LastCompletedPath != "a/eins.txt" { testInstance.Error("der Pruefpunkt ueberlebte den Abbruch nicht") } } // TestSecondRestoreToSameTargetIsRejected prueft den Schutz gegen zwei // gleichzeitige Wiederherstellungen in dasselbe Ziel. // // Sie schrieben sich gegenseitig zu — und zwar ohne dass es auffiele, weil // beide erfolgreich endeten. func TestSecondRestoreToSameTargetIsRejected(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewStore(connectionPool) backupID := seedBackupRow(testInstance, connectionPool, testInstance.TempDir()) sharedTarget := filepath.Join(testInstance.TempDir(), "ziel-"+uuid.NewString()[:8]) if _, createError := testStore.CreateRestore(context.Background(), CreateRequest{ BackupID: backupID, SourceType: "filesystem", TargetType: TargetFilesystem, TargetRef: sharedTarget, }); createError != nil { testInstance.Fatalf("der erste Auftrag liess sich nicht anlegen: %v", createError) } _, duplicateError := testStore.CreateRestore(context.Background(), CreateRequest{ BackupID: backupID, SourceType: "filesystem", TargetType: TargetFilesystem, TargetRef: sharedTarget, }) if !errors.Is(duplicateError, ErrTargetBusy) { testInstance.Fatalf("eine zweite Wiederherstellung in dasselbe Ziel wurde angenommen: %v", duplicateError) } } // TestReclaimStaleRestoreDoesNotRequeue prueft die bewusste Zurueckhaltung. // // Anders als eine Sicherung wird eine verwaiste Wiederherstellung **nicht** // erneut eingereiht: Der Betreiber soll hinsehen, bevor erneut geschrieben wird. func TestReclaimStaleRestoreDoesNotRequeue(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewStore(connectionPool) backupID := seedBackupRow(testInstance, connectionPool, testInstance.TempDir()) createdJob, _ := testStore.CreateRestore(context.Background(), CreateRequest{ BackupID: backupID, SourceType: "filesystem", TargetType: TargetFilesystem, TargetRef: filepath.Join(testInstance.TempDir(), "ziel-"+uuid.NewString()[:8]), }) // Ein Auftrag, dessen Server vor zwei Stunden verschwand. if _, execError := connectionPool.Exec(context.Background(), ` UPDATE restore_jobs SET status = 'running', started_at = now() - interval '2 hours', heartbeat_at = now() - interval '2 hours', scheduler_instance = 'toter-server' WHERE id = $1`, createdJob.ID); execError != nil { testInstance.Fatalf("der verwaiste Auftrag liess sich nicht herstellen: %v", execError) } reclaimedCount, reclaimError := testStore.ReclaimStaleRestores(context.Background(), time.Minute) if reclaimError != nil { testInstance.Fatalf("die Freigabe schlug fehl: %v", reclaimError) } if reclaimedCount != 1 { testInstance.Fatalf("es wurden %d Auftraege freigegeben, erwartet war 1", reclaimedCount) } reclaimedJob, _ := testStore.GetRestore(context.Background(), createdJob.ID) // Gescheitert, nicht erneut eingereiht. if reclaimedJob.Status != RestoreStatusFailed { testInstance.Errorf("der verwaiste Auftrag steht auf %q", reclaimedJob.Status) } if reclaimedJob.ErrorCode != "SCHEDULER_LOST" { testInstance.Errorf("der Abbruch wurde als %q vermerkt", reclaimedJob.ErrorCode) } // Die Sitzung bleibt bestehen, damit eine Fortsetzung moeglich ist. activeSession, _ := testStore.GetActiveSession(context.Background(), createdJob.ID) if activeSession == nil { testInstance.Error("die Sitzung wurde bei der Freigabe geschlossen; eine Fortsetzung ist unmoeglich") } } // TestLoopOptionsRejectShortStaleTimeout prueft die Einstellungspruefung. func TestLoopOptionsRejectShortStaleTimeout(testInstance *testing.T) { _, loopError := NewLoop(nil, nil, LoopOptions{ HeartbeatInterval: time.Minute, StaleRestoreTimeout: 30 * time.Second, }, discardLogger()) if loopError == nil { testInstance.Fatal("eine zu kurze Frist fuer verwaiste Auftraege wurde angenommen") } } // TestNotImplementedRestoreExecutorFailsLoudly prueft den Standard. func TestNotImplementedRestoreExecutorFailsLoudly(testInstance *testing.T) { notImplementedExecutor := &NotImplementedExecutor{} _, executionError := notImplementedExecutor.Execute(context.Background(), ExecutionRequest{}) if executionError == nil { testInstance.Fatal("ein nicht eingerichteter Executor meldete Erfolg") } }