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 12dd0f688..6e9355765 100644 --- a/pkg/repository/udmrepo/kopialib/lib_repo.go +++ b/pkg/repository/udmrepo/kopialib/lib_repo.go @@ -92,6 +92,7 @@ type kopiaObjectWriterEx struct { description string compressor compression.Name splitter string + zeroObject object.ID writeLock sync.Mutex asyncWritesSem chan struct{} asyncWritesGroup sync.WaitGroup @@ -479,6 +480,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, @@ -985,9 +987,119 @@ 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.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) + } + } + + 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 writing 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..fdaeb9f69 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,59 @@ 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 { + t.Helper() + 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) { + t.Helper() + err := kow.getWriteError() + 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) + + 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) { + t.Helper() + err := kow.getWriteError() + require.Error(t, err) + assert.Contains(t, err.Error(), "simulated result error") + }, + }, { name: "success sync write", setupWriter: func(t *testing.T) *kopiaObjectWriterEx { @@ -350,6 +404,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 +463,48 @@ func TestKopiaObjectWriterEx_Result(t *testing.T) { }, expectedID: udmrepo.ID("IIabcdef"), }, + { + name: "write indirect object encoding failure", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() + 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 { + t.Helper() + 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 +615,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 +630,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 +644,652 @@ 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 { + t.Helper() + 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 { + t.Helper() + 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 { + t.Helper() + 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 { + t.Helper() + 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 { + t.Helper() + 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) { + 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) + + 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) { + 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) + 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 { + t.Helper() + 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) { + 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) + assert.Equal(t, int64(1024), kow.entries[1].Start) + }, + }, + { + name: "success write zero length", + setupWriter: func(t *testing.T) *kopiaObjectWriterEx { + t.Helper() + 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) { + 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(), + 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 { + t.Helper() + 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) { + t.Helper() + assert.Len(t, kow.entries, 3) + 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 { + t.Helper() + 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 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) + + 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) { + t.Helper() + err := kow.getWriteError() + require.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 { + require.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) + require.NoError(t, err) + assert.Equal(t, 1024, l) + + l, err = kow.WriteAt(make([]byte, 1024), 2048) + require.NoError(t, err) + assert.Equal(t, 1024, l) + + 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) + 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.NotEmpty(t, kow.entries) +} + +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) + require.NoError(t, err) + assert.Equal(t, int(blockSize), l) + } + + assert.Len(t, kow.entries, blocks) + 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) + require.NoError(t, err) + assert.Equal(t, int(blockSize), l) + + expectedEntries := 5121 + assert.Len(t, kow.entries, expectedEntries) + 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) + require.NoError(t, err) + assert.Equal(t, int(blockSize), l) + + // Entries: [0:1024] + 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) + require.NoError(t, err) + assert.Equal(t, int(blockSize), l) + + // Entries should now be 3: [0:1024, 1024:2048(zero object), 2048:3072] + 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 + 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) + 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.Len(t, kow.entries, 4) + 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++ { + l, err := kow.Write(data) + require.NoError(t, err) + assert.Equal(t, 1024, l) + } + + id, err := kow.Result() + + require.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 + l, err := kow.Write(make([]byte, 1024)) + require.NoError(t, err) + assert.Equal(t, 1024, l) + }() + } + + // 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.NotEmpty(t, kow.entries) +} + +func TestKopiaObjectWriterEx_Checkpoint(t *testing.T) { + kow := &kopiaObjectWriterEx{} + id, err := kow.Checkpoint() + require.Error(t, err) + assert.Equal(t, udmrepo.ID(""), id) + assert.Equal(t, "not supported", err.Error()) +}