From 4d6da21ef73e5c7067451d09bfbb0d783a564911 Mon Sep 17 00:00:00 2001 From: James Fantin-Hardesty <24646452+jfantinhardesty@users.noreply.github.com> Date: Thu, 24 Sep 2026 11:45:06 -0600 Subject: [PATCH 1/3] Add new option to read ahead and cache listings in directories --- component/attr_cache/attr_cache.go | 157 +++++++++++++++++++++--- component/attr_cache/attr_cache_test.go | 83 +++++++++++++ setup/baseConfig.yaml | 1 + 3 files changed, 221 insertions(+), 20 deletions(-) diff --git a/component/attr_cache/attr_cache.go b/component/attr_cache/attr_cache.go index 0688755d9..1f646c241 100644 --- a/component/attr_cache/attr_cache.go +++ b/component/attr_cache/attr_cache.go @@ -60,6 +60,18 @@ type AttrCache struct { cleanupDone chan bool cleanupCtx context.Context cleanupStop context.CancelFunc + + dirPrefetchThreshold uint32 + prefetchLock sync.Mutex + prefetchState map[string]*dirPrefetchState // keyed by directory path +} + +// tracks attribute cache misses in one directory, to decide when to list it +type dirPrefetchState struct { + misses uint32 + windowStart time.Time + listedAt time.Time + done chan struct{} // non-nil while a listing is in flight } // Structure defining your config parameters @@ -73,6 +85,9 @@ type AttrCacheOptions struct { //maximum file attributes overall to be cached MaxFiles int `config:"max-files" yaml:"max-files,omitempty"` + // number of misses in one directory that triggers listing the whole directory (0 = disabled) + DirPrefetchThreshold uint32 `config:"dir-prefetch-threshold" yaml:"dir-prefetch-threshold,omitempty"` + // support v1 CacheOnList bool `config:"cache-on-list"` } @@ -83,6 +98,9 @@ const compName = "attr_cache" // caching more means increased memory usage of the process const defaultMaxFiles = 5000000 // 5 million max files overall to be cached +// maximum number of listing pages fetched by one directory prefetch +const dirPrefetchMaxPages = 10 + // Verification to check satisfaction criteria with Component Interface var _ internal.Component = &AttrCache{} @@ -181,12 +199,15 @@ func (ac *AttrCache) Configure(_ bool) error { ac.enableSymlinks = conf.EnableSymlinks } + ac.dirPrefetchThreshold = conf.DirPrefetchThreshold + log.Crit( - "AttrCache::Configure : cache-timeout %d, enable-symlinks %t, cache-on-list %t, max-files %d", + "AttrCache::Configure : cache-timeout %d, enable-symlinks %t, cache-on-list %t, max-files %d, dir-prefetch-threshold %d", ac.cacheTimeout, ac.enableSymlinks, ac.cacheOnList, ac.maxFiles, + ac.dirPrefetchThreshold, ) return nil @@ -410,6 +431,100 @@ func (ac *AttrCache) cleanupExpiredEntries() { } ac.cacheLock.Unlock() } + + ac.cleanupPrefetchState() +} + +// forget miss counts for directories that have not been touched within the cache timeout +func (ac *AttrCache) cleanupPrefetchState() { + timeout := time.Duration(ac.cacheTimeout) * time.Second + ac.prefetchLock.Lock() + defer ac.prefetchLock.Unlock() + for dirPath, state := range ac.prefetchState { + if state.done == nil && time.Since(state.windowStart) >= timeout && + time.Since(state.listedAt) >= timeout { + delete(ac.prefetchState, dirPath) + } + } +} + +// prefetchDir records an attribute cache miss for the parent directory of name. +// Once one directory accumulates dirPrefetchThreshold misses within the cache timeout, +// the directory is listed once, which caches the attributes of all its entries, and proves +// nonexistence of anything else. Returns true when a listing was fetched (or awaited), +// meaning the cache should be checked again. +func (ac *AttrCache) prefetchDir(name string) bool { + if ac.dirPrefetchThreshold == 0 || !ac.cacheOnList || ac.cacheTimeout == 0 { + return false + } + dirPath := getParentDir(name) + + // only list directories known to exist + ac.cacheLock.RLock() + dir, found := ac.cache.get(dirPath) + isDir := found && dir.exists() && dir.attr.IsDir() + ac.cacheLock.RUnlock() + if !isDir { + return false + } + + now := time.Now() + timeout := time.Duration(ac.cacheTimeout) * time.Second + ac.prefetchLock.Lock() + if ac.prefetchState == nil { + ac.prefetchState = make(map[string]*dirPrefetchState) + } + state, found := ac.prefetchState[dirPath] + if !found { + state = &dirPrefetchState{windowStart: now} + ac.prefetchState[dirPath] = state + } + if done := state.done; done != nil { + // coalesce with the listing already in flight + ac.prefetchLock.Unlock() + <-done + return true + } + if now.Sub(state.listedAt) < timeout { + // listed recently, possibly while this request was waiting - check the cache again + ac.prefetchLock.Unlock() + return true + } + if now.Sub(state.windowStart) >= timeout { + state.misses = 0 + state.windowStart = now + } + state.misses++ + if state.misses < ac.dirPrefetchThreshold { + ac.prefetchLock.Unlock() + return false + } + done := make(chan struct{}) + state.done = done + misses := state.misses + ac.prefetchLock.Unlock() + + log.Debug("AttrCache::prefetchDir : listing %s after %d misses", dirPath, misses) + token := "" + for range dirPrefetchMaxPages { + var err error + _, token, err = ac.StreamDir(internal.StreamDirOptions{Name: dirPath, Token: token}) + if err != nil { + log.Warn("AttrCache::prefetchDir : %s listing failed [%v]", dirPath, err) + break + } + if token == "" { + break + } + } + + ac.prefetchLock.Lock() + state.done = nil + state.misses = 0 + state.listedAt = time.Now() + ac.prefetchLock.Unlock() + close(done) + return true } // ------------------------- Methods implemented by this component ------------------------------------------- @@ -1106,24 +1221,19 @@ func (ac *AttrCache) SyncDir(options internal.SyncDirOptions) error { return err } -// GetAttr : Try to serve the request from the attribute cache, otherwise cache attributes of the path returned by next component -func (ac *AttrCache) GetAttr(options internal.GetAttrOptions) (*internal.ObjAttr, error) { - // Don't log these by default, as it noticeably affects performance - // log.Trace("AttrCache::GetAttr : %s", options.Name) - - // is the answer in the cache? +// lookupAttr returns the cached answer to GetAttr (if any), and whether it is fresh enough to serve +func (ac *AttrCache) lookupAttr(name string) (*internal.ObjAttr, bool, error) { respondFromCache := false var attrFromCache *internal.ObjAttr var errFromCache error ac.cacheLock.RLock() - value, found := ac.cache.get(options.Name) + defer ac.cacheLock.RUnlock() + value, found := ac.cache.get(name) if found && value.valid() { // record cache response if !value.exists() { - // log.Debug("AttrCache::GetAttr : %s found, (ENOENT) served from cache", options.Name) errFromCache = syscall.ENOENT } else { - // log.Debug("AttrCache::GetAttr : %s found, served from cache", options.Name) attrFromCache = value.attr } // only serve this response if it's not expired @@ -1133,25 +1243,32 @@ func (ac *AttrCache) GetAttr(options internal.GetAttrOptions) (*internal.ObjAttr } if !respondFromCache { // drill up for the nearest valid parent directory attribute cache - if parent, found := ac.cache.getCachedParent(options.Name); found { - // Remember, we have no entry for options.Name + if parent, found := ac.cache.getCachedParent(name); found { + // Remember, we have no entry for name // parent is its nearest valid ancestor - // So, if parent doesn't exist, options.Name must not exist + // So, if parent doesn't exist, name must not exist // Or, if parent does exist, and the full list of its contents are cached, - // then since options.Name is *not* in the cache, it must not exist + // then since name is *not* in the cache, it must not exist if !parent.exists() || parent.listingComplete { - // log.Debug( - // "AttrCache::GetAttr : %s not found, but parent exists(%t) or has a complete listing. ENOENT served from cache", - // options.Name, - // parent.exists(), - // ) errFromCache = syscall.ENOENT // only serve this response if it's not expired respondFromCache = time.Since(parent.cachedAt).Seconds() < float64(ac.cacheTimeout) } } } - ac.cacheLock.RUnlock() + return attrFromCache, respondFromCache, errFromCache +} + +// GetAttr : Try to serve the request from the attribute cache, otherwise cache attributes of the path returned by next component +func (ac *AttrCache) GetAttr(options internal.GetAttrOptions) (*internal.ObjAttr, error) { + // Don't log these by default, as it noticeably affects performance + // log.Trace("AttrCache::GetAttr : %s", options.Name) + + // is the answer in the cache? + attrFromCache, respondFromCache, errFromCache := ac.lookupAttr(options.Name) + if !respondFromCache && ac.prefetchDir(options.Name) { + attrFromCache, respondFromCache, errFromCache = ac.lookupAttr(options.Name) + } if respondFromCache { return attrFromCache, errFromCache } diff --git a/component/attr_cache/attr_cache_test.go b/component/attr_cache/attr_cache_test.go index b39b347e6..cdc0a3674 100644 --- a/component/attr_cache/attr_cache_test.go +++ b/component/attr_cache/attr_cache_test.go @@ -2097,6 +2097,89 @@ func (suite *attrCacheTestSuite) TestChown() { } } +func (suite *attrCacheTestSuite) TestDirPrefetchDisabledByDefault() { + defer suite.cleanupTest() + suite.addPathToCache("dir/") + + // every miss goes to cloud storage, and the directory is never listed + for _, attr := range generateListPathAttr("dir", 5) { + suite.mock.EXPECT().GetAttr(internal.GetAttrOptions{Name: attr.Path}).Return(attr, nil) + _, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: attr.Path}) + suite.assert.NoError(err) + } +} + +func (suite *attrCacheTestSuite) TestDirPrefetchOnMisses() { + defer suite.cleanupTest() + suite.cleanupTest() + suite.setupTestHelper("attr_cache:\n dir-prefetch-threshold: 3") + suite.addPathToCache("dir/") + listing := generateListPathAttr("dir", 10) + + // misses below the threshold are fetched individually + for _, attr := range listing[:2] { + suite.mock.EXPECT().GetAttr(internal.GetAttrOptions{Name: attr.Path}).Return(attr, nil) + _, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: attr.Path}) + suite.assert.NoError(err) + } + + // the third miss lists the directory (in two pages) instead + suite.mock.EXPECT(). + StreamDir(internal.StreamDirOptions{Name: "dir"}). + Return(listing[:5], "page2", nil) + suite.mock.EXPECT(). + StreamDir(internal.StreamDirOptions{Name: "dir", Token: "page2"}). + Return(listing[5:], "", nil) + for _, attr := range listing[2:] { + result, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: attr.Path}) + suite.assert.NoError(err) + suite.assert.Equal(attr.Path, result.Path) + } + + // the complete listing proves nonexistence without a cloud request + _, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: "dir/missing"}) + suite.assert.ErrorIs(err, syscall.ENOENT) +} + +func (suite *attrCacheTestSuite) TestDirPrefetchCoalescesConcurrentMisses() { + defer suite.cleanupTest() + suite.cleanupTest() + suite.setupTestHelper("attr_cache:\n dir-prefetch-threshold: 1") + suite.addPathToCache("dir/") + listing := generateListPathAttr("dir", 50) + + suite.mock.EXPECT(). + StreamDir(internal.StreamDirOptions{Name: "dir"}). + DoAndReturn(func(internal.StreamDirOptions) ([]*internal.ObjAttr, string, error) { + time.Sleep(50 * time.Millisecond) + return listing, "", nil + }). + Times(1) + + errs := make(chan error, len(listing)) + for _, attr := range listing { + go func() { + _, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: attr.Path}) + errs <- err + }() + } + for range listing { + suite.assert.NoError(<-errs) + } +} + +func (suite *attrCacheTestSuite) TestDirPrefetchSkipsUncachedDirectory() { + defer suite.cleanupTest() + suite.cleanupTest() + suite.setupTestHelper("attr_cache:\n dir-prefetch-threshold: 1") + + // the parent directory is not known to exist, so it is not listed + attr := getPathAttr("unknown/file", defaultSize, fs.FileMode(defaultMode)) + suite.mock.EXPECT().GetAttr(internal.GetAttrOptions{Name: attr.Path}).Return(attr, nil) + _, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: attr.Path}) + suite.assert.NoError(err) +} + // In order for 'go test' to run this suite, we need to create // a normal test function and pass our suite to suite.Run func TestAttrCacheTestSuite(t *testing.T) { diff --git a/setup/baseConfig.yaml b/setup/baseConfig.yaml index 95790bbde..36c39048c 100644 --- a/setup/baseConfig.yaml +++ b/setup/baseConfig.yaml @@ -128,6 +128,7 @@ attr_cache: no-cache-on-list: true|false enable-symlinks: true|false max-files: + dir-prefetch-threshold: # Loopback configuration loopbackfs: From cef66f3bf22a03a95412d4cb33b3e7623c448142 Mon Sep 17 00:00:00 2001 From: James Fantin-Hardesty <24646452+jfantinhardesty@users.noreply.github.com> Date: Thu, 8 Oct 2026 13:50:48 -0600 Subject: [PATCH 2/3] Simplify implementation and add tests --- component/attr_cache/attr_cache.go | 165 +++++++----------------- component/attr_cache/attr_cache_test.go | 101 +++++++++++++-- component/attr_cache/cacheMap.go | 4 + setup/baseConfig.yaml | 1 - 4 files changed, 138 insertions(+), 133 deletions(-) diff --git a/component/attr_cache/attr_cache.go b/component/attr_cache/attr_cache.go index 1f646c241..2ca978350 100644 --- a/component/attr_cache/attr_cache.go +++ b/component/attr_cache/attr_cache.go @@ -62,16 +62,6 @@ type AttrCache struct { cleanupStop context.CancelFunc dirPrefetchThreshold uint32 - prefetchLock sync.Mutex - prefetchState map[string]*dirPrefetchState // keyed by directory path -} - -// tracks attribute cache misses in one directory, to decide when to list it -type dirPrefetchState struct { - misses uint32 - windowStart time.Time - listedAt time.Time - done chan struct{} // non-nil while a listing is in flight } // Structure defining your config parameters @@ -85,9 +75,6 @@ type AttrCacheOptions struct { //maximum file attributes overall to be cached MaxFiles int `config:"max-files" yaml:"max-files,omitempty"` - // number of misses in one directory that triggers listing the whole directory (0 = disabled) - DirPrefetchThreshold uint32 `config:"dir-prefetch-threshold" yaml:"dir-prefetch-threshold,omitempty"` - // support v1 CacheOnList bool `config:"cache-on-list"` } @@ -98,8 +85,10 @@ const compName = "attr_cache" // caching more means increased memory usage of the process const defaultMaxFiles = 5000000 // 5 million max files overall to be cached -// maximum number of listing pages fetched by one directory prefetch -const dirPrefetchMaxPages = 10 +// number of cache misses in one directory that triggers listing one page of it +// 12 Head Objects requests is about 1 ListObjectsV2 cost via the AWS API. So +// set default to 6 to balance between prefetching and API cost +const defaultDirPrefetchThreshold = 6 // Verification to check satisfaction criteria with Component Interface var _ internal.Component = &AttrCache{} @@ -199,15 +188,14 @@ func (ac *AttrCache) Configure(_ bool) error { ac.enableSymlinks = conf.EnableSymlinks } - ac.dirPrefetchThreshold = conf.DirPrefetchThreshold + ac.dirPrefetchThreshold = defaultDirPrefetchThreshold log.Crit( - "AttrCache::Configure : cache-timeout %d, enable-symlinks %t, cache-on-list %t, max-files %d, dir-prefetch-threshold %d", + "AttrCache::Configure : cache-timeout %d, enable-symlinks %t, cache-on-list %t, max-files %d", ac.cacheTimeout, ac.enableSymlinks, ac.cacheOnList, ac.maxFiles, - ac.dirPrefetchThreshold, ) return nil @@ -431,100 +419,6 @@ func (ac *AttrCache) cleanupExpiredEntries() { } ac.cacheLock.Unlock() } - - ac.cleanupPrefetchState() -} - -// forget miss counts for directories that have not been touched within the cache timeout -func (ac *AttrCache) cleanupPrefetchState() { - timeout := time.Duration(ac.cacheTimeout) * time.Second - ac.prefetchLock.Lock() - defer ac.prefetchLock.Unlock() - for dirPath, state := range ac.prefetchState { - if state.done == nil && time.Since(state.windowStart) >= timeout && - time.Since(state.listedAt) >= timeout { - delete(ac.prefetchState, dirPath) - } - } -} - -// prefetchDir records an attribute cache miss for the parent directory of name. -// Once one directory accumulates dirPrefetchThreshold misses within the cache timeout, -// the directory is listed once, which caches the attributes of all its entries, and proves -// nonexistence of anything else. Returns true when a listing was fetched (or awaited), -// meaning the cache should be checked again. -func (ac *AttrCache) prefetchDir(name string) bool { - if ac.dirPrefetchThreshold == 0 || !ac.cacheOnList || ac.cacheTimeout == 0 { - return false - } - dirPath := getParentDir(name) - - // only list directories known to exist - ac.cacheLock.RLock() - dir, found := ac.cache.get(dirPath) - isDir := found && dir.exists() && dir.attr.IsDir() - ac.cacheLock.RUnlock() - if !isDir { - return false - } - - now := time.Now() - timeout := time.Duration(ac.cacheTimeout) * time.Second - ac.prefetchLock.Lock() - if ac.prefetchState == nil { - ac.prefetchState = make(map[string]*dirPrefetchState) - } - state, found := ac.prefetchState[dirPath] - if !found { - state = &dirPrefetchState{windowStart: now} - ac.prefetchState[dirPath] = state - } - if done := state.done; done != nil { - // coalesce with the listing already in flight - ac.prefetchLock.Unlock() - <-done - return true - } - if now.Sub(state.listedAt) < timeout { - // listed recently, possibly while this request was waiting - check the cache again - ac.prefetchLock.Unlock() - return true - } - if now.Sub(state.windowStart) >= timeout { - state.misses = 0 - state.windowStart = now - } - state.misses++ - if state.misses < ac.dirPrefetchThreshold { - ac.prefetchLock.Unlock() - return false - } - done := make(chan struct{}) - state.done = done - misses := state.misses - ac.prefetchLock.Unlock() - - log.Debug("AttrCache::prefetchDir : listing %s after %d misses", dirPath, misses) - token := "" - for range dirPrefetchMaxPages { - var err error - _, token, err = ac.StreamDir(internal.StreamDirOptions{Name: dirPath, Token: token}) - if err != nil { - log.Warn("AttrCache::prefetchDir : %s listing failed [%v]", dirPath, err) - break - } - if token == "" { - break - } - } - - ac.prefetchLock.Lock() - state.done = nil - state.misses = 0 - state.listedAt = time.Now() - ac.prefetchLock.Unlock() - close(done) - return true } // ------------------------- Methods implemented by this component ------------------------------------------- @@ -1221,8 +1115,9 @@ func (ac *AttrCache) SyncDir(options internal.SyncDirOptions) error { return err } -// lookupAttr returns the cached answer to GetAttr (if any), and whether it is fresh enough to serve -func (ac *AttrCache) lookupAttr(name string) (*internal.ObjAttr, bool, error) { +// lookupAttr returns the cached answer to GetAttr (if any), whether it is fresh enough to serve, +// and the parent directory's cache item (if it is a known existing directory) +func (ac *AttrCache) lookupAttr(name string) (*internal.ObjAttr, bool, *attrCacheItem, error) { respondFromCache := false var attrFromCache *internal.ObjAttr var errFromCache error @@ -1256,7 +1151,39 @@ func (ac *AttrCache) lookupAttr(name string) (*internal.ObjAttr, bool, error) { } } } - return attrFromCache, respondFromCache, errFromCache + dir, found := ac.cache.get(getParentDir(name)) + if !found || !dir.exists() || !dir.attr.IsDir() { + dir = nil + } + return attrFromCache, respondFromCache, dir, errFromCache +} + +// prefetchDirPage lists the next page of dir, to cache the attributes of its entries +func (ac *AttrCache) prefetchDirPage(dir *attrCacheItem) { + ac.cacheLock.RLock() + token, start := dir.prefetchToken, dir.prefetchStart + ac.cacheLock.RUnlock() + // restart a stale chain, so the oldest listed entries don't outlive the directory's listing + if token != "" && time.Since(start) >= time.Duration(ac.cacheTimeout)*time.Second { + token = "" + } + if token == "" { + start = time.Now() + } + log.Debug("AttrCache::prefetchDirPage : %s token=\"%s\"", dir.attr.Path, token) + _, next, err := ac.StreamDir(internal.StreamDirOptions{Name: dir.attr.Path, Token: token}) + ac.cacheLock.Lock() + if err == nil { + dir.prefetchToken = next + dir.prefetchStart = start + // the listing is only as fresh as its first page + if next == "" && token != "" && dir.listingComplete { + dir.cachedAt = start + } + } + ac.cacheLock.Unlock() + // hand off to the next winner after the token write + dir.prefetchMisses.Store(0) } // GetAttr : Try to serve the request from the attribute cache, otherwise cache attributes of the path returned by next component @@ -1265,9 +1192,11 @@ func (ac *AttrCache) GetAttr(options internal.GetAttrOptions) (*internal.ObjAttr // log.Trace("AttrCache::GetAttr : %s", options.Name) // is the answer in the cache? - attrFromCache, respondFromCache, errFromCache := ac.lookupAttr(options.Name) - if !respondFromCache && ac.prefetchDir(options.Name) { - attrFromCache, respondFromCache, errFromCache = ac.lookupAttr(options.Name) + attrFromCache, respondFromCache, dir, errFromCache := ac.lookupAttr(options.Name) + if !respondFromCache && dir != nil && ac.dirPrefetchThreshold > 0 && ac.cacheOnList && + ac.cacheTimeout > 0 && dir.prefetchMisses.Add(1) == ac.dirPrefetchThreshold { + ac.prefetchDirPage(dir) + attrFromCache, respondFromCache, _, errFromCache = ac.lookupAttr(options.Name) } if respondFromCache { return attrFromCache, errFromCache diff --git a/component/attr_cache/attr_cache_test.go b/component/attr_cache/attr_cache_test.go index cdc0a3674..5ab840e42 100644 --- a/component/attr_cache/attr_cache_test.go +++ b/component/attr_cache/attr_cache_test.go @@ -2097,12 +2097,13 @@ func (suite *attrCacheTestSuite) TestChown() { } } -func (suite *attrCacheTestSuite) TestDirPrefetchDisabledByDefault() { +func (suite *attrCacheTestSuite) TestDirPrefetchDisabled() { defer suite.cleanupTest() + suite.attrCache.dirPrefetchThreshold = 0 suite.addPathToCache("dir/") // every miss goes to cloud storage, and the directory is never listed - for _, attr := range generateListPathAttr("dir", 5) { + for _, attr := range generateListPathAttr("dir", 10) { suite.mock.EXPECT().GetAttr(internal.GetAttrOptions{Name: attr.Path}).Return(attr, nil) _, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: attr.Path}) suite.assert.NoError(err) @@ -2111,8 +2112,7 @@ func (suite *attrCacheTestSuite) TestDirPrefetchDisabledByDefault() { func (suite *attrCacheTestSuite) TestDirPrefetchOnMisses() { defer suite.cleanupTest() - suite.cleanupTest() - suite.setupTestHelper("attr_cache:\n dir-prefetch-threshold: 3") + suite.attrCache.dirPrefetchThreshold = 3 suite.addPathToCache("dir/") listing := generateListPathAttr("dir", 10) @@ -2123,13 +2123,10 @@ func (suite *attrCacheTestSuite) TestDirPrefetchOnMisses() { suite.assert.NoError(err) } - // the third miss lists the directory (in two pages) instead + // the third miss lists the directory instead suite.mock.EXPECT(). StreamDir(internal.StreamDirOptions{Name: "dir"}). - Return(listing[:5], "page2", nil) - suite.mock.EXPECT(). - StreamDir(internal.StreamDirOptions{Name: "dir", Token: "page2"}). - Return(listing[5:], "", nil) + Return(listing, "", nil) for _, attr := range listing[2:] { result, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: attr.Path}) suite.assert.NoError(err) @@ -2141,10 +2138,9 @@ func (suite *attrCacheTestSuite) TestDirPrefetchOnMisses() { suite.assert.ErrorIs(err, syscall.ENOENT) } -func (suite *attrCacheTestSuite) TestDirPrefetchCoalescesConcurrentMisses() { +func (suite *attrCacheTestSuite) TestDirPrefetchConcurrentMisses() { defer suite.cleanupTest() - suite.cleanupTest() - suite.setupTestHelper("attr_cache:\n dir-prefetch-threshold: 1") + suite.attrCache.dirPrefetchThreshold = 1 suite.addPathToCache("dir/") listing := generateListPathAttr("dir", 50) @@ -2155,6 +2151,13 @@ func (suite *attrCacheTestSuite) TestDirPrefetchCoalescesConcurrentMisses() { return listing, "", nil }). Times(1) + // misses during the listing go to cloud storage + suite.mock.EXPECT(). + GetAttr(gomock.Any()). + DoAndReturn(func(options internal.GetAttrOptions) (*internal.ObjAttr, error) { + return getPathAttr(options.Name, defaultSize, fs.FileMode(defaultMode)), nil + }). + AnyTimes() errs := make(chan error, len(listing)) for _, attr := range listing { @@ -2168,10 +2171,80 @@ func (suite *attrCacheTestSuite) TestDirPrefetchCoalescesConcurrentMisses() { } } +func (suite *attrCacheTestSuite) TestDirPrefetchMultiPage() { + defer suite.cleanupTest() + suite.attrCache.dirPrefetchThreshold = 1 + suite.addPathToCache("dir/") + listing := generateListPathAttr("dir", 10) + + // each trigger fetches one page, resuming from the stored token + suite.mock.EXPECT(). + StreamDir(internal.StreamDirOptions{Name: "dir"}). + Return(listing[:5], "page2", nil) + result, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: listing[0].Path}) + suite.assert.NoError(err) + suite.assert.Equal(listing[0].Path, result.Path) + + dir, _ := suite.attrCache.cache.get("dir") + suite.assert.Equal("page2", dir.prefetchToken) + chainStart := dir.prefetchStart + + suite.mock.EXPECT(). + StreamDir(internal.StreamDirOptions{Name: "dir", Token: "page2"}). + Return(listing[5:], "", nil) + result, err = suite.attrCache.GetAttr(internal.GetAttrOptions{Name: listing[7].Path}) + suite.assert.NoError(err) + suite.assert.Equal(listing[7].Path, result.Path) + + // the complete listing proves nonexistence, and is dated from its first page + _, err = suite.attrCache.GetAttr(internal.GetAttrOptions{Name: "dir/missing"}) + suite.assert.ErrorIs(err, syscall.ENOENT) + suite.assert.Empty(dir.prefetchToken) + suite.assert.True(dir.listingComplete) + suite.assert.Equal(chainStart, dir.cachedAt) +} + +func (suite *attrCacheTestSuite) TestDirPrefetchStaleChainRestarts() { + defer suite.cleanupTest() + suite.attrCache.dirPrefetchThreshold = 1 + suite.addPathToCache("dir/") + listing := generateListPathAttr("dir", 5) + dir, _ := suite.attrCache.cache.get("dir") + dir.prefetchToken = "page2" + dir.prefetchStart = time.Now().Add(-2 * time.Duration(suite.attrCache.cacheTimeout) * time.Second) + + suite.mock.EXPECT(). + StreamDir(internal.StreamDirOptions{Name: "dir"}). + Return(listing, "", nil) + _, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: listing[0].Path}) + suite.assert.NoError(err) + suite.assert.Empty(dir.prefetchToken) + suite.assert.WithinDuration(time.Now(), dir.prefetchStart, time.Second) +} + +func (suite *attrCacheTestSuite) TestDirPrefetchError() { + defer suite.cleanupTest() + suite.attrCache.dirPrefetchThreshold = 1 + suite.addPathToCache("dir/") + attr := getPathAttr("dir/file", defaultSize, fs.FileMode(defaultMode)) + dir, _ := suite.attrCache.cache.get("dir") + dir.prefetchToken = "page2" + dir.prefetchStart = time.Now() + + // a failed listing falls back to cloud storage, and keeps the token + suite.mock.EXPECT(). + StreamDir(internal.StreamDirOptions{Name: "dir", Token: "page2"}). + Return(nil, "", errors.New("failed to list")) + suite.mock.EXPECT().GetAttr(internal.GetAttrOptions{Name: attr.Path}).Return(attr, nil) + _, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: attr.Path}) + suite.assert.NoError(err) + suite.assert.Equal("page2", dir.prefetchToken) + suite.assert.Zero(dir.prefetchMisses.Load()) +} + func (suite *attrCacheTestSuite) TestDirPrefetchSkipsUncachedDirectory() { defer suite.cleanupTest() - suite.cleanupTest() - suite.setupTestHelper("attr_cache:\n dir-prefetch-threshold: 1") + suite.attrCache.dirPrefetchThreshold = 1 // the parent directory is not known to exist, so it is not listed attr := getPathAttr("unknown/file", defaultSize, fs.FileMode(defaultMode)) diff --git a/component/attr_cache/cacheMap.go b/component/attr_cache/cacheMap.go index 8c9f59f6c..4db804919 100644 --- a/component/attr_cache/cacheMap.go +++ b/component/attr_cache/cacheMap.go @@ -27,6 +27,7 @@ package attr_cache import ( "os" + "sync/atomic" "time" "github.com/Seagate/cloudfuse/common" @@ -60,6 +61,9 @@ type attrCacheItem struct { parent *attrCacheItem listingComplete bool + prefetchMisses atomic.Uint32 // misses since the last prefetch + prefetchToken string // next page to prefetch + prefetchStart time.Time // when the current token chain began } // all cache entries are organized into this structure diff --git a/setup/baseConfig.yaml b/setup/baseConfig.yaml index 36c39048c..95790bbde 100644 --- a/setup/baseConfig.yaml +++ b/setup/baseConfig.yaml @@ -128,7 +128,6 @@ attr_cache: no-cache-on-list: true|false enable-symlinks: true|false max-files: - dir-prefetch-threshold: # Loopback configuration loopbackfs: From 016e07d35e8d190de5c9d81340ea23fa746e450c Mon Sep 17 00:00:00 2001 From: James Fantin-Hardesty <24646452+jfantinhardesty@users.noreply.github.com> Date: Thu, 8 Oct 2026 13:53:49 -0600 Subject: [PATCH 3/3] Fix lint issue --- component/attr_cache/attr_cache_test.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/component/attr_cache/attr_cache_test.go b/component/attr_cache/attr_cache_test.go index 5ab840e42..dbb53930d 100644 --- a/component/attr_cache/attr_cache_test.go +++ b/component/attr_cache/attr_cache_test.go @@ -2211,7 +2211,8 @@ func (suite *attrCacheTestSuite) TestDirPrefetchStaleChainRestarts() { listing := generateListPathAttr("dir", 5) dir, _ := suite.attrCache.cache.get("dir") dir.prefetchToken = "page2" - dir.prefetchStart = time.Now().Add(-2 * time.Duration(suite.attrCache.cacheTimeout) * time.Second) + dir.prefetchStart = time.Now(). + Add(-2 * time.Duration(suite.attrCache.cacheTimeout) * time.Second) suite.mock.EXPECT(). StreamDir(internal.StreamDirOptions{Name: "dir"}).