Compare commits

...
Author SHA1 Message Date
Chris Lu da39ee65a3 mount: add unit tests for refresh event buffering
Test all key scenarios for the cache refresh mechanism:

- Delete during refresh: subscription delete is buffered and replayed,
  preventing ghost entries from the stale snapshot
- Create during refresh: subscription create is buffered and replayed,
  preventing lost entries
- Stale entries removed: entries absent from the filer snapshot are
  cleaned from the cache
- Local delete buffered: direct DeleteEntry during refresh is both
  applied immediately and replayed after commit
- Local create buffered: direct InsertEntry during refresh is preserved
  after snapshot commit
- Cancel discards buffer: CancelRefresh drops buffered events and
  resumes normal event processing
- Concurrent deletes (issue #8442 scenario): 1000 concurrent delete
  events racing with a stale snapshot all resolve correctly
2026-03-06 00:47:37 -08:00
Chris Lu e91d43ef08 mount: use BeginRefresh/CommitRefresh in doEnsureVisited
Replace the incremental batch-insert approach in doEnsureVisited with
the new refresh lifecycle:

1. BeginRefresh - start buffering subscription events
2. ReadDirAllEntries - fetch full listing (no lock held)
3. CommitRefresh - atomically replace cache + replay buffered events

This ensures that any creates, deletes, or updates that arrive via the
subscription handler during the filer listing are not lost. The snapshot
replaces all stale entries, and buffered events are replayed on top to
bring the cache up to date.
2026-03-06 00:46:10 -08:00
Chris Lu 58cf7e9a89 mount: add refresh event buffering to MetaCache
Add infrastructure to buffer subscription events during directory cache
refresh. When a directory is being refreshed from the filer (BeginRefresh),
events from the subscription handler, local deletes, and local creates
are buffered. When the refresh completes (CommitRefresh), the filer
snapshot atomically replaces the cached entries, then buffered events
are replayed to ensure no mutations are lost.

This fixes a race condition where concurrent deletes or creates could be
overwritten by a stale filer snapshot during directory refresh, causing
ghost entries or lost files.

Fixes #8442
2026-03-06 00:45:23 -08:00
3 changed files with 459 additions and 34 deletions
+104 -8
View File
@@ -18,6 +18,20 @@ import (
// need to have logic similar to FilerStoreWrapper
// e.g. fill fileId field for chunks
// bufferedEvent represents a subscription event captured during a directory refresh.
type bufferedEvent struct {
oldPath util.FullPath
newEntry *filer.Entry
}
// refreshState tracks events that arrive while a directory is being refreshed
// from the filer. After the refresh snapshot is applied, buffered events are
// replayed so that creates, deletes, and updates that raced with the snapshot
// are not lost.
type refreshState struct {
events []bufferedEvent
}
type MetaCache struct {
root util.FullPath
localStore filer.VirtualFilerStore
@@ -29,6 +43,7 @@ type MetaCache struct {
invalidateFunc func(fullpath util.FullPath, entry *filer_pb.Entry)
onDirectoryUpdate func(dir util.FullPath)
visitGroup singleflight.Group // deduplicates concurrent EnsureVisited calls for the same path
refreshing map[util.FullPath]*refreshState
}
func NewMetaCache(dbFolder string, uidGidMapper *UidGidMapper, root util.FullPath,
@@ -45,6 +60,7 @@ func NewMetaCache(dbFolder string, uidGidMapper *UidGidMapper, root util.FullPat
invalidateFunc: func(fullpath util.FullPath, entry *filer_pb.Entry) {
invalidateFunc(fullpath, entry)
},
refreshing: make(map[util.FullPath]*refreshState),
}
}
@@ -69,6 +85,11 @@ func openMetaStore(dbFolder string) (*leveldb.LevelDBStore, filer.VirtualFilerSt
func (mc *MetaCache) InsertEntry(ctx context.Context, entry *filer.Entry) error {
mc.Lock()
defer mc.Unlock()
// Buffer the insert if the parent directory is being refreshed
dir, _ := entry.DirAndName()
if state := mc.isRefreshingDir(util.FullPath(dir)); state != nil {
state.events = append(state.events, bufferedEvent{newEntry: entry})
}
return mc.doInsertEntry(ctx, entry)
}
@@ -76,16 +97,33 @@ func (mc *MetaCache) doInsertEntry(ctx context.Context, entry *filer.Entry) erro
return mc.localStore.InsertEntry(ctx, entry)
}
// doBatchInsertEntries inserts multiple entries using LevelDB's batch write.
// This is more efficient than inserting entries one by one.
func (mc *MetaCache) doBatchInsertEntries(ctx context.Context, entries []*filer.Entry) error {
return mc.leveldbStore.BatchInsertEntries(ctx, entries)
}
func (mc *MetaCache) AtomicUpdateEntryFromFiler(ctx context.Context, oldPath util.FullPath, newEntry *filer.Entry) error {
mc.Lock()
defer mc.Unlock()
// If the affected directory is being refreshed, buffer the event
// instead of applying it. It will be replayed after the snapshot commits.
if oldPath != "" {
dir, _ := oldPath.DirAndName()
if state := mc.isRefreshingDir(util.FullPath(dir)); state != nil {
state.events = append(state.events, bufferedEvent{oldPath, newEntry})
return nil
}
}
if newEntry != nil {
newDir, _ := newEntry.DirAndName()
if state := mc.isRefreshingDir(util.FullPath(newDir)); state != nil {
state.events = append(state.events, bufferedEvent{oldPath, newEntry})
return nil
}
}
return mc.doAtomicUpdateEntryFromFiler(ctx, oldPath, newEntry)
}
// doAtomicUpdateEntryFromFiler is the core logic for applying a filer event
// to the local cache. Caller must hold mc.Lock().
func (mc *MetaCache) doAtomicUpdateEntryFromFiler(ctx context.Context, oldPath util.FullPath, newEntry *filer.Entry) error {
entry, err := mc.localStore.FindEntry(ctx, oldPath)
if err != nil && err != filer_pb.ErrNotFound {
glog.Errorf("Metacache: find entry error: %v", err)
@@ -104,8 +142,6 @@ func (mc *MetaCache) AtomicUpdateEntryFromFiler(ctx context.Context, oldPath uti
}
}
}
} else {
// println("unknown old directory:", oldDir)
}
if newEntry != nil {
@@ -143,6 +179,14 @@ func (mc *MetaCache) FindEntry(ctx context.Context, fp util.FullPath) (entry *fi
func (mc *MetaCache) DeleteEntry(ctx context.Context, fp util.FullPath) (err error) {
mc.Lock()
defer mc.Unlock()
// Buffer the delete if the parent directory is being refreshed
dir, _ := fp.DirAndName()
if state := mc.isRefreshingDir(util.FullPath(dir)); state != nil {
state.events = append(state.events, bufferedEvent{oldPath: fp})
}
// Always apply the delete directly as well, so the entry is removed
// immediately for the current node's view. CommitRefresh's snapshot
// may re-insert it, but the buffered event replay will re-delete it.
return mc.localStore.DeleteEntry(ctx, fp)
}
func (mc *MetaCache) DeleteFolderChildren(ctx context.Context, fp util.FullPath) (err error) {
@@ -151,6 +195,58 @@ func (mc *MetaCache) DeleteFolderChildren(ctx context.Context, fp util.FullPath)
return mc.localStore.DeleteFolderChildren(ctx, fp)
}
// BeginRefresh starts buffering subscription events for dirPath.
// While a refresh is active, AtomicUpdateEntryFromFiler will buffer events
// instead of applying them, so they can be replayed after the snapshot is committed.
func (mc *MetaCache) BeginRefresh(dirPath util.FullPath) {
mc.Lock()
defer mc.Unlock()
mc.refreshing[dirPath] = &refreshState{}
}
// CommitRefresh atomically replaces a directory's cached entries with the
// filer snapshot, then replays any subscription events that were buffered
// during the refresh. This ensures no creates, deletes, or updates are lost.
func (mc *MetaCache) CommitRefresh(ctx context.Context, dirPath util.FullPath, entries []*filer.Entry) error {
mc.Lock()
defer mc.Unlock()
// Clear stale entries and insert the fresh snapshot
if err := mc.localStore.DeleteFolderChildren(ctx, dirPath); err != nil {
return err
}
if len(entries) > 0 {
if err := mc.leveldbStore.BatchInsertEntries(ctx, entries); err != nil {
return err
}
}
// Replay buffered events so mutations that raced with the snapshot are applied
state := mc.refreshing[dirPath]
delete(mc.refreshing, dirPath)
if state != nil {
for _, ev := range state.events {
if err := mc.doAtomicUpdateEntryFromFiler(ctx, ev.oldPath, ev.newEntry); err != nil {
glog.Warningf("replay buffered event for %s: %v", dirPath, err)
}
}
}
return nil
}
// CancelRefresh discards the refresh state without replaying buffered events.
func (mc *MetaCache) CancelRefresh(dirPath util.FullPath) {
mc.Lock()
defer mc.Unlock()
delete(mc.refreshing, dirPath)
}
// isRefreshing returns the refresh state for the directory containing fp,
// or nil if no refresh is active. Caller must hold mc.Lock().
func (mc *MetaCache) isRefreshingDir(dirPath util.FullPath) *refreshState {
return mc.refreshing[dirPath]
}
func (mc *MetaCache) ListDirectoryEntries(ctx context.Context, dirPath util.FullPath, startFileName string, includeStartFile bool, limit int64, eachEntryFunc filer.ListEachEntryFunc) error {
mc.RLock()
defer mc.RUnlock()
+15 -26
View File
@@ -49,11 +49,6 @@ func EnsureVisited(mc *MetaCache, client filer_pb.FilerClient, dirPath util.Full
return g.Wait()
}
// batchInsertSize is the number of entries to accumulate before flushing to LevelDB.
// 100 provides a balance between memory usage (~100 Entry pointers) and write efficiency
// (fewer disk syncs). Larger values reduce I/O overhead but increase memory and latency.
const batchInsertSize = 100
func doEnsureVisited(ctx context.Context, mc *MetaCache, client filer_pb.FilerClient, path util.FullPath) error {
// Use singleflight to deduplicate concurrent requests for the same path
_, err, _ := mc.visitGroup.Do(string(path), func() (interface{}, error) {
@@ -69,42 +64,36 @@ func doEnsureVisited(ctx context.Context, mc *MetaCache, client filer_pb.FilerCl
glog.V(4).Infof("ReadDirAllEntries %s ...", path)
// Collect entries in batches for efficient LevelDB writes
var batch []*filer.Entry
// Start buffering subscription events for this directory.
// Any events that arrive during the fetch will be replayed
// after the snapshot is committed, preventing lost mutations.
mc.BeginRefresh(path)
// Collect all entries from the filer. No lock is held during
// network I/O so subscription events can still be buffered.
var allEntries []*filer.Entry
fetchErr := util.Retry("ReadDirAllEntries", func() error {
batch = nil // Reset batch on retry, allow GC of previous entries
allEntries = nil // Reset on retry
return filer_pb.ReadDirAllEntries(ctx, client, path, "", func(pbEntry *filer_pb.Entry, isLast bool) error {
entry := filer.FromPbEntry(string(path), pbEntry)
if IsHiddenSystemEntry(string(path), entry.Name()) {
return nil
}
batch = append(batch, entry)
// Flush batch when it reaches the threshold
// Don't rely on isLast here - hidden entries may cause early return
if len(batch) >= batchInsertSize {
// No lock needed - LevelDB Write() is thread-safe
if err := mc.doBatchInsertEntries(ctx, batch); err != nil {
return fmt.Errorf("batch insert for %s: %w", path, err)
}
// Create new slice to allow GC of flushed entries
batch = make([]*filer.Entry, 0, batchInsertSize)
}
allEntries = append(allEntries, entry)
return nil
})
})
if fetchErr != nil {
mc.CancelRefresh(path)
return nil, fmt.Errorf("list %s: %w", path, fetchErr)
}
// Flush any remaining entries in the batch
if len(batch) > 0 {
if err := mc.doBatchInsertEntries(ctx, batch); err != nil {
return nil, fmt.Errorf("batch insert remaining for %s: %w", path, err)
}
// Atomically replace cached entries with the snapshot and replay
// any buffered events that arrived during the fetch.
if err := mc.CommitRefresh(ctx, path, allEntries); err != nil {
return nil, fmt.Errorf("commit refresh for %s: %w", path, err)
}
mc.markCachedFn(path)
return nil, nil
+340
View File
@@ -0,0 +1,340 @@
package meta_cache
import (
"context"
"fmt"
"path/filepath"
"sync"
"testing"
"github.com/seaweedfs/seaweedfs/weed/filer"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/util"
)
func newTestMetaCache(t *testing.T) (*MetaCache, map[util.FullPath]bool) {
t.Helper()
uidGidMapper, err := NewUidGidMapper("", "")
if err != nil {
t.Fatalf("create uid/gid mapper: %v", err)
}
cached := make(map[util.FullPath]bool)
mc := NewMetaCache(
filepath.Join(t.TempDir(), "meta"),
uidGidMapper,
util.FullPath("/"),
func(path util.FullPath) { cached[path] = true },
func(path util.FullPath) bool { return cached[path] },
func(util.FullPath, *filer_pb.Entry) {},
nil,
)
t.Cleanup(func() { mc.Shutdown() })
return mc, cached
}
func makeEntry(dir, name string) *filer.Entry {
return &filer.Entry{
FullPath: util.NewFullPath(dir, name),
Attr: filer.Attr{Mode: 0644},
}
}
func listEntries(t *testing.T, mc *MetaCache, dir string) []string {
t.Helper()
var names []string
err := mc.ListDirectoryEntries(context.Background(), util.FullPath(dir), "", false, 10000, func(entry *filer.Entry) (bool, error) {
names = append(names, entry.Name())
return true, nil
})
if err != nil {
t.Fatalf("list %s: %v", dir, err)
}
return names
}
// TestCommitRefreshDeleteDuringRefresh verifies that a delete event arriving
// during a directory refresh is not overwritten by the stale snapshot.
func TestCommitRefreshDeleteDuringRefresh(t *testing.T) {
mc, cached := newTestMetaCache(t)
ctx := context.Background()
dir := util.FullPath("/testdir")
cached[dir] = true
// Pre-populate the cache with 3 entries
for _, name := range []string{"a", "b", "c"} {
if err := mc.InsertEntry(ctx, makeEntry("/testdir", name)); err != nil {
t.Fatalf("insert %s: %v", name, err)
}
}
// Start a refresh — simulate fetching a snapshot that includes all 3
mc.BeginRefresh(dir)
// While refresh is in progress, a subscription event deletes "b"
if err := mc.AtomicUpdateEntryFromFiler(ctx, util.NewFullPath("/testdir", "b"), nil); err != nil {
t.Fatalf("atomic delete b: %v", err)
}
// Commit the snapshot (which still includes "b") — the buffered delete
// should be replayed, removing "b" from the final result
snapshot := []*filer.Entry{
makeEntry("/testdir", "a"),
makeEntry("/testdir", "b"),
makeEntry("/testdir", "c"),
}
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
t.Fatalf("commit refresh: %v", err)
}
names := listEntries(t, mc, "/testdir")
expected := map[string]bool{"a": true, "c": true}
if len(names) != len(expected) {
t.Fatalf("expected %v, got %v", expected, names)
}
for _, n := range names {
if !expected[n] {
t.Errorf("unexpected entry %q in listing", n)
}
}
}
// TestCommitRefreshCreateDuringRefresh verifies that a create event arriving
// during a directory refresh is preserved after the snapshot is applied.
func TestCommitRefreshCreateDuringRefresh(t *testing.T) {
mc, cached := newTestMetaCache(t)
ctx := context.Background()
dir := util.FullPath("/testdir")
cached[dir] = true
// Pre-populate with "a"
if err := mc.InsertEntry(ctx, makeEntry("/testdir", "a")); err != nil {
t.Fatalf("insert a: %v", err)
}
// Start refresh — snapshot will have "a" only
mc.BeginRefresh(dir)
// During refresh, a subscription event creates "d"
newEntry := makeEntry("/testdir", "d")
if err := mc.AtomicUpdateEntryFromFiler(ctx, "", newEntry); err != nil {
t.Fatalf("atomic create d: %v", err)
}
// Commit snapshot (only "a") — buffered create of "d" should be replayed
snapshot := []*filer.Entry{makeEntry("/testdir", "a")}
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
t.Fatalf("commit refresh: %v", err)
}
names := listEntries(t, mc, "/testdir")
expected := map[string]bool{"a": true, "d": true}
if len(names) != len(expected) {
t.Fatalf("expected %v, got %v", expected, names)
}
for _, n := range names {
if !expected[n] {
t.Errorf("unexpected entry %q", n)
}
}
}
// TestCommitRefreshStaleEntriesRemoved verifies that entries deleted on the
// filer (absent from the snapshot) are removed from the local cache.
func TestCommitRefreshStaleEntriesRemoved(t *testing.T) {
mc, cached := newTestMetaCache(t)
ctx := context.Background()
dir := util.FullPath("/testdir")
cached[dir] = true
// Pre-populate with a, b, c
for _, name := range []string{"a", "b", "c"} {
if err := mc.InsertEntry(ctx, makeEntry("/testdir", name)); err != nil {
t.Fatalf("insert %s: %v", name, err)
}
}
// Refresh with snapshot that only has "a" and "c" — "b" was deleted on filer
mc.BeginRefresh(dir)
snapshot := []*filer.Entry{
makeEntry("/testdir", "a"),
makeEntry("/testdir", "c"),
}
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
t.Fatalf("commit refresh: %v", err)
}
names := listEntries(t, mc, "/testdir")
for _, n := range names {
if n == "b" {
t.Errorf("stale entry 'b' should have been removed by refresh")
}
}
if len(names) != 2 {
t.Fatalf("expected [a, c], got %v", names)
}
}
// TestCommitRefreshLocalDeleteBuffered verifies that a local DeleteEntry
// during refresh is both applied immediately and replayed after commit.
func TestCommitRefreshLocalDeleteBuffered(t *testing.T) {
mc, cached := newTestMetaCache(t)
ctx := context.Background()
dir := util.FullPath("/testdir")
cached[dir] = true
// Pre-populate
for _, name := range []string{"a", "b"} {
if err := mc.InsertEntry(ctx, makeEntry("/testdir", name)); err != nil {
t.Fatalf("insert %s: %v", name, err)
}
}
mc.BeginRefresh(dir)
// Local delete of "a" — simulates Unlink calling DeleteEntry directly
if err := mc.DeleteEntry(ctx, util.NewFullPath("/testdir", "a")); err != nil {
t.Fatalf("delete a: %v", err)
}
// Verify "a" is gone immediately (before commit)
_, err := mc.FindEntry(ctx, util.NewFullPath("/testdir", "a"))
if err != filer_pb.ErrNotFound {
t.Fatalf("expected ErrNotFound for deleted entry, got: %v", err)
}
// Commit with snapshot that includes "a" — the buffered delete replays
snapshot := []*filer.Entry{
makeEntry("/testdir", "a"),
makeEntry("/testdir", "b"),
}
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
t.Fatalf("commit refresh: %v", err)
}
// "a" should still be gone after commit
names := listEntries(t, mc, "/testdir")
for _, n := range names {
if n == "a" {
t.Errorf("locally deleted entry 'a' reappeared after commit")
}
}
}
// TestCommitRefreshLocalCreateBuffered verifies that a local InsertEntry
// during refresh is preserved after the snapshot is applied.
func TestCommitRefreshLocalCreateBuffered(t *testing.T) {
mc, cached := newTestMetaCache(t)
ctx := context.Background()
dir := util.FullPath("/testdir")
cached[dir] = true
mc.BeginRefresh(dir)
// Local create during refresh
if err := mc.InsertEntry(ctx, makeEntry("/testdir", "new")); err != nil {
t.Fatalf("insert new: %v", err)
}
// Commit with empty snapshot — the buffered create should still appear
if err := mc.CommitRefresh(ctx, dir, nil); err != nil {
t.Fatalf("commit refresh: %v", err)
}
names := listEntries(t, mc, "/testdir")
if len(names) != 1 || names[0] != "new" {
t.Fatalf("expected [new], got %v", names)
}
}
// TestCancelRefreshDiscardsBuffer verifies that CancelRefresh discards
// buffered events and resumes normal event processing.
func TestCancelRefreshDiscardsBuffer(t *testing.T) {
mc, cached := newTestMetaCache(t)
ctx := context.Background()
dir := util.FullPath("/testdir")
cached[dir] = true
if err := mc.InsertEntry(ctx, makeEntry("/testdir", "a")); err != nil {
t.Fatalf("insert a: %v", err)
}
mc.BeginRefresh(dir)
// Buffer a delete event
if err := mc.AtomicUpdateEntryFromFiler(ctx, util.NewFullPath("/testdir", "a"), nil); err != nil {
t.Fatalf("atomic delete: %v", err)
}
// Cancel — buffered events are discarded
mc.CancelRefresh(dir)
// "a" should still be in cache since the buffered delete was discarded
entry, err := mc.FindEntry(ctx, util.NewFullPath("/testdir", "a"))
if err != nil {
t.Fatalf("find a: %v", err)
}
if entry == nil {
t.Fatal("entry 'a' should still exist after cancel")
}
// After cancel, events should apply normally (not buffered)
if err := mc.AtomicUpdateEntryFromFiler(ctx, util.NewFullPath("/testdir", "a"), nil); err != nil {
t.Fatalf("atomic delete after cancel: %v", err)
}
_, err = mc.FindEntry(ctx, util.NewFullPath("/testdir", "a"))
if err != filer_pb.ErrNotFound {
t.Fatalf("expected ErrNotFound after direct delete, got: %v", err)
}
}
// TestConcurrentDeletesDuringRefresh simulates the scenario from issue #8442:
// multiple concurrent deletes racing with a directory refresh.
func TestConcurrentDeletesDuringRefresh(t *testing.T) {
mc, cached := newTestMetaCache(t)
ctx := context.Background()
dir := util.FullPath("/testdir")
cached[dir] = true
const numFiles = 1000
// Pre-populate
for i := 0; i < numFiles; i++ {
name := fmt.Sprintf("file_%04d", i)
if err := mc.InsertEntry(ctx, makeEntry("/testdir", name)); err != nil {
t.Fatalf("insert %s: %v", name, err)
}
}
// Build snapshot (taken before deletes happen)
var snapshot []*filer.Entry
for i := 0; i < numFiles; i++ {
snapshot = append(snapshot, makeEntry("/testdir", fmt.Sprintf("file_%04d", i)))
}
mc.BeginRefresh(dir)
// Concurrently delete all files via subscription events (simulates
// deletes from two mount nodes)
var wg sync.WaitGroup
for i := 0; i < numFiles; i++ {
wg.Add(1)
go func(idx int) {
defer wg.Done()
name := fmt.Sprintf("file_%04d", idx)
fp := util.NewFullPath("/testdir", name)
mc.AtomicUpdateEntryFromFiler(ctx, fp, nil)
}(i)
}
wg.Wait()
// Commit with the stale snapshot — all deletes should be replayed
if err := mc.CommitRefresh(ctx, dir, snapshot); err != nil {
t.Fatalf("commit refresh: %v", err)
}
names := listEntries(t, mc, "/testdir")
if len(names) != 0 {
t.Errorf("expected 0 entries after all deletes, got %d: first few: %v",
len(names), names[:min(5, len(names))])
}
}