diff --git a/weed/filer/reader_cache.go b/weed/filer/reader_cache.go index d4621eb0e..1c2ba7252 100644 --- a/weed/filer/reader_cache.go +++ b/weed/filer/reader_cache.go @@ -25,7 +25,7 @@ type ReaderCache struct { lookupFileIdFn wdclient.LookupFileIdFunctionType cacheInvalidator CacheInvalidator fetchChunkDataFn fetchChunkDataFnType - sync.Mutex + sync.RWMutex downloaders map[string]*SingleChunkCacher limit int budget *ReaderCacheBudget @@ -78,6 +78,13 @@ func (rc *ReaderCache) MaybeCache(chunkViews *Interval[*ChunkView], count int) { count = 1 } + var started []*SingleChunkCacher + defer func() { + for _, cacher := range started { + <-cacher.cacheStartedCh + } + }() + rc.Lock() defer rc.Unlock() @@ -113,9 +120,10 @@ func (rc *ReaderCache) MaybeCache(chunkViews *Interval[*ChunkView], count int) { // cache this chunk if not yet shouldCache := (uint64(chunkView.ViewOffset) + chunkView.ChunkSize) <= rc.chunkCache.GetMaxFilePartSizeInCache() cacher := newSingleChunkCacher(rc, chunkView.FileId, chunkView.CipherKey, chunkView.IsGzipped, int(chunkView.ChunkSize), shouldCache) - go cacher.startCaching() - <-cacher.cacheStartedCh + cacher.wg.Add(1) // the fetch, so destroy() waits even before the goroutine runs rc.downloaders[chunkView.FileId] = cacher + go cacher.startCaching() + started = append(started, cacher) cached++ } @@ -141,6 +149,7 @@ func (rc *ReaderCache) fetchChunkRange(ctx context.Context, buffer []byte, chunk // moves elsewhere, so another stream finishing or leaving the same chunk does // not drop the buffer from under it. type chunkStream struct { + mu sync.Mutex // guards cacher; concurrent ReadAt calls on one ChunkReadAt share the stream cacher *SingleChunkCacher } @@ -150,47 +159,56 @@ func (rc *ReaderCache) ReadChunkAt(ctx context.Context, buffer []byte, fileId st func (rc *ReaderCache) readChunkAt(ctx context.Context, stream *chunkStream, buffer []byte, fileId string, cipherKey []byte, isGzipped bool, offset int64, chunkSize int, shouldCache bool) (int, error) { retry: - rc.Lock() - - for { - if cacher, found := rc.downloaders[fileId]; found { - if cacher.hasCompletedError() { - delete(rc.downloaders, fileId) - rc.Unlock() - cacher.destroy() - rc.Lock() - continue - } - // Count this read on the cacher before releasing the map lock, so a - // concurrent destroy() (error eviction here, LRU, or UnCache) cannot - // start wg.Wait() on a zero counter while this read is about to register. - cacher.wg.Add(1) - atomic.AddInt32(&cacher.readers, 1) - previous := stream.pin(cacher) - rc.Unlock() - rc.unpin(previous) - n, err := cacher.readChunkAt(ctx, buffer, offset) - rc.releaseIfFinished(stream, cacher, offset, n, err, chunkSize) - if n > 0 || err != nil { - return n, err - } - // If n=0 and err=nil, the cacher couldn't provide data for this offset. - // Fall through to try chunkCache. - rc.Lock() + rc.RLock() + if cacher, found := rc.downloaders[fileId]; found { + if cacher.hasCompletedError() { + rc.RUnlock() + rc.remove(cacher) + goto retry } - break + // Count this read on the cacher before releasing the map lock, so a + // concurrent destroy() (error eviction here, LRU, or UnCache) cannot + // start wg.Wait() on a zero counter while this read is about to register. + cacher.wg.Add(1) + atomic.AddInt32(&cacher.readers, 1) + previous := stream.pin(cacher) + rc.RUnlock() + n, err := rc.readFromCacher(ctx, stream, cacher, previous, buffer, offset, chunkSize) + if n > 0 || err != nil { + return n, err + } + // If n=0 and err=nil, the cacher couldn't provide data for this offset. + // Fall through to try chunkCache. + } else { + rc.RUnlock() } + if shouldCache || rc.lookupFileIdFn == nil { n, err := rc.chunkCache.ReadChunkAt(buffer, fileId, uint64(offset)) if n > 0 { // Served from the chunk cache: the stream has left its pinned chunk. - previous := stream.unpinLocked() - rc.Unlock() - rc.unpin(previous) + rc.unpin(stream.detach(nil)) return n, err } } + rc.Lock() + // A downloader may have registered between the read lock and this one. + if cacher, found := rc.downloaders[fileId]; found { + if cacher.hasCompletedError() { + delete(rc.downloaders, fileId) + rc.Unlock() + cacher.destroy() + goto retry + } + cacher.wg.Add(1) + atomic.AddInt32(&cacher.readers, 1) + previous := stream.pin(cacher) + rc.Unlock() + n, err := rc.readFromCacher(ctx, stream, cacher, previous, buffer, offset, chunkSize) + return n, err + } + // 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 { @@ -224,40 +242,57 @@ retry: // glog.V(4).Infof("cache1 %s", fileId) cacher := newSingleChunkCacher(rc, fileId, cipherKey, isGzipped, chunkSize, shouldCache) - go cacher.startCaching() - <-cacher.cacheStartedCh - rc.downloaders[fileId] = cacher - cacher.wg.Add(1) + cacher.wg.Add(2) // the fetch plus this read, so destroy() waits even before the goroutine runs atomic.AddInt32(&cacher.readers, 1) + rc.downloaders[fileId] = cacher previous := stream.pin(cacher) rc.Unlock() - rc.unpin(previous) - - n, err := cacher.readChunkAt(ctx, buffer, offset) - rc.releaseIfFinished(stream, cacher, offset, n, err, chunkSize) + go cacher.startCaching() + <-cacher.cacheStartedCh + n, err := rc.readFromCacher(ctx, stream, cacher, previous, buffer, offset, chunkSize) return n, err } +// readFromCacher serves the read and releases the pin if the stream finished +// the chunk. The caller must have registered the read (wg + readers) and +// pinned the stream under the map lock; previous is the chunk the stream left +// when it pinned. +func (rc *ReaderCache) readFromCacher(ctx context.Context, stream *chunkStream, cacher *SingleChunkCacher, previous *SingleChunkCacher, buffer []byte, offset int64, chunkSize int) (n int, err error) { + rc.unpin(previous) + n, err = cacher.readChunkAt(ctx, buffer, offset) + rc.releaseIfFinished(stream, cacher, offset, n, err, chunkSize) + return +} + // pin makes cacher the stream's current chunk and returns the chunk it was -// pinned to before, which the caller unpins once the ReaderCache lock is -// released. The stream is only touched under the ReaderCache lock, since -// concurrent ReadAt calls on one ChunkReadAt share it. +// pinned to before, which the caller unpins. func (stream *chunkStream) pin(cacher *SingleChunkCacher) (previous *SingleChunkCacher) { - if stream == nil || stream.cacher == cacher { + if stream == nil { return nil } + stream.mu.Lock() + defer stream.mu.Unlock() + if stream.cacher == cacher { + return nil + } + atomic.AddInt32(&cacher.pins, 1) previous = stream.cacher stream.cacher = cacher - atomic.AddInt32(&cacher.pins, 1) return previous } -// unpinLocked detaches the stream from its chunk and returns that chunk for -// the caller to unpin once the ReaderCache lock is released. -func (stream *chunkStream) unpinLocked() (previous *SingleChunkCacher) { +// detach unpins the stream from its chunk and returns that chunk for the +// caller to unpin. A non-nil expected cacher makes the detach conditional on +// the stream still being positioned in it. +func (stream *chunkStream) detach(expected *SingleChunkCacher) (previous *SingleChunkCacher) { if stream == nil { return nil } + stream.mu.Lock() + defer stream.mu.Unlock() + if expected != nil && stream.cacher != expected { + return nil + } previous = stream.cacher stream.cacher = nil return previous @@ -269,24 +304,12 @@ func (rc *ReaderCache) releaseIfFinished(stream *chunkStream, cacher *SingleChun if stream == nil || err != nil || offset+int64(n) < int64(chunkSize) { return } - var previous *SingleChunkCacher - rc.Lock() - if stream.cacher == cacher { - previous = stream.unpinLocked() - } - rc.Unlock() - rc.unpin(previous) + rc.unpin(stream.detach(cacher)) } // releaseStream unpins whatever chunk the stream is positioned in. func (rc *ReaderCache) releaseStream(stream *chunkStream) { - if stream == nil { - return - } - rc.Lock() - previous := stream.unpinLocked() - rc.Unlock() - rc.unpin(previous) + rc.unpin(stream.detach(nil)) } // unpin drops one stream's pin. Once no stream is positioned in the chunk it @@ -341,26 +364,41 @@ func (rc *ReaderCache) removeUnpinned(downloader *SingleChunkCacher) (removed bo return } +// consumable reports whether the cacher could be dropped: its buffer was +// fully read or every stream positioned in it has left, and no reader is +// attached or pinned. All fields are atomics, so callers can skip the +// ReaderCache lock whenever removal is impossible. +func (s *SingleChunkCacher) consumable() bool { + return atomic.LoadInt32(&s.readers) == 0 && + atomic.LoadInt32(&s.pins) == 0 && + (atomic.LoadInt32(&s.consumed) != 0 || atomic.LoadInt32(&s.left) != 0) +} + // 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 // at the same time either wins (the cacher stays and that reader's detach // retries the removal) or misses the map and refetches. func (rc *ReaderCache) removeConsumed(downloader *SingleChunkCacher) { - rc.Lock() - removed := rc.downloaders[downloader.chunkFileId] == downloader && - atomic.LoadInt32(&downloader.readers) == 0 && - atomic.LoadInt32(&downloader.pins) == 0 && - (atomic.LoadInt32(&downloader.consumed) != 0 || atomic.LoadInt32(&downloader.left) != 0) - if removed { - delete(rc.downloaders, downloader.chunkFileId) + if !downloader.consumable() { + return } + rc.Lock() + removed := rc.removeConsumedLocked(downloader) rc.Unlock() if removed { downloader.destroy() } } +func (rc *ReaderCache) removeConsumedLocked(downloader *SingleChunkCacher) (removed bool) { + if rc.downloaders[downloader.chunkFileId] == downloader && downloader.consumable() { + delete(rc.downloaders, downloader.chunkFileId) + return true + } + return false +} + func (rc *ReaderCache) destroy() { rc.Lock() downloaders := rc.downloaders @@ -385,8 +423,9 @@ func newSingleChunkCacher(parent *ReaderCache, fileId string, cipherKey []byte, } // startCaching downloads a chunk shared by concurrent readers. +// The caller must s.wg.Add(1) before publishing the cacher in the +// downloaders map, so a removal can never observe the fetch as absent. func (s *SingleChunkCacher) startCaching() { - s.wg.Add(1) defer func() { close(s.done) s.wg.Done() diff --git a/weed/filer/reader_cache_test.go b/weed/filer/reader_cache_test.go index 672f3d4df..16c3531f5 100644 --- a/weed/filer/reader_cache_test.go +++ b/weed/filer/reader_cache_test.go @@ -724,3 +724,34 @@ func TestSingleChunkCacherOneReaderCancelsOthersContinue(t *testing.T) { t.Error("Other reader did not complete") } } + +// TestReaderCacheBookkeepingOffCacheLock guards the contention fix: pin +// lifecycle calls that cannot remove a cacher must not block on the +// ReaderCache lock. +func TestReaderCacheBookkeepingOffCacheLock(t *testing.T) { + cache := newMockChunkCacheForReaderCache() + rc := NewReaderCache(10, cache, nil, nil) + defer rc.destroy() + + cacher := &SingleChunkCacher{parent: rc, chunkFileId: "pinned-chunk"} + atomic.StoreInt32(&cacher.readers, 1) // a read in flight: not consumable + + rc.Lock() + started := make(chan struct{}) + done := make(chan struct{}) + go func() { + close(started) + defer close(done) + var stream chunkStream + stream.pin(cacher) + rc.releaseStream(&stream) + }() + <-started + select { + case <-done: + case <-time.After(time.Second): + rc.Unlock() + t.Fatal("stream pin lifecycle blocked on the cache lock") + } + rc.Unlock() +}