From 36b25fca7633b78bb2a267e5e679363c74e10f28 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Wed, 22 Jul 2026 15:22:30 -0700 Subject: [PATCH] mount: advance the handle watermark for every acknowledged local mutation MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Flush was the only path stamping the watermark, so a truncate or setattr through saveEntry, a server-side CopyFileRange, or a remote cache download whose response carried no metadata event could still be rolled back by an older queued subscription event in a read-through directory. saveEntry now stamps any open handle for the saved path, and every site tolerates a response without an event: the metadata cache tracks the newest filer log timestamp seen in any applied event, and that position, captured before the RPC, is a safe lower bound for the state the filer served — anything it returns is at least as new as every event already delivered. --- weed/mount/filehandle.go | 13 ++ weed/mount/filehandle_read.go | 3 +- weed/mount/meta_cache/meta_cache.go | 23 +++ weed/mount/weedfs.go | 10 ++ weed/mount/weedfs_file_copy_range.go | 2 + weed/mount/weedfs_file_sync.go | 3 +- .../weedfs_invalidate_open_handle_test.go | 151 ++++++++++++++++++ weed/mount/weedfs_metadata_flush.go | 3 +- weed/mount/wfs_save.go | 9 ++ 9 files changed, 214 insertions(+), 3 deletions(-) diff --git a/weed/mount/filehandle.go b/weed/mount/filehandle.go index 545164089..a199fdbf6 100644 --- a/weed/mount/filehandle.go +++ b/weed/mount/filehandle.go @@ -127,6 +127,19 @@ func (fh *FileHandle) SetEntry(entry *filer_pb.Entry) { fh.invalidateChunkCache() } +// noteFilerAck advances the watermark after a filer RPC acknowledged state +// now reflected in this handle: to the response event's log timestamp when +// one was returned, otherwise to baselineTsNs — the latest filer log position +// known before the RPC was issued, a safe lower bound for the state the +// filer served. +func (fh *FileHandle) noteFilerAck(baselineTsNs int64, event *filer_pb.SubscribeMetadataResponse) { + if tsNs := event.GetTsNs(); tsNs != 0 { + fh.advanceLocalEntryTs(tsNs) + return + } + fh.advanceLocalEntryTs(baselineTsNs) +} + // 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. diff --git a/weed/mount/filehandle_read.go b/weed/mount/filehandle_read.go index 3dda56962..1d40d0ef3 100644 --- a/weed/mount/filehandle_read.go +++ b/weed/mount/filehandle_read.go @@ -188,19 +188,20 @@ func (fh *FileHandle) downloadRemoteEntry(entry *LockedEntry) error { } glog.V(4).Infof("download entry: %v", request) + baselineTsNs := fh.wfs.latestKnownFilerTsNs() resp, err := client.CacheRemoteObjectToLocalCluster(context.Background(), request) if err != nil { return fmt.Errorf("CacheRemoteObjectToLocalCluster file %s: %v", fileFullPath, err) } fh.SetEntry(resp.Entry) + fh.noteFilerAck(baselineTsNs, resp.GetMetadataEvent()) // Async: a sync apply deadlocks against the apply loop's invalidate, which needs this read's file-handle lock. event := resp.GetMetadataEvent() if event == nil { event = metadataUpdateEvent(request.Directory, resp.Entry) } - fh.advanceLocalEntryTs(event.GetTsNs()) fh.wfs.applyLocalMetadataEventAsync(event) return nil diff --git a/weed/mount/meta_cache/meta_cache.go b/weed/mount/meta_cache/meta_cache.go index ff3a17404..16a91505a 100644 --- a/weed/mount/meta_cache/meta_cache.go +++ b/weed/mount/meta_cache/meta_cache.go @@ -6,6 +6,7 @@ import ( "math" "os" "sync" + "sync/atomic" "time" "golang.org/x/sync/singleflight" @@ -42,6 +43,7 @@ type MetaCache struct { buildingDirs map[util.FullPath]*directoryBuildState dedupRing dedupRingBuffer includeSystemEntries bool + latestEventTsNs atomic.Int64 // newest filer log timestamp seen in any applied event // Entry invalidations run on a worker, not inline on the apply loop: // invalidateFunc takes the fh lock, which a flush can hold while waiting on @@ -574,6 +576,7 @@ type metadataResponseSideEffects struct { } func (mc *MetaCache) applyMetadataResponseNow(ctx context.Context, resp *filer_pb.SubscribeMetadataResponse, options MetadataResponseApplyOptions) error { + mc.advanceLatestEventTs(resp.TsNs) if mc.shouldSkipDuplicateEvent(resp) { return nil } @@ -651,6 +654,26 @@ func (mc *MetaCache) WaitForEntryInvalidations() { mc.invalidateWorker.Drain() } +func (mc *MetaCache) advanceLatestEventTs(tsNs int64) { + if tsNs == 0 { + return + } + for { + current := mc.latestEventTsNs.Load() + if tsNs <= current || mc.latestEventTsNs.CompareAndSwap(current, tsNs) { + return + } + } +} + +// LatestEventTsNs returns the newest filer log timestamp seen in any applied +// event. A filer RPC issued after reading this serves state at least this +// new, so it is a safe watermark baseline when the RPC response carries no +// metadata event of its own. +func (mc *MetaCache) LatestEventTsNs() int64 { + return mc.latestEventTsNs.Load() +} + func (mc *MetaCache) applyMetadataResponseLocked(ctx context.Context, resp *filer_pb.SubscribeMetadataResponse, _ MetadataResponseApplyOptions, allowUncachedInsert bool) (metadataResponseSideEffects, error) { message := resp.GetEventNotification() if message == nil { diff --git a/weed/mount/weedfs.go b/weed/mount/weedfs.go index 18c16abe6..a05e845c2 100644 --- a/weed/mount/weedfs.go +++ b/weed/mount/weedfs.go @@ -676,6 +676,16 @@ func (wfs *WFS) lookupEntry(fullpath util.FullPath) (*filer.Entry, fuse.Status) return filer.FromPbEntry(dir, entry), fuse.OK } +// latestKnownFilerTsNs is the newest filer log timestamp seen in any applied +// event — the watermark baseline for filer RPCs whose response carries no +// metadata event. Nil-safe for partially constructed test instances. +func (wfs *WFS) latestKnownFilerTsNs() int64 { + if wfs.metaCache == nil { + return 0 + } + return wfs.metaCache.LatestEventTsNs() +} + // invalidateOpenFileHandle refreshes an open file handle from a metadata // subscription event. No filer lookup happens here: it can fail transiently, // and since the subscription cursor has already advanced past the event, the diff --git a/weed/mount/weedfs_file_copy_range.go b/weed/mount/weedfs_file_copy_range.go index eb9ca4282..1656d0e4f 100644 --- a/weed/mount/weedfs_file_copy_range.go +++ b/weed/mount/weedfs_file_copy_range.go @@ -216,6 +216,7 @@ func (wfs *WFS) tryServerSideWholeFileCopy(cancel <-chan struct{}, in *fuse.Copy glog.V(1).Infof("CopyFileRange server-side copy %s => %s (%d bytes)", copyRequest.srcPath, copyRequest.dstPath, copyRequest.sourceSize) + baselineTsNs := wfs.latestKnownFilerTsNs() entry, outcome, err := performServerSideWholeFileCopy(cancel, wfs, copyRequest) switch outcome { case serverSideWholeFileCopyCommitted: @@ -224,6 +225,7 @@ func (wfs *WFS) tryServerSideWholeFileCopy(cancel <-chan struct{}, in *fuse.Copy } else { glog.V(1).Infof("CopyFileRange server-side copy %s => %s completed (%d bytes)", copyRequest.srcPath, copyRequest.dstPath, copyRequest.sourceSize) } + fhOut.advanceLocalEntryTs(baselineTsNs) wfs.applyServerSideWholeFileCopyResult(fhIn, fhOut, copyRequest.dstPath, entry, copyRequest.sourceSize) return uint32(copyRequest.sourceSize), true, fuse.OK case serverSideWholeFileCopyAmbiguous: diff --git a/weed/mount/weedfs_file_sync.go b/weed/mount/weedfs_file_sync.go index 6ba8fb496..bc5a2953a 100644 --- a/weed/mount/weedfs_file_sync.go +++ b/weed/mount/weedfs_file_sync.go @@ -260,17 +260,18 @@ func (wfs *WFS) flushMetadataToFiler(ctx context.Context, fh *FileHandle, dir, n wfs.mapPbIdFromLocalToFiler(request.Entry) + baselineTsNs := wfs.latestKnownFilerTsNs() resp, err := wfs.streamCreateEntry(ctx, request) if err != nil { glog.Errorf("fh flush create %s: %v", fileFullPath, err) return fmt.Errorf("fh flush create %s: %v", fileFullPath, err) } + fh.noteFilerAck(baselineTsNs, resp.GetMetadataEvent()) event := resp.GetMetadataEvent() 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 20a278ff0..30199f54c 100644 --- a/weed/mount/weedfs_invalidate_open_handle_test.go +++ b/weed/mount/weedfs_invalidate_open_handle_test.go @@ -2,10 +2,13 @@ package mount import ( "context" + "net" "path/filepath" "testing" "time" + "github.com/seaweedfs/go-fuse/v2/fuse" + "github.com/seaweedfs/seaweedfs/weed/filer" "github.com/seaweedfs/seaweedfs/weed/mount/meta_cache" "github.com/seaweedfs/seaweedfs/weed/pb" @@ -339,3 +342,151 @@ func TestQueuedEventOlderThanFlushedStateIsIgnored(t *testing.T) { t.Fatalf("open handle file size = %d, want 200 (event at TsNs 1000 predates the flush at 2000)", size) } } + +type saveEntryTestServer struct { + filer_pb.UnimplementedSeaweedFilerServer +} + +func (s *saveEntryTestServer) UpdateEntry(ctx context.Context, req *filer_pb.UpdateEntryRequest) (*filer_pb.UpdateEntryResponse, error) { + return &filer_pb.UpdateEntryResponse{ + MetadataEvent: &filer_pb.SubscribeMetadataResponse{ + Directory: req.Directory, + TsNs: 2000, + EventNotification: &filer_pb.EventNotification{ + OldEntry: &filer_pb.Entry{Name: req.Entry.Name}, + NewEntry: req.Entry, + NewParentPath: req.Directory, + }, + }, + }, nil +} + +// saveEntry (truncate, setattr) must advance the open handle's watermark from +// the acknowledged mutation's log timestamp, or an older queued event rolls +// the mutation back in a read-through directory. +func TestSaveEntryKeepsOpenHandleAheadOfOlderEvents(t *testing.T) { + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + t.Cleanup(func() { _ = listener.Close() }) + server := pb.NewGrpcServer() + filer_pb.RegisterSeaweedFilerServer(server, &saveEntryTestServer{}) + go server.Serve(listener) + t.Cleanup(server.Stop) + + wfs := newInvalidateTestWFS(t) + wfs.option.FilerAddresses = []pb.ServerAddress{ + pb.NewServerAddressWithGrpcPort("127.0.0.1:1", listener.Addr().(*net.TCPAddr).Port), + } + + 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}, + }) + + 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 truncate-style mutation: the filer acknowledges it at TsNs 2000 and + // the handle takes the new entry. + saved := &filer_pb.Entry{ + Name: "file", + Attributes: &filer_pb.FuseAttributes{FileSize: 200}, + } + if code := wfs.saveEntry(util.FullPath("/dir/file"), saved); code != fuse.OK { + wfs.fhLockTable.ReleaseLock(fh.fh, testLock) + t.Fatalf("saveEntry status = %v, want OK", code) + } + fh.SetEntry(saved) + 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 (saveEntry at TsNs 2000 outranks the queued event at 1000)", size) + } +} + +// When a filer response carries no metadata event, the watermark falls back +// to the latest filer log timestamp already seen: the filer served state at +// least that new, so queued events at or before it must not roll the fresh +// handle state back. +func TestNilAckEventFallsBackToLatestSeenTs(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}, + }) + + // An unrelated event advances the mount's known filer log position. + unrelated := &filer_pb.SubscribeMetadataResponse{ + Directory: "/dir", + TsNs: 1500, + EventNotification: &filer_pb.EventNotification{ + NewEntry: &filer_pb.Entry{ + Name: "other", + Attributes: &filer_pb.FuseAttributes{FileSize: 1}, + }, + }, + } + if err := wfs.metaCache.ApplyMetadataResponse(context.Background(), unrelated, meta_cache.SubscriberMetadataResponseApplyOptions); err != nil { + t.Fatalf("apply unrelated event: %v", err) + } + if got := wfs.metaCache.LatestEventTsNs(); got != 1500 { + t.Fatalf("LatestEventTsNs = %d, want 1500", got) + } + + 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 filer RPC that returned no event (remote cache already populated, + // server-side copy) installs fresh state; the baseline captured before + // the RPC stands in for the missing event timestamp. + baselineTsNs := wfs.metaCache.LatestEventTsNs() + fh.SetEntry(&filer_pb.Entry{ + Name: "file", + Attributes: &filer_pb.FuseAttributes{FileSize: 200}, + }) + fh.noteFilerAck(baselineTsNs, nil) + 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 (queued event at TsNs 1000 predates the known log position 1500)", size) + } +} diff --git a/weed/mount/weedfs_metadata_flush.go b/weed/mount/weedfs_metadata_flush.go index b18f40e96..fc828a42d 100644 --- a/weed/mount/weedfs_metadata_flush.go +++ b/weed/mount/weedfs_metadata_flush.go @@ -171,16 +171,17 @@ func (wfs *WFS) flushFileMetadata(fh *FileHandle) error { wfs.mapPbIdFromLocalToFiler(request.Entry) + baselineTsNs := wfs.latestKnownFilerTsNs() resp, err := wfs.streamCreateEntry(context.Background(), request) if err != nil { return err } + fh.noteFilerAck(baselineTsNs, resp.GetMetadataEvent()) event := resp.GetMetadataEvent() 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/wfs_save.go b/weed/mount/wfs_save.go index 24f692ae2..a8f75f5fe 100644 --- a/weed/mount/wfs_save.go +++ b/weed/mount/wfs_save.go @@ -27,6 +27,7 @@ func (wfs *WFS) saveEntry(path util.FullPath, entry *filer_pb.Entry) (code fuse. glog.V(1).Infof("save entry: %v", request) + baselineTsNs := wfs.latestKnownFilerTsNs() var resp *filer_pb.UpdateEntryResponse err := retryMetadataFlushIf(context.Background(), func() error { var callErr error @@ -51,6 +52,14 @@ func (wfs *WFS) saveEntry(path util.FullPath, entry *filer_pb.Entry) (code fuse. return fuseStatus } + // The mutation is acknowledged; keep any open handle for this path ahead + // of subscription events the filer logged before it. + if inode, found := wfs.inodeToPath.GetInode(path); found { + if fh, fhFound := wfs.fhMap.FindFileHandle(inode); fhFound { + fh.noteFilerAck(baselineTsNs, resp.GetMetadataEvent()) + } + } + event := resp.GetMetadataEvent() if event == nil { event = metadataUpdateEvent(parentDir, entry)