diff --git a/weed/mount/meta_cache/meta_cache.go b/weed/mount/meta_cache/meta_cache.go index f0135a2f4..d866c13f5 100644 --- a/weed/mount/meta_cache/meta_cache.go +++ b/weed/mount/meta_cache/meta_cache.go @@ -3,6 +3,7 @@ package meta_cache import ( "context" "errors" + "math" "os" "sync" "time" @@ -32,7 +33,8 @@ type MetaCache struct { isCachedFn func(fullpath util.FullPath) bool 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 + pinnedChildFn func(util.FullPath) bool // a child a rebuild must not drop (local-only, not yet on the filer); nil disables + visitGroup singleflight.Group // deduplicates concurrent EnsureVisited calls for the same path applyCh chan metadataApplyRequest applyDone chan struct{} applyStateMu sync.Mutex @@ -338,6 +340,42 @@ func (mc *MetaCache) DeleteFolderChildren(ctx context.Context, fp util.FullPath) return mc.localStore.DeleteFolderChildren(ctx, fp) } +// SetPinnedChildFn installs a predicate reporting whether a child holds +// local-only state a rebuild must not discard. See deleteFolderChildrenForRebuild. +func (mc *MetaCache) SetPinnedChildFn(fn func(util.FullPath) bool) { + mc.pinnedChildFn = fn +} + +// deleteFolderChildrenForRebuild clears a directory's cached children ahead of a +// rebuild, but keeps any child flagged pinned by pinnedChildFn — a local-only +// create not yet flushed to the filer. A rebuild refills from a filer listing +// that does not include such a create; a blind wipe would drop it and then +// markCachedFn publishes the directory authoritatively cached without a file the +// client created, so it vanishes from the mount. +func (mc *MetaCache) deleteFolderChildrenForRebuild(ctx context.Context, dirPath util.FullPath) error { + mc.Lock() + defer mc.Unlock() + if mc.pinnedChildFn == nil { + return mc.localStore.DeleteFolderChildren(ctx, dirPath) + } + var pinned []*filer.Entry + if _, err := mc.localStore.ListDirectoryEntries(ctx, dirPath, "", true, math.MaxInt64, func(entry *filer.Entry) (bool, error) { + if mc.pinnedChildFn(entry.FullPath) { + pinned = append(pinned, entry) + } + return true, nil + }); err != nil { + return err + } + if err := mc.localStore.DeleteFolderChildren(ctx, dirPath); err != nil { + return err + } + if len(pinned) > 0 { + return mc.doBatchInsertEntries(ctx, pinned) + } + return nil +} + 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() diff --git a/weed/mount/meta_cache/meta_cache_build_test.go b/weed/mount/meta_cache/meta_cache_build_test.go index 61285f23a..e48f41b9f 100644 --- a/weed/mount/meta_cache/meta_cache_build_test.go +++ b/weed/mount/meta_cache/meta_cache_build_test.go @@ -457,3 +457,95 @@ func TestBufferedRenameUpdatesOtherDirectoryBeforeBuildCompletes(t *testing.T) { t.Fatalf("replayed new path size = %d, want 12", newEntry.FileSize) } } + +// TestEnsureVisitedPreservesLocalOnlyEntry reproduces the residual coherence +// gap behind the FUSE ConcurrentReadWrite ENOENT flake. +// +// A FUSE create on the writeback/deferFilerCreate path inserts the entry into +// the local store directly (weedfs_file_mkrm.go createFile), off the metaCache +// apply loop, before the filer holds it. A concurrent rebuild of the parent — +// triggered when the directory falls out of cache (idle evict, hot-dir +// read-through) — wipes the store and refills it from a filer listing that does +// not yet include the un-flushed local create, then publishes the directory +// authoritatively cached (markCachedFn). The local entry vanishes although the +// client created it: lookupEntry then returns an authoritative ENOENT for it. +func TestEnsureVisitedPreservesLocalOnlyEntry(t *testing.T) { + mc, _, _, _ := newTestMetaCache(t, map[util.FullPath]bool{"/": true}) + defer mc.Shutdown() + + // The mount pins the un-flushed create (open dirty handle / pending flush). + mc.SetPinnedChildFn(func(p util.FullPath) bool { return p == "/dir/pending.txt" }) + + // A deferred local create lands before the rebuild; /dir is not yet cached. + insertCacheEntry(t, mc, "/dir/pending.txt") + + // A concurrent rebuild lists the filer, whose snapshot pre-dates the + // un-flushed create, so it returns only the already-persisted sibling. + accessor := &buildFilerAccessor{client: &buildListClient{ + responses: []*filer_pb.ListEntriesResponse{{ + Entry: &filer_pb.Entry{ + Name: "base.txt", + Attributes: &filer_pb.FuseAttributes{ + Crtime: 1, + Mtime: 1, + FileMode: 0100644, + FileSize: 3, + }, + }, + SnapshotTsNs: 100, + }}, + }} + + if err := EnsureVisited(mc, accessor, util.FullPath("/dir")); err != nil { + t.Fatalf("ensure visited: %v", err) + } + if !mc.IsDirectoryCached(util.FullPath("/dir")) { + t.Fatal("/dir should be cached after build completes") + } + + // base.txt from the listing is present. + if _, err := mc.FindEntry(context.Background(), util.FullPath("/dir/base.txt")); err != nil { + t.Fatalf("listed entry missing after build: %v", err) + } + + // The un-flushed local create must survive the rebuild. With /dir now + // authoritatively cached, losing it is the file-vanishes flake. + if _, err := mc.FindEntry(context.Background(), util.FullPath("/dir/pending.txt")); err != nil { + t.Fatalf("local-only entry lost across concurrent rebuild: %v", err) + } +} + +// TestEnsureVisitedDropsUnpinnedStaleEntry guards the preservation's selectivity: +// a cached child the filer listing no longer returns and that is NOT pinned must +// still be wiped, so the rebuild can't resurrect a deleted/renamed entry. +func TestEnsureVisitedDropsUnpinnedStaleEntry(t *testing.T) { + mc, _, _, _ := newTestMetaCache(t, map[util.FullPath]bool{"/": true}) + defer mc.Shutdown() + + mc.SetPinnedChildFn(func(util.FullPath) bool { return false }) + + // A stale child sits in the cache; the filer no longer has it. + insertCacheEntry(t, mc, "/dir/stale.txt") + + accessor := &buildFilerAccessor{client: &buildListClient{ + responses: []*filer_pb.ListEntriesResponse{{ + Entry: &filer_pb.Entry{ + Name: "base.txt", + Attributes: &filer_pb.FuseAttributes{ + Crtime: 1, + Mtime: 1, + FileMode: 0100644, + FileSize: 3, + }, + }, + SnapshotTsNs: 100, + }}, + }} + + if err := EnsureVisited(mc, accessor, util.FullPath("/dir")); err != nil { + t.Fatalf("ensure visited: %v", err) + } + if entry, err := mc.FindEntry(context.Background(), util.FullPath("/dir/stale.txt")); err != filer_pb.ErrNotFound || entry != nil { + t.Fatalf("unpinned stale entry survived rebuild = %+v, %v; want nil, %v", entry, err, filer_pb.ErrNotFound) + } +} diff --git a/weed/mount/meta_cache/meta_cache_init.go b/weed/mount/meta_cache/meta_cache_init.go index 60af3481c..1ce01b79a 100644 --- a/weed/mount/meta_cache/meta_cache_init.go +++ b/weed/mount/meta_cache/meta_cache_init.go @@ -82,7 +82,7 @@ func doEnsureVisited(ctx context.Context, mc *MetaCache, client filer_pb.FilerCl return } cleanupDone = true - if deleteErr := mc.DeleteFolderChildren(context.Background(), path); deleteErr != nil { + if deleteErr := mc.deleteFolderChildrenForRebuild(context.Background(), path); deleteErr != nil { glog.V(2).Infof("clear %s build %s: %v", reason, path, deleteErr) } if abortErr := mc.AbortDirectoryBuild(context.Background(), path); abortErr != nil { @@ -101,7 +101,7 @@ func doEnsureVisited(ctx context.Context, mc *MetaCache, client filer_pb.FilerCl fetchErr := util.Retry("ReadDirAllEntries", func() error { batch = nil // Reset batch on retry, allow GC of previous entries - if err := mc.DeleteFolderChildren(ctx, path); err != nil { + if err := mc.deleteFolderChildrenForRebuild(ctx, path); err != nil { return fmt.Errorf("clear existing entries for %s: %w", path, err) } var err error diff --git a/weed/mount/weedfs.go b/weed/mount/weedfs.go index cf0f9fe73..bc4bc222f 100644 --- a/weed/mount/weedfs.go +++ b/weed/mount/weedfs.go @@ -320,6 +320,7 @@ func NewSeaweedFileSystem(option *Option) *WFS { wfs.markDirectoryReadThrough(dirPath) } }) + wfs.metaCache.SetPinnedChildFn(wfs.isLocalOnlyEntry) grace.OnInterrupt(func() { // grace calls os.Exit(0) after all hooks, so WaitForAsyncFlush // after server.Serve() would never execute. Drain here first. @@ -549,6 +550,25 @@ func (wfs *WFS) maybeReadEntry(inode uint64) (path util.FullPath, fh *FileHandle return } +// isLocalOnlyEntry reports whether fullpath holds local-only state not yet on the +// filer — an open handle with dirty metadata, or a pending async flush. A +// directory rebuild refills from a filer listing that omits such an entry, so it +// must be preserved across the wipe; this is the same signal lookupEntry trusts +// over a filer ErrNotFound for deferred creates. +func (wfs *WFS) isLocalOnlyEntry(fullpath util.FullPath) bool { + inode, found := wfs.inodeToPath.GetInode(fullpath) + if !found { + return false + } + if fh, fhFound := wfs.fhMap.FindFileHandle(inode); fhFound && fh.dirtyMetadata { + return true + } + wfs.pendingAsyncFlushMu.Lock() + _, pending := wfs.pendingAsyncFlush[inode] + wfs.pendingAsyncFlushMu.Unlock() + return pending +} + func (wfs *WFS) maybeLoadEntry(fullpath util.FullPath) (*filer_pb.Entry, fuse.Status) { // glog.V(3).Infof("read entry cache miss %s", fullpath) _, name := fullpath.DirAndName()