From ed015b0d6b64843997012d53434befcc1ea6a1ae Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Wed, 22 Jul 2026 14:38:27 -0700 Subject: [PATCH] mount: order handle invalidations by filer log timestamp MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Three holes remained after resolving queued invalidations against the local store. A store hit is only trustworthy when the parent directory is children-cached — an uncached parent receives no store writes, so a leftover entry there is stale and would mask the event; gate the store read on the cached flag. A snapshot-covered buffered event got neither a store write nor a replay, yet its immediate invalidation may have run before the listing inserted the newer entry; build completion now re-invalidates every buffered event after publishing the directory. And in a read-through directory nothing reaches the store at all, so a queued event could still roll back a newer local flush: the filer already returns its log-stamped metadata event from CreateEntry, so the handle now keeps a watermark of its last filer-acknowledged local mutation and drops any subscription event at or before it — both sides of that comparison come from the filer clock, so it orders exactly. --- weed/mount/filehandle.go | 23 +++ weed/mount/filehandle_read.go | 1 + weed/mount/filehandle_read_remote_test.go | 2 +- weed/mount/meta_cache/meta_cache.go | 37 +++-- .../mount/meta_cache/meta_cache_apply_test.go | 7 +- .../meta_cache/meta_cache_deadlock_test.go | 2 +- weed/mount/weedfs.go | 38 +++-- weed/mount/weedfs_file_copy_range_test.go | 2 +- weed/mount/weedfs_file_mkrm_test.go | 2 +- weed/mount/weedfs_file_sync.go | 1 + .../weedfs_invalidate_open_handle_test.go | 140 +++++++++++++++--- weed/mount/weedfs_metadata_flush.go | 1 + weed/mount/weedfs_rename_test.go | 2 +- 13 files changed, 197 insertions(+), 61 deletions(-) diff --git a/weed/mount/filehandle.go b/weed/mount/filehandle.go index ba4a2cdb0..545164089 100644 --- a/weed/mount/filehandle.go +++ b/weed/mount/filehandle.go @@ -3,6 +3,7 @@ package mount import ( "os" "sync" + "sync/atomic" "github.com/seaweedfs/go-fuse/v2/fuse" "github.com/seaweedfs/seaweedfs/weed/cluster" @@ -39,6 +40,13 @@ type FileHandle struct { isDeleted bool isRenamed bool // set by Rename before waiting for async flush; skips old-path metadata flush + // lastLocalEntryTsNs is the filer log timestamp of the newest filer- + // acknowledged local mutation reflected in this handle's entry (from the + // metadata event a CreateEntry/CacheRemoteObject response carries). + // Subscription events at or before this timestamp are old news for the + // handle and must not roll it back. + lastLocalEntryTsNs atomic.Int64 + // dlmLock holds the distributed lock for cross-mount write coordination. // Non-nil only when -dlm is enabled and the file was opened for writing. // Acquired in AcquireHandle, released in ReleaseHandle. @@ -119,6 +127,21 @@ func (fh *FileHandle) SetEntry(entry *filer_pb.Entry) { fh.invalidateChunkCache() } +// advanceLocalEntryTs records the filer log timestamp of a filer-acknowledged +// local mutation now reflected in the handle's entry. Monotonic: an older +// timestamp never regresses the watermark. +func (fh *FileHandle) advanceLocalEntryTs(tsNs int64) { + if tsNs == 0 { + return + } + for { + current := fh.lastLocalEntryTsNs.Load() + if tsNs <= current || fh.lastLocalEntryTsNs.CompareAndSwap(current, tsNs) { + return + } + } +} + func (fh *FileHandle) ResetDirtyPages() { fh.dirtyPages.Destroy() fh.dirtyPages = newPageWriter(fh, fh.wfs.option.ChunkSizeLimit) diff --git a/weed/mount/filehandle_read.go b/weed/mount/filehandle_read.go index fac9cc23d..3dda56962 100644 --- a/weed/mount/filehandle_read.go +++ b/weed/mount/filehandle_read.go @@ -200,6 +200,7 @@ func (fh *FileHandle) downloadRemoteEntry(entry *LockedEntry) error { if event == nil { event = metadataUpdateEvent(request.Directory, resp.Entry) } + fh.advanceLocalEntryTs(event.GetTsNs()) fh.wfs.applyLocalMetadataEventAsync(event) return nil diff --git a/weed/mount/filehandle_read_remote_test.go b/weed/mount/filehandle_read_remote_test.go index a6cb0a5cc..3e05126b3 100644 --- a/weed/mount/filehandle_read_remote_test.go +++ b/weed/mount/filehandle_read_remote_test.go @@ -93,7 +93,7 @@ func TestReadUncachedRemoteEntryDoesNotDeadlock(t *testing.T) { func(path util.FullPath) { wfs.inodeToPath.MarkChildrenCached(path) }, func(path util.FullPath) bool { return wfs.inodeToPath.IsChildrenCached(path) }, // Mirror weedfs.go's invalidateFunc: take the file handle exclusive lock. - func(path util.FullPath, _ *filer_pb.Entry) { + func(path util.FullPath, _ *filer_pb.Entry, _ int64) { inode, ok := wfs.inodeToPath.GetInode(path) if !ok { return diff --git a/weed/mount/meta_cache/meta_cache.go b/weed/mount/meta_cache/meta_cache.go index afe9b115b..ff3a17404 100644 --- a/weed/mount/meta_cache/meta_cache.go +++ b/weed/mount/meta_cache/meta_cache.go @@ -31,7 +31,7 @@ type MetaCache struct { uidGidMapper *UidGidMapper markCachedFn func(fullpath util.FullPath) isCachedFn func(fullpath util.FullPath) bool - invalidateFunc func(fullpath util.FullPath, entry *filer_pb.Entry) + invalidateFunc func(fullpath util.FullPath, entry *filer_pb.Entry, eventTsNs int64) onDirectoryUpdate func(dir util.FullPath) pinnedChildFn func(*filer.Entry) 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 @@ -96,7 +96,7 @@ type metadataApplyRequest struct { } func NewMetaCache(dbFolder string, uidGidMapper *UidGidMapper, root util.FullPath, includeSystemEntries bool, - markCachedFn func(path util.FullPath), isCachedFn func(path util.FullPath) bool, invalidateFunc func(util.FullPath, *filer_pb.Entry), onDirectoryUpdate func(dir util.FullPath)) *MetaCache { + markCachedFn func(path util.FullPath), isCachedFn func(path util.FullPath) bool, invalidateFunc func(util.FullPath, *filer_pb.Entry, int64), onDirectoryUpdate func(dir util.FullPath)) *MetaCache { leveldbStore, virtualStore := openMetaStore(dbFolder) mc := &MetaCache{ root: root, @@ -107,8 +107,8 @@ func NewMetaCache(dbFolder string, uidGidMapper *UidGidMapper, root util.FullPat uidGidMapper: uidGidMapper, onDirectoryUpdate: onDirectoryUpdate, includeSystemEntries: includeSystemEntries, - invalidateFunc: func(fullpath util.FullPath, entry *filer_pb.Entry) { - invalidateFunc(fullpath, entry) + invalidateFunc: func(fullpath util.FullPath, entry *filer_pb.Entry, eventTsNs int64) { + invalidateFunc(fullpath, entry, eventTsNs) }, applyCh: make(chan metadataApplyRequest, 128), applyDone: make(chan struct{}), @@ -117,7 +117,7 @@ func NewMetaCache(dbFolder string, uidGidMapper *UidGidMapper, root util.FullPat } mc.invalidateWorker = util.NewAsyncBatchWorker(func(batch []metadataInvalidation) { for _, invalidation := range batch { - mc.invalidateFunc(invalidation.path, invalidation.entry) + mc.invalidateFunc(invalidation.path, invalidation.entry, invalidation.tsNs) } }) go mc.runApplyLoop() @@ -565,6 +565,7 @@ func (mc *MetaCache) handleApplyRequest(req metadataApplyRequest) error { type metadataInvalidation struct { path util.FullPath entry *filer_pb.Entry // entry now at path per the event; nil when the path was vacated (delete, rename away) + tsNs int64 // the event's filer log timestamp; 0 for locally built events } type metadataResponseSideEffects struct { @@ -749,17 +750,21 @@ func (mc *MetaCache) completeDirectoryBuildNow(ctx context.Context, dirPath util if snapshotTsNs != 0 && event.TsNs != 0 && event.TsNs <= snapshotTsNs { continue } - // Re-invalidate on replay: the invalidation enqueued when this event - // arrived resolved against the store mid-build, when the store could - // still hold the older listing snapshot (the event's own store write - // was deferred to this replay). Skipped events need no re-invalidation - // since the listing state is at least as new. - if err := mc.applyMetadataResponseDirect(ctx, event, MetadataResponseApplyOptions{InvalidateEntries: true}, true); err != nil { + if err := mc.applyMetadataResponseDirect(ctx, event, MetadataResponseApplyOptions{}, true); err != nil { return err } } mc.markCachedFn(dirPath) + + // Re-invalidate every buffered event, replayed or snapshot-covered: the + // invalidation issued when an event arrived ran against a mid-build store + // that could miss the path entirely or predate the listing insert, so an + // open handle can hold older state than the completed directory. Enqueued + // after markCachedFn so the refresh resolves against the published store. + for _, event := range state.bufferedEvents { + mc.applyMetadataSideEffects(event, MetadataResponseApplyOptions{InvalidateEntries: true}) + } return nil } @@ -986,11 +991,11 @@ func collectEntryInvalidations(resp *filer_pb.SubscribeMetadataResponse) []metad newDir = message.NewParentPath } if message.OldEntry.Name != message.NewEntry.Name || resp.Directory != newDir { - invalidations = append(invalidations, metadataInvalidation{path: oldKey}) + invalidations = append(invalidations, metadataInvalidation{path: oldKey, tsNs: resp.TsNs}) newKey := util.NewFullPath(newDir, message.NewEntry.Name) - invalidations = append(invalidations, metadataInvalidation{path: newKey, entry: message.NewEntry}) + invalidations = append(invalidations, metadataInvalidation{path: newKey, entry: message.NewEntry, tsNs: resp.TsNs}) } else { - invalidations = append(invalidations, metadataInvalidation{path: oldKey, entry: message.NewEntry}) + invalidations = append(invalidations, metadataInvalidation{path: oldKey, entry: message.NewEntry, tsNs: resp.TsNs}) } return invalidations } @@ -1001,12 +1006,12 @@ func collectEntryInvalidations(resp *filer_pb.SubscribeMetadataResponse) []metad newDir = message.NewParentPath } newKey := util.NewFullPath(newDir, message.NewEntry.Name) - invalidations = append(invalidations, metadataInvalidation{path: newKey, entry: message.NewEntry}) + invalidations = append(invalidations, metadataInvalidation{path: newKey, entry: message.NewEntry, tsNs: resp.TsNs}) } if filer_pb.IsDelete(resp) && message.OldEntry != nil { oldKey := util.NewFullPath(resp.Directory, message.OldEntry.Name) - invalidations = append(invalidations, metadataInvalidation{path: oldKey}) + invalidations = append(invalidations, metadataInvalidation{path: oldKey, tsNs: resp.TsNs}) } return invalidations diff --git a/weed/mount/meta_cache/meta_cache_apply_test.go b/weed/mount/meta_cache/meta_cache_apply_test.go index 2ba85cbfa..a6538730a 100644 --- a/weed/mount/meta_cache/meta_cache_apply_test.go +++ b/weed/mount/meta_cache/meta_cache_apply_test.go @@ -420,6 +420,7 @@ func TestCollectEntryInvalidationsCarryAuthoritativeEntries(t *testing.T) { update := &filer_pb.SubscribeMetadataResponse{ Directory: "/dir", + TsNs: 77, EventNotification: &filer_pb.EventNotification{ OldEntry: &filer_pb.Entry{Name: "file.txt"}, NewEntry: newEntry, @@ -427,8 +428,8 @@ func TestCollectEntryInvalidationsCarryAuthoritativeEntries(t *testing.T) { }, } got := collectEntryInvalidations(update) - if len(got) != 1 || got[0].path != "/dir/file.txt" || got[0].entry != newEntry { - t.Fatalf("in-place update invalidations = %+v, want [{/dir/file.txt NewEntry}]", got) + if len(got) != 1 || got[0].path != "/dir/file.txt" || got[0].entry != newEntry || got[0].tsNs != 77 { + t.Fatalf("in-place update invalidations = %+v, want [{/dir/file.txt NewEntry ts 77}]", got) } rename := &filer_pb.SubscribeMetadataResponse{ @@ -495,7 +496,7 @@ func newTestMetaCache(t *testing.T, cached map[util.FullPath]bool) (*MetaCache, defer cachedMu.Unlock() return cached[path] }, - func(path util.FullPath, entry *filer_pb.Entry) { + func(path util.FullPath, entry *filer_pb.Entry, eventTsNs int64) { invalidations.record(path) }, func(dir util.FullPath) { diff --git a/weed/mount/meta_cache/meta_cache_deadlock_test.go b/weed/mount/meta_cache/meta_cache_deadlock_test.go index ce94bee85..b8726d2ba 100644 --- a/weed/mount/meta_cache/meta_cache_deadlock_test.go +++ b/weed/mount/meta_cache/meta_cache_deadlock_test.go @@ -58,7 +58,7 @@ func TestApplyLoopInvalidateDoesNotDeadlockWithLockHoldingEnqueuer(t *testing.T) defer cachedMu.Unlock() return cached[path] }, - func(path util.FullPath, entry *filer_pb.Entry) { + func(path util.FullPath, entry *filer_pb.Entry, _ int64) { // Mirrors the wfs invalidateFunc: it takes the open file // handle's exclusive lock before refreshing the handle entry. enteredOnce.Do(func() { close(invalidateEntered) }) diff --git a/weed/mount/weedfs.go b/weed/mount/weedfs.go index 82cbdae3b..18c16abe6 100644 --- a/weed/mount/weedfs.go +++ b/weed/mount/weedfs.go @@ -302,8 +302,8 @@ func NewSeaweedFileSystem(option *Option) *WFS { wfs.inodeToPath.MarkChildrenCached(path) }, func(path util.FullPath) bool { return wfs.inodeToPath.IsChildrenCached(path) - }, func(filePath util.FullPath, entry *filer_pb.Entry) { - wfs.invalidateOpenFileHandle(filePath, entry) + }, func(filePath util.FullPath, entry *filer_pb.Entry, eventTsNs int64) { + wfs.invalidateOpenFileHandle(filePath, entry, eventTsNs) }, func(dirPath util.FullPath) { if wfs.inodeToPath.RecordDirectoryUpdate(dirPath, time.Now(), wfs.dirHotWindow, wfs.dirHotThreshold) { wfs.markDirectoryReadThrough(dirPath) @@ -682,7 +682,7 @@ func (wfs *WFS) lookupEntry(fullpath util.FullPath) (*filer.Entry, fuse.Status) // handle would stay pinned to its old entry until an unrelated event arrives. // A nil entry means the path no longer holds one (delete, rename away); the // handle keeps its last entry so unlinked-but-open reads still work. -func (wfs *WFS) invalidateOpenFileHandle(filePath util.FullPath, entry *filer_pb.Entry) { +func (wfs *WFS) invalidateOpenFileHandle(filePath util.FullPath, entry *filer_pb.Entry, eventTsNs int64) { inode, inodeFound := wfs.inodeToPath.GetInode(filePath) if !inodeFound { return @@ -694,22 +694,28 @@ func (wfs *WFS) invalidateOpenFileHandle(filePath util.FullPath, entry *filer_pb fhActiveLock := wfs.fhLockTable.AcquireLock("invalidateFunc", fh.fh, util.ExclusiveLock) defer wfs.fhLockTable.ReleaseLock(fh.fh, fhActiveLock) + // Invalidations apply asynchronously, so this event may be old news by + // now: a local flush can land while the event sits in the queue, and the + // flush's own event is dedup-suppressed, so a rollback would never heal. + // Both timestamps come from the filer log, so this orders exactly. + if eventTsNs != 0 && eventTsNs <= fh.lastLocalEntryTsNs.Load() { + return + } + fh.dirtyPages.Destroy() fh.dirtyPages = newPageWriter(fh, wfs.option.ChunkSizeLimit) - // Invalidations apply asynchronously, so the event entry may be a stale - // snapshot by now: a local flush can install newer state while the event - // sits in the queue, and the flush's own event is dedup-suppressed, so a - // rollback would never heal. The apply loop has already ordered this event - // and any later state into the local store, so prefer the store's entry. - // The store misses for read-through directories and TTL-expired entries; - // falling back to the event entry there can still roll back a racing - // local flush (neither write reaches the store in a read-through - // directory), but it is the best ordered information available without a - // filer round-trip. - if localEntry, findErr := wfs.metaCache.FindEntry(context.Background(), filePath); findErr == nil && localEntry != nil { - fh.SetEntry(localEntry.ToProtoEntry()) - return + // Prefer the local store's entry when the parent directory is cached: the + // apply loop has already ordered this event and anything newer (e.g. a + // local flush) into it. An uncached parent receives no store writes, so a + // hit there could be a stale leftover masking this event — fall through + // to the event entry instead, the freshest information for that case. + dir, _ := filePath.DirAndName() + if wfs.metaCache.IsDirectoryCached(util.FullPath(dir)) { + if localEntry, findErr := wfs.metaCache.FindEntry(context.Background(), filePath); findErr == nil && localEntry != nil { + fh.SetEntry(localEntry.ToProtoEntry()) + return + } } if entry == nil { return diff --git a/weed/mount/weedfs_file_copy_range_test.go b/weed/mount/weedfs_file_copy_range_test.go index 02c34409b..e931084b7 100644 --- a/weed/mount/weedfs_file_copy_range_test.go +++ b/weed/mount/weedfs_file_copy_range_test.go @@ -346,7 +346,7 @@ func newCopyRangeTestWFSWithMetaCache(t *testing.T) *WFS { func(path util.FullPath) bool { return wfs.inodeToPath.IsChildrenCached(path) }, - func(util.FullPath, *filer_pb.Entry) {}, + func(util.FullPath, *filer_pb.Entry, int64) {}, nil, ) t.Cleanup(func() { diff --git a/weed/mount/weedfs_file_mkrm_test.go b/weed/mount/weedfs_file_mkrm_test.go index ba3828a9c..eb2062111 100644 --- a/weed/mount/weedfs_file_mkrm_test.go +++ b/weed/mount/weedfs_file_mkrm_test.go @@ -128,7 +128,7 @@ func newCreateTestWFS(t *testing.T) (*WFS, *createEntryTestServer) { func(path util.FullPath) bool { return wfs.inodeToPath.IsChildrenCached(path) }, - func(util.FullPath, *filer_pb.Entry) {}, + func(util.FullPath, *filer_pb.Entry, int64) {}, nil, ) wfs.inodeToPath.MarkChildrenCached(root) diff --git a/weed/mount/weedfs_file_sync.go b/weed/mount/weedfs_file_sync.go index 57712d8a8..6ba8fb496 100644 --- a/weed/mount/weedfs_file_sync.go +++ b/weed/mount/weedfs_file_sync.go @@ -270,6 +270,7 @@ func (wfs *WFS) flushMetadataToFiler(ctx context.Context, fh *FileHandle, dir, n if event == nil { event = metadataUpdateEvent(string(dir), request.Entry) } + fh.advanceLocalEntryTs(event.GetTsNs()) if applyErr := wfs.applyLocalMetadataEvent(context.Background(), event); applyErr != nil { glog.Warningf("flush %s: best-effort metadata apply failed: %v", fileFullPath, applyErr) wfs.inodeToPath.InvalidateChildrenCache(util.FullPath(dir)) diff --git a/weed/mount/weedfs_invalidate_open_handle_test.go b/weed/mount/weedfs_invalidate_open_handle_test.go index 32ad95af2..20a278ff0 100644 --- a/weed/mount/weedfs_invalidate_open_handle_test.go +++ b/weed/mount/weedfs_invalidate_open_handle_test.go @@ -180,9 +180,10 @@ func TestQueuedEventDoesNotRollBackNewerLocalState(t *testing.T) { // During a directory build, an event touching the building directory is // buffered: its store write is deferred to build completion while its -// invalidation runs immediately, so that refresh can resolve to the older -// listing snapshot. The completion replay must re-invalidate so the handle -// lands on the event's state. +// invalidation runs immediately, against a store that may not reflect the +// listing yet. Build completion must re-invalidate every buffered event — +// including snapshot-covered ones — so the handle lands on the completed +// directory's state. func TestBufferedBuildEventReinvalidatesOnCompletion(t *testing.T) { wfs := newInvalidateTestWFS(t) @@ -198,29 +199,18 @@ func TestBufferedBuildEventReinvalidatesOnCompletion(t *testing.T) { if err := wfs.metaCache.BeginDirectoryBuild(context.Background(), util.FullPath("/dir")); err != nil { t.Fatalf("begin build: %v", err) } - // The in-progress listing inserts the pre-event snapshot. - if err := wfs.metaCache.InsertEntry(context.Background(), &filer.Entry{ - FullPath: "/dir/file", - Attr: filer.Attr{ - Crtime: time.Unix(1, 0), - Mtime: time.Unix(1, 0), - Mode: 0100644, - FileSize: 100, - }, - }); err != nil { - t.Fatalf("insert listing entry: %v", err) - } - // Postdates the listing snapshot, so it is buffered; its immediate - // invalidation resolves to the older listing entry. + // Covered by the upcoming listing snapshot (TsNs 900 <= snapshot 1000); + // its immediate invalidation runs before the listing inserts the newer + // entry, so the handle picks up the event's state. event := &filer_pb.SubscribeMetadataResponse{ Directory: "/dir", - TsNs: 2000, + TsNs: 900, EventNotification: &filer_pb.EventNotification{ OldEntry: &filer_pb.Entry{Name: "file"}, NewEntry: &filer_pb.Entry{ Name: "file", - Attributes: &filer_pb.FuseAttributes{FileSize: 300}, + Attributes: &filer_pb.FuseAttributes{FileSize: 100}, }, NewParentPath: "/dir", }, @@ -230,7 +220,20 @@ func TestBufferedBuildEventReinvalidatesOnCompletion(t *testing.T) { } wfs.metaCache.WaitForEntryInvalidations() if size := fh.GetEntry().GetEntry().Attributes.FileSize; size != 100 { - t.Fatalf("open handle file size mid-build = %d, want 100 (listing snapshot)", size) + t.Fatalf("open handle file size mid-build = %d, want 100 (event state)", size) + } + + // The listing then inserts the newer entry the snapshot already covers. + if err := wfs.metaCache.InsertEntry(context.Background(), &filer.Entry{ + FullPath: "/dir/file", + Attr: filer.Attr{ + Crtime: time.Unix(1, 0), + Mtime: time.Unix(1, 0), + Mode: 0100644, + FileSize: 300, + }, + }); err != nil { + t.Fatalf("insert listing entry: %v", err) } if err := wfs.metaCache.CompleteDirectoryBuild(context.Background(), util.FullPath("/dir"), 1000); err != nil { @@ -238,6 +241,101 @@ func TestBufferedBuildEventReinvalidatesOnCompletion(t *testing.T) { } wfs.metaCache.WaitForEntryInvalidations() if size := fh.GetEntry().GetEntry().Attributes.FileSize; size != 300 { - t.Fatalf("open handle file size after build completion = %d, want 300 (buffered event must re-invalidate)", size) + t.Fatalf("open handle file size after build completion = %d, want 300 (snapshot-covered event must re-invalidate)", size) + } +} + +// A hit in the local store only resolves an invalidation when the parent +// directory is cached. An uncached parent receives no store writes, so a +// leftover entry there is stale and must not mask the event. +func TestUncachedDirStaleStoreEntryDoesNotMaskEvent(t *testing.T) { + wfs := newInvalidateTestWFS(t) + + inode := wfs.inodeToPath.Lookup(util.FullPath("/dir/file"), time.Now().Unix(), false, false, 0, false) + fh := wfs.fhMap.AcquireFileHandle(wfs, inode, &filer_pb.Entry{ + Name: "file", + Attributes: &filer_pb.FuseAttributes{FileSize: 88}, + }) + + // Leftover store entry under a parent that is not children-cached. + if err := wfs.metaCache.InsertEntry(context.Background(), &filer.Entry{ + FullPath: "/dir/file", + Attr: filer.Attr{ + Crtime: time.Unix(1, 0), + Mtime: time.Unix(1, 0), + Mode: 0100644, + FileSize: 88, + }, + }); err != nil { + t.Fatalf("insert stale entry: %v", err) + } + + updateResp := &filer_pb.SubscribeMetadataResponse{ + Directory: "/dir", + TsNs: 1000, + EventNotification: &filer_pb.EventNotification{ + OldEntry: &filer_pb.Entry{Name: "file"}, + NewEntry: &filer_pb.Entry{ + Name: "file", + Attributes: &filer_pb.FuseAttributes{FileSize: 180020}, + }, + NewParentPath: "/dir", + }, + } + if err := wfs.metaCache.ApplyMetadataResponse(context.Background(), updateResp, meta_cache.SubscriberMetadataResponseApplyOptions); err != nil { + t.Fatalf("apply update event: %v", err) + } + wfs.metaCache.WaitForEntryInvalidations() + + if size := fh.GetEntry().GetEntry().Attributes.FileSize; size != 180020 { + t.Fatalf("open handle file size = %d, want 180020 (stale store entry must not mask the event)", size) + } +} + +// In a read-through directory neither a local flush nor the event reaches the +// local store, so ordering falls to the filer log timestamps: an event at or +// before the handle's last filer-acknowledged local mutation is old news and +// must not roll the handle back. +func TestQueuedEventOlderThanFlushedStateIsIgnored(t *testing.T) { + wfs := newInvalidateTestWFS(t) + + inode := wfs.inodeToPath.Lookup(util.FullPath("/dir/file"), time.Now().Unix(), false, false, 0, false) + fh := wfs.fhMap.AcquireFileHandle(wfs, inode, &filer_pb.Entry{ + Name: "file", + Attributes: &filer_pb.FuseAttributes{FileSize: 88}, + }) + + // Hold the handle lock so the queued invalidation cannot apply yet. + testLock := wfs.fhLockTable.AcquireLock("test", fh.fh, util.ExclusiveLock) + older := &filer_pb.SubscribeMetadataResponse{ + Directory: "/dir", + TsNs: 1000, + EventNotification: &filer_pb.EventNotification{ + OldEntry: &filer_pb.Entry{Name: "file"}, + NewEntry: &filer_pb.Entry{ + Name: "file", + Attributes: &filer_pb.FuseAttributes{FileSize: 100}, + }, + NewParentPath: "/dir", + }, + } + if err := wfs.metaCache.ApplyMetadataResponse(context.Background(), older, meta_cache.SubscriberMetadataResponseApplyOptions); err != nil { + wfs.fhLockTable.ReleaseLock(fh.fh, testLock) + t.Fatalf("apply subscriber event: %v", err) + } + + // A local flush lands: the filer acknowledged it with a later log + // timestamp than the queued event. + fh.SetEntry(&filer_pb.Entry{ + Name: "file", + Attributes: &filer_pb.FuseAttributes{FileSize: 200}, + }) + fh.advanceLocalEntryTs(2000) + wfs.fhLockTable.ReleaseLock(fh.fh, testLock) + + wfs.metaCache.WaitForEntryInvalidations() + + if size := fh.GetEntry().GetEntry().Attributes.FileSize; size != 200 { + t.Fatalf("open handle file size = %d, want 200 (event at TsNs 1000 predates the flush at 2000)", size) } } diff --git a/weed/mount/weedfs_metadata_flush.go b/weed/mount/weedfs_metadata_flush.go index 621a60c45..b18f40e96 100644 --- a/weed/mount/weedfs_metadata_flush.go +++ b/weed/mount/weedfs_metadata_flush.go @@ -180,6 +180,7 @@ func (wfs *WFS) flushFileMetadata(fh *FileHandle) error { if event == nil { event = metadataUpdateEvent(string(dir), request.Entry) } + fh.advanceLocalEntryTs(event.GetTsNs()) if applyErr := wfs.applyLocalMetadataEvent(context.Background(), event); applyErr != nil { glog.Warningf("flushFileMetadata %s: best-effort metadata apply failed: %v", fileFullPath, applyErr) wfs.inodeToPath.InvalidateChildrenCache(util.FullPath(dir)) diff --git a/weed/mount/weedfs_rename_test.go b/weed/mount/weedfs_rename_test.go index 75b1e19eb..50ae0df73 100644 --- a/weed/mount/weedfs_rename_test.go +++ b/weed/mount/weedfs_rename_test.go @@ -30,7 +30,7 @@ func TestHandleRenameResponseLeavesUncachedTargetOutOfCache(t *testing.T) { func(path util.FullPath) bool { return inodeToPath.IsChildrenCached(path) }, - func(util.FullPath, *filer_pb.Entry) {}, + func(util.FullPath, *filer_pb.Entry, int64) {}, nil, ) defer mc.Shutdown()