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>
518 lines
16 KiB
Go
518 lines
16 KiB
Go
package filer
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/glog"
|
|
"github.com/seaweedfs/seaweedfs/weed/util/chunk_cache"
|
|
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
|
|
"github.com/seaweedfs/seaweedfs/weed/util/mem"
|
|
"github.com/seaweedfs/seaweedfs/weed/wdclient"
|
|
)
|
|
|
|
type CacheInvalidator interface {
|
|
InvalidateCache(fileId string)
|
|
}
|
|
|
|
type fetchChunkDataFnType func(ctx context.Context, buffer []byte, urlStrings []string, cipherKey []byte, isGzipped bool, isFullChunk bool, offset int64, fileId string, refreshUrls util_http.RefreshUrlsFunc) (n int, err error)
|
|
|
|
type ReaderCache struct {
|
|
chunkCache chunk_cache.ChunkCache
|
|
lookupFileIdFn wdclient.LookupFileIdFunctionType
|
|
cacheInvalidator CacheInvalidator
|
|
fetchChunkDataFn fetchChunkDataFnType
|
|
sync.Mutex
|
|
downloaders map[string]*SingleChunkCacher
|
|
limit int
|
|
budget *ReaderCacheBudget
|
|
}
|
|
|
|
type SingleChunkCacher struct {
|
|
completedTimeNew int64
|
|
readers int32
|
|
consumed int32
|
|
pins int32 // streams currently positioned inside this chunk
|
|
left int32 // set once the last pinning stream moved on
|
|
sync.Mutex
|
|
parent *ReaderCache
|
|
chunkFileId string
|
|
data []byte
|
|
err error
|
|
cipherKey []byte
|
|
isGzipped bool
|
|
chunkSize int
|
|
shouldCache bool
|
|
wg sync.WaitGroup
|
|
cacheStartedCh chan struct{}
|
|
done chan struct{} // signals when download is complete
|
|
}
|
|
|
|
func NewReaderCache(limit int, chunkCache chunk_cache.ChunkCache, lookupFileIdFn wdclient.LookupFileIdFunctionType, cacheInvalidator CacheInvalidator, budgets ...*ReaderCacheBudget) *ReaderCache {
|
|
var budget *ReaderCacheBudget
|
|
if len(budgets) > 0 {
|
|
budget = budgets[0]
|
|
}
|
|
return &ReaderCache{
|
|
limit: limit,
|
|
budget: budget,
|
|
chunkCache: chunkCache,
|
|
lookupFileIdFn: lookupFileIdFn,
|
|
cacheInvalidator: cacheInvalidator,
|
|
fetchChunkDataFn: util_http.RetriedFetchChunkData,
|
|
downloaders: make(map[string]*SingleChunkCacher),
|
|
}
|
|
}
|
|
|
|
// MaybeCache prefetches up to 'count' chunks ahead in parallel.
|
|
// This improves read throughput for sequential reads by keeping the
|
|
// network pipeline full with parallel chunk fetches.
|
|
func (rc *ReaderCache) MaybeCache(chunkViews *Interval[*ChunkView], count int) {
|
|
if rc.lookupFileIdFn == nil {
|
|
return
|
|
}
|
|
if count <= 0 {
|
|
count = 1
|
|
}
|
|
|
|
rc.Lock()
|
|
defer rc.Unlock()
|
|
|
|
if len(rc.downloaders) >= rc.limit {
|
|
return
|
|
}
|
|
|
|
cached := 0
|
|
for x := chunkViews; x != nil && cached < count; x = x.Next {
|
|
chunkView := x.Value
|
|
if _, found := rc.downloaders[chunkView.FileId]; found {
|
|
continue
|
|
}
|
|
if rc.chunkCache.IsInCache(chunkView.FileId, true) {
|
|
glog.V(4).Infof("%s is in cache", chunkView.FileId)
|
|
continue
|
|
}
|
|
|
|
if len(rc.downloaders) >= rc.limit {
|
|
// abort when slots are filled
|
|
return
|
|
}
|
|
|
|
// glog.V(4).Infof("prefetch %s offset %d", chunkView.FileId, chunkView.ViewOffset)
|
|
// 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
|
|
rc.downloaders[chunkView.FileId] = cacher
|
|
cached++
|
|
}
|
|
|
|
return
|
|
}
|
|
|
|
// 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
|
|
// not drop the buffer from under it.
|
|
type chunkStream struct {
|
|
cacher *SingleChunkCacher
|
|
}
|
|
|
|
func (rc *ReaderCache) ReadChunkAt(ctx context.Context, buffer []byte, fileId string, cipherKey []byte, isGzipped bool, offset int64, chunkSize int, shouldCache bool) (int, error) {
|
|
return rc.readChunkAt(ctx, nil, buffer, fileId, cipherKey, isGzipped, offset, chunkSize, shouldCache)
|
|
}
|
|
|
|
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()
|
|
}
|
|
break
|
|
}
|
|
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)
|
|
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 {
|
|
oldestFid, oldestTime := "", time.Now().UnixNano()
|
|
pinnedFid, pinnedTime := "", int64(0)
|
|
for fid, downloader := range rc.downloaders {
|
|
completedTime := atomic.LoadInt64(&downloader.completedTimeNew)
|
|
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]
|
|
delete(rc.downloaders, oldestFid)
|
|
rc.Unlock()
|
|
oldDownloader.destroy()
|
|
goto 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)
|
|
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)
|
|
return n, err
|
|
}
|
|
|
|
// 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.
|
|
func (stream *chunkStream) pin(cacher *SingleChunkCacher) (previous *SingleChunkCacher) {
|
|
if stream == nil || stream.cacher == cacher {
|
|
return nil
|
|
}
|
|
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) {
|
|
if stream == nil {
|
|
return nil
|
|
}
|
|
previous = stream.cacher
|
|
stream.cacher = nil
|
|
return previous
|
|
}
|
|
|
|
// releaseIfFinished unpins the stream's chunk once the stream has read it to
|
|
// the end, since the stream will not come back to it.
|
|
func (rc *ReaderCache) releaseIfFinished(stream *chunkStream, cacher *SingleChunkCacher, offset int64, n int, err error, chunkSize int) {
|
|
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)
|
|
}
|
|
|
|
// 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)
|
|
}
|
|
|
|
// unpin drops one stream's pin. Once no stream is positioned in the chunk it
|
|
// is dropped like UnCache would, as soon as no read is in progress either:
|
|
// here if none is, otherwise by the last read's removeConsumed.
|
|
func (rc *ReaderCache) unpin(downloader *SingleChunkCacher) {
|
|
if downloader == nil {
|
|
return
|
|
}
|
|
if atomic.AddInt32(&downloader.pins, -1) == 0 {
|
|
atomic.StoreInt32(&downloader.left, 1)
|
|
}
|
|
rc.removeConsumed(downloader)
|
|
}
|
|
|
|
func (rc *ReaderCache) UnCache(fileId string) {
|
|
rc.Lock()
|
|
downloader := rc.downloaders[fileId]
|
|
delete(rc.downloaders, fileId)
|
|
rc.Unlock()
|
|
if downloader != nil {
|
|
downloader.destroy()
|
|
}
|
|
}
|
|
|
|
func (rc *ReaderCache) remove(downloader *SingleChunkCacher) {
|
|
rc.Lock()
|
|
removed := rc.downloaders[downloader.chunkFileId] == downloader
|
|
if removed {
|
|
delete(rc.downloaders, downloader.chunkFileId)
|
|
}
|
|
rc.Unlock()
|
|
if removed {
|
|
downloader.destroy()
|
|
}
|
|
}
|
|
|
|
// 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
|
|
// 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)
|
|
}
|
|
rc.Unlock()
|
|
if removed {
|
|
downloader.destroy()
|
|
}
|
|
}
|
|
|
|
func (rc *ReaderCache) destroy() {
|
|
rc.Lock()
|
|
downloaders := rc.downloaders
|
|
rc.downloaders = make(map[string]*SingleChunkCacher)
|
|
rc.Unlock()
|
|
for _, downloader := range downloaders {
|
|
downloader.destroy()
|
|
}
|
|
}
|
|
|
|
func newSingleChunkCacher(parent *ReaderCache, fileId string, cipherKey []byte, isGzipped bool, chunkSize int, shouldCache bool) *SingleChunkCacher {
|
|
return &SingleChunkCacher{
|
|
parent: parent,
|
|
chunkFileId: fileId,
|
|
cipherKey: cipherKey,
|
|
isGzipped: isGzipped,
|
|
chunkSize: chunkSize,
|
|
shouldCache: shouldCache,
|
|
cacheStartedCh: make(chan struct{}),
|
|
done: make(chan struct{}),
|
|
}
|
|
}
|
|
|
|
// startCaching downloads a chunk shared by concurrent readers.
|
|
func (s *SingleChunkCacher) startCaching() {
|
|
s.wg.Add(1)
|
|
defer func() {
|
|
close(s.done)
|
|
s.wg.Done()
|
|
if s.hasCompletedError() {
|
|
s.parent.remove(s)
|
|
} else {
|
|
s.parent.budget.complete(s)
|
|
}
|
|
}()
|
|
|
|
s.cacheStartedCh <- struct{}{}
|
|
if err := s.parent.budget.reserve(s); err != nil {
|
|
s.setError(err)
|
|
return
|
|
}
|
|
|
|
// Intentionally use context.Background(), not a request-specific context.
|
|
// The downloaded chunk is a shared resource: multiple concurrent readers may
|
|
// wait on this same download via s.done. A request-scoped context that got
|
|
// cancelled would abort the download and error every other waiting reader.
|
|
// The download always runs to completion once started; readers that cancel
|
|
// individually drop out via readChunkAt's select on ctx.Done().
|
|
urlStrings, err := s.parent.lookupFileIdFn(context.Background(), s.chunkFileId)
|
|
if err != nil {
|
|
s.setError(fmt.Errorf("operation LookupFileId %s failed, err: %v", s.chunkFileId, err))
|
|
return
|
|
}
|
|
if len(urlStrings) == 0 {
|
|
s.setError(fmt.Errorf("operation LookupFileId %s failed, err: urls not found", s.chunkFileId))
|
|
return
|
|
}
|
|
|
|
data, fetchErr := s.fetchChunkData(context.Background(), urlStrings)
|
|
if fetchErr != nil {
|
|
data, fetchErr = s.retryFetchAfterCacheInvalidation(context.Background(), urlStrings, fetchErr)
|
|
}
|
|
|
|
// Now acquire lock to update state
|
|
s.Lock()
|
|
atomic.StoreInt64(&s.completedTimeNew, time.Now().UnixNano())
|
|
if fetchErr != nil {
|
|
s.err = fetchErr
|
|
} else {
|
|
s.data = data
|
|
if s.shouldCache {
|
|
s.parent.chunkCache.SetChunk(s.chunkFileId, s.data)
|
|
}
|
|
}
|
|
s.Unlock()
|
|
}
|
|
|
|
func (s *SingleChunkCacher) setError(err error) {
|
|
s.Lock()
|
|
defer s.Unlock()
|
|
s.err = err
|
|
atomic.StoreInt64(&s.completedTimeNew, time.Now().UnixNano())
|
|
}
|
|
|
|
func (s *SingleChunkCacher) hasCompletedError() bool {
|
|
if atomic.LoadInt64(&s.completedTimeNew) == 0 {
|
|
return false
|
|
}
|
|
s.Lock()
|
|
defer s.Unlock()
|
|
return s.err != nil
|
|
}
|
|
|
|
func (s *SingleChunkCacher) fetchChunkData(ctx context.Context, urlStrings []string) ([]byte, error) {
|
|
// Allocate buffer and download without holding the lock.
|
|
// This allows multiple downloads to proceed in parallel.
|
|
data := mem.Allocate(s.chunkSize)
|
|
_, fetchErr := s.parent.fetchChunkDataFn(ctx, data, urlStrings, s.cipherKey, s.isGzipped, true, 0, s.chunkFileId, refreshUrls(ctx, s.parent.cacheInvalidator, s.parent.lookupFileIdFn, s.chunkFileId))
|
|
if fetchErr != nil {
|
|
mem.Free(data)
|
|
return nil, fetchErr
|
|
}
|
|
return data, nil
|
|
}
|
|
|
|
func (s *SingleChunkCacher) retryFetchAfterCacheInvalidation(ctx context.Context, oldUrlStrings []string, originalErr error) ([]byte, error) {
|
|
var data []byte
|
|
err := retryFetchWithFreshLocations(ctx, s.parent.cacheInvalidator, s.parent.lookupFileIdFn, s.chunkFileId, oldUrlStrings, originalErr, func(newUrls []string) error {
|
|
var fetchErr error
|
|
data, fetchErr = s.fetchChunkData(ctx, newUrls)
|
|
return fetchErr
|
|
})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return data, nil
|
|
}
|
|
|
|
func (s *SingleChunkCacher) destroy() {
|
|
// wait for all reads to finish before destroying the data
|
|
s.wg.Wait()
|
|
s.Lock()
|
|
if s.data != nil {
|
|
mem.Free(s.data)
|
|
s.data = nil
|
|
}
|
|
s.Unlock()
|
|
s.parent.budget.release(s)
|
|
}
|
|
|
|
// readChunkAt reads data from the cached chunk.
|
|
// It waits for the download to complete if it's still in progress.
|
|
// The ctx parameter allows the reader to cancel its wait (but the download continues
|
|
// for other readers - see comment in startCaching about shared resource semantics).
|
|
// The caller must s.wg.Add(1) under the ReaderCache lock before calling; this only releases it.
|
|
func (s *SingleChunkCacher) readChunkAt(ctx context.Context, buf []byte, offset int64) (n int, err error) {
|
|
defer func() {
|
|
s.wg.Done()
|
|
atomic.AddInt32(&s.readers, -1)
|
|
s.parent.removeConsumed(s)
|
|
}()
|
|
|
|
// Wait for download to complete, but allow reader cancellation.
|
|
// Prioritize checking done first - if data is already available,
|
|
// return it even if context is also cancelled.
|
|
select {
|
|
case <-s.done:
|
|
// Download already completed, proceed immediately
|
|
default:
|
|
// Download not complete, wait for it or context cancellation
|
|
select {
|
|
case <-s.done:
|
|
// Download completed
|
|
case <-ctx.Done():
|
|
// Reader cancelled while waiting - download continues for other readers
|
|
return 0, ctx.Err()
|
|
}
|
|
}
|
|
|
|
s.Lock()
|
|
defer s.Unlock()
|
|
|
|
if s.err != nil {
|
|
return 0, s.err
|
|
}
|
|
|
|
if len(s.data) <= int(offset) {
|
|
return 0, nil
|
|
}
|
|
|
|
n = copy(buf, s.data[offset:])
|
|
if offset+int64(n) == int64(len(s.data)) {
|
|
atomic.StoreInt32(&s.consumed, 1)
|
|
}
|
|
return n, nil
|
|
}
|