Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
82 changes: 64 additions & 18 deletions component/attr_cache/attr_cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,8 @@ type AttrCache struct {
cleanupDone chan bool
cleanupCtx context.Context
cleanupStop context.CancelFunc

dirPrefetchThreshold uint32
}

// Structure defining your config parameters
Expand All @@ -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{}

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand All @@ -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
}
Expand Down
157 changes: 157 additions & 0 deletions component/attr_cache/attr_cache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
4 changes: 4 additions & 0 deletions component/attr_cache/cacheMap.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ package attr_cache

import (
"os"
"sync/atomic"
"time"

"github.com/Seagate/cloudfuse/common"
Expand Down Expand Up @@ -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
Expand Down
Loading