From 6ec84c30f88055398840932eeb5ad45b5ca1cb1f Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Thu, 3 Sep 2026 15:20:18 +0800 Subject: [PATCH 1/9] enhance CBT retrievement to indicate the result in error message Signed-off-by: Lyndon-Li --- pkg/uploader/block/snapshot.go | 12 +++----- pkg/uploader/cbt/bitmap.go | 13 +++++++++ pkg/uploader/cbt/set.go | 52 +++++++++++++++++++++++++++------ pkg/uploader/cbt/types/types.go | 9 ++++++ 4 files changed, 69 insertions(+), 17 deletions(-) diff --git a/pkg/uploader/block/snapshot.go b/pkg/uploader/block/snapshot.go index 595a80fd5..a0823de8c 100644 --- a/pkg/uploader/block/snapshot.go +++ b/pkg/uploader/block/snapshot.go @@ -110,10 +110,10 @@ func snapshotSource( bitmap := cbt.NewBitmap(blockSize, uint64(source.size), cbtSource.Snapshot, parentBackup.changeID, parentBackup.volumeID) - err := cbt.SetBitmapOrFull(ctx, cbtService, bitmap) + err := cbt.SetBitmapOrFull(ctx, cbtService, bitmap, false) if err != nil { parentBackup.parentObject = "" - log.WithError(err).Warnf("Failed to create CBT with source %v, fallback to real full backup", cbtSource) + log.WithError(err).Warnf("Failed to create CBT with source %v", cbtSource) } snap, backupSize, err := u.Backup(source, parentBackup.parentObject, bitmap.Iterator(), uploaderCfg) @@ -218,16 +218,12 @@ func Restore(ctx context.Context, blkUp Uploader, rep udmrepo.BackupRepo, snapsh if incremental { if snapshot.Tags == nil { log.Warnf("No tag from snapshot %s, fallback to full restore", snapshotID) - incremental = false } else if snapshot.Tags[uploader.CBTChangeIDTag] == "" { log.Warnf("No ChangeID tag from snapshot %s, fallback to full restore", snapshotID) - incremental = false } else if snapshot.Tags[uploader.CBTVolumeIDTag] == "" { log.Warnf("No VolumeID tag from snapshot %s, fallback to full restore", snapshotID) - incremental = false } else if snapshot.Tags[uploader.CBTVolumeIDTag] != cbtSource.VolumeID { log.Warnf("VolumeID %s from snapshot %s is not expected as %s, fallback to full restore", snapshot.Tags[uploader.CBTVolumeIDTag], snapshotID, cbtSource.VolumeID) - incremental = false } else { volumeSnapshot = cbtSource.Snapshot changeID = snapshot.Tags[uploader.CBTChangeIDTag] @@ -237,8 +233,8 @@ func Restore(ctx context.Context, blkUp Uploader, rep udmrepo.BackupRepo, snapsh bitmap := cbt.NewBitmap(blockSize, uint64(snapshot.TotalSize), volumeSnapshot, changeID, volumeID) if incremental { - if err = cbt.SetBitmapOrFull(ctx, cbtService, bitmap); err != nil { - log.WithError(err).Warnf("Failed to create CBT with source %v, fallback to full restore", cbtSource) + if err = cbt.SetBitmapOrFull(ctx, cbtService, bitmap, true); err != nil { + log.WithError(err).Warnf("Failed to create CBT with source %v", cbtSource) } } else { bitmap.SetFull() diff --git a/pkg/uploader/cbt/bitmap.go b/pkg/uploader/cbt/bitmap.go index f26cb22b7..f1b4b3c7d 100644 --- a/pkg/uploader/cbt/bitmap.go +++ b/pkg/uploader/cbt/bitmap.go @@ -36,6 +36,7 @@ type bitmapImpl struct { snapshot string changeID string volumeID string + cbtError error } type bitmapIterator struct { @@ -89,6 +90,14 @@ func (c *bitmapImpl) VolumeID() string { return c.volumeID } +func (c *bitmapImpl) SetError(err error) { + c.cbtError = err +} + +func (c *bitmapImpl) Error() error { + return c.cbtError +} + func (c *bitmapImpl) Iterator() types.Iterator { if c.bitmap == nil { return nil @@ -115,3 +124,7 @@ func (c *bitmapIterator) Count() uint64 { func (c *bitmapIterator) BlockSize() uint { return c.blockSize } + +func (c *bitmapIterator) Error() error { + return c.cbtError +} diff --git a/pkg/uploader/cbt/set.go b/pkg/uploader/cbt/set.go index 11361cf77..be039d0c7 100644 --- a/pkg/uploader/cbt/set.go +++ b/pkg/uploader/cbt/set.go @@ -26,36 +26,70 @@ import ( ) // SetBitmapOrFull translates the allocated/changed blocks from CBT service to the given bitmap or set the bitmap to full when error happens -func SetBitmapOrFull(ctx context.Context, service cbtservice.Service, bitmap types.Bitmap) (err error) { +func SetBitmapOrFull(ctx context.Context, service cbtservice.Service, bitmap types.Bitmap, incOnly bool) (ret error) { + setFull := false + defer func() { - if err != nil { + bitmap.SetError(ret) + + if setFull { bitmap.SetFull() } }() if service == nil { - return errors.New("CBT service is absent") + setFull = true + return errors.New("CBT service is absent, fallback to real full") } if bitmap.Snapshot() == "" { - return errors.New("invalid snapshot") + setFull = true + return errors.New("invalid snapshot, fallback to real full") } - if bitmap.ChangeID() == "" { - return errors.Wrapf(service.GetAllocatedBlocks(ctx, bitmap.Snapshot(), func(blocks []cbtservice.Range) error { + if incOnly && bitmap.ChangeID() == "" { + setFull = true + return errors.New("invalid changeID, fallback to real full") + } + + var changedErr error + if bitmap.ChangeID() != "" { + err := service.GetChangedBlocks(ctx, bitmap.Snapshot(), bitmap.ChangeID(), func(blocks []cbtservice.Range) error { for _, b := range blocks { bitmap.Set(b.Offset, b.Length) } return nil - }), "error getting allocated blocks from CBT service") + }) + + if err == nil { + return nil + } + + if incOnly { + setFull = true + return errors.Wrap(err, "error getting changed blocks from CBT service, fallback to real full") + } + + changedErr = err } - return errors.Wrapf(service.GetChangedBlocks(ctx, bitmap.Snapshot(), bitmap.ChangeID(), func(blocks []cbtservice.Range) error { + err := service.GetAllocatedBlocks(ctx, bitmap.Snapshot(), func(blocks []cbtservice.Range) error { for _, b := range blocks { bitmap.Set(b.Offset, b.Length) } return nil - }), "error getting changed blocks from CBT service") + }) + + if err != nil { + setFull = true + return errors.Wrap(err, "error getting allocated blocks from CBT service, fallback to real full") + } + + if changedErr != nil { + return errors.Wrap(changedErr, "error getting changed blocks from CBT service, fallback to full") + } + + return nil } diff --git a/pkg/uploader/cbt/types/types.go b/pkg/uploader/cbt/types/types.go index 9b13669b2..447c3f9f3 100644 --- a/pkg/uploader/cbt/types/types.go +++ b/pkg/uploader/cbt/types/types.go @@ -35,6 +35,12 @@ type Bitmap interface { // Iterator returns the iterator for the CBT Bitmap Iterator() Iterator + + // SetError sets CBT error when preparing this bitmap + SetError(error) + + // Error returns the CBT error when preparing this bitmap + Error() error } // Iterator defines the methods to iterate the CBT bitmap and query the associated information @@ -56,4 +62,7 @@ type Iterator interface { // Next returns the offset of the next set block and whether it comes to the end of the iteration Next() (uint64, bool) + + // Error returns the CBT error when preparing this bitmap + Error() error } From 7e968e10d829419f24230c291953afde58b68b0f Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Fri, 4 Sep 2026 13:48:17 +0800 Subject: [PATCH 2/9] UT for new bitmap implementation Signed-off-by: Lyndon-Li --- pkg/uploader/cbt/set_test.go | 38 ++++++++++++++++++++---- pkg/uploader/cbt/types/mocks/Bitmap.go | 20 +++++++++++++ pkg/uploader/cbt/types/mocks/Iterator.go | 15 ++++++++++ 3 files changed, 67 insertions(+), 6 deletions(-) diff --git a/pkg/uploader/cbt/set_test.go b/pkg/uploader/cbt/set_test.go index 55701dedd..8b40c1143 100644 --- a/pkg/uploader/cbt/set_test.go +++ b/pkg/uploader/cbt/set_test.go @@ -33,6 +33,7 @@ func TestSetBitmapOrFull(t *testing.T) { tests := []struct { name string nilService bool + incOnly bool setupMocks func(*cbtservicemocks.Service, *cbtmocks.Bitmap) expectedErrStr string }{ @@ -42,7 +43,7 @@ func TestSetBitmapOrFull(t *testing.T) { setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { bmp.On("SetFull").Return() }, - expectedErrStr: "CBT service is absent", + expectedErrStr: "CBT service is absent, fallback to real full", }, { name: "invalid snapshot", @@ -50,7 +51,17 @@ func TestSetBitmapOrFull(t *testing.T) { bmp.On("Snapshot").Return("") bmp.On("SetFull").Return() }, - expectedErrStr: "invalid snapshot", + expectedErrStr: "invalid snapshot, fallback to real full", + }, + { + name: "invalid changeID", + incOnly: true, + setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { + bmp.On("Snapshot").Return("snap-1") + bmp.On("ChangeID").Return("") + bmp.On("SetFull").Return() + }, + expectedErrStr: "invalid changeID, fallback to real full", }, { name: "allocated blocks success", @@ -79,7 +90,7 @@ func TestSetBitmapOrFull(t *testing.T) { svc.On("GetAllocatedBlocks", mock.Anything, "snap-1", mock.Anything).Return(errors.New("mock alloc error")) bmp.On("SetFull").Return() }, - expectedErrStr: "error getting allocated blocks from CBT service: mock alloc error", + expectedErrStr: "error getting allocated blocks from CBT service, fallback to real full: mock alloc error", }, { name: "changed blocks success", @@ -98,7 +109,8 @@ func TestSetBitmapOrFull(t *testing.T) { }, }, { - name: "changed blocks error", + name: "changed blocks error with incOnly", + incOnly: true, setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { bmp.On("Snapshot").Return("snap-1") bmp.On("ChangeID").Return("change-1") @@ -106,7 +118,19 @@ func TestSetBitmapOrFull(t *testing.T) { svc.On("GetChangedBlocks", mock.Anything, "snap-1", "change-1", mock.Anything).Return(errors.New("mock changed error")) bmp.On("SetFull").Return() }, - expectedErrStr: "error getting changed blocks from CBT service: mock changed error", + expectedErrStr: "error getting changed blocks from CBT service, fallback to real full: mock changed error", + }, + { + name: "changed blocks error fallback to full", + setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { + bmp.On("Snapshot").Return("snap-1") + bmp.On("ChangeID").Return("change-1") + + svc.On("GetChangedBlocks", mock.Anything, "snap-1", "change-1", mock.Anything).Return(errors.New("mock changed error")) + + svc.On("GetAllocatedBlocks", mock.Anything, "snap-1", mock.Anything).Return(nil) + }, + expectedErrStr: "error getting changed blocks from CBT service, fallback to full: mock changed error", }, } @@ -124,7 +148,9 @@ func TestSetBitmapOrFull(t *testing.T) { svc = svcMock } - err := SetBitmapOrFull(context.Background(), svc, bmpMock) + bmpMock.On("SetError", mock.Anything).Return() + + err := SetBitmapOrFull(context.Background(), svc, bmpMock, tt.incOnly) if tt.expectedErrStr != "" { require.Error(t, err) diff --git a/pkg/uploader/cbt/types/mocks/Bitmap.go b/pkg/uploader/cbt/types/mocks/Bitmap.go index faa4ee242..514e5c28f 100644 --- a/pkg/uploader/cbt/types/mocks/Bitmap.go +++ b/pkg/uploader/cbt/types/mocks/Bitmap.go @@ -292,3 +292,23 @@ func (_c *Bitmap_VolumeID_Call) RunAndReturn(run func() string) *Bitmap_VolumeID _c.Call.Return(run) return _c } +// SetError provides a mock function for the type Bitmap +func (_mock *Bitmap) SetError(err error) { + _mock.Called(err) +} + +// Error provides a mock function for the type Bitmap +func (_mock *Bitmap) Error() error { + ret := _mock.Called() + if len(ret) == 0 { + panic("no return value specified for Error") + } + + var r0 error + if returnFunc, ok := ret.Get(0).(func() error); ok { + r0 = returnFunc() + } else { + r0 = ret.Error(0) + } + return r0 +} diff --git a/pkg/uploader/cbt/types/mocks/Iterator.go b/pkg/uploader/cbt/types/mocks/Iterator.go index eeafd0a59..00fc163d9 100644 --- a/pkg/uploader/cbt/types/mocks/Iterator.go +++ b/pkg/uploader/cbt/types/mocks/Iterator.go @@ -307,3 +307,18 @@ func (_c *Iterator_VolumeID_Call) RunAndReturn(run func() string) *Iterator_Volu _c.Call.Return(run) return _c } +// Error provides a mock function for the type Iterator +func (_mock *Iterator) Error() error { + ret := _mock.Called() + if len(ret) == 0 { + panic("no return value specified for Error") + } + + var r0 error + if returnFunc, ok := ret.Get(0).(func() error); ok { + r0 = returnFunc() + } else { + r0 = ret.Error(0) + } + return r0 +} From 67923bca6f9f97967448fba7daadad83c1552462 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Fri, 4 Sep 2026 14:17:28 +0800 Subject: [PATCH 3/9] uploader report incremental fallback Signed-off-by: Lyndon-Li --- pkg/datapath/data_path.go | 6 +++++- pkg/uploader/block/uploader.go | 8 ++++++++ pkg/uploader/kopia/snapshot.go | 25 +++++++++++++++---------- pkg/uploader/provider/kopia.go | 2 +- pkg/uploader/types.go | 5 +++-- 5 files changed, 32 insertions(+), 14 deletions(-) diff --git a/pkg/datapath/data_path.go b/pkg/datapath/data_path.go index 0513619be..2a4f4e18b 100644 --- a/pkg/datapath/data_path.go +++ b/pkg/datapath/data_path.go @@ -278,7 +278,11 @@ func (dp *generalDataPath) StartRestore(snapshotID string, target AccessPoint, u // UpdateProgress which implement ProgressUpdater interface to update progress status 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}) + dp.callbacks.OnProgress(context.Background(), dp.namespace, dp.jobName, &uploader.Progress{ + TotalBytes: p.TotalBytes, + BytesDone: p.BytesDone, + Message: p.Message, + }) } } diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index e2d464872..4b69de59f 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -86,6 +86,10 @@ func (blkup *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, b return udmrepo.Snapshot{}, 0, errors.New("bitmap is not available") } + if bitmap.Error() != nil { + blkup.progress.UpdateProgress(&uploader.Progress{BytesDone: -1, TotalBytes: -1, Message: bitmap.Error().Error()}) + } + backupMode := udmrepo.ObjectDataBackupModeInc if parentObject == "" { backupMode = udmrepo.ObjectDataBackupModeFull @@ -153,6 +157,10 @@ func (blkup *blockUploader) Restore(snapshot udmrepo.Snapshot, dest destInfo, bi return 0, 0, errors.New("bitmap is not available") } + if bitmap.Error() != nil { + blkup.progress.UpdateProgress(&uploader.Progress{BytesDone: -1, TotalBytes: -1, Message: bitmap.Error().Error()}) + } + meta, err := blkup.repoWriter.ReadMetadata(blkup.ctx, snapshot.RootObject.ID) if err != nil { return 0, 0, errors.Wrapf(err, "error reading snapshot metadata for %s", snapshot.Description) diff --git a/pkg/uploader/kopia/snapshot.go b/pkg/uploader/kopia/snapshot.go index fae7a517c..682b0f314 100644 --- a/pkg/uploader/kopia/snapshot.go +++ b/pkg/uploader/kopia/snapshot.go @@ -153,7 +153,8 @@ func setupPolicy(ctx context.Context, rep repo.RepositoryWriter, sourceInfo snap // Backup backup specific sourcePath and update progress func Backup(ctx context.Context, fsUploader SnapshotUploader, repoWriter repo.RepositoryWriter, sourcePath string, realSource string, - forceFull bool, parentSnapshot string, volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, tags map[string]string, log logrus.FieldLogger) (*uploader.SnapshotInfo, bool, error) { + forceFull bool, parentSnapshot string, volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, tags map[string]string, + updater uploader.ProgressUpdater, log logrus.FieldLogger) (*uploader.SnapshotInfo, bool, error) { if fsUploader == nil { return nil, false, errors.New("get empty kopia uploader") } @@ -189,7 +190,7 @@ func Backup(ctx context.Context, fsUploader SnapshotUploader, repoWriter repo.Re kopiaCtx := kopia.SetupKopiaLog(ctx, log) - snapID, snapshotSize, err := SnapshotSource(kopiaCtx, repoWriter, fsUploader, sourceInfo, sourceEntry, forceFull, parentSnapshot, tags, uploaderCfg, log, "Kopia Uploader") + snapID, snapshotSize, err := SnapshotSource(kopiaCtx, repoWriter, fsUploader, sourceInfo, sourceEntry, forceFull, parentSnapshot, tags, uploaderCfg, updater, log, "Kopia Uploader") snapshotInfo := &uploader.SnapshotInfo{ ID: snapID, Size: snapshotSize, @@ -237,6 +238,7 @@ func SnapshotSource( parentSnapshot string, snapshotTags map[string]string, uploaderCfg map[string]string, + updater uploader.ProgressUpdater, log logrus.FieldLogger, description string, ) (string, int64, error) { @@ -248,21 +250,24 @@ func SnapshotSource( if parentSnapshot != "" { log.Infof("Using provided parent snapshot %s", parentSnapshot) - mani, err := loadSnapshotFunc(ctx, rep, manifest.ID(parentSnapshot)) - if err != nil { - log.WithError(err).Warnf("Failed to load previous snapshot %v from kopia, fallback to full backup", parentSnapshot) + if mani, err := loadSnapshotFunc(ctx, rep, manifest.ID(parentSnapshot)); err != nil { + msg := fmt.Sprintf("Failed to load previous snapshot %v from kopia, fallback to full backup", parentSnapshot) + log.WithError(err).Warn(msg) + updater.UpdateProgress(&uploader.Progress{BytesDone: -1, TotalBytes: -1, Message: msg}) } else { previous = append(previous, mani) } } else { log.Infof("Searching for parent snapshot") - pre, err := findPreviousSnapshotManifest(ctx, rep, sourceInfo, snapshotTags, nil, log) - if err != nil { - return "", 0, errors.Wrapf(err, "Failed to find previous kopia snapshot manifests for si %v", sourceInfo) - } + if pre, err := findPreviousSnapshotManifest(ctx, rep, sourceInfo, snapshotTags, nil, log); err != nil { + log.WithError(err).Warnf("Failed to find previous kopia snapshot manifests for si %v, fallback to full backup", sourceInfo) - previous = pre + msg := fmt.Sprint("Failed to find previous kopia snapshot manifests, fallback to full backup") + updater.UpdateProgress(&uploader.Progress{BytesDone: -1, TotalBytes: -1, Message: msg}) + } else { + previous = pre + } } } else { log.Info("Forcing full snapshot") diff --git a/pkg/uploader/provider/kopia.go b/pkg/uploader/provider/kopia.go index 1cd1ecd45..39abbbe8f 100644 --- a/pkg/uploader/provider/kopia.go +++ b/pkg/uploader/provider/kopia.go @@ -166,7 +166,7 @@ func (kp *kopiaProvider) RunBackup( uploaderCfg[kopia.UploaderConfigMultipartKey] = "true" } - snapshotInfo, _, err := kopiaBackupFunc(ctx, kpUploader, repoWriter, path, realSource, forceFull, parentSnapshot, volMode, uploaderCfg, tags, log) + snapshotInfo, _, err := kopiaBackupFunc(ctx, kpUploader, repoWriter, path, realSource, forceFull, parentSnapshot, volMode, uploaderCfg, tags, updater, log) if err != nil { snapshotID := "" if snapshotInfo != nil { diff --git a/pkg/uploader/types.go b/pkg/uploader/types.go index 9c700193f..a608add06 100644 --- a/pkg/uploader/types.go +++ b/pkg/uploader/types.go @@ -58,8 +58,9 @@ type SnapshotInfo struct { // Progress which defined two variables to record progress type Progress struct { - TotalBytes int64 `json:"totalBytes,omitempty"` - BytesDone int64 `json:"doneBytes,omitempty"` + TotalBytes int64 `json:"totalBytes,omitempty"` + BytesDone int64 `json:"doneBytes,omitempty"` + Message string `json:"message,omitempty"` } // UploaderProgress which defined generic interface to update progress From 1048f26c208ad696f7d7e9019642b0181a653fb3 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Fri, 4 Sep 2026 14:28:31 +0800 Subject: [PATCH 4/9] controllers support message in progress Signed-off-by: Lyndon-Li --- pkg/controller/data_download_controller.go | 15 +++++++++++++-- pkg/controller/data_upload_controller.go | 15 +++++++++++++-- pkg/controller/pod_volume_backup_controller.go | 15 +++++++++++++-- pkg/controller/pod_volume_restore_controller.go | 15 +++++++++++++-- 4 files changed, 52 insertions(+), 8 deletions(-) diff --git a/pkg/controller/data_download_controller.go b/pkg/controller/data_download_controller.go index e2a62830e..673b82c02 100644 --- a/pkg/controller/data_download_controller.go +++ b/pkg/controller/data_download_controller.go @@ -42,7 +42,6 @@ import ( "sigs.k8s.io/controller-runtime/pkg/predicate" "sigs.k8s.io/controller-runtime/pkg/reconcile" - "github.com/vmware-tanzu/velero/pkg/apis/velero/shared" velerov1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v1" velerov2alpha1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v2alpha1" "github.com/vmware-tanzu/velero/pkg/constant" @@ -609,7 +608,19 @@ func (r *DataDownloadReconciler) OnDataDownloadProgress(ctx context.Context, nam log := r.logger.WithField("datadownload", ddName) if err := UpdateDataDownloadWithRetry(ctx, r.client, types.NamespacedName{Namespace: namespace, Name: ddName}, log, func(dd *velerov2alpha1api.DataDownload) bool { - dd.Status.Progress = shared.DataMoveOperationProgress{TotalBytes: progress.TotalBytes, BytesDone: progress.BytesDone} + if progress.TotalBytes != -1 { + dd.Status.Progress.TotalBytes = progress.TotalBytes + } + + if progress.BytesDone != -1 { + dd.Status.Progress.BytesDone = progress.BytesDone + } + + if progress.Message != "" { + dd.Status.Message += progress.Message + dd.Status.Message += ";" + } + return true }); err != nil { log.WithError(err).Error("Failed to update progress") diff --git a/pkg/controller/data_upload_controller.go b/pkg/controller/data_upload_controller.go index ae50a7741..8c7795765 100644 --- a/pkg/controller/data_upload_controller.go +++ b/pkg/controller/data_upload_controller.go @@ -42,7 +42,6 @@ import ( "sigs.k8s.io/controller-runtime/pkg/predicate" "sigs.k8s.io/controller-runtime/pkg/reconcile" - "github.com/vmware-tanzu/velero/pkg/apis/velero/shared" velerov1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v1" velerov2alpha1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v2alpha1" "github.com/vmware-tanzu/velero/pkg/constant" @@ -634,7 +633,19 @@ func (r *DataUploadReconciler) OnDataUploadProgress(ctx context.Context, namespa log := r.logger.WithField("dataupload", duName) if err := UpdateDataUploadWithRetry(ctx, r.client, types.NamespacedName{Namespace: namespace, Name: duName}, log, func(du *velerov2alpha1api.DataUpload) bool { - du.Status.Progress = shared.DataMoveOperationProgress{TotalBytes: progress.TotalBytes, BytesDone: progress.BytesDone} + if progress.TotalBytes != -1 { + du.Status.Progress.TotalBytes = progress.TotalBytes + } + + if progress.BytesDone != -1 { + du.Status.Progress.BytesDone = progress.BytesDone + } + + if progress.Message != "" { + du.Status.Message += progress.Message + du.Status.Message += ";" + } + return true }); err != nil { log.WithError(err).Error("Failed to update progress") diff --git a/pkg/controller/pod_volume_backup_controller.go b/pkg/controller/pod_volume_backup_controller.go index c4e68ce33..13d0dd851 100644 --- a/pkg/controller/pod_volume_backup_controller.go +++ b/pkg/controller/pod_volume_backup_controller.go @@ -41,7 +41,6 @@ import ( "sigs.k8s.io/controller-runtime/pkg/predicate" "sigs.k8s.io/controller-runtime/pkg/reconcile" - veleroapishared "github.com/vmware-tanzu/velero/pkg/apis/velero/shared" velerov1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v1" "github.com/vmware-tanzu/velero/pkg/constant" "github.com/vmware-tanzu/velero/pkg/datapath" @@ -628,7 +627,19 @@ func (r *PodVolumeBackupReconciler) OnDataPathProgress(ctx context.Context, name log := r.logger.WithField("pvb", pvbName) if err := UpdatePVBWithRetry(ctx, r.client, types.NamespacedName{Namespace: namespace, Name: pvbName}, log, func(pvb *velerov1api.PodVolumeBackup) bool { - pvb.Status.Progress = veleroapishared.DataMoveOperationProgress{TotalBytes: progress.TotalBytes, BytesDone: progress.BytesDone} + if progress.TotalBytes != -1 { + pvb.Status.Progress.TotalBytes = progress.TotalBytes + } + + if progress.BytesDone != -1 { + pvb.Status.Progress.BytesDone = progress.BytesDone + } + + if progress.Message != "" { + pvb.Status.Message += progress.Message + pvb.Status.Message += ";" + } + return true }); err != nil { log.WithError(err).Error("Failed to update progress") diff --git a/pkg/controller/pod_volume_restore_controller.go b/pkg/controller/pod_volume_restore_controller.go index 4fe9bdaa3..1ca274017 100644 --- a/pkg/controller/pod_volume_restore_controller.go +++ b/pkg/controller/pod_volume_restore_controller.go @@ -44,7 +44,6 @@ import ( "sigs.k8s.io/controller-runtime/pkg/predicate" "sigs.k8s.io/controller-runtime/pkg/reconcile" - veleroapishared "github.com/vmware-tanzu/velero/pkg/apis/velero/shared" velerov1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v1" "github.com/vmware-tanzu/velero/pkg/constant" "github.com/vmware-tanzu/velero/pkg/datapath" @@ -905,7 +904,19 @@ func (r *PodVolumeRestoreReconciler) OnDataPathProgress(ctx context.Context, nam log := r.logger.WithField("PVR", pvrName) if err := UpdatePVRWithRetry(ctx, r.client, types.NamespacedName{Namespace: namespace, Name: pvrName}, log, func(pvr *velerov1api.PodVolumeRestore) bool { - pvr.Status.Progress = veleroapishared.DataMoveOperationProgress{TotalBytes: progress.TotalBytes, BytesDone: progress.BytesDone} + if progress.TotalBytes != -1 { + pvr.Status.Progress.TotalBytes = progress.TotalBytes + } + + if progress.BytesDone != -1 { + pvr.Status.Progress.BytesDone = progress.BytesDone + } + + if progress.Message != "" { + pvr.Status.Message += progress.Message + pvr.Status.Message += ";" + } + return true }); err != nil { log.WithError(err).Error("Failed to update progress") From 5146992b5b20e4877b0b2de816b6183ecf2df260 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Fri, 4 Sep 2026 15:55:08 +0800 Subject: [PATCH 5/9] add UT for progress message Signed-off-by: Lyndon-Li --- .../data_download_controller_test.go | 33 ++++++++++++++----- pkg/controller/data_upload_controller_test.go | 29 ++++++++++++---- .../pod_volume_backup_controller_test.go | 29 ++++++++++++---- .../pod_volume_restore_controller_test.go | 29 ++++++++++++---- pkg/uploader/block/uploader_test.go | 8 +++++ pkg/uploader/kopia/snapshot_test.go | 6 ++-- pkg/uploader/provider/kopia_test.go | 8 ++--- 7 files changed, 109 insertions(+), 33 deletions(-) diff --git a/pkg/controller/data_download_controller_test.go b/pkg/controller/data_download_controller_test.go index 72d51167b..c9026eafa 100644 --- a/pkg/controller/data_download_controller_test.go +++ b/pkg/controller/data_download_controller_test.go @@ -787,6 +787,15 @@ func TestOnDataDownloadProgress(t *testing.T) { BytesDone: bytesDone, }, }, + { + name: "patch in progress phase with negative progress values and message", + dd: dataDownloadBuilder().Result(), + progress: uploader.Progress{ + TotalBytes: -1, + BytesDone: -1, + Message: "some warning message", + }, + }, { name: "failed to get datadownload", dd: dataDownloadBuilder().Result(), @@ -815,20 +824,28 @@ func TestOnDataDownloadProgress(t *testing.T) { require.NoError(t, r.client.Create(t.Context(), dd)) // Create a Progress object - progress := &uploader.Progress{ - TotalBytes: totalBytes, - BytesDone: bytesDone, - } + progress := &test.progress // Call the OnDataDownloadProgress function r.OnDataDownloadProgress(ctx, namespace, duName, progress) if len(test.needErrs) != 0 && !test.needErrs[0] { // Get the updated DataDownload object from the fake client - updatedDu := &velerov2alpha1api.DataDownload{} - require.NoError(t, r.client.Get(ctx, types.NamespacedName{Name: duName, Namespace: namespace}, updatedDu)) + updatedDd := &velerov2alpha1api.DataDownload{} + require.NoError(t, r.client.Get(ctx, types.NamespacedName{Name: duName, Namespace: namespace}, updatedDd)) // Assert that the DataDownload object has been updated with the progress - assert.Equal(t, test.progress.TotalBytes, updatedDu.Status.Progress.TotalBytes) - assert.Equal(t, test.progress.BytesDone, updatedDu.Status.Progress.BytesDone) + if progress.TotalBytes != -1 { + assert.Equal(t, test.progress.TotalBytes, updatedDd.Status.Progress.TotalBytes) + } else { + assert.Equal(t, int64(0), updatedDd.Status.Progress.TotalBytes) // assuming default or original value + } + if progress.BytesDone != -1 { + assert.Equal(t, test.progress.BytesDone, updatedDd.Status.Progress.BytesDone) + } else { + assert.Equal(t, int64(0), updatedDd.Status.Progress.BytesDone) // assuming default or original value + } + if progress.Message != "" { + assert.Contains(t, updatedDd.Status.Message, progress.Message) + } } }) } diff --git a/pkg/controller/data_upload_controller_test.go b/pkg/controller/data_upload_controller_test.go index 30b5926ac..d50ceb749 100644 --- a/pkg/controller/data_upload_controller_test.go +++ b/pkg/controller/data_upload_controller_test.go @@ -809,6 +809,15 @@ func TestOnDataUploadProgress(t *testing.T) { BytesDone: bytesDone, }, }, + { + name: "patch in progress phase with negative progress values and message", + du: dataUploadBuilder().Result(), + progress: uploader.Progress{ + TotalBytes: -1, + BytesDone: -1, + Message: "some warning message", + }, + }, { name: "failed to get dataupload", du: dataUploadBuilder().Result(), @@ -837,10 +846,7 @@ func TestOnDataUploadProgress(t *testing.T) { require.NoError(t, r.client.Create(t.Context(), du)) // Create a Progress object - progress := &uploader.Progress{ - TotalBytes: totalBytes, - BytesDone: bytesDone, - } + progress := &test.progress // Call the OnDataUploadProgress function r.OnDataUploadProgress(ctx, namespace, duName, progress) @@ -849,8 +855,19 @@ func TestOnDataUploadProgress(t *testing.T) { updatedDu := &velerov2alpha1api.DataUpload{} require.NoError(t, r.client.Get(ctx, types.NamespacedName{Name: duName, Namespace: namespace}, updatedDu)) // Assert that the DataUpload object has been updated with the progress - assert.Equal(t, test.progress.TotalBytes, updatedDu.Status.Progress.TotalBytes) - assert.Equal(t, test.progress.BytesDone, updatedDu.Status.Progress.BytesDone) + if progress.TotalBytes != -1 { + assert.Equal(t, test.progress.TotalBytes, updatedDu.Status.Progress.TotalBytes) + } else { + assert.Equal(t, int64(0), updatedDu.Status.Progress.TotalBytes) // assuming default or original value + } + if progress.BytesDone != -1 { + assert.Equal(t, test.progress.BytesDone, updatedDu.Status.Progress.BytesDone) + } else { + assert.Equal(t, int64(0), updatedDu.Status.Progress.BytesDone) // assuming default or original value + } + if progress.Message != "" { + assert.Contains(t, updatedDu.Status.Message, progress.Message) + } } }) } diff --git a/pkg/controller/pod_volume_backup_controller_test.go b/pkg/controller/pod_volume_backup_controller_test.go index 21e30d5db..e74a1c269 100644 --- a/pkg/controller/pod_volume_backup_controller_test.go +++ b/pkg/controller/pod_volume_backup_controller_test.go @@ -625,6 +625,15 @@ func TestOnPVBProgress(t *testing.T) { BytesDone: bytesDone, }, }, + { + name: "patch in progress phase with negative progress values and message", + pvb: pvbBuilder().Result(), + progress: uploader.Progress{ + TotalBytes: -1, + BytesDone: -1, + Message: "some warning message", + }, + }, { name: "failed to get pvb", pvb: pvbBuilder().Result(), @@ -653,17 +662,25 @@ func TestOnPVBProgress(t *testing.T) { require.NoError(t, r.client.Create(t.Context(), pvb)) // Create a Progress object - progress := &uploader.Progress{ - TotalBytes: totalBytes, - BytesDone: bytesDone, - } + progress := &test.progress r.OnDataPathProgress(ctx, namespace, pvbName, progress) if len(test.needErrs) != 0 && !test.needErrs[0] { updatedPvb := &velerov1api.PodVolumeBackup{} require.NoError(t, r.client.Get(ctx, types.NamespacedName{Name: pvbName, Namespace: namespace}, updatedPvb)) - assert.Equal(t, test.progress.TotalBytes, updatedPvb.Status.Progress.TotalBytes) - assert.Equal(t, test.progress.BytesDone, updatedPvb.Status.Progress.BytesDone) + if progress.TotalBytes != -1 { + assert.Equal(t, test.progress.TotalBytes, updatedPvb.Status.Progress.TotalBytes) + } else { + assert.Equal(t, int64(0), updatedPvb.Status.Progress.TotalBytes) // assuming default or original value + } + if progress.BytesDone != -1 { + assert.Equal(t, test.progress.BytesDone, updatedPvb.Status.Progress.BytesDone) + } else { + assert.Equal(t, int64(0), updatedPvb.Status.Progress.BytesDone) // assuming default or original value + } + if progress.Message != "" { + assert.Contains(t, updatedPvb.Status.Message, progress.Message) + } } }) } diff --git a/pkg/controller/pod_volume_restore_controller_test.go b/pkg/controller/pod_volume_restore_controller_test.go index 73167c76f..6bd68e7e1 100644 --- a/pkg/controller/pod_volume_restore_controller_test.go +++ b/pkg/controller/pod_volume_restore_controller_test.go @@ -1470,6 +1470,15 @@ func TestOnPodVolumeRestoreProgress(t *testing.T) { BytesDone: bytesDone, }, }, + { + name: "patch in progress phase with negative progress values and message", + pvr: pvrBuilder().Result(), + progress: uploader.Progress{ + TotalBytes: -1, + BytesDone: -1, + Message: "some warning message", + }, + }, { name: "failed to get pvr", pvr: pvrBuilder().Result(), @@ -1498,17 +1507,25 @@ func TestOnPodVolumeRestoreProgress(t *testing.T) { require.NoError(t, r.client.Create(t.Context(), pvr)) // Create a Progress object - progress := &uploader.Progress{ - TotalBytes: totalBytes, - BytesDone: bytesDone, - } + progress := &test.progress r.OnDataPathProgress(ctx, namespace, pvrName, progress) if len(test.needErrs) != 0 && !test.needErrs[0] { updatedPVR := &velerov1api.PodVolumeRestore{} require.NoError(t, r.client.Get(ctx, types.NamespacedName{Name: pvrName, Namespace: namespace}, updatedPVR)) - assert.Equal(t, test.progress.TotalBytes, updatedPVR.Status.Progress.TotalBytes) - assert.Equal(t, test.progress.BytesDone, updatedPVR.Status.Progress.BytesDone) + if progress.TotalBytes != -1 { + assert.Equal(t, test.progress.TotalBytes, updatedPVR.Status.Progress.TotalBytes) + } else { + assert.Equal(t, int64(0), updatedPVR.Status.Progress.TotalBytes) // assuming default or original value + } + if progress.BytesDone != -1 { + assert.Equal(t, test.progress.BytesDone, updatedPVR.Status.Progress.BytesDone) + } else { + assert.Equal(t, int64(0), updatedPVR.Status.Progress.BytesDone) // assuming default or original value + } + if progress.Message != "" { + assert.Contains(t, updatedPVR.Status.Message, progress.Message) + } } }) } diff --git a/pkg/uploader/block/uploader_test.go b/pkg/uploader/block/uploader_test.go index bb94eb3de..f011219e4 100644 --- a/pkg/uploader/block/uploader_test.go +++ b/pkg/uploader/block/uploader_test.go @@ -295,6 +295,8 @@ func TestBlockUploaderBackup(t *testing.T) { var iterator cbt.Iterator if !tc.nilBitmap { iterMock := cbtmocks.NewIterator(t) + iterMock.On("Error").Return(nil).Maybe() + iterMock.On("Error").Return(nil).Maybe() iterator = iterMock backupMode := udmrepo.ObjectDataBackupModeInc @@ -573,6 +575,7 @@ func TestRestoreData(t *testing.T) { reader := bytes.NewReader(data) iterMock := cbtmocks.NewIterator(t) + iterMock.On("Error").Return(nil).Maybe() iterMock.On("Count").Return(uint64(1)) iterMock.On("Next").Return(uint64(0), true).Once() iterMock.On("Next").Return(uint64(0), false) @@ -603,6 +606,7 @@ func TestRestoreData(t *testing.T) { reader := &errReader{err: errors.New("read error")} iterMock := cbtmocks.NewIterator(t) + iterMock.On("Error").Return(nil).Maybe() iterMock.On("Count").Return(uint64(1)) iterMock.On("Next").Return(uint64(0), true).Once() iterMock.On("Next").Return(uint64(0), false) @@ -623,6 +627,7 @@ func TestBlockUploaderRestore(t *testing.T) { repoWriter.On("ReadMetadata", mock.Anything, udmrepo.ID("root-id")).Return(nil, errors.New("meta not found")) iterMock := cbtmocks.NewIterator(t) + iterMock.On("Error").Return(nil).Maybe() _, _, err := blkup.Restore(udmrepo.Snapshot{RootObject: udmrepo.ObjectMetadata{ID: "root-id"}}, destInfo{}, iterMock, nil) require.Error(t, err) assert.Contains(t, err.Error(), "meta not found") @@ -680,6 +685,7 @@ func TestBlockUploaderRestore(t *testing.T) { } iterMock := cbtmocks.NewIterator(t) + iterMock.On("Error").Return(nil).Maybe() iterMock.On("Count").Return(uint64(1)) iterMock.On("Next").Return(uint64(0), true).Once() iterMock.On("Next").Return(uint64(0), false) @@ -708,6 +714,7 @@ func TestBlockUploaderRestore(t *testing.T) { } dest := destInfo{size: 4194304, path: "/dev/target"} iterMock := cbtmocks.NewIterator(t) + iterMock.On("Error").Return(nil).Maybe() _, _, err := blkup.Restore(snap, dest, iterMock, nil) require.Error(t, err) @@ -732,6 +739,7 @@ func TestBlockUploaderRestore(t *testing.T) { } dest := destInfo{size: 512, path: "/dev/small"} iterMock := cbtmocks.NewIterator(t) + iterMock.On("Error").Return(nil).Maybe() _, _, err := blkup.Restore(snap, dest, iterMock, nil) require.Error(t, err) diff --git a/pkg/uploader/kopia/snapshot_test.go b/pkg/uploader/kopia/snapshot_test.go index e58c2bb88..08c34befd 100644 --- a/pkg/uploader/kopia/snapshot_test.go +++ b/pkg/uploader/kopia/snapshot_test.go @@ -200,7 +200,7 @@ func TestSnapshotSource(t *testing.T) { t.Run(tc.name, func(t *testing.T) { s := injectSnapshotFuncs() MockFuncs(s, tc.args) - _, _, err = SnapshotSource(ctx, s.repoWriterMock, s.uploderMock, sourceInfo, rootDir, false, "/", nil, tc.uploaderCfg, log, "TestSnapshotSource") + _, _, err = SnapshotSource(ctx, s.repoWriterMock, s.uploderMock, sourceInfo, rootDir, false, "/", nil, tc.uploaderCfg, &fakeProgressUpdater{}, log, "TestSnapshotSource") if tc.notError { assert.NoError(t, err) } else { @@ -648,9 +648,9 @@ func TestBackup(t *testing.T) { var snapshotInfo *uploader.SnapshotInfo var err error if tc.isEmptyUploader { - snapshotInfo, isSnapshotEmpty, err = Backup(t.Context(), nil, s.repoWriterMock, tc.sourcePath, "", tc.forceFull, tc.parentSnapshot, tc.volMode, map[string]string{}, tc.tags, &logrus.Logger{}) + snapshotInfo, isSnapshotEmpty, err = Backup(t.Context(), nil, s.repoWriterMock, tc.sourcePath, "", tc.forceFull, tc.parentSnapshot, tc.volMode, map[string]string{}, tc.tags, &fakeProgressUpdater{}, &logrus.Logger{}) } else { - snapshotInfo, isSnapshotEmpty, err = Backup(t.Context(), s.uploderMock, s.repoWriterMock, tc.sourcePath, "", tc.forceFull, tc.parentSnapshot, tc.volMode, map[string]string{}, tc.tags, &logrus.Logger{}) + snapshotInfo, isSnapshotEmpty, err = Backup(t.Context(), s.uploderMock, s.repoWriterMock, tc.sourcePath, "", tc.forceFull, tc.parentSnapshot, tc.volMode, map[string]string{}, tc.tags, &fakeProgressUpdater{}, &logrus.Logger{}) } // Check if the returned error matches the expected error if tc.expectedError != nil { diff --git a/pkg/uploader/provider/kopia_test.go b/pkg/uploader/provider/kopia_test.go index ca1cf8f5a..fd5629151 100644 --- a/pkg/uploader/provider/kopia_test.go +++ b/pkg/uploader/provider/kopia_test.go @@ -65,27 +65,27 @@ func (f *FakeRestoreProgressUpdater) UpdateProgress(p *uploader.Progress) {} func TestRunBackup(t *testing.T) { testCases := []struct { name string - hookBackupFunc func(ctx context.Context, fsUploader kopia.SnapshotUploader, repoWriter repo.RepositoryWriter, sourcePath string, realSource string, forceFull bool, parentSnapshot string, volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, tags map[string]string, log logrus.FieldLogger) (*uploader.SnapshotInfo, bool, error) + hookBackupFunc func(ctx context.Context, fsUploader kopia.SnapshotUploader, repoWriter repo.RepositoryWriter, sourcePath string, realSource string, forceFull bool, parentSnapshot string, volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, tags map[string]string, updater uploader.ProgressUpdater, log logrus.FieldLogger) (*uploader.SnapshotInfo, bool, error) volMode uploader.PersistentVolumeMode notError bool }{ { name: "success to backup", - hookBackupFunc: func(ctx context.Context, fsUploader kopia.SnapshotUploader, repoWriter repo.RepositoryWriter, sourcePath string, realSource string, forceFull bool, parentSnapshot string, volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, tags map[string]string, log logrus.FieldLogger) (*uploader.SnapshotInfo, bool, error) { + hookBackupFunc: func(ctx context.Context, fsUploader kopia.SnapshotUploader, repoWriter repo.RepositoryWriter, sourcePath string, realSource string, forceFull bool, parentSnapshot string, volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, tags map[string]string, updater uploader.ProgressUpdater, log logrus.FieldLogger) (*uploader.SnapshotInfo, bool, error) { return &uploader.SnapshotInfo{}, false, nil }, notError: true, }, { name: "get error to backup", - hookBackupFunc: func(ctx context.Context, fsUploader kopia.SnapshotUploader, repoWriter repo.RepositoryWriter, sourcePath string, realSource string, forceFull bool, parentSnapshot string, volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, tags map[string]string, log logrus.FieldLogger) (*uploader.SnapshotInfo, bool, error) { + hookBackupFunc: func(ctx context.Context, fsUploader kopia.SnapshotUploader, repoWriter repo.RepositoryWriter, sourcePath string, realSource string, forceFull bool, parentSnapshot string, volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, tags map[string]string, updater uploader.ProgressUpdater, log logrus.FieldLogger) (*uploader.SnapshotInfo, bool, error) { return &uploader.SnapshotInfo{}, false, errors.New("failed to backup") }, notError: false, }, { name: "success to backup block mode volume", - hookBackupFunc: func(ctx context.Context, fsUploader kopia.SnapshotUploader, repoWriter repo.RepositoryWriter, sourcePath string, realSource string, forceFull bool, parentSnapshot string, volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, tags map[string]string, log logrus.FieldLogger) (*uploader.SnapshotInfo, bool, error) { + hookBackupFunc: func(ctx context.Context, fsUploader kopia.SnapshotUploader, repoWriter repo.RepositoryWriter, sourcePath string, realSource string, forceFull bool, parentSnapshot string, volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, tags map[string]string, updater uploader.ProgressUpdater, log logrus.FieldLogger) (*uploader.SnapshotInfo, bool, error) { return &uploader.SnapshotInfo{}, false, nil }, volMode: uploader.PersistentVolumeBlock, From b93c24c58a8755cfc6538c5ac9384b8aaeb34032 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Fri, 4 Sep 2026 16:10:54 +0800 Subject: [PATCH 6/9] report incremental fallback message Signed-off-by: Lyndon-Li --- changelogs/unreleased/10479-Lyndon-Li | 1 + 1 file changed, 1 insertion(+) create mode 100644 changelogs/unreleased/10479-Lyndon-Li diff --git a/changelogs/unreleased/10479-Lyndon-Li b/changelogs/unreleased/10479-Lyndon-Li new file mode 100644 index 000000000..bdd241963 --- /dev/null +++ b/changelogs/unreleased/10479-Lyndon-Li @@ -0,0 +1 @@ +Enhance data mover progress update to include a message field so that critical activity messages such as incremental fallback could be saved to the message field of DU/DD/PVB/PVR \ No newline at end of file From df709d39d675df91a2334c041c4bd0087e2f5bbf Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Fri, 4 Sep 2026 16:35:39 +0800 Subject: [PATCH 7/9] report incremental fallback message Signed-off-by: Lyndon-Li --- pkg/controller/data_download_controller.go | 6 +- pkg/controller/data_upload_controller.go | 6 +- .../pod_volume_backup_controller.go | 6 +- .../pod_volume_restore_controller.go | 6 +- pkg/uploader/cbt/set.go | 7 +- pkg/uploader/cbt/set_test.go | 206 +++++++++++------- pkg/uploader/kopia/snapshot.go | 17 +- 7 files changed, 156 insertions(+), 98 deletions(-) diff --git a/pkg/controller/data_download_controller.go b/pkg/controller/data_download_controller.go index 673b82c02..053a083e7 100644 --- a/pkg/controller/data_download_controller.go +++ b/pkg/controller/data_download_controller.go @@ -617,8 +617,10 @@ func (r *DataDownloadReconciler) OnDataDownloadProgress(ctx context.Context, nam } if progress.Message != "" { - dd.Status.Message += progress.Message - dd.Status.Message += ";" + message := progress.Message + ";" + if !strings.HasSuffix(dd.Status.Message, message) { + dd.Status.Message += message + } } return true diff --git a/pkg/controller/data_upload_controller.go b/pkg/controller/data_upload_controller.go index 8c7795765..dbbf21013 100644 --- a/pkg/controller/data_upload_controller.go +++ b/pkg/controller/data_upload_controller.go @@ -642,8 +642,10 @@ func (r *DataUploadReconciler) OnDataUploadProgress(ctx context.Context, namespa } if progress.Message != "" { - du.Status.Message += progress.Message - du.Status.Message += ";" + message := progress.Message + ";" + if !strings.HasSuffix(du.Status.Message, message) { + du.Status.Message += message + } } return true diff --git a/pkg/controller/pod_volume_backup_controller.go b/pkg/controller/pod_volume_backup_controller.go index 13d0dd851..3e343b073 100644 --- a/pkg/controller/pod_volume_backup_controller.go +++ b/pkg/controller/pod_volume_backup_controller.go @@ -636,8 +636,10 @@ func (r *PodVolumeBackupReconciler) OnDataPathProgress(ctx context.Context, name } if progress.Message != "" { - pvb.Status.Message += progress.Message - pvb.Status.Message += ";" + message := progress.Message + ";" + if !strings.HasSuffix(pvb.Status.Message, message) { + pvb.Status.Message += message + } } return true diff --git a/pkg/controller/pod_volume_restore_controller.go b/pkg/controller/pod_volume_restore_controller.go index 1ca274017..d8ffb0e1a 100644 --- a/pkg/controller/pod_volume_restore_controller.go +++ b/pkg/controller/pod_volume_restore_controller.go @@ -913,8 +913,10 @@ func (r *PodVolumeRestoreReconciler) OnDataPathProgress(ctx context.Context, nam } if progress.Message != "" { - pvr.Status.Message += progress.Message - pvr.Status.Message += ";" + message := progress.Message + ";" + if !strings.HasSuffix(pvr.Status.Message, message) { + pvr.Status.Message += message + } } return true diff --git a/pkg/uploader/cbt/set.go b/pkg/uploader/cbt/set.go index be039d0c7..0e50cb048 100644 --- a/pkg/uploader/cbt/set.go +++ b/pkg/uploader/cbt/set.go @@ -84,7 +84,12 @@ func SetBitmapOrFull(ctx context.Context, service cbtservice.Service, bitmap typ if err != nil { setFull = true - return errors.Wrap(err, "error getting allocated blocks from CBT service, fallback to real full") + + if changedErr != nil { + return errors.Wrap(err, "error getting both changed and allocated blocks from CBT service, fallback to real full") + } else { + return errors.Wrap(err, "error getting allocated blocks from CBT service, fallback to real full") + } } if changedErr != nil { diff --git a/pkg/uploader/cbt/set_test.go b/pkg/uploader/cbt/set_test.go index 8b40c1143..6c1d4a579 100644 --- a/pkg/uploader/cbt/set_test.go +++ b/pkg/uploader/cbt/set_test.go @@ -21,126 +21,148 @@ import ( "errors" "testing" + "github.com/stretchr/testify/assert" "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" "github.com/vmware-tanzu/velero/pkg/cbtservice" cbtservicemocks "github.com/vmware-tanzu/velero/pkg/cbtservice/mocks" - cbtmocks "github.com/vmware-tanzu/velero/pkg/uploader/cbt/types/mocks" ) func TestSetBitmapOrFull(t *testing.T) { + const mb = 1024 * 1024 tests := []struct { name string nilService bool incOnly bool - setupMocks func(*cbtservicemocks.Service, *cbtmocks.Bitmap) + snapshotID string + changeID string + setupMocks func(*cbtservicemocks.Service) expectedErrStr string + expectedCount uint64 + expectedNext []uint64 }{ { - name: "nil service", - nilService: true, - setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { - bmp.On("SetFull").Return() - }, + name: "nil service", + nilService: true, + snapshotID: "snap-1", + changeID: "change-1", + setupMocks: func(svc *cbtservicemocks.Service) {}, expectedErrStr: "CBT service is absent, fallback to real full", + expectedCount: 3, + expectedNext: []uint64{0, mb, 2 * mb}, }, { - name: "invalid snapshot", - setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { - bmp.On("Snapshot").Return("") - bmp.On("SetFull").Return() - }, + name: "invalid snapshot", + snapshotID: "", + setupMocks: func(svc *cbtservicemocks.Service) {}, expectedErrStr: "invalid snapshot, fallback to real full", + expectedCount: 3, + expectedNext: []uint64{0, mb, 2 * mb}, }, { - name: "invalid changeID", - incOnly: true, - setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { - bmp.On("Snapshot").Return("snap-1") - bmp.On("ChangeID").Return("") - bmp.On("SetFull").Return() - }, + name: "invalid changeID", + incOnly: true, + snapshotID: "snap-1", + changeID: "", + setupMocks: func(svc *cbtservicemocks.Service) {}, expectedErrStr: "invalid changeID, fallback to real full", + expectedCount: 3, + expectedNext: []uint64{0, mb, 2 * mb}, }, { - name: "allocated blocks success", - setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { - bmp.On("Snapshot").Return("snap-1") - bmp.On("ChangeID").Return("") + name: "allocated blocks success", + snapshotID: "snap-1", + changeID: "", + setupMocks: func(svc *cbtservicemocks.Service) { + svc.On("GetAllocatedBlocks", mock.Anything, "snap-1", mock.Anything).Run(func(args mock.Arguments) { + record := args.Get(2).(func([]cbtservice.Range) error) + record([]cbtservice.Range{ + {Offset: 0, Length: uint64(mb)}, + {Offset: uint64(2 * mb), Length: uint64(mb)}, + }) + }).Return(nil) + }, + expectedCount: 2, + expectedNext: []uint64{0, 2 * mb}, + }, + { + name: "allocated blocks error", + snapshotID: "snap-1", + changeID: "", + setupMocks: func(svc *cbtservicemocks.Service) { + svc.On("GetAllocatedBlocks", mock.Anything, "snap-1", mock.Anything).Return(errors.New("mock alloc error")) + }, + expectedErrStr: "error getting allocated blocks from CBT service, fallback to real full: mock alloc error", + expectedCount: 3, + expectedNext: []uint64{0, mb, 2 * mb}, + }, + { + name: "changed blocks success", + snapshotID: "snap-1", + changeID: "change-1", + setupMocks: func(svc *cbtservicemocks.Service) { + svc.On("GetChangedBlocks", mock.Anything, "snap-1", "change-1", mock.Anything).Run(func(args mock.Arguments) { + record := args.Get(3).(func([]cbtservice.Range) error) + record([]cbtservice.Range{ + {Offset: uint64(mb), Length: uint64(mb)}, + }) + }).Return(nil) + }, + expectedCount: 1, + expectedNext: []uint64{mb}, + }, + { + name: "changed blocks error with incOnly", + incOnly: true, + snapshotID: "snap-1", + changeID: "change-1", + setupMocks: func(svc *cbtservicemocks.Service) { + svc.On("GetChangedBlocks", mock.Anything, "snap-1", "change-1", mock.Anything).Return(errors.New("mock changed error")) + }, + expectedErrStr: "error getting changed blocks from CBT service, fallback to real full: mock changed error", + expectedCount: 3, + expectedNext: []uint64{0, mb, 2 * mb}, + }, + { + name: "both changed blocks error and allocated blocks error", + snapshotID: "snap-1", + changeID: "change-1", + setupMocks: func(svc *cbtservicemocks.Service) { + svc.On("GetChangedBlocks", mock.Anything, "snap-1", "change-1", mock.Anything).Return(errors.New("mock changed error")) + svc.On("GetAllocatedBlocks", mock.Anything, "snap-1", mock.Anything).Return(errors.New("mock alloc error")) + }, + expectedErrStr: "error getting both changed and allocated blocks from CBT service, fallback to real full: mock alloc error", + expectedCount: 3, + expectedNext: []uint64{0, mb, 2 * mb}, + }, + { + name: "changed blocks error fallback to full", + snapshotID: "snap-1", + changeID: "change-1", + setupMocks: func(svc *cbtservicemocks.Service) { + svc.On("GetChangedBlocks", mock.Anything, "snap-1", "change-1", mock.Anything).Return(errors.New("mock changed error")) svc.On("GetAllocatedBlocks", mock.Anything, "snap-1", mock.Anything).Run(func(args mock.Arguments) { record := args.Get(2).(func([]cbtservice.Range) error) record([]cbtservice.Range{ - {Offset: 0, Length: 4096}, - {Offset: 8192, Length: 4096}, + {Offset: 0, Length: uint64(mb)}, + {Offset: uint64(2 * mb), Length: uint64(mb)}, }) }).Return(nil) - - bmp.On("Set", uint64(0), uint64(4096)).Return() - bmp.On("Set", uint64(8192), uint64(4096)).Return() - }, - }, - { - name: "allocated blocks error", - setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { - bmp.On("Snapshot").Return("snap-1") - bmp.On("ChangeID").Return("") - - svc.On("GetAllocatedBlocks", mock.Anything, "snap-1", mock.Anything).Return(errors.New("mock alloc error")) - bmp.On("SetFull").Return() - }, - expectedErrStr: "error getting allocated blocks from CBT service, fallback to real full: mock alloc error", - }, - { - name: "changed blocks success", - setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { - bmp.On("Snapshot").Return("snap-1") - bmp.On("ChangeID").Return("change-1") - - svc.On("GetChangedBlocks", mock.Anything, "snap-1", "change-1", mock.Anything).Run(func(args mock.Arguments) { - record := args.Get(3).(func([]cbtservice.Range) error) - record([]cbtservice.Range{ - {Offset: 4096, Length: 4096}, - }) - }).Return(nil) - - bmp.On("Set", uint64(4096), uint64(4096)).Return() - }, - }, - { - name: "changed blocks error with incOnly", - incOnly: true, - setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { - bmp.On("Snapshot").Return("snap-1") - bmp.On("ChangeID").Return("change-1") - - svc.On("GetChangedBlocks", mock.Anything, "snap-1", "change-1", mock.Anything).Return(errors.New("mock changed error")) - bmp.On("SetFull").Return() - }, - expectedErrStr: "error getting changed blocks from CBT service, fallback to real full: mock changed error", - }, - { - name: "changed blocks error fallback to full", - setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { - bmp.On("Snapshot").Return("snap-1") - bmp.On("ChangeID").Return("change-1") - - svc.On("GetChangedBlocks", mock.Anything, "snap-1", "change-1", mock.Anything).Return(errors.New("mock changed error")) - - svc.On("GetAllocatedBlocks", mock.Anything, "snap-1", mock.Anything).Return(nil) }, expectedErrStr: "error getting changed blocks from CBT service, fallback to full: mock changed error", + expectedCount: 2, + expectedNext: []uint64{0, 2 * mb}, }, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { svcMock := new(cbtservicemocks.Service) - bmpMock := new(cbtmocks.Bitmap) if tt.setupMocks != nil { - tt.setupMocks(svcMock, bmpMock) + tt.setupMocks(svcMock) } var svc cbtservice.Service @@ -148,9 +170,9 @@ func TestSetBitmapOrFull(t *testing.T) { svc = svcMock } - bmpMock.On("SetError", mock.Anything).Return() + bmp := NewBitmap(mb, 3*mb, tt.snapshotID, tt.changeID, "vol-1") - err := SetBitmapOrFull(context.Background(), svc, bmpMock, tt.incOnly) + err := SetBitmapOrFull(context.Background(), svc, bmp, tt.incOnly) if tt.expectedErrStr != "" { require.Error(t, err) @@ -162,7 +184,25 @@ func TestSetBitmapOrFull(t *testing.T) { if !tt.nilService { svcMock.AssertExpectations(t) } - bmpMock.AssertExpectations(t) + + iter := bmp.Iterator() + require.NotNil(t, iter) + assert.Equal(t, tt.expectedCount, iter.Count()) + + var actualOffsets []uint64 + for { + offset, hasNext := iter.Next() + if !hasNext { + break + } + actualOffsets = append(actualOffsets, offset) + } + + if len(tt.expectedNext) > 0 { + assert.Equal(t, tt.expectedNext, actualOffsets) + } else { + assert.Empty(t, actualOffsets) + } }) } } diff --git a/pkg/uploader/kopia/snapshot.go b/pkg/uploader/kopia/snapshot.go index 682b0f314..c4c3251c9 100644 --- a/pkg/uploader/kopia/snapshot.go +++ b/pkg/uploader/kopia/snapshot.go @@ -251,9 +251,12 @@ func SnapshotSource( log.Infof("Using provided parent snapshot %s", parentSnapshot) if mani, err := loadSnapshotFunc(ctx, rep, manifest.ID(parentSnapshot)); err != nil { - msg := fmt.Sprintf("Failed to load previous snapshot %v from kopia, fallback to full backup", parentSnapshot) - log.WithError(err).Warn(msg) - updater.UpdateProgress(&uploader.Progress{BytesDone: -1, TotalBytes: -1, Message: msg}) + log.WithError(err).Warnf("Failed to load previous snapshot %v from kopia, fallback to full backup", parentSnapshot) + updater.UpdateProgress(&uploader.Progress{ + BytesDone: -1, + TotalBytes: -1, + Message: fmt.Sprintf("Failed to load previous snapshot %v, fallback to full backup. Err: %v", parentSnapshot, err), + }) } else { previous = append(previous, mani) } @@ -262,9 +265,11 @@ func SnapshotSource( if pre, err := findPreviousSnapshotManifest(ctx, rep, sourceInfo, snapshotTags, nil, log); err != nil { log.WithError(err).Warnf("Failed to find previous kopia snapshot manifests for si %v, fallback to full backup", sourceInfo) - - msg := fmt.Sprint("Failed to find previous kopia snapshot manifests, fallback to full backup") - updater.UpdateProgress(&uploader.Progress{BytesDone: -1, TotalBytes: -1, Message: msg}) + updater.UpdateProgress(&uploader.Progress{ + BytesDone: -1, + TotalBytes: -1, + Message: fmt.Sprintf("Failed to find previous snapshots, fallback to full backup. Err: %v", err), + }) } else { previous = pre } From 9d85334d4073883526dd42c5455e3c44a4288759 Mon Sep 17 00:00:00 2001 From: Yonghui Li Date: Tue, 8 Sep 2026 15:52:21 +0800 Subject: [PATCH 8/9] refactor CBT retrievement to report the concrete error Signed-off-by: Yonghui Li --- pkg/uploader/block/snapshot.go | 176 +++++++++++++++++++------------- pkg/uploader/block/uploader.go | 12 ++- pkg/uploader/cbt/bitmap.go | 19 ++-- pkg/uploader/cbt/types/types.go | 11 +- 4 files changed, 130 insertions(+), 88 deletions(-) diff --git a/pkg/uploader/block/snapshot.go b/pkg/uploader/block/snapshot.go index a74529d35..a2d96eade 100644 --- a/pkg/uploader/block/snapshot.go +++ b/pkg/uploader/block/snapshot.go @@ -40,6 +40,10 @@ type parentBackupInfo struct { volumeID string } +type backupInfo struct { + changeID string +} + // Backup backup specific sourcePath and update progress func Backup(ctx context.Context, blkUp Uploader, repoWriter udmrepo.BackupRepo, sourcePath string, realSource string, cbtSource cbtservice.SourceInfo, forceFull bool, parentSnapshot string, cbtService cbtservice.Service, uploaderCfg map[string]string, tags map[string]string, log logrus.FieldLogger) (uploader.SnapshotInfo, bool, error) { @@ -106,11 +110,17 @@ func snapshotSource( log.Info("Start to snapshot...") snapshotStartTime := time.Now() - parentBackup := getParentBackupInfo(ctx, rep, forceFull, parentSnapshot, cbtSource.VolumeID, source.realSource, snapshotTags, log) + bitmap := cbt.NewBitmap(blockSize, uint64(source.size), cbtSource.Snapshot, cbtSource.VolumeID) - bitmap := cbt.NewBitmap(blockSize, uint64(source.size), cbtSource.Snapshot, parentBackup.changeID, parentBackup.volumeID) + parentBackup, err := getParentBackupInfo(ctx, rep, forceFull, parentSnapshot, cbtSource.VolumeID, source.realSource, snapshotTags, log) + if err != nil { + log.WithError(err).Warn("Failed to get parent backup info, fallback to full backup") + bitmap.SetError(errors.Wrap(err, "error getting parent backup info, fallback to full backup")) + } else { + bitmap.SetChangeID(parentBackup.changeID) + } - err := cbt.SetBitmapOrFull(ctx, cbtService, bitmap, false) + err = cbt.SetBitmapOrFull(ctx, cbtService, bitmap, false) if err != nil { parentBackup.parentObject = "" log.WithError(err).Warnf("Failed to create CBT with source %v", cbtSource) @@ -147,61 +157,67 @@ func snapshotSource( return string(snapID), backupSize, nil } -func getParentBackupInfo(ctx context.Context, rep udmrepo.BackupRepo, forceFull bool, parentSnapshot string, volumeID string, realSource string, snapshotTags map[string]string, log logrus.FieldLogger) parentBackupInfo { - var previous *udmrepo.Snapshot - - // parentID names whichever snapshot ended up being the parent. On the discovery - // branch the parentSnapshot parameter is empty by definition, so logging it there - // produces messages that describe a decision without naming the object it was about. - parentID := parentSnapshot - - if !forceFull { - if parentSnapshot != "" { - snap, err := rep.GetSnapshot(ctx, udmrepo.ID(parentSnapshot)) - if err != nil { - log.WithError(err).Warn("Failed to load previous snapshot, fallback to full backup") - } else { - previous = &snap - log.Infof("Using provided parent snapshot %s", parentSnapshot) - } - } else { - log.Infof("Searching for parent snapshot") - - snap, err := findPreviousSnapshot(ctx, rep, realSource, snapshotTags, nil, log) - if err != nil { - log.WithError(err).Warn("Failed to search previous snapshot, fallback to full backup") - } else { - previous = &snap - parentID = string(snap.RootObject.ID) - log.Infof("Using previous snapshot %s", snap.RootObject.ID) - } - } - } else { +func getParentBackupInfo(ctx context.Context, rep udmrepo.BackupRepo, forceFull bool, parentSnapshot string, volumeID string, + realSource string, snapshotTags map[string]string, log logrus.FieldLogger) (parentBackupInfo, error) { + if forceFull { log.Info("Forcing full snapshot") + return parentBackupInfo{}, nil } - parentInfo := parentBackupInfo{} - if previous != nil { - if previous.Tags == nil { - log.Warnf("No tag from parent snapshot %s, fallback to full backup", parentID) - } else if previous.Tags[uploader.CBTChangeIDTag] == "" { - log.Warnf("No ChangeID tag from parent snapshot %s, fallback to full backup", parentID) - } else if previous.Tags[uploader.CBTVolumeIDTag] == "" { - log.Warnf("No VolumeID tag from parent snapshot %s, fallback to full backup", parentID) - } else if previous.Tags[uploader.CBTVolumeIDTag] != volumeID { - log.Warnf("VolumeID %s from parent snapshot %s is not expected as %s, fallback to full backup", previous.Tags[uploader.CBTVolumeIDTag], parentID, volumeID) - } else if obj, err := loadObjectFromSnapshot(ctx, rep, previous); err != nil { - log.WithError(err).Warnf("Failed to load object from parent snapshot %s, fallback to full backup", parentID) - } else { - parentInfo.parentObject = obj - parentInfo.changeID = previous.Tags[uploader.CBTChangeIDTag] - parentInfo.volumeID = previous.Tags[uploader.CBTVolumeIDTag] + if volumeID == "" { + return parentBackupInfo{}, errors.New("volumeID is not provided from the volume snapshot") + } - log.Infof("Using parent snapshot %s, start time %v, end time %v, description %s", parentID, previous.StartTime, previous.EndTime, previous.Description) + var previous *udmrepo.Snapshot + if parentSnapshot != "" { + log.Infof("Loading provided parent snapshot %s", parentSnapshot) + + snap, err := rep.GetSnapshot(ctx, udmrepo.ID(parentSnapshot)) + if err != nil { + return parentBackupInfo{}, errors.Wrapf(err, "error loading previous snapshot") } + + previous = &snap + + } else { + log.Infof("Searching for parent snapshot") + + snap, err := findPreviousSnapshot(ctx, rep, realSource, snapshotTags, nil, log) + if err != nil { + return parentBackupInfo{}, errors.Wrapf(err, "error searching previous snapshot") + } + + previous = &snap } - return parentInfo + if previous.Tags == nil { + return parentBackupInfo{}, errors.Errorf("no tag from parent snapshot %s", previous.ID) + } + + if previous.Tags[uploader.CBTChangeIDTag] == "" { + return parentBackupInfo{}, errors.Errorf("no ChangeID tag from parent snapshot %s", previous.ID) + } + + if previous.Tags[uploader.CBTVolumeIDTag] == "" { + return parentBackupInfo{}, errors.Errorf("no VolumeID tag from parent snapshot %s", previous.ID) + } + + if previous.Tags[uploader.CBTVolumeIDTag] != volumeID { + return parentBackupInfo{}, errors.Errorf("VolumeID %s from parent snapshot %s is not expected as %s", previous.Tags[uploader.CBTVolumeIDTag], previous.ID, volumeID) + } + + obj, err := loadObjectFromSnapshot(ctx, rep, previous) + if err != nil { + return parentBackupInfo{}, errors.Errorf("error loading object from parent snapshot %s", previous.ID) + } + + log.Infof("Using parent snapshot %s, start time %v, end time %v, description %s", previous.ID, previous.StartTime, previous.EndTime, previous.Description) + + return parentBackupInfo{ + parentObject: obj, + changeID: previous.Tags[uploader.CBTChangeIDTag], + volumeID: previous.Tags[uploader.CBTVolumeIDTag], + }, nil } // Restore restore specific sourcePath with given snapshotID and update progress @@ -214,29 +230,19 @@ func Restore(ctx context.Context, blkUp Uploader, rep udmrepo.BackupRepo, snapsh } log.Infof("Restore from snapshot %s, incremental %v, cbt source %v, description %s, created time %v, tags %v", snapshotID, incremental, cbtSource, snapshot.Description, snapshot.EndTime, snapshot.Tags) - var volumeSnapshot, changeID, volumeID string - if incremental { - if snapshot.Tags == nil { - log.Warnf("No tag from snapshot %s, fallback to full restore", snapshotID) - } else if snapshot.Tags[uploader.CBTChangeIDTag] == "" { - log.Warnf("No ChangeID tag from snapshot %s, fallback to full restore", snapshotID) - } else if snapshot.Tags[uploader.CBTVolumeIDTag] == "" { - log.Warnf("No VolumeID tag from snapshot %s, fallback to full restore", snapshotID) - } else if cbtSource.VolumeID == "" { - log.Warnf("No VolumeID in cbt source %v, fallback to full restore", cbtSource) - } else if snapshot.Tags[uploader.CBTVolumeIDTag] != cbtSource.VolumeID { - log.Warnf("VolumeID %s from snapshot %s is not expected as %s, fallback to full restore", snapshot.Tags[uploader.CBTVolumeIDTag], snapshotID, cbtSource.VolumeID) - } else { - volumeSnapshot = cbtSource.Snapshot - changeID = snapshot.Tags[uploader.CBTChangeIDTag] - volumeID = snapshot.Tags[uploader.CBTVolumeIDTag] - } - } + bitmap := cbt.NewBitmap(blockSize, uint64(snapshot.TotalSize), cbtSource.Snapshot, cbtSource.VolumeID) - bitmap := cbt.NewBitmap(blockSize, uint64(snapshot.TotalSize), volumeSnapshot, changeID, volumeID) if incremental { - if err = cbt.SetBitmapOrFull(ctx, cbtService, bitmap, true); err != nil { - log.WithError(err).Warnf("Failed to create CBT with source %v", cbtSource) + if bkInfo, err := getBackupInfo(snapshot, cbtSource.VolumeID); err != nil { + log.WithError(err).Warn("Failed to get backup info, fallback to full restore") + + bitmap.SetError(errors.Wrap(err, "error getting backup info, fallback to full restore")) + bitmap.SetFull() + } else { + bitmap.SetChangeID(bkInfo.changeID) + if err = cbt.SetBitmapOrFull(ctx, cbtService, bitmap, true); err != nil { + log.WithError(err).Warnf("Failed to create CBT with source %v", cbtSource) + } } } else { bitmap.SetFull() @@ -274,6 +280,32 @@ func Restore(ctx context.Context, blkUp Uploader, rep udmrepo.BackupRepo, snapsh return incrementalBytes, totalSize, nil } +func getBackupInfo(snapshot udmrepo.Snapshot, volumeID string) (backupInfo, error) { + if snapshot.Tags == nil { + return backupInfo{}, errors.Errorf("no tag from snapshot %s", snapshot.ID) + } + + if snapshot.Tags[uploader.CBTChangeIDTag] == "" { + return backupInfo{}, errors.Errorf("no ChangeID tag from snapshot %s", snapshot.ID) + } + + if snapshot.Tags[uploader.CBTVolumeIDTag] == "" { + return backupInfo{}, errors.Errorf("no VolumeID tag from snapshot %s", snapshot.ID) + } + + if volumeID == "" { + return backupInfo{}, errors.New("no VolumeID tag from the volume snapshot") + } + + if snapshot.Tags[uploader.CBTVolumeIDTag] != volumeID { + return backupInfo{}, errors.Errorf("volumeID %s from snapshot %s is not expected as %s", snapshot.Tags[uploader.CBTVolumeIDTag], snapshot.ID, volumeID) + } + + return backupInfo{ + changeID: snapshot.Tags[uploader.CBTChangeIDTag], + }, nil +} + func findPreviousSnapshot(ctx context.Context, rep udmrepo.BackupRepo, path string, snapshotTags map[string]string, noLaterThan *time.Time, log logrus.FieldLogger) (udmrepo.Snapshot, error) { snaps, err := rep.ListSnapshot(ctx, path) if err != nil { diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index 4b69de59f..283526672 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -86,8 +86,10 @@ func (blkup *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, b return udmrepo.Snapshot{}, 0, errors.New("bitmap is not available") } - if bitmap.Error() != nil { - blkup.progress.UpdateProgress(&uploader.Progress{BytesDone: -1, TotalBytes: -1, Message: bitmap.Error().Error()}) + if bitmap.Errors() != nil { + for _, err := range bitmap.Errors() { + blkup.progress.UpdateProgress(&uploader.Progress{BytesDone: -1, TotalBytes: -1, Message: err.Error()}) + } } backupMode := udmrepo.ObjectDataBackupModeInc @@ -157,8 +159,10 @@ func (blkup *blockUploader) Restore(snapshot udmrepo.Snapshot, dest destInfo, bi return 0, 0, errors.New("bitmap is not available") } - if bitmap.Error() != nil { - blkup.progress.UpdateProgress(&uploader.Progress{BytesDone: -1, TotalBytes: -1, Message: bitmap.Error().Error()}) + if bitmap.Errors() != nil { + for _, err := range bitmap.Errors() { + blkup.progress.UpdateProgress(&uploader.Progress{BytesDone: -1, TotalBytes: -1, Message: err.Error()}) + } } meta, err := blkup.repoWriter.ReadMetadata(blkup.ctx, snapshot.RootObject.ID) diff --git a/pkg/uploader/cbt/bitmap.go b/pkg/uploader/cbt/bitmap.go index f1b4b3c7d..29da0c190 100644 --- a/pkg/uploader/cbt/bitmap.go +++ b/pkg/uploader/cbt/bitmap.go @@ -36,7 +36,7 @@ type bitmapImpl struct { snapshot string changeID string volumeID string - cbtError error + cbtErrors []error } type bitmapIterator struct { @@ -44,14 +44,13 @@ type bitmapIterator struct { iterator roaring.IntPeekable } -func NewBitmap(blockSize uint, length uint64, snapshot string, changeID string, volumeID string) types.Bitmap { +func NewBitmap(blockSize uint, length uint64, snapshot string, volumeID string) types.Bitmap { return &bitmapImpl{ bitmap: roaring.New(), blockSize: blockSize, blockSizeLog: bits.Len(blockSize) - 1, length: length, snapshot: snapshot, - changeID: changeID, volumeID: volumeID, } } @@ -82,6 +81,10 @@ func (c *bitmapImpl) Snapshot() string { return c.snapshot } +func (c *bitmapImpl) SetChangeID(id string) { + c.changeID = id +} + func (c *bitmapImpl) ChangeID() string { return c.changeID } @@ -91,11 +94,11 @@ func (c *bitmapImpl) VolumeID() string { } func (c *bitmapImpl) SetError(err error) { - c.cbtError = err + c.cbtErrors = append(c.cbtErrors, err) } -func (c *bitmapImpl) Error() error { - return c.cbtError +func (c *bitmapImpl) Errors() []error { + return c.cbtErrors } func (c *bitmapImpl) Iterator() types.Iterator { @@ -125,6 +128,6 @@ func (c *bitmapIterator) BlockSize() uint { return c.blockSize } -func (c *bitmapIterator) Error() error { - return c.cbtError +func (c *bitmapIterator) Errors() []error { + return c.cbtErrors } diff --git a/pkg/uploader/cbt/types/types.go b/pkg/uploader/cbt/types/types.go index 447c3f9f3..025d822dc 100644 --- a/pkg/uploader/cbt/types/types.go +++ b/pkg/uploader/cbt/types/types.go @@ -39,8 +39,11 @@ type Bitmap interface { // SetError sets CBT error when preparing this bitmap SetError(error) - // Error returns the CBT error when preparing this bitmap - Error() error + // Errors returns the CBT errors when preparing this bitmap + Errors() []error + + // SetChangeID sets the changeID of the bitmap + SetChangeID(string) } // Iterator defines the methods to iterate the CBT bitmap and query the associated information @@ -63,6 +66,6 @@ type Iterator interface { // Next returns the offset of the next set block and whether it comes to the end of the iteration Next() (uint64, bool) - // Error returns the CBT error when preparing this bitmap - Error() error + // Errors returns the CBT errors when preparing this bitmap + Errors() []error } From e52bf08c99be279e4870278b6379767183a6e276 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Wed, 9 Sep 2026 14:34:36 +0800 Subject: [PATCH 9/9] refactor CBT retrievement to report the concrete error Signed-off-by: Lyndon-Li --- pkg/uploader/block/snapshot.go | 75 +++++----- pkg/uploader/block/snapshot_test.go | 174 +++++++++++++++++++++-- pkg/uploader/block/uploader_test.go | 15 +- pkg/uploader/cbt/bitmap.go | 4 + pkg/uploader/cbt/bitmap_test.go | 23 ++- pkg/uploader/cbt/set_test.go | 3 +- pkg/uploader/cbt/types/mocks/Bitmap.go | 19 ++- pkg/uploader/cbt/types/mocks/Iterator.go | 14 +- 8 files changed, 249 insertions(+), 78 deletions(-) diff --git a/pkg/uploader/block/snapshot.go b/pkg/uploader/block/snapshot.go index 9abdff3c1..d75e94f3f 100644 --- a/pkg/uploader/block/snapshot.go +++ b/pkg/uploader/block/snapshot.go @@ -169,52 +169,49 @@ func getParentBackupInfo(ctx context.Context, rep udmrepo.BackupRepo, forceFull } var previous *udmrepo.Snapshot + if parentSnapshot != "" { + log.Infof("Loading provided parent snapshot %s", parentSnapshot) - if !forceFull { - if parentSnapshot != "" { - snap, err := rep.GetSnapshot(ctx, udmrepo.ID(parentSnapshot)) - if err != nil { - log.WithError(err).Warn("Failed to load previous snapshot, fallback to full backup") - } else { - previous = &snap - log.Infof("Using provided parent snapshot %s", parentSnapshot) - } - } else { - log.Infof("Searching for parent snapshot") - - snap, err := findPreviousSnapshot(ctx, rep, realSource, snapshotTags, nil, log) - if err != nil { - log.WithError(err).Warn("Failed to search previous snapshot, fallback to full backup") - } else { - previous = &snap - log.Infof("Using previous snapshot %s", snap.ID) - } + snap, err := rep.GetSnapshot(ctx, udmrepo.ID(parentSnapshot)) + if err != nil { + return parentBackupInfo{}, errors.Wrapf(err, "error loading previous snapshot") } + + previous = &snap } else { - log.Info("Forcing full snapshot") - } + log.Infof("Searching for parent snapshot") - parentInfo := parentBackupInfo{} - if previous != nil { - if previous.Tags == nil { - log.Warnf("No tag from parent snapshot %s, fallback to full backup", previous.ID) - } else if previous.Tags[uploader.CBTChangeIDTag] == "" { - log.Warnf("No ChangeID tag from parent snapshot %s, fallback to full backup", previous.ID) - } else if previous.Tags[uploader.CBTVolumeIDTag] == "" { - log.Warnf("No VolumeID tag from parent snapshot %s, fallback to full backup", previous.ID) - } else if previous.Tags[uploader.CBTVolumeIDTag] != volumeID { - log.Warnf("VolumeID %s from parent snapshot %s is not expected as %s, fallback to full backup", previous.Tags[uploader.CBTVolumeIDTag], previous.ID, volumeID) - } else if obj, err := loadObjectFromSnapshot(ctx, rep, previous); err != nil { - log.WithError(err).Warnf("Failed to load object from parent snapshot %s, fallback to full backup", previous.ID) - } else { - parentInfo.parentObject = obj - parentInfo.changeID = previous.Tags[uploader.CBTChangeIDTag] - parentInfo.volumeID = previous.Tags[uploader.CBTVolumeIDTag] - - log.Infof("Using parent snapshot %s, start time %v, end time %v, description %s", previous.ID, previous.StartTime, previous.EndTime, previous.Description) + snap, err := findPreviousSnapshot(ctx, rep, realSource, snapshotTags, nil, log) + if err != nil { + return parentBackupInfo{}, errors.Wrapf(err, "error searching previous snapshot") } + + previous = &snap } + if previous.Tags == nil { + return parentBackupInfo{}, errors.Errorf("no tag from parent snapshot %s", previous.ID) + } + + if previous.Tags[uploader.CBTChangeIDTag] == "" { + return parentBackupInfo{}, errors.Errorf("no ChangeID tag from parent snapshot %s", previous.ID) + } + + if previous.Tags[uploader.CBTVolumeIDTag] == "" { + return parentBackupInfo{}, errors.Errorf("no VolumeID tag from parent snapshot %s", previous.ID) + } + + if previous.Tags[uploader.CBTVolumeIDTag] != volumeID { + return parentBackupInfo{}, errors.Errorf("VolumeID %s from parent snapshot %s is not expected as %s", previous.Tags[uploader.CBTVolumeIDTag], previous.ID, volumeID) + } + + obj, err := loadObjectFromSnapshot(ctx, rep, previous) + if err != nil { + return parentBackupInfo{}, errors.Wrapf(err, "error loading object from parent snapshot %s", previous.ID) + } + + log.Infof("Using parent snapshot %s, start time %v, end time %v, description %s", previous.ID, previous.StartTime, previous.EndTime, previous.Description) + return parentBackupInfo{ parentObject: obj, changeID: previous.Tags[uploader.CBTChangeIDTag], diff --git a/pkg/uploader/block/snapshot_test.go b/pkg/uploader/block/snapshot_test.go index 260f94b82..fb5eaf4f7 100644 --- a/pkg/uploader/block/snapshot_test.go +++ b/pkg/uploader/block/snapshot_test.go @@ -379,11 +379,12 @@ func TestGetParentBackupInfoLogsDiscoveredParentID(t *testing.T) { SubObjects: []udmrepo.ObjectMetadata{{ID: udmrepo.ID("parent-obj")}}, }, nil) - info := getParentBackupInfo( + info, err := getParentBackupInfo( context.Background(), repo, false, "", // no explicit parent -> discovery branch volumeID, realSource, snapshotTags, logger, ) + require.NoError(t, err) require.Equal(t, udmrepo.ID("parent-obj"), info.parentObject) @@ -408,6 +409,7 @@ func TestGetParentBackupInfo(t *testing.T) { } validSnap := udmrepo.Snapshot{ + ID: "snap-valid", RootObject: udmrepo.ObjectMetadata{ID: "root-obj"}, Tags: map[string]string{ uploader.CBTChangeIDTag: "cid-abc", @@ -421,7 +423,10 @@ func TestGetParentBackupInfo(t *testing.T) { name string forceFull bool parentSnapshot string + emptyVolID bool setupMocks func(repo *udmrepomocks.BackupRepo) + expectErr bool + expectedErrStr string expectEmpty bool expectedParent udmrepo.ID expectedCID string @@ -432,6 +437,13 @@ func TestGetParentBackupInfo(t *testing.T) { forceFull: true, expectEmpty: true, }, + { + name: "volumeID not provided", + emptyVolID: true, + expectEmpty: true, + expectErr: true, + expectedErrStr: "volumeID is not provided from the volume snapshot", + }, { name: "GetSnapshot fails — falls back to full", parentSnapshot: "snap-parent", @@ -439,46 +451,69 @@ func TestGetParentBackupInfo(t *testing.T) { repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-parent")). Return(udmrepo.Snapshot{}, errors.New("not found")) }, - expectEmpty: true, + expectEmpty: true, + expectErr: true, + expectedErrStr: "error loading previous snapshot", }, { name: "parent snapshot has nil tags — falls back to full", parentSnapshot: "snap-notags", setupMocks: func(repo *udmrepomocks.BackupRepo) { repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-notags")). - Return(udmrepo.Snapshot{Tags: nil}, nil) + Return(udmrepo.Snapshot{ID: "snap-notags", Tags: nil}, nil) }, - expectEmpty: true, + expectEmpty: true, + expectErr: true, + expectedErrStr: "no tag from parent snapshot snap-notags", }, { name: "parent snapshot missing ChangeID tag — falls back to full", parentSnapshot: "snap-nocid", setupMocks: func(repo *udmrepomocks.BackupRepo) { repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-nocid")). - Return(udmrepo.Snapshot{Tags: map[string]string{uploader.CBTVolumeIDTag: volumeID}}, nil) + Return(udmrepo.Snapshot{ID: "snap-nocid", Tags: map[string]string{uploader.CBTVolumeIDTag: volumeID}}, nil) }, - expectEmpty: true, + expectEmpty: true, + expectErr: true, + expectedErrStr: "no ChangeID tag from parent snapshot snap-nocid", }, { name: "parent snapshot missing VolumeID tag — falls back to full", parentSnapshot: "snap-novid", setupMocks: func(repo *udmrepomocks.BackupRepo) { repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-novid")). - Return(udmrepo.Snapshot{Tags: map[string]string{uploader.CBTChangeIDTag: "cid"}}, nil) + Return(udmrepo.Snapshot{ID: "snap-novid", Tags: map[string]string{uploader.CBTChangeIDTag: "cid"}}, nil) }, - expectEmpty: true, + expectEmpty: true, + expectErr: true, + expectedErrStr: "no VolumeID tag from parent snapshot snap-novid", }, { name: "parent snapshot VolumeID mismatch — falls back to full", parentSnapshot: "snap-vidmismatch", setupMocks: func(repo *udmrepomocks.BackupRepo) { repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-vidmismatch")). - Return(udmrepo.Snapshot{Tags: map[string]string{ + Return(udmrepo.Snapshot{ID: "snap-vidmismatch", Tags: map[string]string{ uploader.CBTChangeIDTag: "cid", uploader.CBTVolumeIDTag: "different-vol", }}, nil) }, - expectEmpty: true, + expectEmpty: true, + expectErr: true, + expectedErrStr: "VolumeID different-vol from parent snapshot snap-vidmismatch is not expected as vol-123", + }, + { + name: "loadObjectFromSnapshot fails — falls back to full", + parentSnapshot: "snap-valid", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-valid")). + Return(validSnap, nil) + repo.On("ReadMetadata", mock.Anything, udmrepo.ID("root-obj")). + Return(nil, errors.New("read error")) + }, + expectEmpty: true, + expectErr: true, + expectedErrStr: "error loading object from parent snapshot snap-valid", }, { name: "valid parent snapshot — returns parent info", @@ -499,7 +534,9 @@ func TestGetParentBackupInfo(t *testing.T) { repo.On("ListSnapshot", mock.Anything, realSource). Return(nil, errors.New("list error")) }, - expectEmpty: true, + expectEmpty: true, + expectErr: true, + expectedErrStr: "error searching previous snapshot", }, { name: "no parentSnapshot — no matching snapshot — falls back to full", @@ -507,7 +544,9 @@ func TestGetParentBackupInfo(t *testing.T) { repo.On("ListSnapshot", mock.Anything, realSource). Return([]udmrepo.Snapshot{{Tags: map[string]string{"other": "tag"}}}, nil) }, - expectEmpty: true, + expectEmpty: true, + expectErr: true, + expectedErrStr: "error searching previous snapshot", }, { name: "no parentSnapshot — matching snapshot found — returns parent info", @@ -532,7 +571,21 @@ func TestGetParentBackupInfo(t *testing.T) { tc.setupMocks(mockRepo) } - info := getParentBackupInfo(ctx, mockRepo, tc.forceFull, tc.parentSnapshot, volumeID, realSource, snapshotTags, testLog()) + volID := volumeID + if tc.emptyVolID { + volID = "" + } + + info, err := getParentBackupInfo(ctx, mockRepo, tc.forceFull, tc.parentSnapshot, volID, realSource, snapshotTags, testLog()) + + if tc.expectErr { + require.Error(t, err) + if tc.expectedErrStr != "" { + assert.Contains(t, err.Error(), tc.expectedErrStr) + } + } else { + require.NoError(t, err) + } if tc.expectEmpty { assert.Empty(t, info.parentObject) @@ -547,6 +600,101 @@ func TestGetParentBackupInfo(t *testing.T) { } } +func TestGetBackupInfo(t *testing.T) { + const volumeID = "vol-123" + + validSnap := udmrepo.Snapshot{ + ID: "snap-valid", + Tags: map[string]string{ + uploader.CBTChangeIDTag: "cid-abc", + uploader.CBTVolumeIDTag: volumeID, + }, + } + + testCases := []struct { + name string + snapshot udmrepo.Snapshot + volumeID string + expectErr bool + expectedErrStr string + expectedCID string + }{ + { + name: "nil tags", + snapshot: udmrepo.Snapshot{ID: "snap-nil-tags"}, + volumeID: volumeID, + expectErr: true, + expectedErrStr: "no tag from snapshot snap-nil-tags", + }, + { + name: "missing ChangeID tag", + snapshot: udmrepo.Snapshot{ + ID: "snap-no-cid", + Tags: map[string]string{uploader.CBTVolumeIDTag: volumeID}, + }, + volumeID: volumeID, + expectErr: true, + expectedErrStr: "no ChangeID tag from snapshot snap-no-cid", + }, + { + name: "missing VolumeID tag", + snapshot: udmrepo.Snapshot{ + ID: "snap-no-vid", + Tags: map[string]string{uploader.CBTChangeIDTag: "cid-abc"}, + }, + volumeID: volumeID, + expectErr: true, + expectedErrStr: "no VolumeID tag from snapshot snap-no-vid", + }, + { + name: "empty volumeID parameter", + snapshot: udmrepo.Snapshot{ + ID: "snap-valid", + Tags: map[string]string{ + uploader.CBTChangeIDTag: "cid-abc", + uploader.CBTVolumeIDTag: volumeID, + }, + }, + volumeID: "", + expectErr: true, + expectedErrStr: "no VolumeID tag from the volume snapshot", + }, + { + name: "volumeID mismatch", + snapshot: udmrepo.Snapshot{ + ID: "snap-vid-mismatch", + Tags: map[string]string{ + uploader.CBTChangeIDTag: "cid-abc", + uploader.CBTVolumeIDTag: "other-vol", + }, + }, + volumeID: volumeID, + expectErr: true, + expectedErrStr: "volumeID other-vol from snapshot snap-vid-mismatch is not expected as vol-123", + }, + { + name: "valid snapshot", + snapshot: validSnap, + volumeID: volumeID, + expectErr: false, + expectedCID: "cid-abc", + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + info, err := getBackupInfo(tc.snapshot, tc.volumeID) + if tc.expectErr { + require.Error(t, err) + assert.Contains(t, err.Error(), tc.expectedErrStr) + } else { + require.NoError(t, err) + assert.Equal(t, tc.expectedCID, info.changeID) + } + }) + } +} + func TestFindPreviousSnapshot(t *testing.T) { snapshotTags := map[string]string{ uploader.SnapshotRequesterTag: "test-requester", diff --git a/pkg/uploader/block/uploader_test.go b/pkg/uploader/block/uploader_test.go index f011219e4..b3621ed33 100644 --- a/pkg/uploader/block/uploader_test.go +++ b/pkg/uploader/block/uploader_test.go @@ -295,8 +295,7 @@ func TestBlockUploaderBackup(t *testing.T) { var iterator cbt.Iterator if !tc.nilBitmap { iterMock := cbtmocks.NewIterator(t) - iterMock.On("Error").Return(nil).Maybe() - iterMock.On("Error").Return(nil).Maybe() + iterMock.On("Errors").Return(nil).Maybe() iterator = iterMock backupMode := udmrepo.ObjectDataBackupModeInc @@ -575,7 +574,7 @@ func TestRestoreData(t *testing.T) { reader := bytes.NewReader(data) iterMock := cbtmocks.NewIterator(t) - iterMock.On("Error").Return(nil).Maybe() + iterMock.On("Errors").Return(nil).Maybe() iterMock.On("Count").Return(uint64(1)) iterMock.On("Next").Return(uint64(0), true).Once() iterMock.On("Next").Return(uint64(0), false) @@ -606,7 +605,7 @@ func TestRestoreData(t *testing.T) { reader := &errReader{err: errors.New("read error")} iterMock := cbtmocks.NewIterator(t) - iterMock.On("Error").Return(nil).Maybe() + iterMock.On("Errors").Return(nil).Maybe() iterMock.On("Count").Return(uint64(1)) iterMock.On("Next").Return(uint64(0), true).Once() iterMock.On("Next").Return(uint64(0), false) @@ -627,7 +626,7 @@ func TestBlockUploaderRestore(t *testing.T) { repoWriter.On("ReadMetadata", mock.Anything, udmrepo.ID("root-id")).Return(nil, errors.New("meta not found")) iterMock := cbtmocks.NewIterator(t) - iterMock.On("Error").Return(nil).Maybe() + iterMock.On("Errors").Return(nil).Maybe() _, _, err := blkup.Restore(udmrepo.Snapshot{RootObject: udmrepo.ObjectMetadata{ID: "root-id"}}, destInfo{}, iterMock, nil) require.Error(t, err) assert.Contains(t, err.Error(), "meta not found") @@ -685,7 +684,7 @@ func TestBlockUploaderRestore(t *testing.T) { } iterMock := cbtmocks.NewIterator(t) - iterMock.On("Error").Return(nil).Maybe() + iterMock.On("Errors").Return(nil).Maybe() iterMock.On("Count").Return(uint64(1)) iterMock.On("Next").Return(uint64(0), true).Once() iterMock.On("Next").Return(uint64(0), false) @@ -714,7 +713,7 @@ func TestBlockUploaderRestore(t *testing.T) { } dest := destInfo{size: 4194304, path: "/dev/target"} iterMock := cbtmocks.NewIterator(t) - iterMock.On("Error").Return(nil).Maybe() + iterMock.On("Errors").Return(nil).Maybe() _, _, err := blkup.Restore(snap, dest, iterMock, nil) require.Error(t, err) @@ -739,7 +738,7 @@ func TestBlockUploaderRestore(t *testing.T) { } dest := destInfo{size: 512, path: "/dev/small"} iterMock := cbtmocks.NewIterator(t) - iterMock.On("Error").Return(nil).Maybe() + iterMock.On("Errors").Return(nil).Maybe() _, _, err := blkup.Restore(snap, dest, iterMock, nil) require.Error(t, err) diff --git a/pkg/uploader/cbt/bitmap.go b/pkg/uploader/cbt/bitmap.go index 29da0c190..d861ea6da 100644 --- a/pkg/uploader/cbt/bitmap.go +++ b/pkg/uploader/cbt/bitmap.go @@ -94,6 +94,10 @@ func (c *bitmapImpl) VolumeID() string { } func (c *bitmapImpl) SetError(err error) { + if err == nil { + return + } + c.cbtErrors = append(c.cbtErrors, err) } diff --git a/pkg/uploader/cbt/bitmap_test.go b/pkg/uploader/cbt/bitmap_test.go index 9671ad3d9..8e428f7bd 100644 --- a/pkg/uploader/cbt/bitmap_test.go +++ b/pkg/uploader/cbt/bitmap_test.go @@ -17,6 +17,7 @@ limitations under the License. package cbt import ( + "errors" "testing" "github.com/stretchr/testify/assert" @@ -24,10 +25,18 @@ import ( ) func TestBitmapProperties(t *testing.T) { - b := NewBitmap(1024*1024, 10000*1024*1024, "snap-1", "change-1", "vol-1") + b := NewBitmap(1024*1024, 10000*1024*1024, "snap-1", "vol-1") assert.Equal(t, "snap-1", b.Snapshot()) - assert.Equal(t, "change-1", b.ChangeID()) + assert.Empty(t, b.ChangeID()) assert.Equal(t, "vol-1", b.VolumeID()) + assert.Empty(t, b.Errors()) + + b.SetChangeID("change-1") + assert.Equal(t, "change-1", b.ChangeID()) + + err := errors.New("test error") + b.SetError(err) + assert.Equal(t, []error{err}, b.Errors()) } func TestBitmapSet(t *testing.T) { @@ -138,7 +147,7 @@ func TestBitmapSet(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - b := NewBitmap(tt.blockSize, tt.totalLength, "snap-1", "change-1", "vol-1") + b := NewBitmap(tt.blockSize, tt.totalLength, "snap-1", "vol-1") for _, call := range tt.setCalls { b.Set(call.offset, call.length) @@ -174,7 +183,7 @@ func TestBitmapSetFull(t *testing.T) { // block 0: 0 - 1MB // block 1: 1MB - 2MB // block 2: 2MB - 3MB - b := NewBitmap(mb, 3*mb, "snap-1", "change-1", "vol-1") + b := NewBitmap(mb, 3*mb, "snap-1", "vol-1") b.SetFull() iter := b.Iterator() @@ -199,7 +208,10 @@ func TestBitmapIterator(t *testing.T) { const mb = 1024 * 1024 const gb = 1024 * 1024 * 1024 - b := NewBitmap(mb, 10*gb, "snap-1", "change-1", "vol-1") + b := NewBitmap(mb, 10*gb, "snap-1", "vol-1") + b.SetChangeID("change-1") + err := errors.New("test error") + b.SetError(err) // Set multiple ranges to test contiguous iteration b.Set(mb, 100) // Block 1 @@ -214,6 +226,7 @@ func TestBitmapIterator(t *testing.T) { assert.Equal(t, "change-1", iter.ChangeID()) assert.Equal(t, "vol-1", iter.VolumeID()) assert.Equal(t, uint(mb), iter.BlockSize()) + assert.Equal(t, []error{err}, iter.Errors()) assert.Equal(t, uint64(7), iter.Count()) // 1 + 5 + 1 = 7 blocks expectedOffsets := []uint64{ diff --git a/pkg/uploader/cbt/set_test.go b/pkg/uploader/cbt/set_test.go index 6c1d4a579..e9a6d1a52 100644 --- a/pkg/uploader/cbt/set_test.go +++ b/pkg/uploader/cbt/set_test.go @@ -170,7 +170,8 @@ func TestSetBitmapOrFull(t *testing.T) { svc = svcMock } - bmp := NewBitmap(mb, 3*mb, tt.snapshotID, tt.changeID, "vol-1") + bmp := NewBitmap(mb, 3*mb, tt.snapshotID, "vol-1") + bmp.SetChangeID(tt.changeID) err := SetBitmapOrFull(context.Background(), svc, bmp, tt.incOnly) diff --git a/pkg/uploader/cbt/types/mocks/Bitmap.go b/pkg/uploader/cbt/types/mocks/Bitmap.go index 514e5c28f..3ed29917b 100644 --- a/pkg/uploader/cbt/types/mocks/Bitmap.go +++ b/pkg/uploader/cbt/types/mocks/Bitmap.go @@ -292,23 +292,30 @@ func (_c *Bitmap_VolumeID_Call) RunAndReturn(run func() string) *Bitmap_VolumeID _c.Call.Return(run) return _c } +// SetChangeID provides a mock function for the type Bitmap +func (_mock *Bitmap) SetChangeID(id string) { + _mock.Called(id) +} + // SetError provides a mock function for the type Bitmap func (_mock *Bitmap) SetError(err error) { _mock.Called(err) } -// Error provides a mock function for the type Bitmap -func (_mock *Bitmap) Error() error { +// Errors provides a mock function for the type Bitmap +func (_mock *Bitmap) Errors() []error { ret := _mock.Called() if len(ret) == 0 { - panic("no return value specified for Error") + panic("no return value specified for Errors") } - var r0 error - if returnFunc, ok := ret.Get(0).(func() error); ok { + var r0 []error + if returnFunc, ok := ret.Get(0).(func() []error); ok { r0 = returnFunc() } else { - r0 = ret.Error(0) + if ret.Get(0) != nil { + r0 = ret.Get(0).([]error) + } } return r0 } diff --git a/pkg/uploader/cbt/types/mocks/Iterator.go b/pkg/uploader/cbt/types/mocks/Iterator.go index 00fc163d9..1ddd59f76 100644 --- a/pkg/uploader/cbt/types/mocks/Iterator.go +++ b/pkg/uploader/cbt/types/mocks/Iterator.go @@ -307,18 +307,20 @@ func (_c *Iterator_VolumeID_Call) RunAndReturn(run func() string) *Iterator_Volu _c.Call.Return(run) return _c } -// Error provides a mock function for the type Iterator -func (_mock *Iterator) Error() error { +// Errors provides a mock function for the type Iterator +func (_mock *Iterator) Errors() []error { ret := _mock.Called() if len(ret) == 0 { - panic("no return value specified for Error") + panic("no return value specified for Errors") } - var r0 error - if returnFunc, ok := ret.Get(0).(func() error); ok { + var r0 []error + if returnFunc, ok := ret.Get(0).(func() []error); ok { r0 = returnFunc() } else { - r0 = ret.Error(0) + if ret.Get(0) != nil { + r0 = ret.Get(0).([]error) + } } return r0 }