From df21463629fd31c103aa6c2a2c553e2898b57662 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Thu, 14 May 2026 10:23:08 +0800 Subject: [PATCH 1/5] block uploader backup implementation Signed-off-by: Lyndon-Li --- pkg/repository/udmrepo/kopialib/lib_repo.go | 2 +- .../udmrepo/kopialib/lib_repo_ex_test.go | 2 +- pkg/uploader/block/uploader.go | 267 ++++++++++++- pkg/uploader/block/uploader_test.go | 353 +++++++++++++++++- .../kopialib => util}/freelist/freelist.go | 0 .../freelist/freelist_test.go | 0 6 files changed, 617 insertions(+), 7 deletions(-) rename pkg/{repository/udmrepo/kopialib => util}/freelist/freelist.go (100%) rename pkg/{repository/udmrepo/kopialib => util}/freelist/freelist_test.go (100%) diff --git a/pkg/repository/udmrepo/kopialib/lib_repo.go b/pkg/repository/udmrepo/kopialib/lib_repo.go index 29ff02eea..151bf1cb2 100644 --- a/pkg/repository/udmrepo/kopialib/lib_repo.go +++ b/pkg/repository/udmrepo/kopialib/lib_repo.go @@ -44,7 +44,7 @@ import ( "github.com/vmware-tanzu/velero/pkg/kopia" "github.com/vmware-tanzu/velero/pkg/repository/udmrepo" "github.com/vmware-tanzu/velero/pkg/repository/udmrepo/kopialib/backend" - "github.com/vmware-tanzu/velero/pkg/repository/udmrepo/kopialib/freelist" + "github.com/vmware-tanzu/velero/pkg/util/freelist" ) type kopiaRepoService struct { diff --git a/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go b/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go index fdaeb9f69..a42c02c14 100644 --- a/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go +++ b/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go @@ -15,8 +15,8 @@ import ( "github.com/vmware-tanzu/velero/pkg/repository/udmrepo" repomocks "github.com/vmware-tanzu/velero/pkg/repository/udmrepo/kopialib/backend/mocks" - "github.com/vmware-tanzu/velero/pkg/repository/udmrepo/kopialib/freelist" velerotest "github.com/vmware-tanzu/velero/pkg/test" + "github.com/vmware-tanzu/velero/pkg/util/freelist" ) type mockDirectRepository struct { diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index 233d72a17..318e4a7f3 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -18,7 +18,10 @@ package block import ( "context" + "io" "os" + "runtime" + "strings" "github.com/cockroachdb/errors" "github.com/sirupsen/logrus" @@ -26,12 +29,14 @@ import ( "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" + "github.com/vmware-tanzu/velero/pkg/util/freelist" ) var ErrCanceled = errors.New("uploader is canceled") const ( - blockSize = (1 << 20) + blockSize = (1 << 20) + bufferSize = 100 << 20 ) type sourceInfo struct { @@ -50,9 +55,265 @@ type Uploader interface { Restore(udmrepo.Snapshot, destInfo, cbt.Iterator, map[string]string) (int64, error) } -// implement in following PRs +type blockUploader struct { + ctx context.Context + repoWriter udmrepo.BackupRepo + progress uploader.ProgressUpdater + log logrus.FieldLogger +} + func NewUploader(ctx context.Context, repoWriter udmrepo.BackupRepo, progress uploader.ProgressUpdater, log logrus.FieldLogger) Uploader { - return nil + return &blockUploader{ + ctx: ctx, + repoWriter: repoWriter, + progress: progress, + log: log, + } +} + +func (bu *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, bitmap cbt.Iterator, configs map[string]string) (udmrepo.Snapshot, int64, error) { + snapStart := bu.repoWriter.Time() + + if bitmap == nil { + return udmrepo.Snapshot{}, 0, errors.New("bitmap is not available") + } + + backupMode := udmrepo.ObjectDataBackupModeInc + if parentObject == "" { + backupMode = udmrepo.ObjectDataBackupModeFull + } + + destObj, err := bu.repoWriter.NewObjectWriter(bu.ctx, udmrepo.ObjectWriteOptions{ + Description: "BDEV:" + getObjectName(source.realSource), + DataType: udmrepo.ObjectDataTypeData, + AccessMode: udmrepo.ObjectDataAccessModeBlock, + ParentObject: parentObject, + BackupMode: backupMode, + AsyncWrites: runtime.NumCPU(), + }) + if err != nil { + return udmrepo.Snapshot{}, 0, errors.Wrap(err, "error creating object writer") + } + + defer destObj.Close() + + id, backupSize, objectSize, err := bu.backupObject(source.dev, destObj, bitmap, source.size) + if err != nil { + return udmrepo.Snapshot{}, 0, errors.Wrap(err, "error to backup file with incremental") + } + + entryId, err := bu.repoWriter.WriteMetadata(bu.ctx, &udmrepo.Metadata{ + SubObjects: []udmrepo.ObjectMetadata{ + { + ID: id, + Name: getObjectName(source.realSource), + Type: udmrepo.ObjectDataTypeData, + Size: objectSize, + Permissions: 0o777, + }, + }, + }, + udmrepo.ObjectWriteOptions{ + Description: "bdev-root", + }) + if err != nil { + return udmrepo.Snapshot{}, 0, errors.Wrap(err, "error to write metadata") + } + + snapEnd := bu.repoWriter.Time() + + return udmrepo.Snapshot{ + Source: source.realSource, + StartTime: snapStart, + EndTime: snapEnd, + Description: source.realSource, + RootObject: udmrepo.ObjectMetadata{ + ID: entryId, + Name: "bdev-root", + Type: udmrepo.ObjectDataTypeMetadata, + Permissions: 0o777, + }, + }, backupSize, nil +} + +// TODO implement in following PRs +func (bu *blockUploader) Restore(snapshot udmrepo.Snapshot, dest destInfo, bitmap cbt.Iterator, configs map[string]string) (int64, error) { + return 0, nil +} + +func (bu *blockUploader) backupObject(dev *os.File, dest udmrepo.ObjectWriter, bitmap cbt.Iterator, totalLength int64) (udmrepo.ID, int64, int64, error) { + backupSize, objectSize, err := bu.backupData(dev, dest, bitmap, totalLength) + if err != nil { + return "", backupSize, objectSize, errors.Wrap(err, "error copying file data incremental") + } + + id, err := dest.Result() + return id, backupSize, objectSize, err +} + +type readResult struct { + buffer []byte + offset int64 + err error +} + +func (r *readResult) resetBuffer(list *freelist.FreeList) { + if r.buffer != nil { + list.Return(r.buffer) + r.buffer = nil + } +} + +func (bu *blockUploader) backupData(reader io.ReaderAt, writer udmrepo.ObjectWriter, bitmap cbt.Iterator, totalLength int64) (int64, int64, error) { + blockSize := bitmap.BlockSize() + list := freelist.New(bufferSize, int(blockSize)) + resultChan := make(chan readResult, list.Capacity()) + totalCount := bitmap.Count() + aligned := (totalLength + int64(blockSize) - 1) / int64(blockSize) * int64(blockSize) + + quit := make(chan struct{}) + defer close(quit) + + go func() { + defer close(resultChan) + + offset, valid := bitmap.Next() + var buffer []byte + for valid { + select { + case <-bu.ctx.Done(): + return + case <-quit: + return + case buffer = <-list.Chunks(): + } + + length := blockSize + if offset+uint64(length) > uint64(totalLength) { + length = uint(uint64(totalLength) - offset) + clear(buffer) + } + + readBytes, err := reader.ReadAt(buffer[:length], int64(offset)) + if err == nil && readBytes <= 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 + } + + offset, valid = bitmap.Next() + } + }() + + var lastPos int64 + var result readResult + var written int64 + var curCount int64 + var writeErr error + var readerRunning bool + + for curCount < int64(totalCount) { + select { + case <-bu.ctx.Done(): + writeErr = ErrCanceled + case result, readerRunning = <-resultChan: + if !readerRunning { + if bu.ctx.Err() != nil { + writeErr = ErrCanceled + } else { + writeErr = io.ErrUnexpectedEOF + } + } + } + + if writeErr != nil { + break + } + + if result.err != nil { + writeErr = result.err + break + } + + n, err := writer.WriteAt(result.buffer, result.offset) + if err != nil { + writeErr = err + break + } + + if blockSize != uint(n) { + writeErr = io.ErrShortWrite + break + } + + written += int64(blockSize) + lastPos = result.offset + int64(blockSize) + result.resetBuffer(list) + curCount++ + + bu.progress.UpdateProgress(&uploader.Progress{BytesDone: lastPos, TotalBytes: aligned}) + } + + result.resetBuffer(list) + + if writeErr != nil { + return written, aligned, writeErr + } + + if lastPos < aligned { + s, err := copyTailData(reader, writer, totalLength, int64(blockSize)) + if err != nil { + return written, aligned, errors.Wrapf(err, "unable to write tail data at %v", lastPos) + } + + written += s + + bu.progress.UpdateProgress(&uploader.Progress{BytesDone: aligned, TotalBytes: aligned}) + } + + return written, aligned, nil +} + +func copyTailData(source io.ReaderAt, writer udmrepo.ObjectWriter, totalLength int64, blockSize int64) (int64, error) { + roundUp := (totalLength + blockSize - 1) / blockSize * blockSize + roundDown := totalLength / blockSize * blockSize + length := totalLength - roundDown + + if length == 0 { + if _, err := writer.WriteAt(nil, roundUp); err != nil { + return -1, errors.Wrapf(err, "error writing sparse to %v", roundUp) + } + } else { + buffer := make([]byte, blockSize) + if _, err := source.ReadAt(buffer[:length], roundDown); err != nil { + return -1, errors.Wrapf(err, "error reading tail data with length %v", length) + } + + if _, err := writer.WriteAt(buffer, roundDown); err != nil { + return -1, errors.Wrapf(err, "error writing tail data at %v", roundDown) + } + } + + return length, nil +} + +func getObjectName(source string) string { + s := strings.ReplaceAll(source, "/", "-") + return strings.ReplaceAll(s, "\\", "-") } func loadObjectFromSnapshot(ctx context.Context, rep udmrepo.BackupRepo, snapshot *udmrepo.Snapshot) (udmrepo.ID, error) { diff --git a/pkg/uploader/block/uploader_test.go b/pkg/uploader/block/uploader_test.go index 8209569e1..d6e2e3d90 100644 --- a/pkg/uploader/block/uploader_test.go +++ b/pkg/uploader/block/uploader_test.go @@ -5,7 +5,7 @@ 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 + 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, @@ -17,18 +17,367 @@ limitations under the License. package block import ( + "bytes" "context" + "io" + "os" "testing" + "time" - "github.com/cockroachdb/errors" + "github.com/pkg/errors" + "github.com/sirupsen/logrus" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" "github.com/vmware-tanzu/velero/pkg/repository/udmrepo" udmrepomocks "github.com/vmware-tanzu/velero/pkg/repository/udmrepo/mocks" + "github.com/vmware-tanzu/velero/pkg/uploader" + cbt "github.com/vmware-tanzu/velero/pkg/uploader/cbt/types" + cbtmocks "github.com/vmware-tanzu/velero/pkg/uploader/cbt/types/mocks" ) +type mockProgressUpdater struct { + mock.Mock +} + +func (m *mockProgressUpdater) UpdateProgress(p *uploader.Progress) { + m.Called(p) +} + +func TestNewUploader(t *testing.T) { + ctx := context.Background() + repoWriter := udmrepomocks.NewBackupRepo(t) + progress := &mockProgressUpdater{} + log := logrus.New() + + uploader := NewUploader(ctx, repoWriter, progress, log) + + bu, ok := uploader.(*blockUploader) + assert.True(t, ok) + assert.Equal(t, ctx, bu.ctx) + assert.Equal(t, repoWriter, bu.repoWriter) + assert.Equal(t, progress, bu.progress) + assert.Equal(t, log, bu.log) +} + +func TestGetObjectName(t *testing.T) { + testCases := []struct { + name string + source string + expected string + }{ + { + name: "no slashes", + source: "test", + expected: "test", + }, + { + name: "unix path", + source: "/var/lib/kubelet/pods/uuid/volumes/test", + expected: "-var-lib-kubelet-pods-uuid-volumes-test", + }, + { + name: "windows path", + source: `c:\var\lib\kubelet\pods\uuid\volumes\test`, + expected: `c:-var-lib-kubelet-pods-uuid-volumes-test`, + }, + { + name: "mixed slashes", + source: `c:\var/lib\kubelet/pods`, + expected: `c:-var-lib-kubelet-pods`, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + result := getObjectName(tc.source) + assert.Equal(t, tc.expected, result) + }) + } +} + +func TestCopyTailData(t *testing.T) { + testCases := []struct { + name string + totalLength int64 + blockSize int64 + sourceData []byte + writeErr error + readErr error + expected int64 + expectErr bool + }{ + { + name: "tail length 0", + totalLength: 2048, + blockSize: 1024, + expected: 0, + }, + { + name: "tail length 512 with 1024 block size", + totalLength: 1536, + blockSize: 1024, + sourceData: make([]byte, 1536), + expected: 512, + }, + { + name: "tail length with write error", + totalLength: 1536, + blockSize: 1024, + sourceData: make([]byte, 1536), + writeErr: errors.New("write error"), + expectErr: true, + }, + { + name: "tail length 0 with sparse write error", + totalLength: 2048, + blockSize: 1024, + writeErr: errors.New("write error"), + expectErr: true, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + writer := udmrepomocks.NewObjectWriter(t) + var source io.ReaderAt + + if tc.totalLength%tc.blockSize == 0 { + writer.On("WriteAt", []byte(nil), tc.totalLength).Return(0, tc.writeErr) + } else { + length := tc.totalLength - (tc.totalLength/tc.blockSize)*tc.blockSize + paddedData := make([]byte, tc.blockSize) + copy(paddedData[:length], tc.sourceData) + + source = bytes.NewReader(tc.sourceData) + writer.On("WriteAt", paddedData, (tc.totalLength/tc.blockSize)*tc.blockSize).Return(int(tc.blockSize), tc.writeErr) + } + + n, err := copyTailData(source, writer, tc.totalLength, tc.blockSize) + if tc.expectErr { + assert.Error(t, err) + } else { + assert.NoError(t, err) + assert.Equal(t, tc.expected, n) + } + }) + } +} + +func TestBlockUploaderBackup(t *testing.T) { + testCases := []struct { + name string + nilBitmap bool + createObjErr error + writeMetaErr error + writeObjErr error + parentObj udmrepo.ID + cancelCtx bool + cancelInProgress bool + readDataErr bool + shortWrite bool + fewerBlocks bool + expectErr bool + expectErrStr string + }{ + { + name: "nil bitmap", + nilBitmap: true, + expectErr: true, + }, + { + name: "canceled context", + cancelCtx: true, + expectErr: true, + expectErrStr: "uploader is canceled", + }, + { + name: "canceled in progress", + cancelInProgress: true, + expectErr: true, + expectErrStr: "error copying file data incremental: uploader is canceled", + }, + { + name: "create object writer err", + createObjErr: errors.New("create obj err"), + expectErr: true, + }, + { + name: "read data err", + readDataErr: true, + expectErr: true, + expectErrStr: "EOF", + }, + { + name: "short write err", + shortWrite: true, + expectErr: true, + expectErrStr: "short write", + }, + { + name: "unexpected EOF fewer blocks", + fewerBlocks: true, + expectErr: true, + expectErrStr: "unexpected EOF", + }, + { + name: "write meta err", + writeMetaErr: errors.New("write meta err"), + expectErr: true, + }, + { + name: "success full backup", + parentObj: "", + expectErr: false, + }, + { + name: "success inc backup", + parentObj: "parent-01", + expectErr: false, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ctx := context.Background() + var cancel context.CancelFunc + ctx, cancel = context.WithCancel(ctx) + + if tc.cancelCtx { + cancel() + } else if tc.cancelInProgress { + go func() { + time.Sleep(100 * time.Millisecond) + cancel() + }() + } else { + defer cancel() + } + + repoWriter := udmrepomocks.NewBackupRepo(t) + progress := &mockProgressUpdater{} + progress.On("UpdateProgress", mock.Anything).Return() + log := logrus.New() + log.Out = io.Discard + + bu := NewUploader(ctx, repoWriter, progress, log) + + f, err := os.CreateTemp("", "blktest-*") + require.NoError(t, err) + defer os.Remove(f.Name()) + defer f.Close() + + if tc.cancelInProgress { + require.NoError(t, f.Truncate(2*1048576)) + } else if tc.readDataErr { + // Don't truncate so that reading hits EOF immediately + } else { + require.NoError(t, f.Truncate(1048576)) + } + + fi, err := f.Stat() + require.NoError(t, err) + + srcInfo := sourceInfo{ + dev: f, + realSource: "/data/volume1", + size: fi.Size(), + } + + if tc.readDataErr { + srcInfo.size = 1048576 + } + + repoWriter.On("Time").Return(time.Now()) + + var iterator cbt.Iterator + if !tc.nilBitmap { + iterMock := cbtmocks.NewIterator(t) + iterator = iterMock + + backupMode := udmrepo.ObjectDataBackupModeInc + if tc.parentObj == "" { + backupMode = udmrepo.ObjectDataBackupModeFull + } + + objWriter := udmrepomocks.NewObjectWriter(t) + if tc.createObjErr == nil { + objWriter.On("Close").Return(nil) + + if tc.cancelInProgress { + iterMock.On("BlockSize").Return(uint(1048576)) + iterMock.On("Count").Return(uint64(1000)) + iterMock.On("Next").Return(uint64(0), true) + + objWriter.On("WriteAt", mock.Anything, mock.Anything).Run(func(args mock.Arguments) { + <-ctx.Done() + }).Return(1048576, nil) + objWriter.On("Result").Return(udmrepo.ID(""), errors.New("write failed")).Maybe() + } else if tc.cancelCtx { + iterMock.On("BlockSize").Return(uint(1048576)) + iterMock.On("Count").Return(uint64(1)) + iterMock.On("Next").Return(uint64(0), true) + + objWriter.On("Result").Return(udmrepo.ID(""), errors.New("write failed")).Maybe() + } else if tc.shortWrite { + iterMock.On("BlockSize").Return(uint(1048576)) + iterMock.On("Count").Return(uint64(1)) + iterMock.On("Next").Return(uint64(0), true) + + objWriter.On("WriteAt", mock.Anything, mock.Anything).Return(512, nil) + objWriter.On("Result").Return(udmrepo.ID(""), errors.New("write failed")).Maybe() + } else if tc.fewerBlocks { + iterMock.On("BlockSize").Return(uint(1048576)) + iterMock.On("Count").Return(uint64(5)) + iterMock.On("Next").Return(uint64(0), false) + + objWriter.On("Result").Return(udmrepo.ID(""), errors.New("write failed")).Maybe() + } else if tc.readDataErr { + iterMock.On("BlockSize").Return(uint(1048576)) + iterMock.On("Count").Return(uint64(1)) + iterMock.On("Next").Return(uint64(0), true) + + objWriter.On("Result").Return(udmrepo.ID(""), errors.New("write failed")).Maybe() + } else { + // Setup backupData sequence: next returns false immediately + iterMock.On("BlockSize").Return(uint(1048576)) + iterMock.On("Count").Return(uint64(0)) + iterMock.On("Next").Return(uint64(0), false) + + if tc.writeObjErr != nil { + objWriter.On("WriteAt", mock.Anything, mock.Anything).Return(0, tc.writeObjErr) + objWriter.On("Result").Return(udmrepo.ID(""), errors.New("write failed")) + } else { + objWriter.On("WriteAt", mock.Anything, mock.Anything).Return(1048576, nil) + objWriter.On("Result").Return(udmrepo.ID("obj-01"), nil) + repoWriter.On("WriteMetadata", mock.Anything, mock.Anything, mock.Anything).Return(udmrepo.ID("meta-01"), tc.writeMetaErr) + } + } + } + + repoWriter.On("NewObjectWriter", mock.Anything, mock.MatchedBy(func(opt udmrepo.ObjectWriteOptions) bool { + return opt.Description == "BDEV:-data-volume1" && opt.BackupMode == backupMode + })).Return(objWriter, tc.createObjErr) + } + + snap, size, err := bu.Backup(srcInfo, tc.parentObj, iterator, nil) + + if tc.expectErr { + assert.Error(t, err) + if tc.expectErrStr != "" { + assert.Contains(t, err.Error(), tc.expectErrStr) + } + } else { + assert.NoError(t, err) + assert.Equal(t, "/data/volume1", snap.Source) + assert.Equal(t, udmrepo.ID("meta-01"), snap.RootObject.ID) + assert.Equal(t, int64(0), size) + } + }) + } +} + func TestLoadObjectFromSnapshot(t *testing.T) { testCases := []struct { name string diff --git a/pkg/repository/udmrepo/kopialib/freelist/freelist.go b/pkg/util/freelist/freelist.go similarity index 100% rename from pkg/repository/udmrepo/kopialib/freelist/freelist.go rename to pkg/util/freelist/freelist.go diff --git a/pkg/repository/udmrepo/kopialib/freelist/freelist_test.go b/pkg/util/freelist/freelist_test.go similarity index 100% rename from pkg/repository/udmrepo/kopialib/freelist/freelist_test.go rename to pkg/util/freelist/freelist_test.go From 39c745ef612eb81b77ebe64d854f8d27573b5e53 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Tue, 2 Jun 2026 17:25:41 +0800 Subject: [PATCH 2/5] set totalSize from uploader Signed-off-by: Lyndon-Li --- pkg/uploader/block/uploader.go | 1 + 1 file changed, 1 insertion(+) diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index 318e4a7f3..d58b004b9 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -127,6 +127,7 @@ func (bu *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, bitm StartTime: snapStart, EndTime: snapEnd, Description: source.realSource, + TotalSize: objectSize, RootObject: udmrepo.ObjectMetadata{ ID: entryId, Name: "bdev-root", From 502ec5f08666429c41011ffc4fc90bde3b4c344d Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Thu, 4 Jun 2026 15:56:57 +0800 Subject: [PATCH 3/5] remove leading and trailing separator Signed-off-by: Lyndon-Li --- pkg/uploader/block/uploader.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index d58b004b9..bb8195e31 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -314,7 +314,8 @@ func copyTailData(source io.ReaderAt, writer udmrepo.ObjectWriter, totalLength i func getObjectName(source string) string { s := strings.ReplaceAll(source, "/", "-") - return strings.ReplaceAll(s, "\\", "-") + s = strings.ReplaceAll(s, "\\", "-") + return strings.Trim(s, "-") } func loadObjectFromSnapshot(ctx context.Context, rep udmrepo.BackupRepo, snapshot *udmrepo.Snapshot) (udmrepo.ID, error) { From 8d23c7e813183f3db67cd04f9a047265d5e93638 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Tue, 30 Jun 2026 16:27:18 +0800 Subject: [PATCH 4/5] block uploader backup implementation Signed-off-by: Lyndon-Li --- .../udmrepo/kopialib/lib_repo_ex_test.go | 16 ++++++++++++++++ pkg/uploader/block/uploader.go | 10 +++++----- pkg/uploader/block/uploader_test.go | 6 +++--- 3 files changed, 24 insertions(+), 8 deletions(-) diff --git a/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go b/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go index a42c02c14..3294063a6 100644 --- a/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go +++ b/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go @@ -1,3 +1,19 @@ +/* +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 kopialib import ( diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index bb8195e31..9d4dde9cb 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -5,7 +5,7 @@ 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 + 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, @@ -99,7 +99,7 @@ func (bu *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, bitm id, backupSize, objectSize, err := bu.backupObject(source.dev, destObj, bitmap, source.size) if err != nil { - return udmrepo.Snapshot{}, 0, errors.Wrap(err, "error to backup file with incremental") + return udmrepo.Snapshot{}, 0, errors.Wrapf(err, "error backing up bdev %s", source.realSource) } entryId, err := bu.repoWriter.WriteMetadata(bu.ctx, &udmrepo.Metadata{ @@ -117,7 +117,7 @@ func (bu *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, bitm Description: "bdev-root", }) if err != nil { - return udmrepo.Snapshot{}, 0, errors.Wrap(err, "error to write metadata") + return udmrepo.Snapshot{}, 0, errors.Wrap(err, "error writing metadata") } snapEnd := bu.repoWriter.Time() @@ -139,13 +139,13 @@ func (bu *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, bitm // TODO implement in following PRs func (bu *blockUploader) Restore(snapshot udmrepo.Snapshot, dest destInfo, bitmap cbt.Iterator, configs map[string]string) (int64, error) { - return 0, nil + return 0, errors.New("not implemented") } func (bu *blockUploader) backupObject(dev *os.File, dest udmrepo.ObjectWriter, bitmap cbt.Iterator, totalLength int64) (udmrepo.ID, int64, int64, error) { backupSize, objectSize, err := bu.backupData(dev, dest, bitmap, totalLength) if err != nil { - return "", backupSize, objectSize, errors.Wrap(err, "error copying file data incremental") + return "", backupSize, objectSize, err } id, err := dest.Result() diff --git a/pkg/uploader/block/uploader_test.go b/pkg/uploader/block/uploader_test.go index d6e2e3d90..2032e31ad 100644 --- a/pkg/uploader/block/uploader_test.go +++ b/pkg/uploader/block/uploader_test.go @@ -75,7 +75,7 @@ func TestGetObjectName(t *testing.T) { { name: "unix path", source: "/var/lib/kubelet/pods/uuid/volumes/test", - expected: "-var-lib-kubelet-pods-uuid-volumes-test", + expected: "var-lib-kubelet-pods-uuid-volumes-test", }, { name: "windows path", @@ -196,7 +196,7 @@ func TestBlockUploaderBackup(t *testing.T) { name: "canceled in progress", cancelInProgress: true, expectErr: true, - expectErrStr: "error copying file data incremental: uploader is canceled", + expectErrStr: "error backing up bdev /data/volume1: uploader is canceled", }, { name: "create object writer err", @@ -357,7 +357,7 @@ func TestBlockUploaderBackup(t *testing.T) { } repoWriter.On("NewObjectWriter", mock.Anything, mock.MatchedBy(func(opt udmrepo.ObjectWriteOptions) bool { - return opt.Description == "BDEV:-data-volume1" && opt.BackupMode == backupMode + return opt.Description == "BDEV:data-volume1" && opt.BackupMode == backupMode })).Return(objWriter, tc.createObjErr) } From 84bee825758ef06c288fd3744c528023f615833e Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Thu, 9 Jul 2026 11:49:18 +0800 Subject: [PATCH 5/5] block uploader backup implementation Signed-off-by: Lyndon-Li --- changelogs/unreleased/9979-Lyndon-Li | 1 + pkg/uploader/block/uploader.go | 32 ++++++++++++++-------------- pkg/uploader/block/uploader_test.go | 26 +++++++++++----------- 3 files changed, 30 insertions(+), 29 deletions(-) create mode 100644 changelogs/unreleased/9979-Lyndon-Li diff --git a/changelogs/unreleased/9979-Lyndon-Li b/changelogs/unreleased/9979-Lyndon-Li new file mode 100644 index 000000000..78134da35 --- /dev/null +++ b/changelogs/unreleased/9979-Lyndon-Li @@ -0,0 +1 @@ +Add the backup implementation for block data mover \ No newline at end of file diff --git a/pkg/uploader/block/uploader.go b/pkg/uploader/block/uploader.go index 9d4dde9cb..75e913cb7 100644 --- a/pkg/uploader/block/uploader.go +++ b/pkg/uploader/block/uploader.go @@ -71,8 +71,8 @@ func NewUploader(ctx context.Context, repoWriter udmrepo.BackupRepo, progress up } } -func (bu *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, bitmap cbt.Iterator, configs map[string]string) (udmrepo.Snapshot, int64, error) { - snapStart := bu.repoWriter.Time() +func (blkup *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, bitmap cbt.Iterator, configs map[string]string) (udmrepo.Snapshot, int64, error) { + snapStart := blkup.repoWriter.Time() if bitmap == nil { return udmrepo.Snapshot{}, 0, errors.New("bitmap is not available") @@ -83,7 +83,7 @@ func (bu *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, bitm backupMode = udmrepo.ObjectDataBackupModeFull } - destObj, err := bu.repoWriter.NewObjectWriter(bu.ctx, udmrepo.ObjectWriteOptions{ + destObj, err := blkup.repoWriter.NewObjectWriter(blkup.ctx, udmrepo.ObjectWriteOptions{ Description: "BDEV:" + getObjectName(source.realSource), DataType: udmrepo.ObjectDataTypeData, AccessMode: udmrepo.ObjectDataAccessModeBlock, @@ -97,12 +97,12 @@ func (bu *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, bitm defer destObj.Close() - id, backupSize, objectSize, err := bu.backupObject(source.dev, destObj, bitmap, source.size) + id, backupSize, objectSize, err := blkup.backupObject(source.dev, destObj, bitmap, source.size) if err != nil { return udmrepo.Snapshot{}, 0, errors.Wrapf(err, "error backing up bdev %s", source.realSource) } - entryId, err := bu.repoWriter.WriteMetadata(bu.ctx, &udmrepo.Metadata{ + entryID, err := blkup.repoWriter.WriteMetadata(blkup.ctx, &udmrepo.Metadata{ SubObjects: []udmrepo.ObjectMetadata{ { ID: id, @@ -120,7 +120,7 @@ func (bu *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, bitm return udmrepo.Snapshot{}, 0, errors.Wrap(err, "error writing metadata") } - snapEnd := bu.repoWriter.Time() + snapEnd := blkup.repoWriter.Time() return udmrepo.Snapshot{ Source: source.realSource, @@ -129,7 +129,7 @@ func (bu *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, bitm Description: source.realSource, TotalSize: objectSize, RootObject: udmrepo.ObjectMetadata{ - ID: entryId, + ID: entryID, Name: "bdev-root", Type: udmrepo.ObjectDataTypeMetadata, Permissions: 0o777, @@ -138,12 +138,12 @@ func (bu *blockUploader) Backup(source sourceInfo, parentObject udmrepo.ID, bitm } // TODO implement in following PRs -func (bu *blockUploader) Restore(snapshot udmrepo.Snapshot, dest destInfo, bitmap cbt.Iterator, configs map[string]string) (int64, error) { +func (blkup *blockUploader) Restore(snapshot udmrepo.Snapshot, dest destInfo, bitmap cbt.Iterator, configs map[string]string) (int64, error) { return 0, errors.New("not implemented") } -func (bu *blockUploader) backupObject(dev *os.File, dest udmrepo.ObjectWriter, bitmap cbt.Iterator, totalLength int64) (udmrepo.ID, int64, int64, error) { - backupSize, objectSize, err := bu.backupData(dev, dest, bitmap, totalLength) +func (blkup *blockUploader) backupObject(dev *os.File, dest udmrepo.ObjectWriter, bitmap cbt.Iterator, totalLength int64) (udmrepo.ID, int64, int64, error) { + backupSize, objectSize, err := blkup.backupData(dev, dest, bitmap, totalLength) if err != nil { return "", backupSize, objectSize, err } @@ -165,7 +165,7 @@ func (r *readResult) resetBuffer(list *freelist.FreeList) { } } -func (bu *blockUploader) backupData(reader io.ReaderAt, writer udmrepo.ObjectWriter, bitmap cbt.Iterator, totalLength int64) (int64, int64, error) { +func (blkup *blockUploader) backupData(reader io.ReaderAt, writer udmrepo.ObjectWriter, bitmap cbt.Iterator, totalLength int64) (int64, int64, error) { blockSize := bitmap.BlockSize() list := freelist.New(bufferSize, int(blockSize)) resultChan := make(chan readResult, list.Capacity()) @@ -182,7 +182,7 @@ func (bu *blockUploader) backupData(reader io.ReaderAt, writer udmrepo.ObjectWri var buffer []byte for valid { select { - case <-bu.ctx.Done(): + case <-blkup.ctx.Done(): return case <-quit: return @@ -229,11 +229,11 @@ func (bu *blockUploader) backupData(reader io.ReaderAt, writer udmrepo.ObjectWri for curCount < int64(totalCount) { select { - case <-bu.ctx.Done(): + case <-blkup.ctx.Done(): writeErr = ErrCanceled case result, readerRunning = <-resultChan: if !readerRunning { - if bu.ctx.Err() != nil { + if blkup.ctx.Err() != nil { writeErr = ErrCanceled } else { writeErr = io.ErrUnexpectedEOF @@ -266,7 +266,7 @@ func (bu *blockUploader) backupData(reader io.ReaderAt, writer udmrepo.ObjectWri result.resetBuffer(list) curCount++ - bu.progress.UpdateProgress(&uploader.Progress{BytesDone: lastPos, TotalBytes: aligned}) + blkup.progress.UpdateProgress(&uploader.Progress{BytesDone: lastPos, TotalBytes: aligned}) } result.resetBuffer(list) @@ -283,7 +283,7 @@ func (bu *blockUploader) backupData(reader io.ReaderAt, writer udmrepo.ObjectWri written += s - bu.progress.UpdateProgress(&uploader.Progress{BytesDone: aligned, TotalBytes: aligned}) + blkup.progress.UpdateProgress(&uploader.Progress{BytesDone: aligned, TotalBytes: aligned}) } return written, aligned, nil diff --git a/pkg/uploader/block/uploader_test.go b/pkg/uploader/block/uploader_test.go index 2032e31ad..88fd4771e 100644 --- a/pkg/uploader/block/uploader_test.go +++ b/pkg/uploader/block/uploader_test.go @@ -24,7 +24,7 @@ import ( "testing" "time" - "github.com/pkg/errors" + "github.com/cockroachdb/errors" "github.com/sirupsen/logrus" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/mock" @@ -53,12 +53,12 @@ func TestNewUploader(t *testing.T) { uploader := NewUploader(ctx, repoWriter, progress, log) - bu, ok := uploader.(*blockUploader) + blkup, ok := uploader.(*blockUploader) assert.True(t, ok) - assert.Equal(t, ctx, bu.ctx) - assert.Equal(t, repoWriter, bu.repoWriter) - assert.Equal(t, progress, bu.progress) - assert.Equal(t, log, bu.log) + assert.Equal(t, ctx, blkup.ctx) + assert.Equal(t, repoWriter, blkup.repoWriter) + assert.Equal(t, progress, blkup.progress) + assert.Equal(t, log, blkup.log) } func TestGetObjectName(t *testing.T) { @@ -158,7 +158,7 @@ func TestCopyTailData(t *testing.T) { if tc.expectErr { assert.Error(t, err) } else { - assert.NoError(t, err) + require.NoError(t, err) assert.Equal(t, tc.expected, n) } }) @@ -261,9 +261,9 @@ func TestBlockUploaderBackup(t *testing.T) { log := logrus.New() log.Out = io.Discard - bu := NewUploader(ctx, repoWriter, progress, log) + blkup := NewUploader(ctx, repoWriter, progress, log) - f, err := os.CreateTemp("", "blktest-*") + f, err := os.CreateTemp(t.TempDir(), "blktest-*") require.NoError(t, err) defer os.Remove(f.Name()) defer f.Close() @@ -317,7 +317,7 @@ func TestBlockUploaderBackup(t *testing.T) { } else if tc.cancelCtx { iterMock.On("BlockSize").Return(uint(1048576)) iterMock.On("Count").Return(uint64(1)) - iterMock.On("Next").Return(uint64(0), true) + iterMock.On("Next").Return(uint64(0), true).Maybe() objWriter.On("Result").Return(udmrepo.ID(""), errors.New("write failed")).Maybe() } else if tc.shortWrite { @@ -361,15 +361,15 @@ func TestBlockUploaderBackup(t *testing.T) { })).Return(objWriter, tc.createObjErr) } - snap, size, err := bu.Backup(srcInfo, tc.parentObj, iterator, nil) + snap, size, err := blkup.Backup(srcInfo, tc.parentObj, iterator, nil) if tc.expectErr { - assert.Error(t, err) + require.Error(t, err) if tc.expectErrStr != "" { assert.Contains(t, err.Error(), tc.expectErrStr) } } else { - assert.NoError(t, err) + require.NoError(t, err) assert.Equal(t, "/data/volume1", snap.Source) assert.Equal(t, udmrepo.ID("meta-01"), snap.RootObject.ID) assert.Equal(t, int64(0), size)