From 912b116bdb2c6d4d5fa7ca520a3252dd8b6bc000 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Thu, 2 Jan 2025 17:01:44 +0800 Subject: [PATCH] always use job's time Signed-off-by: Lyndon-Li --- .../backup_repository_controller.go | 204 +++++---- .../backup_repository_controller_test.go | 9 +- pkg/repository/maintenance.go | 89 ++-- pkg/repository/maintenance_test.go | 419 ++++++++++++------ pkg/repository/manager/manager.go | 33 +- pkg/repository/mocks/Manager.go | 20 +- 6 files changed, 483 insertions(+), 291 deletions(-) diff --git a/pkg/controller/backup_repository_controller.go b/pkg/controller/backup_repository_controller.go index 46aef5755..535fc41b6 100644 --- a/pkg/controller/backup_repository_controller.go +++ b/pkg/controller/backup_repository_controller.go @@ -208,7 +208,7 @@ func (r *BackupRepoReconciler) Reconcile(ctx context.Context, req ctrl.Request) } fallthrough case velerov1api.BackupRepositoryPhaseReady: - if err := r.processUnrecordedMaintenance(ctx, backupRepo, log); err != nil { + if err := r.recallMaintenance(ctx, backupRepo, log); err != nil { return ctrl.Result{}, errors.Wrap(err, "error handling incomplete repo maintenance jobs") } @@ -218,85 +218,6 @@ func (r *BackupRepoReconciler) Reconcile(ctx context.Context, req ctrl.Request) return ctrl.Result{}, nil } -func (r *BackupRepoReconciler) processUnrecordedMaintenance(ctx context.Context, req *velerov1api.BackupRepository, log logrus.FieldLogger) error { - history, err := repository.WaitIncompleteMaintenance(ctx, r.Client, req, defaultMaintenanceStatusQueueLength, log) - if err != nil { - return errors.Wrapf(err, "error waiting incomplete repo maintenance job for repo %s", req.Name) - } - - consolidated := consolidateHistory(history, req.Status.RecentMaintenance) - if consolidated == nil { - return nil - } - - log.Warn("Updating backup repository because of unrecorded histories") - - return r.patchBackupRepository(ctx, req, func(rr *velerov1api.BackupRepository) { - rr.Status.RecentMaintenance = consolidated - }) -} - -func consolidateHistory(coming, cur []velerov1api.BackupRepositoryMaintenanceStatus) []velerov1api.BackupRepositoryMaintenanceStatus { - if len(coming) == 0 { - return nil - } - - if isIdenticalHistories(coming, cur) { - return nil - } - - truncated := []velerov1api.BackupRepositoryMaintenanceStatus{} - i := len(cur) - 1 - j := len(coming) - 1 - for i >= 0 || j >= 0 { - if len(truncated) == defaultMaintenanceStatusQueueLength { - break - } - - if i >= 0 && j >= 0 { - if isEarlierHistory(cur[i], coming[j]) { - truncated = append(truncated, coming[j]) - j-- - } else { - truncated = append(truncated, cur[i]) - i-- - } - } else if i >= 0 { - truncated = append(truncated, cur[i]) - i-- - } else { - truncated = append(truncated, coming[j]) - j-- - } - } - - slices.Reverse(truncated) - - if isIdenticalHistories(truncated, cur) { - return nil - } - - return truncated -} - -func isIdenticalHistories(a, b []velerov1api.BackupRepositoryMaintenanceStatus) bool { - if len(a) != len(b) { - return false - } - - for i := 0; i < len(a); i++ { - if !a[i].CompleteTimestamp.Equal(b[i].CompleteTimestamp) { - return false - } - } - - return true -} - -func isEarlierHistory(a, b velerov1api.BackupRepositoryMaintenanceStatus) bool { - return a.CompleteTimestamp.Before(b.CompleteTimestamp) -} - func (r *BackupRepoReconciler) getIdentiferByBSL(ctx context.Context, req *velerov1api.BackupRepository) (string, error) { loc := &velerov1api.BackupStorageLocation{} @@ -384,10 +305,109 @@ func ensureRepo(repo *velerov1api.BackupRepository, repoManager repomanager.Mana return repoManager.PrepareRepo(repo) } -func (r *BackupRepoReconciler) runMaintenanceIfDue(ctx context.Context, req *velerov1api.BackupRepository, log logrus.FieldLogger) error { - startTime := r.clock.Now() +func (r *BackupRepoReconciler) recallMaintenance(ctx context.Context, req *velerov1api.BackupRepository, log logrus.FieldLogger) error { + history, err := repository.WaitIncompleteMaintenance(ctx, r.Client, req, defaultMaintenanceStatusQueueLength, log) + if err != nil { + return errors.Wrapf(err, "error waiting incomplete repo maintenance job for repo %s", req.Name) + } - if !dueForMaintenance(req, startTime) { + consolidated := consolidateHistory(history, req.Status.RecentMaintenance) + if consolidated == nil { + return nil + } + + lastMaintenanceTime := getLastMaintenanceTimeFromHistory(consolidated) + + log.Warn("Updating backup repository because of unrecorded histories") + + if lastMaintenanceTime.After(req.Status.LastMaintenanceTime.Time) { + log.Warnf("Updating backup repository last maintenance time (%v) from history (%v)", req.Status.LastMaintenanceTime.Time, lastMaintenanceTime.Time) + } + + return r.patchBackupRepository(ctx, req, func(rr *velerov1api.BackupRepository) { + if lastMaintenanceTime.After(rr.Status.LastMaintenanceTime.Time) { + rr.Status.LastMaintenanceTime = lastMaintenanceTime + } + + rr.Status.RecentMaintenance = consolidated + }) +} + +func consolidateHistory(coming, cur []velerov1api.BackupRepositoryMaintenanceStatus) []velerov1api.BackupRepositoryMaintenanceStatus { + if len(coming) == 0 { + return nil + } + + if isIdenticalHistory(coming, cur) { + return nil + } + + truncated := []velerov1api.BackupRepositoryMaintenanceStatus{} + i := len(cur) - 1 + j := len(coming) - 1 + for i >= 0 || j >= 0 { + if len(truncated) == defaultMaintenanceStatusQueueLength { + break + } + + if i >= 0 && j >= 0 { + if isEarlierMaintenanceStatus(cur[i], coming[j]) { + truncated = append(truncated, coming[j]) + j-- + } else { + truncated = append(truncated, cur[i]) + i-- + } + } else if i >= 0 { + truncated = append(truncated, cur[i]) + i-- + } else { + truncated = append(truncated, coming[j]) + j-- + } + } + + slices.Reverse(truncated) + + if isIdenticalHistory(truncated, cur) { + return nil + } + + return truncated +} + +func getLastMaintenanceTimeFromHistory(history []velerov1api.BackupRepositoryMaintenanceStatus) *metav1.Time { + time := history[0].CompleteTimestamp + + for i := range history { + if time.Before(history[i].CompleteTimestamp) { + time = history[i].CompleteTimestamp + } + } + + return time +} + +func isIdenticalHistory(a, b []velerov1api.BackupRepositoryMaintenanceStatus) bool { + if len(a) != len(b) { + return false + } + + for i := 0; i < len(a); i++ { + if !a[i].StartTimestamp.Equal(b[i].StartTimestamp) { + return false + } + } + + return true +} + +func isEarlierMaintenanceStatus(a, b velerov1api.BackupRepositoryMaintenanceStatus) bool { + return a.StartTimestamp.Before(b.StartTimestamp) +} + +func (r *BackupRepoReconciler) runMaintenanceIfDue(ctx context.Context, req *velerov1api.BackupRepository, log logrus.FieldLogger) error { + if !dueForMaintenance(req, r.clock.Now()) { log.Debug("not due for maintenance") return nil } @@ -398,17 +418,23 @@ func (r *BackupRepoReconciler) runMaintenanceIfDue(ctx context.Context, req *vel // should not cause the repo to move to `NotReady`. log.Debug("Pruning repo") - if err := r.repositoryManager.PruneRepo(req); err != nil { - log.WithError(err).Warn("error pruning repository") + // when PruneRepo fails, the maintenance result will be left temporarily + // If the maintenenance still completes later, recallMaintenance recalls the left onces and update LastMaintenanceTime and history + status, err := r.repositoryManager.PruneRepo(req) + if err != nil { + return errors.Wrapf(err, "error pruning repository") + } + + if status.Result == velerov1api.BackupRepositoryMaintenanceFailed { + log.WithError(err).Warn("Pruning repository failed") return r.patchBackupRepository(ctx, req, func(rr *velerov1api.BackupRepository) { - updateRepoMaintenanceHistory(rr, velerov1api.BackupRepositoryMaintenanceFailed, startTime, r.clock.Now(), err.Error()) + updateRepoMaintenanceHistory(rr, velerov1api.BackupRepositoryMaintenanceFailed, status.StartTimestamp.Time, status.CompleteTimestamp.Time, status.Message) }) } return r.patchBackupRepository(ctx, req, func(rr *velerov1api.BackupRepository) { - completionTime := r.clock.Now() - rr.Status.LastMaintenanceTime = &metav1.Time{Time: completionTime} - updateRepoMaintenanceHistory(rr, velerov1api.BackupRepositoryMaintenanceSucceeded, startTime, completionTime, "") + rr.Status.LastMaintenanceTime = &metav1.Time{Time: status.CompleteTimestamp.Time} + updateRepoMaintenanceHistory(rr, velerov1api.BackupRepositoryMaintenanceSucceeded, status.StartTimestamp.Time, status.CompleteTimestamp.Time, status.Message) }) } diff --git a/pkg/controller/backup_repository_controller_test.go b/pkg/controller/backup_repository_controller_test.go index 376b17ce8..9b5dd4c4a 100644 --- a/pkg/controller/backup_repository_controller_test.go +++ b/pkg/controller/backup_repository_controller_test.go @@ -37,11 +37,11 @@ import ( const testMaintenanceFrequency = 10 * time.Minute -func mockBackupRepoReconciler(t *testing.T, mockOn string, arg interface{}, ret interface{}) *BackupRepoReconciler { +func mockBackupRepoReconciler(t *testing.T, mockOn string, arg interface{}, ret ...interface{}) *BackupRepoReconciler { t.Helper() mgr := &repomokes.Manager{} if mockOn != "" { - mgr.On(mockOn, arg).Return(ret) + mgr.On(mockOn, arg).Return(ret...) } return NewBackupRepoReconciler( velerov1api.DefaultNamespace, @@ -106,7 +106,10 @@ func TestCheckNotReadyRepo(t *testing.T) { func TestRunMaintenanceIfDue(t *testing.T) { rr := mockBackupRepositoryCR() - reconciler := mockBackupRepoReconciler(t, "PruneRepo", rr, nil) + reconciler := mockBackupRepoReconciler(t, "PruneRepo", rr, velerov1api.BackupRepositoryMaintenanceStatus{ + StartTimestamp: &metav1.Time{}, + CompleteTimestamp: &metav1.Time{}, + }, nil) err := reconciler.Client.Create(context.TODO(), rr) assert.NoError(t, err) lastTm := rr.Status.LastMaintenanceTime diff --git a/pkg/repository/maintenance.go b/pkg/repository/maintenance.go index 96e08c147..5a1ec1f61 100644 --- a/pkg/repository/maintenance.go +++ b/pkg/repository/maintenance.go @@ -86,22 +86,26 @@ func DeleteOldMaintenanceJobs(cli client.Client, repo string, keep int) error { return nil } -func WaitForJobComplete(ctx context.Context, client client.Client, job *batchv1.Job) error { - return wait.PollUntilContextCancel(ctx, time.Second, true, func(ctx context.Context) (bool, error) { - err := client.Get(ctx, types.NamespacedName{Namespace: job.Namespace, Name: job.Name}, job) +// WaitForJobComplete wait for completion of the specified job and update the latest job object +func WaitForJobComplete(ctx context.Context, client client.Client, ns string, job string) (*batchv1.Job, error) { + updated := &batchv1.Job{} + err := wait.PollUntilContextCancel(ctx, time.Second, true, func(ctx context.Context) (bool, error) { + err := client.Get(ctx, types.NamespacedName{Namespace: ns, Name: job}, updated) if err != nil && !apierrors.IsNotFound(err) { return false, err } - if job.Status.Succeeded > 0 { + if updated.Status.Succeeded > 0 { return true, nil } - if job.Status.Failed > 0 { - return true, fmt.Errorf("maintenance job %s/%s failed", job.Namespace, job.Name) + if updated.Status.Failed > 0 { + return true, fmt.Errorf("maintenance job %s/%s failed", job, job) } return false, nil }) + + return updated, err } func GetMaintenanceResultFromJob(cli client.Client, job *batchv1.Job) (string, error) { @@ -269,43 +273,52 @@ func WaitIncompleteMaintenance(ctx context.Context, cli client.Client, repo *vel return nil, nil } - history := []velerov1api.BackupRepositoryMaintenanceStatus{} - - for _, job := range jobList.Items { - if job.Status.Succeeded == 0 && job.Status.Failed == 0 { - log.Infof("Waiting for maintenance job %s to complete", job.Name) - - if err := WaitForJobComplete(ctx, cli, &job); err != nil { - return nil, errors.Wrapf(err, "error waiting maintenance job[%s] complete", job.Name) - } - } - - result := velerov1api.BackupRepositoryMaintenanceSucceeded - if job.Status.Failed > 0 { - result = velerov1api.BackupRepositoryMaintenanceFailed - } - - message, err := GetMaintenanceResultFromJob(cli, &job) - if err != nil { - return nil, errors.Wrapf(err, "error getting maintenance job[%s] result", job.Name) - } - - history = append(history, velerov1api.BackupRepositoryMaintenanceStatus{ - Result: result, - StartTimestamp: &metav1.Time{Time: job.Status.StartTime.Time}, - CompleteTimestamp: &metav1.Time{Time: job.Status.CompletionTime.Time}, - Message: message, - }) - } - - sort.Slice(history, func(i, j int) bool { - return history[i].CompleteTimestamp.Time.After(history[j].CompleteTimestamp.Time) + sort.Slice(jobList.Items, func(i, j int) bool { + return jobList.Items[i].CreationTimestamp.Time.Before(jobList.Items[j].CreationTimestamp.Time) }) + history := []velerov1api.BackupRepositoryMaintenanceStatus{} + startPos := len(history) - limit if startPos < 0 { startPos = 0 } - return history[startPos:], nil + for i := startPos; i < len(jobList.Items); i++ { + job := &jobList.Items[i] + + if job.Status.Succeeded == 0 && job.Status.Failed == 0 { + log.Infof("Waiting for maintenance job %s to complete", job.Name) + + updated, err := WaitForJobComplete(ctx, cli, job.Namespace, job.Name) + if err != nil { + return nil, errors.Wrapf(err, "error waiting maintenance job[%s] complete", job.Name) + } + + job = updated + } + + message, err := GetMaintenanceResultFromJob(cli, job) + if err != nil { + return nil, errors.Wrapf(err, "error getting maintenance job[%s] result", job.Name) + } + + history = append(history, ComposeMaintenanceStatusFromJob(job, message)) + } + + return history, nil +} + +func ComposeMaintenanceStatusFromJob(job *batchv1.Job, message string) velerov1api.BackupRepositoryMaintenanceStatus { + result := velerov1api.BackupRepositoryMaintenanceSucceeded + if job.Status.Failed > 0 { + result = velerov1api.BackupRepositoryMaintenanceFailed + } + + return velerov1api.BackupRepositoryMaintenanceStatus{ + Result: result, + StartTimestamp: &metav1.Time{Time: job.CreationTimestamp.Time}, + CompleteTimestamp: &metav1.Time{Time: job.Status.CompletionTime.Time}, + Message: message, + } } diff --git a/pkg/repository/maintenance_test.go b/pkg/repository/maintenance_test.go index 7917ac65d..6cde95c95 100644 --- a/pkg/repository/maintenance_test.go +++ b/pkg/repository/maintenance_test.go @@ -34,6 +34,7 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client/fake" velerov1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v1" + "github.com/vmware-tanzu/velero/pkg/builder" velerotest "github.com/vmware-tanzu/velero/pkg/test" "github.com/vmware-tanzu/velero/pkg/util/kube" ) @@ -167,7 +168,7 @@ func TestWaitForJobComplete(t *testing.T) { // Create a fake Kubernetes client cli := fake.NewClientBuilder().WithObjects(job).Build() // Call the function - err := WaitForJobComplete(context.Background(), cli, job) + _, err := WaitForJobComplete(context.Background(), cli, job.Namespace, job.Name) // Check if the error matches the expectation if tc.expectError { @@ -459,143 +460,291 @@ func TestGetMaintenanceJobConfig(t *testing.T) { } } -// func TestWaitIncompleteMaintenance(t *testing.T) { -// ctx, cancel := context.WithCancel(context.Background()) +func TestWaitIncompleteMaintenance(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), time.Second*2) -// veleroNamespace := "velero" -// repo := &velerov1api.BackupRepository{ -// ObjectMeta: metav1.ObjectMeta{ -// Namespace: veleroNamespace, -// Name: "fake-repo", -// }, -// Spec: velerov1api.BackupRepositorySpec{ -// BackupStorageLocation: "default", -// RepositoryType: "kopia", -// VolumeNamespace: "test", -// }, -// } + veleroNamespace := "velero" + repo := &velerov1api.BackupRepository{ + ObjectMeta: metav1.ObjectMeta{ + Namespace: veleroNamespace, + Name: "fake-repo", + }, + Spec: velerov1api.BackupRepositorySpec{ + BackupStorageLocation: "default", + RepositoryType: "kopia", + VolumeNamespace: "test", + }, + } -// scheme := runtime.NewScheme() -// batchv1.AddToScheme(scheme) + now := time.Now().Round(time.Second) -// testCases := []struct { -// name string -// ctx context.Context -// kubeClientObj []runtime.Object -// runtimeScheme *runtime.Scheme -// expectedStatus []velerov1api.BackupRepositoryMaintenanceStatus -// expectedError string -// }{ -// { -// name: "list job error", -// expectedError: "error listing maintenance job for repo fake-repo", -// }, -// { -// name: "job not exist", -// runtimeScheme: scheme, -// }, -// { -// name: "no matching job", -// runtimeScheme: scheme, -// kubeClientObj: []runtime.Object{ -// jobOtherLabel, -// }, -// }, -// { -// name: "wait complete error", -// ctx: context.WithTimeout(context.TODO(), time.Second), -// runtimeScheme: scheme, -// kubeClientObj: []runtime.Object{ -// jobIncomplete, -// }, -// expectedError: nil, -// }, -// { -// name: "Find config specific for global", -// repoJobConfig: &v1.ConfigMap{ -// ObjectMeta: metav1.ObjectMeta{ -// Namespace: veleroNamespace, -// Name: repoMaintenanceJobConfig, -// }, -// Data: map[string]string{ -// GlobalKeyForRepoMaintenanceJobCM: "{\"podResources\":{\"cpuRequest\":\"50m\",\"cpuLimit\":\"100m\",\"memoryRequest\":\"50Mi\",\"memoryLimit\":\"100Mi\"},\"loadAffinity\":[{\"nodeSelector\":{\"matchExpressions\":[{\"key\":\"cloud.google.com/machine-family\",\"operator\":\"In\",\"values\":[\"n2\"]}]}}]}", -// }, -// }, -// expectedConfig: &JobConfigs{ -// PodResources: &kube.PodResources{ -// CPURequest: "50m", -// CPULimit: "100m", -// MemoryRequest: "50Mi", -// MemoryLimit: "100Mi", -// }, -// LoadAffinities: []*kube.LoadAffinity{ -// { -// NodeSelector: metav1.LabelSelector{ -// MatchExpressions: []metav1.LabelSelectorRequirement{ -// { -// Key: "cloud.google.com/machine-family", -// Operator: metav1.LabelSelectorOpIn, -// Values: []string{"n2"}, -// }, -// }, -// }, -// }, -// }, -// }, -// expectedError: nil, -// }, -// { -// name: "Specific config supersede global config", -// repoJobConfig: &v1.ConfigMap{ -// ObjectMeta: metav1.ObjectMeta{ -// Namespace: veleroNamespace, -// Name: repoMaintenanceJobConfig, -// }, -// Data: map[string]string{ -// GlobalKeyForRepoMaintenanceJobCM: "{\"podResources\":{\"cpuRequest\":\"50m\",\"cpuLimit\":\"100m\",\"memoryRequest\":\"50Mi\",\"memoryLimit\":\"100Mi\"},\"loadAffinity\":[{\"nodeSelector\":{\"matchExpressions\":[{\"key\":\"cloud.google.com/machine-family\",\"operator\":\"In\",\"values\":[\"n2\"]}]}}]}", -// "test-default-kopia": "{\"podResources\":{\"cpuRequest\":\"100m\",\"cpuLimit\":\"200m\",\"memoryRequest\":\"100Mi\",\"memoryLimit\":\"200Mi\"},\"loadAffinity\":[{\"nodeSelector\":{\"matchExpressions\":[{\"key\":\"cloud.google.com/machine-family\",\"operator\":\"In\",\"values\":[\"e2\"]}]}}]}", -// }, -// }, -// expectedConfig: &JobConfigs{ -// PodResources: &kube.PodResources{ -// CPURequest: "100m", -// CPULimit: "200m", -// MemoryRequest: "100Mi", -// MemoryLimit: "200Mi", -// }, -// LoadAffinities: []*kube.LoadAffinity{ -// { -// NodeSelector: metav1.LabelSelector{ -// MatchExpressions: []metav1.LabelSelectorRequirement{ -// { -// Key: "cloud.google.com/machine-family", -// Operator: metav1.LabelSelectorOpIn, -// Values: []string{"e2"}, -// }, -// }, -// }, -// }, -// }, -// }, -// expectedError: nil, -// }, -// } + jobOtherLabel := &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{ + Name: "job1", + Namespace: veleroNamespace, + Labels: map[string]string{RepositoryNameLabel: "other-repo"}, + CreationTimestamp: metav1.Time{Time: now}, + }, + } -// for _, test := range testCases { -// t.Run(test.name, func(t *testing.T) { -// fakeClientBuilder := fake.NewClientBuilder() -// if test.runtimeScheme != nil { -// fakeClientBuilder = fakeClientBuilder.WithScheme(test.runtimeScheme) -// } + jobIncomplete := &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{ + Name: "job1", + Namespace: veleroNamespace, + Labels: map[string]string{RepositoryNameLabel: "fake-repo"}, + CreationTimestamp: metav1.Time{Time: now}, + }, + } -// fakeClient := fakeClientBuilder.WithRuntimeObjects(test.kubeClientObj...).Build() + jobSucceeded1 := &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{ + Name: "job1", + Namespace: veleroNamespace, + Labels: map[string]string{RepositoryNameLabel: "fake-repo"}, + CreationTimestamp: metav1.Time{Time: now}, + }, + Status: batchv1.JobStatus{ + StartTime: &metav1.Time{Time: now}, + CompletionTime: &metav1.Time{Time: now.Add(time.Hour)}, + Succeeded: 1, + }, + } -// if tc.expectedError != nil { -// require.ErrorContains(t, err, tc.expectedError.Error()) -// } else { -// require.NoError(t, err) -// } -// require.Equal(t, tc.expectedConfig, jobConfig) -// }) -// } -// } + jobPodSucceeded1 := builder.ForPod(veleroNamespace, "job1").Labels(map[string]string{"job-name": "job1"}).ContainerStatuses(&v1.ContainerStatus{ + State: v1.ContainerState{ + Terminated: &v1.ContainerStateTerminated{}, + }, + }).Result() + + jobFailed1 := &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{ + Name: "job2", + Namespace: veleroNamespace, + Labels: map[string]string{RepositoryNameLabel: "fake-repo"}, + CreationTimestamp: metav1.Time{Time: now.Add(time.Hour)}, + }, + Status: batchv1.JobStatus{ + StartTime: &metav1.Time{Time: now.Add(time.Hour)}, + CompletionTime: &metav1.Time{Time: now.Add(time.Hour * 2)}, + Failed: 1, + }, + } + + jobPodFailed1 := builder.ForPod(veleroNamespace, "job2").Labels(map[string]string{"job-name": "job2"}).ContainerStatuses(&v1.ContainerStatus{ + State: v1.ContainerState{ + Terminated: &v1.ContainerStateTerminated{ + Message: "fake-message-2", + }, + }, + }).Result() + + jobSucceeded2 := &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{ + Name: "job3", + Namespace: veleroNamespace, + Labels: map[string]string{RepositoryNameLabel: "fake-repo"}, + CreationTimestamp: metav1.Time{Time: now.Add(time.Hour * 2)}, + }, + Status: batchv1.JobStatus{ + StartTime: &metav1.Time{Time: now.Add(time.Hour * 2)}, + CompletionTime: &metav1.Time{Time: now.Add(time.Hour * 3)}, + Succeeded: 1, + }, + } + + jobPodSucceeded2 := builder.ForPod(veleroNamespace, "job3").Labels(map[string]string{"job-name": "job3"}).ContainerStatuses(&v1.ContainerStatus{ + State: v1.ContainerState{ + Terminated: &v1.ContainerStateTerminated{}, + }, + }).Result() + + jobSucceeded3 := &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{ + Name: "job4", + Namespace: veleroNamespace, + Labels: map[string]string{RepositoryNameLabel: "fake-repo"}, + CreationTimestamp: metav1.Time{Time: now.Add(time.Hour * 3)}, + }, + Status: batchv1.JobStatus{ + StartTime: &metav1.Time{Time: now.Add(time.Hour * 3)}, + CompletionTime: &metav1.Time{Time: now.Add(time.Hour * 4)}, + Succeeded: 1, + }, + } + + jobPodSucceeded3 := builder.ForPod(veleroNamespace, "job4").Labels(map[string]string{"job-name": "job4"}).ContainerStatuses(&v1.ContainerStatus{ + State: v1.ContainerState{ + Terminated: &v1.ContainerStateTerminated{}, + }, + }).Result() + + schemeFail := runtime.NewScheme() + + scheme := runtime.NewScheme() + batchv1.AddToScheme(scheme) + v1.AddToScheme(scheme) + + testCases := []struct { + name string + ctx context.Context + kubeClientObj []runtime.Object + runtimeScheme *runtime.Scheme + expectedStatus []velerov1api.BackupRepositoryMaintenanceStatus + expectedError string + }{ + { + name: "list job error", + runtimeScheme: schemeFail, + expectedError: "error listing maintenance job for repo fake-repo: no kind is registered for the type v1.JobList in scheme \"pkg/runtime/scheme.go:100\"", + }, + { + name: "job not exist", + runtimeScheme: scheme, + }, + { + name: "no matching job", + runtimeScheme: scheme, + kubeClientObj: []runtime.Object{ + jobOtherLabel, + }, + }, + { + name: "wait complete error", + ctx: ctx, + runtimeScheme: scheme, + kubeClientObj: []runtime.Object{ + jobIncomplete, + }, + expectedError: "error waiting maintenance job[job1] complete: context deadline exceeded", + }, + { + name: "get result error", + ctx: context.TODO(), + runtimeScheme: scheme, + kubeClientObj: []runtime.Object{ + jobSucceeded1, + }, + expectedError: "error getting maintenance job[job1] result: no pod found for job job1", + }, + { + name: "less than limit", + ctx: context.TODO(), + runtimeScheme: scheme, + kubeClientObj: []runtime.Object{ + jobFailed1, + jobSucceeded1, + jobPodSucceeded1, + jobPodFailed1, + }, + expectedStatus: []velerov1api.BackupRepositoryMaintenanceStatus{ + { + Result: velerov1api.BackupRepositoryMaintenanceSucceeded, + StartTimestamp: &metav1.Time{Time: now}, + CompleteTimestamp: &metav1.Time{Time: now.Add(time.Hour)}, + }, + { + Result: velerov1api.BackupRepositoryMaintenanceFailed, + StartTimestamp: &metav1.Time{Time: now.Add(time.Hour)}, + CompleteTimestamp: &metav1.Time{Time: now.Add(time.Hour * 2)}, + Message: "fake-message-2", + }, + }, + }, + { + name: "equal to limit", + ctx: context.TODO(), + runtimeScheme: scheme, + kubeClientObj: []runtime.Object{ + jobSucceeded2, + jobFailed1, + jobSucceeded1, + jobPodSucceeded1, + jobPodFailed1, + jobPodSucceeded2, + }, + expectedStatus: []velerov1api.BackupRepositoryMaintenanceStatus{ + { + Result: velerov1api.BackupRepositoryMaintenanceSucceeded, + StartTimestamp: &metav1.Time{Time: now}, + CompleteTimestamp: &metav1.Time{Time: now.Add(time.Hour)}, + }, + { + Result: velerov1api.BackupRepositoryMaintenanceFailed, + StartTimestamp: &metav1.Time{Time: now.Add(time.Hour)}, + CompleteTimestamp: &metav1.Time{Time: now.Add(time.Hour * 2)}, + Message: "fake-message-2", + }, + { + Result: velerov1api.BackupRepositoryMaintenanceSucceeded, + StartTimestamp: &metav1.Time{Time: now.Add(time.Hour * 2)}, + CompleteTimestamp: &metav1.Time{Time: now.Add(time.Hour * 3)}, + }, + }, + }, + { + name: "more than limit", + ctx: context.TODO(), + runtimeScheme: scheme, + kubeClientObj: []runtime.Object{ + jobSucceeded3, + jobSucceeded2, + jobFailed1, + jobSucceeded1, + jobPodSucceeded1, + jobPodFailed1, + jobPodSucceeded2, + jobPodSucceeded3, + }, + expectedStatus: []velerov1api.BackupRepositoryMaintenanceStatus{ + { + Result: velerov1api.BackupRepositoryMaintenanceSucceeded, + StartTimestamp: &metav1.Time{Time: now}, + CompleteTimestamp: &metav1.Time{Time: now.Add(time.Hour)}, + }, + { + Result: velerov1api.BackupRepositoryMaintenanceFailed, + StartTimestamp: &metav1.Time{Time: now.Add(time.Hour)}, + CompleteTimestamp: &metav1.Time{Time: now.Add(time.Hour * 2)}, + Message: "fake-message-2", + }, + { + Result: velerov1api.BackupRepositoryMaintenanceSucceeded, + StartTimestamp: &metav1.Time{Time: now.Add(time.Hour * 2)}, + CompleteTimestamp: &metav1.Time{Time: now.Add(time.Hour * 3)}, + }, + { + Result: velerov1api.BackupRepositoryMaintenanceSucceeded, + StartTimestamp: &metav1.Time{Time: now.Add(time.Hour * 3)}, + CompleteTimestamp: &metav1.Time{Time: now.Add(time.Hour * 4)}, + }, + }, + }, + } + + for _, test := range testCases { + t.Run(test.name, func(t *testing.T) { + fakeClientBuilder := fake.NewClientBuilder() + fakeClientBuilder = fakeClientBuilder.WithScheme(test.runtimeScheme) + + fakeClient := fakeClientBuilder.WithRuntimeObjects(test.kubeClientObj...).Build() + + history, err := WaitIncompleteMaintenance(test.ctx, fakeClient, repo, 3, velerotest.NewLogger()) + + if test.expectedError != "" { + assert.EqualError(t, err, test.expectedError) + } else { + require.NoError(t, err) + } + + assert.Len(t, history, len(test.expectedStatus)) + for i := 0; i < len(test.expectedStatus); i++ { + assert.Equal(t, test.expectedStatus[i].Result, history[i].Result) + assert.Equal(t, test.expectedStatus[i].Message, history[i].Message) + assert.Equal(t, test.expectedStatus[i].StartTimestamp.Time, history[i].StartTimestamp.Time) + assert.Equal(t, test.expectedStatus[i].CompleteTimestamp.Time, history[i].CompleteTimestamp.Time) + } + }) + } + + cancel() +} diff --git a/pkg/repository/manager/manager.go b/pkg/repository/manager/manager.go index f590f2b14..3853f2d3c 100644 --- a/pkg/repository/manager/manager.go +++ b/pkg/repository/manager/manager.go @@ -54,7 +54,7 @@ type Manager interface { PrepareRepo(repo *velerov1api.BackupRepository) error // PruneRepo deletes unused data from a repo. - PruneRepo(repo *velerov1api.BackupRepository) error + PruneRepo(repo *velerov1api.BackupRepository) (velerov1api.BackupRepositoryMaintenanceStatus, error) // UnlockRepo removes stale locks from a repo. UnlockRepo(repo *velerov1api.BackupRepository) error @@ -172,13 +172,13 @@ func (m *manager) PrepareRepo(repo *velerov1api.BackupRepository) error { return prd.PrepareRepo(context.Background(), param) } -func (m *manager) PruneRepo(repo *velerov1api.BackupRepository) error { +func (m *manager) PruneRepo(repo *velerov1api.BackupRepository) (velerov1api.BackupRepositoryMaintenanceStatus, error) { m.repoLocker.LockExclusive(repo.Name) defer m.repoLocker.UnlockExclusive(repo.Name) param, err := m.assembleRepoParam(repo) if err != nil { - return errors.WithStack(err) + return velerov1api.BackupRepositoryMaintenanceStatus{}, errors.WithStack(err) } log := m.log.WithFields(logrus.Fields{ @@ -190,12 +190,12 @@ func (m *manager) PruneRepo(repo *velerov1api.BackupRepository) error { job, err := repository.GetLatestMaintenanceJob(m.client, m.namespace) if err != nil { - return errors.WithStack(err) + return velerov1api.BackupRepositoryMaintenanceStatus{}, errors.WithStack(err) } if job != nil && job.Status.Succeeded == 0 && job.Status.Failed == 0 { log.Debugf("There already has a unfinished maintenance job %s/%s for repository %s, please wait for it to complete", job.Namespace, job.Name, param.BackupRepo.Name) - return nil + return velerov1api.BackupRepositoryMaintenanceStatus{}, nil } jobConfig, err := repository.GetMaintenanceJobConfig( @@ -220,13 +220,13 @@ func (m *manager) PruneRepo(repo *velerov1api.BackupRepository) error { param, ) if err != nil { - return errors.Wrap(err, "error to build maintenance job") + return velerov1api.BackupRepositoryMaintenanceStatus{}, errors.Wrap(err, "error to build maintenance job") } log = log.WithField("job", fmt.Sprintf("%s/%s", maintenanceJob.Namespace, maintenanceJob.Name)) if err := m.client.Create(context.TODO(), maintenanceJob); err != nil { - return errors.Wrap(err, "error to create maintenance job") + return velerov1api.BackupRepositoryMaintenanceStatus{}, errors.Wrap(err, "error to create maintenance job") } log.Debug("Creating maintenance job") @@ -240,27 +240,18 @@ func (m *manager) PruneRepo(repo *velerov1api.BackupRepository) error { } }() - var jobErr error - if err := repository.WaitForJobComplete(context.TODO(), m.client, maintenanceJob); err != nil { - log.WithError(err).Error("Error to wait for maintenance job complete") - jobErr = err // we won't return here for job may failed by maintenance failure, we want return the actual error + maintenanceJob, err = repository.WaitForJobComplete(context.TODO(), m.client, maintenanceJob.Namespace, maintenanceJob.Name) + if err != nil { + return velerov1api.BackupRepositoryMaintenanceStatus{}, errors.Wrap(err, "error to wait for maintenance job complete") } result, err := repository.GetMaintenanceResultFromJob(m.client, maintenanceJob) if err != nil { - return errors.Wrap(err, "error to get maintenance job result") - } - - if result != "" { - return errors.New(fmt.Sprintf("Maintenance job %s failed: %s", maintenanceJob.Name, result)) - } - - if jobErr != nil { - return errors.Wrap(jobErr, "error to wait for maintenance job complete") + return velerov1api.BackupRepositoryMaintenanceStatus{}, errors.Wrap(err, "error to get maintenance job result") } log.Info("Maintenance repo complete") - return nil + return repository.ComposeMaintenanceStatusFromJob(maintenanceJob, result), nil } func (m *manager) UnlockRepo(repo *velerov1api.BackupRepository) error { diff --git a/pkg/repository/mocks/Manager.go b/pkg/repository/mocks/Manager.go index 226411775..c987693b4 100644 --- a/pkg/repository/mocks/Manager.go +++ b/pkg/repository/mocks/Manager.go @@ -138,21 +138,31 @@ func (_m *Manager) PrepareRepo(repo *v1.BackupRepository) error { } // PruneRepo provides a mock function with given fields: repo -func (_m *Manager) PruneRepo(repo *v1.BackupRepository) error { +func (_m *Manager) PruneRepo(repo *v1.BackupRepository) (v1.BackupRepositoryMaintenanceStatus, error) { ret := _m.Called(repo) if len(ret) == 0 { panic("no return value specified for PruneRepo") } - var r0 error - if rf, ok := ret.Get(0).(func(*v1.BackupRepository) error); ok { + var r0 v1.BackupRepositoryMaintenanceStatus + var r1 error + if rf, ok := ret.Get(0).(func(*v1.BackupRepository) (v1.BackupRepositoryMaintenanceStatus, error)); ok { + return rf(repo) + } + if rf, ok := ret.Get(0).(func(*v1.BackupRepository) v1.BackupRepositoryMaintenanceStatus); ok { r0 = rf(repo) } else { - r0 = ret.Error(0) + r0 = ret.Get(0).(v1.BackupRepositoryMaintenanceStatus) } - return r0 + if rf, ok := ret.Get(1).(func(*v1.BackupRepository) error); ok { + r1 = rf(repo) + } else { + r1 = ret.Error(1) + } + + return r0, r1 } // UnlockRepo provides a mock function with given fields: repo