From 22be7e321831b8f1df6bad9158a9b48875e2ccf4 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Mon, 11 May 2026 13:53:49 +0800 Subject: [PATCH 1/2] data path naming adjustment Signed-off-by: Lyndon-Li --- pkg/datamover/backup_micro_service.go | 6 +- pkg/datamover/backup_micro_service_test.go | 2 +- pkg/datamover/restore_micro_service.go | 8 +- pkg/datamover/restore_micro_service_test.go | 2 +- pkg/datapath/{file_system.go => data_path.go} | 136 +++++++++--------- ...{file_system_test.go => data_path_test.go} | 24 ++-- pkg/datapath/manager.go | 8 +- pkg/datapath/manager_test.go | 8 +- pkg/podvolume/backup_micro_service.go | 6 +- pkg/podvolume/backup_micro_service_test.go | 2 +- pkg/podvolume/restore_micro_service.go | 4 +- pkg/podvolume/restore_micro_service_test.go | 2 +- pkg/uploader/provider/kopia.go | 9 +- pkg/uploader/provider/kopia_test.go | 4 +- 14 files changed, 110 insertions(+), 111 deletions(-) rename pkg/datapath/{file_system.go => data_path.go} (54%) rename pkg/datapath/{file_system_test.go => data_path_test.go} (89%) diff --git a/pkg/datamover/backup_micro_service.go b/pkg/datamover/backup_micro_service.go index 61c308f15..4edd24c97 100644 --- a/pkg/datamover/backup_micro_service.go +++ b/pkg/datamover/backup_micro_service.go @@ -170,14 +170,14 @@ func (r *BackupMicroService) RunCancelableDataPath(ctx context.Context) (string, OnProgress: r.OnDataUploadProgress, } - fsBackup, err := r.dataPathMgr.CreateFileSystemBR(du.Name, dataUploadDownloadRequestor, ctx, r.client, du.Namespace, callbacks, log) + dp, err := r.dataPathMgr.CreateGenericDataPath(du.Name, dataUploadDownloadRequestor, ctx, r.client, du.Namespace, callbacks, log) if err != nil { return "", errors.Wrap(err, "error to create data path") } log.Debug("Async fs br created") - if err := fsBackup.Init(ctx, &datapath.FSBRInitParam{ + if err := dp.Init(ctx, &datapath.InitParam{ BSLName: du.Spec.BackupStorageLocation, SourceNamespace: du.Spec.SourceNamespace, UploaderType: GetUploaderType(du.Spec.DataMover), @@ -195,7 +195,7 @@ func (r *BackupMicroService) RunCancelableDataPath(ctx context.Context) (string, velerov1api.AsyncOperationIDLabel: du.Labels[velerov1api.AsyncOperationIDLabel], } - if err := fsBackup.StartBackup(r.sourceTargetPath, du.Spec.DataMoverConfig, &datapath.FSBRStartParam{ + if err := dp.StartBackup(r.sourceTargetPath, du.Spec.DataMoverConfig, &datapath.BackupStartParam{ RealSource: GetRealSource(du.Spec.SourceNamespace, du.Spec.SourcePVC), ParentSnapshot: "", ForceFull: false, diff --git a/pkg/datamover/backup_micro_service_test.go b/pkg/datamover/backup_micro_service_test.go index c5a9a4487..79b4a834a 100644 --- a/pkg/datamover/backup_micro_service_test.go +++ b/pkg/datamover/backup_micro_service_test.go @@ -403,7 +403,7 @@ func TestRunCancelableDataPath(t *testing.T) { bs.dataPathMgr = test.dataPathMgr } - datapath.FSBRCreator = func(string, string, kbclient.Client, string, datapath.Callbacks, logrus.FieldLogger) datapath.AsyncBR { + datapath.VGDPCreator = func(string, string, kbclient.Client, string, datapath.Callbacks, logrus.FieldLogger) datapath.AsyncBR { fsBR := datapathmockes.NewAsyncBR(t) if test.initErr != nil { fsBR.On("Init", mock.Anything, mock.Anything).Return(test.initErr) diff --git a/pkg/datamover/restore_micro_service.go b/pkg/datamover/restore_micro_service.go index 0cb5f18ec..e1c95ee3d 100644 --- a/pkg/datamover/restore_micro_service.go +++ b/pkg/datamover/restore_micro_service.go @@ -159,14 +159,14 @@ func (r *RestoreMicroService) RunCancelableDataPath(ctx context.Context) (string OnProgress: r.OnDataDownloadProgress, } - fsRestore, err := r.dataPathMgr.CreateFileSystemBR(dd.Name, dataUploadDownloadRequestor, ctx, r.client, dd.Namespace, callbacks, log) + dp, err := r.dataPathMgr.CreateGenericDataPath(dd.Name, dataUploadDownloadRequestor, ctx, r.client, dd.Namespace, callbacks, log) if err != nil { return "", errors.Wrap(err, "error to create data path") } log.Debug("Found volume path") - if err := fsRestore.Init(ctx, - &datapath.FSBRInitParam{ + if err := dp.Init(ctx, + &datapath.InitParam{ BSLName: dd.Spec.BackupStorageLocation, SourceNamespace: dd.Spec.SourceNamespace, UploaderType: GetUploaderType(dd.Spec.DataMover), @@ -180,7 +180,7 @@ func (r *RestoreMicroService) RunCancelableDataPath(ctx context.Context) (string } log.Info("fs init") - if err := fsRestore.StartRestore(dd.Spec.SnapshotID, r.sourceTargetPath, dd.Spec.DataMoverConfig); err != nil { + if err := dp.StartRestore(dd.Spec.SnapshotID, r.sourceTargetPath, dd.Spec.DataMoverConfig); err != nil { return "", errors.Wrap(err, "error starting data path restore") } diff --git a/pkg/datamover/restore_micro_service_test.go b/pkg/datamover/restore_micro_service_test.go index a53e4d6d5..bce4d94b7 100644 --- a/pkg/datamover/restore_micro_service_test.go +++ b/pkg/datamover/restore_micro_service_test.go @@ -347,7 +347,7 @@ func TestRunCancelableRestore(t *testing.T) { rs.dataPathMgr = test.dataPathMgr } - datapath.FSBRCreator = func(string, string, kbclient.Client, string, datapath.Callbacks, logrus.FieldLogger) datapath.AsyncBR { + datapath.VGDPCreator = func(string, string, kbclient.Client, string, datapath.Callbacks, logrus.FieldLogger) datapath.AsyncBR { fsBR := datapathmockes.NewAsyncBR(t) if test.initErr != nil { fsBR.On("Init", mock.Anything, mock.Anything).Return(test.initErr) diff --git a/pkg/datapath/file_system.go b/pkg/datapath/data_path.go similarity index 54% rename from pkg/datapath/file_system.go rename to pkg/datapath/data_path.go index 61fac1e47..6ea9fc9a9 100644 --- a/pkg/datapath/file_system.go +++ b/pkg/datapath/data_path.go @@ -34,8 +34,8 @@ import ( "github.com/vmware-tanzu/velero/pkg/util/filesystem" ) -// FSBRInitParam define the input param for FSBR init -type FSBRInitParam struct { +// InitParam define the input param for data path init +type InitParam struct { BSLName string SourceNamespace string UploaderType string @@ -47,15 +47,15 @@ type FSBRInitParam struct { CacheDir string } -// FSBRStartParam define the input param for FSBR start -type FSBRStartParam struct { +// BackupStartParam define the input param for backup start +type BackupStartParam struct { RealSource string ParentSnapshot string ForceFull bool Tags map[string]string } -type fileSystemBR struct { +type genearalDataPath struct { ctx context.Context cancel context.CancelFunc backupRepo *velerov1api.BackupRepository @@ -72,8 +72,8 @@ type fileSystemBR struct { dataPathLock sync.Mutex } -func newFileSystemBR(jobName string, requestorType string, client client.Client, namespace string, callbacks Callbacks, log logrus.FieldLogger) AsyncBR { - fs := &fileSystemBR{ +func newGeneralDataPath(jobName string, requestorType string, client client.Client, namespace string, callbacks Callbacks, log logrus.FieldLogger) AsyncBR { + dp := &genearalDataPath{ jobName: jobName, requestorType: requestorType, client: client, @@ -83,151 +83,151 @@ func newFileSystemBR(jobName string, requestorType string, client client.Client, log: log, } - return fs + return dp } -func (fs *fileSystemBR) Init(ctx context.Context, param any) error { - initParam := param.(*FSBRInitParam) +func (dp *genearalDataPath) Init(ctx context.Context, param any) error { + initParam := param.(*InitParam) var err error defer func() { if err != nil { - fs.Close(ctx) + dp.Close(ctx) } }() - fs.ctx, fs.cancel = context.WithCancel(ctx) + dp.ctx, dp.cancel = context.WithCancel(ctx) backupLocation := &velerov1api.BackupStorageLocation{} - if err = fs.client.Get(ctx, client.ObjectKey{ - Namespace: fs.namespace, + if err = dp.client.Get(ctx, client.ObjectKey{ + Namespace: dp.namespace, Name: initParam.BSLName, }, backupLocation); err != nil { return errors.Wrapf(err, "error getting backup storage location %s", initParam.BSLName) } - fs.backupLocation = backupLocation + dp.backupLocation = backupLocation - fs.backupRepo, err = initParam.RepositoryEnsurer.EnsureRepo(ctx, fs.namespace, initParam.SourceNamespace, initParam.BSLName, initParam.RepositoryType) + dp.backupRepo, err = initParam.RepositoryEnsurer.EnsureRepo(ctx, dp.namespace, initParam.SourceNamespace, initParam.BSLName, initParam.RepositoryType) if err != nil { return errors.Wrapf(err, "error to ensure backup repository %s-%s-%s", initParam.BSLName, initParam.SourceNamespace, initParam.RepositoryType) } - err = fs.boostRepoConnect(ctx, initParam.RepositoryType, initParam.CredentialGetter, initParam.CacheDir) + err = dp.boostRepoConnect(ctx, initParam.RepositoryType, initParam.CredentialGetter, initParam.CacheDir) if err != nil { return errors.Wrapf(err, "error to boost backup repository connection %s-%s-%s", initParam.BSLName, initParam.SourceNamespace, initParam.RepositoryType) } - fs.uploaderProv, err = provider.NewUploaderProvider(ctx, fs.client, initParam.UploaderType, fs.requestorType, initParam.RepoIdentifier, - fs.backupLocation, fs.backupRepo, initParam.CredentialGetter, repokey.RepoKeySelector(), fs.log) + dp.uploaderProv, err = provider.NewUploaderProvider(ctx, dp.client, initParam.UploaderType, dp.requestorType, initParam.RepoIdentifier, + dp.backupLocation, dp.backupRepo, initParam.CredentialGetter, repokey.RepoKeySelector(), dp.log) if err != nil { return errors.Wrapf(err, "error creating uploader %s", initParam.UploaderType) } - fs.initialized = true + dp.initialized = true - fs.log.WithFields( + dp.log.WithFields( logrus.Fields{ - "jobName": fs.jobName, + "jobName": dp.jobName, "bsl": initParam.BSLName, "source namespace": initParam.SourceNamespace, "uploader": initParam.UploaderType, "repository": initParam.RepositoryType, - }).Info("FileSystemBR is initialized") + }).Info("Data path is initialized") return nil } -func (fs *fileSystemBR) Close(ctx context.Context) { - if fs.cancel != nil { - fs.cancel() +func (dp *genearalDataPath) Close(ctx context.Context) { + if dp.cancel != nil { + dp.cancel() } - fs.log.WithField("user", fs.jobName).Info("Closing FileSystemBR") + dp.log.WithField("user", dp.jobName).Info("Closing data path") - fs.wgDataPath.Wait() + dp.wgDataPath.Wait() - fs.close(ctx) + dp.close(ctx) - fs.log.WithField("user", fs.jobName).Info("FileSystemBR is closed") + dp.log.WithField("user", dp.jobName).Info("Data path is closed") } -func (fs *fileSystemBR) close(ctx context.Context) { - fs.dataPathLock.Lock() - defer fs.dataPathLock.Unlock() +func (dp *genearalDataPath) close(ctx context.Context) { + dp.dataPathLock.Lock() + defer dp.dataPathLock.Unlock() - if fs.uploaderProv != nil { - if err := fs.uploaderProv.Close(ctx); err != nil { - fs.log.Errorf("failed to close uploader provider with error %v", err) + if dp.uploaderProv != nil { + if err := dp.uploaderProv.Close(ctx); err != nil { + dp.log.Errorf("failed to close uploader provider with error %v", err) } - fs.uploaderProv = nil + dp.uploaderProv = nil } } -func (fs *fileSystemBR) StartBackup(source AccessPoint, uploaderConfig map[string]string, param any) error { - if !fs.initialized { +func (dp *genearalDataPath) StartBackup(source AccessPoint, uploaderConfig map[string]string, param any) error { + if !dp.initialized { return errors.New("file system data path is not initialized") } - fs.wgDataPath.Add(1) + dp.wgDataPath.Add(1) - backupParam := param.(*FSBRStartParam) + backupParam := param.(*BackupStartParam) go func() { - fs.log.Info("Start data path backup") + dp.log.Info("Start data path backup") defer func() { - fs.close(context.Background()) - fs.wgDataPath.Done() + dp.close(context.Background()) + dp.wgDataPath.Done() }() - snapshotID, emptySnapshot, totalBytes, incrementalBytes, err := fs.uploaderProv.RunBackup(fs.ctx, source.ByPath, backupParam.RealSource, backupParam.Tags, backupParam.ForceFull, - backupParam.ParentSnapshot, source.VolMode, uploaderConfig, fs) + snapshotID, emptySnapshot, totalBytes, incrementalBytes, err := dp.uploaderProv.RunBackup(dp.ctx, source.ByPath, backupParam.RealSource, backupParam.Tags, backupParam.ForceFull, + backupParam.ParentSnapshot, source.VolMode, uploaderConfig, dp) if err == provider.ErrorCanceled { - fs.callbacks.OnCancelled(context.Background(), fs.namespace, fs.jobName) + dp.callbacks.OnCancelled(context.Background(), dp.namespace, dp.jobName) } else if err != nil { dataPathErr := DataPathError{ snapshotID: snapshotID, err: err, } - fs.callbacks.OnFailed(context.Background(), fs.namespace, fs.jobName, dataPathErr) + dp.callbacks.OnFailed(context.Background(), dp.namespace, dp.jobName, dataPathErr) } else { - fs.callbacks.OnCompleted(context.Background(), fs.namespace, fs.jobName, Result{Backup: BackupResult{snapshotID, emptySnapshot, source, totalBytes, incrementalBytes}}) + dp.callbacks.OnCompleted(context.Background(), dp.namespace, dp.jobName, Result{Backup: BackupResult{snapshotID, emptySnapshot, source, totalBytes, incrementalBytes}}) } }() return nil } -func (fs *fileSystemBR) StartRestore(snapshotID string, target AccessPoint, uploaderConfigs map[string]string) error { - if !fs.initialized { - return errors.New("file system data path is not initialized") +func (dp *genearalDataPath) StartRestore(snapshotID string, target AccessPoint, uploaderConfigs map[string]string) error { + if !dp.initialized { + return errors.New("data path is not initialized") } - fs.wgDataPath.Add(1) + dp.wgDataPath.Add(1) go func() { - fs.log.Info("Start data path restore") + dp.log.Info("Start data path restore") defer func() { - fs.close(context.Background()) - fs.wgDataPath.Done() + dp.close(context.Background()) + dp.wgDataPath.Done() }() - totalBytes, err := fs.uploaderProv.RunRestore(fs.ctx, snapshotID, target.ByPath, target.VolMode, uploaderConfigs, fs) + totalBytes, err := dp.uploaderProv.RunRestore(dp.ctx, snapshotID, target.ByPath, target.VolMode, uploaderConfigs, dp) if err == provider.ErrorCanceled { - fs.callbacks.OnCancelled(context.Background(), fs.namespace, fs.jobName) + dp.callbacks.OnCancelled(context.Background(), dp.namespace, dp.jobName) } else if err != nil { dataPathErr := DataPathError{ snapshotID: snapshotID, err: err, } - fs.callbacks.OnFailed(context.Background(), fs.namespace, fs.jobName, dataPathErr) + dp.callbacks.OnFailed(context.Background(), dp.namespace, dp.jobName, dataPathErr) } else { - fs.callbacks.OnCompleted(context.Background(), fs.namespace, fs.jobName, Result{Restore: RestoreResult{Target: target, TotalBytes: totalBytes}}) + dp.callbacks.OnCompleted(context.Background(), dp.namespace, dp.jobName, Result{Restore: RestoreResult{Target: target, TotalBytes: totalBytes}}) } }() @@ -235,20 +235,20 @@ func (fs *fileSystemBR) StartRestore(snapshotID string, target AccessPoint, uplo } // UpdateProgress which implement ProgressUpdater interface to update progress status -func (fs *fileSystemBR) UpdateProgress(p *uploader.Progress) { - if fs.callbacks.OnProgress != nil { - fs.callbacks.OnProgress(context.Background(), fs.namespace, fs.jobName, &uploader.Progress{TotalBytes: p.TotalBytes, BytesDone: p.BytesDone}) +func (dp *genearalDataPath) UpdateProgress(p *uploader.Progress) { + if dp.callbacks.OnProgress != nil { + dp.callbacks.OnProgress(context.Background(), dp.namespace, dp.jobName, &uploader.Progress{TotalBytes: p.TotalBytes, BytesDone: p.BytesDone}) } } -func (fs *fileSystemBR) Cancel() { - fs.cancel() - fs.log.WithField("user", fs.jobName).Info("FileSystemBR is canceled") +func (dp *genearalDataPath) Cancel() { + dp.cancel() + dp.log.WithField("user", dp.jobName).Info("FileSystemBR is canceled") } -func (fs *fileSystemBR) boostRepoConnect(ctx context.Context, repositoryType string, credentialGetter *credentials.CredentialGetter, cacheDir string) error { +func (dp *genearalDataPath) boostRepoConnect(ctx context.Context, repositoryType string, credentialGetter *credentials.CredentialGetter, cacheDir string) error { if repositoryType == velerov1api.BackupRepositoryTypeKopia { - if err := repoProvider.NewUnifiedRepoProvider(*credentialGetter, repositoryType, fs.log).BoostRepoConnect(ctx, repoProvider.RepoParam{BackupLocation: fs.backupLocation, BackupRepo: fs.backupRepo, CacheDir: cacheDir}); err != nil { + if err := repoProvider.NewUnifiedRepoProvider(*credentialGetter, repositoryType, dp.log).BoostRepoConnect(ctx, repoProvider.RepoParam{BackupLocation: dp.backupLocation, BackupRepo: dp.backupRepo, CacheDir: cacheDir}); err != nil { return err } diff --git a/pkg/datapath/file_system_test.go b/pkg/datapath/data_path_test.go similarity index 89% rename from pkg/datapath/file_system_test.go rename to pkg/datapath/data_path_test.go index 3887a82e3..ab1213a51 100644 --- a/pkg/datapath/file_system_test.go +++ b/pkg/datapath/data_path_test.go @@ -94,21 +94,21 @@ func TestAsyncBackup(t *testing.T) { for _, test := range tests { t.Run(test.name, func(t *testing.T) { - fs := newFileSystemBR("job-1", "test", nil, "velero", Callbacks{}, velerotest.NewLogger()).(*fileSystemBR) + dp := newGeneralDataPath("job-1", "test", nil, "velero", Callbacks{}, velerotest.NewLogger()).(*genearalDataPath) mockProvider := providerMock.NewProvider(t) mockProvider.On("RunBackup", mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return(test.result.Backup.SnapshotID, test.result.Backup.EmptySnapshot, test.result.Backup.TotalBytes, test.result.Backup.IncrementalBytes, test.err) mockProvider.On("Close", mock.Anything).Return(nil) - fs.uploaderProv = mockProvider - fs.initialized = true - fs.callbacks = test.callbacks + dp.uploaderProv = mockProvider + dp.initialized = true + dp.callbacks = test.callbacks - err := fs.StartBackup(AccessPoint{ByPath: test.path}, map[string]string{}, &FSBRStartParam{}) + err := dp.StartBackup(AccessPoint{ByPath: test.path}, map[string]string{}, &BackupStartParam{}) require.NoError(t, err) <-finish // Ensure the goroutine finishes so deferred fs.close executes, satisfying mock expectations. - fs.wgDataPath.Wait() + dp.wgDataPath.Wait() assert.Equal(t, test.err, asyncErr) assert.Equal(t, test.result, asyncResult) @@ -182,21 +182,21 @@ func TestAsyncRestore(t *testing.T) { for _, test := range tests { t.Run(test.name, func(t *testing.T) { - fs := newFileSystemBR("job-1", "test", nil, "velero", Callbacks{}, velerotest.NewLogger()).(*fileSystemBR) + dp := newGeneralDataPath("job-1", "test", nil, "velero", Callbacks{}, velerotest.NewLogger()).(*genearalDataPath) mockProvider := providerMock.NewProvider(t) mockProvider.On("RunRestore", mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return(test.result.Restore.TotalBytes, test.err) mockProvider.On("Close", mock.Anything).Return(nil) - fs.uploaderProv = mockProvider - fs.initialized = true - fs.callbacks = test.callbacks + dp.uploaderProv = mockProvider + dp.initialized = true + dp.callbacks = test.callbacks - err := fs.StartRestore(test.snapshot, AccessPoint{ByPath: test.path}, map[string]string{}) + err := dp.StartRestore(test.snapshot, AccessPoint{ByPath: test.path}, map[string]string{}) require.NoError(t, err) <-finish // Ensure the goroutine finishes so deferred fs.close executes, satisfying mock expectations. - fs.wgDataPath.Wait() + dp.wgDataPath.Wait() assert.Equal(t, asyncErr, test.err) assert.Equal(t, asyncResult, test.result) diff --git a/pkg/datapath/manager.go b/pkg/datapath/manager.go index 0b790a5cc..79e692b08 100644 --- a/pkg/datapath/manager.go +++ b/pkg/datapath/manager.go @@ -28,7 +28,7 @@ import ( ) var ConcurrentLimitExceed error = errors.New("Concurrent number exceeds") -var FSBRCreator = newFileSystemBR +var VGDPCreator = newGeneralDataPath var MicroServiceBRWatcherCreator = newMicroServiceBRWatcher type Manager struct { @@ -45,8 +45,8 @@ func NewManager(cocurrentNum int) *Manager { } } -// CreateFileSystemBR creates a new file system backup/restore data path instance -func (m *Manager) CreateFileSystemBR(jobName string, requestorType string, ctx context.Context, client client.Client, namespace string, callbacks Callbacks, log logrus.FieldLogger) (AsyncBR, error) { +// CreateGenericDataPath creates a new generic data path instance +func (m *Manager) CreateGenericDataPath(jobName string, requestorType string, ctx context.Context, client client.Client, namespace string, callbacks Callbacks, log logrus.FieldLogger) (AsyncBR, error) { m.trackerLock.Lock() defer m.trackerLock.Unlock() @@ -54,7 +54,7 @@ func (m *Manager) CreateFileSystemBR(jobName string, requestorType string, ctx c return nil, ConcurrentLimitExceed } - m.tracker[jobName] = FSBRCreator(jobName, requestorType, client, namespace, callbacks, log) + m.tracker[jobName] = VGDPCreator(jobName, requestorType, client, namespace, callbacks, log) return m.tracker[jobName], nil } diff --git a/pkg/datapath/manager_test.go b/pkg/datapath/manager_test.go index c91517bac..8ee2b7c1d 100644 --- a/pkg/datapath/manager_test.go +++ b/pkg/datapath/manager_test.go @@ -23,16 +23,16 @@ import ( "github.com/stretchr/testify/require" ) -func TestCreateFileSystemBR(t *testing.T) { +func TestCreateGenericDataPath(t *testing.T) { m := NewManager(2) - async_job_1, err := m.CreateFileSystemBR("job-1", "test", t.Context(), nil, "velero", Callbacks{}, nil) + async_job_1, err := m.CreateGenericDataPath("job-1", "test", t.Context(), nil, "velero", Callbacks{}, nil) require.NoError(t, err) - _, err = m.CreateFileSystemBR("job-2", "test", t.Context(), nil, "velero", Callbacks{}, nil) + _, err = m.CreateGenericDataPath("job-2", "test", t.Context(), nil, "velero", Callbacks{}, nil) require.NoError(t, err) - _, err = m.CreateFileSystemBR("job-3", "test", t.Context(), nil, "velero", Callbacks{}, nil) + _, err = m.CreateGenericDataPath("job-3", "test", t.Context(), nil, "velero", Callbacks{}, nil) assert.Equal(t, ConcurrentLimitExceed, err) ret := m.GetAsyncBR("job-0") diff --git a/pkg/podvolume/backup_micro_service.go b/pkg/podvolume/backup_micro_service.go index 11ca66676..d393b2bcb 100644 --- a/pkg/podvolume/backup_micro_service.go +++ b/pkg/podvolume/backup_micro_service.go @@ -169,14 +169,14 @@ func (r *BackupMicroService) RunCancelableDataPath(ctx context.Context) (string, OnProgress: r.OnDataPathProgress, } - fsBackup, err := r.dataPathMgr.CreateFileSystemBR(pvb.Name, podVolumeRequestor, ctx, r.client, pvb.Namespace, callbacks, log) + fsBackup, err := r.dataPathMgr.CreateGenericDataPath(pvb.Name, podVolumeRequestor, ctx, r.client, pvb.Namespace, callbacks, log) if err != nil { return "", errors.Wrap(err, "error to create data path") } log.Debug("Async fs br created") - if err := fsBackup.Init(ctx, &datapath.FSBRInitParam{ + if err := fsBackup.Init(ctx, &datapath.InitParam{ BSLName: pvb.Spec.BackupStorageLocation, SourceNamespace: pvb.Spec.Pod.Namespace, UploaderType: pvb.Spec.UploaderType, @@ -192,7 +192,7 @@ func (r *BackupMicroService) RunCancelableDataPath(ctx context.Context) (string, tags := map[string]string{} - if err := fsBackup.StartBackup(r.sourceTargetPath, pvb.Spec.UploaderSettings, &datapath.FSBRStartParam{ + if err := fsBackup.StartBackup(r.sourceTargetPath, pvb.Spec.UploaderSettings, &datapath.BackupStartParam{ RealSource: GetRealSource(pvb), ParentSnapshot: "", ForceFull: false, diff --git a/pkg/podvolume/backup_micro_service_test.go b/pkg/podvolume/backup_micro_service_test.go index a5b6cd1f4..bec46f353 100644 --- a/pkg/podvolume/backup_micro_service_test.go +++ b/pkg/podvolume/backup_micro_service_test.go @@ -402,7 +402,7 @@ func TestRunCancelableDataPath(t *testing.T) { bs.dataPathMgr = test.dataPathMgr } - datapath.FSBRCreator = func(string, string, kbclient.Client, string, datapath.Callbacks, logrus.FieldLogger) datapath.AsyncBR { + datapath.VGDPCreator = func(string, string, kbclient.Client, string, datapath.Callbacks, logrus.FieldLogger) datapath.AsyncBR { fsBR := datapathmockes.NewAsyncBR(t) if test.initErr != nil { fsBR.On("Init", mock.Anything, mock.Anything).Return(test.initErr) diff --git a/pkg/podvolume/restore_micro_service.go b/pkg/podvolume/restore_micro_service.go index decf94b1b..77e84eb04 100644 --- a/pkg/podvolume/restore_micro_service.go +++ b/pkg/podvolume/restore_micro_service.go @@ -161,7 +161,7 @@ func (r *RestoreMicroService) RunCancelableDataPath(ctx context.Context) (string OnProgress: r.OnPvrProgress, } - fsRestore, err := r.dataPathMgr.CreateFileSystemBR(pvr.Name, podVolumeRequestor, ctx, r.client, pvr.Namespace, callbacks, log) + fsRestore, err := r.dataPathMgr.CreateGenericDataPath(pvr.Name, podVolumeRequestor, ctx, r.client, pvr.Namespace, callbacks, log) if err != nil { return "", errors.Wrap(err, "error to create data path") } @@ -169,7 +169,7 @@ func (r *RestoreMicroService) RunCancelableDataPath(ctx context.Context) (string log.Debug("Async fs br created") if err := fsRestore.Init(ctx, - &datapath.FSBRInitParam{ + &datapath.InitParam{ BSLName: pvr.Spec.BackupStorageLocation, SourceNamespace: pvr.Spec.SourceNamespace, UploaderType: pvr.Spec.UploaderType, diff --git a/pkg/podvolume/restore_micro_service_test.go b/pkg/podvolume/restore_micro_service_test.go index 2c94ddf46..1d25e342f 100644 --- a/pkg/podvolume/restore_micro_service_test.go +++ b/pkg/podvolume/restore_micro_service_test.go @@ -428,7 +428,7 @@ func TestRunCancelableDataPathRestore(t *testing.T) { rs.dataPathMgr = test.dataPathMgr } - datapath.FSBRCreator = func(string, string, kbclient.Client, string, datapath.Callbacks, logrus.FieldLogger) datapath.AsyncBR { + datapath.VGDPCreator = func(string, string, kbclient.Client, string, datapath.Callbacks, logrus.FieldLogger) datapath.AsyncBR { fsBR := datapathmockes.NewAsyncBR(t) if test.initErr != nil { fsBR.On("Init", mock.Anything, mock.Anything).Return(test.initErr) diff --git a/pkg/uploader/provider/kopia.go b/pkg/uploader/provider/kopia.go index 16d8aeb52..e51786c87 100644 --- a/pkg/uploader/provider/kopia.go +++ b/pkg/uploader/provider/kopia.go @@ -36,9 +36,8 @@ import ( "github.com/vmware-tanzu/velero/pkg/repository/udmrepo/service" ) -// BackupFunc mainly used to make testing more convenient -var BackupFunc = kopia.Backup -var RestoreFunc = kopia.Restore +var kopiaBackupFunc = kopia.Backup +var kopiaRestoreFunc = kopia.Restore var BackupRepoServiceCreateFunc = service.Create // kopiaProvider recorded info related with kopiaProvider @@ -165,7 +164,7 @@ func (kp *kopiaProvider) RunBackup( uploaderCfg[kopia.UploaderConfigMultipartKey] = "true" } - snapshotInfo, _, err := BackupFunc(ctx, kpUploader, repoWriter, path, realSource, forceFull, parentSnapshot, volMode, uploaderCfg, tags, log) + snapshotInfo, _, err := kopiaBackupFunc(ctx, kpUploader, repoWriter, path, realSource, forceFull, parentSnapshot, volMode, uploaderCfg, tags, log) if err != nil { snapshotID := "" if snapshotInfo != nil { @@ -233,7 +232,7 @@ func (kp *kopiaProvider) RunRestore( // We use the cancel channel to control the restore cancel, so don't pass a context with cancel to Kopia restore. // Otherwise, Kopia restore will not response to the cancel control but return an arbitrary error. // Kopia restore cancel is not designed as well as Kopia backup which uses the context to control backup cancel all the way. - size, fileCount, err := RestoreFunc(context.Background(), repoWriter, progress, snapshotID, volumePath, volMode, uploaderCfg, log, restoreCancel) + size, fileCount, err := kopiaRestoreFunc(context.Background(), repoWriter, progress, snapshotID, volumePath, volMode, uploaderCfg, log, restoreCancel) if err != nil { return 0, errors.Wrapf(err, "Failed to run kopia restore") diff --git a/pkg/uploader/provider/kopia_test.go b/pkg/uploader/provider/kopia_test.go index 734bdb176..c4e50d8fb 100644 --- a/pkg/uploader/provider/kopia_test.go +++ b/pkg/uploader/provider/kopia_test.go @@ -105,7 +105,7 @@ func TestRunBackup(t *testing.T) { if tc.volMode == "" { tc.volMode = uploader.PersistentVolumeFilesystem } - BackupFunc = tc.hookBackupFunc + kopiaBackupFunc = tc.hookBackupFunc _, _, _, _, err := kp.RunBackup(t.Context(), "var", "", nil, false, "", tc.volMode, map[string]string{}, &updater) if tc.notError { assert.NoError(t, err) @@ -156,7 +156,7 @@ func TestRunRestore(t *testing.T) { if tc.volMode == "" { tc.volMode = uploader.PersistentVolumeFilesystem } - RestoreFunc = tc.hookRestoreFunc + kopiaRestoreFunc = tc.hookRestoreFunc _, err := kp.RunRestore(t.Context(), "", "/var", tc.volMode, map[string]string{}, &updater) if tc.notError { assert.NoError(t, err) From 53f25cde2b6a4e6defe221a020f77a49b1750c76 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Thu, 21 May 2026 14:46:20 +0800 Subject: [PATCH 2/2] data path naming adjustment for block data mover Signed-off-by: Lyndon-Li --- changelogs/unreleased/9819-Lyndon-Li | 1 + pkg/datapath/data_path.go | 20 ++++++++++---------- pkg/datapath/data_path_test.go | 4 ++-- 3 files changed, 13 insertions(+), 12 deletions(-) create mode 100644 changelogs/unreleased/9819-Lyndon-Li diff --git a/changelogs/unreleased/9819-Lyndon-Li b/changelogs/unreleased/9819-Lyndon-Li new file mode 100644 index 000000000..1766476a8 --- /dev/null +++ b/changelogs/unreleased/9819-Lyndon-Li @@ -0,0 +1 @@ +Data path naming adjustment for block data mover \ No newline at end of file diff --git a/pkg/datapath/data_path.go b/pkg/datapath/data_path.go index 6ea9fc9a9..12cc00ac6 100644 --- a/pkg/datapath/data_path.go +++ b/pkg/datapath/data_path.go @@ -55,7 +55,7 @@ type BackupStartParam struct { Tags map[string]string } -type genearalDataPath struct { +type generalDataPath struct { ctx context.Context cancel context.CancelFunc backupRepo *velerov1api.BackupRepository @@ -73,7 +73,7 @@ type genearalDataPath struct { } func newGeneralDataPath(jobName string, requestorType string, client client.Client, namespace string, callbacks Callbacks, log logrus.FieldLogger) AsyncBR { - dp := &genearalDataPath{ + dp := &generalDataPath{ jobName: jobName, requestorType: requestorType, client: client, @@ -86,7 +86,7 @@ func newGeneralDataPath(jobName string, requestorType string, client client.Clie return dp } -func (dp *genearalDataPath) Init(ctx context.Context, param any) error { +func (dp *generalDataPath) Init(ctx context.Context, param any) error { initParam := param.(*InitParam) var err error @@ -138,7 +138,7 @@ func (dp *genearalDataPath) Init(ctx context.Context, param any) error { return nil } -func (dp *genearalDataPath) Close(ctx context.Context) { +func (dp *generalDataPath) Close(ctx context.Context) { if dp.cancel != nil { dp.cancel() } @@ -152,7 +152,7 @@ func (dp *genearalDataPath) Close(ctx context.Context) { dp.log.WithField("user", dp.jobName).Info("Data path is closed") } -func (dp *genearalDataPath) close(ctx context.Context) { +func (dp *generalDataPath) close(ctx context.Context) { dp.dataPathLock.Lock() defer dp.dataPathLock.Unlock() @@ -165,7 +165,7 @@ func (dp *genearalDataPath) close(ctx context.Context) { } } -func (dp *genearalDataPath) StartBackup(source AccessPoint, uploaderConfig map[string]string, param any) error { +func (dp *generalDataPath) StartBackup(source AccessPoint, uploaderConfig map[string]string, param any) error { if !dp.initialized { return errors.New("file system data path is not initialized") } @@ -201,7 +201,7 @@ func (dp *genearalDataPath) StartBackup(source AccessPoint, uploaderConfig map[s return nil } -func (dp *genearalDataPath) StartRestore(snapshotID string, target AccessPoint, uploaderConfigs map[string]string) error { +func (dp *generalDataPath) StartRestore(snapshotID string, target AccessPoint, uploaderConfigs map[string]string) error { if !dp.initialized { return errors.New("data path is not initialized") } @@ -235,18 +235,18 @@ func (dp *genearalDataPath) StartRestore(snapshotID string, target AccessPoint, } // UpdateProgress which implement ProgressUpdater interface to update progress status -func (dp *genearalDataPath) UpdateProgress(p *uploader.Progress) { +func (dp *generalDataPath) UpdateProgress(p *uploader.Progress) { if dp.callbacks.OnProgress != nil { dp.callbacks.OnProgress(context.Background(), dp.namespace, dp.jobName, &uploader.Progress{TotalBytes: p.TotalBytes, BytesDone: p.BytesDone}) } } -func (dp *genearalDataPath) Cancel() { +func (dp *generalDataPath) Cancel() { dp.cancel() dp.log.WithField("user", dp.jobName).Info("FileSystemBR is canceled") } -func (dp *genearalDataPath) boostRepoConnect(ctx context.Context, repositoryType string, credentialGetter *credentials.CredentialGetter, cacheDir string) error { +func (dp *generalDataPath) boostRepoConnect(ctx context.Context, repositoryType string, credentialGetter *credentials.CredentialGetter, cacheDir string) error { if repositoryType == velerov1api.BackupRepositoryTypeKopia { if err := repoProvider.NewUnifiedRepoProvider(*credentialGetter, repositoryType, dp.log).BoostRepoConnect(ctx, repoProvider.RepoParam{BackupLocation: dp.backupLocation, BackupRepo: dp.backupRepo, CacheDir: cacheDir}); err != nil { return err diff --git a/pkg/datapath/data_path_test.go b/pkg/datapath/data_path_test.go index ab1213a51..36f995501 100644 --- a/pkg/datapath/data_path_test.go +++ b/pkg/datapath/data_path_test.go @@ -94,7 +94,7 @@ func TestAsyncBackup(t *testing.T) { for _, test := range tests { t.Run(test.name, func(t *testing.T) { - dp := newGeneralDataPath("job-1", "test", nil, "velero", Callbacks{}, velerotest.NewLogger()).(*genearalDataPath) + dp := newGeneralDataPath("job-1", "test", nil, "velero", Callbacks{}, velerotest.NewLogger()).(*generalDataPath) mockProvider := providerMock.NewProvider(t) mockProvider.On("RunBackup", mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return(test.result.Backup.SnapshotID, test.result.Backup.EmptySnapshot, test.result.Backup.TotalBytes, test.result.Backup.IncrementalBytes, test.err) mockProvider.On("Close", mock.Anything).Return(nil) @@ -182,7 +182,7 @@ func TestAsyncRestore(t *testing.T) { for _, test := range tests { t.Run(test.name, func(t *testing.T) { - dp := newGeneralDataPath("job-1", "test", nil, "velero", Callbacks{}, velerotest.NewLogger()).(*genearalDataPath) + dp := newGeneralDataPath("job-1", "test", nil, "velero", Callbacks{}, velerotest.NewLogger()).(*generalDataPath) mockProvider := providerMock.NewProvider(t) mockProvider.On("RunRestore", mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return(test.result.Restore.TotalBytes, test.err) mockProvider.On("Close", mock.Anything).Return(nil)