diff --git a/weed/filer/filechunk_manifest.go b/weed/filer/filechunk_manifest.go index 015bbe790..f800333c9 100644 --- a/weed/filer/filechunk_manifest.go +++ b/weed/filer/filechunk_manifest.go @@ -380,15 +380,6 @@ func fetchWholeChunk(ctx context.Context, bytesBuffer *bytes.Buffer, lookupFileI }) } -func fetchChunkRange(ctx context.Context, buffer []byte, lookupFileIdFn wdclient.LookupFileIdFunctionType, fileId string, cipherKey []byte, isGzipped bool, offset int64, refreshUrls util_http.RefreshUrlsFunc) (int, error) { - urlStrings, err := lookupFileIdFn(ctx, fileId) - if err != nil { - glog.ErrorfCtx(ctx, "operation LookupFileId %s failed, err: %v", fileId, err) - return 0, err - } - return util_http.RetriedFetchChunkData(ctx, buffer, urlStrings, cipherKey, isGzipped, false, offset, fileId, refreshUrls) -} - // retriedStreamFetchChunkData streams a chunk from the first location that // answers. refreshUrls may be nil; when a location failed and a later one // answered, it is called so the reads that follow start from a fresh list. diff --git a/weed/filer/filechunks.go b/weed/filer/filechunks.go index 8cd8bb843..954547296 100644 --- a/weed/filer/filechunks.go +++ b/weed/filer/filechunks.go @@ -185,6 +185,15 @@ func (cv *ChunkView) IsFullChunk() bool { return cv.OffsetInChunk == 0 && cv.ViewSize == cv.ChunkSize } +// CanRangeFetch reports whether fetching just the view's byte range avoids +// reading more than the view needs. Ciphered and compressed chunks are +// stored and served whole — a range on either still costs a full read plus +// decrypt or decompress on the volume server — so partial views of them +// take the shared whole-chunk path instead. +func (cv *ChunkView) CanRangeFetch() bool { + return cv.CipherKey == nil && !cv.IsGzipped +} + func ViewFromChunks(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, chunks []*filer_pb.FileChunk, offset int64, size int64) (chunkViews *IntervalList[*ChunkView]) { visibles, _ := NonOverlappingVisibleIntervals(ctx, lookupFileIdFn, chunks, offset, offset+size) diff --git a/weed/filer/reader_at.go b/weed/filer/reader_at.go index dfc57e88d..c8c1f78de 100644 --- a/weed/filer/reader_at.go +++ b/weed/filer/reader_at.go @@ -349,14 +349,23 @@ 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() { + // A view clipped to part of its chunk (e.g. the edge of a ranged GET, + // whose views ViewFromVisibleIntervals clips to the request) only ever + // needs that part: fetch it as a range no matter the detected pattern. + // Fetching the chunk whole would multiply volume-server reads. Ciphered + // and compressed chunks are the exception: the volume server reads the + // whole blob to serve a range, so they take the shared whole-chunk path + // where one download serves every buffer — unless the whole chunk cannot + // even fit the reader budget, in which case a range fetch is the only + // way to serve the request. + rangeFetch := chunkView.CanRangeFetch() || !c.readerCache.budget.canFit(int(chunkView.ChunkSize)) + if rangeFetch && (!chunkView.IsFullChunk() || c.readerPattern.IsRandomMode()) { c.readerCache.releaseStream(&c.stream) n, err := c.readerCache.chunkCache.ReadChunkAt(buffer, chunkView.FileId, offset) if n > 0 { return n, err } - return fetchChunkRange(ctx, buffer, c.readerCache.lookupFileIdFn, chunkView.FileId, chunkView.CipherKey, chunkView.IsGzipped, int64(offset), - refreshUrls(ctx, c.readerCache.cacheInvalidator, c.readerCache.lookupFileIdFn, chunkView.FileId)) + return c.readerCache.fetchChunkRange(ctx, buffer, chunkView, int64(offset)) } shouldCache := (uint64(chunkView.ViewOffset) + chunkView.ChunkSize) <= c.readerCache.chunkCache.GetMaxFilePartSizeInCache() @@ -380,6 +389,13 @@ func (c *ChunkReadAt) readChunkSliceAt(ctx context.Context, buffer []byte, chunk // readChunkSliceAtForParallel is a simplified version for parallel chunk fetching // It doesn't update lastChunkFid or trigger prefetch (handled by the caller) func (c *ChunkReadAt) readChunkSliceAtForParallel(ctx context.Context, buffer []byte, chunkView *ChunkView, offset uint64) (n int, err error) { + if (chunkView.CanRangeFetch() || !c.readerCache.budget.canFit(int(chunkView.ChunkSize))) && !chunkView.IsFullChunk() { + n, err = c.readerCache.chunkCache.ReadChunkAt(buffer, chunkView.FileId, offset) + if n > 0 { + return n, err + } + return c.readerCache.fetchChunkRange(ctx, buffer, chunkView, int64(offset)) + } shouldCache := (uint64(chunkView.ViewOffset) + chunkView.ChunkSize) <= c.readerCache.chunkCache.GetMaxFilePartSizeInCache() return c.readerCache.ReadChunkAt(ctx, buffer, chunkView.FileId, chunkView.CipherKey, chunkView.IsGzipped, int64(offset), int(chunkView.ChunkSize), shouldCache) } diff --git a/weed/filer/reader_at_shared_test.go b/weed/filer/reader_at_shared_test.go index 86e113d56..a29316390 100644 --- a/weed/filer/reader_at_shared_test.go +++ b/weed/filer/reader_at_shared_test.go @@ -208,6 +208,201 @@ func TestChunkStreamConcurrentReadsOnOneReader(t *testing.T) { } } +type recordedFetch struct { + fileId string + isFullChunk bool + offset int64 + size int +} + +// fetchRecorder stubs the volume fetch and records how each chunk was +// requested: isFullChunk=false is a range fetch of just the view's slice, +// isFullChunk=true is a whole-chunk download into the shared cache. +func fetchRecorder(rc *ReaderCache) (fetches *[]recordedFetch) { + var mu sync.Mutex + recorded := &[]recordedFetch{} + rc.fetchChunkDataFn = func(_ context.Context, buffer []byte, _ []string, _ []byte, _ bool, isFullChunk bool, offset int64, fileId string, _ util_http.RefreshUrlsFunc) (int, error) { + mu.Lock() + *recorded = append(*recorded, recordedFetch{fileId, isFullChunk, offset, len(buffer)}) + mu.Unlock() + for i := range buffer { + buffer[i] = fileId[len(fileId)-1] + } + return len(buffer), nil + } + return recorded +} + +// A reader whose views are clipped to a request window — how the S3 gateway +// builds a ranged GET — must fetch only the covered part of each chunk: +// clipped edge views take range fetches, a fully covered chunk keeps the +// shared whole-chunk path. This is what keeps a ranged GET larger than a +// buffer from multiplying volume-server reads (issue #11564), without giving +// up whole-chunk caching where the whole chunk is actually wanted. +func TestChunkReadAtClippedViewsFetchOnlyCoveredParts(t *testing.T) { + const chunkSize = 64 << 10 + + rc := NewReaderCache(64, (*chunk_cache.TieredChunkCache)(nil), func(context.Context, string) ([]string, error) { + return []string{"unused"}, nil + }, nil) + defer rc.destroy() + fetches := fetchRecorder(rc) + + // Window [56KiB, 152KiB): tail of chunk0, all of chunk1, head of + // chunk2, head of ciphered chunk3, head of compressed chunk4 (file + // chunks need not be aligned). + views := NewIntervalList[*ChunkView]() + views.AppendInterval(&Interval[*ChunkView]{ + StartOffset: chunkSize - 8<<10, + StopOffset: chunkSize, + Value: &ChunkView{FileId: "chunk0", OffsetInChunk: chunkSize - 8<<10, ViewSize: 8 << 10, ViewOffset: chunkSize - 8<<10, ChunkSize: chunkSize}, + }) + views.AppendInterval(&Interval[*ChunkView]{ + StartOffset: chunkSize, + StopOffset: 2 * chunkSize, + Value: &ChunkView{FileId: "chunk1", ViewSize: chunkSize, ViewOffset: chunkSize, ChunkSize: chunkSize}, + }) + views.AppendInterval(&Interval[*ChunkView]{ + StartOffset: 2 * chunkSize, + StopOffset: 2*chunkSize + 8<<10, + Value: &ChunkView{FileId: "chunk2", ViewSize: 8 << 10, ViewOffset: 2 * chunkSize, ChunkSize: chunkSize}, + }) + views.AppendInterval(&Interval[*ChunkView]{ + StartOffset: 2*chunkSize + 8<<10, + StopOffset: 2*chunkSize + 16<<10, + Value: &ChunkView{FileId: "chunk3", ViewSize: 8 << 10, ViewOffset: 2*chunkSize + 8<<10, ChunkSize: chunkSize, CipherKey: []byte("key")}, + }) + views.AppendInterval(&Interval[*ChunkView]{ + StartOffset: 2*chunkSize + 16<<10, + StopOffset: 2*chunkSize + 24<<10, + Value: &ChunkView{FileId: "chunk4", ViewSize: 8 << 10, ViewOffset: 2*chunkSize + 16<<10, ChunkSize: chunkSize, IsGzipped: true}, + }) + + reader := NewChunkReaderAtFromClient(context.Background(), rc, views, 4*chunkSize, 0) + buf := make([]byte, chunkSize+32<<10) + if n, err := reader.ReadAt(buf, chunkSize-8<<10); err != nil || n != len(buf) { + t.Fatalf("window read: n=%d err=%v", n, err) + } + // buf holds [56KiB, 152KiB): chunk0's tail, chunk1, and the heads of + // chunk2, chunk3 and chunk4. + for i, b := range buf { + want := byte('1') + if i < 8<<10 { + want = '0' + } else if i >= 24<<10+chunkSize { + want = '4' + } else if i >= 16<<10+chunkSize { + want = '3' + } else if i >= 8<<10+chunkSize { + want = '2' + } + if b != want { + t.Fatalf("buf[%d]=%q, want %q", i, b, want) + } + } + + want := []recordedFetch{ + {fileId: "chunk0", isFullChunk: false, offset: chunkSize - 8<<10, size: 8 << 10}, + {fileId: "chunk1", isFullChunk: true, offset: 0, size: chunkSize}, + {fileId: "chunk2", isFullChunk: false, offset: 0, size: 8 << 10}, + // partial views, but ciphered and compressed chunks download whole + // either way and the shared path decrypts/decompresses once for + // every buffer + {fileId: "chunk3", isFullChunk: true, offset: 0, size: chunkSize}, + {fileId: "chunk4", isFullChunk: true, offset: 0, size: chunkSize}, + } + got := map[string]recordedFetch{} + for _, f := range *fetches { + if _, dup := got[f.fileId]; dup { + t.Fatalf("chunk %s fetched more than once: %+v", f.fileId, *fetches) + } + got[f.fileId] = f + } + for _, w := range want { + if g, ok := got[w.fileId]; !ok { + t.Fatalf("chunk %s never fetched: %+v", w.fileId, *fetches) + } else if g != w { + t.Fatalf("chunk %s fetched as %+v, want %+v", w.fileId, g, w) + } + } +} + +// The regression from issue #11564: a ranged GET sitting inside one big +// chunk. Every buffer of the request must stay a range fetch — none may +// escalate into a whole-chunk download once the reads look sequential. +func TestChunkReadAtRangeInsideOneChunkStaysRangeFetch(t *testing.T) { + const chunkSize = 1 << 20 + const sliceSize = 16 << 10 + + rc := NewReaderCache(64, (*chunk_cache.TieredChunkCache)(nil), func(context.Context, string) ([]string, error) { + return []string{"unused"}, nil + }, nil) + defer rc.destroy() + fetches := fetchRecorder(rc) + + // Range [32KiB, 96KiB) inside one 1MiB chunk: a single clipped view. + views := NewIntervalList[*ChunkView]() + views.AppendInterval(&Interval[*ChunkView]{ + StartOffset: 32 << 10, + StopOffset: 96 << 10, + Value: &ChunkView{FileId: "chunk0", OffsetInChunk: 32 << 10, ViewSize: 64 << 10, ViewOffset: 32 << 10, ChunkSize: chunkSize}, + }) + + reader := NewChunkReaderAtFromClient(context.Background(), rc, views, chunkSize, 0) + for offset := int64(32 << 10); offset < 96<<10; offset += sliceSize { + buf := make([]byte, sliceSize) + if n, err := reader.ReadAt(buf, offset); err != nil || n != sliceSize { + t.Fatalf("read at %d: n=%d err=%v", offset, n, err) + } + } + + if len(*fetches) != 4 { + t.Fatalf("got %d fetches, want 4 range fetches: %+v", len(*fetches), *fetches) + } + for i, f := range *fetches { + wantOffset := int64(32<<10) + int64(i)*sliceSize + if f.isFullChunk || f.offset != wantOffset || f.size != sliceSize { + t.Fatalf("fetch %d = %+v, want range fetch offset=%d size=%d", i, f, wantOffset, sliceSize) + } + } +} + +// A compressed chunk larger than the reader cache budget can never be +// downloaded whole — the budget rejects the buffer — so its partial view +// must fall back to a range fetch even though each range costs a full +// decompress server-side. The alternative is a failed GET. +func TestChunkReadAtOversizedCompressedChunkFallsBackToRange(t *testing.T) { + const chunkSize = 1 << 20 + const sliceSize = 16 << 10 + + budget := NewReaderCacheBudget(64 << 10) // smaller than the chunk + rc := NewReaderCache(64, (*chunk_cache.TieredChunkCache)(nil), func(context.Context, string) ([]string, error) { + return []string{"unused"}, nil + }, nil, budget) + defer rc.destroy() + fetches := fetchRecorder(rc) + + views := NewIntervalList[*ChunkView]() + views.AppendInterval(&Interval[*ChunkView]{ + StartOffset: 32 << 10, + StopOffset: 64 << 10, + Value: &ChunkView{FileId: "chunk0", OffsetInChunk: 32 << 10, ViewSize: 32 << 10, ViewOffset: 32 << 10, ChunkSize: chunkSize, IsGzipped: true}, + }) + + reader := NewChunkReaderAtFromClient(context.Background(), rc, views, chunkSize, 0) + buf := make([]byte, 32<<10) + if n, err := reader.ReadAt(buf, 32<<10); err != nil || n != len(buf) { + t.Fatalf("read: n=%d err=%v", n, err) + } + + if len(*fetches) != 1 { + t.Fatalf("got %d fetches, want 1 range fetch: %+v", len(*fetches), *fetches) + } + if f := (*fetches)[0]; f.isFullChunk || f.offset != 32<<10 || f.size != 32<<10 { + t.Fatalf("fetch = %+v, want range fetch offset=%d size=%d", f, 32<<10, 32<<10) + } +} + // 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) { diff --git a/weed/filer/reader_cache.go b/weed/filer/reader_cache.go index 2cd046ea1..d4621eb0e 100644 --- a/weed/filer/reader_cache.go +++ b/weed/filer/reader_cache.go @@ -100,6 +100,14 @@ func (rc *ReaderCache) MaybeCache(chunkViews *Interval[*ChunkView], count int) { // abort when slots are filled return } + if (chunkView.CanRangeFetch() || !rc.budget.canFit(int(chunkView.ChunkSize))) && !chunkView.IsFullChunk() { + // the view is clipped to part of the chunk and will be + // range-fetched, so prefetching it whole would download bytes + // nobody needs; a ciphered or compressed partial view needs + // the whole blob anyway and is worth prefetching, but not when + // it cannot fit the budget at all + continue + } // glog.V(4).Infof("prefetch %s offset %d", chunkView.FileId, chunkView.ViewOffset) // cache this chunk if not yet @@ -114,6 +122,20 @@ func (rc *ReaderCache) MaybeCache(chunkViews *Interval[*ChunkView], count int) { return } +// fetchChunkRange downloads only [offset, offset+len(buffer)) of a chunk, +// for views clipped to part of their chunk and for random-mode reads. It +// goes through fetchChunkDataFn so tests observe range fetches the same way +// they observe whole-chunk downloads. +func (rc *ReaderCache) fetchChunkRange(ctx context.Context, buffer []byte, chunkView *ChunkView, offset int64) (int, error) { + urlStrings, err := rc.lookupFileIdFn(ctx, chunkView.FileId) + if err != nil { + glog.ErrorfCtx(ctx, "operation LookupFileId %s failed, err: %v", chunkView.FileId, err) + return 0, err + } + return rc.fetchChunkDataFn(ctx, buffer, urlStrings, chunkView.CipherKey, chunkView.IsGzipped, false, offset, chunkView.FileId, + refreshUrls(ctx, rc.cacheInvalidator, rc.lookupFileIdFn, chunkView.FileId)) +} + // 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 diff --git a/weed/filer/reader_cache_budget.go b/weed/filer/reader_cache_budget.go index 9e2e91f3f..b048e26d9 100644 --- a/weed/filer/reader_cache_budget.go +++ b/weed/filer/reader_cache_budget.go @@ -93,6 +93,13 @@ func (b *ReaderCacheBudget) reserve(s *SingleChunkCacher) error { } } +// canFit reports whether a whole-chunk buffer of this size can ever be +// reserved. A chunk bigger than the budget cannot be read through the +// whole-chunk path at all, so callers must fall back to range fetches. +func (b *ReaderCacheBudget) canFit(size int) bool { + return b == nil || int64(mem.AllocationSize(size)) <= b.limit +} + func (b *ReaderCacheBudget) complete(s *SingleChunkCacher) { if b == nil { return diff --git a/weed/filer/reader_cache_memory_test.go b/weed/filer/reader_cache_memory_test.go index a723ba3e2..d363a5b4d 100644 --- a/weed/filer/reader_cache_memory_test.go +++ b/weed/filer/reader_cache_memory_test.go @@ -72,7 +72,7 @@ func TestReaderCacheBudgetInFlight(t *testing.T) { return len(buffer), nil } if prefetch { - rc.MaybeCache(&Interval[*ChunkView]{Value: &ChunkView{FileId: "chunk", ChunkSize: 3 << 10}}, 1) + rc.MaybeCache(&Interval[*ChunkView]{Value: &ChunkView{FileId: "chunk", ViewSize: 3 << 10, ChunkSize: 3 << 10}}, 1) } else { readers.Add(1) go func() { @@ -180,7 +180,7 @@ func TestReaderCacheFailedPrefetchReleasesBudget(t *testing.T) { rc.fetchChunkDataFn = func(_ context.Context, _ []byte, _ []string, _ []byte, _ bool, _ bool, _ int64, _ string, _ util_http.RefreshUrlsFunc) (int, error) { return 0, fmt.Errorf("fetch failed") } - rc.MaybeCache(&Interval[*ChunkView]{Value: &ChunkView{FileId: "failed", ChunkSize: 1024}}, 1) + rc.MaybeCache(&Interval[*ChunkView]{Value: &ChunkView{FileId: "failed", ViewSize: 1024, ChunkSize: 1024}}, 1) deadline := time.Now().Add(5 * time.Second) for { rc.Lock() @@ -278,7 +278,7 @@ func TestReaderCachePrefetchBufferDroppedAfterRead(t *testing.T) { buffer[0] = 42 return len(buffer), nil } - rc.MaybeCache(&Interval[*ChunkView]{Value: &ChunkView{FileId: "chunk", ChunkSize: 4 << 10}}, 1) + rc.MaybeCache(&Interval[*ChunkView]{Value: &ChunkView{FileId: "chunk", ViewSize: 4 << 10, ChunkSize: 4 << 10}}, 1) buf := make([]byte, 4<<10) n, err := rc.ReadChunkAt(context.Background(), buf, "chunk", nil, false, 0, 4<<10, false) diff --git a/weed/filer/reader_pattern.go b/weed/filer/reader_pattern.go index 6a9a95f83..65184722f 100644 --- a/weed/filer/reader_pattern.go +++ b/weed/filer/reader_pattern.go @@ -54,10 +54,13 @@ func (rp *ReaderPattern) MonitorReadAt(offset int64, size int) { if counter < ModeChangeLimit { atomic.AddInt64(&rp.isSequentialCounter, 1) } + } else if counter <= 0 { + // Entering random mode is a strong verdict: drop to the bottom of + // the window so the contiguous tail of one ranged request cannot + // flip it back on the next buffer read and pay a whole-chunk fetch. + atomic.StoreInt64(&rp.isSequentialCounter, -ModeChangeLimit) } else { - if counter > -ModeChangeLimit { - atomic.AddInt64(&rp.isSequentialCounter, -1) - } + atomic.AddInt64(&rp.isSequentialCounter, -1) } } diff --git a/weed/filer/reader_pattern_test.go b/weed/filer/reader_pattern_test.go index 54b896ab4..147d70cfd 100644 --- a/weed/filer/reader_pattern_test.go +++ b/weed/filer/reader_pattern_test.go @@ -102,3 +102,23 @@ func TestReaderPatternRecoversFromRandom(t *testing.T) { t.Fatal("sustained near reads must recover sequential mode") } } + +// A ranged request's first read lands far from the frontier, but its +// remaining buffer reads are contiguous. The random verdict must stick for +// them — otherwise the tail of every range >256KiB pays a whole-chunk fetch. +func TestReaderPatternRangedReadStaysRandom(t *testing.T) { + rp := NewReaderPattern() + rp.MonitorReadAt(500*mb, 256*1024) // far first read -> -ModeChangeLimit + for i := 1; i <= 2; i++ { + rp.MonitorReadAt(500*mb+int64(i)*256*1024, 256*1024) + if !rp.IsRandomMode() { + t.Fatalf("contiguous read %d of a ranged request flipped back to sequential", i+1) + } + } + for i := 3; i < 10; i++ { + rp.MonitorReadAt(500*mb+int64(i)*256*1024, 256*1024) + } + if rp.IsRandomMode() { + t.Fatal("sustained sequential reads should restore sequential mode") + } +}