diff --git a/weed/filer/reader_at.go b/weed/filer/reader_at.go index 9e31228a7..9ac7ecad2 100644 --- a/weed/filer/reader_at.go +++ b/weed/filer/reader_at.go @@ -32,7 +32,9 @@ type ChunkReadAt struct { fileSize int64 readerCache *ReaderCache readerPattern *ReaderPattern + lastChunkMu sync.Mutex // guards lastChunkFid; mount issues concurrent ReadAt calls lastChunkFid string + stream chunkStream // chunk this reader is positioned in, pinned in the shared readerCache prefetchCount int // Number of chunks to prefetch ahead during sequential reads ctx context.Context // Context used for cancellation during chunk read operations } @@ -341,6 +343,7 @@ func (c *ChunkReadAt) doReadAt(ctx context.Context, p []byte, offset int64) (n i func (c *ChunkReadAt) readChunkSliceAt(ctx context.Context, buffer []byte, chunkView *ChunkView, nextChunkViews *Interval[*ChunkView], offset uint64) (n int, err error) { if c.readerPattern.IsRandomMode() { + c.readerCache.releaseStream(&c.stream) n, err := c.readerCache.chunkCache.ReadChunkAt(buffer, chunkView.FileId, offset) if n > 0 { return n, err @@ -350,20 +353,20 @@ func (c *ChunkReadAt) readChunkSliceAt(ctx context.Context, buffer []byte, chunk } shouldCache := (uint64(chunkView.ViewOffset) + chunkView.ChunkSize) <= c.readerCache.chunkCache.GetMaxFilePartSizeInCache() - n, err = c.readerCache.ReadChunkAt(ctx, buffer, chunkView.FileId, chunkView.CipherKey, chunkView.IsGzipped, int64(offset), int(chunkView.ChunkSize), shouldCache) - if c.lastChunkFid != chunkView.FileId { - if chunkView.OffsetInChunk == 0 { // start of a new chunk - if c.lastChunkFid != "" { - c.readerCache.UnCache(c.lastChunkFid) - } - if nextChunkViews != nil && c.prefetchCount > 0 { - // Prefetch multiple chunks ahead for better sequential read throughput - // This keeps the network pipeline full with parallel chunk fetches - c.readerCache.MaybeCache(nextChunkViews, c.prefetchCount) - } + // The previous chunk is released through the stream pin rather than + // UnCache: the buffer is shared, and other streams may still be reading it. + n, err = c.readerCache.readChunkAt(ctx, &c.stream, buffer, chunkView.FileId, chunkView.CipherKey, chunkView.IsGzipped, int64(offset), int(chunkView.ChunkSize), shouldCache) + c.lastChunkMu.Lock() + enteredChunk := c.lastChunkFid != chunkView.FileId + c.lastChunkFid = chunkView.FileId + c.lastChunkMu.Unlock() + if enteredChunk && chunkView.OffsetInChunk == 0 { // start of a new chunk + if nextChunkViews != nil && c.prefetchCount > 0 { + // Prefetch multiple chunks ahead for better sequential read throughput + // This keeps the network pipeline full with parallel chunk fetches + c.readerCache.MaybeCache(nextChunkViews, c.prefetchCount) } } - c.lastChunkFid = chunkView.FileId return } diff --git a/weed/filer/reader_at_shared_test.go b/weed/filer/reader_at_shared_test.go new file mode 100644 index 000000000..22937a3b8 --- /dev/null +++ b/weed/filer/reader_at_shared_test.go @@ -0,0 +1,209 @@ +package filer + +import ( + "context" + "fmt" + "io" + "sync" + "sync/atomic" + "testing" + + "github.com/seaweedfs/seaweedfs/weed/util/chunk_cache" + util_http "github.com/seaweedfs/seaweedfs/weed/util/http" +) + +// Two streams (e.g. two S3 GETs) reading the same object through one shared +// ReaderCache, interleaved in small slices, must share each chunk download: +// one stream finishing a chunk, or moving on to the next one, must not drop a +// buffer the other stream is still reading. +func TestChunkReadAtConcurrentStreamsShareChunks(t *testing.T) { + const chunkSize = 64 << 10 + const chunkCount = 3 + const sliceSize = 16 << 10 + + var mu sync.Mutex + fetches := map[string]int{} + rc := NewReaderCache(64, (*chunk_cache.TieredChunkCache)(nil), func(context.Context, string) ([]string, error) { + return []string{"unused"}, nil + }, nil) + defer rc.destroy() + rc.fetchChunkDataFn = func(_ context.Context, buffer []byte, _ []string, _ []byte, _ bool, _ bool, _ int64, fileId string, _ util_http.RefreshUrlsFunc) (int, error) { + mu.Lock() + fetches[fileId]++ + mu.Unlock() + for i := range buffer { + buffer[i] = fileId[len(fileId)-1] + } + return len(buffer), nil + } + + newStream := func() *ChunkReadAt { + views := NewIntervalList[*ChunkView]() + for i := 0; i < chunkCount; i++ { + views.AppendInterval(&Interval[*ChunkView]{ + StartOffset: int64(i * chunkSize), + StopOffset: int64((i + 1) * chunkSize), + Value: &ChunkView{ + FileId: fmt.Sprintf("chunk%d", i), + ViewSize: chunkSize, + ViewOffset: int64(i * chunkSize), + ChunkSize: chunkSize, + }, + }) + } + return NewChunkReaderAtFromClient(context.Background(), rc, views, chunkSize*chunkCount, 0) + } + streams := []*ChunkReadAt{newStream(), newStream()} + + // The streams alternate slice by slice, the second one lagging one slice + // behind, so the leader always finishes a chunk while the follower is + // still inside it. + offsets := make([]int64, len(streams)) + offsets[1] = -sliceSize + for offsets[1] < chunkSize*chunkCount { + for i, stream := range streams { + if offsets[i] < 0 || offsets[i] >= chunkSize*chunkCount { + offsets[i] += sliceSize + continue + } + buf := make([]byte, sliceSize) + n, err := stream.ReadAt(buf, offsets[i]) + if (err != nil && err != io.EOF) || n != sliceSize { + t.Fatalf("stream %d at %d: n=%d err=%v", i, offsets[i], n, err) + } + if want := byte('0' + offsets[i]/chunkSize); buf[0] != want || buf[n-1] != want { + t.Fatalf("stream %d at %d: got %q, want %q", i, offsets[i], buf[0], want) + } + offsets[i] += sliceSize + } + } + + mu.Lock() + defer mu.Unlock() + for i := 0; i < chunkCount; i++ { + fileId := fmt.Sprintf("chunk%d", i) + if fetches[fileId] != 1 { + t.Errorf("%s fetched %d times for two concurrent streams, want 1", fileId, fetches[fileId]) + } + } +} + +func newStreamTestReaderCache(chunkCache chunk_cache.ChunkCache) *ReaderCache { + rc := NewReaderCache(64, chunkCache, func(context.Context, string) ([]string, error) { + return []string{"unused"}, nil + }, nil) + rc.fetchChunkDataFn = func(_ context.Context, buffer []byte, _ []string, _ []byte, _ bool, _ bool, _ int64, _ string, _ util_http.RefreshUrlsFunc) (int, error) { + return len(buffer), nil + } + return rc +} + +func isRetained(rc *ReaderCache, fileId string) bool { + rc.Lock() + defer rc.Unlock() + _, found := rc.downloaders[fileId] + return found +} + +// A stream moving on to a chunk served from the chunk cache must still +// release the chunk it was positioned in before. +func TestChunkStreamReleasesPreviousChunkOnCacheHit(t *testing.T) { + cache := newMockChunkCacheForReaderCache() + cache.SetChunk("chunk1", make([]byte, 4<<10)) + rc := newStreamTestReaderCache(cache) + defer rc.destroy() + + stream := &chunkStream{} + if _, err := rc.readChunkAt(context.Background(), stream, make([]byte, 1<<10), "chunk0", nil, false, 0, 4<<10, false); err != nil { + t.Fatal(err) + } + if n, err := rc.readChunkAt(context.Background(), stream, make([]byte, 1<<10), "chunk1", nil, false, 0, 4<<10, true); err != nil || n == 0 { + t.Fatalf("cache hit read: n=%d err=%v", n, err) + } + if isRetained(rc, "chunk0") { + t.Fatal("chunk0 still retained after the stream moved on to a cached chunk") + } +} + +// A chunk the stream leaves while another read is in flight must be dropped +// once that read ends, even if it did not read the chunk to the end. +func TestChunkStreamDropsLeftChunkAfterInFlightRead(t *testing.T) { + rc := newStreamTestReaderCache(newMockChunkCacheForReaderCache()) + defer rc.destroy() + + stream := &chunkStream{} + if _, err := rc.readChunkAt(context.Background(), stream, make([]byte, 1<<10), "chunk0", nil, false, 0, 4<<10, false); err != nil { + t.Fatal(err) + } + // Another reader is in the middle of a partial read of chunk0. + rc.Lock() + other := rc.downloaders["chunk0"] + other.wg.Add(1) + atomic.AddInt32(&other.readers, 1) + rc.Unlock() + + if _, err := rc.readChunkAt(context.Background(), stream, make([]byte, 1<<10), "chunk1", nil, false, 0, 4<<10, false); err != nil { + t.Fatal(err) + } + if !isRetained(rc, "chunk0") { + t.Fatal("chunk0 dropped while a read was still in flight") + } + + // The in-flight read ends without reaching the end of the chunk. + other.wg.Done() + atomic.AddInt32(&other.readers, -1) + rc.removeConsumed(other) + if isRetained(rc, "chunk0") { + t.Fatal("chunk0 retained after the stream left it and the last read ended") + } +} + +// Concurrent ReadAt calls on one ChunkReadAt (as mount does) share its stream +// pin; they must neither race on it nor leak or double-release pins. +func TestChunkStreamConcurrentReadsOnOneReader(t *testing.T) { + const chunkSize = 16 << 10 + const chunkCount = 4 + const sliceSize = 4 << 10 + rc := newStreamTestReaderCache((*chunk_cache.TieredChunkCache)(nil)) + defer rc.destroy() + + views := NewIntervalList[*ChunkView]() + for i := 0; i < chunkCount; i++ { + views.AppendInterval(&Interval[*ChunkView]{ + StartOffset: int64(i * chunkSize), + StopOffset: int64((i + 1) * chunkSize), + Value: &ChunkView{ + FileId: fmt.Sprintf("chunk%d", i), + ViewSize: chunkSize, + ViewOffset: int64(i * chunkSize), + ChunkSize: chunkSize, + }, + }) + } + reader := NewChunkReaderAtFromClient(context.Background(), rc, views, chunkSize*chunkCount, 0) + + var wg sync.WaitGroup + for g := 0; g < 8; g++ { + wg.Add(1) + go func() { + defer wg.Done() + for offset := int64(0); offset < chunkSize*chunkCount; offset += sliceSize { + if _, err := reader.ReadAt(make([]byte, sliceSize), offset); err != nil && err != io.EOF { + t.Error(err) + return + } + } + }() + } + wg.Wait() + + // Finishing the last chunk releases the stream's final pin. + if _, err := reader.ReadAt(make([]byte, sliceSize), chunkSize*chunkCount-sliceSize); err != nil && err != io.EOF { + t.Fatal(err) + } + rc.Lock() + defer rc.Unlock() + for fileId, cacher := range rc.downloaders { + t.Errorf("%s retained after all reads finished: pins=%d readers=%d", fileId, atomic.LoadInt32(&cacher.pins), atomic.LoadInt32(&cacher.readers)) + } +} diff --git a/weed/filer/reader_cache.go b/weed/filer/reader_cache.go index b54ff6ecf..93b7883e2 100644 --- a/weed/filer/reader_cache.go +++ b/weed/filer/reader_cache.go @@ -35,6 +35,8 @@ type SingleChunkCacher struct { completedTimeNew int64 readers int32 consumed int32 + pins int32 // streams currently positioned inside this chunk + left int32 // set once the last pinning stream moved on sync.Mutex parent *ReaderCache chunkFileId string @@ -112,7 +114,19 @@ func (rc *ReaderCache) MaybeCache(chunkViews *Interval[*ChunkView], count int) { return } +// chunkStream is one sequential reader's position in a shared ReaderCache. +// The chunk it is reading stays pinned until the stream reads it to the end or +// moves elsewhere, so another stream finishing or leaving the same chunk does +// not drop the buffer from under it. +type chunkStream struct { + cacher *SingleChunkCacher +} + func (rc *ReaderCache) ReadChunkAt(ctx context.Context, buffer []byte, fileId string, cipherKey []byte, isGzipped bool, offset int64, chunkSize int, shouldCache bool) (int, error) { + return rc.readChunkAt(ctx, nil, buffer, fileId, cipherKey, isGzipped, offset, chunkSize, shouldCache) +} + +func (rc *ReaderCache) readChunkAt(ctx context.Context, stream *chunkStream, buffer []byte, fileId string, cipherKey []byte, isGzipped bool, offset int64, chunkSize int, shouldCache bool) (int, error) { retry: rc.Lock() @@ -130,8 +144,11 @@ retry: // start wg.Wait() on a zero counter while this read is about to register. cacher.wg.Add(1) atomic.AddInt32(&cacher.readers, 1) + previous := stream.pin(cacher) rc.Unlock() + rc.unpin(previous) n, err := cacher.readChunkAt(ctx, buffer, offset) + rc.releaseIfFinished(stream, cacher, offset, n, err, chunkSize) if n > 0 || err != nil { return n, err } @@ -144,7 +161,10 @@ retry: if shouldCache || rc.lookupFileIdFn == nil { n, err := rc.chunkCache.ReadChunkAt(buffer, fileId, uint64(offset)) if n > 0 { + // Served from the chunk cache: the stream has left its pinned chunk. + previous := stream.unpinLocked() rc.Unlock() + rc.unpin(previous) return n, err } } @@ -175,9 +195,77 @@ retry: rc.downloaders[fileId] = cacher cacher.wg.Add(1) atomic.AddInt32(&cacher.readers, 1) + previous := stream.pin(cacher) rc.Unlock() + rc.unpin(previous) - return cacher.readChunkAt(ctx, buffer, offset) + n, err := cacher.readChunkAt(ctx, buffer, offset) + rc.releaseIfFinished(stream, cacher, offset, n, err, chunkSize) + return n, err +} + +// pin makes cacher the stream's current chunk and returns the chunk it was +// pinned to before, which the caller unpins once the ReaderCache lock is +// released. The stream is only touched under the ReaderCache lock, since +// concurrent ReadAt calls on one ChunkReadAt share it. +func (stream *chunkStream) pin(cacher *SingleChunkCacher) (previous *SingleChunkCacher) { + if stream == nil || stream.cacher == cacher { + return nil + } + previous = stream.cacher + stream.cacher = cacher + atomic.AddInt32(&cacher.pins, 1) + return previous +} + +// unpinLocked detaches the stream from its chunk and returns that chunk for +// the caller to unpin once the ReaderCache lock is released. +func (stream *chunkStream) unpinLocked() (previous *SingleChunkCacher) { + if stream == nil { + return nil + } + previous = stream.cacher + stream.cacher = nil + return previous +} + +// releaseIfFinished unpins the stream's chunk once the stream has read it to +// the end, since the stream will not come back to it. +func (rc *ReaderCache) releaseIfFinished(stream *chunkStream, cacher *SingleChunkCacher, offset int64, n int, err error, chunkSize int) { + if stream == nil || err != nil || offset+int64(n) < int64(chunkSize) { + return + } + var previous *SingleChunkCacher + rc.Lock() + if stream.cacher == cacher { + previous = stream.unpinLocked() + } + rc.Unlock() + rc.unpin(previous) +} + +// releaseStream unpins whatever chunk the stream is positioned in. +func (rc *ReaderCache) releaseStream(stream *chunkStream) { + if stream == nil { + return + } + rc.Lock() + previous := stream.unpinLocked() + rc.Unlock() + rc.unpin(previous) +} + +// unpin drops one stream's pin. Once no stream is positioned in the chunk it +// is dropped like UnCache would, as soon as no read is in progress either: +// here if none is, otherwise by the last read's removeConsumed. +func (rc *ReaderCache) unpin(downloader *SingleChunkCacher) { + if downloader == nil { + return + } + if atomic.AddInt32(&downloader.pins, -1) == 0 { + atomic.StoreInt32(&downloader.left, 1) + } + rc.removeConsumed(downloader) } func (rc *ReaderCache) UnCache(fileId string) { @@ -202,15 +290,17 @@ func (rc *ReaderCache) remove(downloader *SingleChunkCacher) { } } -// removeConsumed drops a cacher once its buffer was fully read and no -// readers remain attached. The checks run under the ReaderCache lock so a -// reader attaching at the same time either wins (the cacher stays and that -// reader's detach retries the removal) or misses the map and refetches. +// removeConsumed drops a cacher once its buffer was fully read, or the +// streams positioned in it have left, and no readers remain attached or +// pinned. The checks run under the ReaderCache lock so a reader attaching +// at the same time either wins (the cacher stays and that reader's detach +// retries the removal) or misses the map and refetches. func (rc *ReaderCache) removeConsumed(downloader *SingleChunkCacher) { rc.Lock() removed := rc.downloaders[downloader.chunkFileId] == downloader && atomic.LoadInt32(&downloader.readers) == 0 && - atomic.LoadInt32(&downloader.consumed) != 0 + atomic.LoadInt32(&downloader.pins) == 0 && + (atomic.LoadInt32(&downloader.consumed) != 0 || atomic.LoadInt32(&downloader.left) != 0) if removed { delete(rc.downloaders, downloader.chunkFileId) }