diff --git a/changelogs/unreleased/9862-Lyndon-Li b/changelogs/unreleased/9862-Lyndon-Li new file mode 100644 index 000000000..04601d969 --- /dev/null +++ b/changelogs/unreleased/9862-Lyndon-Li @@ -0,0 +1 @@ +Enhance backup exposer for block data mover \ No newline at end of file diff --git a/pkg/controller/data_download_controller.go b/pkg/controller/data_download_controller.go index 36fec450b..738334ceb 100644 --- a/pkg/controller/data_download_controller.go +++ b/pkg/controller/data_download_controller.go @@ -150,7 +150,7 @@ func (r *DataDownloadReconciler) Reconcile(ctx context.Context, req ctrl.Request return ctrl.Result{}, err } - if !datamover.IsBuiltInUploader(dd.Spec.DataMover) { + if !datamover.IsBuiltInDataMover(dd.Spec.DataMover) { log.WithField("data mover", dd.Spec.DataMover).Info("it is not one built-in data mover which is not supported by Velero") return ctrl.Result{}, nil } diff --git a/pkg/controller/data_upload_controller.go b/pkg/controller/data_upload_controller.go index d7248e7c3..9aaf3d653 100644 --- a/pkg/controller/data_upload_controller.go +++ b/pkg/controller/data_upload_controller.go @@ -159,7 +159,7 @@ func (r *DataUploadReconciler) Reconcile(ctx context.Context, req ctrl.Request) return ctrl.Result{}, errors.Wrap(err, "getting DataUpload") } - if !datamover.IsBuiltInUploader(du.Spec.DataMover) { + if !datamover.IsBuiltInDataMover(du.Spec.DataMover) { log.WithField("Data mover", du.Spec.DataMover).Debug("it is not one built-in data mover which is not supported by Velero") return ctrl.Result{}, nil } @@ -939,7 +939,12 @@ func (r *DataUploadReconciler) setupExposeParam(du *velerov2alpha1api.DataUpload return nil, errors.Wrapf(err, "failed to get source PV %s", pvc.Spec.VolumeName) } - nodeOS := kube.GetPVCAttachingNodeOS(pvc, r.kubeClient.CoreV1(), r.kubeClient.StorageV1(), log) + nodeOS := "" + if du.Spec.DataMover == datamover.DataMoverTypeVeleroBlock { + nodeOS = kube.NodeOSLinux + } else { + nodeOS = kube.GetPVCAttachingNodeOS(pvc, r.kubeClient.CoreV1(), r.kubeClient.StorageV1(), log) + } if err := kube.HasNodeWithOS(context.Background(), nodeOS, r.kubeClient.CoreV1()); err != nil { return nil, errors.Wrapf(err, "no appropriate node to run data upload for PVC %s/%s", du.Spec.SourceNamespace, du.Spec.SourcePVC) @@ -1013,6 +1018,7 @@ func (r *DataUploadReconciler) setupExposeParam(du *velerov2alpha1api.DataUpload Resources: r.podResources, NodeOS: nodeOS, PriorityClassName: r.dataMovePriorityClass, + DataMover: du.Spec.DataMover, SnapshotMetadataServiceConfigs: r.snapshotMetadataServiceConfigs, }, nil } diff --git a/pkg/datamover/dataupload_delete_action.go b/pkg/datamover/dataupload_delete_action.go index 6b36e1068..46c62a1f1 100644 --- a/pkg/datamover/dataupload_delete_action.go +++ b/pkg/datamover/dataupload_delete_action.go @@ -88,7 +88,7 @@ func (d *DataUploadDeleteAction) Execute(input *velero.DeleteItemActionExecuteIn // generate the configmap which is to be created and used as a way to communicate the snapshot info to the backup deletion controller func genConfigmap(bak *velerov1.Backup, du velerov2alpha1.DataUpload) *corev1api.ConfigMap { - if !IsBuiltInUploader(du.Spec.DataMover) || du.Status.SnapshotID == "" { + if !IsBuiltInDataMover(du.Spec.DataMover) || du.Status.SnapshotID == "" { return nil } snapshot := repotypes.SnapshotIdentifier{ diff --git a/pkg/datamover/util.go b/pkg/datamover/util.go index e4097f07e..7e37695b6 100644 --- a/pkg/datamover/util.go +++ b/pkg/datamover/util.go @@ -18,6 +18,11 @@ package datamover import "fmt" +const ( + DataMoverTypeVeleroFs string = "velero-fs" + DataMoverTypeVeleroBlock string = "velero-block" +) + func GetUploaderType(dataMover string) string { if dataMover == "" || dataMover == "velero" { return "kopia" @@ -26,7 +31,7 @@ func GetUploaderType(dataMover string) string { } } -func IsBuiltInUploader(dataMover string) bool { +func IsBuiltInDataMover(dataMover string) bool { return dataMover == "" || dataMover == "velero" } diff --git a/pkg/datamover/util_test.go b/pkg/datamover/util_test.go index d4e3e6efe..80e2f4e16 100644 --- a/pkg/datamover/util_test.go +++ b/pkg/datamover/util_test.go @@ -30,7 +30,7 @@ func TestIsBuiltInUploader(t *testing.T) { } for _, tc := range testcases { t.Run(tc.name, func(tt *testing.T) { - assert.Equal(tt, tc.want, IsBuiltInUploader(tc.dataMover)) + assert.Equal(tt, tc.want, IsBuiltInDataMover(tc.dataMover)) }) } } diff --git a/pkg/exposer/csi_snapshot.go b/pkg/exposer/csi_snapshot.go index de1bf83c8..33dafff99 100644 --- a/pkg/exposer/csi_snapshot.go +++ b/pkg/exposer/csi_snapshot.go @@ -19,6 +19,7 @@ package exposer import ( "context" "fmt" + "maps" "time" snapshotv1api "github.com/kubernetes-csi/external-snapshotter/client/v8/apis/volumesnapshot/v1" @@ -33,6 +34,7 @@ import ( "k8s.io/client-go/kubernetes" "sigs.k8s.io/controller-runtime/pkg/client" + "github.com/vmware-tanzu/velero/pkg/datamover" "github.com/vmware-tanzu/velero/pkg/nodeagent" velerotypes "github.com/vmware-tanzu/velero/pkg/types" "github.com/vmware-tanzu/velero/pkg/util" @@ -94,6 +96,9 @@ type CSISnapshotExposeParam struct { // PriorityClassName is the priority class name for the data mover pod PriorityClassName string + // DataMover is the data mover type, e.g., velero-fs, velero-block + DataMover string + // SnapshotMetadataServiceConfigs is the config for CSI snapshot metadata service SnapshotMetadataServiceConfigs *velerotypes.CSISnapshotMetadataService } @@ -237,7 +242,7 @@ func (e *csiSnapshotExposer) Expose(ctx context.Context, ownerObject corev1api.O } } - backupPVC, err := e.createBackupPVC(ctx, ownerObject, backupVS.Name, backupPVCStorageClass, csiExposeParam.AccessMode, volumeSize, backupPVCReadOnly, backupPVCAnnotations) + backupPVC, err := e.createBackupPVC(ctx, ownerObject, backupVS.Name, backupPVCStorageClass, csiExposeParam.AccessMode, volumeSize, backupPVCReadOnly, backupPVCAnnotations, csiExposeParam.DataMover) if err != nil { return errors.Wrap(err, "error to create backup pvc") } @@ -454,7 +459,11 @@ func (e *csiSnapshotExposer) CleanUp(ctx context.Context, ownerObject corev1api. csi.DeleteVolumeSnapshotIfAny(ctx, e.csiSnapshotClient, vsName, sourceNamespace, e.log) } -func getVolumeModeByAccessMode(accessMode string) (corev1api.PersistentVolumeMode, error) { +func getVolumeModeByAccessMode(accessMode string, dataMover string) (corev1api.PersistentVolumeMode, error) { + if dataMover == datamover.DataMoverTypeVeleroBlock { + return corev1api.PersistentVolumeBlock, nil + } + switch accessMode { case AccessModeFileSystem: return corev1api.PersistentVolumeFilesystem, nil @@ -492,10 +501,14 @@ func (e *csiSnapshotExposer) createBackupVS(ctx context.Context, ownerObject cor func (e *csiSnapshotExposer) createBackupVSC(ctx context.Context, ownerObject corev1api.ObjectReference, snapshotVSC *snapshotv1api.VolumeSnapshotContent, vs *snapshotv1api.VolumeSnapshot) (*snapshotv1api.VolumeSnapshotContent, error) { backupVSCName := ownerObject.Name + anno := make(map[string]string) + maps.Copy(anno, snapshotVSC.Annotations) + anno[kube.KubeAnnAllowVolumeModeChange] = "true" + vsc := &snapshotv1api.VolumeSnapshotContent{ ObjectMeta: metav1.ObjectMeta{ Name: backupVSCName, - Annotations: snapshotVSC.Annotations, + Annotations: anno, Labels: map[string]string{}, }, Spec: snapshotv1api.VolumeSnapshotContentSpec{ @@ -528,10 +541,10 @@ func (e *csiSnapshotExposer) createBackupVSC(ctx context.Context, ownerObject co return e.csiSnapshotClient.VolumeSnapshotContents().Create(ctx, vsc, metav1.CreateOptions{}) } -func (e *csiSnapshotExposer) createBackupPVC(ctx context.Context, ownerObject corev1api.ObjectReference, backupVS, storageClass, accessMode string, resource resource.Quantity, readOnly bool, annotations map[string]string) (*corev1api.PersistentVolumeClaim, error) { +func (e *csiSnapshotExposer) createBackupPVC(ctx context.Context, ownerObject corev1api.ObjectReference, backupVS, storageClass, accessMode string, resource resource.Quantity, readOnly bool, annotations map[string]string, dataMover string) (*corev1api.PersistentVolumeClaim, error) { backupPVCName := ownerObject.Name - volumeMode, err := getVolumeModeByAccessMode(accessMode) + volumeMode, err := getVolumeModeByAccessMode(accessMode, dataMover) if err != nil { return nil, err } diff --git a/pkg/exposer/csi_snapshot_test.go b/pkg/exposer/csi_snapshot_test.go index e1a9860eb..8e00654f9 100644 --- a/pkg/exposer/csi_snapshot_test.go +++ b/pkg/exposer/csi_snapshot_test.go @@ -18,6 +18,7 @@ package exposer import ( "fmt" + "maps" "testing" "time" @@ -1056,21 +1057,25 @@ func TestExpose(t *testing.T) { backupPVC, err := exposer.kubeClient.CoreV1().PersistentVolumeClaims(ownerObject.Namespace).Get(t.Context(), ownerObject.Name, metav1.GetOptions{}) require.NoError(t, err) - expectedVS, err := exposer.csiSnapshotClient.VolumeSnapshots(ownerObject.Namespace).Get(t.Context(), ownerObject.Name, metav1.GetOptions{}) + backupVS, err := exposer.csiSnapshotClient.VolumeSnapshots(ownerObject.Namespace).Get(t.Context(), ownerObject.Name, metav1.GetOptions{}) require.NoError(t, err) - expectedVSC, err := exposer.csiSnapshotClient.VolumeSnapshotContents().Get(t.Context(), ownerObject.Name, metav1.GetOptions{}) + backupVSC, err := exposer.csiSnapshotClient.VolumeSnapshotContents().Get(t.Context(), ownerObject.Name, metav1.GetOptions{}) require.NoError(t, err) - assert.Equal(t, expectedVS.Annotations, vsObject.Annotations) - assert.Equal(t, *expectedVS.Spec.VolumeSnapshotClassName, *vsObject.Spec.VolumeSnapshotClassName) - assert.Equal(t, expectedVSC.Name, *expectedVS.Spec.Source.VolumeSnapshotContentName) + assert.Equal(t, vsObject.Annotations, backupVS.Annotations) + assert.Equal(t, *vsObject.Spec.VolumeSnapshotClassName, *backupVS.Spec.VolumeSnapshotClassName) + assert.Equal(t, *backupVS.Spec.Source.VolumeSnapshotContentName, backupVSC.Name) - assert.Equal(t, expectedVSC.Annotations, vscObj.Annotations) - assert.Equal(t, expectedVSC.Labels, vscObj.Labels) - assert.Equal(t, expectedVSC.Spec.DeletionPolicy, vscObj.Spec.DeletionPolicy) - assert.Equal(t, expectedVSC.Spec.Driver, vscObj.Spec.Driver) - assert.Equal(t, *expectedVSC.Spec.VolumeSnapshotClassName, *vscObj.Spec.VolumeSnapshotClassName) + anno := make(map[string]string) + maps.Copy(anno, vscObj.Annotations) + anno[kube.KubeAnnAllowVolumeModeChange] = "true" + + assert.Equal(t, anno, backupVSC.Annotations) + assert.Equal(t, vscObj.Labels, backupVSC.Labels) + assert.Equal(t, vscObj.Spec.DeletionPolicy, backupVSC.Spec.DeletionPolicy) + assert.Equal(t, vscObj.Spec.Driver, backupVSC.Spec.Driver) + assert.Equal(t, *vscObj.Spec.VolumeSnapshotClassName, *backupVSC.Spec.VolumeSnapshotClassName) if test.expectedVolumeSize != nil { assert.Equal(t, *test.expectedVolumeSize, backupPVC.Spec.Resources.Requests[corev1api.ResourceStorage]) @@ -1514,7 +1519,7 @@ func Test_csiSnapshotExposer_createBackupPVC(t *testing.T) { APIVersion: tt.ownerBackup.APIVersion, } } - got, err := e.createBackupPVC(t.Context(), ownerObject, tt.backupVS, tt.storageClass, tt.accessMode, tt.resource, tt.readOnly, map[string]string{}) + got, err := e.createBackupPVC(t.Context(), ownerObject, tt.backupVS, tt.storageClass, tt.accessMode, tt.resource, tt.readOnly, map[string]string{}, "") if !tt.wantErr(t, err, fmt.Sprintf("createBackupPVC(%v, %v, %v, %v, %v, %v)", ownerObject, tt.backupVS, tt.storageClass, tt.accessMode, tt.resource, tt.readOnly)) { return } diff --git a/pkg/util/kube/utils.go b/pkg/util/kube/utils.go index 5e5e97603..d93effd7f 100644 --- a/pkg/util/kube/utils.go +++ b/pkg/util/kube/utils.go @@ -53,6 +53,7 @@ const ( KubeAnnDynamicallyProvisioned = "pv.kubernetes.io/provisioned-by" KubeAnnMigratedTo = "pv.kubernetes.io/migrated-to" KubeAnnSelectedNode = "volume.kubernetes.io/selected-node" + KubeAnnAllowVolumeModeChange = "snapshot.storage.kubernetes.io/allow-volume-mode-change" ) // VolumeSnapshotContentManagedByLabel is applied by the snapshot controller