From 30960d1edd3b31bd2d6bdbde36e083e9974a2b3b Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Fri, 8 May 2026 17:50:53 +0800 Subject: [PATCH 1/2] incremental aware object writer - writeat Signed-off-by: Lyndon-Li --- pkg/repository/udmrepo/kopialib/lib_repo.go | 112 ++- .../udmrepo/kopialib/lib_repo_ex_test.go | 723 +++++++++++++++++- 2 files changed, 830 insertions(+), 5 deletions(-) diff --git a/pkg/repository/udmrepo/kopialib/lib_repo.go b/pkg/repository/udmrepo/kopialib/lib_repo.go index 4fe7d6f35..5b2efbd69 100644 --- a/pkg/repository/udmrepo/kopialib/lib_repo.go +++ b/pkg/repository/udmrepo/kopialib/lib_repo.go @@ -92,6 +92,8 @@ type kopiaObjectWriterEx struct { description string compressor compression.Name splitter string + zeroBuffer []byte + zeroObject object.ID writeLock sync.Mutex asyncWritesSem chan struct{} asyncWritesGroup sync.WaitGroup @@ -479,6 +481,7 @@ func (kr *kopiaRepository) NewObjectWriter(ctx context.Context, opt udmrepo.Obje description: opt.Description, compressor: getCompressorForObject(opt), blockSize: fixedBlockSize, + zeroObject: object.EmptyID, splitter: fixedSplitter1M, asyncWritesSem: asyncWritesSem, asyncBuffer: asyncBuffer, @@ -977,9 +980,114 @@ func (kow *kopiaObjectWriterEx) writeObjectAsync(objName string, entryID int, p } } -// TODO add implementation in following PRs +func (kow *kopiaObjectWriterEx) writeZeroObject(objName string, entryID int) error { + if kow.zeroObject == object.EmptyID { + zeroBuffer := make([]byte, kow.blockSize) + objectID, err := kow.writeObject(objName, zeroBuffer) + if err != nil { + return err + } + + kow.zeroObject = objectID + } + + kow.entryLock.Lock() + kow.entries[entryID].Object = kow.zeroObject + kow.entryLock.Unlock() + + return nil +} + func (kow *kopiaObjectWriterEx) WriteAt(p []byte, offset int64) (int, error) { - return 0, errors.New("not implemented") + kow.writeLock.Lock() + defer kow.writeLock.Unlock() + + if kow.rawRepoWriter == nil { + return 0, errors.New("object writer is closed or not open") + } + + if err := kow.getWriteError(); err != nil { + return 0, errors.Wrapf(err, "error happened during writing object") + } + + if offset%kow.blockSize != 0 { + return 0, errors.Errorf("invalid offset %v", offset) + } + + length := len(p) + if int64(length)%kow.blockSize != 0 { + return 0, errors.Errorf("invalid length %v", length) + } + + kow.entryLock.Lock() + curPos := int64(len(kow.entries)) * kow.blockSize + kow.entryLock.Unlock() + + if offset < curPos { + return 0, errors.Errorf("cannot write back, cur pos %v", curPos) + } + + if offset > curPos && kow.parentEntries != nil { + startEntry := int(curPos / kow.blockSize) + endEntry := int(offset / kow.blockSize) + if startEntry < len(kow.parentEntries) { + if len(kow.parentEntries) < endEntry { + endEntry = len(kow.parentEntries) + } + + for i := startEntry; i < endEntry; i++ { + e := kow.parentEntries[i] + if e.Length != kow.blockSize { + return 0, errors.Errorf("parent entry %v length %v does not match child block size %v", i, e.Length, kow.blockSize) + } + } + + kow.entryLock.Lock() + kow.entries = append(kow.entries, kow.parentEntries[startEntry:endEntry]...) + curPos = int64(len(kow.entries)) * kow.blockSize + kow.entryLock.Unlock() + } + } + + entryID := 0 + for curPos < offset { + kow.entryLock.Lock() + entryID = len(kow.entries) + kow.entries = append(kow.entries, object.IndirectObjectEntry{ + Start: curPos, + Length: kow.blockSize, + }) + kow.entryLock.Unlock() + + objName := fmt.Sprintf("%s-b%v", kow.description, entryID) + if err := kow.writeZeroObject(objName, entryID); err != nil { + return 0, errors.Wrapf(err, "error writting zero object for %s", objName) + } + + curPos += kow.blockSize + } + + if length == 0 { + return length, nil + } + + for curPos < offset+int64(length) { + kow.entryLock.Lock() + entryID = len(kow.entries) + kow.entries = append(kow.entries, object.IndirectObjectEntry{ + Start: curPos, + Length: kow.blockSize, + }) + kow.entryLock.Unlock() + + buffOffset := curPos - offset + objName := fmt.Sprintf("%s-b%v", kow.description, entryID) + kow.writeObjectAsync(objName, entryID, p[buffOffset:buffOffset+kow.blockSize]) + + curPos += kow.blockSize + } + + return length, nil } func (kow *kopiaObjectWriterEx) Checkpoint() (udmrepo.ID, error) { diff --git a/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go b/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go index 6d9c5fc98..428ed0f11 100644 --- a/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go +++ b/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go @@ -216,6 +216,7 @@ func TestKopiaObjectWriterEx_Write(t *testing.T) { inputData []byte expectedErr string expectedLen int + verify func(t *testing.T, kow *kopiaObjectWriterEx) }{ { name: "writer is closed", @@ -254,6 +255,55 @@ func TestKopiaObjectWriterEx_Write(t *testing.T) { inputData: make([]byte, 1023), expectedErr: "invalid length 1023", }, + { + name: "write object returns nil writer", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(nil) + + return &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + logger: velerotest.NewLogger(), + } + }, + inputData: make([]byte, 1024), + expectedLen: 1024, + verify: func(t *testing.T, kow *kopiaObjectWriterEx) { + err := kow.getWriteError() + assert.Error(t, err) + assert.Contains(t, err.Error(), "error openning writer for -b0") + }, + }, + { + name: "write object result error", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(1024, nil) + mockWriter.On("Close").Return(nil) + + mockWriter.On("Result").Return(object.EmptyID, errors.New("simulated result error")) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + return &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + logger: velerotest.NewLogger(), + } + }, + inputData: make([]byte, 1024), + expectedLen: 1024, + verify: func(t *testing.T, kow *kopiaObjectWriterEx) { + err := kow.getWriteError() + assert.Error(t, err) + assert.Contains(t, err.Error(), "simulated result error") + }, + }, { name: "success sync write", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { @@ -350,6 +400,9 @@ func TestKopiaObjectWriterEx_Write(t *testing.T) { } else { require.NoError(t, err) assert.Equal(t, tc.expectedLen, l) + if tc.verify != nil { + tc.verify(t, kow) + } } }) } @@ -406,6 +459,46 @@ func TestKopiaObjectWriterEx_Result(t *testing.T) { }, expectedID: udmrepo.ID("IIabcdef"), }, + { + name: "write indirect object encoding failure", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(0, errors.New("json encoding failed")) + mockWriter.On("Close").Return(nil) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + return &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + logger: velerotest.NewLogger(), + } + }, + expectedErr: "error to write indirect object: unable to write indirect object index: json encoding failed", + }, + { + name: "write indirect object result failure", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(100, nil) + mockWriter.On("Close").Return(nil) + + mockWriter.On("Result").Return(object.EmptyID, errors.New("result generation failed")) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + return &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + logger: velerotest.NewLogger(), + } + }, + expectedErr: "error to write indirect object: result generation failed", + }, } for _, tc := range testCases { @@ -516,7 +609,6 @@ func TestKopiaObjectWriterEx_MultipleWrites(t *testing.T) { mockRepoWriter := repomocks.NewMockRepositoryWriter(t) mockWriter := repomocks.NewWriter(t) - // Since we are writing 3 blocks, Write should be called 3 times and Close 3 times mockWriter.On("Write", mock.Anything).Return(1024, nil) mockWriter.On("Close").Return(nil) @@ -532,12 +624,10 @@ func TestKopiaObjectWriterEx_MultipleWrites(t *testing.T) { logger: velerotest.NewLogger(), } - // Write 1st block l, err := kow.Write(make([]byte, 1024)) require.NoError(t, err) assert.Equal(t, 1024, l) - // Write 2nd and 3rd block l, err = kow.Write(make([]byte, 2048)) require.NoError(t, err) assert.Equal(t, 2048, l) @@ -548,3 +638,630 @@ func TestKopiaObjectWriterEx_MultipleWrites(t *testing.T) { assert.Equal(t, int64(1024), kow.entries[1].Start) assert.Equal(t, int64(2048), kow.entries[2].Start) } + +func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { + testCases := []struct { + name string + setupWriter func(t *testing.T) *kopiaObjectWriterEx + inputData []byte + offset int64 + expectedErr string + expectedLen int + verify func(t *testing.T, kow *kopiaObjectWriterEx) + }{ + { + name: "writer is closed", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + return &kopiaObjectWriterEx{ + rawRepoWriter: nil, + } + }, + inputData: make([]byte, 1024), + offset: 0, + expectedErr: "object writer is closed or not open", + }, + { + name: "invalid offset", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + return &kopiaObjectWriterEx{ + rawRepoWriter: repomocks.NewMockRepositoryWriter(t), + blockSize: 1024, + } + }, + inputData: make([]byte, 1024), + offset: 1023, + expectedErr: "invalid offset 1023", + }, + { + name: "invalid length", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + return &kopiaObjectWriterEx{ + rawRepoWriter: repomocks.NewMockRepositoryWriter(t), + blockSize: 1024, + } + }, + inputData: make([]byte, 1023), + offset: 0, + expectedErr: "invalid length 1023", + }, + { + name: "cannot write back", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + return &kopiaObjectWriterEx{ + rawRepoWriter: repomocks.NewMockRepositoryWriter(t), + blockSize: 1024, + entries: []object.IndirectObjectEntry{ + {Start: 0, Length: 1024}, + }, + } + }, + inputData: make([]byte, 1024), + offset: 0, + expectedErr: "cannot write back, cur pos 1024", + }, + { + name: "success write at cur pos", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(1024, nil) + mockWriter.On("Close").Return(nil) + + id, _ := object.ParseID("I12345") + mockWriter.On("Result").Return(id, nil) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + return &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + logger: velerotest.NewLogger(), + } + }, + inputData: make([]byte, 1024), + offset: 0, + expectedLen: 1024, + verify: func(t *testing.T, kow *kopiaObjectWriterEx) { + assert.Equal(t, 1, len(kow.entries)) + assert.Equal(t, int64(0), kow.entries[0].Start) + }, + }, + { + name: "success write with gap filling zeros", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(1024, nil) + mockWriter.On("Close").Return(nil) + + id, _ := object.ParseID("I12345") + mockWriter.On("Result").Return(id, nil) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + return &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + zeroObject: object.EmptyID, + logger: velerotest.NewLogger(), + } + }, + inputData: make([]byte, 1024), + offset: 1024, + expectedLen: 1024, + verify: func(t *testing.T, kow *kopiaObjectWriterEx) { + assert.Equal(t, 2, len(kow.entries)) + assert.Equal(t, int64(0), kow.entries[0].Start) + id, _ := object.ParseID("I12345") + assert.Equal(t, id, kow.entries[0].Object) + assert.Equal(t, id, kow.zeroObject) + assert.Equal(t, int64(1024), kow.entries[1].Start) + }, + }, + { + name: "success write with gap filling from parent", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(1024, nil) + mockWriter.On("Close").Return(nil) + + id, _ := object.ParseID("I12345") + mockWriter.On("Result").Return(id, nil) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + parentID, _ := object.ParseID("Iparent") + return &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + parentEntries: []object.IndirectObjectEntry{ + {Start: 0, Length: 1024, Object: parentID}, + }, + logger: velerotest.NewLogger(), + } + }, + inputData: make([]byte, 1024), + offset: 1024, + expectedLen: 1024, + verify: func(t *testing.T, kow *kopiaObjectWriterEx) { + assert.Equal(t, 2, len(kow.entries)) + assert.Equal(t, int64(0), kow.entries[0].Start) + parentID, _ := object.ParseID("Iparent") + assert.Equal(t, parentID, kow.entries[0].Object) + assert.Equal(t, int64(1024), kow.entries[1].Start) + }, + }, + { + name: "success write zero length", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + return &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + logger: velerotest.NewLogger(), + } + }, + inputData: []byte{}, + offset: 0, + expectedLen: 0, + verify: func(t *testing.T, kow *kopiaObjectWriterEx) { + assert.Equal(t, 0, len(kow.entries)) + }, + }, + { + name: "gap filling with invalid parent entry length", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + return &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + parentEntries: []object.IndirectObjectEntry{ + {Start: 0, Length: 512, Object: object.EmptyID}, + }, + logger: velerotest.NewLogger(), + } + }, + inputData: make([]byte, 1024), + offset: 1024, + expectedErr: "parent entry 0 length 512 does not match child block size 1024", + }, + { + name: "gap filling partially with parent and rest with zeros", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(1024, nil) + mockWriter.On("Close").Return(nil) + + id, _ := object.ParseID("I12345") + mockWriter.On("Result").Return(id, nil) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + parentID, _ := object.ParseID("Iparent") + return &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + zeroObject: object.EmptyID, + parentEntries: []object.IndirectObjectEntry{ + {Start: 0, Length: 1024, Object: parentID}, + }, + logger: velerotest.NewLogger(), + } + }, + inputData: make([]byte, 1024), + offset: 2048, + expectedLen: 1024, + verify: func(t *testing.T, kow *kopiaObjectWriterEx) { + assert.Equal(t, 3, len(kow.entries)) + assert.Equal(t, int64(0), kow.entries[0].Start) + + parentID, _ := object.ParseID("Iparent") + assert.Equal(t, parentID, kow.entries[0].Object) + + zeroID, _ := object.ParseID("I12345") + assert.Equal(t, int64(1024), kow.entries[1].Start) + assert.Equal(t, zeroID, kow.entries[1].Object) + + assert.Equal(t, int64(2048), kow.entries[2].Start) + }, + }, + { + name: "writeZeroObject failure", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(0, errors.New("simulated zero object write error")) + mockWriter.On("Close").Return(nil) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + return &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + zeroObject: object.EmptyID, + logger: velerotest.NewLogger(), + } + }, + inputData: make([]byte, 1024), + offset: 1024, + expectedErr: "error writting zero object for -b0: error writting for -b0: simulated zero object write error", + }, + { + name: "writeObject short write", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(512, nil) + mockWriter.On("Close").Return(nil) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + return &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + logger: velerotest.NewLogger(), + } + }, + inputData: make([]byte, 1024), + offset: 0, + expectedLen: 1024, + verify: func(t *testing.T, kow *kopiaObjectWriterEx) { + err := kow.getWriteError() + assert.Error(t, err) + assert.Contains(t, err.Error(), "short write for -b0") + }, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + kow := tc.setupWriter(t) + l, err := kow.WriteAt(tc.inputData, tc.offset) + + if kow.asyncWritesSem != nil { + kow.asyncWritesGroup.Wait() + } + + if tc.expectedErr != "" { + assert.EqualError(t, err, tc.expectedErr) + } else { + assert.NoError(t, err) + assert.Equal(t, tc.expectedLen, l) + if tc.verify != nil { + tc.verify(t, kow) + } + } + }) + } +} + +func TestKopiaObjectWriterEx_MultipleWriteAt(t *testing.T) { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(1024, nil) + mockWriter.On("Close").Return(nil) + + id, _ := object.ParseID("I12345") + mockWriter.On("Result").Return(id, nil) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + kow := &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + zeroObject: object.EmptyID, + logger: velerotest.NewLogger(), + } + + l, err := kow.WriteAt(make([]byte, 1024), 0) + assert.NoError(t, err) + assert.Equal(t, 1024, l) + + l, err = kow.WriteAt(make([]byte, 1024), 2048) + assert.NoError(t, err) + assert.Equal(t, 1024, l) + + assert.Equal(t, 3, len(kow.entries)) + assert.Equal(t, int64(0), kow.entries[0].Start) + assert.Equal(t, int64(1024), kow.entries[1].Start) + assert.Equal(t, id, kow.entries[1].Object) + assert.Equal(t, int64(2048), kow.entries[2].Start) +} + +func TestKopiaObjectWriterEx_ConcurrentWriteAt(t *testing.T) { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(1024, nil) + mockWriter.On("Close").Return(nil) + + id, _ := object.ParseID("I12345") + mockWriter.On("Result").Return(id, nil) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + kow := &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + logger: velerotest.NewLogger(), + } + + numGoroutines := 10 + var wg sync.WaitGroup + + start := make(chan struct{}) + + for i := 0; i < numGoroutines; i++ { + wg.Add(1) + go func(offset int64) { + defer wg.Done() + <-start + + data := make([]byte, 1024) + _, err := kow.WriteAt(data, offset) + + if err != nil { + assert.Contains(t, err.Error(), "cannot write back") + } + }(int64(i * 1024)) + } + + close(start) + wg.Wait() + + assert.Greater(t, len(kow.entries), 0) +} + +type dummyObjectWriter struct { + writtenBytes int +} + +func (dw *dummyObjectWriter) Write(p []byte) (int, error) { + dw.writtenBytes += len(p) + return len(p), nil +} + +func (dw *dummyObjectWriter) Close() error { + return nil +} + +func (dw *dummyObjectWriter) Result() (object.ID, error) { + id, _ := object.ParseID("I12345") + return id, nil +} + +func (dw *dummyObjectWriter) Checkpoint() (object.ID, error) { + return dw.Result() +} + +type dummyRepoWriter struct { + repo.RepositoryWriter +} + +func (drw *dummyRepoWriter) NewObjectWriter(ctx context.Context, opt object.WriterOptions) object.Writer { + return &dummyObjectWriter{} +} + +func TestKopiaObjectWriterEx_LargeSequentialWrite(t *testing.T) { + mockRepoWriter := &dummyRepoWriter{} + + blockSize := int64(1 << 20) + + kow := &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: blockSize, + logger: velerotest.NewLogger(), + } + + data := make([]byte, blockSize) + blocks := 5120 + + for i := 0; i < blocks; i++ { + l, err := kow.Write(data) + assert.NoError(t, err) + assert.Equal(t, int(blockSize), l) + } + + assert.Equal(t, blocks, len(kow.entries)) + assert.Equal(t, int64(blocks-1)*blockSize, kow.entries[blocks-1].Start) +} + +func TestKopiaObjectWriterEx_LargeSparseWriteAt(t *testing.T) { + mockRepoWriter := &dummyRepoWriter{} + + blockSize := int64(1 << 20) + + kow := &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: blockSize, + zeroObject: object.EmptyID, + logger: velerotest.NewLogger(), + } + + var offset int64 = 5 * 1024 * 1024 * 1024 + + data := make([]byte, blockSize) + l, err := kow.WriteAt(data, offset) + assert.NoError(t, err) + assert.Equal(t, int(blockSize), l) + + expectedEntries := 5121 + assert.Equal(t, expectedEntries, len(kow.entries)) + assert.Equal(t, int64(0), kow.entries[0].Start) + assert.Equal(t, offset, kow.entries[expectedEntries-1].Start) +} + +func TestKopiaObjectWriterEx_MixedWriteAndWriteAt(t *testing.T) { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + blockSize := int64(1024) + + mockWriter.On("Write", mock.Anything).Return(int(blockSize), nil) + mockWriter.On("Close").Return(nil) + + id, _ := object.ParseID("I12345") + mockWriter.On("Result").Return(id, nil) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + kow := &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: blockSize, + zeroObject: object.EmptyID, + logger: velerotest.NewLogger(), + } + + // 1. Write 1 block sequentially + data1 := make([]byte, blockSize) + l, err := kow.Write(data1) + assert.NoError(t, err) + assert.Equal(t, int(blockSize), l) + + // Entries: [0:1024] + assert.Equal(t, 1, len(kow.entries)) + assert.Equal(t, int64(0), kow.entries[0].Start) + + // 2. WriteAt with gap (offset = 2048). This creates a gap block at 1024 + data2 := make([]byte, blockSize) + l, err = kow.WriteAt(data2, 2048) + assert.NoError(t, err) + assert.Equal(t, int(blockSize), l) + + // Entries should now be 3: [0:1024, 1024:2048(zero object), 2048:3072] + assert.Equal(t, 3, len(kow.entries)) + assert.Equal(t, int64(0), kow.entries[0].Start) + assert.Equal(t, int64(1024), kow.entries[1].Start) + assert.Equal(t, id, kow.entries[1].Object) // filled with zero block + assert.Equal(t, int64(2048), kow.entries[2].Start) + + // 3. Write another block sequentially. It should append at 3072. + data3 := make([]byte, blockSize) + l, err = kow.Write(data3) + assert.NoError(t, err) + assert.Equal(t, int(blockSize), l) + + // Entries should now be 4: [0:1024, 1024:2048(zero object), 2048:3072, 3072:4096] + assert.Equal(t, 4, len(kow.entries)) + assert.Equal(t, int64(3072), kow.entries[3].Start) +} + +func TestKopiaObjectWriterEx_ConcurrentAsyncErrors(t *testing.T) { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(0, errors.New("simulated async error")) + mockWriter.On("Close").Return(nil) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + sem := make(chan struct{}, 10) + buf := freelist.New(10*1024, 1024) + + kow := &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + asyncWritesSem: sem, + asyncBuffer: buf, + logger: velerotest.NewLogger(), + } + + data := make([]byte, 1024) + + // Issue multiple writes so they all spawn async goroutines + // First few writes shouldn't fail immediately until getWriteError catches the asynchronous fault + for i := 0; i < 10; i++ { + kow.Write(data) + } + + id, err := kow.Result() + + assert.Error(t, err) + assert.Contains(t, err.Error(), "simulated async error") + assert.Equal(t, udmrepo.ID(""), id) +} + +func TestKopiaObjectWriterEx_ConcurrentWriteAndWriteAt(t *testing.T) { + mockRepoWriter := repomocks.NewMockRepositoryWriter(t) + mockWriter := repomocks.NewWriter(t) + + mockWriter.On("Write", mock.Anything).Return(1024, nil) + mockWriter.On("Close").Return(nil) + + id, _ := object.ParseID("I12345") + mockWriter.On("Result").Return(id, nil) + + mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(mockWriter) + + kow := &kopiaObjectWriterEx{ + ctx: context.Background(), + rawRepoWriter: mockRepoWriter, + blockSize: 1024, + zeroObject: object.EmptyID, + logger: velerotest.NewLogger(), + } + + var wg sync.WaitGroup + start := make(chan struct{}) + + for i := 0; i < 5; i++ { + wg.Add(1) + go func() { + defer wg.Done() + <-start + kow.Write(make([]byte, 1024)) + }() + } + + // Fire multiple sparse WriteAts alongside them + // Note: Because order is totally random and WriteAt strictly demands monotonic offsets, + // some will hit the legitimate "cannot write back" error, which we safely expect. + for i := 0; i < 5; i++ { + wg.Add(1) + go func(offset int64) { + defer wg.Done() + <-start + _, err := kow.WriteAt(make([]byte, 1024), offset) + if err != nil { + assert.Contains(t, err.Error(), "cannot write back") + } + }(int64(i * 2048)) + } + + close(start) + wg.Wait() + + // We only care that the locking effectively mitigated a panic or slice data corruption + assert.Greater(t, len(kow.entries), 0) +} + +func TestKopiaObjectWriterEx_Checkpoint(t *testing.T) { + kow := &kopiaObjectWriterEx{} + id, err := kow.Checkpoint() + assert.Error(t, err) + assert.Equal(t, udmrepo.ID(""), id) + assert.Equal(t, "not supported", err.Error()) +} From f474e313fa25a25cef34c0be577696668bd0c509 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Fri, 5 Jun 2026 13:32:55 +0800 Subject: [PATCH 2/2] incremental object aware write at Signed-off-by: Lyndon-Li --- changelogs/unreleased/9887-Lyndon-Li | 1 + pkg/repository/udmrepo/kopialib/lib_repo.go | 8 +- .../udmrepo/kopialib/lib_repo_ex_test.go | 88 ++++++++++++------- 3 files changed, 65 insertions(+), 32 deletions(-) create mode 100644 changelogs/unreleased/9887-Lyndon-Li diff --git a/changelogs/unreleased/9887-Lyndon-Li b/changelogs/unreleased/9887-Lyndon-Li new file mode 100644 index 000000000..a44418c9c --- /dev/null +++ b/changelogs/unreleased/9887-Lyndon-Li @@ -0,0 +1 @@ +Add WriteAt implementation for Incremental aware object writer 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 bb70e1b18..6e9355765 100644 --- a/pkg/repository/udmrepo/kopialib/lib_repo.go +++ b/pkg/repository/udmrepo/kopialib/lib_repo.go @@ -92,7 +92,6 @@ type kopiaObjectWriterEx struct { description string compressor compression.Name splitter string - zeroBuffer []byte zeroObject object.ID writeLock sync.Mutex asyncWritesSem chan struct{} @@ -1045,6 +1044,11 @@ func (kow *kopiaObjectWriterEx) WriteAt(p []byte, offset int64) (int, error) { for i := startEntry; i < endEntry; i++ { e := kow.parentEntries[i] + + if e.Start != int64(i)*kow.blockSize { + return 0, errors.Errorf("parent entry %v start %v does not match expected start %v", i, e.Start, int64(i)*kow.blockSize) + } + if e.Length != kow.blockSize { return 0, errors.Errorf("parent entry %v length %v does not match child block size %v", i, e.Length, kow.blockSize) } @@ -1069,7 +1073,7 @@ func (kow *kopiaObjectWriterEx) WriteAt(p []byte, offset int64) (int, error) { objName := fmt.Sprintf("%s-b%v", kow.description, entryID) if err := kow.writeZeroObject(objName, entryID); err != nil { - return 0, errors.Wrapf(err, "error writting zero object for %s", objName) + return 0, errors.Wrapf(err, "error writing zero object for %s", objName) } curPos += kow.blockSize diff --git a/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go b/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go index 428ed0f11..fdaeb9f69 100644 --- a/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go +++ b/pkg/repository/udmrepo/kopialib/lib_repo_ex_test.go @@ -258,6 +258,7 @@ func TestKopiaObjectWriterEx_Write(t *testing.T) { { name: "write object returns nil writer", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() mockRepoWriter := repomocks.NewMockRepositoryWriter(t) mockRepoWriter.On("NewObjectWriter", mock.Anything, mock.Anything).Return(nil) @@ -271,14 +272,16 @@ func TestKopiaObjectWriterEx_Write(t *testing.T) { inputData: make([]byte, 1024), expectedLen: 1024, verify: func(t *testing.T, kow *kopiaObjectWriterEx) { + t.Helper() err := kow.getWriteError() - assert.Error(t, err) - assert.Contains(t, err.Error(), "error openning writer for -b0") + require.Error(t, err) + assert.Contains(t, err.Error(), "error opening writer for -b0") }, }, { name: "write object result error", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() mockRepoWriter := repomocks.NewMockRepositoryWriter(t) mockWriter := repomocks.NewWriter(t) @@ -299,8 +302,9 @@ func TestKopiaObjectWriterEx_Write(t *testing.T) { inputData: make([]byte, 1024), expectedLen: 1024, verify: func(t *testing.T, kow *kopiaObjectWriterEx) { + t.Helper() err := kow.getWriteError() - assert.Error(t, err) + require.Error(t, err) assert.Contains(t, err.Error(), "simulated result error") }, }, @@ -462,6 +466,7 @@ func TestKopiaObjectWriterEx_Result(t *testing.T) { { name: "write indirect object encoding failure", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() mockRepoWriter := repomocks.NewMockRepositoryWriter(t) mockWriter := repomocks.NewWriter(t) @@ -481,6 +486,7 @@ func TestKopiaObjectWriterEx_Result(t *testing.T) { { name: "write indirect object result failure", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() mockRepoWriter := repomocks.NewMockRepositoryWriter(t) mockWriter := repomocks.NewWriter(t) @@ -652,6 +658,7 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { { name: "writer is closed", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() return &kopiaObjectWriterEx{ rawRepoWriter: nil, } @@ -663,6 +670,7 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { { name: "invalid offset", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() return &kopiaObjectWriterEx{ rawRepoWriter: repomocks.NewMockRepositoryWriter(t), blockSize: 1024, @@ -675,6 +683,7 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { { name: "invalid length", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() return &kopiaObjectWriterEx{ rawRepoWriter: repomocks.NewMockRepositoryWriter(t), blockSize: 1024, @@ -687,6 +696,7 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { { name: "cannot write back", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() return &kopiaObjectWriterEx{ rawRepoWriter: repomocks.NewMockRepositoryWriter(t), blockSize: 1024, @@ -702,6 +712,7 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { { name: "success write at cur pos", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() mockRepoWriter := repomocks.NewMockRepositoryWriter(t) mockWriter := repomocks.NewWriter(t) @@ -724,13 +735,15 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { offset: 0, expectedLen: 1024, verify: func(t *testing.T, kow *kopiaObjectWriterEx) { - assert.Equal(t, 1, len(kow.entries)) + t.Helper() + assert.Len(t, kow.entries, 1) assert.Equal(t, int64(0), kow.entries[0].Start) }, }, { name: "success write with gap filling zeros", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() mockRepoWriter := repomocks.NewMockRepositoryWriter(t) mockWriter := repomocks.NewWriter(t) @@ -754,7 +767,8 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { offset: 1024, expectedLen: 1024, verify: func(t *testing.T, kow *kopiaObjectWriterEx) { - assert.Equal(t, 2, len(kow.entries)) + t.Helper() + assert.Len(t, kow.entries, 2) assert.Equal(t, int64(0), kow.entries[0].Start) id, _ := object.ParseID("I12345") assert.Equal(t, id, kow.entries[0].Object) @@ -765,6 +779,7 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { { name: "success write with gap filling from parent", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() mockRepoWriter := repomocks.NewMockRepositoryWriter(t) mockWriter := repomocks.NewWriter(t) @@ -791,7 +806,8 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { offset: 1024, expectedLen: 1024, verify: func(t *testing.T, kow *kopiaObjectWriterEx) { - assert.Equal(t, 2, len(kow.entries)) + t.Helper() + assert.Len(t, kow.entries, 2) assert.Equal(t, int64(0), kow.entries[0].Start) parentID, _ := object.ParseID("Iparent") assert.Equal(t, parentID, kow.entries[0].Object) @@ -801,6 +817,7 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { { name: "success write zero length", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() mockRepoWriter := repomocks.NewMockRepositoryWriter(t) return &kopiaObjectWriterEx{ ctx: context.Background(), @@ -813,12 +830,14 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { offset: 0, expectedLen: 0, verify: func(t *testing.T, kow *kopiaObjectWriterEx) { - assert.Equal(t, 0, len(kow.entries)) + t.Helper() + assert.Empty(t, kow.entries) }, }, { name: "gap filling with invalid parent entry length", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() mockRepoWriter := repomocks.NewMockRepositoryWriter(t) return &kopiaObjectWriterEx{ ctx: context.Background(), @@ -837,6 +856,7 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { { name: "gap filling partially with parent and rest with zeros", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() mockRepoWriter := repomocks.NewMockRepositoryWriter(t) mockWriter := repomocks.NewWriter(t) @@ -864,7 +884,8 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { offset: 2048, expectedLen: 1024, verify: func(t *testing.T, kow *kopiaObjectWriterEx) { - assert.Equal(t, 3, len(kow.entries)) + t.Helper() + assert.Len(t, kow.entries, 3) assert.Equal(t, int64(0), kow.entries[0].Start) parentID, _ := object.ParseID("Iparent") @@ -880,6 +901,7 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { { name: "writeZeroObject failure", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() mockRepoWriter := repomocks.NewMockRepositoryWriter(t) mockWriter := repomocks.NewWriter(t) @@ -898,11 +920,12 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { }, inputData: make([]byte, 1024), offset: 1024, - expectedErr: "error writting zero object for -b0: error writting for -b0: simulated zero object write error", + expectedErr: "error writing zero object for -b0: error writing for -b0: simulated zero object write error", }, { name: "writeObject short write", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() mockRepoWriter := repomocks.NewMockRepositoryWriter(t) mockWriter := repomocks.NewWriter(t) @@ -922,8 +945,9 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { offset: 0, expectedLen: 1024, verify: func(t *testing.T, kow *kopiaObjectWriterEx) { + t.Helper() err := kow.getWriteError() - assert.Error(t, err) + require.Error(t, err) assert.Contains(t, err.Error(), "short write for -b0") }, }, @@ -941,7 +965,7 @@ func TestKopiaObjectWriterEx_WriteAt(t *testing.T) { if tc.expectedErr != "" { assert.EqualError(t, err, tc.expectedErr) } else { - assert.NoError(t, err) + require.NoError(t, err) assert.Equal(t, tc.expectedLen, l) if tc.verify != nil { tc.verify(t, kow) @@ -972,14 +996,14 @@ func TestKopiaObjectWriterEx_MultipleWriteAt(t *testing.T) { } l, err := kow.WriteAt(make([]byte, 1024), 0) - assert.NoError(t, err) + require.NoError(t, err) assert.Equal(t, 1024, l) l, err = kow.WriteAt(make([]byte, 1024), 2048) - assert.NoError(t, err) + require.NoError(t, err) assert.Equal(t, 1024, l) - assert.Equal(t, 3, len(kow.entries)) + assert.Len(t, kow.entries, 3) assert.Equal(t, int64(0), kow.entries[0].Start) assert.Equal(t, int64(1024), kow.entries[1].Start) assert.Equal(t, id, kow.entries[1].Object) @@ -1028,7 +1052,7 @@ func TestKopiaObjectWriterEx_ConcurrentWriteAt(t *testing.T) { close(start) wg.Wait() - assert.Greater(t, len(kow.entries), 0) + assert.NotEmpty(t, kow.entries) } type dummyObjectWriter struct { @@ -1078,11 +1102,11 @@ func TestKopiaObjectWriterEx_LargeSequentialWrite(t *testing.T) { for i := 0; i < blocks; i++ { l, err := kow.Write(data) - assert.NoError(t, err) + require.NoError(t, err) assert.Equal(t, int(blockSize), l) } - assert.Equal(t, blocks, len(kow.entries)) + assert.Len(t, kow.entries, blocks) assert.Equal(t, int64(blocks-1)*blockSize, kow.entries[blocks-1].Start) } @@ -1103,11 +1127,11 @@ func TestKopiaObjectWriterEx_LargeSparseWriteAt(t *testing.T) { data := make([]byte, blockSize) l, err := kow.WriteAt(data, offset) - assert.NoError(t, err) + require.NoError(t, err) assert.Equal(t, int(blockSize), l) expectedEntries := 5121 - assert.Equal(t, expectedEntries, len(kow.entries)) + assert.Len(t, kow.entries, expectedEntries) assert.Equal(t, int64(0), kow.entries[0].Start) assert.Equal(t, offset, kow.entries[expectedEntries-1].Start) } @@ -1137,21 +1161,21 @@ func TestKopiaObjectWriterEx_MixedWriteAndWriteAt(t *testing.T) { // 1. Write 1 block sequentially data1 := make([]byte, blockSize) l, err := kow.Write(data1) - assert.NoError(t, err) + require.NoError(t, err) assert.Equal(t, int(blockSize), l) // Entries: [0:1024] - assert.Equal(t, 1, len(kow.entries)) + assert.Len(t, kow.entries, 1) assert.Equal(t, int64(0), kow.entries[0].Start) // 2. WriteAt with gap (offset = 2048). This creates a gap block at 1024 data2 := make([]byte, blockSize) l, err = kow.WriteAt(data2, 2048) - assert.NoError(t, err) + require.NoError(t, err) assert.Equal(t, int(blockSize), l) // Entries should now be 3: [0:1024, 1024:2048(zero object), 2048:3072] - assert.Equal(t, 3, len(kow.entries)) + assert.Len(t, kow.entries, 3) assert.Equal(t, int64(0), kow.entries[0].Start) assert.Equal(t, int64(1024), kow.entries[1].Start) assert.Equal(t, id, kow.entries[1].Object) // filled with zero block @@ -1160,11 +1184,11 @@ func TestKopiaObjectWriterEx_MixedWriteAndWriteAt(t *testing.T) { // 3. Write another block sequentially. It should append at 3072. data3 := make([]byte, blockSize) l, err = kow.Write(data3) - assert.NoError(t, err) + require.NoError(t, err) assert.Equal(t, int(blockSize), l) // Entries should now be 4: [0:1024, 1024:2048(zero object), 2048:3072, 3072:4096] - assert.Equal(t, 4, len(kow.entries)) + assert.Len(t, kow.entries, 4) assert.Equal(t, int64(3072), kow.entries[3].Start) } @@ -1194,12 +1218,14 @@ func TestKopiaObjectWriterEx_ConcurrentAsyncErrors(t *testing.T) { // Issue multiple writes so they all spawn async goroutines // First few writes shouldn't fail immediately until getWriteError catches the asynchronous fault for i := 0; i < 10; i++ { - kow.Write(data) + l, err := kow.Write(data) + require.NoError(t, err) + assert.Equal(t, 1024, l) } id, err := kow.Result() - assert.Error(t, err) + require.Error(t, err) assert.Contains(t, err.Error(), "simulated async error") assert.Equal(t, udmrepo.ID(""), id) } @@ -1232,7 +1258,9 @@ func TestKopiaObjectWriterEx_ConcurrentWriteAndWriteAt(t *testing.T) { go func() { defer wg.Done() <-start - kow.Write(make([]byte, 1024)) + l, err := kow.Write(make([]byte, 1024)) + require.NoError(t, err) + assert.Equal(t, 1024, l) }() } @@ -1255,13 +1283,13 @@ func TestKopiaObjectWriterEx_ConcurrentWriteAndWriteAt(t *testing.T) { wg.Wait() // We only care that the locking effectively mitigated a panic or slice data corruption - assert.Greater(t, len(kow.entries), 0) + assert.NotEmpty(t, kow.entries) } func TestKopiaObjectWriterEx_Checkpoint(t *testing.T) { kow := &kopiaObjectWriterEx{} id, err := kow.Checkpoint() - assert.Error(t, err) + require.Error(t, err) assert.Equal(t, udmrepo.ID(""), id) assert.Equal(t, "not supported", err.Error()) }