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/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..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 ( @@ -15,8 +31,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..75e913cb7 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, @@ -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,267 @@ 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 (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") + } + + backupMode := udmrepo.ObjectDataBackupModeInc + if parentObject == "" { + backupMode = udmrepo.ObjectDataBackupModeFull + } + + destObj, err := blkup.repoWriter.NewObjectWriter(blkup.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 := 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 := blkup.repoWriter.WriteMetadata(blkup.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 writing metadata") + } + + snapEnd := blkup.repoWriter.Time() + + return udmrepo.Snapshot{ + Source: source.realSource, + StartTime: snapStart, + EndTime: snapEnd, + Description: source.realSource, + TotalSize: objectSize, + RootObject: udmrepo.ObjectMetadata{ + ID: entryID, + Name: "bdev-root", + Type: udmrepo.ObjectDataTypeMetadata, + Permissions: 0o777, + }, + }, 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") +} + +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 + } + + 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 (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()) + 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 <-blkup.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 <-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 + } + + 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++ + + blkup.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 + + blkup.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, "/", "-") + s = strings.ReplaceAll(s, "\\", "-") + return strings.Trim(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..88fd4771e 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/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) + + blkup, ok := uploader.(*blockUploader) + assert.True(t, ok) + 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) { + 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 { + require.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 backing up bdev /data/volume1: 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 + + blkup := NewUploader(ctx, repoWriter, progress, log) + + f, err := os.CreateTemp(t.TempDir(), "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).Maybe() + + 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 := blkup.Backup(srcInfo, tc.parentObj, iterator, nil) + + if tc.expectErr { + require.Error(t, err) + if tc.expectErrStr != "" { + assert.Contains(t, err.Error(), tc.expectErrStr) + } + } else { + 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) + } + }) + } +} + 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