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 new file mode 100644 index 000000000..85b378c55 --- /dev/null +++ b/pkg/uploader/block/dev_linux.go @@ -0,0 +1,31 @@ +//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/cockroachdb/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/dev_other.go b/pkg/uploader/block/dev_other.go new file mode 100644 index 000000000..c8a55cab2 --- /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 non-linux platforms") +} diff --git a/pkg/uploader/block/snapshot.go b/pkg/uploader/block/snapshot.go new file mode 100644 index 000000000..30626da53 --- /dev/null +++ b/pkg/uploader/block/snapshot.go @@ -0,0 +1,272 @@ +/* +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/cockroachdb/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 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 = obj + parentInfo.changeID = previous.Tags[uploader.CBTChangeIDTag] + parentInfo.volumeID = previous.Tags[uploader.CBTVolumeIDTag] + + log.Infof("Using parent snapshot %s, start time %v, end time %v, description %s", 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) + + 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) + } + + 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}, bitmap.Iterator(), 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..8f6338311 --- /dev/null +++ b/pkg/uploader/block/snapshot_test.go @@ -0,0 +1,634 @@ +/* +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/cockroachdb/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, iter cbttypes.Iterator, cfg map[string]string) (int64, error) { + args := m.Called(snap, dest, iter, 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(t.TempDir(), "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 { + t.Helper() + 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 { + t.Helper() + 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) { + t.Helper() + assert.Equal(t, "snap-001", info.ID) + assert.Equal(t, int64(8), info.IncrementalSize) + assert.Positive(t, info.Size) + }, + }, + } + + 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) + require.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) + require.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) + 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", + 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) + 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", + 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) + require.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, mock.Anything). + 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", + }, + { + 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, mock.Anything). + Return(int64(4096), nil) + }, + setupOpenDev: func(t *testing.T) *os.File { + t.Helper() + 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) + require.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..233d72a17 --- /dev/null +++ b/pkg/uploader/block/uploader.go @@ -0,0 +1,73 @@ +/* +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/cockroachdb/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, 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 +} + +func loadObjectFromSnapshot(ctx context.Context, rep udmrepo.BackupRepo, snapshot *udmrepo.Snapshot) (udmrepo.ID, error) { + if snapshot == nil { + return "", errors.New("snapshot is empty") + } + + 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(meta.SubObjects) != 1 { + return "", errors.Errorf("unexpected number of bdev object (%d) for snapshot %s", len(meta.SubObjects), snapshot.Description) + } + + return meta.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..8209569e1 --- /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/cockroachdb/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) + require.ErrorContains(t, err, tc.expectedErrStr) + assert.Empty(t, id) + } else { + require.NoError(t, err) + assert.Equal(t, tc.expectedID, id) + } + }) + } +} diff --git a/pkg/uploader/provider/block.go b/pkg/uploader/provider/block.go index dc5028040..9135bb67b 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,57 @@ 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("backup progress updater is invalid") + } + + if path == "" { + return "", false, 0, 0, errors.New("path is empty") + } + + log := bp.log.WithFields(logrus.Fields{ + "path": path, + "realSource": realSource, + "parentSnapshot": parentSnapshot, + }) + + log.Infof("Run block backup, CBT source info: %v", cbtParam.Source) + + 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 +162,35 @@ func (bp *blockProvider) RunRestore( volMode uploader.PersistentVolumeMode, uploaderCfg map[string]string, updater uploader.ProgressUpdater) (int64, error) { - return 0, errors.New("block restore not implemented") + if updater == nil { + return 0, errors.New("restore progress updater is invalid") + } + + 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..ad8f68b52 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,286 @@ 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: "backup progress updater is invalid", + 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) { + t.Helper() + 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) { + t.Helper() + 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) { + t.Helper() + assert.Empty(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 != "" { + require.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: "nil updater returns error", + updater: nil, + expectError: true, + expectedErrStr: "restore progress updater is invalid", + }, + { + 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) { + t.Helper() + 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 != "" { + require.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