block uploader backup implementation

Signed-off-by: Lyndon-Li <lyonghui@vmware.com>
This commit is contained in:
Lyndon-Li
2026-06-30 16:21:59 +08:00
parent 08385354db
commit df21463629
6 changed files with 617 additions and 7 deletions
+1 -1
View File
@@ -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 {
@@ -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 {
+264 -3
View File
@@ -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) {
+351 -2
View File
@@ -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