package jobs import ( "context" "errors" "os" "strings" "sync" "testing" "time" "github.com/google/uuid" "github.com/jackc/pgx/v5/pgxpool" "github.com/syncova/syncova/packages/scheduler" ) // connectTestDatabase öffnet die Testdatenbank. // // Ohne erreichbare Datenbank werden die Tests übersprungen statt zu scheitern: // Die Rechenlogik ist ohne sie prüfbar, und ein rot leuchtender Testlauf auf // einem Rechner ohne Docker gewöhnt nur das Wegschauen an. func connectTestDatabase(testInstance *testing.T) *pgxpool.Pool { testInstance.Helper() connectionString := os.Getenv("SYNCOVA_TEST_DATABASE_URL") if connectionString == "" { connectionString = os.Getenv("SYNCOVA_DATABASE_URL") } if connectionString == "" { testInstance.Skip("SYNCOVA_TEST_DATABASE_URL ist nicht gesetzt; die Datenbanktests werden übersprungen") } 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 } // createTestRepository legt ein Repository für die Tests an. func createTestRepository(testInstance *testing.T, connectionPool *pgxpool.Pool) uuid.UUID { testInstance.Helper() repositoryName := "test-" + uuid.NewString() var repositoryID uuid.UUID if scanError := connectionPool.QueryRow(context.Background(), `INSERT INTO repositories (name, location) VALUES ($1, $2) RETURNING id`, repositoryName, "/tmp/"+repositoryName).Scan(&repositoryID); scanError != nil { testInstance.Fatalf("das Test-Repository ließ sich nicht anlegen: %v", scanError) } testInstance.Cleanup(func() { _, _ = connectionPool.Exec(context.Background(), `DELETE FROM repositories WHERE id = $1`, repositoryID) }) return repositoryID } // buildStoredJob liefert einen ablegbaren Auftrag. func buildStoredJob(repositoryID uuid.UUID) *Job { return &Job{ Name: "auftrag-" + uuid.NewString(), Description: "Ein Test", Status: JobStatusActive, Priority: scheduler.PriorityHigh, Schedule: scheduler.Schedule{ScheduleType: scheduler.ScheduleTypeDaily, Hour: 2, TimeZone: "Europe/Berlin"}, RepositoryID: repositoryID, Sources: []JobSource{ {SourceType: SourceTypeFilesystem, SourceID: "/daten", SourceName: "Dateiserver", ExcludePatterns: []string{"*.tmp", "cache"}}, // Bewusst keine Proxmox-Quelle: Die braucht seit Phase 7 einen // Verbund, und dieser Auftrag soll ohne Virtualisierungsumgebung // ablegbar bleiben. Der Verbundbezug wird eigens geprüft // (TestProxmoxSourceRequiresCluster). {SourceType: SourceTypeLinuxSystem, SourceID: "srv-01", SourceName: "Linux-Server"}, }, MaximumConcurrency: 1, RetryPolicy: scheduler.DefaultRetryPolicy(), BandwidthLimitBytesPerSecond: 50 << 20, RecoveryTimeObjective: 2 * time.Hour, } } // cleanupJob entfernt einen Auftrag samt Läufen nach dem Test. func cleanupJob(testInstance *testing.T, connectionPool *pgxpool.Pool, jobIdentifier uuid.UUID) { testInstance.Cleanup(func() { _, _ = connectionPool.Exec(context.Background(), `DELETE FROM backup_job_runs WHERE job_id = $1`, jobIdentifier) _, _ = connectionPool.Exec(context.Background(), `DELETE FROM backup_jobs WHERE id = $1`, jobIdentifier) }) } // TestCreateAndReadJobRoundTrip prüft die Ablage und das Wiederlesen. func TestCreateAndReadJobRoundTrip(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewPostgresStore(connectionPool) repositoryID := createTestRepository(testInstance, connectionPool) originalJob := buildStoredJob(repositoryID) createdJobID, createError := testStore.CreateJob(context.Background(), originalJob) if createError != nil { testInstance.Fatalf("der Auftrag ließ sich nicht anlegen: %v", createError) } cleanupJob(testInstance, connectionPool, createdJobID) readJob, readError := testStore.GetJob(context.Background(), createdJobID) if readError != nil { testInstance.Fatalf("der Auftrag ließ sich nicht lesen: %v", readError) } if readJob.Name != originalJob.Name { testInstance.Errorf("der Name kam als %q zurück", readJob.Name) } // Der Zeitplan muss vollständig überleben — besonders die Zeitzone. // Ginge sie verloren, liefe der Auftrag nach dem Neustart zu einer anderen // Uhrzeit. if readJob.Schedule.TimeZone != "Europe/Berlin" || readJob.Schedule.Hour != 2 { testInstance.Errorf("der Zeitplan kam verändert zurück: %+v", readJob.Schedule) } if readJob.Priority != scheduler.PriorityHigh { testInstance.Errorf("die Dringlichkeit kam als %q zurück", readJob.Priority) } if readJob.BandwidthLimitBytesPerSecond != 50<<20 { testInstance.Errorf("die Bandbreitengrenze kam als %d zurück", readJob.BandwidthLimitBytesPerSecond) } if len(readJob.Sources) != 2 { testInstance.Fatalf("es kamen %d Quellen zurück, abgelegt waren 2", len(readJob.Sources)) } // Die Ausschlussregeln sind der Teil, der beim Ablegen als JSON am leichtesten // verloren geht. var foundFilesystemSource bool for _, readSource := range readJob.Sources { if readSource.SourceType != SourceTypeFilesystem { continue } foundFilesystemSource = true if len(readSource.ExcludePatterns) != 2 || readSource.ExcludePatterns[0] != "*.tmp" { testInstance.Errorf("die Ausschlussregeln kamen als %v zurück", readSource.ExcludePatterns) } } if !foundFilesystemSource { testInstance.Error("die Dateisystemquelle fehlte") } } // TestDuplicateJobNameIsRejected prüft die Eindeutigkeit des Namens. func TestDuplicateJobNameIsRejected(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewPostgresStore(connectionPool) repositoryID := createTestRepository(testInstance, connectionPool) firstJob := buildStoredJob(repositoryID) createdJobID, createError := testStore.CreateJob(context.Background(), firstJob) if createError != nil { testInstance.Fatalf("der erste Auftrag ließ sich nicht anlegen: %v", createError) } cleanupJob(testInstance, connectionPool, createdJobID) secondJob := buildStoredJob(repositoryID) secondJob.Name = firstJob.Name if _, duplicateError := testStore.CreateJob(context.Background(), secondJob); !errors.Is(duplicateError, ErrJobNameTaken) { testInstance.Fatalf("ein doppelter Name wurde angenommen: %v", duplicateError) } } // TestCreateJobIsAtomic prüft die Transaktion. // // Ein Auftrag ohne seine Quellen wäre ein Auftrag, der erfolgreich durchliefe, // ohne etwas zu sichern. func TestCreateJobIsAtomic(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewPostgresStore(connectionPool) repositoryID := createTestRepository(testInstance, connectionPool) // Eine Abhängigkeit auf einen nicht vorhandenen Auftrag lässt den // Fremdschlüssel scheitern — nachdem der Auftrag selbst bereits eingefügt // wurde. brokenJob := buildStoredJob(repositoryID) brokenJob.DependsOnJobIDs = []uuid.UUID{uuid.New()} if _, createError := testStore.CreateJob(context.Background(), brokenJob); createError == nil { testInstance.Fatal("eine Abhängigkeit auf einen unbekannten Auftrag wurde angenommen") } // Der Auftrag darf nicht zurückgeblieben sein. var remainingCount int if scanError := connectionPool.QueryRow(context.Background(), `SELECT COUNT(*) FROM backup_jobs WHERE name = $1`, brokenJob.Name).Scan(&remainingCount); scanError != nil { testInstance.Fatalf("die Zählung schlug fehl: %v", scanError) } if remainingCount != 0 { testInstance.Fatal("der gescheiterte Anlagevorgang hinterließ einen Auftrag ohne Abhängigkeiten") } } // TestClaimDueJobsNeverHandsOutSameJobTwice ist der wichtigste Datenbanktest. // // Zwei Control-Server sehen denselben fälligen Auftrag. Ohne Sperre starten // beide ihn — und es entstehen zwei gleichzeitige Sicherungen derselben Quelle // auf dasselbe Repository. func TestClaimDueJobsNeverHandsOutSameJobTwice(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewPostgresStore(connectionPool) repositoryID := createTestRepository(testInstance, connectionPool) dueTime := time.Now().UTC().Add(-time.Minute) // Zehn fällige Aufträge. const jobCount = 10 for jobIndex := 0; jobIndex < jobCount; jobIndex++ { dueJob := buildStoredJob(repositoryID) 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) } // Vier Server greifen gleichzeitig zu. const schedulerCount = 4 claimedByScheduler := make([][]uuid.UUID, schedulerCount) var waitGroup sync.WaitGroup waitGroup.Add(schedulerCount) for schedulerIndex := 0; schedulerIndex < schedulerCount; schedulerIndex++ { go func(instanceIndex int) { defer waitGroup.Done() claimedJobs, claimError := testStore.ClaimDueJobs(context.Background(), time.Now().UTC(), jobCount, "server-"+uuid.NewString()) if claimError != nil { testInstance.Errorf("die Übernahme schlug fehl: %v", claimError) return } claimedByScheduler[instanceIndex] = claimedJobs }(schedulerIndex) } waitGroup.Wait() // Der Kern: Kein Auftrag darf zwei Läufe haben. var runsPerJob int if scanError := connectionPool.QueryRow(context.Background(), ` SELECT COALESCE(MAX(run_count), 0) FROM ( SELECT COUNT(*) AS run_count FROM backup_job_runs WHERE job_id IN (SELECT id FROM backup_jobs WHERE repository_id = $1) GROUP BY job_id ) counted`, repositoryID).Scan(&runsPerJob); scanError != nil { testInstance.Fatalf("die Zählung schlug fehl: %v", scanError) } if runsPerJob > 1 { testInstance.Fatalf("ein Auftrag wurde %d mal gleichzeitig gestartet", runsPerJob) } // Und alle Aufträge müssen genau einmal übernommen worden sein. var totalRuns int if scanError := connectionPool.QueryRow(context.Background(), ` SELECT COUNT(*) FROM backup_job_runs WHERE job_id IN (SELECT id FROM backup_jobs WHERE repository_id = $1)`, repositoryID).Scan(&totalRuns); scanError != nil { testInstance.Fatalf("die Zählung schlug fehl: %v", scanError) } if totalRuns != jobCount { testInstance.Errorf("es entstanden %d Läufe für %d Aufträge", totalRuns, jobCount) } } // TestClaimDueJobsSkipsRunningJobs prüft, dass ein laufender Auftrag ruht. func TestClaimDueJobsSkipsRunningJobs(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewPostgresStore(connectionPool) repositoryID := createTestRepository(testInstance, connectionPool) dueTime := time.Now().UTC().Add(-time.Minute) dueJob := buildStoredJob(repositoryID) 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) firstClaim, firstError := testStore.ClaimDueJobs(context.Background(), time.Now().UTC(), 10, "server-1") if firstError != nil { testInstance.Fatalf("die erste Übernahme schlug fehl: %v", firstError) } if len(firstClaim) != 1 { testInstance.Fatalf("die erste Übernahme lieferte %d Aufträge", len(firstClaim)) } // Der zweite Durchgang darf nichts mehr finden: Der Lauf steht noch auf // queued. secondClaim, secondError := testStore.ClaimDueJobs(context.Background(), time.Now().UTC(), 10, "server-2") if secondError != nil { testInstance.Fatalf("die zweite Übernahme schlug fehl: %v", secondError) } if len(secondClaim) != 0 { testInstance.Fatalf("ein bereits laufender Auftrag wurde erneut übernommen: %v", secondClaim) } } // TestClaimDueJobsIgnoresPausedJobs prüft die Aussetzung. func TestClaimDueJobsIgnoresPausedJobs(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewPostgresStore(connectionPool) repositoryID := createTestRepository(testInstance, connectionPool) dueTime := time.Now().UTC().Add(-time.Minute) pausedJob := buildStoredJob(repositoryID) pausedJob.NextRunAt = &dueTime createdJobID, createError := testStore.CreateJob(context.Background(), pausedJob) if createError != nil { testInstance.Fatalf("der Auftrag ließ sich nicht anlegen: %v", createError) } cleanupJob(testInstance, connectionPool, createdJobID) if statusError := testStore.SetJobStatus(context.Background(), createdJobID, JobStatusPaused, nil); statusError != nil { testInstance.Fatalf("die Aussetzung schlug fehl: %v", statusError) } claimedJobs, claimError := testStore.ClaimDueJobs(context.Background(), time.Now().UTC(), 10, "server-1") if claimError != nil { testInstance.Fatalf("die Übernahme schlug fehl: %v", claimError) } for _, claimedJobID := range claimedJobs { if claimedJobID == createdJobID { testInstance.Fatal("ein ausgesetzter Auftrag wurde übernommen") } } // Nach dem Fortsetzen muss er wieder laufen. if statusError := testStore.SetJobStatus(context.Background(), createdJobID, JobStatusActive, nil); statusError != nil { testInstance.Fatalf("das Fortsetzen schlug fehl: %v", statusError) } resumedClaim, resumeError := testStore.ClaimDueJobs(context.Background(), time.Now().UTC(), 10, "server-1") if resumeError != nil { testInstance.Fatalf("die Übernahme schlug fehl: %v", resumeError) } var wasResumed bool for _, claimedJobID := range resumedClaim { if claimedJobID == createdJobID { wasResumed = true } } if !wasResumed { testInstance.Fatal("ein fortgesetzter Auftrag wurde nicht übernommen") } } // TestSoftDeleteKeepsHistory prüft die weiche Löschung. // // Ein hart gelöschter Auftrag risse seine Läufe mit — und damit den Nachweis, // dass gesichert wurde. func TestSoftDeleteKeepsHistory(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewPostgresStore(connectionPool) repositoryID := createTestRepository(testInstance, connectionPool) deletableJob := buildStoredJob(repositoryID) createdJobID, createError := testStore.CreateJob(context.Background(), deletableJob) if createError != nil { testInstance.Fatalf("der Auftrag ließ sich nicht anlegen: %v", createError) } cleanupJob(testInstance, connectionPool, createdJobID) // Ein abgeschlossener Lauf als Historie. if _, execError := connectionPool.Exec(context.Background(), ` INSERT INTO backup_job_runs (job_id, status, completed_at, correlation_id) VALUES ($1, 'succeeded', now(), gen_random_uuid())`, createdJobID); execError != nil { testInstance.Fatalf("der Lauf ließ sich nicht anlegen: %v", execError) } if deleteError := testStore.SoftDeleteJob(context.Background(), createdJobID); deleteError != nil { testInstance.Fatalf("die Löschung schlug fehl: %v", deleteError) } // Der Auftrag ist weg … if _, readError := testStore.GetJob(context.Background(), createdJobID); !errors.Is(readError, ErrJobNotFound) { testInstance.Errorf("ein gelöschter Auftrag war noch lesbar: %v", readError) } // … seine Historie aber nicht. var remainingRuns int if scanError := connectionPool.QueryRow(context.Background(), `SELECT COUNT(*) FROM backup_job_runs WHERE job_id = $1`, createdJobID).Scan(&remainingRuns); scanError != nil { testInstance.Fatalf("die Zählung schlug fehl: %v", scanError) } if remainingRuns != 1 { testInstance.Fatalf("die Löschung nahm die Historie mit: %d Läufe übrig", remainingRuns) } } // TestListJobsPaginatesAndFilters prüft Liste und Filter. func TestListJobsPaginatesAndFilters(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewPostgresStore(connectionPool) repositoryID := createTestRepository(testInstance, connectionPool) const jobCount = 5 for jobIndex := 0; jobIndex < jobCount; jobIndex++ { listedJob := buildStoredJob(repositoryID) createdJobID, createError := testStore.CreateJob(context.Background(), listedJob) if createError != nil { testInstance.Fatalf("der Auftrag ließ sich nicht anlegen: %v", createError) } cleanupJob(testInstance, connectionPool, createdJobID) } listedJobs, totalCount, listError := testStore.ListJobs(context.Background(), ListFilter{ RepositoryID: &repositoryID, Page: 1, PageSize: 2, }) if listError != nil { testInstance.Fatalf("die Liste schlug fehl: %v", listError) } if len(listedJobs) != 2 { testInstance.Errorf("die Seite enthielt %d Aufträge, angefordert waren 2", len(listedJobs)) } // Die Gesamtzahl muss den gesamten Bestand nennen, nicht die Seite. if totalCount != jobCount { testInstance.Errorf("die Gesamtzahl war %d, erwartet waren %d", totalCount, jobCount) } // Die Quellen müssen auch in der Liste mitkommen. if len(listedJobs[0].Sources) == 0 { testInstance.Error("die Aufträge der Liste kamen ohne ihre Quellen") } } // TestListJobsCapsPageSize prüft die Obergrenze. // // Ohne sie könnte ein Aufrufer mit page_size=1000000 die gesamte Tabelle in den // Speicher des Dienstes ziehen. func TestListJobsCapsPageSize(testInstance *testing.T) { unboundedFilter := ListFilter{Page: 0, PageSize: 1_000_000} unboundedFilter.normalize() if unboundedFilter.PageSize != maximumPageSize { testInstance.Errorf("die Seitengröße wurde auf %d begrenzt, erwartet waren %d", unboundedFilter.PageSize, maximumPageSize) } if unboundedFilter.Page != 1 { testInstance.Errorf("die Seitennummer wurde auf %d gesetzt, erwartet war 1", unboundedFilter.Page) } } // TestGetMissingJobReportsNotFound prüft die Fehlermeldung. func TestGetMissingJobReportsNotFound(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewPostgresStore(connectionPool) _, readError := testStore.GetJob(context.Background(), uuid.New()) if !errors.Is(readError, ErrJobNotFound) { testInstance.Fatalf("ein fehlender Auftrag wurde nicht als solcher gemeldet: %v", readError) } if !strings.Contains(readError.Error(), "nicht gefunden") { testInstance.Errorf("die Fehlermeldung ist unverständlich: %v", readError) } } // TestProxmoxSourceRequiresCluster prueft den Verbundbezug einer Gastquelle. // // Bei genau einem eingerichteten Verbund liesse sich der Bezug raten — bei // zweien sicherte der Lauf die falsche Maschine. Die Regel steht als CHECK in // der Datenbank, weil im Code jede schreibende Stelle sie einhalten muesste und // eine es vergisst. func TestProxmoxSourceRequiresCluster(testInstance *testing.T) { connectionPool := connectTestDatabase(testInstance) testStore := NewPostgresStore(connectionPool) repositoryID := createTestRepository(testInstance, connectionPool) jobWithoutCluster := buildStoredJob(repositoryID) jobWithoutCluster.Sources = []JobSource{ {SourceType: SourceTypeProxmoxVM, SourceID: "qemu/100", SourceName: "web-01"}, } if _, createError := testStore.CreateJob(context.Background(), jobWithoutCluster); createError == nil { testInstance.Fatal("eine Proxmox-Quelle ohne Verbund wurde angelegt") } clusterIdentifier := createTestCluster(testInstance, connectionPool) jobWithCluster := buildStoredJob(repositoryID) jobWithCluster.Sources = []JobSource{ {SourceType: SourceTypeProxmoxVM, SourceID: "qemu/100", SourceName: "web-01", ClusterID: &clusterIdentifier}, } createdIdentifier, createError := testStore.CreateJob(context.Background(), jobWithCluster) if createError != nil { testInstance.Fatalf("der Auftrag ließ sich nicht anlegen: %v", createError) } cleanupJob(testInstance, connectionPool, createdIdentifier) readJob, readError := testStore.GetJob(context.Background(), createdIdentifier) if readError != nil { testInstance.Fatalf("der Auftrag ließ sich nicht lesen: %v", readError) } if len(readJob.Sources) != 1 || readJob.Sources[0].ClusterID == nil { testInstance.Fatalf("der Verbundbezug ging beim Lesen verloren: %+v", readJob.Sources) } if *readJob.Sources[0].ClusterID != clusterIdentifier { testInstance.Errorf("der Verbund ist %s statt %s", *readJob.Sources[0].ClusterID, clusterIdentifier) } } // createTestCluster legt einen Virtualisierungsverbund fuer die Tests an. // // Ueber rohes SQL statt ueber das Paket hypervisor: Dieses Paket darf nicht von // ihm abhaengen — die Abhaengigkeit laeuft in die andere Richtung. func createTestCluster(testInstance *testing.T, connectionPool *pgxpool.Pool) uuid.UUID { testInstance.Helper() var clusterIdentifier uuid.UUID if scanError := connectionPool.QueryRow(context.Background(), ` INSERT INTO proxmox_clusters (name, api_endpoint, api_token_id, api_token_ciphertext, api_token_key_version, backup_storage_id, archive_transport, archive_mount_roots) VALUES ($1, 'https://pve.example:8006', 'syncova@pve!backup', '\x00', 'v1', 'local', 'local', '{"local":"/var/lib/vz"}') RETURNING id`, "verbund-"+uuid.NewString()).Scan(&clusterIdentifier); scanError != nil { testInstance.Fatalf("der Verbund ließ sich nicht anlegen: %v", scanError) } testInstance.Cleanup(func() { _, _ = connectionPool.Exec(context.Background(), `DELETE FROM proxmox_clusters WHERE id = $1`, clusterIdentifier) }) return clusterIdentifier }