mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-11 09:05:50 +00:00
Compare commits
8
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
89e7753ed5 | ||
|
|
495e24a476 | ||
|
|
1f15e7865a | ||
|
|
521a19c567 | ||
|
|
dd75add62f | ||
|
|
58aa9acfb0 | ||
|
|
8227ae3f49 | ||
|
|
3fd5e4c329 |
@@ -761,7 +761,7 @@ impl Store {
|
||||
Err(e) => writes
|
||||
.iter()
|
||||
.map(|_| match e {
|
||||
VolumeError::ReadOnly => Err(VolumeError::ReadOnly),
|
||||
VolumeError::ReadOnly(vid) => Err(VolumeError::ReadOnly(vid)),
|
||||
_ => Err(VolumeError::NotFound),
|
||||
})
|
||||
.collect(),
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
+19
-3
@@ -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)
|
||||
}
|
||||
|
||||
@@ -208,6 +208,200 @@ 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
|
||||
|
||||
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) {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user