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 diff --git a/pkg/controller/data_download_controller.go b/pkg/controller/data_download_controller.go index e2a62830e..053a083e7 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,21 @@ 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 != "" { + message := progress.Message + ";" + if !strings.HasSuffix(dd.Status.Message, message) { + dd.Status.Message += message + } + } + return true }); err != nil { log.WithError(err).Error("Failed to update progress") 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.go b/pkg/controller/data_upload_controller.go index ae50a7741..dbbf21013 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,21 @@ 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 != "" { + message := progress.Message + ";" + if !strings.HasSuffix(du.Status.Message, message) { + du.Status.Message += message + } + } + return true }); err != nil { log.WithError(err).Error("Failed to update progress") 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.go b/pkg/controller/pod_volume_backup_controller.go index c4e68ce33..3e343b073 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,21 @@ 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 != "" { + message := progress.Message + ";" + if !strings.HasSuffix(pvb.Status.Message, message) { + pvb.Status.Message += message + } + } + return true }); err != nil { log.WithError(err).Error("Failed to update progress") 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.go b/pkg/controller/pod_volume_restore_controller.go index 4fe9bdaa3..d8ffb0e1a 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,21 @@ 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 != "" { + message := progress.Message + ";" + if !strings.HasSuffix(pvr.Status.Message, message) { + pvr.Status.Message += message + } + } + return true }); err != nil { log.WithError(err).Error("Failed to update progress") 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/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/snapshot.go b/pkg/uploader/block/snapshot.go index 370737bd6..d75e94f3f 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,14 +110,20 @@ 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) + 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) @@ -147,55 +157,66 @@ 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 - - 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) - } - } - } 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", 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] + 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", previous.ID, 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.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], + volumeID: previous.Tags[uploader.CBTVolumeIDTag], + }, nil } // Restore restore specific sourcePath with given snapshotID and update progress @@ -208,34 +229,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) - 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 cbtSource.VolumeID == "" { - log.Warnf("No VolumeID in cbt source %v, fallback to full restore", cbtSource) - 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] - 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); err != nil { - log.WithError(err).Warnf("Failed to create CBT with source %v, fallback to full restore", 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() @@ -273,6 +279,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/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.go b/pkg/uploader/block/uploader.go index beccb3972..1a08a70fa 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -86,6 +86,12 @@ func (blkup *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, b return udmrepo.Snapshot{}, 0, errors.New("bitmap is not available") } + if bitmap.Errors() != nil { + for _, err := range bitmap.Errors() { + blkup.progress.UpdateProgress(&uploader.Progress{BytesDone: -1, TotalBytes: -1, Message: err.Error()}) + } + } + backupMode := udmrepo.ObjectDataBackupModeInc if parentObject == "" { backupMode = udmrepo.ObjectDataBackupModeFull @@ -153,6 +159,12 @@ func (blkup *blockUploader) Restore(snapshot udmrepo.Snapshot, dest destInfo, bi return 0, 0, errors.New("bitmap is not available") } + 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) if err != nil { return 0, 0, errors.Wrapf(err, "error reading snapshot metadata for %s", snapshot.ID) diff --git a/pkg/uploader/block/uploader_test.go b/pkg/uploader/block/uploader_test.go index bb94eb3de..b3621ed33 100644 --- a/pkg/uploader/block/uploader_test.go +++ b/pkg/uploader/block/uploader_test.go @@ -295,6 +295,7 @@ func TestBlockUploaderBackup(t *testing.T) { var iterator cbt.Iterator if !tc.nilBitmap { iterMock := cbtmocks.NewIterator(t) + iterMock.On("Errors").Return(nil).Maybe() iterator = iterMock backupMode := udmrepo.ObjectDataBackupModeInc @@ -573,6 +574,7 @@ func TestRestoreData(t *testing.T) { reader := bytes.NewReader(data) iterMock := cbtmocks.NewIterator(t) + 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) @@ -603,6 +605,7 @@ func TestRestoreData(t *testing.T) { reader := &errReader{err: errors.New("read error")} iterMock := cbtmocks.NewIterator(t) + 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) @@ -623,6 +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("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") @@ -680,6 +684,7 @@ func TestBlockUploaderRestore(t *testing.T) { } iterMock := cbtmocks.NewIterator(t) + 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) @@ -708,6 +713,7 @@ func TestBlockUploaderRestore(t *testing.T) { } dest := destInfo{size: 4194304, path: "/dev/target"} iterMock := cbtmocks.NewIterator(t) + iterMock.On("Errors").Return(nil).Maybe() _, _, err := blkup.Restore(snap, dest, iterMock, nil) require.Error(t, err) @@ -732,6 +738,7 @@ func TestBlockUploaderRestore(t *testing.T) { } dest := destInfo{size: 512, path: "/dev/small"} iterMock := cbtmocks.NewIterator(t) + 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 f26cb22b7..d861ea6da 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 + cbtErrors []error } type bitmapIterator struct { @@ -43,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, } } @@ -81,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 } @@ -89,6 +93,18 @@ func (c *bitmapImpl) VolumeID() string { return c.volumeID } +func (c *bitmapImpl) SetError(err error) { + if err == nil { + return + } + + c.cbtErrors = append(c.cbtErrors, err) +} + +func (c *bitmapImpl) Errors() []error { + return c.cbtErrors +} + func (c *bitmapImpl) Iterator() types.Iterator { if c.bitmap == nil { return nil @@ -115,3 +131,7 @@ func (c *bitmapIterator) Count() uint64 { func (c *bitmapIterator) BlockSize() uint { return c.blockSize } + +func (c *bitmapIterator) Errors() []error { + return c.cbtErrors +} 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.go b/pkg/uploader/cbt/set.go index 11361cf77..0e50cb048 100644 --- a/pkg/uploader/cbt/set.go +++ b/pkg/uploader/cbt/set.go @@ -26,36 +26,75 @@ 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 + + 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 { + return errors.Wrap(changedErr, "error getting changed blocks from CBT service, fallback to full") + } + + return nil } diff --git a/pkg/uploader/cbt/set_test.go b/pkg/uploader/cbt/set_test.go index 55701dedd..e9a6d1a52 100644 --- a/pkg/uploader/cbt/set_test.go +++ b/pkg/uploader/cbt/set_test.go @@ -21,102 +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 - setupMocks func(*cbtservicemocks.Service, *cbtmocks.Bitmap) + incOnly bool + 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() - }, - expectedErrStr: "CBT service is absent", + 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() - }, - expectedErrStr: "invalid snapshot", + 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: "allocated blocks success", - setupMocks: func(svc *cbtservicemocks.Service, bmp *cbtmocks.Bitmap) { - bmp.On("Snapshot").Return("snap-1") - bmp.On("ChangeID").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", + 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: 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", - 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: mock changed error", + 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 @@ -124,7 +170,10 @@ func TestSetBitmapOrFull(t *testing.T) { svc = svcMock } - err := SetBitmapOrFull(context.Background(), svc, bmpMock) + bmp := NewBitmap(mb, 3*mb, tt.snapshotID, "vol-1") + bmp.SetChangeID(tt.changeID) + + err := SetBitmapOrFull(context.Background(), svc, bmp, tt.incOnly) if tt.expectedErrStr != "" { require.Error(t, err) @@ -136,7 +185,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/cbt/types/mocks/Bitmap.go b/pkg/uploader/cbt/types/mocks/Bitmap.go index faa4ee242..3ed29917b 100644 --- a/pkg/uploader/cbt/types/mocks/Bitmap.go +++ b/pkg/uploader/cbt/types/mocks/Bitmap.go @@ -292,3 +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) +} + +// 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 Errors") + } + + var r0 []error + if returnFunc, ok := ret.Get(0).(func() []error); ok { + r0 = returnFunc() + } else { + 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 eeafd0a59..1ddd59f76 100644 --- a/pkg/uploader/cbt/types/mocks/Iterator.go +++ b/pkg/uploader/cbt/types/mocks/Iterator.go @@ -307,3 +307,20 @@ func (_c *Iterator_VolumeID_Call) RunAndReturn(run func() string) *Iterator_Volu _c.Call.Return(run) return _c } +// 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 Errors") + } + + var r0 []error + if returnFunc, ok := ret.Get(0).(func() []error); ok { + r0 = returnFunc() + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).([]error) + } + } + return r0 +} diff --git a/pkg/uploader/cbt/types/types.go b/pkg/uploader/cbt/types/types.go index 9b13669b2..025d822dc 100644 --- a/pkg/uploader/cbt/types/types.go +++ b/pkg/uploader/cbt/types/types.go @@ -35,6 +35,15 @@ type Bitmap interface { // Iterator returns the iterator for the CBT Bitmap Iterator() Iterator + + // SetError sets CBT error when preparing this bitmap + SetError(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 @@ -56,4 +65,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) + + // Errors returns the CBT errors when preparing this bitmap + Errors() []error } diff --git a/pkg/uploader/kopia/snapshot.go b/pkg/uploader/kopia/snapshot.go index fae7a517c..c4c3251c9 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,29 @@ func SnapshotSource( if parentSnapshot != "" { log.Infof("Using provided parent snapshot %s", parentSnapshot) - mani, err := loadSnapshotFunc(ctx, rep, manifest.ID(parentSnapshot)) - if err != nil { + if mani, err := loadSnapshotFunc(ctx, rep, manifest.ID(parentSnapshot)); err != nil { 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) } } 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) + 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 } - - previous = pre } } else { log.Info("Forcing full snapshot") 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.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/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, 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