diff --git a/pkg/storage/indexheader/encoding/bucket_async_reader.go b/pkg/storage/indexheader/encoding/bucket_async_reader.go new file mode 100644 index 00000000000..b82346d53e8 --- /dev/null +++ b/pkg/storage/indexheader/encoding/bucket_async_reader.go @@ -0,0 +1,397 @@ +// SPDX-License-Identifier: AGPL-3.0-only + +package encoding + +import ( + "bufio" + "context" + "errors" + "fmt" + "io" + "sync" + + "github.com/thanos-io/objstore" + "golang.org/x/sync/errgroup" +) + +const ( + ReadAheadFactor = 2 +) + +type BufPromise struct { + bufioReader *bufio.Reader + + cancel context.CancelCauseFunc + eg *errgroup.Group +} + +func NewBufPromise( + ctx context.Context, + bkt objstore.BucketReader, + name string, + base int, + length int, + bufioReader *bufio.Reader, +) *BufPromise { + // Create handle to cancel an inflight bucket read + ctx, cancel := context.WithCancelCause(ctx) + bucketReader := NewBucketReader(ctx, bkt, name, base, length) + bufioReader.Reset(bucketReader) + + // Propagate cancellable context to errgroup. + eg, _ := errgroup.WithContext(ctx) + bp := &BufPromise{ + cancel: cancel, + eg: eg, + bufioReader: bufioReader, + } + + bp.eg.Go(func() error { + // Send the fill to the background. + // Peek reads length bytes into the buffer, and does not consume them. + b, err := bp.bufioReader.Peek(length) + if errors.Is(err, io.EOF) || errors.Is(err, io.ErrUnexpectedEOF) { + return fmt.Errorf( + "%w reading %d bytes at offset %d of %s (got %d bytes): %s", + ErrInvalidSize, length, base, name, len(b), err, + ) + } + return err + }) + return bp +} + +func (bp *BufPromise) Read(dst []byte) (n int, err error) { + if err := bp.eg.Wait(); err != nil { + // Return any error in the underlying bucket read from the initial fill-via-Peek. + // A Read after a bad Peek overwrites the bucket read error with a generic short read error. + return 0, err + } + return bp.bufioReader.Read(dst) +} + +func (bp *BufPromise) Buffered() (n int, err error) { + if err := bp.eg.Wait(); err != nil { + // Return any error in the underlying bucket read from the initial fill-via-Peek. + // A Read after a bad Peek overwrites the bucket read error with a generic short read error. + return 0, err + } + return bp.bufioReader.Buffered(), nil +} + +// Discard consumes and drops at most n bytes from the promise. +func (bp *BufPromise) Discard(n int) (discarded int, err error) { + if err := bp.eg.Wait(); err != nil { + // Return any error in the underlying bucket read from the initial fill-via-Peek. + return 0, err + } + return bp.bufioReader.Discard(n) +} + +func (bp *BufPromise) Release(bufioPool *sync.Pool, cancelCause error) { + bp.cancel(cancelCause) + _ = bp.eg.Wait() // Ensure any write to the buffer completes. + bufioPool.Put(bp.bufioReader) + bp.bufioReader = nil +} + +type BucketAsyncBufReader struct { + ctx context.Context + bkt objstore.BucketReader + name string + base int + length int + readOffset int + + bufSize int + + // peekBuf holds peeked bytes until they are read or discarded. + // Peek is intended to return a slice of bytes without an extra allocation, + // but we cannot return a slice which spans two underlying promise buffers. + // We allocate peekBuf with a large capacity and Peek returns subslices of it. + // peekBuf's slice bounds slides forward through the backing array until it runs out of capacity, + // then it is compacted by copying remaining elements back to the start of the array. + peekBuf []byte + peekBufBase []byte // Holds the reference to the start of the slice for compaction + + bufIdx int + bufPromises []*BufPromise + bufferedOffset int + bufioPool *sync.Pool +} + +func NewBucketAsyncBufReader( + ctx context.Context, + bkt objstore.BucketReader, + name string, base int, length int, + +) *BucketAsyncBufReader { + return newBucketAsyncBufReader( + ctx, bkt, name, base, length, 0, + &bucketBufioPool, ReadBufferSize, ReadAheadFactor, + ) +} + +func newBucketAsyncBufReader( + ctx context.Context, + bkt objstore.BucketReader, + name string, + base int, + length int, + startOffset int, + bufioPool *sync.Pool, + bufSize int, + maxBufCount int, +) *BucketAsyncBufReader { + bufsForLength := (length + bufSize - 1) / bufSize + numBufs := min(maxBufCount, bufsForLength) + bufPromises := make([]*BufPromise, numBufs) + + iBase := base + startOffset + bufferedOffset := startOffset + for i := range numBufs { + bufioReader := bufioPool.Get().(*bufio.Reader) + bufLen := min(length-bufferedOffset, bufSize) + bufPromises[i] = NewBufPromise(ctx, bkt, name, iBase, bufLen, bufioReader) + iBase += bufLen + bufferedOffset += bufLen + } + + peekBufBase := make([]byte, 0, bufSize) + return &BucketAsyncBufReader{ + ctx: ctx, + bkt: bkt, + name: name, + base: base, + length: length, + readOffset: startOffset, + bufSize: bufSize, + peekBufBase: peekBufBase, + peekBuf: peekBufBase[:0], + bufPromises: bufPromises, + bufferedOffset: bufferedOffset, + bufioPool: bufioPool, + } +} + +// rotateHead releases the exhausted head promise back to the pool +// and queues a promise to buffer the next read range in its place. +func (bbar *BucketAsyncBufReader) rotateHead() { + // No need to clean up buffer, we reset when we retrieve it from the pool + bbar.bufPromises[bbar.bufIdx].Release(bbar.bufioPool, nil) + bufioReader := bbar.bufioPool.Get().(*bufio.Reader) + + bufLen := min(bbar.length-bbar.bufferedOffset, bbar.bufSize) + bbar.bufPromises[bbar.bufIdx] = NewBufPromise( + bbar.ctx, bbar.bkt, bbar.name, bbar.base+bbar.bufferedOffset, bufLen, bufioReader, + ) + bbar.bufferedOffset += bufLen + + // Advance current buffer queue index - modulo wraps to the front of the slice. + bbar.bufIdx = (bbar.bufIdx + 1) % len(bbar.bufPromises) +} + +func (bbar *BucketAsyncBufReader) Reset() error { + return bbar.ResetAt(0) +} + +func (bbar *BucketAsyncBufReader) ResetAt(off int) error { + if off > bbar.length { + return ErrInvalidSize + } + + if dist := off - bbar.readOffset; dist > 0 && dist < bbar.Buffered() { + // Reset via Skip to avoid discarding all buffered bytes. + return bbar.Skip(dist) + } + + bbar.Close() + newBbar := newBucketAsyncBufReader( + bbar.ctx, bbar.bkt, bbar.name, bbar.base, bbar.length, off, + bbar.bufioPool, bbar.bufSize, ReadAheadFactor, + ) + *bbar = *newBbar + return nil +} + +func (bbar *BucketAsyncBufReader) Skip(l int) error { + if l > bbar.Len() { + return ErrInvalidSize + } + + // Start with any previously-peeked bytes. + bytesSkipped := min(len(bbar.peekBuf), l) + bbar.readOffset += bytesSkipped + + // Advance past the consumed bytes without moving them, + // so a slice returned by an earlier Peek stays valid. + bbar.peekBuf = bbar.peekBuf[bytesSkipped:] + + // Move on to the promises if we have not satisfied the skip yet. + // Promises are consumed to discard the data and rotated if exhausted. + for bytesSkipped < l { + headPromise := bbar.bufPromises[bbar.bufIdx] + headPromiseBuffered, err := headPromise.Buffered() + if err != nil { + return err + } + + toSkip := min(l-bytesSkipped, headPromiseBuffered) + skipN, err := headPromise.Discard(toSkip) + if err != nil { + return err + } + bbar.readOffset += skipN + bytesSkipped += skipN + + headPromiseBuffered, err = headPromise.Buffered() + if err != nil { + return err + } + + if headPromiseBuffered <= 0 { + bbar.rotateHead() + } + } + + return nil +} + +func (bbar *BucketAsyncBufReader) Peek(n int) ([]byte, error) { + // Clamp n to the lesser of the buffer size or the length of the section. + n = min(n, bbar.bufSize, bbar.Len()) + if n > cap(bbar.peekBuf) { + // Slide remaining peeked bytes + bbar.peekBuf = append(bbar.peekBufBase[:0], bbar.peekBuf...) + } + + // Start with any previously-peeked bytes. + // Any data remaining in peekBuf is assumed to still be the valid start of a Peek. + // Read/ReadInto, Reset/ResetAt, and Skip are required to update peekBuf + // to discard any previously-peeked bytes which were consumed by the read. + peekBytesAvailable := len(bbar.peekBuf) + if n > peekBytesAvailable { + bbar.peekBuf = bbar.peekBuf[:n] // Grow length; will not allocate. + } + peekBytesWritten := min(n, peekBytesAvailable) + + // Move on to the promises if we have not satisfied the peek yet. + // Promises are consumed to copy into the peekBuf and rotated if exhausted. + for peekBytesWritten < n { + headPromise := bbar.bufPromises[bbar.bufIdx] + headPromiseBuffered, err := headPromise.Buffered() + if err != nil { + return nil, err + } + + toRead := min(n-peekBytesWritten, headPromiseBuffered) + readN, err := io.ReadFull(headPromise, bbar.peekBuf[peekBytesWritten:peekBytesWritten+toRead]) + peekBytesWritten += readN + if err != nil { + return nil, err + } + + headPromiseBuffered, err = headPromise.Buffered() + if err != nil { + return nil, err + } + + if headPromiseBuffered <= 0 { + bbar.rotateHead() + } + } + + if peekBytesWritten == 0 { + return nil, nil + } + return bbar.peekBuf[:peekBytesWritten], nil +} + +func (bbar *BucketAsyncBufReader) Read(n int) ([]byte, error) { + b := make([]byte, n) + + err := bbar.ReadInto(b) + if err != nil { + return nil, err + } + + return b, nil +} + +func (bbar *BucketAsyncBufReader) ReadInto(dst []byte) error { + if len(dst) > bbar.Len() { + if err := bbar.Skip(bbar.Len()); err != nil { + return err + } + return ErrInvalidSize + } + + // Start with any previously-peeked bytes. + dstBytesWritten := copy(dst, bbar.peekBuf) + bbar.readOffset += dstBytesWritten + // Advance past the consumed bytes without moving the + // so a slice returned by an earlier Peek stays valid. + bbar.peekBuf = bbar.peekBuf[dstBytesWritten:] + + // Move on to the promises if we have not satisfied the read yet. + // Promises are consumed to copy into dst and rotated if exhausted. + for dstBytesWritten < len(dst) { + headPromise := bbar.bufPromises[bbar.bufIdx] + headPromiseBuffered, err := headPromise.Buffered() + if err != nil { + return err + } + + toRead := min(len(dst)-dstBytesWritten, headPromiseBuffered) + readN, err := io.ReadFull(headPromise, dst[dstBytesWritten:dstBytesWritten+toRead]) + bbar.readOffset += readN + dstBytesWritten += readN + if err != nil { + return err + } + + headPromiseBuffered, err = headPromise.Buffered() + if err != nil { + return err + } + + if headPromiseBuffered <= 0 { + bbar.rotateHead() + } + } + return nil +} + +func (bbar *BucketAsyncBufReader) Size() int { + // Reported capacity of peekBuf changes as its referenced window slides, + // but we will compact to make use of its full underlying allocated size if needed. + return bbar.bufSize +} + +func (bbar *BucketAsyncBufReader) Len() int { + return bbar.length - bbar.readOffset +} + +func (bbar *BucketAsyncBufReader) Offset() int { + return bbar.readOffset +} + +func (bbar *BucketAsyncBufReader) Buffered() int { + return bbar.bufferedOffset - bbar.readOffset +} + +// errBufPromiseReleased is the cancellation cause when release stops a fill that is still in flight. +var errBufReaderClosed = errors.New("BufReader closed") + +// Close cancels all promises and releases buffers back to the pool. +func (bbar *BucketAsyncBufReader) Close() error { + for i, bufPromise := range bbar.bufPromises { + // No need to clean up buffer, we reset when we retrieve it from the pool. + bufPromise.Release(bbar.bufioPool, errBufReaderClosed) + bbar.bufPromises[i] = nil + } + + // The BucketReader of a promise does not need a Close call. + // It closes the reader created by bkt.GetRange on each Read call. + return nil +} diff --git a/pkg/storage/indexheader/encoding/bucket_async_reader_test.go b/pkg/storage/indexheader/encoding/bucket_async_reader_test.go new file mode 100644 index 00000000000..89ebef5e6bb --- /dev/null +++ b/pkg/storage/indexheader/encoding/bucket_async_reader_test.go @@ -0,0 +1,330 @@ +package encoding + +import ( + "context" + "testing" + + "github.com/stretchr/testify/require" +) + +// testBucketContentsLong is a 64-byte payload for tests that need more data than +// the read-ahead window holds. Every byte is distinct, so a read at a wrong offset +// gives a different result. +var testBucketContentsLong = []byte("abcdefghijklmnopqrstuvwxyz0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZ+/") + +func newTestAsyncBufReader(t *testing.T, base, length int) (*BucketAsyncBufReader, *trackingBucket) { + t.Helper() + ctx := context.Background() + + objectData := make([]byte, 0, length) + objectData = append(objectData, testBucketContents...) + bkt := newTrackingBucket(t, objectData) + + return newBucketAsyncBufReader( + ctx, bkt, testBucketObjectName, base, length, 0, + &testBucketBufPool, testBufPoolSize, 4, + ), bkt +} + +// newTestAsyncBufReaderWithData builds a reader over the given object data, +// with an explicit read-ahead buffer count. +// A test that needs a rotation of the read-ahead window must use this helper, +// because a rotation needs a length greater than maxBufCount*testBufPoolSize bytes. +func newTestAsyncBufReaderWithData( + t *testing.T, objectData []byte, base, length, maxBufCount int, +) (*BucketAsyncBufReader, *trackingBucket) { + t.Helper() + + bkt := newTrackingBucket(t, objectData) + + return newBucketAsyncBufReader( + t.Context(), bkt, testBucketObjectName, base, length, 0, + &testBucketBufPool, testBufPoolSize, maxBufCount, + ), bkt +} + +func TestBucketAsyncBufReader_Peek(t *testing.T) { + r, _ := newTestAsyncBufReader(t, 0, len(testBucketContents)) + + peek1, err := r.Peek(5) + require.NoError(t, err) + require.Equal(t, testBucketContents[:5], peek1) + require.Equal(t, 0, r.Offset(), "Peek does not consume") + + // Same-size Peek returns the same bytes. + peek2, err := r.Peek(5) + require.NoError(t, err) + require.Equal(t, testBucketContents[:5], peek2) + + // Smaller Peek returns the same starting bytes. + peek3, err := r.Peek(3) + require.NoError(t, err) + require.Equal(t, testBucketContents[:3], peek3) + + // Larger Peek returns the original bytes plus more + peek4, err := r.Peek(10) + require.NoError(t, err) + require.Equal(t, testBucketContents[:10], peek4) +} + +func TestBucketAsyncBufReader_Peek_Skip(t *testing.T) { + r, _ := newTestAsyncBufReader(t, 0, len(testBucketContents)) + + peek1, err := r.Peek(5) + require.NoError(t, err) + require.Equal(t, testBucketContents[:5], peek1) + require.Equal(t, 0, r.Offset()) + + // Skip the same bytes. + require.NoError(t, r.Skip(5)) + require.Equal(t, 5, r.Offset()) + + // Another Peek returns the next bytes after Skip consumed previous bytes. + peek2, err := r.Peek(5) + require.NoError(t, err) + require.Equal(t, testBucketContents[5:10], peek2) + require.Equal(t, 5, r.Offset()) + + // A Skip short of the peeked bytes consumes some of them. + require.NoError(t, r.Skip(3)) + require.Equal(t, 8, r.Offset()) + + // A Skip beyond the remaining peeked bytes consumes them and more. + require.NoError(t, r.Skip(8)) + require.Equal(t, 16, r.Offset()) +} + +func TestBucketAsyncBufReader_Peek_Read(t *testing.T) { + r, _ := newTestAsyncBufReader(t, 0, len(testBucketContents)) + + peek1, err := r.Peek(5) + require.NoError(t, err) + require.Equal(t, testBucketContents[:5], peek1) + require.Equal(t, 0, r.Offset()) + + // Read returns the same bytes. + read1, err := r.Read(5) + require.NoError(t, err) + require.Equal(t, testBucketContents[:5], read1) + require.Equal(t, 5, r.Offset()) + + // Another Peek returns the next bytes after Read consumed previous bytes. + peek2, err := r.Peek(5) + require.NoError(t, err) + require.Equal(t, testBucketContents[5:10], peek2) + require.Equal(t, 5, r.Offset()) + + // A Read short of the peeked bytes consumes some of them. + read2, err := r.Read(3) + require.NoError(t, err) + require.Equal(t, testBucketContents[5:8], read2) + require.Equal(t, 8, r.Offset()) + + // A Read beyond the remaining peeked bytes consumes them and more. + read3, err := r.Read(8) + require.NoError(t, err) + require.Equal(t, testBucketContents[8:16], read3) + require.Equal(t, 16, r.Offset()) +} + +func TestBucketAsyncBufReader_Peek_AcrossPromiseBoundary(t *testing.T) { + // Each buffer promise holds testBufPoolSize bytes. + // A peek of 8 bytes at offset 12 takes 4 bytes from the first promise + // and 4 bytes from the second promise. + r, _ := newTestAsyncBufReader(t, 0, len(testBucketContents)) + require.NoError(t, r.Skip(testBufPoolSize-4)) + + b, err := r.Peek(8) + require.NoError(t, err) + require.Equal(t, testBucketContents[testBufPoolSize-4:testBufPoolSize+4], b) + require.Equal(t, testBufPoolSize-4, r.Offset(), "Peek does not consume") +} + +func TestBucketAsyncBufReader_Peek_PastSegmentEnd(t *testing.T) { + // Peek must return the bytes read and suppress the EOF error + // when peeking past the configured length or peeking beyond the true end of the object. + const sectionLen = 5 + r, _ := newTestAsyncBufReader(t, 0, sectionLen) + + b, err := r.Peek(sectionLen + 5) + require.NoError(t, err) + require.Equal(t, testBucketContents[:sectionLen], b) + require.Equal(t, 0, r.Offset()) +} + +func TestBucketAsyncBufReader_Peek_AtEnd(t *testing.T) { + const sectionLen = 5 + r, _ := newTestAsyncBufReader(t, 0, sectionLen) + require.NoError(t, r.Skip(sectionLen)) + + b, err := r.Peek(1) + require.NoError(t, err) + require.Nil(t, b) +} + +func TestBucketAsyncBufReader_Read_Sequential(t *testing.T) { + r, _ := newTestAsyncBufReader(t, 0, len(testBucketContents)) + + b, err := r.Read(5) + require.NoError(t, err) + require.Equal(t, testBucketContents[:5], b) + require.Equal(t, 5, r.Offset()) + + b, err = r.Read(5) + require.NoError(t, err) + require.Equal(t, testBucketContents[5:10], b) + require.Equal(t, 10, r.Offset()) +} + +func TestBucketAsyncBufReader_Read_ExactLength(t *testing.T) { + r, _ := newTestAsyncBufReader(t, 0, len(testBucketContents)) + + b, err := r.Read(len(testBucketContents)) + require.NoError(t, err) + require.Equal(t, testBucketContents, b) + require.Equal(t, len(testBucketContents), r.Offset()) + require.Equal(t, 0, r.Len()) +} + +func TestBucketAsyncBufReader_Read_BeyondEnd(t *testing.T) { + const sectionLen = 10 + r, _ := newTestAsyncBufReader(t, 0, sectionLen) + + b, err := r.Read(sectionLen + 1) + require.ErrorIs(t, err, ErrInvalidSize) + require.Nil(t, b) + require.Equal(t, sectionLen, r.Offset()) +} + +//func TestBucketAsyncBufReader_Peek_BeyondPeekBuffer(t *testing.T) { +// r, _ := newTestAsyncBufReader(t, 0, len(testBucketContents)) +// +// b, err := r.Peek(r.Size() + 1) +// require.ErrorIs(t, err, ErrInvalidSize) +// require.Nil(t, b) +//} + +// TestBucketAsyncBufReader_Peek_ThenSkip covers the access pattern of +// Decbuf.UnsafeUvarintBytes, which peeks bytes and then skips exactly those bytes. +// The skip drains the head promise and returns the buffer of that promise to the pool. +// The reader then takes that same buffer back and fills it with the third chunk. +// Peek copies into peekBuf, so the peeked bytes stay valid through that refill. +func TestBucketAsyncBufReader_Peek_ThenSkip(t *testing.T) { + // The read-ahead window holds maxBufCount*testBufPoolSize = 32 bytes of the 48-byte segment, + // so the reader must refill a buffer to reach the third chunk. + const ( + length = 48 + maxBufCount = 2 + ) + r, _ := newTestAsyncBufReaderWithData(t, testBucketContentsLong, 0, length, maxBufCount) + + // Peek and skip the whole head promise, which drains and releases it. + b, err := r.Peek(testBufPoolSize) + require.NoError(t, err) + require.NoError(t, r.Skip(len(b))) + require.Equal(t, testBufPoolSize, r.Offset()) + + // Read the second chunk, then one byte of the third chunk. + // The read of the third chunk waits for the refill of the released buffer. + got, err := r.Read(testBufPoolSize) + require.NoError(t, err) + require.Equal(t, testBucketContentsLong[testBufPoolSize:2*testBufPoolSize], got) + + got, err = r.Read(1) + require.NoError(t, err) + require.Equal(t, testBucketContentsLong[2*testBufPoolSize:2*testBufPoolSize+1], got) + + require.Equal(t, testBucketContentsLong[:testBufPoolSize], b, "peeked bytes survive the refill") +} + +// TestBucketAsyncBufReader_Peek_ThenSkip_AcrossPromiseBoundary makes sure that a skip +// after a straddling peek discards the bytes that the peek returned. +// Peek does not consume, so the cursor still points at the first peeked byte. +func TestBucketAsyncBufReader_Peek_ThenSkip_AcrossPromiseBoundary(t *testing.T) { + // A peek of 8 bytes at offset 12 takes 4 bytes from the first promise + // and 4 bytes from the second promise, so it goes through peekBuf. + const peekAt = testBufPoolSize - 4 + r, _ := newTestAsyncBufReader(t, 0, len(testBucketContents)) + require.NoError(t, r.Skip(peekAt)) + + b, err := r.Peek(8) + require.NoError(t, err) + require.Equal(t, testBucketContents[peekAt:peekAt+8], b) + + require.NoError(t, r.Skip(len(b))) + require.Equal(t, peekAt+8, r.Offset(), "the skip consumes the peeked bytes, not the bytes after them") + require.Equal(t, testBucketContents[peekAt:peekAt+8], b, "the skip does not disturb peekBuf") + + // The next read starts one byte past the peeked bytes. + got, err := r.Read(4) + require.NoError(t, err) + require.Equal(t, testBucketContents[peekAt+8:peekAt+12], got) +} + +func TestBucketAsyncBufReader_Skip_Basic(t *testing.T) { + r, _ := newTestAsyncBufReader(t, 0, len(testBucketContents)) + + require.NoError(t, r.Skip(10)) + require.Equal(t, 10, r.Offset()) + require.Equal(t, len(testBucketContents)-10, r.Len()) + + b, err := r.Read(3) + require.NoError(t, err) + require.Equal(t, testBucketContents[10:13], b) +} + +func TestBucketAsyncBufReader_Skip_AcrossPromises(t *testing.T) { + // A skip of 20 bytes drains the first promise of testBufPoolSize bytes, + // then takes the remainder from the second promise. + r, _ := newTestAsyncBufReader(t, 0, len(testBucketContents)) + + require.NoError(t, r.Skip(testBufPoolSize+4)) + require.Equal(t, testBufPoolSize+4, r.Offset()) + + b, err := r.Read(3) + require.NoError(t, err) + require.Equal(t, testBucketContents[testBufPoolSize+4:testBufPoolSize+7], b) +} + +func TestBucketAsyncBufReader_Skip_ToEnd(t *testing.T) { + r, _ := newTestAsyncBufReader(t, 0, len(testBucketContents)) + + require.NoError(t, r.Skip(len(testBucketContents))) + require.Equal(t, len(testBucketContents), r.Offset()) + require.Equal(t, 0, r.Len()) +} + +func TestBucketAsyncBufReader_Skip_BeyondEnd(t *testing.T) { + const sectionLen = 10 + r, _ := newTestAsyncBufReader(t, 0, sectionLen) + + require.ErrorIs(t, r.Skip(sectionLen+1), ErrInvalidSize) +} + +//func TestBucketAsyncBufReader_Skip_Negative(t *testing.T) { +// // Decbuf.SkipUvarintBytes converts a uint64 from the object to an int, +// // so a corrupt length can arrive as a negative number. +// r, _ := newTestAsyncBufReader(t, 0, len(testBucketContents)) +// +// require.ErrorIs(t, r.Skip(-1), ErrInvalidSize) +// require.Equal(t, 0, r.Offset()) +//} + +func TestBucketAsyncBufReader_Read_ExactLength_NonZeroBaseWithRotation(t *testing.T) { + // The read-ahead window holds maxBufCount*testBufPoolSize = 32 bytes. + // A length of 48 forces the reader to release one buffer promise and refill it. + // The refilled promise must add base to the buffered offset. + // Without base, the promise reads the object 4 bytes too early. + const ( + base = 4 + length = 48 + maxBufCount = 2 + ) + r, _ := newTestAsyncBufReaderWithData(t, testBucketContentsLong, base, length, maxBufCount) + + b, err := r.Read(length) + require.NoError(t, err) + require.Equal(t, testBucketContentsLong[base:base+length], b) + require.Equal(t, length, r.Offset()) + require.Equal(t, 0, r.Len()) +} diff --git a/pkg/storage/indexheader/encoding/bucket_factory.go b/pkg/storage/indexheader/encoding/bucket_factory.go index 7088b227f99..063598b57bb 100644 --- a/pkg/storage/indexheader/encoding/bucket_factory.go +++ b/pkg/storage/indexheader/encoding/bucket_factory.go @@ -88,7 +88,9 @@ func (bf *BucketDecbufFactory) NewDecbufAtChecked(ctx context.Context, offset in bufferLength := numLenBytes + contentLength + crc32.Size - bufReader := NewBucketBufReader(ctx, bf.bkt, bf.objectPath, offset, bufferLength) + //bufReader := NewBucketBufReader(ctx, bf.bkt, bf.objectPath, offset, bufferLength) + bufReader := NewBucketAsyncBufReader(ctx, bf.bkt, bf.objectPath, offset, bufferLength) + // bufReader is expected start at base offset + 4 after consuming length bytes err = bufReader.Skip(numLenBytes) if err != nil { @@ -128,7 +130,7 @@ func (bf *BucketDecbufFactory) NewDecbufInSection(ctx context.Context, tableOffs if sectionLength <= 0 { return Decbuf{E: fmt.Errorf("section length must be greater than 0")} } - bufReader := NewBucketBufReader( + bufReader := NewBucketAsyncBufReader( ctx, bf.bkt, bf.objectPath, @@ -151,7 +153,7 @@ func (bf *BucketDecbufFactory) NewRawDecbuf(ctx context.Context) Decbuf { return Decbuf{E: fmt.Errorf("get size from %s: %w", bf.objectPath, err)} } // Create reader from full file range - r := NewBucketBufReader( + r := NewBucketAsyncBufReader( ctx, bf.bkt, bf.objectPath, offset, int(attrs.Size), ) d := Decbuf{r: r} diff --git a/pkg/storage/indexheader/encoding/bucket_reader.go b/pkg/storage/indexheader/encoding/bucket_reader.go index 30e973131db..21c90be2128 100644 --- a/pkg/storage/indexheader/encoding/bucket_reader.go +++ b/pkg/storage/indexheader/encoding/bucket_reader.go @@ -13,6 +13,16 @@ import ( "github.com/thanos-io/objstore" ) +const ( + ReadBufferSize = 1 << 20 // 1 MiB +) + +var bucketBufioPool = sync.Pool{ + New: func() any { + return bufio.NewReaderSize(nil, ReadBufferSize) + }, +} + type BucketReader struct { ctx context.Context bkt objstore.BucketReader @@ -34,14 +44,14 @@ func NewBucketReader( } } -func (r *BucketReader) Read(p []byte) (n int, err error) { - if len(p) == 0 { +func (r *BucketReader) Read(dst []byte) (n int, err error) { + if len(dst) == 0 { return 0, nil } if r.off >= r.length { return 0, io.EOF } - toRead := len(p) + toRead := len(dst) remaining := r.length - r.off if toRead > remaining { toRead = remaining @@ -51,7 +61,7 @@ func (r *BucketReader) Read(p []byte) (n int, err error) { return 0, err } defer rc.Close() - n, err = io.ReadFull(rc, p[:toRead]) + n, err = io.ReadFull(rc, dst[:toRead]) r.off += n if errors.Is(err, io.ErrUnexpectedEOF) { err = io.EOF @@ -70,44 +80,23 @@ func (r *BucketReader) Seek(offset int64, whence int) (int64, error) { return offset, nil } -var bucketBufPool = sync.Pool{ - New: func() any { - // 1MiB buffer chosen as starting point; - // we could make this configurable and benchmark. - return bufio.NewReaderSize(nil, 1<<20) - }, -} - type BucketBufReader struct { - ctx context.Context - bkt objstore.BucketReader - name string - base int - length int - off int - r *BucketReader - resetReader func(off int) error - buf *bufio.Reader - // Hold a reference to the pool for returning on Close - allows tests to use different pool. - bufPool *sync.Pool -} - -func resetReaderFunc(bufReader *BucketBufReader) func(off int) error { - return func(off int) error { - r := NewBucketReader(bufReader.ctx, bufReader.bkt, bufReader.name, bufReader.base, bufReader.length) - _, err := r.Seek(int64(off), io.SeekStart) - if err != nil { - return err - } - bufReader.r = r - return nil - } + ctx context.Context + bkt objstore.BucketReader + name string + base int + length int + off int + + r *BucketReader + buf *bufio.Reader + bufPool *sync.Pool // Reference to return to on Close } func NewBucketBufReader( ctx context.Context, bkt objstore.BucketReader, name string, base int, length int, ) *BucketBufReader { - return newBucketBufReader(ctx, &bucketBufPool, bkt, name, base, length) + return newBucketBufReader(ctx, &bucketBufioPool, bkt, name, base, length) } func newBucketBufReader( @@ -128,7 +117,6 @@ func newBucketBufReader( bufPool: bufioPool, } - bufReader.resetReader = resetReaderFunc(bufReader) return bufReader } @@ -142,15 +130,17 @@ func (bbr *BucketBufReader) ResetAt(off int) error { } if dist := off - bbr.off; dist > 0 && dist < bbr.Buffered() { - // skip ahead by discarding the distance bytes + // Reset via Skip to avoid discarding all buffered bytes. return bbr.Skip(dist) } - if err := bbr.resetReader(off); err != nil { + r := NewBucketReader(bbr.ctx, bbr.bkt, bbr.name, bbr.base, bbr.length) + _, err := r.Seek(int64(off), io.SeekStart) + if err != nil { return err } - - bbr.buf.Reset(bbr.r) + bbr.r = r + bbr.buf.Reset(r) bbr.off = off return nil @@ -162,9 +152,7 @@ func (bbr *BucketBufReader) Skip(l int) error { } n, err := bbr.buf.Discard(l) - if n > 0 { - bbr.off += n - } + bbr.off += n return err } @@ -195,14 +183,14 @@ func (bbr *BucketBufReader) Read(n int) ([]byte, error) { return b, nil } -func (bbr *BucketBufReader) ReadInto(b []byte) error { - n, err := io.ReadFull(bbr.buf, b) +func (bbr *BucketBufReader) ReadInto(dst []byte) error { + n, err := io.ReadFull(bbr.buf, dst) if n > 0 { bbr.off += n } if errors.Is(err, io.EOF) || errors.Is(err, io.ErrUnexpectedEOF) { - return fmt.Errorf("%w reading %d bytes: %s", ErrInvalidSize, len(b), err) + return fmt.Errorf("%w reading %d bytes: %s", ErrInvalidSize, len(dst), err) } else if err != nil { return err } @@ -227,10 +215,9 @@ func (bbr *BucketBufReader) Buffered() int { } func (bbr *BucketBufReader) Close() error { - // Note that we don't do anything to clean up the buffer before returning it to the pool here: - // we reset the buffer when we retrieve it from the pool instead. + // No need to clean up buffer, we reset when we retrieve it from the pool bbr.bufPool.Put(bbr.buf) - // The BucketReader does not need closed - - // it closes the reader generated from bkt.GetRange on each Read call. + // The BucketReader of a promise does not need a Close call. + // It closes the reader created by bkt.GetRange on each Read call. return nil } diff --git a/pkg/storage/indexheader/encoding/bucket_reader_test.go b/pkg/storage/indexheader/encoding/bucket_reader_test.go index 5310428eb4d..e6a688dcd66 100644 --- a/pkg/storage/indexheader/encoding/bucket_reader_test.go +++ b/pkg/storage/indexheader/encoding/bucket_reader_test.go @@ -8,6 +8,7 @@ import ( "context" "errors" "io" + "slices" "sync" "testing" @@ -278,12 +279,12 @@ func TestBucketBufReader_GetRangeCalls_Buffering(t *testing.T) { _, err := r.Read(1) require.NoError(t, err) } - require.Len(t, bkt.calls, 1) + require.Len(t, bkt.rangeCalls(), 1) // Buffer depleted; next read operation triggers another bufio fill and GetRange call. _, err := r.Read(1) require.NoError(t, err) - require.Len(t, bkt.calls, 2) + require.Len(t, bkt.rangeCalls(), 2) } func TestBucketBufReader_GetRangeCalls_ResetRefetches(t *testing.T) { @@ -291,13 +292,13 @@ func TestBucketBufReader_GetRangeCalls_ResetRefetches(t *testing.T) { _, err := r.Read(1) require.NoError(t, err) - require.Len(t, bkt.calls, 1) + require.Len(t, bkt.rangeCalls(), 1) // After Reset the buffer is discarded; the next read must refetch from the bucket. require.NoError(t, r.Reset()) _, err = r.Read(1) require.NoError(t, err) - require.Len(t, bkt.calls, 2) + require.Len(t, bkt.rangeCalls(), 2) } func TestBucketBufReader_Read_GetRangeError(t *testing.T) { @@ -317,8 +318,12 @@ func TestBucketBufReader_ReadInto_GetRangeError(t *testing.T) { } // trackingBucket wraps an InstrumentedBucketReader and records every GetRange call. +// The read-ahead reader fills its buffer promises from several goroutines at the same time, +// so the mutex protects the record of the calls. type trackingBucket struct { objstore.InstrumentedBucketReader + + mtx sync.Mutex calls []rangeCall } @@ -328,10 +333,19 @@ type rangeCall struct { } func (b *trackingBucket) GetRange(ctx context.Context, name string, off, length int64) (io.ReadCloser, error) { + b.mtx.Lock() b.calls = append(b.calls, rangeCall{off, length}) + b.mtx.Unlock() return b.InstrumentedBucketReader.GetRange(ctx, name, off, length) } +// rangeCalls returns a copy of the recorded GetRange calls. +func (b *trackingBucket) rangeCalls() []rangeCall { + b.mtx.Lock() + defer b.mtx.Unlock() + return slices.Clone(b.calls) +} + func newTrackingBucket(t *testing.T, objectData []byte) *trackingBucket { t.Helper() inmem := objstore.NewInMemBucket() diff --git a/pkg/storage/indexheader/encoding/file_reader.go b/pkg/storage/indexheader/encoding/file_reader.go index 5fe6bdb1d83..0ae0abba2cc 100644 --- a/pkg/storage/indexheader/encoding/file_reader.go +++ b/pkg/storage/indexheader/encoding/file_reader.go @@ -110,14 +110,14 @@ func (f *FileReader) Read(n int) ([]byte, error) { return b, nil } -func (f *FileReader) ReadInto(b []byte) error { - r, err := io.ReadFull(f.buf, b) +func (f *FileReader) ReadInto(dst []byte) error { + r, err := io.ReadFull(f.buf, dst) if r > 0 { f.off += r } if errors.Is(err, io.EOF) || errors.Is(err, io.ErrUnexpectedEOF) { - return fmt.Errorf("%w reading %d bytes: %s", ErrInvalidSize, len(b), err) + return fmt.Errorf("%w reading %d bytes: %s", ErrInvalidSize, len(dst), err) } else if err != nil { return err } @@ -143,8 +143,7 @@ func (f *FileReader) Buffered() int { // Close cleans up the underlying resources used by this FileReader. func (f *FileReader) Close() error { - // Note that we don't do anything to clean up the buffer before returning it to the pool here: - // we reset the buffer when we retrieve it from the pool instead. + // No need to clean up buffer, we reset when we retrieve it from the pool bufferPool.Put(f.buf) // File handles are pooled, so we don't actually close the handle here, just return it. return f.closer.Put(f.file) diff --git a/pkg/storage/indexheader/encoding/reader.go b/pkg/storage/indexheader/encoding/reader.go index 9c3222b1ab4..1646e40bb0b 100644 --- a/pkg/storage/indexheader/encoding/reader.go +++ b/pkg/storage/indexheader/encoding/reader.go @@ -19,14 +19,24 @@ type BufReader interface { ResetAt(off int) error // Skip advances the cursor by the given number of bytes in the data segment. - // Attempting to skip to the end of the data segment is valid. - // Attempting to skip _beyond_ the end of the data segment will return an error. + // It is valid to skip to exactly the end of the data segment. + // It is NOT valid to skip beyond the end of the data segment; + // in this case implementations MUST return an ErrInvalidSize error, + // but MUST NOT advance the cursor or consume any remaining bytes. Skip(l int) error // Peek returns at most the given number of bytes from the data segment, without consuming them. - // The byte slice returned becomes invalid at the next read. // It is valid to Peek beyond the end of the data segment; // in this case implementations MUST return the available bytes up to the end and a nil error. + // + // The byte slice returned MUST remain valid for one subsequent Skip of the returned byte length; + // callers use a Peek-Skip pattern in place of Read to avoid a slice allocation. + // It is NOT valid to read the returned byte slice after any subsequent read operation: + // Peek, Read, ReadInto, Reset, ResetAt, and Skip. + // + // Peek is limited to and only must support reads up to the underlying buffer length. + // Caller checks Size first to see if the next read operation length fits in the underlying buffer, + // then a Peek-Skip pattern is used to avoid the slice allocation which must occur in Read. Peek(n int) ([]byte, error) // Read returns the given number of bytes from the data segment, consuming them. @@ -35,11 +45,11 @@ type BufReader interface { // and the remaining bytes MUST be consumed. Read(n int) ([]byte, error) - // ReadInto reads len(b) bytes from the data segment into b, consuming them. + // ReadInto reads len(dst) bytes from the data segment into dst, consuming them. // It is NOT valid to read beyond the end of the data segment; // in this case implementations MUST return a nil byte slice and an ErrInvalidSize error, // and the remaining bytes MUST be consumed. - ReadInto(b []byte) error + ReadInto(dst []byte) error // Size returns the length of the underlying buffer in bytes. Size() int