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/snapshot.go b/pkg/uploader/block/snapshot.go index e30f5c1bb..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 { @@ -222,7 +225,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..8aa58bf96 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -17,10 +17,12 @@ limitations under the License. package block import ( + "bytes" "context" "io" "os" "runtime" + "strconv" "strings" "github.com/cockroachdb/errors" @@ -35,8 +37,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 +51,7 @@ type sourceInfo struct { type destInfo struct { dev *os.File path string + size int64 } type Uploader interface { @@ -134,12 +138,52 @@ 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 } -// 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 reading snapshot metadata for %s", snapshot.Description) + } + + if len(meta.SubObjects) != 1 { + return 0, errors.Errorf("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) { @@ -318,6 +362,202 @@ 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) { + 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(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 (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 + } + + 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 { + 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 +} + +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") diff --git a/pkg/uploader/block/uploader_test.go b/pkg/uploader/block/uploader_test.go index 88fd4771e..79c7be954 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{ + bdevSourceSizeTag: "abc", + }, + }, + expectErr: true, + }, + { + name: "valid tag value", + snapshot: udmrepo.Snapshot{ + Tags: map[string]string{ + bdevSourceSizeTag: "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 { + require.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(t.TempDir(), "zerotest-*") + require.NoError(t, err) + defer os.Remove(f.Name()) + defer f.Close() + + require.NoError(t, f.Truncate(2048)) + + blkup := &blockUploader{ + log: logrus.New(), + } + blkup.log.(*logrus.Logger).Out = io.Discard + + zeroBlock := make([]byte, 1024) + err = blkup.flushZeroBlocks(f, 0, 2048, zeroBlock, f.Name()) + + require.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() + blkup := &blockUploader{ + ctx: ctx, + progress: progress, + log: logrus.New(), + } + + f, err := os.CreateTemp(t.TempDir(), "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 := blkup.restoreData(reader, f, iterMock, 1048576, f.Name()) + require.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() + blkup := &blockUploader{ + ctx: ctx, + log: logrus.New(), + } + + f, err := os.CreateTemp(t.TempDir(), "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 = blkup.restoreData(reader, f, iterMock, 1048576, f.Name()) + require.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) + 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 := 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") + }) + + t.Run("success", func(t *testing.T) { + ctx := context.Background() + repoWriter := udmrepomocks.NewBackupRepo(t) + progress := &mockProgressUpdater{} + progress.On("UpdateProgress", mock.Anything).Return() + + blkup := NewUploader(ctx, repoWriter, progress, logrus.New()) + + f, err := os.CreateTemp(t.TempDir(), "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{ + bdevSourceSizeTag: "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 := blkup.Restore(snap, dest, iterMock, nil) + require.NoError(t, err) + assert.Equal(t, int64(1048576), written) + }) +}