mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-30 19:55:48 +00:00
* 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>
284 lines
9.7 KiB
Go
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")
|
|
}
|
|
}
|