diff --git a/weed/filer/reader_at.go b/weed/filer/reader_at.go index 9ac7ecad2..dfc57e88d 100644 --- a/weed/filer/reader_at.go +++ b/weed/filer/reader_at.go @@ -169,10 +169,17 @@ func (c *ChunkReadAt) Size() int64 { } func (c *ChunkReadAt) Close() error { + c.ReleaseStream() c.readerCache.destroy() return nil } +// ReleaseStream drops this reader's hold on the chunk it is positioned in. +// Unlike Close it leaves the (possibly shared) ReaderCache intact. +func (c *ChunkReadAt) ReleaseStream() { + c.readerCache.releaseStream(&c.stream) +} + func (c *ChunkReadAt) ReadAt(p []byte, offset int64) (n int, err error) { c.readerPattern.MonitorReadAt(offset, len(p)) diff --git a/weed/filer/reader_at_shared_test.go b/weed/filer/reader_at_shared_test.go index 22937a3b8..86e113d56 100644 --- a/weed/filer/reader_at_shared_test.go +++ b/weed/filer/reader_at_shared_test.go @@ -207,3 +207,77 @@ func TestChunkStreamConcurrentReadsOnOneReader(t *testing.T) { 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") + } +} diff --git a/weed/filer/reader_cache.go b/weed/filer/reader_cache.go index 93b7883e2..2cd046ea1 100644 --- a/weed/filer/reader_cache.go +++ b/weed/filer/reader_cache.go @@ -169,14 +169,26 @@ retry: } } - // clean up old downloaders + // clean up old downloaders; prefer one no stream is positioned in, but + // fall back to a pinned one so abandoned pins cannot bypass the limit if len(rc.downloaders) >= rc.limit { oldestFid, oldestTime := "", time.Now().UnixNano() + pinnedFid, pinnedTime := "", int64(0) for fid, downloader := range rc.downloaders { completedTime := atomic.LoadInt64(&downloader.completedTimeNew) - if completedTime > 0 && completedTime < oldestTime { - oldestFid, oldestTime = fid, completedTime + if completedTime <= 0 { + continue } + if atomic.LoadInt32(&downloader.pins) == 0 { + if completedTime < oldestTime { + oldestFid, oldestTime = fid, completedTime + } + } else if pinnedFid == "" || completedTime < pinnedTime { + pinnedFid, pinnedTime = fid, completedTime + } + } + if oldestFid == "" { + oldestFid = pinnedFid } if oldestFid != "" { oldDownloader := rc.downloaders[oldestFid] @@ -290,6 +302,23 @@ func (rc *ReaderCache) remove(downloader *SingleChunkCacher) { } } +// removeUnpinned drops a cacher only while no stream is positioned in it. +// Budget eviction picks its victim under the budget lock, so the pin check +// and the map removal must happen together under the ReaderCache lock. +func (rc *ReaderCache) removeUnpinned(downloader *SingleChunkCacher) (removed bool) { + rc.Lock() + removed = rc.downloaders[downloader.chunkFileId] == downloader && + atomic.LoadInt32(&downloader.pins) == 0 + if removed { + delete(rc.downloaders, downloader.chunkFileId) + } + rc.Unlock() + if removed { + downloader.destroy() + } + return +} + // 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 diff --git a/weed/filer/reader_cache_budget.go b/weed/filer/reader_cache_budget.go index c932235b0..9e2e91f3f 100644 --- a/weed/filer/reader_cache_budget.go +++ b/weed/filer/reader_cache_budget.go @@ -4,6 +4,7 @@ import ( "container/list" "fmt" "sync" + "sync/atomic" "github.com/seaweedfs/seaweedfs/weed/util/mem" ) @@ -51,12 +52,39 @@ func (b *ReaderCacheBudget) reserve(s *SingleChunkCacher) error { b.Unlock() return nil } - if entry := b.idle.Front(); entry != nil { - victim := entry.Value.(*SingleChunkCacher) + // Prefer evicting an idle chunk no stream is positioned in; fall back + // to the oldest pinned one so abandoned pins cannot block the budget. + var victim *SingleChunkCacher + var entry *list.Element + pinnedVictim := false + for e := b.idle.Front(); e != nil; e = e.Next() { + c := e.Value.(*SingleChunkCacher) + if atomic.LoadInt32(&c.pins) == 0 { + victim, entry = c, e + pinnedVictim = false + break + } + if victim == nil { + victim, entry = c, e + pinnedVictim = true + } + } + if entry != nil { b.idle.Remove(entry) delete(b.idleEntries, victim) b.Unlock() - victim.parent.remove(victim) + if pinnedVictim { + victim.parent.remove(victim) + } else if !victim.parent.removeUnpinned(victim) { + // The victim was pinned between selection and removal: keep it + // evictable so a pin abandoned in that gap cannot wedge the + // budget, then retry the selection. + b.Lock() + if _, ok := b.reservations[victim]; ok && b.idleEntries[victim] == nil { + b.idleEntries[victim] = b.idle.PushBack(victim) + } + b.Unlock() + } continue } changed := b.changed diff --git a/weed/s3api/s3api_object_handlers.go b/weed/s3api/s3api_object_handlers.go index acc1ddcb2..36b19d166 100644 --- a/weed/s3api/s3api_object_handlers.go +++ b/weed/s3api/s3api_object_handlers.go @@ -1103,6 +1103,7 @@ func (s3a *S3ApiServer) streamFromVolumeServers(w http.ResponseWriter, r *http.R tStreamPrep := time.Now() chunkViews := filer.ViewFromVisibleIntervals(visibleIntervals, offset, size) reader := filer.NewChunkReaderAtFromClient(ctx, s3a.readerCache, chunkViews, totalSize, filer.DefaultPrefetchCount) + defer reader.ReleaseStream() streamPrepTime = time.Since(tStreamPrep) // A cached chunk whose volume server is down, or whose needle was evicted and diff --git a/weed/server/webdav_server.go b/weed/server/webdav_server.go index f2939243a..1289a1321 100644 --- a/weed/server/webdav_server.go +++ b/weed/server/webdav_server.go @@ -130,7 +130,7 @@ type WebDavFile struct { off int64 entry *filer_pb.Entry visibleIntervals *filer.IntervalList[*filer.VisibleInterval] - reader io.ReaderAt + reader *filer.ChunkReadAt bufWriter *buffered_writer.BufferedWriteCloser ctx context.Context } @@ -537,6 +537,9 @@ func (f *WebDavFile) Write(buf []byte) (int, error) { func (f *WebDavFile) Close() error { glog.V(2).Infof("WebDavFileSystem.Close %v", f.name) + if f.reader != nil { + f.reader.ReleaseStream() + } if f.bufWriter == nil { return nil }