Files
seaweedfs/weed/filer/reader_at_shared_test.go
T
Chris LuGitHubDevin <158243242+devin-ai-integration[bot]@users.noreply.github.com>Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
9b3b12c607 filer: pin-aware reader cache eviction and stream release (#11503)
* filer: synchronize stream pins and release them on transitions

Guard chunkStream.cacher with the ReaderCache lock everywhere: mount
sections share one ChunkReadAt across concurrent reads, and unsynchronized
release could double-unpin. Reads served from the chunk cache now detach
the stream's pin instead of retaining the previous chunk. Eviction prefers
unpinned downloaders so a pinned buffer is not dropped mid-stream. A new
ReleaseStream lets callers drop their pin without destroying the shared
cache; S3 and WebDAV readers use it. lastChunkFid becomes atomic since
concurrent mount reads can update it.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: keep eviction bounded when every downloader is pinned

Both eviction paths still fall back to a pinned victim when no unpinned
one exists, so abandoned stream pins cannot bypass the downloader limit
or stall the memory budget. Budget eviction also rechecks the pin under
the ReaderCache lock at removal time: a stream that pinned the selected
victim in between keeps it mapped and the selection retries.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: restore budget bookkeeping when a victim gets pinned mid-eviction

removeUnpinned losing the pin race left the victim out of the idle list
while still holding its reservation, making it unevictable even as the
pinned fallback. Push it back when the reservation is still live.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-09-28 22:37:04 +08:00

284 lines
9.7 KiB
Go

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))
}
}
// A chunk a stream is positioned in must outlast downloader-limit eviction:
// otherwise a busy cache drops the buffer mid-stream and forces a refetch.
func TestChunkReadAtPinnedChunkSurvivesEviction(t *testing.T) {
const chunkSize = 64 << 10
var fetches int32
rc := NewReaderCache(2, newMockChunkCacheForReaderCache(), 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) {
if fileId == "chunk0" {
atomic.AddInt32(&fetches, 1)
}
return len(buffer), nil
}
stream := &chunkStream{}
buf := make([]byte, 16<<10)
// One slice in: the stream is positioned in chunk0 but has not left it.
if _, err := rc.readChunkAt(context.Background(), stream, buf, "chunk0", nil, false, 0, chunkSize, false); err != nil {
t.Fatal(err)
}
// Fill the downloader map past its limit with unpinned chunks.
for _, fileId := range []string{"chunk1", "chunk2", "chunk3"} {
if _, err := rc.ReadChunkAt(context.Background(), buf, fileId, nil, false, 0, chunkSize, true); err != nil {
t.Fatalf("read %s: %v", fileId, err)
}
}
// The stream's next slice must come from the still-pinned buffer.
if _, err := rc.readChunkAt(context.Background(), stream, buf, "chunk0", nil, false, 16<<10, chunkSize, false); err != nil {
t.Fatal(err)
}
if got := atomic.LoadInt32(&fetches); got != 1 {
t.Errorf("chunk0 fetched %d times, want 1", got)
}
}
// When every downloader is pinned the limit still applies: the oldest pinned
// buffer is evicted so abandoned streams cannot grow memory past the limit.
func TestChunkReadAtPinnedEvictionFallsBackWhenAllPinned(t *testing.T) {
const chunkSize = 64 << 10
rc := NewReaderCache(2, newMockChunkCacheForReaderCache(), 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, _ string, _ util_http.RefreshUrlsFunc) (int, error) {
return len(buffer), nil
}
buf := make([]byte, 16<<10)
streamA := &chunkStream{}
streamB := &chunkStream{}
if _, err := rc.readChunkAt(context.Background(), streamA, buf, "chunk0", nil, false, 0, chunkSize, false); err != nil {
t.Fatal(err)
}
if _, err := rc.readChunkAt(context.Background(), streamB, buf, "chunk1", nil, false, 0, chunkSize, false); err != nil {
t.Fatal(err)
}
// Every downloader is now pinned; the next chunk must still get in.
if _, err := rc.readChunkAt(context.Background(), &chunkStream{}, buf, "chunk2", nil, false, 0, chunkSize, false); err != nil {
t.Fatal(err)
}
if isRetained(rc, "chunk0") {
t.Fatal("oldest pinned downloader was not evicted past the limit")
}
if !isRetained(rc, "chunk2") {
t.Fatal("new downloader missing after pinned fallback eviction")
}
}