Skip to content
Merged
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
27 changes: 19 additions & 8 deletions internal/fs/inode/dir.go
Original file line number Diff line number Diff line change
Expand Up @@ -265,7 +265,7 @@ type dirInode struct {
// (via kernel) the directory listing from the filesystem.
// Specially used when kernelListCacheTTL > 0 that means kernel list-cache is
// enabled.
prevDirListingTimeStamp time.Time
prevDirListingTimeStamp atomic.Int64
Comment thread
vedantdas-source marked this conversation as resolved.

metadataCacheTtlSecs int64

Expand Down Expand Up @@ -634,7 +634,6 @@ func (d *dirInode) CancelCurrDirPrefetcher() {
if d.prefetcher != nil {
d.prefetcher.Cancel()
}
return
}

// UpdateSize is a no-op for implicit directories. These directories are not
Expand Down Expand Up @@ -1105,7 +1104,12 @@ func (d *dirInode) ReadEntryCores(ctx context.Context, tok string) (cores map[Na
return
}

d.prevDirListingTimeStamp = d.cacheClock.Now()
if now := d.cacheClock.Now(); !now.IsZero() {
d.prevDirListingTimeStamp.Store(now.UnixNano())
} else {
d.prevDirListingTimeStamp.Store(0)
}

return
}

Expand Down Expand Up @@ -1508,12 +1512,19 @@ func (d *dirInode) LocalFileEntries(localFileInodes map[Name]Inode) (localEntrie
func (d *dirInode) ShouldInvalidateKernelListCache(ttl time.Duration) bool {
// prevDirListingTimeStamp.IsZero() true means listing has not happened yet, and we should
// invalidate for clean start.
if d.prevDirListingTimeStamp.IsZero() {
prevNS := d.prevDirListingTimeStamp.Load()
if prevNS == 0 {
return true
}

cachedDuration := d.cacheClock.Now().Sub(d.prevDirListingTimeStamp)
return cachedDuration >= ttl
now := d.cacheClock.Now()
if now.IsZero() {
return true
}
nowNS := now.UnixNano()
if nowNS < prevNS {
return true
}
return time.Duration(nowNS-prevNS) >= ttl
}

// LOCKS_REQUIRED(d)
Expand Down Expand Up @@ -1571,7 +1582,7 @@ func (d *dirInode) RenameFolder(ctx context.Context, folderName string, destinat

func (d *dirInode) InvalidateKernelListCache() {
// Set prevDirListingTimeStamp to Zero time so that cache is invalidated.
d.prevDirListingTimeStamp = time.Time{}
d.prevDirListingTimeStamp.Store(0)
}

func (d *dirInode) isBucketHierarchical() bool {
Expand Down
76 changes: 55 additions & 21 deletions internal/fs/inode/dir_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import (
"os"
"path"
"sort"
"sync"
"testing"
"time"

Expand Down Expand Up @@ -933,13 +934,13 @@ func (t *DirTest) TestReadDescendants_NonEmpty() {
func (t *DirTest) TestReadEntries_Empty() {
d := t.in.(*dirInode)
require.NotNil(t.T(), d)
require.True(t.T(), d.prevDirListingTimeStamp.IsZero())
require.Zero(t.T(), d.prevDirListingTimeStamp.Load())
entries, err := t.readAllEntries()

require.NoError(t.T(), err)
assert.ElementsMatch(t.T(), []fuseutil.Dirent{}, entries)
// Make sure prevDirListingTimeStamp is initialized.
require.False(t.T(), d.prevDirListingTimeStamp.IsZero())
require.NotZero(t.T(), d.prevDirListingTimeStamp.Load())
}

func (t *DirTest) TestReadEntries_NonEmpty_ImplicitDirsDisabled() {
Expand All @@ -966,7 +967,7 @@ func (t *DirTest) TestReadEntries_NonEmpty_ImplicitDirsDisabled() {
// Nil prevDirListingTimeStamp
d := t.in.(*dirInode)
require.NotNil(t.T(), d)
require.True(t.T(), d.prevDirListingTimeStamp.IsZero())
require.Zero(t.T(), d.prevDirListingTimeStamp.Load())

// Read entries.
entries, err := t.readAllEntries()
Expand Down Expand Up @@ -1003,7 +1004,7 @@ func (t *DirTest) TestReadEntries_NonEmpty_ImplicitDirsDisabled() {
}

// Make sure prevDirListingTimeStamp is initialized.
require.False(t.T(), d.prevDirListingTimeStamp.IsZero())
require.NotZero(t.T(), d.prevDirListingTimeStamp.Load())
}

func (t *DirTest) TestReadEntries_NonEmpty_ImplicitDirsEnabled() {
Expand Down Expand Up @@ -1033,7 +1034,7 @@ func (t *DirTest) TestReadEntries_NonEmpty_ImplicitDirsEnabled() {
// Nil prevDirListingTimeStamp
d := t.in.(*dirInode)
require.NotNil(t.T(), d)
require.True(t.T(), d.prevDirListingTimeStamp.IsZero())
require.Zero(t.T(), d.prevDirListingTimeStamp.Load())

// Read entries.
entries, err := t.readAllEntries()
Expand Down Expand Up @@ -1077,7 +1078,7 @@ func (t *DirTest) TestReadEntries_NonEmpty_ImplicitDirsEnabled() {
}

// Make sure prevDirListingTimeStamp is initialized.
require.False(t.T(), d.prevDirListingTimeStamp.IsZero())
require.NotZero(t.T(), d.prevDirListingTimeStamp.Load())
}

func (t *DirTest) TestReadEntries_TypeCaching() {
Expand All @@ -1097,7 +1098,7 @@ func (t *DirTest) TestReadEntries_TypeCaching() {
// Nil prevDirListingTimeStamp
d := t.in.(*dirInode)
require.NotNil(t.T(), d)
require.True(t.T(), d.prevDirListingTimeStamp.IsZero())
require.Zero(t.T(), d.prevDirListingTimeStamp.Load())

// Read the directory, priming the type cache.
_, err = t.readAllEntries()
Expand Down Expand Up @@ -1136,21 +1137,21 @@ func (t *DirTest) TestReadEntries_TypeCaching() {
assert.Equal(t.T(), dirObjName, result.MinObject.Name)

// Make sure prevDirListingTimeStamp is initialized.
require.False(t.T(), d.prevDirListingTimeStamp.IsZero())
require.NotZero(t.T(), d.prevDirListingTimeStamp.Load())
}

func (t *DirTest) TestReadEntryCores_Empty() {
d := t.in.(*dirInode)
require.NotNil(t.T(), d)
require.True(t.T(), d.prevDirListingTimeStamp.IsZero())
require.Zero(t.T(), d.prevDirListingTimeStamp.Load())

cores, unsupportedPaths, err := t.readAllEntryCores()

require.NoError(t.T(), err)
assert.Equal(t.T(), 0, len(cores))
assert.Equal(t.T(), 0, len(unsupportedPaths))
// Make sure prevDirListingTimeStamp is initialized.
require.False(t.T(), d.prevDirListingTimeStamp.IsZero())
require.NotZero(t.T(), d.prevDirListingTimeStamp.Load())
}

func (t *DirTest) TestReadEntryCores_NonEmpty_ImplicitDirsDisabled() {
Expand Down Expand Up @@ -1184,7 +1185,7 @@ func (t *DirTest) TestReadEntryCores_NonEmpty_ImplicitDirsDisabled() {
// Nil prevDirListingTimeStamp
d := t.in.(*dirInode)
require.NotNil(t.T(), d)
require.True(t.T(), d.prevDirListingTimeStamp.IsZero())
require.Zero(t.T(), d.prevDirListingTimeStamp.Load())

// Read cores.
cores, _, _, err = t.in.ReadEntryCores(t.ctx, "")
Expand All @@ -1196,7 +1197,7 @@ func (t *DirTest) TestReadEntryCores_NonEmpty_ImplicitDirsDisabled() {
t.validateCore(cores, "file", false, metadata.RegularFileType, testFileName)
t.validateCore(cores, "symlink", false, metadata.SymlinkType, symlinkName)
// Make sure prevDirListingTimeStamp is initialized.
require.False(t.T(), d.prevDirListingTimeStamp.IsZero())
require.NotZero(t.T(), d.prevDirListingTimeStamp.Load())
}

func (t *DirTest) TestReadEntryCores_NonEmpty_ImplicitDirsEnabled() {
Expand Down Expand Up @@ -1238,7 +1239,7 @@ func (t *DirTest) TestReadEntryCores_NonEmpty_ImplicitDirsEnabled() {
// Nil prevDirListingTimeStamp
d := t.in.(*dirInode)
require.NotNil(t.T(), d)
require.True(t.T(), d.prevDirListingTimeStamp.IsZero())
require.Zero(t.T(), d.prevDirListingTimeStamp.Load())

// Read cores.
cores, unsupportedPaths, err = t.readAllEntryCores()
Expand All @@ -1253,7 +1254,7 @@ func (t *DirTest) TestReadEntryCores_NonEmpty_ImplicitDirsEnabled() {
t.validateCore(cores, "symlink", false, metadata.SymlinkType, symlinkName)
assert.ElementsMatch(t.T(), []string{dirInodeName + "../", dirInodeName + "/"}, unsupportedPaths)
// Make sure prevDirListingTimeStamp is initialized.
require.False(t.T(), d.prevDirListingTimeStamp.IsZero())
require.NotZero(t.T(), d.prevDirListingTimeStamp.Load())
}

func (t *DirTest) TestCreateChildFile_DoesntExist() {
Expand Down Expand Up @@ -2023,7 +2024,7 @@ func (t *DirTest) TestLocalFileEntriesWithUnlinkedLocalChildFiles() {

func (t *DirTest) Test_ShouldInvalidateKernelListCache_ListingNotHappenedYet() {
d := t.in.(*dirInode)
d.prevDirListingTimeStamp = time.Time{}
d.prevDirListingTimeStamp.Store(0)

// Irrespective of the ttl value, this should always return true.
shouldInvalidate := t.in.ShouldInvalidateKernelListCache(util.MaxTimeDuration)
Expand All @@ -2033,7 +2034,7 @@ func (t *DirTest) Test_ShouldInvalidateKernelListCache_ListingNotHappenedYet() {

func (t *DirTest) Test_ShouldInvalidateKernelListCache_WithinTtl() {
d := t.in.(*dirInode)
d.prevDirListingTimeStamp = d.cacheClock.Now()
d.prevDirListingTimeStamp.Store(d.cacheClock.Now().UnixNano())
ttl := time.Second * 10
t.clock.AdvanceTime(ttl / 2)

Expand All @@ -2044,7 +2045,7 @@ func (t *DirTest) Test_ShouldInvalidateKernelListCache_WithinTtl() {

func (t *DirTest) Test_ShouldInvalidateKernelListCache_ExpiredTtl() {
d := t.in.(*dirInode)
d.prevDirListingTimeStamp = d.cacheClock.Now()
d.prevDirListingTimeStamp.Store(d.cacheClock.Now().UnixNano())
ttl := 10 * time.Second
t.clock.AdvanceTime(ttl + time.Second)

Expand All @@ -2055,7 +2056,7 @@ func (t *DirTest) Test_ShouldInvalidateKernelListCache_ExpiredTtl() {

func (t *DirTest) Test_ShouldInvalidateKernelListCache_ZeroTtl() {
d := t.in.(*dirInode)
d.prevDirListingTimeStamp = d.cacheClock.Now()
d.prevDirListingTimeStamp.Store(d.cacheClock.Now().UnixNano())
ttl := time.Duration(0)

shouldInvalidate := t.in.ShouldInvalidateKernelListCache(ttl)
Expand All @@ -2065,12 +2066,12 @@ func (t *DirTest) Test_ShouldInvalidateKernelListCache_ZeroTtl() {

func (t *DirTest) Test_InvalidateKernelListCache() {
d := t.in.(*dirInode)
d.prevDirListingTimeStamp = d.cacheClock.Now()
assert.False(t.T(), d.prevDirListingTimeStamp.IsZero())
d.prevDirListingTimeStamp.Store(d.cacheClock.Now().UnixNano())
assert.False(t.T(), d.prevDirListingTimeStamp.Load() == 0)

t.in.InvalidateKernelListCache()

assert.True(t.T(), d.prevDirListingTimeStamp.IsZero())
assert.True(t.T(), d.prevDirListingTimeStamp.Load() == 0)
}

func (t *DirTest) Test_ReadObjectsUnlocked() {
Expand Down Expand Up @@ -2402,3 +2403,36 @@ func (t *DirTest) TestMetadataPrefetcher_InitializationGuards() {
})
}
}

func (t *DirTest) Test_ShouldInvalidateKernelListCache_RaceCondition() {
Comment thread
vedantdas-source marked this conversation as resolved.
// This test demonstrates the concurrency data race that existed when using time.Time.
// If run with `go test -race`, it would fail prior to the atomic.Int64 refactor.
var wg sync.WaitGroup
wg.Add(3)

go func() {
defer wg.Done()
for range 1000 {
t.in.InvalidateKernelListCache()
}
}()

go func() {
defer wg.Done()
for range 1000 {
// Simulate concurrent lock-free reads
_ = t.in.ShouldInvalidateKernelListCache(time.Second)
}
}()

go func() {
defer wg.Done()
d := t.in.(*dirInode)
for range 1000 {
// Simulate concurrent writes (like what ReadEntryCores does)
d.prevDirListingTimeStamp.Store(d.cacheClock.Now().UnixNano())
}
}()

wg.Wait()
}
Loading