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
157 changes: 137 additions & 20 deletions component/attr_cache/attr_cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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"`
Comment on lines +88 to +89

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I honestly think this is so deep in the weeds that we shouldn't expose this option to the user until someone asks us to.


// support v1
CacheOnList bool `config:"cache-on-list"`
}
Expand All @@ -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{}

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 -------------------------------------------
Expand Down Expand Up @@ -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
Expand All @@ -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
}
Expand Down
83 changes: 83 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,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) {
Expand Down
1 change: 1 addition & 0 deletions setup/baseConfig.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -128,6 +128,7 @@ attr_cache:
no-cache-on-list: true|false <do not cache attributes or directory contents during listing. Enabling may cause performance problems.>
enable-symlinks: true|false <enable symlink support. When false, symlinks will be treated like regular files. Enabling may cause performance problems.>
max-files: <maximum number of files in the attribute cache at a time. Default - 5000000>
dir-prefetch-threshold: <number of attribute cache misses in one directory (within timeout-sec) after which the whole directory is listed, caching the attributes of all its entries. Speeds up stat-heavy workloads. 0 disables. Default - 0>

# Loopback configuration
loopbackfs:
Expand Down
Loading