From 7bbd172684dc6bb92801a0d1f91b7047b85980d8 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Wed, 22 Jul 2026 17:52:01 +0800 Subject: [PATCH 1/7] persist source dev size Signed-off-by: Lyndon-Li --- pkg/uploader/block/snapshot.go | 12 +++++++++++- pkg/uploader/block/uploader.go | 10 ++++++++-- 2 files changed, 19 insertions(+), 3 deletions(-) diff --git a/pkg/uploader/block/snapshot.go b/pkg/uploader/block/snapshot.go index e30f5c1bb..b185f4e15 100644 --- a/pkg/uploader/block/snapshot.go +++ b/pkg/uploader/block/snapshot.go @@ -222,7 +222,17 @@ func Restore(ctx context.Context, blkUp Uploader, rep udmrepo.BackupRepo, snapsh defer destDev.Close() - size, err := blkUp.Restore(snapshot, destInfo{dev: destDev, path: destPath}, bitmap.Iterator(), uploaderCfg) + destSize, err := destDev.Seek(0, io.SeekEnd) + if err != nil { + return 0, errors.Wrapf(err, "error getting length of block device %s", dest) + } + + _, err = destDev.Seek(0, io.SeekStart) + if err != nil { + return 0, errors.Wrapf(err, "error reset pos of block device %s", dest) + } + + size, err := blkUp.Restore(snapshot, destInfo{dev: destDev, path: destPath, size: destSize}, bitmap.Iterator(), uploaderCfg) if err != nil { return 0, errors.Wrapf(err, "error restoring to block dev %s", destPath) } diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index 75e913cb7..1e4982483 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -21,6 +21,7 @@ import ( "io" "os" "runtime" + "strconv" "strings" "github.com/cockroachdb/errors" @@ -35,8 +36,9 @@ import ( var ErrCanceled = errors.New("uploader is canceled") const ( - blockSize = (1 << 20) - bufferSize = 100 << 20 + blockSize = (1 << 20) + bufferSize = 100 << 20 + bdevSourceSizeTag = "bdev-source-size" ) type sourceInfo struct { @@ -48,6 +50,7 @@ type sourceInfo struct { type destInfo struct { dev *os.File path string + size int64 } type Uploader interface { @@ -134,6 +137,9 @@ func (blkup *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, b Type: udmrepo.ObjectDataTypeMetadata, Permissions: 0o777, }, + Tags: map[string]string{ + bdevSourceSizeTag: strconv.FormatInt(source.size, 10), + }, }, backupSize, nil } From bec292e738fce49f221ba16a2e3a6247acb37563 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Wed, 22 Jul 2026 17:58:02 +0800 Subject: [PATCH 2/7] block uploader restore data Signed-off-by: Lyndon-Li --- pkg/uploader/block/uploader.go | 63 ++++++++++++++++++++++++++++++++-- 1 file changed, 61 insertions(+), 2 deletions(-) diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index 1e4982483..0567edb8c 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -143,9 +143,46 @@ func (blkup *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, b }, backupSize, nil } -// TODO implement in following PRs func (blkup *blockUploader) Restore(snapshot udmrepo.Snapshot, dest destInfo, bitmap cbt.Iterator, configs map[string]string) (int64, error) { - return 0, errors.New("not implemented") + if bitmap == nil { + return 0, errors.New("bitmap is not available") + } + + meta, err := blkup.repoWriter.ReadMetadata(blkup.ctx, snapshot.RootObject.ID) + if err != nil { + return 0, errors.Wrapf(err, "error readding snapshot metadata for %s", snapshot.Description) + } + + if len(meta.SubObjects) != 1 { + return 0, errors.Wrapf(err, "unexpected number of bdev object (%d) for snapshot %s", len(meta.SubObjects), snapshot.Description) + } + + sourceSize, err := getSourceSize(snapshot) + if err != nil { + sourceSize = meta.SubObjects[0].Size + blkup.log.Warnf("Failed to get source size from snapshot %s, use backup size %v", snapshot.Description, sourceSize) + } + + if sourceSize > meta.SubObjects[0].Size { + return 0, errors.Wrapf(err, "unexpected size (%v vs. %v) for bdev object %s", meta.SubObjects[0].Size, sourceSize, meta.SubObjects[0].Name) + } + + if sourceSize > dest.size { + return 0, errors.Wrapf(err, "dest dev(%s) size is too small (%v vs. %v)", dest.path, dest.size, sourceSize) + } + + reader, err := blkup.repoWriter.OpenObject(blkup.ctx, meta.SubObjects[0].ID) + if err != nil { + return 0, errors.Wrapf(err, "error opening bdev object %v", meta.SubObjects[0].Name) + } + defer reader.Close() + + size, err := blkup.restoreData(reader, dest.dev, bitmap, sourceSize, dest.path) + if err != nil { + return 0, errors.Wrapf(err, "error restoring bdev object %s to volume %s", meta.SubObjects[0].Name, dest.path) + } + + return size, nil } func (blkup *blockUploader) backupObject(dev *os.File, dest udmrepo.ObjectWriter, bitmap cbt.Iterator, totalLength int64) (udmrepo.ID, int64, int64, error) { @@ -324,6 +361,28 @@ func getObjectName(source string) string { return strings.Trim(s, "-") } +func (blkup *blockUploader) restoreData(reader io.ReadSeeker, dest *os.File, bitmap cbt.Iterator, totalLength int64, destPath string) (int64, error) { + return 0, nil +} + +func getSourceSize(snapshot udmrepo.Snapshot) (int64, error) { + if snapshot.Tags == nil { + return 0, errors.New("source size tag is empty") + } + + s, found := snapshot.Tags[bdevSourceSizeTag] + if !found { + return 0, errors.New("source size tag is missing") + } + + size, err := strconv.ParseInt(s, 10, 64) + if err != nil { + return 0, errors.Wrapf(err, "error parsing size from %s", s) + } + + return size, nil +} + func loadObjectFromSnapshot(ctx context.Context, rep udmrepo.BackupRepo, snapshot *udmrepo.Snapshot) (udmrepo.ID, error) { if snapshot == nil { return "", errors.New("snapshot is empty") From 9c1d8bee71a529b52f1e0a3d3c710b5b18836464 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Wed, 22 Jul 2026 18:00:13 +0800 Subject: [PATCH 3/7] block uploader restore data Signed-off-by: Lyndon-Li --- pkg/uploader/block/uploader.go | 154 ++++++++++++++++++++++++++++++++- 1 file changed, 153 insertions(+), 1 deletion(-) diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index 0567edb8c..30908b396 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -17,6 +17,7 @@ limitations under the License. package block import ( + "bytes" "context" "io" "os" @@ -362,7 +363,158 @@ func getObjectName(source string) string { } func (blkup *blockUploader) restoreData(reader io.ReadSeeker, dest *os.File, bitmap cbt.Iterator, totalLength int64, destPath string) (int64, error) { - return 0, nil + list := freelist.New(bufferSize, blockSize) + resultChan := make(chan readResult, list.Capacity()) + zeroBlock := make([]byte, blockSize) + totalCount := bitmap.Count() + + quit := make(chan struct{}) + defer close(quit) + + go func() { + defer close(resultChan) + + offset, valid := bitmap.Next() + var buffer []byte + var nextPos uint64 = uint64(0) + for valid { + select { + case <-blkup.ctx.Done(): + return + case <-quit: + return + case buffer = <-list.Chunks(): + } + + var err error + + if nextPos != offset { + _, err = reader.Seek(int64(offset), io.SeekStart) + } + + if err == nil { + var length int + length, err = io.ReadFull(reader, buffer) + if err == nil && length <= 0 { + err = io.ErrUnexpectedEOF + } + } + + r := readResult{ + buffer: buffer, + offset: int64(offset), + err: err, + } + + if r.err != nil { + r.resetBuffer(list) + } + + resultChan <- r + + if r.err != nil { + return + } + + nextPos = offset + uint64(blockSize) + offset, valid = bitmap.Next() + } + }() + + var written int64 + var result readResult + var writeErr error + var readerRunning bool + var zeroStart int64 = -1 + var zeroLength int64 + var curCount int64 + + for curCount < int64(totalCount) { + select { + case <-blkup.ctx.Done(): + writeErr = ErrCanceled + case result, readerRunning = <-resultChan: + if !readerRunning { + if blkup.ctx.Err() != nil { + writeErr = ErrCanceled + } else { + writeErr = io.ErrUnexpectedEOF + } + } + } + + if writeErr != nil { + break + } + + if result.err != nil { + writeErr = result.err + break + } + + length := min(int64(blockSize), totalLength-result.offset) + if bytes.Equal(result.buffer, zeroBlock) { + if zeroStart == -1 { + zeroStart = result.offset + zeroLength = length + } else if result.offset == zeroStart+zeroLength { + zeroLength += length + } else { + if err := blkup.flushZeroBlocks(dest, zeroStart, zeroLength, zeroBlock, destPath); err != nil { + writeErr = errors.Wrapf(err, "error flushing zero blocks from %v, length %v", zeroStart, zeroLength) + break + } + zeroStart = result.offset + zeroLength = length + } + } else { + if zeroStart != -1 { + if err := blkup.flushZeroBlocks(dest, zeroStart, zeroLength, zeroBlock, destPath); err != nil { + writeErr = errors.Wrapf(err, "error flushing zero blocks from %v, length %v", zeroStart, zeroLength) + break + } + + zeroStart = -1 + zeroLength = 0 + } + + n, err := dest.WriteAt(result.buffer[:length], result.offset) + if err != nil { + writeErr = err + break + } + + if length != int64(n) { + writeErr = io.ErrShortWrite + break + } + } + + written += length + curCount++ + + result.resetBuffer(list) + + blkup.progress.UpdateProgress(&uploader.Progress{BytesDone: written, TotalBytes: totalLength}) + } + + result.resetBuffer(list) + + if writeErr != nil { + return written, writeErr + } + + if zeroStart != -1 { + if err := blkup.flushZeroBlocks(dest, zeroStart, zeroLength, zeroBlock, destPath); err != nil { + return written, errors.Wrapf(err, "error flushing zero blocks from %v, length %v", zeroStart, zeroLength) + } + } + + return written, nil +} + +func (bu *blockUploader) flushZeroBlocks(dest *os.File, start int64, length int64, zeroBlock []byte, destPath string) error { + return nil } func getSourceSize(snapshot udmrepo.Snapshot) (int64, error) { From 97978ed9b767a0e14b1ed0b8dc50c1d406cc1ef5 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Wed, 22 Jul 2026 18:01:57 +0800 Subject: [PATCH 4/7] block uploader flush zero blocks Signed-off-by: Lyndon-Li --- pkg/uploader/block/uploader.go | 23 +++++++++++++++++++++++ 1 file changed, 23 insertions(+) diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index 30908b396..7717a7e7e 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -514,6 +514,29 @@ func (blkup *blockUploader) restoreData(reader io.ReadSeeker, dest *os.File, bit } func (bu *blockUploader) flushZeroBlocks(dest *os.File, start int64, length int64, zeroBlock []byte, destPath string) error { + err := blkZeroOut(dest, start, length) + if err == nil { + return nil + } + + bu.log.WithError(err).Warnf("Failed to call zero out from dev %s, start %v, length %v. Fallback to conservative way", destPath, start, length) + + var written int64 + for written < length { + writeSize := min(len(zeroBlock), int(length-written)) + + n, err := dest.WriteAt(zeroBlock[:writeSize], start+written) + if err != nil { + return errors.Wrapf(err, "error writing zero buffer at %v, length %v", start+written, writeSize) + } + + if writeSize != n { + return errors.Wrapf(err, "short write zero buffer at %v, length %v", start+written, writeSize) + } + + written += int64(writeSize) + } + return nil } From 236c45b4436229c0012d28fdccbafdc720578a71 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Fri, 24 Jul 2026 13:37:40 +0800 Subject: [PATCH 5/7] block uploader restore UT Signed-off-by: Lyndon-Li --- pkg/uploader/block/uploader_test.go | 228 ++++++++++++++++++++++++++++ 1 file changed, 228 insertions(+) diff --git a/pkg/uploader/block/uploader_test.go b/pkg/uploader/block/uploader_test.go index 88fd4771e..69b5efcb5 100644 --- a/pkg/uploader/block/uploader_test.go +++ b/pkg/uploader/block/uploader_test.go @@ -459,3 +459,231 @@ func TestLoadObjectFromSnapshot(t *testing.T) { }) } } + +func TestGetSourceSize(t *testing.T) { + testCases := []struct { + name string + snapshot udmrepo.Snapshot + expectErr bool + expected int64 + }{ + { + name: "nil tags", + snapshot: udmrepo.Snapshot{}, + expectErr: true, + }, + { + name: "missing tag", + snapshot: udmrepo.Snapshot{ + Tags: map[string]string{}, + }, + expectErr: true, + }, + { + name: "invalid tag value", + snapshot: udmrepo.Snapshot{ + Tags: map[string]string{ + "bdev-source-size": "abc", + }, + }, + expectErr: true, + }, + { + name: "valid tag value", + snapshot: udmrepo.Snapshot{ + Tags: map[string]string{ + "bdev-source-size": "1048576", + }, + }, + expectErr: false, + expected: 1048576, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + size, err := getSourceSize(tc.snapshot) + if tc.expectErr { + assert.Error(t, err) + } else { + assert.NoError(t, err) + assert.Equal(t, tc.expected, size) + } + }) + } +} + +func TestFlushZeroBlocks(t *testing.T) { + t.Run("success via write fallback", func(t *testing.T) { + f, err := os.CreateTemp("", "zerotest-*") + require.NoError(t, err) + defer os.Remove(f.Name()) + defer f.Close() + + require.NoError(t, f.Truncate(2048)) + + bu := &blockUploader{ + log: logrus.New(), + } + bu.log.(*logrus.Logger).Out = io.Discard + + zeroBlock := make([]byte, 1024) + err = bu.flushZeroBlocks(f, 0, 2048, zeroBlock, f.Name()) + + assert.NoError(t, err) + + data, err := os.ReadFile(f.Name()) + require.NoError(t, err) + assert.Equal(t, make([]byte, 2048), data) + }) +} + +type errReader struct { + err error +} + +func (r *errReader) Read(p []byte) (n int, err error) { + return 0, r.err +} + +func (r *errReader) Seek(offset int64, whence int) (int64, error) { + return 0, nil +} + +func TestRestoreData(t *testing.T) { + t.Run("success", func(t *testing.T) { + ctx := context.Background() + progress := &mockProgressUpdater{} + progress.On("UpdateProgress", mock.Anything).Return() + bu := &blockUploader{ + ctx: ctx, + progress: progress, + log: logrus.New(), + } + + f, err := os.CreateTemp("", "restoretest-*") + require.NoError(t, err) + defer os.Remove(f.Name()) + defer f.Close() + + data := make([]byte, 1048576) + for i := range data { + data[i] = 1 + } + reader := bytes.NewReader(data) + + iterMock := cbtmocks.NewIterator(t) + iterMock.On("Count").Return(uint64(1)) + iterMock.On("Next").Return(uint64(0), true).Once() + iterMock.On("Next").Return(uint64(0), false) + + written, err := bu.restoreData(reader, f, iterMock, 1048576, f.Name()) + assert.NoError(t, err) + assert.Equal(t, int64(1048576), written) + + f.Seek(0, 0) + writtenData, err := io.ReadAll(f) + require.NoError(t, err) + assert.Equal(t, data, writtenData) + }) + + t.Run("read err", func(t *testing.T) { + ctx := context.Background() + bu := &blockUploader{ + ctx: ctx, + log: logrus.New(), + } + + f, err := os.CreateTemp("", "restoretest-*") + require.NoError(t, err) + defer os.Remove(f.Name()) + defer f.Close() + + reader := &errReader{err: errors.New("read error")} + + iterMock := cbtmocks.NewIterator(t) + iterMock.On("Count").Return(uint64(1)) + iterMock.On("Next").Return(uint64(0), true).Once() + iterMock.On("Next").Return(uint64(0), false) + + _, err = bu.restoreData(reader, f, iterMock, 1048576, f.Name()) + assert.Error(t, err) + assert.Contains(t, err.Error(), "read error") + }) +} + +func TestBlockUploaderRestore(t *testing.T) { + t.Run("missing metadata", func(t *testing.T) { + ctx := context.Background() + repoWriter := udmrepomocks.NewBackupRepo(t) + bu := NewUploader(ctx, repoWriter, nil, logrus.New()) + + repoWriter.On("ReadMetadata", mock.Anything, udmrepo.ID("root-id")).Return(nil, errors.New("meta not found")) + + iterMock := cbtmocks.NewIterator(t) + _, err := bu.Restore(udmrepo.Snapshot{RootObject: udmrepo.ObjectMetadata{ID: "root-id"}}, destInfo{}, iterMock, nil) + assert.Error(t, err) + assert.Contains(t, err.Error(), "meta not found") + }) + + t.Run("success", func(t *testing.T) { + ctx := context.Background() + repoWriter := udmrepomocks.NewBackupRepo(t) + progress := &mockProgressUpdater{} + progress.On("UpdateProgress", mock.Anything).Return() + + bu := NewUploader(ctx, repoWriter, progress, logrus.New()) + + f, err := os.CreateTemp("", "restoretest-*") + require.NoError(t, err) + defer os.Remove(f.Name()) + defer f.Close() + + meta := &udmrepo.Metadata{ + SubObjects: []udmrepo.ObjectMetadata{ + { + ID: "data-id", + Name: "bdev", + Size: 1048576, + }, + }, + } + + repoWriter.On("ReadMetadata", mock.Anything, udmrepo.ID("root-id")).Return(meta, nil) + + objReader := udmrepomocks.NewObjectReader(t) + objReader.On("Read", mock.Anything).Run(func(args mock.Arguments) { + p := args.Get(0).([]byte) + for i := range p { + p[i] = 1 + } + }).Return(1048576, io.EOF).Once() + objReader.On("Read", mock.Anything).Return(0, io.EOF) + objReader.On("Close").Return(nil) + + repoWriter.On("OpenObject", mock.Anything, udmrepo.ID("data-id")).Return(objReader, nil) + + snap := udmrepo.Snapshot{ + Description: "test snapshot", + RootObject: udmrepo.ObjectMetadata{ID: "root-id"}, + Tags: map[string]string{ + "bdev-source-size": "1048576", + }, + } + + dest := destInfo{ + dev: f, + size: 2048576, + path: f.Name(), + } + + iterMock := cbtmocks.NewIterator(t) + iterMock.On("Count").Return(uint64(1)) + iterMock.On("Next").Return(uint64(0), true).Once() + iterMock.On("Next").Return(uint64(0), false) + + written, err := bu.Restore(snap, dest, iterMock, nil) + assert.NoError(t, err) + assert.Equal(t, int64(1048576), written) + }) +} From 00f1626f7aaf318f4b5fb436e70889abbd05363b Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Fri, 24 Jul 2026 13:47:39 +0800 Subject: [PATCH 6/7] block uploader restore implementation Signed-off-by: Lyndon-Li --- changelogs/unreleased/10071-Lyndon-Li | 1 + pkg/uploader/block/uploader.go | 8 ++--- pkg/uploader/block/uploader_test.go | 42 +++++++++++++-------------- 3 files changed, 26 insertions(+), 25 deletions(-) create mode 100644 changelogs/unreleased/10071-Lyndon-Li diff --git a/changelogs/unreleased/10071-Lyndon-Li b/changelogs/unreleased/10071-Lyndon-Li new file mode 100644 index 000000000..dd3454a4d --- /dev/null +++ b/changelogs/unreleased/10071-Lyndon-Li @@ -0,0 +1 @@ +Fix issue #9828, add implementation for block uploader restore \ No newline at end of file diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index 7717a7e7e..0378f4f5b 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -151,7 +151,7 @@ func (blkup *blockUploader) Restore(snapshot udmrepo.Snapshot, dest destInfo, bi meta, err := blkup.repoWriter.ReadMetadata(blkup.ctx, snapshot.RootObject.ID) if err != nil { - return 0, errors.Wrapf(err, "error readding snapshot metadata for %s", snapshot.Description) + return 0, errors.Wrapf(err, "error reading snapshot metadata for %s", snapshot.Description) } if len(meta.SubObjects) != 1 { @@ -376,7 +376,7 @@ func (blkup *blockUploader) restoreData(reader io.ReadSeeker, dest *os.File, bit offset, valid := bitmap.Next() var buffer []byte - var nextPos uint64 = uint64(0) + var nextPos = uint64(0) for valid { select { case <-blkup.ctx.Done(): @@ -513,13 +513,13 @@ func (blkup *blockUploader) restoreData(reader io.ReadSeeker, dest *os.File, bit return written, nil } -func (bu *blockUploader) flushZeroBlocks(dest *os.File, start int64, length int64, zeroBlock []byte, destPath string) error { +func (blkup *blockUploader) flushZeroBlocks(dest *os.File, start int64, length int64, zeroBlock []byte, destPath string) error { err := blkZeroOut(dest, start, length) if err == nil { return nil } - bu.log.WithError(err).Warnf("Failed to call zero out from dev %s, start %v, length %v. Fallback to conservative way", destPath, start, length) + blkup.log.WithError(err).Warnf("Failed to call zero out from dev %s, start %v, length %v. Fallback to conservative way", destPath, start, length) var written int64 for written < length { diff --git a/pkg/uploader/block/uploader_test.go b/pkg/uploader/block/uploader_test.go index 69b5efcb5..bb7c79c5a 100644 --- a/pkg/uploader/block/uploader_test.go +++ b/pkg/uploader/block/uploader_test.go @@ -506,7 +506,7 @@ func TestGetSourceSize(t *testing.T) { if tc.expectErr { assert.Error(t, err) } else { - assert.NoError(t, err) + require.NoError(t, err) assert.Equal(t, tc.expected, size) } }) @@ -515,22 +515,22 @@ func TestGetSourceSize(t *testing.T) { func TestFlushZeroBlocks(t *testing.T) { t.Run("success via write fallback", func(t *testing.T) { - f, err := os.CreateTemp("", "zerotest-*") + f, err := os.CreateTemp(t.TempDir(), "zerotest-*") require.NoError(t, err) defer os.Remove(f.Name()) defer f.Close() require.NoError(t, f.Truncate(2048)) - bu := &blockUploader{ + blkup := &blockUploader{ log: logrus.New(), } - bu.log.(*logrus.Logger).Out = io.Discard + blkup.log.(*logrus.Logger).Out = io.Discard zeroBlock := make([]byte, 1024) - err = bu.flushZeroBlocks(f, 0, 2048, zeroBlock, f.Name()) + err = blkup.flushZeroBlocks(f, 0, 2048, zeroBlock, f.Name()) - assert.NoError(t, err) + require.NoError(t, err) data, err := os.ReadFile(f.Name()) require.NoError(t, err) @@ -555,13 +555,13 @@ func TestRestoreData(t *testing.T) { ctx := context.Background() progress := &mockProgressUpdater{} progress.On("UpdateProgress", mock.Anything).Return() - bu := &blockUploader{ + blkup := &blockUploader{ ctx: ctx, progress: progress, log: logrus.New(), } - f, err := os.CreateTemp("", "restoretest-*") + f, err := os.CreateTemp(t.TempDir(), "restoretest-*") require.NoError(t, err) defer os.Remove(f.Name()) defer f.Close() @@ -577,8 +577,8 @@ func TestRestoreData(t *testing.T) { iterMock.On("Next").Return(uint64(0), true).Once() iterMock.On("Next").Return(uint64(0), false) - written, err := bu.restoreData(reader, f, iterMock, 1048576, f.Name()) - assert.NoError(t, err) + written, err := blkup.restoreData(reader, f, iterMock, 1048576, f.Name()) + require.NoError(t, err) assert.Equal(t, int64(1048576), written) f.Seek(0, 0) @@ -589,12 +589,12 @@ func TestRestoreData(t *testing.T) { t.Run("read err", func(t *testing.T) { ctx := context.Background() - bu := &blockUploader{ + blkup := &blockUploader{ ctx: ctx, log: logrus.New(), } - f, err := os.CreateTemp("", "restoretest-*") + f, err := os.CreateTemp(t.TempDir(), "restoretest-*") require.NoError(t, err) defer os.Remove(f.Name()) defer f.Close() @@ -606,8 +606,8 @@ func TestRestoreData(t *testing.T) { iterMock.On("Next").Return(uint64(0), true).Once() iterMock.On("Next").Return(uint64(0), false) - _, err = bu.restoreData(reader, f, iterMock, 1048576, f.Name()) - assert.Error(t, err) + _, err = blkup.restoreData(reader, f, iterMock, 1048576, f.Name()) + require.Error(t, err) assert.Contains(t, err.Error(), "read error") }) } @@ -616,13 +616,13 @@ func TestBlockUploaderRestore(t *testing.T) { t.Run("missing metadata", func(t *testing.T) { ctx := context.Background() repoWriter := udmrepomocks.NewBackupRepo(t) - bu := NewUploader(ctx, repoWriter, nil, logrus.New()) + blkup := NewUploader(ctx, repoWriter, nil, logrus.New()) repoWriter.On("ReadMetadata", mock.Anything, udmrepo.ID("root-id")).Return(nil, errors.New("meta not found")) iterMock := cbtmocks.NewIterator(t) - _, err := bu.Restore(udmrepo.Snapshot{RootObject: udmrepo.ObjectMetadata{ID: "root-id"}}, destInfo{}, iterMock, nil) - assert.Error(t, err) + _, err := blkup.Restore(udmrepo.Snapshot{RootObject: udmrepo.ObjectMetadata{ID: "root-id"}}, destInfo{}, iterMock, nil) + require.Error(t, err) assert.Contains(t, err.Error(), "meta not found") }) @@ -632,9 +632,9 @@ func TestBlockUploaderRestore(t *testing.T) { progress := &mockProgressUpdater{} progress.On("UpdateProgress", mock.Anything).Return() - bu := NewUploader(ctx, repoWriter, progress, logrus.New()) + blkup := NewUploader(ctx, repoWriter, progress, logrus.New()) - f, err := os.CreateTemp("", "restoretest-*") + f, err := os.CreateTemp(t.TempDir(), "restoretest-*") require.NoError(t, err) defer os.Remove(f.Name()) defer f.Close() @@ -682,8 +682,8 @@ func TestBlockUploaderRestore(t *testing.T) { iterMock.On("Next").Return(uint64(0), true).Once() iterMock.On("Next").Return(uint64(0), false) - written, err := bu.Restore(snap, dest, iterMock, nil) - assert.NoError(t, err) + written, err := blkup.Restore(snap, dest, iterMock, nil) + require.NoError(t, err) assert.Equal(t, int64(1048576), written) }) } From bc596da38666c8e813953b91db7f2fe94ebfe0f9 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Mon, 27 Jul 2026 14:19:22 +0800 Subject: [PATCH 7/7] block uploader restore implementation Signed-off-by: Lyndon-Li --- pkg/uploader/block/snapshot.go | 5 ++++- pkg/uploader/block/uploader.go | 2 +- pkg/uploader/block/uploader_test.go | 6 +++--- 3 files changed, 8 insertions(+), 5 deletions(-) diff --git a/pkg/uploader/block/snapshot.go b/pkg/uploader/block/snapshot.go index b185f4e15..53e7e7f14 100644 --- a/pkg/uploader/block/snapshot.go +++ b/pkg/uploader/block/snapshot.go @@ -121,7 +121,10 @@ func snapshotSource( return "", 0, errors.Wrapf(err, "Failed to run uploader backup for si %v", source) } - snap.Tags = make(map[string]string) + if snap.Tags == nil { + snap.Tags = make(map[string]string) + } + snap.Tags[uploader.CBTChangeIDTag] = cbtSource.ChangeID snap.Tags[uploader.CBTVolumeIDTag] = cbtSource.VolumeID if snapshotTags != nil { diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index 0378f4f5b..8aa58bf96 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -155,7 +155,7 @@ func (blkup *blockUploader) Restore(snapshot udmrepo.Snapshot, dest destInfo, bi } if len(meta.SubObjects) != 1 { - return 0, errors.Wrapf(err, "unexpected number of bdev object (%d) for snapshot %s", len(meta.SubObjects), snapshot.Description) + return 0, errors.Errorf("unexpected number of bdev object (%d) for snapshot %s", len(meta.SubObjects), snapshot.Description) } sourceSize, err := getSourceSize(snapshot) diff --git a/pkg/uploader/block/uploader_test.go b/pkg/uploader/block/uploader_test.go index bb7c79c5a..79c7be954 100644 --- a/pkg/uploader/block/uploader_test.go +++ b/pkg/uploader/block/uploader_test.go @@ -483,7 +483,7 @@ func TestGetSourceSize(t *testing.T) { name: "invalid tag value", snapshot: udmrepo.Snapshot{ Tags: map[string]string{ - "bdev-source-size": "abc", + bdevSourceSizeTag: "abc", }, }, expectErr: true, @@ -492,7 +492,7 @@ func TestGetSourceSize(t *testing.T) { name: "valid tag value", snapshot: udmrepo.Snapshot{ Tags: map[string]string{ - "bdev-source-size": "1048576", + bdevSourceSizeTag: "1048576", }, }, expectErr: false, @@ -667,7 +667,7 @@ func TestBlockUploaderRestore(t *testing.T) { Description: "test snapshot", RootObject: udmrepo.ObjectMetadata{ID: "root-id"}, Tags: map[string]string{ - "bdev-source-size": "1048576", + bdevSourceSizeTag: "1048576", }, }