Merge pull request #9945 from Lyndon-Li/block-uploader-snapshot-operations

Block uploader snapshot operations
This commit is contained in:
lyndon-li
2026-06-30 15:46:31 +08:00
committed by GitHub
10 changed files with 1539 additions and 6 deletions
+1
View File
@@ -0,0 +1 @@
Add snapshot operations for block uploader
+31
View File
@@ -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")
}
+29
View File
@@ -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")
}
+272
View File
@@ -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
}
+634
View File
@@ -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)
})
}
}
+73
View File
@@ -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
}
+112
View File
@@ -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)
}
})
}
}
+85 -4
View File
@@ -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
}
+297
View File
@@ -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)
}
})
}
}
+5 -2
View File
@@ -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