diff --git a/component/attr_cache/attr_cache.go b/component/attr_cache/attr_cache.go index 0688755d9..2ca978350 100644 --- a/component/attr_cache/attr_cache.go +++ b/component/attr_cache/attr_cache.go @@ -60,6 +60,8 @@ type AttrCache struct { cleanupDone chan bool cleanupCtx context.Context cleanupStop context.CancelFunc + + dirPrefetchThreshold uint32 } // Structure defining your config parameters @@ -83,6 +85,11 @@ const compName = "attr_cache" // caching more means increased memory usage of the process const defaultMaxFiles = 5000000 // 5 million max files overall to be cached +// 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{} @@ -181,6 +188,8 @@ func (ac *AttrCache) Configure(_ bool) error { ac.enableSymlinks = conf.EnableSymlinks } + ac.dirPrefetchThreshold = defaultDirPrefetchThreshold + log.Crit( "AttrCache::Configure : cache-timeout %d, enable-symlinks %t, cache-on-list %t, max-files %d", ac.cacheTimeout, @@ -1106,24 +1115,20 @@ 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), 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 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 +1138,66 @@ 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) } } } + 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 +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, 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 b39b347e6..dbb53930d 100644 --- a/component/attr_cache/attr_cache_test.go +++ b/component/attr_cache/attr_cache_test.go @@ -2097,6 +2097,163 @@ func (suite *attrCacheTestSuite) TestChown() { } } +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", 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) + } +} + +func (suite *attrCacheTestSuite) TestDirPrefetchOnMisses() { + defer suite.cleanupTest() + suite.attrCache.dirPrefetchThreshold = 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 instead + suite.mock.EXPECT(). + StreamDir(internal.StreamDirOptions{Name: "dir"}). + Return(listing, "", 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) TestDirPrefetchConcurrentMisses() { + defer suite.cleanupTest() + suite.attrCache.dirPrefetchThreshold = 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) + // 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 { + go func() { + _, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: attr.Path}) + errs <- err + }() + } + for range listing { + suite.assert.NoError(<-errs) + } +} + +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.attrCache.dirPrefetchThreshold = 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/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