always use job's time

Signed-off-by: Lyndon-Li <lyonghui@vmware.com>
This commit is contained in:
Lyndon-Li
2025-01-03 16:50:35 +08:00
parent 6ff0aa32e3
commit 912b116bdb
6 changed files with 483 additions and 291 deletions
+115 -89
View File
@@ -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)
})
}
@@ -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
+51 -38
View File
@@ -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,
}
}
+284 -135
View File
@@ -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()
}
+12 -21
View File
@@ -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 {
+15 -5
View File
@@ -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