From 845bd2dc4f8a8b17aa0c1fc107008d11d7081bac Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Tue, 12 May 2026 17:57:31 +0800 Subject: [PATCH 1/6] block uploader snapshot implementation Signed-off-by: Lyndon-Li --- pkg/uploader/block/dev_linux.go | 30 ++ pkg/uploader/block/dev_other.go | 29 ++ pkg/uploader/block/snapshot.go | 266 ++++++++++++ pkg/uploader/block/snapshot_test.go | 625 ++++++++++++++++++++++++++++ pkg/uploader/block/uploader.go | 54 +++ pkg/uploader/provider/block.go | 83 +++- pkg/uploader/provider/block_test.go | 287 +++++++++++++ pkg/uploader/types.go | 7 +- 8 files changed, 1375 insertions(+), 6 deletions(-) create mode 100644 pkg/uploader/block/dev_linux.go create mode 100644 pkg/uploader/block/dev_other.go create mode 100644 pkg/uploader/block/snapshot.go create mode 100644 pkg/uploader/block/snapshot_test.go create mode 100644 pkg/uploader/block/uploader.go diff --git a/pkg/uploader/block/dev_linux.go b/pkg/uploader/block/dev_linux.go new file mode 100644 index 000000000..4d49442b3 --- /dev/null +++ b/pkg/uploader/block/dev_linux.go @@ -0,0 +1,30 @@ +//go:build linux +// +build linux + +/* +Copyright The Velero Contributors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package block + +import ( + "os" + + "github.com/pkg/errors" +) + +func openBlockDevice(path string, read bool) (*os.File, error) { + return nil, errors.New("Not implemented") +} diff --git a/pkg/uploader/block/dev_other.go b/pkg/uploader/block/dev_other.go new file mode 100644 index 000000000..60689a3d6 --- /dev/null +++ b/pkg/uploader/block/dev_other.go @@ -0,0 +1,29 @@ +//go:build !linux +// +build !linux + +/* +Copyright The Velero Contributors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package block + +import ( + "fmt" + "os" +) + +func openBlockDevice(_ string, _ bool) (*os.File, error) { + return nil, fmt.Errorf("block mode is not supported for Windows") +} diff --git a/pkg/uploader/block/snapshot.go b/pkg/uploader/block/snapshot.go new file mode 100644 index 000000000..272b6dd16 --- /dev/null +++ b/pkg/uploader/block/snapshot.go @@ -0,0 +1,266 @@ +/* +Copyright The Velero Contributors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + +http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package block + +import ( + "context" + "io" + "maps" + "path/filepath" + "time" + + "github.com/pkg/errors" + "github.com/sirupsen/logrus" + "github.com/vmware-tanzu/velero/pkg/cbtservice" + "github.com/vmware-tanzu/velero/pkg/repository/udmrepo" + "github.com/vmware-tanzu/velero/pkg/uploader" + "github.com/vmware-tanzu/velero/pkg/uploader/cbt" +) + +var openBlockDeviceFunc = openBlockDevice + +type parentBackupInfo struct { + parentObject udmrepo.ID + changeID string + volumeID 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) { + if blkup == nil { + return uploader.SnapshotInfo{}, false, errors.New("get empty block uploader") + } + + source, err := filepath.Abs(sourcePath) + if err != nil { + return uploader.SnapshotInfo{}, false, errors.Wrapf(err, "invalid source path %s", sourcePath) + } + + source = filepath.Clean(source) + + sourceInfo := sourceInfo{ + realSource: filepath.Clean(realSource), + } + + if realSource == "" { + sourceInfo.realSource = source + } + + sourceInfo.dev, err = openBlockDeviceFunc(source, true) + if err != nil { + return uploader.SnapshotInfo{}, false, errors.Wrapf(err, "error opening block device %s", source) + } + + sourceInfo.size, err = sourceInfo.dev.Seek(0, io.SeekEnd) + if err != nil { + return uploader.SnapshotInfo{}, false, errors.Wrapf(err, "error getting length of block device %s", source) + } + + _, err = sourceInfo.dev.Seek(0, io.SeekStart) + if err != nil { + return uploader.SnapshotInfo{}, false, errors.Wrapf(err, "error reset pos of block device %s", source) + } + + snapID, backupSize, err := snapshotSource(ctx, repoWriter, blkup, sourceInfo, forceFull, parentSnapshot, cbtSource, cbtservice, tags, uploaderCfg, log, "Block Uploader") + snapshotInfo := uploader.SnapshotInfo{ + ID: snapID, + Size: sourceInfo.size, + IncrementalSize: backupSize, + } + + return snapshotInfo, false, err +} + +func snapshotSource( + ctx context.Context, + rep udmrepo.BackupRepo, + u Uploader, + source sourceInfo, + forceFull bool, + parentSnapshot string, + cbtSource cbtservice.SourceInfo, + cbtservice cbtservice.Service, + snapshotTags map[string]string, + uploaderCfg map[string]string, + log logrus.FieldLogger, + description string, +) (string, int64, error) { + 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, parentBackup.changeID, parentBackup.volumeID) + + err := cbt.SetBitmapOrFull(ctx, cbtservice, bitmap) + if err != nil { + parentBackup.parentObject = "" + log.WithError(err).Warnf("Failed to create CBT with source %v, fallback to real full backup", cbtSource) + } + + snap, backupSize, err := u.Backup(source, parentBackup.parentObject, bitmap.Iterator(), uploaderCfg) + if err != nil { + return "", 0, errors.Wrapf(err, "Failed to run uploader backup for si %v", source) + } + + snap.Tags = make(map[string]string) + snap.Tags[uploader.CBTChangeIDTag] = cbtSource.ChangeID + snap.Tags[uploader.CBTVolumeIDTag] = cbtSource.VolumeID + if snapshotTags != nil { + maps.Copy(snap.Tags, snapshotTags) + } + + snap.Description = description + + snapID, err := rep.SaveSnapshot(ctx, snap) + if err != nil { + return "", 0, errors.Wrapf(err, "Failed to save snapshot %v", snap) + } + + if err = rep.Flush(ctx); err != nil { + return "", 0, errors.Wrapf(err, "Failed to flush repository") + } + + log.Infof("Created snapshot with root %v and ID %v in %v", snap.RootObject, snapID, time.Since(snapshotStartTime).Truncate(time.Second)) + + 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.RootObject.ID) + } + } + } else { + log.Info("Forcing full snapshot") + } + + parentInfo := parentBackupInfo{} + if previous != nil { + if previous.Tags == nil { + log.Warnf("No tag from parent snapshot %s, fallback to full backup", parentSnapshot) + } else if previous.Tags[uploader.CBTChangeIDTag] == "" { + log.Warnf("No ChangeID tag from parent snapshot %s, fallback to full backup", parentSnapshot) + } else if previous.Tags[uploader.CBTVolumeIDTag] == "" { + log.Warnf("No VolumeID tag from parent snapshot %s, fallback to full backup", parentSnapshot) + } 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], parentSnapshot, volumeID) + } else { + parentInfo.parentObject = previous.RootObject.ID + 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", parentSnapshot, previous.StartTime, previous.EndTime, previous.Description) + } + } + + return parentInfo +} + +// Restore restore specific sourcePath with given snapshotID and update progress +func Restore(ctx context.Context, blkup Uploader, rep udmrepo.BackupRepo, snapshotID, dest string, uploaderCfg map[string]string, log logrus.FieldLogger) (int64, error) { + log.Info("Start to restore...") + + snapshot, err := rep.GetSnapshot(ctx, udmrepo.ID(snapshotID)) + if err != nil { + return 0, errors.Wrapf(err, "Unable to load snapshot %v", snapshotID) + } + + log.Infof("Restore from snapshot %s, description %s, created time %v, tags %v", snapshotID, snapshot.Description, snapshot.EndTime, snapshot.Tags) + + destPath, err := filepath.Abs(dest) + if err != nil { + return 0, errors.Wrapf(err, "invalid dest path '%s'", dest) + } + + destPath = filepath.Clean(destPath) + + destDev, err := openBlockDeviceFunc(destPath, false) + if err != nil { + return 0, errors.Wrapf(err, "error opening block device '%s'", destPath) + } + + size, err := blkup.Restore(snapshot, destInfo{dev: destDev, path: destPath}, uploaderCfg) + if err != nil { + return 0, errors.Wrapf(err, "error restoring to block dev %s", destPath) + } + + return size, 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 { + return udmrepo.Snapshot{}, errors.Wrapf(err, "error list snapshots for %s", path) + } + + var previous *udmrepo.Snapshot + + for _, snap := range snaps { + log.Debugf("Found one snapshot %s, start time %v, tags %v", snap.RootObject.ID, snap.StartTime, snap.Tags) + + requester, found := snap.Tags[uploader.SnapshotRequesterTag] + if !found { + continue + } + + if requester != snapshotTags[uploader.SnapshotRequesterTag] { + continue + } + + uploaderName, found := snap.Tags[uploader.SnapshotUploaderTag] + if !found { + continue + } + + if uploaderName != uploader.BlockType { + continue + } + + if noLaterThan != nil && snap.StartTime.After(*noLaterThan) { + continue + } + + if previous == nil || snap.StartTime.After(previous.StartTime) { + previous = &snap + } + } + + if previous == nil { + return udmrepo.Snapshot{}, errors.Errorf("no matching snapshot found for source %s", path) + } + + return *previous, nil +} diff --git a/pkg/uploader/block/snapshot_test.go b/pkg/uploader/block/snapshot_test.go new file mode 100644 index 000000000..1e609eb2f --- /dev/null +++ b/pkg/uploader/block/snapshot_test.go @@ -0,0 +1,625 @@ +/* +Copyright The Velero Contributors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + +http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Tests live in package block (not block_test) so they can access unexported +// types sourceInfo and destInfo, which appear in the Uploader interface. +package block + +import ( + "context" + "os" + "testing" + "time" + + "github.com/pkg/errors" + "github.com/sirupsen/logrus" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" + + "github.com/vmware-tanzu/velero/pkg/cbtservice" + "github.com/vmware-tanzu/velero/pkg/repository/udmrepo" + udmrepomocks "github.com/vmware-tanzu/velero/pkg/repository/udmrepo/mocks" + "github.com/vmware-tanzu/velero/pkg/uploader" + cbttypes "github.com/vmware-tanzu/velero/pkg/uploader/cbt/types" +) + +type mockUploader struct { + mock.Mock +} + +func (m *mockUploader) Backup(src sourceInfo, parent udmrepo.ID, iter cbttypes.Iterator, cfg map[string]string) (udmrepo.Snapshot, int64, error) { + args := m.Called(src, parent, iter, cfg) + return args.Get(0).(udmrepo.Snapshot), args.Get(1).(int64), args.Error(2) +} + +func (m *mockUploader) Restore(snap udmrepo.Snapshot, dest destInfo, cfg map[string]string) (int64, error) { + args := m.Called(snap, dest, cfg) + return args.Get(0).(int64), args.Error(1) +} + +func testLog() logrus.FieldLogger { + l := logrus.New() + l.SetLevel(logrus.DebugLevel) + return l +} + +func tempFile(t *testing.T, content string) *os.File { + t.Helper() + f, err := os.CreateTemp("", "blktest-*") + require.NoError(t, err) + if content != "" { + _, err = f.WriteString(content) + require.NoError(t, err) + } + t.Cleanup(func() { + f.Close() + os.Remove(f.Name()) + }) + return f +} + +func TestBackup(t *testing.T) { + testCases := []struct { + name string + useNilBlkup bool + setupOpenDev func(t *testing.T) *os.File + setupMocks func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) + expectedErrStr string + checkInfo func(*testing.T, uploader.SnapshotInfo) + }{ + { + name: "nil uploader returns error", + useNilBlkup: true, + expectedErrStr: "get empty block uploader", + }, + { + name: "openBlockDevice error", + expectedErrStr: "error opening block device", + }, + { + name: "SnapshotSource error propagates", + setupOpenDev: func(t *testing.T) *os.File { + return tempFile(t, "") + }, + setupMocks: func(blkup *mockUploader, _ *udmrepomocks.BackupRepo) { + blkup.On("Backup", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(udmrepo.Snapshot{}, int64(0), errors.New("I/O error")) + }, + expectedErrStr: "Failed to run uploader backup", + }, + { + name: "success returns correct SnapshotInfo", + setupOpenDev: func(t *testing.T) *os.File { + return tempFile(t, "test-block-data") + }, + setupMocks: func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) { + blkup.On("Backup", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(udmrepo.Snapshot{RootObject: udmrepo.ObjectMetadata{ID: "root"}}, int64(8), nil) + repo.On("SaveSnapshot", mock.Anything, mock.Anything).Return(udmrepo.ID("snap-001"), nil) + repo.On("Flush", mock.Anything).Return(nil) + }, + checkInfo: func(t *testing.T, info uploader.SnapshotInfo) { + assert.Equal(t, "snap-001", info.ID) + assert.Equal(t, int64(8), info.IncrementalSize) + assert.Greater(t, info.Size, int64(0)) + }, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ctx := context.Background() + mockBlkup := &mockUploader{} + mockRepo := udmrepomocks.NewBackupRepo(t) + + var blkup Uploader + if !tc.useNilBlkup { + blkup = mockBlkup + } + + if tc.setupOpenDev != nil { + f := tc.setupOpenDev(t) + openBlockDeviceFunc = func(_ string, _ bool) (*os.File, error) { + return f, nil + } + } else { + openBlockDeviceFunc = func(_ string, _ bool) (*os.File, error) { + return nil, errors.New("device not available") + } + } + + if tc.setupMocks != nil { + tc.setupMocks(mockBlkup, mockRepo) + } + + info, isEmpty, err := Backup( + ctx, blkup, mockRepo, + "/dev/sda", "", + cbtservice.SourceInfo{}, + true, "", nil, + map[string]string{}, map[string]string{}, + testLog(), + ) + + if tc.expectedErrStr != "" { + require.Error(t, err) + assert.ErrorContains(t, err, tc.expectedErrStr) + } else { + require.NoError(t, err) + assert.False(t, isEmpty) + } + + if tc.checkInfo != nil { + tc.checkInfo(t, info) + } + + mockBlkup.AssertExpectations(t) + }) + } +} + +func TestSnapshotSource(t *testing.T) { + baseSource := sourceInfo{realSource: "/test/vol", size: 1024} + + testCases := []struct { + name string + setupMocks func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) + expectedErrStr string + expectedSnapID string + expectedSize int64 + }{ + { + name: "uploader Backup error", + setupMocks: func(blkup *mockUploader, _ *udmrepomocks.BackupRepo) { + blkup.On("Backup", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(udmrepo.Snapshot{}, int64(0), errors.New("uploader error")) + }, + expectedErrStr: "Failed to run uploader backup", + }, + { + name: "SaveSnapshot error", + setupMocks: func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) { + blkup.On("Backup", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(udmrepo.Snapshot{}, int64(0), nil) + repo.On("SaveSnapshot", mock.Anything, mock.Anything). + Return(udmrepo.ID(""), errors.New("save failed")) + }, + expectedErrStr: "Failed to save snapshot", + }, + { + name: "Flush error", + setupMocks: func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) { + blkup.On("Backup", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(udmrepo.Snapshot{}, int64(0), nil) + repo.On("SaveSnapshot", mock.Anything, mock.Anything).Return(udmrepo.ID("snap-001"), nil) + repo.On("Flush", mock.Anything).Return(errors.New("flush failed")) + }, + expectedErrStr: "Failed to flush repository", + }, + { + name: "success with nil cbtService falls back to full bitmap", + setupMocks: func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) { + blkup.On("Backup", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(udmrepo.Snapshot{RootObject: udmrepo.ObjectMetadata{ID: "root"}}, int64(512), nil) + repo.On("SaveSnapshot", mock.Anything, mock.Anything).Return(udmrepo.ID("snap-success"), nil) + repo.On("Flush", mock.Anything).Return(nil) + }, + expectedSnapID: "snap-success", + expectedSize: 512, + }, + { + name: "tags from cbtSource and snapshotTags are merged onto snapshot", + setupMocks: func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) { + blkup.On("Backup", mock.Anything, mock.Anything, mock.Anything, mock.Anything). + Return(udmrepo.Snapshot{}, int64(0), nil) + repo.On("SaveSnapshot", mock.Anything, mock.MatchedBy(func(snap udmrepo.Snapshot) bool { + return snap.Tags[uploader.CBTChangeIDTag] == "cid-1" && + snap.Tags[uploader.CBTVolumeIDTag] == "vid-1" && + snap.Tags["custom"] == "val" && + snap.Description == "Block Uploader" + })).Return(udmrepo.ID("snap-tags"), nil) + repo.On("Flush", mock.Anything).Return(nil) + }, + expectedSnapID: "snap-tags", + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ctx := context.Background() + mockBlkup := &mockUploader{} + mockRepo := udmrepomocks.NewBackupRepo(t) + + tc.setupMocks(mockBlkup, mockRepo) + + cbtSrc := cbtservice.SourceInfo{ChangeID: "cid-1", VolumeID: "vid-1"} + snapshotTags := map[string]string{"custom": "val"} + + snapID, size, err := snapshotSource( + ctx, mockRepo, mockBlkup, + baseSource, + true, "", + cbtSrc, nil, + snapshotTags, map[string]string{}, + testLog(), "Block Uploader", + ) + + if tc.expectedErrStr != "" { + require.Error(t, err) + assert.ErrorContains(t, err, tc.expectedErrStr) + } else { + require.NoError(t, err) + assert.Equal(t, tc.expectedSnapID, snapID) + assert.Equal(t, tc.expectedSize, size) + } + + mockBlkup.AssertExpectations(t) + }) + } +} + +func TestGetParentBackupInfo(t *testing.T) { + const volumeID = "vol-123" + const realSource = "/test/source" + + snapshotTags := map[string]string{ + uploader.SnapshotRequesterTag: "test-requester", + uploader.SnapshotUploaderTag: uploader.BlockType, + } + + validSnap := udmrepo.Snapshot{ + RootObject: udmrepo.ObjectMetadata{ID: "root-obj"}, + Tags: map[string]string{ + uploader.CBTChangeIDTag: "cid-abc", + uploader.CBTVolumeIDTag: volumeID, + uploader.SnapshotRequesterTag: "test-requester", + uploader.SnapshotUploaderTag: uploader.BlockType, + }, + } + + testCases := []struct { + name string + forceFull bool + parentSnapshot string + setupMocks func(repo *udmrepomocks.BackupRepo) + expectEmpty bool + expectedParent udmrepo.ID + expectedCID string + expectedVID string + }{ + { + name: "forceFull skips all parent lookup", + forceFull: true, + expectEmpty: true, + }, + { + name: "GetSnapshot fails — falls back to full", + parentSnapshot: "snap-parent", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-parent")). + Return(udmrepo.Snapshot{}, errors.New("not found")) + }, + expectEmpty: true, + }, + { + 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) + }, + expectEmpty: true, + }, + { + 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) + }, + expectEmpty: true, + }, + { + 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) + }, + expectEmpty: true, + }, + { + 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{ + uploader.CBTChangeIDTag: "cid", + uploader.CBTVolumeIDTag: "different-vol", + }}, nil) + }, + expectEmpty: true, + }, + { + name: "valid parent snapshot — returns parent info", + parentSnapshot: "snap-valid", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-valid")). + Return(validSnap, nil) + }, + expectedParent: "root-obj", + expectedCID: "cid-abc", + expectedVID: volumeID, + }, + { + name: "no parentSnapshot — ListSnapshot fails — falls back to full", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ListSnapshot", mock.Anything, realSource). + Return(nil, errors.New("list error")) + }, + expectEmpty: true, + }, + { + name: "no parentSnapshot — no matching snapshot — falls back to full", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ListSnapshot", mock.Anything, realSource). + Return([]udmrepo.Snapshot{{Tags: map[string]string{"other": "tag"}}}, nil) + }, + expectEmpty: true, + }, + { + name: "no parentSnapshot — matching snapshot found — returns parent info", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ListSnapshot", mock.Anything, realSource). + Return([]udmrepo.Snapshot{validSnap}, nil) + }, + expectedParent: "root-obj", + expectedCID: "cid-abc", + expectedVID: volumeID, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ctx := context.Background() + mockRepo := udmrepomocks.NewBackupRepo(t) + + if tc.setupMocks != nil { + tc.setupMocks(mockRepo) + } + + info := getParentBackupInfo(ctx, mockRepo, tc.forceFull, tc.parentSnapshot, volumeID, realSource, snapshotTags, testLog()) + + if tc.expectEmpty { + assert.Empty(t, info.parentObject) + assert.Empty(t, info.changeID) + assert.Empty(t, info.volumeID) + } else { + assert.Equal(t, tc.expectedParent, info.parentObject) + assert.Equal(t, tc.expectedCID, info.changeID) + assert.Equal(t, tc.expectedVID, info.volumeID) + } + }) + } +} + +func TestFindPreviousSnapshot(t *testing.T) { + snapshotTags := map[string]string{ + uploader.SnapshotRequesterTag: "test-requester", + uploader.SnapshotUploaderTag: uploader.BlockType, + } + + matchingSnap := func(id string, start time.Time) udmrepo.Snapshot { + return udmrepo.Snapshot{ + RootObject: udmrepo.ObjectMetadata{ID: udmrepo.ID(id)}, + StartTime: start, + Tags: map[string]string{ + uploader.SnapshotRequesterTag: "test-requester", + uploader.SnapshotUploaderTag: uploader.BlockType, + }, + } + } + + testCases := []struct { + name string + setupMocks func(repo *udmrepomocks.BackupRepo) + expectedErrStr string + expectedID string + }{ + { + name: "ListSnapshot error", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ListSnapshot", mock.Anything, "source"). + Return(nil, errors.New("list error")) + }, + expectedErrStr: "error list snapshots", + }, + { + name: "empty snapshot list — no match", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ListSnapshot", mock.Anything, "source"). + Return([]udmrepo.Snapshot{}, nil) + }, + expectedErrStr: "no matching snapshot found", + }, + { + name: "snapshots without matching tags are filtered", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ListSnapshot", mock.Anything, "source"). + Return([]udmrepo.Snapshot{ + {Tags: map[string]string{"unrelated": "tag"}}, + {Tags: nil}, + }, nil) + }, + expectedErrStr: "no matching snapshot found", + }, + { + name: "snapshot with wrong requester tag is filtered", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ListSnapshot", mock.Anything, "source"). + Return([]udmrepo.Snapshot{{ + Tags: map[string]string{ + uploader.SnapshotRequesterTag: "other-requester", + uploader.SnapshotUploaderTag: uploader.BlockType, + }, + }}, nil) + }, + expectedErrStr: "no matching snapshot found", + }, + { + name: "snapshot with wrong uploader tag is filtered", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ListSnapshot", mock.Anything, "source"). + Return([]udmrepo.Snapshot{{ + Tags: map[string]string{ + uploader.SnapshotRequesterTag: "test-requester", + uploader.SnapshotUploaderTag: "kopia", + }, + }}, nil) + }, + expectedErrStr: "no matching snapshot found", + }, + { + name: "single matching snapshot is returned", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ListSnapshot", mock.Anything, "source"). + Return([]udmrepo.Snapshot{matchingSnap("snap-a", time.Now())}, nil) + }, + expectedID: "snap-a", + }, + { + name: "most recent of multiple matching snapshots is returned", + setupMocks: func(repo *udmrepomocks.BackupRepo) { + now := time.Now() + repo.On("ListSnapshot", mock.Anything, "source"). + Return([]udmrepo.Snapshot{ + matchingSnap("snap-old", now.Add(-2*time.Hour)), + matchingSnap("snap-new", now.Add(-time.Minute)), + matchingSnap("snap-mid", now.Add(-time.Hour)), + }, nil) + }, + expectedID: "snap-new", + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ctx := context.Background() + mockRepo := udmrepomocks.NewBackupRepo(t) + tc.setupMocks(mockRepo) + + snap, err := findPreviousSnapshot(ctx, mockRepo, "source", snapshotTags, nil, testLog()) + + if tc.expectedErrStr != "" { + require.Error(t, err) + assert.ErrorContains(t, err, tc.expectedErrStr) + } else { + require.NoError(t, err) + assert.Equal(t, udmrepo.ID(tc.expectedID), snap.RootObject.ID) + } + }) + } +} + +func TestRestore(t *testing.T) { + storedSnap := udmrepo.Snapshot{Description: "test snapshot"} + + testCases := []struct { + name string + setupMocks func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) + setupOpenDev func(t *testing.T) *os.File + expectedErrStr string + expectedSize int64 + }{ + { + name: "GetSnapshot error", + setupMocks: func(_ *mockUploader, repo *udmrepomocks.BackupRepo) { + repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-001")). + Return(udmrepo.Snapshot{}, errors.New("not found")) + }, + expectedErrStr: "Unable to load snapshot", + }, + { + name: "openBlockDevice error", + setupMocks: func(_ *mockUploader, repo *udmrepomocks.BackupRepo) { + repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-001")). + Return(storedSnap, nil) + }, + expectedErrStr: "error opening block device", + }, + { + name: "Restore error", + setupMocks: func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) { + repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-001")). + Return(storedSnap, nil) + blkup.On("Restore", mock.Anything, mock.Anything, mock.Anything). + Return(int64(0), errors.New("restore I/O error")) + }, + setupOpenDev: func(t *testing.T) *os.File { + return tempFile(t, "") + }, + expectedErrStr: "error restoring to block dev", + }, + { + name: "success returns size", + setupMocks: func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) { + repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-001")). + Return(storedSnap, nil) + blkup.On("Restore", mock.Anything, mock.Anything, mock.Anything). + Return(int64(4096), nil) + }, + setupOpenDev: func(t *testing.T) *os.File { + return tempFile(t, "") + }, + expectedSize: 4096, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ctx := context.Background() + mockBlkup := &mockUploader{} + mockRepo := udmrepomocks.NewBackupRepo(t) + + tc.setupMocks(mockBlkup, mockRepo) + + if tc.setupOpenDev != nil { + f := tc.setupOpenDev(t) + openBlockDeviceFunc = func(_ string, _ bool) (*os.File, error) { + return f, nil + } + } else { + openBlockDeviceFunc = func(_ string, _ bool) (*os.File, error) { + return nil, errors.New("device not available") + } + } + + size, err := Restore(ctx, mockBlkup, mockRepo, "snap-001", "/dev/sdb", map[string]string{}, testLog()) + + if tc.expectedErrStr != "" { + require.Error(t, err) + assert.ErrorContains(t, err, tc.expectedErrStr) + assert.Equal(t, int64(0), size) + } else { + require.NoError(t, err) + assert.Equal(t, tc.expectedSize, size) + } + + mockBlkup.AssertExpectations(t) + }) + } +} diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go new file mode 100644 index 000000000..118a09713 --- /dev/null +++ b/pkg/uploader/block/uploader.go @@ -0,0 +1,54 @@ +/* +Copyright The Velero Contributors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + +http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package block + +import ( + "context" + "os" + + "github.com/pkg/errors" + "github.com/sirupsen/logrus" + "github.com/vmware-tanzu/velero/pkg/repository/udmrepo" + "github.com/vmware-tanzu/velero/pkg/uploader" + cbt "github.com/vmware-tanzu/velero/pkg/uploader/cbt/types" +) + +var ErrCanceled = errors.New("uploader is canceled") + +const ( + blockSize = (1 << 20) +) + +type sourceInfo struct { + dev *os.File + realSource string + size int64 +} + +type destInfo struct { + dev *os.File + path string +} + +type Uploader interface { + Backup(sourceInfo, udmrepo.ID, cbt.Iterator, map[string]string) (udmrepo.Snapshot, int64, error) + Restore(udmrepo.Snapshot, destInfo, map[string]string) (int64, error) +} + +func NewUploader(ctx context.Context, repoWriter udmrepo.BackupRepo, progress uploader.ProgressUpdater, log logrus.FieldLogger) Uploader { + return nil +} diff --git a/pkg/uploader/provider/block.go b/pkg/uploader/provider/block.go index dc5028040..427d3fae3 100644 --- a/pkg/uploader/provider/block.go +++ b/pkg/uploader/provider/block.go @@ -18,6 +18,7 @@ package provider import ( "context" + "fmt" "strings" "github.com/cockroachdb/errors" @@ -28,8 +29,12 @@ import ( repokeys "github.com/vmware-tanzu/velero/pkg/repository/keys" "github.com/vmware-tanzu/velero/pkg/repository/udmrepo" "github.com/vmware-tanzu/velero/pkg/uploader" + "github.com/vmware-tanzu/velero/pkg/uploader/block" ) +var blockBackupFunc = block.Backup +var blockRestoreFunc = block.Restore + type blockProvider struct { requestorType string bkRepo udmrepo.BackupRepo @@ -88,7 +93,6 @@ func (bp *blockProvider) GetPassword(param any) (string, error) { return strings.TrimSpace(rawPass), nil } -// TODO: implement in the following PRs func (bp *blockProvider) RunBackup( ctx context.Context, path string, @@ -100,10 +104,55 @@ func (bp *blockProvider) RunBackup( volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, updater uploader.ProgressUpdater) (string, bool, int64, int64, error) { - return "", false, 0, 0, errors.New("block backup not implemented") + if updater == nil { + return "", false, 0, 0, errors.New("Need to initial backup progress updater first") + } + + if path == "" { + return "", false, 0, 0, errors.New("path is empty") + } + + log := bp.log.WithFields(logrus.Fields{ + "path": path, + "realSource": realSource, + "parentSnapshot": parentSnapshot, + }) + + blkUploader := block.NewUploader(ctx, bp.bkRepo, updater, log) + + if tags == nil { + tags = make(map[string]string) + } + tags[uploader.SnapshotRequesterTag] = bp.requestorType + tags[uploader.SnapshotUploaderTag] = uploader.BlockType + + if realSource != "" { + realSource = fmt.Sprintf("%s/%s/%s", bp.requestorType, uploader.BlockType, realSource) + } + + snapshotInfo, _, err := blockBackupFunc(ctx, blkUploader, bp.bkRepo, path, realSource, cbtParam.Source, forceFull, parentSnapshot, cbtParam.Service, uploaderCfg, tags, log) + + if err == block.ErrCanceled { + log.Warn("Block backup is canceled") + return snapshotInfo.ID, false, snapshotInfo.Size, snapshotInfo.IncrementalSize, ErrorCanceled + } + + if err != nil { + return snapshotInfo.ID, false, snapshotInfo.Size, snapshotInfo.IncrementalSize, errors.Wrapf(err, "Failed to run block backup") + } + + updater.UpdateProgress( + &uploader.Progress{ + TotalBytes: snapshotInfo.Size, + BytesDone: snapshotInfo.Size, + }, + ) + + log.Infof("Block backup finished, snapshot ID %s, backup size %d", snapshotInfo.ID, snapshotInfo.Size) + + return snapshotInfo.ID, false, snapshotInfo.Size, snapshotInfo.IncrementalSize, nil } -// TODO: implement in the following PRs func (bp *blockProvider) RunRestore( ctx context.Context, snapshotID string, @@ -111,5 +160,31 @@ func (bp *blockProvider) RunRestore( volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, updater uploader.ProgressUpdater) (int64, error) { - return 0, errors.New("block restore not implemented") + log := bp.log.WithFields(logrus.Fields{ + "snapshotID": snapshotID, + "volumePath": volumePath, + }) + log.Info("Starting restore") + + blkUploader := block.NewUploader(ctx, bp.bkRepo, updater, log) + + size, err := blockRestoreFunc(ctx, blkUploader, bp.bkRepo, snapshotID, volumePath, uploaderCfg, log) + + if err == block.ErrCanceled { + log.Warn("Block restore is canceled") + return 0, ErrorCanceled + } + + if err != nil { + return 0, errors.Wrapf(err, "Failed to run block restore") + } + + updater.UpdateProgress(&uploader.Progress{ + TotalBytes: size, + BytesDone: size, + }) + + log.Infof("Block restore finished, restore size %v", size) + + return size, nil } diff --git a/pkg/uploader/provider/block_test.go b/pkg/uploader/provider/block_test.go index 1c180513e..e7af93855 100644 --- a/pkg/uploader/provider/block_test.go +++ b/pkg/uploader/provider/block_test.go @@ -17,6 +17,7 @@ limitations under the License. package provider import ( + "context" "testing" "github.com/cockroachdb/errors" @@ -29,9 +30,12 @@ import ( "github.com/vmware-tanzu/velero/internal/credentials" "github.com/vmware-tanzu/velero/internal/credentials/mocks" velerov1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v1" + "github.com/vmware-tanzu/velero/pkg/cbtservice" "github.com/vmware-tanzu/velero/pkg/repository" "github.com/vmware-tanzu/velero/pkg/repository/udmrepo" udmrepomocks "github.com/vmware-tanzu/velero/pkg/repository/udmrepo/mocks" + "github.com/vmware-tanzu/velero/pkg/uploader" + "github.com/vmware-tanzu/velero/pkg/uploader/block" ) func TestNewBlockUploaderProvider(t *testing.T) { @@ -125,6 +129,16 @@ func TestBlockProviderClose(t *testing.T) { mockBRepo.AssertExpectations(t) } +type blockMockProgressUpdater struct { + lastProgress *uploader.Progress + callCount int +} + +func (u *blockMockProgressUpdater) UpdateProgress(p *uploader.Progress) { + u.lastProgress = p + u.callCount++ +} + func TestBlockProviderGetPassword(t *testing.T) { testCases := []struct { name string @@ -185,3 +199,276 @@ func TestBlockProviderGetPassword(t *testing.T) { }) } } + +func TestBlockProviderRunBackup(t *testing.T) { + const requestorType = "test-requestor" + + testCases := []struct { + name string + path string + realSource string + tags map[string]string + updater uploader.ProgressUpdater + mockBackupResult uploader.SnapshotInfo + mockBackupErr error + expectedID string + expectedSize int64 + expectedIncrSize int64 + expectError bool + expectedErrStr string + skipMock bool + checkCaptures func(*testing.T, string, map[string]string) + }{ + { + name: "nil updater returns error", + path: "/dev/sda", + updater: nil, + expectError: true, + expectedErrStr: "Need to initial backup progress updater first", + skipMock: true, + }, + { + name: "empty path returns error", + path: "", + updater: &FakeBackupProgressUpdater{}, + expectError: true, + expectedErrStr: "path is empty", + skipMock: true, + }, + { + name: "success returns correct snapshot info and updates progress", + path: "/dev/sda", + updater: &blockMockProgressUpdater{}, + mockBackupResult: uploader.SnapshotInfo{ + ID: "snap-001", + Size: 1024, + IncrementalSize: 512, + }, + expectedID: "snap-001", + expectedSize: 1024, + expectedIncrSize: 512, + }, + { + name: "canceled backup returns ErrorCanceled with partial snapshot info", + path: "/dev/sda", + updater: &FakeBackupProgressUpdater{}, + mockBackupResult: uploader.SnapshotInfo{ + ID: "snap-canceled", + Size: 2048, + IncrementalSize: 1024, + }, + mockBackupErr: block.ErrCanceled, + expectedID: "snap-canceled", + expectedSize: 2048, + expectedIncrSize: 1024, + expectError: true, + expectedErrStr: "uploader is canceled", + }, + { + name: "generic backup error is wrapped", + path: "/dev/sda", + updater: &FakeBackupProgressUpdater{}, + mockBackupErr: errors.New("disk I/O error"), + expectError: true, + expectedErrStr: "Failed to run block backup", + }, + { + name: "nil tags are initialized with required tags", + path: "/dev/sda", + tags: nil, + updater: &FakeBackupProgressUpdater{}, + mockBackupResult: uploader.SnapshotInfo{ID: "snap-tags"}, + expectedID: "snap-tags", + checkCaptures: func(t *testing.T, _ string, tags map[string]string) { + assert.Equal(t, requestorType, tags[uploader.SnapshotRequesterTag]) + assert.Equal(t, uploader.BlockType, tags[uploader.SnapshotUploaderTag]) + }, + }, + { + name: "non-empty realSource is prefixed with requestorType and BlockType", + path: "/dev/sda", + realSource: "my-volume", + updater: &FakeBackupProgressUpdater{}, + mockBackupResult: uploader.SnapshotInfo{ID: "snap-source"}, + expectedID: "snap-source", + checkCaptures: func(t *testing.T, realSource string, _ map[string]string) { + assert.Equal(t, requestorType+"/"+uploader.BlockType+"/my-volume", realSource) + }, + }, + { + name: "empty realSource is passed through unchanged", + path: "/dev/sda", + realSource: "", + updater: &FakeBackupProgressUpdater{}, + mockBackupResult: uploader.SnapshotInfo{ID: "snap-nosource"}, + expectedID: "snap-nosource", + checkCaptures: func(t *testing.T, realSource string, _ map[string]string) { + assert.Equal(t, "", realSource) + }, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + mockBRepo := udmrepomocks.NewBackupRepo(t) + + var capturedRealSrc string + var capturedTags map[string]string + + if !tc.skipMock { + blockBackupFunc = func(_ context.Context, _ block.Uploader, _ udmrepo.BackupRepo, _ string, realSource string, _ cbtservice.SourceInfo, _ bool, _ string, _ cbtservice.Service, _ map[string]string, tags map[string]string, _ logrus.FieldLogger) (uploader.SnapshotInfo, bool, error) { + capturedRealSrc = realSource + capturedTags = tags + return tc.mockBackupResult, false, tc.mockBackupErr + } + } + + bp := &blockProvider{ + requestorType: requestorType, + bkRepo: mockBRepo, + log: logrus.New(), + } + + snapshotID, isEmpty, size, incrSize, err := bp.RunBackup( + t.Context(), + tc.path, + tc.realSource, + tc.tags, + false, + "", + CBTParam{}, + uploader.PersistentVolumeBlock, + map[string]string{}, + tc.updater, + ) + + assert.Equal(t, tc.expectedID, snapshotID) + assert.Equal(t, tc.expectedSize, size) + assert.Equal(t, tc.expectedIncrSize, incrSize) + + if tc.expectError { + require.Error(t, err) + if tc.expectedErrStr != "" { + assert.ErrorContains(t, err, tc.expectedErrStr) + } + } else { + require.NoError(t, err) + assert.False(t, isEmpty) + if mu, ok := tc.updater.(*blockMockProgressUpdater); ok { + assert.Equal(t, 1, mu.callCount) + require.NotNil(t, mu.lastProgress) + assert.Equal(t, tc.expectedSize, mu.lastProgress.TotalBytes) + assert.Equal(t, tc.expectedSize, mu.lastProgress.BytesDone) + } + } + + if tc.checkCaptures != nil { + tc.checkCaptures(t, capturedRealSrc, capturedTags) + } + }) + } +} + +func TestBlockProviderRunRestore(t *testing.T) { + testCases := []struct { + name string + snapshotID string + volumePath string + updater uploader.ProgressUpdater + mockRestoreSize int64 + mockRestoreErr error + expectedSize int64 + expectError bool + expectedErrStr string + checkCaptures func(*testing.T, string, string) + }{ + { + name: "success returns size and updates progress", + snapshotID: "snap-001", + volumePath: "/dev/sdb", + updater: &blockMockProgressUpdater{}, + mockRestoreSize: 4096, + expectedSize: 4096, + }, + { + name: "canceled restore returns ErrorCanceled", + snapshotID: "snap-canceled", + volumePath: "/dev/sdb", + updater: &FakeRestoreProgressUpdater{}, + mockRestoreErr: block.ErrCanceled, + expectError: true, + expectedErrStr: "uploader is canceled", + }, + { + name: "generic restore error is wrapped", + snapshotID: "snap-error", + volumePath: "/dev/sdb", + updater: &FakeRestoreProgressUpdater{}, + mockRestoreErr: errors.New("disk read error"), + expectError: true, + expectedErrStr: "Failed to run block restore", + }, + { + name: "snapshotID and volumePath are forwarded to restore func", + snapshotID: "snap-fwd", + volumePath: "/dev/sdc", + updater: &FakeRestoreProgressUpdater{}, + mockRestoreSize: 512, + expectedSize: 512, + checkCaptures: func(t *testing.T, snapshotID, volumePath string) { + assert.Equal(t, "snap-fwd", snapshotID) + assert.Equal(t, "/dev/sdc", volumePath) + }, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + mockBRepo := udmrepomocks.NewBackupRepo(t) + + var capturedSnapshotID string + var capturedVolumePath string + + blockRestoreFunc = func(_ context.Context, _ block.Uploader, _ udmrepo.BackupRepo, snapshotID string, volumePath string, _ map[string]string, _ logrus.FieldLogger) (int64, error) { + capturedSnapshotID = snapshotID + capturedVolumePath = volumePath + return tc.mockRestoreSize, tc.mockRestoreErr + } + + bp := &blockProvider{ + bkRepo: mockBRepo, + log: logrus.New(), + } + + size, err := bp.RunRestore( + t.Context(), + tc.snapshotID, + tc.volumePath, + uploader.PersistentVolumeBlock, + map[string]string{}, + tc.updater, + ) + + if tc.expectError { + require.Error(t, err) + if tc.expectedErrStr != "" { + assert.ErrorContains(t, err, tc.expectedErrStr) + } + assert.Equal(t, int64(0), size) + } else { + require.NoError(t, err) + assert.Equal(t, tc.expectedSize, size) + if mu, ok := tc.updater.(*blockMockProgressUpdater); ok { + assert.Equal(t, 1, mu.callCount) + require.NotNil(t, mu.lastProgress) + assert.Equal(t, tc.expectedSize, mu.lastProgress.TotalBytes) + assert.Equal(t, tc.expectedSize, mu.lastProgress.BytesDone) + } + } + + if tc.checkCaptures != nil { + tc.checkCaptures(t, capturedSnapshotID, capturedVolumePath) + } + }) + } +} diff --git a/pkg/uploader/types.go b/pkg/uploader/types.go index 12ff1dc52..9c700193f 100644 --- a/pkg/uploader/types.go +++ b/pkg/uploader/types.go @@ -26,6 +26,8 @@ const ( BlockType = "velero-block" SnapshotRequesterTag = "snapshot-requester" SnapshotUploaderTag = "snapshot-uploader" + CBTChangeIDTag = "cbt-change-id" + CBTVolumeIDTag = "cbt-volume-id" ) type PersistentVolumeMode string @@ -49,8 +51,9 @@ func ValidateUploaderType(t string) (string, error) { } type SnapshotInfo struct { - ID string `json:"id"` - Size int64 `json:"Size"` + ID string + Size int64 + IncrementalSize int64 } // Progress which defined two variables to record progress From f4f897f669ae10065599c4f2454c52766db2d27e Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Fri, 22 May 2026 18:19:02 +0800 Subject: [PATCH 2/6] load object from snapshot Signed-off-by: Lyndon-Li --- pkg/uploader/block/snapshot.go | 4 +++- pkg/uploader/block/uploader.go | 17 +++++++++++++++++ 2 files changed, 20 insertions(+), 1 deletion(-) diff --git a/pkg/uploader/block/snapshot.go b/pkg/uploader/block/snapshot.go index 272b6dd16..a3bb92431 100644 --- a/pkg/uploader/block/snapshot.go +++ b/pkg/uploader/block/snapshot.go @@ -177,8 +177,10 @@ func getParentBackupInfo(ctx context.Context, rep udmrepo.BackupRepo, forceFull log.Warnf("No VolumeID tag from parent snapshot %s, fallback to full backup", parentSnapshot) } 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], parentSnapshot, 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", parentSnapshot) } else { - parentInfo.parentObject = previous.RootObject.ID + parentInfo.parentObject = obj parentInfo.changeID = previous.Tags[uploader.CBTChangeIDTag] parentInfo.volumeID = previous.Tags[uploader.CBTVolumeIDTag] diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index 118a09713..f487b39cb 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -52,3 +52,20 @@ type Uploader interface { func NewUploader(ctx context.Context, repoWriter udmrepo.BackupRepo, progress uploader.ProgressUpdater, log logrus.FieldLogger) Uploader { return nil } + +func loadObjectFromSnapshot(ctx context.Context, rep udmrepo.BackupRepo, snapshot *udmrepo.Snapshot) (udmrepo.ID, error) { + if snapshot == nil { + return "", errors.New("snapshot is empty") + } + + parentMeta, err := rep.ReadMetadata(ctx, snapshot.RootObject.ID) + if err != nil { + return "", errors.Wrapf(err, "error readding snapshot metadata for %s", snapshot.Description) + } + + if len(parentMeta.SubObjects) != 1 { + return "", errors.Wrapf(err, "unexpected number of bdev object (%d) for snapshot %s", len(parentMeta.SubObjects), snapshot.Description) + } + + return parentMeta.SubObjects[0].ID, nil +} From 5bef38dc9578221f5eb37e3dd045ba25d6171142 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Thu, 25 Jun 2026 15:42:08 +0800 Subject: [PATCH 3/6] block uploader snapshot implementation Signed-off-by: Lyndon-Li --- changelogs/unreleased/9945-Lyndon-Li | 1 + pkg/uploader/block/dev_linux.go | 1 + pkg/uploader/block/snapshot.go | 6 +- pkg/uploader/block/snapshot_test.go | 12 ++- pkg/uploader/block/uploader.go | 8 +- pkg/uploader/block/uploader_test.go | 112 +++++++++++++++++++++++++++ 6 files changed, 132 insertions(+), 8 deletions(-) create mode 100644 changelogs/unreleased/9945-Lyndon-Li create mode 100644 pkg/uploader/block/uploader_test.go diff --git a/changelogs/unreleased/9945-Lyndon-Li b/changelogs/unreleased/9945-Lyndon-Li new file mode 100644 index 000000000..bdb4d8d5e --- /dev/null +++ b/changelogs/unreleased/9945-Lyndon-Li @@ -0,0 +1 @@ +Add snapshot operations for block uploader \ No newline at end of file diff --git a/pkg/uploader/block/dev_linux.go b/pkg/uploader/block/dev_linux.go index 4d49442b3..6383060fb 100644 --- a/pkg/uploader/block/dev_linux.go +++ b/pkg/uploader/block/dev_linux.go @@ -25,6 +25,7 @@ import ( "github.com/pkg/errors" ) +// implement in following PRs func openBlockDevice(path string, read bool) (*os.File, error) { return nil, errors.New("Not implemented") } diff --git a/pkg/uploader/block/snapshot.go b/pkg/uploader/block/snapshot.go index a3bb92431..41d42ba36 100644 --- a/pkg/uploader/block/snapshot.go +++ b/pkg/uploader/block/snapshot.go @@ -25,6 +25,7 @@ import ( "github.com/pkg/errors" "github.com/sirupsen/logrus" + "github.com/vmware-tanzu/velero/pkg/cbtservice" "github.com/vmware-tanzu/velero/pkg/repository/udmrepo" "github.com/vmware-tanzu/velero/pkg/uploader" @@ -202,6 +203,9 @@ func Restore(ctx context.Context, blkup Uploader, rep udmrepo.BackupRepo, snapsh log.Infof("Restore from snapshot %s, description %s, created time %v, tags %v", snapshotID, snapshot.Description, snapshot.EndTime, snapshot.Tags) + bitmap := cbt.NewBitmap(blockSize, uint64(snapshot.TotalSize), "", "", "") + bitmap.SetFull() + destPath, err := filepath.Abs(dest) if err != nil { return 0, errors.Wrapf(err, "invalid dest path '%s'", dest) @@ -214,7 +218,7 @@ func Restore(ctx context.Context, blkup Uploader, rep udmrepo.BackupRepo, snapsh return 0, errors.Wrapf(err, "error opening block device '%s'", destPath) } - size, err := blkup.Restore(snapshot, destInfo{dev: destDev, path: destPath}, uploaderCfg) + size, err := blkup.Restore(snapshot, destInfo{dev: destDev, path: destPath}, bitmap.Iterator(), uploaderCfg) if err != nil { return 0, errors.Wrapf(err, "error restoring to block dev %s", destPath) } diff --git a/pkg/uploader/block/snapshot_test.go b/pkg/uploader/block/snapshot_test.go index 1e609eb2f..5d17ee0f4 100644 --- a/pkg/uploader/block/snapshot_test.go +++ b/pkg/uploader/block/snapshot_test.go @@ -46,8 +46,8 @@ func (m *mockUploader) Backup(src sourceInfo, parent udmrepo.ID, iter cbttypes.I return args.Get(0).(udmrepo.Snapshot), args.Get(1).(int64), args.Error(2) } -func (m *mockUploader) Restore(snap udmrepo.Snapshot, dest destInfo, cfg map[string]string) (int64, error) { - args := m.Called(snap, dest, cfg) +func (m *mockUploader) Restore(snap udmrepo.Snapshot, dest destInfo, iter cbttypes.Iterator, cfg map[string]string) (int64, error) { + args := m.Called(snap, dest, iter, cfg) return args.Get(0).(int64), args.Error(1) } @@ -360,6 +360,8 @@ func TestGetParentBackupInfo(t *testing.T) { 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(&udmrepo.Metadata{SubObjects: []udmrepo.ObjectMetadata{{ID: "root-obj"}}}, nil) }, expectedParent: "root-obj", expectedCID: "cid-abc", @@ -386,6 +388,8 @@ func TestGetParentBackupInfo(t *testing.T) { setupMocks: func(repo *udmrepomocks.BackupRepo) { repo.On("ListSnapshot", mock.Anything, realSource). Return([]udmrepo.Snapshot{validSnap}, nil) + repo.On("ReadMetadata", mock.Anything, udmrepo.ID("root-obj")). + Return(&udmrepo.Metadata{SubObjects: []udmrepo.ObjectMetadata{{ID: "root-obj"}}}, nil) }, expectedParent: "root-obj", expectedCID: "cid-abc", @@ -566,7 +570,7 @@ func TestRestore(t *testing.T) { setupMocks: func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) { repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-001")). Return(storedSnap, nil) - blkup.On("Restore", mock.Anything, mock.Anything, mock.Anything). + blkup.On("Restore", mock.Anything, mock.Anything, mock.Anything, mock.Anything). Return(int64(0), errors.New("restore I/O error")) }, setupOpenDev: func(t *testing.T) *os.File { @@ -579,7 +583,7 @@ func TestRestore(t *testing.T) { setupMocks: func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) { repo.On("GetSnapshot", mock.Anything, udmrepo.ID("snap-001")). Return(storedSnap, nil) - blkup.On("Restore", mock.Anything, mock.Anything, mock.Anything). + blkup.On("Restore", mock.Anything, mock.Anything, mock.Anything, mock.Anything). Return(int64(4096), nil) }, setupOpenDev: func(t *testing.T) *os.File { diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index f487b39cb..7f089bd01 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -22,6 +22,7 @@ import ( "github.com/pkg/errors" "github.com/sirupsen/logrus" + "github.com/vmware-tanzu/velero/pkg/repository/udmrepo" "github.com/vmware-tanzu/velero/pkg/uploader" cbt "github.com/vmware-tanzu/velero/pkg/uploader/cbt/types" @@ -46,9 +47,10 @@ type destInfo struct { type Uploader interface { Backup(sourceInfo, udmrepo.ID, cbt.Iterator, map[string]string) (udmrepo.Snapshot, int64, error) - Restore(udmrepo.Snapshot, destInfo, map[string]string) (int64, error) + Restore(udmrepo.Snapshot, destInfo, cbt.Iterator, map[string]string) (int64, error) } +// implement in following PRs func NewUploader(ctx context.Context, repoWriter udmrepo.BackupRepo, progress uploader.ProgressUpdater, log logrus.FieldLogger) Uploader { return nil } @@ -60,11 +62,11 @@ func loadObjectFromSnapshot(ctx context.Context, rep udmrepo.BackupRepo, snapsho parentMeta, err := rep.ReadMetadata(ctx, snapshot.RootObject.ID) if err != nil { - return "", errors.Wrapf(err, "error readding snapshot metadata for %s", snapshot.Description) + return "", errors.Wrapf(err, "error reading snapshot metadata for %s", snapshot.Description) } if len(parentMeta.SubObjects) != 1 { - return "", errors.Wrapf(err, "unexpected number of bdev object (%d) for snapshot %s", len(parentMeta.SubObjects), snapshot.Description) + return "", errors.Errorf("unexpected number of bdev object (%d) for snapshot %s", len(parentMeta.SubObjects), snapshot.Description) } return parentMeta.SubObjects[0].ID, nil diff --git a/pkg/uploader/block/uploader_test.go b/pkg/uploader/block/uploader_test.go new file mode 100644 index 000000000..ea1986197 --- /dev/null +++ b/pkg/uploader/block/uploader_test.go @@ -0,0 +1,112 @@ +/* +Copyright The Velero Contributors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + +http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package block + +import ( + "context" + "testing" + + "github.com/pkg/errors" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" + + "github.com/vmware-tanzu/velero/pkg/repository/udmrepo" + udmrepomocks "github.com/vmware-tanzu/velero/pkg/repository/udmrepo/mocks" +) + +func TestLoadObjectFromSnapshot(t *testing.T) { + testCases := []struct { + name string + snapshot *udmrepo.Snapshot + setupMocks func(repo *udmrepomocks.BackupRepo) + expectedErrStr string + expectedID udmrepo.ID + }{ + { + name: "nil snapshot", + snapshot: nil, + expectedErrStr: "snapshot is empty", + }, + { + name: "ReadMetadata error", + snapshot: &udmrepo.Snapshot{ + RootObject: udmrepo.ObjectMetadata{ID: "root-obj"}, + }, + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ReadMetadata", mock.Anything, udmrepo.ID("root-obj")). + Return(nil, errors.New("read error")) + }, + expectedErrStr: "error reading snapshot metadata", + }, + { + name: "unexpected number of subobjects (0)", + snapshot: &udmrepo.Snapshot{ + RootObject: udmrepo.ObjectMetadata{ID: "root-obj"}, + }, + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ReadMetadata", mock.Anything, udmrepo.ID("root-obj")). + Return(&udmrepo.Metadata{SubObjects: []udmrepo.ObjectMetadata{}}, nil) + }, + expectedErrStr: "unexpected number of bdev object", + }, + { + name: "unexpected number of subobjects (2)", + snapshot: &udmrepo.Snapshot{ + RootObject: udmrepo.ObjectMetadata{ID: "root-obj"}, + }, + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ReadMetadata", mock.Anything, udmrepo.ID("root-obj")). + Return(&udmrepo.Metadata{SubObjects: []udmrepo.ObjectMetadata{{ID: "obj-1"}, {ID: "obj-2"}}}, nil) + }, + expectedErrStr: "unexpected number of bdev object", + }, + { + name: "success", + snapshot: &udmrepo.Snapshot{ + RootObject: udmrepo.ObjectMetadata{ID: "root-obj"}, + }, + setupMocks: func(repo *udmrepomocks.BackupRepo) { + repo.On("ReadMetadata", mock.Anything, udmrepo.ID("root-obj")). + Return(&udmrepo.Metadata{SubObjects: []udmrepo.ObjectMetadata{{ID: "bdev-obj"}}}, nil) + }, + expectedID: "bdev-obj", + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ctx := context.Background() + mockRepo := udmrepomocks.NewBackupRepo(t) + + if tc.setupMocks != nil { + tc.setupMocks(mockRepo) + } + + id, err := loadObjectFromSnapshot(ctx, mockRepo, tc.snapshot) + + if tc.expectedErrStr != "" { + require.Error(t, err) + assert.ErrorContains(t, err, tc.expectedErrStr) + assert.Empty(t, id) + } else { + require.NoError(t, err) + assert.Equal(t, tc.expectedID, id) + } + }) + } +} From f38bc20a0aec2a9e7dac3506061698b298236c0a Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Thu, 25 Jun 2026 17:01:17 +0800 Subject: [PATCH 4/6] block uploader snapshot implementation Signed-off-by: Lyndon-Li --- pkg/uploader/block/dev_linux.go | 2 +- pkg/uploader/block/dev_other.go | 2 +- pkg/uploader/block/snapshot.go | 18 +++++++++--------- pkg/uploader/block/snapshot_test.go | 19 ++++++++++++------- pkg/uploader/block/uploader.go | 10 +++++----- pkg/uploader/block/uploader_test.go | 4 ++-- pkg/uploader/provider/block_test.go | 10 +++++++--- 7 files changed, 37 insertions(+), 28 deletions(-) diff --git a/pkg/uploader/block/dev_linux.go b/pkg/uploader/block/dev_linux.go index 6383060fb..85b378c55 100644 --- a/pkg/uploader/block/dev_linux.go +++ b/pkg/uploader/block/dev_linux.go @@ -22,7 +22,7 @@ package block import ( "os" - "github.com/pkg/errors" + "github.com/cockroachdb/errors" ) // implement in following PRs diff --git a/pkg/uploader/block/dev_other.go b/pkg/uploader/block/dev_other.go index 60689a3d6..c8a55cab2 100644 --- a/pkg/uploader/block/dev_other.go +++ b/pkg/uploader/block/dev_other.go @@ -25,5 +25,5 @@ import ( ) func openBlockDevice(_ string, _ bool) (*os.File, error) { - return nil, fmt.Errorf("block mode is not supported for Windows") + return nil, fmt.Errorf("block mode is not supported for non-linux platforms") } diff --git a/pkg/uploader/block/snapshot.go b/pkg/uploader/block/snapshot.go index 41d42ba36..30626da53 100644 --- a/pkg/uploader/block/snapshot.go +++ b/pkg/uploader/block/snapshot.go @@ -23,7 +23,7 @@ import ( "path/filepath" "time" - "github.com/pkg/errors" + "github.com/cockroachdb/errors" "github.com/sirupsen/logrus" "github.com/vmware-tanzu/velero/pkg/cbtservice" @@ -41,9 +41,9 @@ type parentBackupInfo struct { } // 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) { - if blkup == nil { +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) { + if blkUp == nil { return uploader.SnapshotInfo{}, false, errors.New("get empty block uploader") } @@ -77,7 +77,7 @@ func Backup(ctx context.Context, blkup Uploader, repoWriter udmrepo.BackupRepo, return uploader.SnapshotInfo{}, false, errors.Wrapf(err, "error reset pos of block device %s", source) } - snapID, backupSize, err := snapshotSource(ctx, repoWriter, blkup, sourceInfo, forceFull, parentSnapshot, cbtSource, cbtservice, tags, uploaderCfg, log, "Block Uploader") + snapID, backupSize, err := snapshotSource(ctx, repoWriter, blkUp, sourceInfo, forceFull, parentSnapshot, cbtSource, cbtService, tags, uploaderCfg, log, "Block Uploader") snapshotInfo := uploader.SnapshotInfo{ ID: snapID, Size: sourceInfo.size, @@ -95,7 +95,7 @@ func snapshotSource( forceFull bool, parentSnapshot string, cbtSource cbtservice.SourceInfo, - cbtservice cbtservice.Service, + cbtService cbtservice.Service, snapshotTags map[string]string, uploaderCfg map[string]string, log logrus.FieldLogger, @@ -108,7 +108,7 @@ 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) if err != nil { parentBackup.parentObject = "" log.WithError(err).Warnf("Failed to create CBT with source %v, fallback to real full backup", cbtSource) @@ -193,7 +193,7 @@ func getParentBackupInfo(ctx context.Context, rep udmrepo.BackupRepo, forceFull } // Restore restore specific sourcePath with given snapshotID and update progress -func Restore(ctx context.Context, blkup Uploader, rep udmrepo.BackupRepo, snapshotID, dest string, uploaderCfg map[string]string, log logrus.FieldLogger) (int64, error) { +func Restore(ctx context.Context, blkUp Uploader, rep udmrepo.BackupRepo, snapshotID, dest string, uploaderCfg map[string]string, log logrus.FieldLogger) (int64, error) { log.Info("Start to restore...") snapshot, err := rep.GetSnapshot(ctx, udmrepo.ID(snapshotID)) @@ -218,7 +218,7 @@ func Restore(ctx context.Context, blkup Uploader, rep udmrepo.BackupRepo, snapsh return 0, errors.Wrapf(err, "error opening block device '%s'", destPath) } - size, err := blkup.Restore(snapshot, destInfo{dev: destDev, path: destPath}, bitmap.Iterator(), uploaderCfg) + size, err := blkUp.Restore(snapshot, destInfo{dev: destDev, path: destPath}, bitmap.Iterator(), uploaderCfg) if err != nil { return 0, errors.Wrapf(err, "error restoring to block dev %s", destPath) } diff --git a/pkg/uploader/block/snapshot_test.go b/pkg/uploader/block/snapshot_test.go index 5d17ee0f4..8f6338311 100644 --- a/pkg/uploader/block/snapshot_test.go +++ b/pkg/uploader/block/snapshot_test.go @@ -24,7 +24,7 @@ import ( "testing" "time" - "github.com/pkg/errors" + "github.com/cockroachdb/errors" "github.com/sirupsen/logrus" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/mock" @@ -59,7 +59,7 @@ func testLog() logrus.FieldLogger { func tempFile(t *testing.T, content string) *os.File { t.Helper() - f, err := os.CreateTemp("", "blktest-*") + f, err := os.CreateTemp(t.TempDir(), "blktest-*") require.NoError(t, err) if content != "" { _, err = f.WriteString(content) @@ -93,6 +93,7 @@ func TestBackup(t *testing.T) { { name: "SnapshotSource error propagates", setupOpenDev: func(t *testing.T) *os.File { + t.Helper() return tempFile(t, "") }, setupMocks: func(blkup *mockUploader, _ *udmrepomocks.BackupRepo) { @@ -104,6 +105,7 @@ func TestBackup(t *testing.T) { { name: "success returns correct SnapshotInfo", setupOpenDev: func(t *testing.T) *os.File { + t.Helper() return tempFile(t, "test-block-data") }, setupMocks: func(blkup *mockUploader, repo *udmrepomocks.BackupRepo) { @@ -113,9 +115,10 @@ func TestBackup(t *testing.T) { repo.On("Flush", mock.Anything).Return(nil) }, checkInfo: func(t *testing.T, info uploader.SnapshotInfo) { + t.Helper() assert.Equal(t, "snap-001", info.ID) assert.Equal(t, int64(8), info.IncrementalSize) - assert.Greater(t, info.Size, int64(0)) + assert.Positive(t, info.Size) }, }, } @@ -157,7 +160,7 @@ func TestBackup(t *testing.T) { if tc.expectedErrStr != "" { require.Error(t, err) - assert.ErrorContains(t, err, tc.expectedErrStr) + require.ErrorContains(t, err, tc.expectedErrStr) } else { require.NoError(t, err) assert.False(t, isEmpty) @@ -260,7 +263,7 @@ func TestSnapshotSource(t *testing.T) { if tc.expectedErrStr != "" { require.Error(t, err) - assert.ErrorContains(t, err, tc.expectedErrStr) + require.ErrorContains(t, err, tc.expectedErrStr) } else { require.NoError(t, err) assert.Equal(t, tc.expectedSnapID, snapID) @@ -530,7 +533,7 @@ func TestFindPreviousSnapshot(t *testing.T) { if tc.expectedErrStr != "" { require.Error(t, err) - assert.ErrorContains(t, err, tc.expectedErrStr) + require.ErrorContains(t, err, tc.expectedErrStr) } else { require.NoError(t, err) assert.Equal(t, udmrepo.ID(tc.expectedID), snap.RootObject.ID) @@ -574,6 +577,7 @@ func TestRestore(t *testing.T) { Return(int64(0), errors.New("restore I/O error")) }, setupOpenDev: func(t *testing.T) *os.File { + t.Helper() return tempFile(t, "") }, expectedErrStr: "error restoring to block dev", @@ -587,6 +591,7 @@ func TestRestore(t *testing.T) { Return(int64(4096), nil) }, setupOpenDev: func(t *testing.T) *os.File { + t.Helper() return tempFile(t, "") }, expectedSize: 4096, @@ -616,7 +621,7 @@ func TestRestore(t *testing.T) { if tc.expectedErrStr != "" { require.Error(t, err) - assert.ErrorContains(t, err, tc.expectedErrStr) + require.ErrorContains(t, err, tc.expectedErrStr) assert.Equal(t, int64(0), size) } else { require.NoError(t, err) diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index 7f089bd01..233d72a17 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -20,7 +20,7 @@ import ( "context" "os" - "github.com/pkg/errors" + "github.com/cockroachdb/errors" "github.com/sirupsen/logrus" "github.com/vmware-tanzu/velero/pkg/repository/udmrepo" @@ -60,14 +60,14 @@ func loadObjectFromSnapshot(ctx context.Context, rep udmrepo.BackupRepo, snapsho return "", errors.New("snapshot is empty") } - parentMeta, err := rep.ReadMetadata(ctx, snapshot.RootObject.ID) + meta, err := rep.ReadMetadata(ctx, snapshot.RootObject.ID) if err != nil { return "", errors.Wrapf(err, "error reading snapshot metadata for %s", snapshot.Description) } - if len(parentMeta.SubObjects) != 1 { - return "", errors.Errorf("unexpected number of bdev object (%d) for snapshot %s", len(parentMeta.SubObjects), snapshot.Description) + if len(meta.SubObjects) != 1 { + return "", errors.Errorf("unexpected number of bdev object (%d) for snapshot %s", len(meta.SubObjects), snapshot.Description) } - return parentMeta.SubObjects[0].ID, nil + return meta.SubObjects[0].ID, nil } diff --git a/pkg/uploader/block/uploader_test.go b/pkg/uploader/block/uploader_test.go index ea1986197..8209569e1 100644 --- a/pkg/uploader/block/uploader_test.go +++ b/pkg/uploader/block/uploader_test.go @@ -20,7 +20,7 @@ import ( "context" "testing" - "github.com/pkg/errors" + "github.com/cockroachdb/errors" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" @@ -101,7 +101,7 @@ func TestLoadObjectFromSnapshot(t *testing.T) { if tc.expectedErrStr != "" { require.Error(t, err) - assert.ErrorContains(t, err, tc.expectedErrStr) + require.ErrorContains(t, err, tc.expectedErrStr) assert.Empty(t, id) } else { require.NoError(t, err) diff --git a/pkg/uploader/provider/block_test.go b/pkg/uploader/provider/block_test.go index e7af93855..8ec445168 100644 --- a/pkg/uploader/provider/block_test.go +++ b/pkg/uploader/provider/block_test.go @@ -280,6 +280,7 @@ func TestBlockProviderRunBackup(t *testing.T) { mockBackupResult: uploader.SnapshotInfo{ID: "snap-tags"}, expectedID: "snap-tags", checkCaptures: func(t *testing.T, _ string, tags map[string]string) { + t.Helper() assert.Equal(t, requestorType, tags[uploader.SnapshotRequesterTag]) assert.Equal(t, uploader.BlockType, tags[uploader.SnapshotUploaderTag]) }, @@ -292,6 +293,7 @@ func TestBlockProviderRunBackup(t *testing.T) { mockBackupResult: uploader.SnapshotInfo{ID: "snap-source"}, expectedID: "snap-source", checkCaptures: func(t *testing.T, realSource string, _ map[string]string) { + t.Helper() assert.Equal(t, requestorType+"/"+uploader.BlockType+"/my-volume", realSource) }, }, @@ -303,7 +305,8 @@ func TestBlockProviderRunBackup(t *testing.T) { mockBackupResult: uploader.SnapshotInfo{ID: "snap-nosource"}, expectedID: "snap-nosource", checkCaptures: func(t *testing.T, realSource string, _ map[string]string) { - assert.Equal(t, "", realSource) + t.Helper() + assert.Empty(t, realSource) }, }, } @@ -349,7 +352,7 @@ func TestBlockProviderRunBackup(t *testing.T) { if tc.expectError { require.Error(t, err) if tc.expectedErrStr != "" { - assert.ErrorContains(t, err, tc.expectedErrStr) + require.ErrorContains(t, err, tc.expectedErrStr) } } else { require.NoError(t, err) @@ -416,6 +419,7 @@ func TestBlockProviderRunRestore(t *testing.T) { mockRestoreSize: 512, expectedSize: 512, checkCaptures: func(t *testing.T, snapshotID, volumePath string) { + t.Helper() assert.Equal(t, "snap-fwd", snapshotID) assert.Equal(t, "/dev/sdc", volumePath) }, @@ -452,7 +456,7 @@ func TestBlockProviderRunRestore(t *testing.T) { if tc.expectError { require.Error(t, err) if tc.expectedErrStr != "" { - assert.ErrorContains(t, err, tc.expectedErrStr) + require.ErrorContains(t, err, tc.expectedErrStr) } assert.Equal(t, int64(0), size) } else { From 82dbef2cd9bee892ee42d4c6c039932d9eb5f4f7 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Tue, 30 Jun 2026 11:06:46 +0800 Subject: [PATCH 5/6] block uploader snapshot implementation Signed-off-by: Lyndon-Li --- pkg/uploader/provider/block.go | 6 +++++- pkg/uploader/provider/block_test.go | 8 +++++++- 2 files changed, 12 insertions(+), 2 deletions(-) diff --git a/pkg/uploader/provider/block.go b/pkg/uploader/provider/block.go index 427d3fae3..4bc26f9e1 100644 --- a/pkg/uploader/provider/block.go +++ b/pkg/uploader/provider/block.go @@ -105,7 +105,7 @@ func (bp *blockProvider) RunBackup( uploaderCfg map[string]string, updater uploader.ProgressUpdater) (string, bool, int64, int64, error) { if updater == nil { - return "", false, 0, 0, errors.New("Need to initial backup progress updater first") + return "", false, 0, 0, errors.New("backup progress updater is invalid") } if path == "" { @@ -160,6 +160,10 @@ func (bp *blockProvider) RunRestore( volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, updater uploader.ProgressUpdater) (int64, error) { + if updater == nil { + return 0, errors.New("restore progress updater is invalid") + } + log := bp.log.WithFields(logrus.Fields{ "snapshotID": snapshotID, "volumePath": volumePath, diff --git a/pkg/uploader/provider/block_test.go b/pkg/uploader/provider/block_test.go index 8ec445168..ad8f68b52 100644 --- a/pkg/uploader/provider/block_test.go +++ b/pkg/uploader/provider/block_test.go @@ -224,7 +224,7 @@ func TestBlockProviderRunBackup(t *testing.T) { path: "/dev/sda", updater: nil, expectError: true, - expectedErrStr: "Need to initial backup progress updater first", + expectedErrStr: "backup progress updater is invalid", skipMock: true, }, { @@ -385,6 +385,12 @@ func TestBlockProviderRunRestore(t *testing.T) { expectedErrStr string checkCaptures func(*testing.T, string, string) }{ + { + name: "nil updater returns error", + updater: nil, + expectError: true, + expectedErrStr: "restore progress updater is invalid", + }, { name: "success returns size and updates progress", snapshotID: "snap-001", From 8df8709a8a522c9477b615eee609e2cd67834fc6 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Tue, 30 Jun 2026 14:14:23 +0800 Subject: [PATCH 6/6] block uploader snapshot implementation Signed-off-by: Lyndon-Li --- pkg/uploader/provider/block.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/pkg/uploader/provider/block.go b/pkg/uploader/provider/block.go index 4bc26f9e1..9135bb67b 100644 --- a/pkg/uploader/provider/block.go +++ b/pkg/uploader/provider/block.go @@ -118,6 +118,8 @@ func (bp *blockProvider) RunBackup( "parentSnapshot": parentSnapshot, }) + log.Infof("Run block backup, CBT source info: %v", cbtParam.Source) + blkUploader := block.NewUploader(ctx, bp.bkRepo, updater, log) if tags == nil {