Delete the unneeded pvRestorer action in

handleSkippedPVHasRetainPolicy

According to comment, calling executePVAction aims to reset PV's
claimRef, but the reset logic was moved into resetVolumeBindingInfo
since release-1.4.

Signed-off-by: Xun Jiang <blackpigletbruce@gmail.com>
This commit is contained in:
Xun Jiang
2024-03-29 14:12:12 +08:00
parent b06d7a467f
commit 5462035469
8 changed files with 251 additions and 97 deletions
+103 -8
View File
@@ -18,6 +18,7 @@ package backup
import (
"archive/tar"
"bytes"
"compress/gzip"
"context"
"encoding/json"
@@ -45,6 +46,8 @@ import (
"github.com/vmware-tanzu/velero/pkg/discovery"
"github.com/vmware-tanzu/velero/pkg/itemoperation"
"github.com/vmware-tanzu/velero/pkg/kuberesource"
"github.com/vmware-tanzu/velero/pkg/persistence"
"github.com/vmware-tanzu/velero/pkg/plugin/clientmgmt"
"github.com/vmware-tanzu/velero/pkg/plugin/framework"
"github.com/vmware-tanzu/velero/pkg/plugin/velero"
biav2 "github.com/vmware-tanzu/velero/pkg/plugin/velero/backupitemaction/v2"
@@ -90,7 +93,6 @@ type Backupper interface {
outBackupFile io.Writer,
backupItemActionResolver framework.BackupItemActionResolverV2,
asyncBIAOperations []*itemoperation.BackupOperation,
volumeInfos []*volume.VolumeInfo,
) error
}
@@ -105,6 +107,8 @@ type kubernetesBackupper struct {
defaultVolumesToFsBackup bool
clientPageSize int
uploaderType string
pluginManager func(logrus.FieldLogger) clientmgmt.Manager
backupStoreGetter persistence.ObjectBackupStoreGetter
}
func (i *itemKey) String() string {
@@ -132,6 +136,8 @@ func NewKubernetesBackupper(
defaultVolumesToFsBackup bool,
clientPageSize int,
uploaderType string,
pluginManager func(logrus.FieldLogger) clientmgmt.Manager,
backupStoreGetter persistence.ObjectBackupStoreGetter,
) (Backupper, error) {
return &kubernetesBackupper{
kbClient: kbClient,
@@ -143,6 +149,8 @@ func NewKubernetesBackupper(
defaultVolumesToFsBackup: defaultVolumesToFsBackup,
clientPageSize: clientPageSize,
uploaderType: uploaderType,
pluginManager: pluginManager,
backupStoreGetter: backupStoreGetter,
}, nil
}
@@ -584,7 +592,6 @@ func (kb *kubernetesBackupper) FinalizeBackup(
outBackupFile io.Writer,
backupItemActionResolver framework.BackupItemActionResolverV2,
asyncBIAOperations []*itemoperation.BackupOperation,
volumeInfos []*volume.VolumeInfo,
) error {
gzw := gzip.NewWriter(outBackupFile)
defer gzw.Close()
@@ -648,6 +655,8 @@ func (kb *kubernetesBackupper) FinalizeBackup(
updateFiles := make(map[string]FileForArchive)
backedUpGroupResources := map[schema.GroupResource]bool{}
unstructuredDataUploads := make([]unstructured.Unstructured, 0)
for i, item := range items {
log.WithFields(map[string]interface{}{
"progress": "",
@@ -674,7 +683,9 @@ func (kb *kubernetesBackupper) FinalizeBackup(
return
}
updateVolumeInfos(volumeInfos, unstructured, item.groupResource, log)
if item.groupResource == kuberesource.DataUploads {
unstructuredDataUploads = append(unstructuredDataUploads, unstructured)
}
backedUp, itemFiles := kb.finalizeItem(log, item.groupResource, itemBackupper, &unstructured, item.preferredGVR)
if backedUp {
@@ -697,6 +708,22 @@ func (kb *kubernetesBackupper) FinalizeBackup(
}).Infof("Updated %d items out of an estimated total of %d (estimate will change throughout the backup finalizer)", len(backupRequest.BackedUpItems), totalItems)
}
backupStore, volumeInfos, err := kb.getVolumeInfos(*backupRequest.Backup, log)
if err != nil {
log.WithError(err).Errorf("fail to get the backup VolumeInfos for backup %s", backupRequest.Name)
return err
}
if err := updateVolumeInfos(volumeInfos, unstructuredDataUploads, asyncBIAOperations, log); err != nil {
log.WithError(err).Errorf("fail to update VolumeInfos for backup %s", backupRequest.Name)
return err
}
if err := putVolumeInfos(backupRequest.Name, volumeInfos, backupStore); err != nil {
log.WithError(err).Errorf("fail to put the VolumeInfos for backup %s", backupRequest.Name)
return err
}
// write new tar archive replacing files in original with content updateFiles for matches
if err := buildFinalTarball(tr, tw, updateFiles); err != nil {
log.Errorf("Error building final tarball: %s", err.Error())
@@ -765,19 +792,48 @@ type tarWriter interface {
WriteHeader(*tar.Header) error
}
func (kb *kubernetesBackupper) getVolumeInfos(
backup velerov1api.Backup,
log logrus.FieldLogger,
) (persistence.BackupStore, []*volume.VolumeInfo, error) {
location := &velerov1api.BackupStorageLocation{}
if err := kb.kbClient.Get(context.Background(), kbclient.ObjectKey{
Namespace: backup.Namespace,
Name: backup.Spec.StorageLocation,
}, location); err != nil {
return nil, nil, errors.WithStack(err)
}
pluginManager := kb.pluginManager(log)
defer pluginManager.CleanupClients()
backupStore, storeErr := kb.backupStoreGetter.Get(location, pluginManager, log)
if storeErr != nil {
return nil, nil, storeErr
}
volumeInfos, err := backupStore.GetBackupVolumeInfos(backup.Name)
if err != nil {
return nil, nil, err
}
return backupStore, volumeInfos, nil
}
// updateVolumeInfos update the VolumeInfos according to the AsyncOperations
func updateVolumeInfos(
volumeInfos []*volume.VolumeInfo,
unstructured unstructured.Unstructured,
groupResource schema.GroupResource,
unstructuredItems []unstructured.Unstructured,
operations []*itemoperation.BackupOperation,
log logrus.FieldLogger,
) {
switch groupResource.String() {
case kuberesource.DataUploads.String():
) error {
for _, unstructured := range unstructuredItems {
var dataUpload velerov2alpha1.DataUpload
err := runtime.DefaultUnstructuredConverter.FromUnstructured(unstructured.UnstructuredContent(), &dataUpload)
if err != nil {
log.WithError(err).Errorf("fail to convert DataUpload: %s/%s",
unstructured.GetNamespace(), unstructured.GetName())
return err
}
for index := range volumeInfos {
@@ -792,4 +848,43 @@ func updateVolumeInfos(
}
}
}
// Update CSI snapshot VolumeInfo's CompletionTimestamp by the operation update time.
for volumeIndex := range volumeInfos {
if volumeInfos[volumeIndex].BackupMethod == volume.CSISnapshot &&
volumeInfos[volumeIndex].CSISnapshotInfo != nil {
for opIndex := range operations {
if volumeInfos[volumeIndex].CSISnapshotInfo.OperationID == operations[opIndex].Spec.OperationID {
// The VolumeSnapshot and VolumeSnapshotContent don't have a completion timestamp,
// so use the operation.Status.Updated as the alternative. It is not the exact time
// when the snapshot turns ready, but the operation controller periodically watch the
// VSC and VS status. When the controller finds they reach to the ReadyToUse state,
// The operation.Status.Updated is set as the found time.
volumeInfos[volumeIndex].CompletionTimestamp = operations[opIndex].Status.Updated
}
}
}
}
return nil
}
func putVolumeInfos(
backupName string,
volumeInfos []*volume.VolumeInfo,
backupStore persistence.BackupStore,
) error {
backupVolumeInfoBuf := new(bytes.Buffer)
gzw := gzip.NewWriter(backupVolumeInfoBuf)
defer gzw.Close()
if err := json.NewEncoder(gzw).Encode(volumeInfos); err != nil {
return errors.Wrap(err, "error encoding restore results to JSON")
}
if err := gzw.Close(); err != nil {
return errors.Wrap(err, "error closing gzip writer")
}
return backupStore.PutBackupVolumeInfos(backupName, backupVolumeInfoBuf)
}
+134 -45
View File
@@ -52,6 +52,10 @@ import (
"github.com/vmware-tanzu/velero/pkg/features"
"github.com/vmware-tanzu/velero/pkg/itemoperation"
"github.com/vmware-tanzu/velero/pkg/kuberesource"
"github.com/vmware-tanzu/velero/pkg/persistence"
persistencemocks "github.com/vmware-tanzu/velero/pkg/persistence/mocks"
"github.com/vmware-tanzu/velero/pkg/plugin/clientmgmt"
pluginmocks "github.com/vmware-tanzu/velero/pkg/plugin/mocks"
"github.com/vmware-tanzu/velero/pkg/plugin/velero"
biav2 "github.com/vmware-tanzu/velero/pkg/plugin/velero/backupitemaction/v2"
vsv1 "github.com/vmware-tanzu/velero/pkg/plugin/velero/volumesnapshotter/v1"
@@ -4436,58 +4440,143 @@ func TestBackupNamespaces(t *testing.T) {
}
}
func TestUpdateVolumeInfos(t *testing.T) {
logger := logrus.StandardLogger()
// The unstructured conversion will loose the time precision to second
// level. To make test pass. Set the now precision at second at the
// beginning.
now := metav1.Now().Rfc3339Copy()
volumeInfos := []*volume.VolumeInfo{
{
PVCName: "pvc1",
PVCNamespace: "ns1",
SnapshotDataMovementInfo: &volume.SnapshotDataMovementInfo{},
},
}
dataUpload := velerov2alpha1.DataUpload{
ObjectMeta: metav1.ObjectMeta{
Name: "du1",
Namespace: "velero",
},
Spec: velerov2alpha1.DataUploadSpec{
SourcePVC: "pvc1",
SourceNamespace: "ns1",
CSISnapshot: &velerov2alpha1.CSISnapshotSpec{
VolumeSnapshot: "vs1",
},
},
Status: velerov2alpha1.DataUploadStatus{
CompletionTimestamp: &now,
SnapshotID: "snapshot1",
Progress: shared.DataMoveOperationProgress{
TotalBytes: 10000,
},
},
}
duMap, err := runtime.DefaultUnstructuredConverter.ToUnstructured(&dataUpload)
func TestGetVolumeInfos(t *testing.T) {
h := newHarness(t)
pluginManager := new(pluginmocks.Manager)
backupStore := new(persistencemocks.BackupStore)
h.backupper.pluginManager = func(logrus.FieldLogger) clientmgmt.Manager { return pluginManager }
h.backupper.backupStoreGetter = NewFakeSingleObjectBackupStoreGetter(backupStore)
backupStore.On("GetBackupVolumeInfos", "backup-01").Return([]*volume.VolumeInfo{}, nil)
pluginManager.On("CleanupClients").Return()
backup := builder.ForBackup("velero", "backup-01").StorageLocation("default").Result()
bsl := builder.ForBackupStorageLocation("velero", "default").Result()
require.NoError(t, h.backupper.kbClient.Create(context.Background(), bsl))
_, _, err := h.backupper.getVolumeInfos(*backup, h.log)
require.NoError(t, err)
}
expectedVolumeInfos := []*volume.VolumeInfo{
func TestUpdateVolumeInfos(t *testing.T) {
timeExample := time.Date(2014, 6, 5, 11, 56, 45, 0, time.Local)
now := metav1.NewTime(timeExample)
logger := logrus.StandardLogger()
tests := []struct {
name string
operations []*itemoperation.BackupOperation
dataUpload *velerov2alpha1.DataUpload
volumeInfos []*volume.VolumeInfo
expectedVolumeInfos []*volume.VolumeInfo
}{
{
PVCName: "pvc1",
PVCNamespace: "ns1",
CompletionTimestamp: &now,
SnapshotDataMovementInfo: &volume.SnapshotDataMovementInfo{
SnapshotHandle: "snapshot1",
Size: 10000,
RetainedSnapshot: "vs1",
name: "CSISnapshot VolumeInfo update",
operations: []*itemoperation.BackupOperation{
{
Spec: itemoperation.BackupOperationSpec{
OperationID: "test-operation",
},
Status: itemoperation.OperationStatus{
Updated: &now,
},
},
},
volumeInfos: []*volume.VolumeInfo{
{
BackupMethod: volume.CSISnapshot,
CompletionTimestamp: &metav1.Time{},
CSISnapshotInfo: &volume.CSISnapshotInfo{
OperationID: "test-operation",
},
},
},
expectedVolumeInfos: []*volume.VolumeInfo{
{
BackupMethod: volume.CSISnapshot,
CompletionTimestamp: &now,
CSISnapshotInfo: &volume.CSISnapshotInfo{
OperationID: "test-operation",
},
},
},
},
{
name: "DataUpload VolumeInfo update",
operations: []*itemoperation.BackupOperation{},
dataUpload: builder.ForDataUpload("velero", "du-1").
CompletionTimestamp(&now).
CSISnapshot(&velerov2alpha1.CSISnapshotSpec{VolumeSnapshot: "vs-1"}).
SnapshotID("snapshot-id").
Progress(shared.DataMoveOperationProgress{TotalBytes: 1000}).
SourceNamespace("ns-1").
SourcePVC("pvc-1").
Result(),
volumeInfos: []*volume.VolumeInfo{
{
PVCName: "pvc-1",
PVCNamespace: "ns-1",
CompletionTimestamp: &metav1.Time{},
SnapshotDataMovementInfo: &volume.SnapshotDataMovementInfo{
DataMover: "velero",
},
},
},
expectedVolumeInfos: []*volume.VolumeInfo{
{
PVCName: "pvc-1",
PVCNamespace: "ns-1",
CompletionTimestamp: &now,
SnapshotDataMovementInfo: &volume.SnapshotDataMovementInfo{
DataMover: "velero",
RetainedSnapshot: "vs-1",
SnapshotHandle: "snapshot-id",
Size: 1000,
},
},
},
},
}
updateVolumeInfos(volumeInfos, unstructured.Unstructured{Object: duMap}, kuberesource.DataUploads, logger)
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
unstructures := []unstructured.Unstructured{}
if tc.dataUpload != nil {
duMap, error := runtime.DefaultUnstructuredConverter.ToUnstructured(tc.dataUpload)
require.NoError(t, error)
unstructures = append(unstructures,
unstructured.Unstructured{
Object: duMap,
},
)
}
if len(expectedVolumeInfos) > 0 {
require.Equal(t, expectedVolumeInfos[0].SnapshotDataMovementInfo, volumeInfos[0].SnapshotDataMovementInfo)
require.NoError(t, updateVolumeInfos(tc.volumeInfos, unstructures, tc.operations, logger))
require.Equal(t, tc.expectedVolumeInfos[0].CompletionTimestamp, tc.volumeInfos[0].CompletionTimestamp)
require.Equal(t, tc.expectedVolumeInfos[0].SnapshotDataMovementInfo, tc.volumeInfos[0].SnapshotDataMovementInfo)
})
}
}
func TestPutVolumeInfos(t *testing.T) {
backupName := "backup-01"
backupStore := new(persistencemocks.BackupStore)
backupStore.On("PutBackupVolumeInfos", mock.Anything, mock.Anything).Return(nil)
require.NoError(t, putVolumeInfos(backupName, []*volume.VolumeInfo{}, backupStore))
}
type fakeSingleObjectBackupStoreGetter struct {
store persistence.BackupStore
}
func (f *fakeSingleObjectBackupStoreGetter) Get(*velerov1.BackupStorageLocation, persistence.ObjectStoreGetter, logrus.FieldLogger) (persistence.BackupStore, error) {
return f.store, nil
}
// NewFakeSingleObjectBackupStoreGetter returns an ObjectBackupStoreGetter
// that will return only the given BackupStore.
func NewFakeSingleObjectBackupStoreGetter(store persistence.BackupStore) persistence.ObjectBackupStoreGetter {
return &fakeSingleObjectBackupStoreGetter{store: store}
}
+6
View File
@@ -19,6 +19,7 @@ package builder
import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"github.com/vmware-tanzu/velero/pkg/apis/velero/shared"
velerov2alpha1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v2alpha1"
)
@@ -131,3 +132,8 @@ func (d *DataUploadBuilder) Labels(labels map[string]string) *DataUploadBuilder
d.object.Labels = labels
return d
}
func (d *DataUploadBuilder) Progress(progress shared.DataMoveOperationProgress) *DataUploadBuilder {
d.object.Status.Progress = progress
return d
}
+4
View File
@@ -780,6 +780,8 @@ func (s *server) runControllers(defaultVolumeSnapshotLocations map[string]string
s.config.defaultVolumesToFsBackup,
s.config.clientPageSize,
s.config.uploaderType,
newPluginManager,
backupStoreGetter,
)
cmd.CheckError(err)
if err := controller.NewBackupReconciler(
@@ -860,6 +862,8 @@ func (s *server) runControllers(defaultVolumeSnapshotLocations map[string]string
s.config.defaultVolumesToFsBackup,
s.config.clientPageSize,
s.config.uploaderType,
newPluginManager,
backupStoreGetter,
)
cmd.CheckError(err)
r := controller.NewBackupFinalizerReconciler(
-2
View File
@@ -44,7 +44,6 @@ import (
kbclient "sigs.k8s.io/controller-runtime/pkg/client"
fakeClient "sigs.k8s.io/controller-runtime/pkg/client/fake"
"github.com/vmware-tanzu/velero/internal/volume"
velerov1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v1"
pkgbackup "github.com/vmware-tanzu/velero/pkg/backup"
"github.com/vmware-tanzu/velero/pkg/builder"
@@ -86,7 +85,6 @@ func (b *fakeBackupper) FinalizeBackup(
outBackupFile io.Writer,
backupItemActionResolver framework.BackupItemActionResolverV2,
asyncBIAOperations []*itemoperation.BackupOperation,
volumeInfos []*volume.VolumeInfo,
) error {
args := b.Called(logger, backup, inBackupFile, outBackupFile, backupItemActionResolver, asyncBIAOperations)
return args.Error(0)
@@ -18,9 +18,7 @@ package controller
import (
"bytes"
"compress/gzip"
"context"
"encoding/json"
"os"
"github.com/pkg/errors"
@@ -31,7 +29,6 @@ import (
ctrl "sigs.k8s.io/controller-runtime"
kbclient "sigs.k8s.io/controller-runtime/pkg/client"
"github.com/vmware-tanzu/velero/internal/volume"
velerov1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v1"
pkgbackup "github.com/vmware-tanzu/velero/pkg/backup"
"github.com/vmware-tanzu/velero/pkg/itemoperation"
@@ -157,14 +154,7 @@ func (r *backupFinalizerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
SkippedPVTracker: pkgbackup.NewSkipPVTracker(),
}
var outBackupFile *os.File
var volumeInfos []*volume.VolumeInfo
if len(operations) > 0 {
volumeInfos, err = backupStore.GetBackupVolumeInfos(backup.Name)
if err != nil {
log.WithError(err).Error("error getting backup VolumeInfos")
return ctrl.Result{}, errors.WithStack(err)
}
log.Info("Setting up finalized backup temp file")
inBackupFile, err := downloadToTempFile(backup.Name, backupStore, log)
if err != nil {
@@ -194,7 +184,6 @@ func (r *backupFinalizerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
outBackupFile,
backupItemActionsResolver,
operations,
volumeInfos,
)
if err != nil {
log.WithError(err).Error("error finalizing Backup")
@@ -232,24 +221,6 @@ func (r *backupFinalizerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
if err != nil {
return ctrl.Result{}, errors.Wrap(err, "error uploading backup final contents")
}
// Update the backup's VolumeInfos
backupVolumeInfoBuf := new(bytes.Buffer)
gzw := gzip.NewWriter(backupVolumeInfoBuf)
defer gzw.Close()
if err := json.NewEncoder(gzw).Encode(volumeInfos); err != nil {
return ctrl.Result{}, errors.Wrap(err, "error encoding restore results to JSON")
}
if err := gzw.Close(); err != nil {
return ctrl.Result{}, errors.Wrap(err, "error closing gzip writer")
}
err = backupStore.PutBackupVolumeInfos(backup.Name, backupVolumeInfoBuf)
if err != nil {
return ctrl.Result{}, errors.Wrap(err, "fail to upload backup VolumeInfos")
}
}
return ctrl.Result{}, nil
}
+1 -1
View File
@@ -118,7 +118,7 @@ func TestDeleteOldMaintenanceJobs(t *testing.T) {
assert.NoError(t, err)
// We expect the number of jobs to be equal to 'keep'
assert.Equal(t, keep, len(jobList.Items))
assert.Len(t, jobList.Items, keep)
// We expect that the oldest jobs were deleted
// Job3 should not be present in the remaining list
+3 -12
View File
@@ -1268,7 +1268,7 @@ func (ctx *restoreContext) restoreItem(obj *unstructured.Unstructured, groupReso
// want to dynamically re-provision it.
return warnings, errs, itemExists
} else {
obj, err = ctx.handleSkippedPVHasRetainPolicy(obj, resourceID, restoreLogger)
obj, err = ctx.handleSkippedPVHasRetainPolicy(obj, restoreLogger)
if err != nil {
errs.Add(namespace, err)
return warnings, errs, itemExists
@@ -1322,7 +1322,7 @@ func (ctx *restoreContext) restoreItem(obj *unstructured.Unstructured, groupReso
return warnings, errs, itemExists
default:
obj, err = ctx.handleSkippedPVHasRetainPolicy(obj, resourceID, restoreLogger)
obj, err = ctx.handleSkippedPVHasRetainPolicy(obj, restoreLogger)
if err != nil {
errs.Add(namespace, err)
return warnings, errs, itemExists
@@ -2497,7 +2497,6 @@ func (ctx *restoreContext) handlePVHasNativeSnapshot(obj *unstructured.Unstructu
func (ctx *restoreContext) handleSkippedPVHasRetainPolicy(
obj *unstructured.Unstructured,
resourceID string,
logger logrus.FieldLogger,
) (*unstructured.Unstructured, error) {
logger.Infof("Restoring persistent volume as-is because it doesn't have a snapshot and its reclaim policy is not Delete.")
@@ -2508,13 +2507,5 @@ func (ctx *restoreContext) handleSkippedPVHasRetainPolicy(
}
obj = resetVolumeBindingInfo(obj)
// We call the pvRestorer here to clear out the PV's claimRef.UID,
// so it can be re-claimed when its PVC is restored and gets a new UID.
updatedObj, err := ctx.pvRestorer.executePVAction(obj)
if err != nil {
return nil, fmt.Errorf("error executing PVAction for %s: %v", resourceID, err)
}
return updatedObj, nil
return obj, nil
}