diff --git a/weed/filer/filer.go b/weed/filer/filer.go index 39b7f75b0..0eab4ef7e 100644 --- a/weed/filer/filer.go +++ b/weed/filer/filer.go @@ -7,6 +7,7 @@ import ( "os" "sort" "strings" + "sync/atomic" "time" "github.com/seaweedfs/seaweedfs/weed/remote_storage" @@ -73,6 +74,11 @@ type Filer struct { EmptyFolderCleanupDelay time.Duration persistedLogCache *persistedLogCache metaLogInflight metaLogInflight + remoteTombstones *remoteDeletionTombstones + // remoteTombstonesDone, when non-nil, is closed once the startup tombstone + // rebuild finishes; lazy remote reads wait on it so a pending delete + // cannot resurrect in the gap. + remoteTombstonesDone atomic.Pointer[chan struct{}] } func NewFiler(masters pb.ServerDiscovery, grpcDialOption grpc.DialOption, filerHost pb.ServerAddress, filerGroup string, collection string, replication string, dataCenter string, maxFilenameLength uint32, notifyFn func()) *Filer { @@ -88,6 +94,7 @@ func NewFiler(masters pb.ServerDiscovery, grpcDialOption grpc.DialOption, filerH deletionQuit: make(chan struct{}), DeletionRetryQueue: NewDeletionRetryQueue(), persistedLogCache: newPersistedLogCache(persistedLogCacheMaxBytes), + remoteTombstones: newRemoteDeletionTombstones(), } if f.UniqueFilerId < 0 { f.UniqueFilerId = -f.UniqueFilerId diff --git a/weed/filer/filer_delete_entry.go b/weed/filer/filer_delete_entry.go index 6c86b2a4f..245f9381c 100644 --- a/weed/filer/filer_delete_entry.go +++ b/weed/filer/filer_delete_entry.go @@ -125,6 +125,15 @@ func (f *Filer) DeleteEntryMetaAndData(ctx context.Context, p util.FullPath, isR func (f *Filer) doBatchDeleteFolderMetaAndData(ctx context.Context, entry *Entry, isRecursive, ignoreRecursiveError, shouldDeleteChunks, isDeletingBucket, isFromOtherCluster bool, signatures []int32, onHardLinkIdsFn OnHardLinkIdsFunc) (err error) { + var dirTombstoneTs int64 + if isRecursive { + // Tombstone the directory before its children: when the store drops + // the subtree without listing it, or a child error aborts the walk, + // the ancestor tombstone still covers every descendant. + dirTombstoneTs = time.Now().UnixNano() + f.noteRemoteDeletion(entry.FullPath, true, dirTombstoneTs) + } + //collect all the chunks of this layer and delete them together at the end var chunksToDelete []*filer_pb.FileChunk lastFileName := "" @@ -134,6 +143,9 @@ func (f *Filer) doBatchDeleteFolderMetaAndData(ctx context.Context, entry *Entry for { entries, _, err := f.ListDirectoryEntries(ctx, entry.FullPath, lastFileName, includeLastFile, PaginationSize, "", "", "") if err != nil { + // nothing was deleted; a leftover tombstone would hide the + // still-existing remote children + f.unnoteRemoteDeletion(entry.FullPath, dirTombstoneTs) glog.ErrorfCtx(ctx, "list folder %s: %v", entry.FullPath, err) return fmt.Errorf("list folder %s: %v", entry.FullPath, err) } @@ -145,6 +157,7 @@ func (f *Filer) doBatchDeleteFolderMetaAndData(ctx context.Context, entry *Entry for _, sub := range entries { lastFileName = sub.Name() + f.noteRemoteDeletion(sub.FullPath, sub.IsDirectory(), time.Now().UnixNano()) if sub.IsDirectory() { subIsDeletingBucket := f.IsBucket(sub) err = f.doBatchDeleteFolderMetaAndData(ctx, sub, isRecursive, ignoreRecursiveError, shouldDeleteChunks, subIsDeletingBucket, isFromOtherCluster, nil, onHardLinkIdsFn) @@ -207,6 +220,8 @@ func (f *Filer) doDeleteEntryMetaAndData(ctx context.Context, entry *Entry, shou } } + f.noteRemoteDeletion(entry.FullPath, entry.IsDirectory(), time.Now().UnixNano()) + if storeDeletionErr := f.Store.DeleteOneEntry(ctx, entry); storeDeletionErr != nil { return fmt.Errorf("filer store delete: %w", storeDeletionErr) } diff --git a/weed/filer/filer_lazy_remote.go b/weed/filer/filer_lazy_remote.go index 1b6983be3..7c642cd6a 100644 --- a/weed/filer/filer_lazy_remote.go +++ b/weed/filer/filer_lazy_remote.go @@ -7,7 +7,10 @@ import ( "strings" "time" + "google.golang.org/protobuf/proto" + "github.com/seaweedfs/seaweedfs/weed/glog" + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" "github.com/seaweedfs/seaweedfs/weed/pb/remote_pb" "github.com/seaweedfs/seaweedfs/weed/remote_storage" "github.com/seaweedfs/seaweedfs/weed/util" @@ -43,6 +46,21 @@ func (f *Filer) maybeLazyFetchFromRemote(ctx context.Context, p util.FullPath) ( return nil, nil } + // A startup tombstone rebuild may still be replaying the meta log; wait + // for it so a pending delete cannot resurrect here. + if done := f.remoteTombstonesDone.Load(); done != nil { + select { + case <-*done: + case <-ctx.Done(): + return nil, ctx.Err() + } + } + + if f.isRemoteDeletionPending(ctx, p, mountDir) { + glog.V(2).InfofCtx(ctx, "maybeLazyFetchFromRemote: %s deleted locally, remote delete pending", p) + return nil, nil + } + remoteConf, found := f.RemoteStorage.FindRemoteStorageConf(p) if !found { return nil, nil @@ -103,10 +121,27 @@ func (f *Filer) maybeLazyFetchFromRemote(ctx context.Context, p util.FullPath) ( persistBaseCtx, cancelPersist := context.WithTimeout(context.Background(), 30*time.Second) defer cancelPersist() persistCtx := context.WithValue(persistBaseCtx, lazyFetchContextKey{}, true) + // A delete may have landed while StatFile was in flight; re-check so + // the fetched object cannot resurrect a path whose delete is pending. + if f.isRemoteDeletionPending(persistCtx, p, mountDir) { + glog.V(2).InfofCtx(ctx, "maybeLazyFetchFromRemote: %s deleted during remote stat", p) + return lazyFetchResult{nil}, nil + } saveErr := f.CreateEntry(persistCtx, entry, nil, false, false, nil, true, f.MaxFilenameLength) if saveErr != nil { glog.Warningf("maybeLazyFetchFromRemote: failed to persist filer entry for %s: %v", p, saveErr) f.lazyFetchGroup.Forget(key) + return lazyFetchResult{entry}, nil + } + + // A delete records its tombstone before removing the entry, so a + // tombstone visible now means the insert raced a delete that already + // ran: retract the persisted entry so the path stays deleted. + if f.isRemoteDeletionPending(persistCtx, p, mountDir) { + glog.V(2).InfofCtx(ctx, "maybeLazyFetchFromRemote: %s deleted while persisting", p) + f.lazyFetchGroup.Forget(key) + f.retractLazyRemoteEntry(persistCtx, entry) + return lazyFetchResult{nil}, nil } return lazyFetchResult{entry}, nil @@ -122,6 +157,25 @@ func (f *Filer) maybeLazyFetchFromRemote(ctx context.Context, p util.FullPath) ( return result.entry, nil } +// retractLazyRemoteEntry deletes the entry at entry.FullPath only when it is +// still the entry a lazy remote read just materialized — a concurrent write +// may have replaced it, and deleting by path alone would take that write down. +func (f *Filer) retractLazyRemoteEntry(ctx context.Context, entry *Entry) { + existing, findErr := f.FindEntry(ctx, entry.FullPath) + if findErr != nil || existing == nil { + return + } + // The stored entry must still be exactly what this read materialized — + // an intervening write (appended chunks, touched attributes) means a + // real update owns the path now. + if !proto.Equal(existing.ToProtoEntry(), entry.ToProtoEntry()) { + return + } + if err := f.doDeleteEntryMetaAndData(ctx, existing, false, false, nil); err != nil && !errors.Is(err, filer_pb.ErrNotFound) { + glog.Warningf("retractLazyRemoteEntry %s: %v", entry.FullPath, err) + } +} + func (f *Filer) maybeDeleteFromRemote(ctx context.Context, entry *Entry) (bool, error) { if entry == nil || f.RemoteStorage == nil { return false, nil diff --git a/weed/filer/filer_lazy_remote_listing.go b/weed/filer/filer_lazy_remote_listing.go index e292c208e..d692908a8 100644 --- a/weed/filer/filer_lazy_remote_listing.go +++ b/weed/filer/filer_lazy_remote_listing.go @@ -56,6 +56,14 @@ func (f *Filer) maybeLazyListFromRemote(ctx context.Context, p util.FullPath) { } } + if done := f.remoteTombstonesDone.Load(); done != nil { + select { + case <-*done: + case <-ctx.Done(): + return + } + } + // Lazy listing is opt-in: disabled when TTL is 0 if remoteLoc.ListingCacheTtlSeconds <= 0 { return @@ -109,6 +117,10 @@ func (f *Filer) maybeLazyListFromRemote(ctx context.Context, p util.FullPath) { return nil } + if existingEntry == nil && f.isRemoteDeletionPending(persistCtx, childPath, mountDir) { + return nil + } + if existingEntry != nil { // Merge: update remote metadata while preserving local state // (Chunks, Extended, Uid/Gid/Mode, etc.) @@ -161,6 +173,9 @@ func (f *Filer) maybeLazyListFromRemote(ctx context.Context, p util.FullPath) { } if saveErr := f.CreateEntry(persistCtx, entry, nil, false, false, nil, true, f.MaxFilenameLength); saveErr != nil { glog.Warningf("maybeLazyListFromRemote: persist %s: %v", childPath, saveErr) + } else if f.isRemoteDeletionPending(persistCtx, childPath, mountDir) { + // a delete landed between the check above and the insert + f.retractLazyRemoteEntry(persistCtx, entry) } } return nil diff --git a/weed/filer/filer_lazy_remote_test.go b/weed/filer/filer_lazy_remote_test.go index 15cfe8c40..468142613 100644 --- a/weed/filer/filer_lazy_remote_test.go +++ b/weed/filer/filer_lazy_remote_test.go @@ -283,6 +283,7 @@ func newTestFiler(t *testing.T, store *stubFilerStore, rs *FilerRemoteStorage) * MasterClient: mc, FileIdDeletionQueue: util.NewUnboundedQueue(), deletionQuit: make(chan struct{}), + remoteTombstones: newRemoteDeletionTombstones(), LocalMetaLogBuffer: log_buffer.NewLogBuffer("test", time.Minute, func(*log_buffer.LogBuffer, time.Time, time.Time, []byte, int64, int64) {}, nil, func() {}), } diff --git a/weed/filer/filer_on_meta_event.go b/weed/filer/filer_on_meta_event.go index 264240548..33acc5b2d 100644 --- a/weed/filer/filer_on_meta_event.go +++ b/weed/filer/filer_on_meta_event.go @@ -15,6 +15,7 @@ func (f *Filer) onMetadataChangeEvent(event *filer_pb.SubscribeMetadataResponse) f.maybeReloadRemoteStorageConfigurationAndMapping(event) f.onBucketEvents(event) f.onEmptyFolderCleanupEvents(event) + f.onRemoteDeletionEvents(event) } func (f *Filer) onBucketEvents(event *filer_pb.SubscribeMetadataResponse) { diff --git a/weed/filer/filer_remote_tombstone.go b/weed/filer/filer_remote_tombstone.go new file mode 100644 index 000000000..bf7b4b187 --- /dev/null +++ b/weed/filer/filer_remote_tombstone.go @@ -0,0 +1,423 @@ +package filer + +import ( + "context" + "math" + "strings" + "sync" + "time" + + "github.com/seaweedfs/seaweedfs/weed/glog" + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/remote_storage" + "github.com/seaweedfs/seaweedfs/weed/util" + "github.com/seaweedfs/seaweedfs/weed/util/log_buffer" + "google.golang.org/protobuf/proto" +) + +const ( + // remoteDeletionTombstoneTTL bounds a tombstone when no write-back sync + // offset ever confirms the remote delete. A remote object re-created + // outside the filer at a deleted path stays hidden for this long. + remoteDeletionTombstoneTTL = 24 * time.Hour + // remoteDeletionTombstoneLimit bounds tracked paths; past it new + // tombstones are dropped after an expired sweep still leaves no room. + remoteDeletionTombstoneLimit = 1 << 16 +) + +// remoteDeletionTombstones tracks paths deleted under a remote mount whose +// remote object may still exist because the write-back daemon has not +// consumed the delete event yet. A lazy remote fetch or listing must not +// resurrect them. A tombstone lifts when the path is written again, when the +// mount's persisted sync offset passes the delete event, or on TTL. +// +// Tombstones recorded before their metadata event lands (the synchronous +// delete path) are marked pending: the sync offset orders against event +// timestamps, which only the event itself knows, so a pending tombstone can +// only be lifted by the event confirming it or by TTL. Once the event stamps +// the real timestamp the tombstone is releasable by the offset. +type remoteDeletionTombstones struct { + mu sync.Mutex + files map[string]int64 // file path -> delete event TsNs + dirs map[string]int64 // deleted directory path -> event TsNs; covers its subtree + pending map[string]bool // tombstone path recorded ahead of its event +} + +func newRemoteDeletionTombstones() *remoteDeletionTombstones { + return &remoteDeletionTombstones{ + files: make(map[string]int64), + dirs: make(map[string]int64), + pending: make(map[string]bool), + } +} + +// add records a tombstone ahead of its metadata event — the timestamp is the +// local delete time, a lower bound the sync offset cannot order against. +func (t *remoteDeletionTombstones) add(path string, isDir bool, tsNs int64) { + t.upsert(path, isDir, tsNs, false) +} + +// addFromEvent records a tombstone stamped by the delete event itself, so the +// mount's sync offset can release it once the daemon passes that event. +func (t *remoteDeletionTombstones) addFromEvent(path string, isDir bool, tsNs int64) { + t.upsert(path, isDir, tsNs, true) +} + +func (t *remoteDeletionTombstones) upsert(path string, isDir bool, tsNs int64, fromEvent bool) { + if t == nil { + return + } + t.mu.Lock() + defer t.mu.Unlock() + // An ancestor directory tombstone at least as new already covers the + // path; recording it again only spends capacity. + for p := path; ; { + i := strings.LastIndexByte(p, '/') + if i <= 0 { + break + } + p = p[:i] + if ancestorTs, ok := t.dirs[p]; ok && ancestorTs >= tsNs { + return + } + } + m := t.files + if isDir { + m = t.dirs + } + if cur, ok := m[path]; ok { + if tsNs > cur { + m[path] = tsNs + if !fromEvent { + // a newer local delete restamps the tombstone ahead of its + // event — the offset cannot vouch for it until the event lands + t.pending[path] = true + } + } + if fromEvent && m[path] == tsNs { + delete(t.pending, path) + } + return + } + if len(t.files)+len(t.dirs) >= remoteDeletionTombstoneLimit { + t.evictExpiredLocked(time.Now().UnixNano()) + if len(t.files)+len(t.dirs) >= remoteDeletionTombstoneLimit { + glog.V(0).Infof("remote deletion tombstones full (%d), skipping %s", remoteDeletionTombstoneLimit, path) + return + } + } + m[path] = tsNs + if fromEvent { + delete(t.pending, path) + } else { + t.pending[path] = true + } + if isDir { + t.dropCoveredLocked(path, tsNs) + } +} + +// dropCoveredLocked removes descendant tombstones a new directory tombstone +// subsumes: their deletes predate it, so the ancestor already hides those +// remote objects. Descendants deleted later keep their own tombstone. +// Caller must hold t.mu. +func (t *remoteDeletionTombstones) dropCoveredLocked(dirPath string, tsNs int64) { + prefix := dirPath + "/" + for p, ts := range t.files { + if ts <= tsNs && strings.HasPrefix(p, prefix) { + delete(t.files, p) + delete(t.pending, p) + } + } + for p, ts := range t.dirs { + if ts <= tsNs && strings.HasPrefix(p, prefix) { + delete(t.dirs, p) + delete(t.pending, p) + } + } +} + +// drop removes the exact tombstone recorded for path, e.g. when the delete +// that recorded it fails before touching anything. +func (t *remoteDeletionTombstones) drop(path string, tsNs int64) { + if t == nil { + return + } + t.mu.Lock() + defer t.mu.Unlock() + if cur, ok := t.dirs[path]; ok && cur <= tsNs { + delete(t.dirs, path) + delete(t.pending, path) + } + if cur, ok := t.files[path]; ok && cur <= tsNs { + delete(t.files, path) + delete(t.pending, path) + } +} + +// clear drops a file tombstone when a write at the path is at least as new as +// the delete; a replayed older create must not lift a newer delete. +func (t *remoteDeletionTombstones) clear(path string, tsNs int64) { + if t == nil { + return + } + t.mu.Lock() + defer t.mu.Unlock() + if cur, ok := t.files[path]; ok && tsNs >= cur { + delete(t.files, path) + delete(t.pending, path) + } +} + +// blockedSince returns the newest delete timestamp governing path — its own +// file tombstone or one from a deleted ancestor directory — and whether that +// tombstone is still waiting for its event. 0 means clear. +func (t *remoteDeletionTombstones) blockedSince(path string) (tsNs int64, pending bool) { + if t == nil { + return 0, false + } + t.mu.Lock() + defer t.mu.Unlock() + if ts, ok := t.files[path]; ok { + tsNs = ts + pending = t.pending[path] + } + for p := path; ; { + if ts, ok := t.dirs[p]; ok && ts > tsNs { + tsNs = ts + pending = t.pending[p] + } + i := strings.LastIndexByte(p, '/') + if i <= 0 { + break + } + p = p[:i] + } + return tsNs, pending +} + +// releaseThrough drops the tombstones governing path that are no newer than +// tsNs, once their remote deletes are confirmed consumed. +func (t *remoteDeletionTombstones) releaseThrough(path string, tsNs int64) { + if t == nil { + return + } + t.mu.Lock() + defer t.mu.Unlock() + if cur, ok := t.files[path]; ok && cur <= tsNs { + delete(t.files, path) + delete(t.pending, path) + } + for p := path; ; { + if cur, ok := t.dirs[p]; ok && cur <= tsNs { + delete(t.dirs, p) + delete(t.pending, p) + } + i := strings.LastIndexByte(p, '/') + if i <= 0 { + break + } + p = p[:i] + } +} + +func (t *remoteDeletionTombstones) evictExpiredLocked(nowNs int64) { + for p, ts := range t.files { + if nowNs-ts >= int64(remoteDeletionTombstoneTTL) { + delete(t.files, p) + delete(t.pending, p) + } + } + for p, ts := range t.dirs { + if nowNs-ts >= int64(remoteDeletionTombstoneTTL) { + delete(t.dirs, p) + delete(t.pending, p) + } + } +} + +// noteRemoteDeletion records a delete of path under a remote mount so lazy +// remote reads skip it until the remote delete is confirmed. +func (f *Filer) noteRemoteDeletion(p util.FullPath, isDir bool, tsNs int64) { + if f.RemoteStorage == nil || f.remoteTombstones == nil { + return + } + if _, remoteLoc := f.RemoteStorage.FindMountDirectory(p); remoteLoc == nil { + return + } + f.remoteTombstones.add(string(p), isDir, tsNs) +} + +// unnoteRemoteDeletion retracts a tombstone when the delete that recorded it +// fails before touching anything under path. +func (f *Filer) unnoteRemoteDeletion(p util.FullPath, tsNs int64) { + if f.remoteTombstones == nil { + return + } + f.remoteTombstones.drop(string(p), tsNs) +} + +// isRemoteDeletionPending reports whether a remote write-back delete for p is +// still owed: p was deleted under mountDir and neither a rewrite, the mount's +// sync offset, nor the TTL has lifted the tombstone. +func (f *Filer) isRemoteDeletionPending(ctx context.Context, p util.FullPath, mountDir util.FullPath) bool { + if f.remoteTombstones == nil { + return false + } + tsNs, pending := f.remoteTombstones.blockedSince(string(p)) + if tsNs == 0 { + return false + } + if pending { + // Recorded ahead of its delete event — the write-back offset cannot + // vouch for it yet; only the TTL lifts it. + if f.remoteDeletionExpired(tsNs) { + f.remoteTombstones.releaseThrough(string(p), tsNs) + return false + } + return true + } + if f.remoteDeletionConsumed(ctx, mountDir, tsNs) { + f.remoteTombstones.releaseThrough(string(p), tsNs) + return false + } + return true +} + +func (f *Filer) remoteDeletionExpired(tsNs int64) bool { + return time.Now().UnixNano()-tsNs >= int64(remoteDeletionTombstoneTTL) +} + +func (f *Filer) remoteDeletionConsumed(ctx context.Context, mountDir util.FullPath, tsNs int64) bool { + if f.remoteDeletionExpired(tsNs) { + return true + } + offset, err := f.readRemoteSyncOffset(ctx, mountDir) + return err == nil && offset >= tsNs +} + +// readRemoteSyncOffset reads the write-back daemon's persisted watermark for +// mountDir straight from the local store: every event at or below it has been +// applied to the remote. +func (f *Filer) readRemoteSyncOffset(ctx context.Context, mountDir util.FullPath) (int64, error) { + value, err := f.Store.KvGet(ctx, remote_storage.SyncOffsetKey(string(mountDir))) + if err != nil { + return 0, err + } + if len(value) < 8 { + return 0, nil + } + return int64(util.BytesToUint64(value)), nil +} + +// RebuildRemoteDeletionTombstones gates lazy remote reads and replays the +// persisted metadata log from the oldest write-back offset across mounts, +// restoring tombstones for deletes committed before a restart but not yet +// applied to the remote. Every filer writes its log files under the same +// directory, so pending peer deletes replay too; only events still inside +// the unflushed buffer tail are missed. The gate is set synchronously so no +// lazy read can slip in before replay starts, and it stays closed until a +// replay succeeds. +func (f *Filer) RebuildRemoteDeletionTombstones(ctx context.Context) { + if f.RemoteStorage == nil || f.remoteTombstones == nil { + return + } + mounts := f.RemoteStorage.MountedDirectories() + if len(mounts) == 0 { + return + } + done := make(chan struct{}) + f.remoteTombstonesDone.Store(&done) + go f.rebuildRemoteDeletionTombstones(ctx, mounts, done) +} + +func (f *Filer) rebuildRemoteDeletionTombstones(ctx context.Context, mounts []util.FullPath, done chan struct{}) { + // the replay itself lists directories; do not let it wait on its own gate + ctx = context.WithValue(ctx, lazyFetchContextKey{}, true) + startTsNs := f.remoteDeletionRebuildStartTsNs(ctx, mounts) + openGate := func() { + close(done) + f.remoteTombstonesDone.Store(nil) + } + // Deletes made after startup are recorded through the live delete and + // event paths, so replay can only be missing deletes committed before + // the restart — and those have all crossed the tombstone TTL once this + // deadline passes. Holding the gate longer protects nothing. + replayDeadline := time.Now().Add(remoteDeletionTombstoneTTL) + backoff := 2 * time.Second + for { + _, _, err := f.ReadPersistedLogBuffer(ctx, log_buffer.NewMessagePosition(startTsNs, 0), 0, + func(logEntry *filer_pb.LogEntry) (bool, error) { + event := &filer_pb.SubscribeMetadataResponse{} + if err := proto.Unmarshal(logEntry.Data, event); err != nil { + return false, nil + } + f.onRemoteDeletionEvents(event) + return false, nil + }) + if err == nil { + openGate() + return + } + glog.WarningfCtx(ctx, "rebuild remote deletion tombstones: %v", err) + if !time.Now().Before(replayDeadline) { + glog.ErrorfCtx(ctx, "rebuild remote deletion tombstones: giving up after %v", remoteDeletionTombstoneTTL) + openGate() + return + } + select { + case <-ctx.Done(): + return + case <-time.After(backoff): + } + if backoff < time.Minute { + backoff *= 2 + } + } +} + +// remoteDeletionRebuildStartTsNs returns the oldest write-back offset across +// mounts — the earliest event the daemon may not have applied. Mounts without +// a recorded offset replay from the TTL floor: events older than it would +// build tombstones that are already expired. +func (f *Filer) remoteDeletionRebuildStartTsNs(ctx context.Context, mounts []util.FullPath) int64 { + startTsNs := int64(math.MaxInt64) + for _, dir := range mounts { + offset, err := f.readRemoteSyncOffset(ctx, dir) + if err != nil { + glog.WarningfCtx(ctx, "read remote sync offset for %s: %v", dir, err) + offset = 0 + } + if offset < startTsNs { + startTsNs = offset + } + } + ttlFloor := time.Now().Add(-remoteDeletionTombstoneTTL).UnixNano() + if startTsNs == int64(math.MaxInt64) || startTsNs < ttlFloor { + return ttlFloor + } + return startTsNs +} + +// onRemoteDeletionEvents folds peer and local metadata events into the +// tombstone set: a delete or rename source is tombstoned at the event +// timestamp, a create/update/rename target lifts a file tombstone. +func (f *Filer) onRemoteDeletionEvents(event *filer_pb.SubscribeMetadataResponse) { + message := event.EventNotification + if message == nil { + return + } + if message.OldEntry != nil { + sourcePath := filer_pb.MetadataEventSourceFullPath(event) + if message.NewEntry == nil || sourcePath != filer_pb.MetadataEventTargetFullPath(event) { + if f.RemoteStorage != nil && f.remoteTombstones != nil { + if _, remoteLoc := f.RemoteStorage.FindMountDirectory(util.FullPath(sourcePath)); remoteLoc != nil { + f.remoteTombstones.addFromEvent(sourcePath, message.OldEntry.IsDirectory, event.TsNs) + } + } + } + } + if message.NewEntry != nil { + f.remoteTombstones.clear(filer_pb.MetadataEventTargetFullPath(event), event.TsNs) + } +} diff --git a/weed/filer/filer_remote_tombstone_test.go b/weed/filer/filer_remote_tombstone_test.go new file mode 100644 index 000000000..3d989d4c3 --- /dev/null +++ b/weed/filer/filer_remote_tombstone_test.go @@ -0,0 +1,437 @@ +package filer + +import ( + "context" + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/pb/remote_pb" + "github.com/seaweedfs/seaweedfs/weed/remote_storage" + "github.com/seaweedfs/seaweedfs/weed/util" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func newMountedTestFiler(t *testing.T, storageType string, stub remote_storage.RemoteStorageClient, listingTtlSeconds int32) (*Filer, *stubFilerStore) { + t.Helper() + if stub != nil { + t.Cleanup(registerStubMaker(t, storageType, stub)) + } + conf := &remote_pb.RemoteConf{Name: "tombstonestore", Type: storageType} + rs := NewFilerRemoteStorage() + rs.storageNameToConf[conf.Name] = conf + rs.mapDirectoryToRemoteStorage("/buckets/mybucket", &remote_pb.RemoteStorageLocation{ + Name: "tombstonestore", + Bucket: "mybucket", + Path: "/", + ListingCacheTtlSeconds: listingTtlSeconds, + }) + store := newStubFilerStore() + return newTestFiler(t, store, rs), store +} + +func putRemoteSyncOffset(t *testing.T, store *stubFilerStore, dir string, offsetTsNs int64) { + t.Helper() + buf := make([]byte, 8) + util.Uint64toBytes(buf, uint64(offsetTsNs)) + require.NoError(t, store.KvPut(context.Background(), remote_storage.SyncOffsetKey(dir), buf)) +} + +func TestRemoteDeletionTombstones_BlocksAndReleases(t *testing.T) { + tombs := newRemoteDeletionTombstones() + + tombs.addFromEvent("/m/a.txt", false, 100) + ts, _ := tombs.blockedSince("/m/a.txt") + assert.Equal(t, int64(100), ts) + ts, _ = tombs.blockedSince("/m/b.txt") + assert.Zero(t, ts) + + // a delete of a directory blocks its whole subtree + tombs.addFromEvent("/m/dir", true, 200) + ts, _ = tombs.blockedSince("/m/dir") + assert.Equal(t, int64(200), ts) + ts, _ = tombs.blockedSince("/m/dir/deep/x.txt") + assert.Equal(t, int64(200), ts) + ts, _ = tombs.blockedSince("/m/dirx/y.txt") + assert.Zero(t, ts) + + // a rewrite at or after the delete lifts only that file's tombstone + tombs.clear("/m/a.txt", 100) + ts, _ = tombs.blockedSince("/m/a.txt") + assert.Zero(t, ts) + tombs.addFromEvent("/m/a.txt", false, 300) + tombs.clear("/m/a.txt", 250) + ts, _ = tombs.blockedSince("/m/a.txt") + assert.Equal(t, int64(300), ts, "older create must not lift newer delete") + + tombs.releaseThrough("/m/dir/deep/x.txt", 200) + ts, _ = tombs.blockedSince("/m/dir/deep/x.txt") + assert.Zero(t, ts) + ts, _ = tombs.blockedSince("/m/dir") + assert.Zero(t, ts) +} + +func TestRemoteDeletionTombstones_AncestorSubsumesAndCovers(t *testing.T) { + tombs := newRemoteDeletionTombstones() + + // a child tombstone recorded before its ancestor is dropped once the + // ancestor's newer delete covers the whole subtree + tombs.addFromEvent("/m/dir/a.txt", false, 100) + tombs.addFromEvent("/m/dir", true, 200) + ts, _ := tombs.blockedSince("/m/dir") + assert.Equal(t, int64(200), ts) + ts, _ = tombs.blockedSince("/m/dir/a.txt") + assert.Equal(t, int64(200), ts) + _, exists := tombs.files["/m/dir/a.txt"] + assert.False(t, exists, "descendant tombstone is subsumed by the ancestor") + + // adds under the covered subtree are skipped while the ancestor stands + tombs.addFromEvent("/m/dir/b.txt", false, 150) + _, exists = tombs.files["/m/dir/b.txt"] + assert.False(t, exists) + // ...but a child deleted after the ancestor still records its own tombstone + tombs.addFromEvent("/m/dir/c.txt", false, 300) + ts, _ = tombs.blockedSince("/m/dir/c.txt") + assert.Equal(t, int64(300), ts) +} + +func TestRemoteDeletionTombstones_PendingIgnoresOffset(t *testing.T) { + f, store := newMountedTestFiler(t, "stub_tomb_pending_offset", nil, 0) + + // recorded before its event lands: a later event may already have moved + // the mount's watermark past the local timestamp, so the offset cannot + // vouch for this delete yet + filePath := "/buckets/mybucket/dir/a.txt" + now := time.Now().UnixNano() + f.noteRemoteDeletion(util.FullPath(filePath), false, now) + putRemoteSyncOffset(t, store, "/buckets/mybucket", now+100) + assert.True(t, f.isRemoteDeletionPending(context.Background(), util.FullPath(filePath), "/buckets/mybucket")) + + // once the delete event stamps the real timestamp, the watermark releases it + f.onMetadataChangeEvent(&filer_pb.SubscribeMetadataResponse{ + Directory: "/buckets/mybucket/dir", + TsNs: now + 50, + EventNotification: &filer_pb.EventNotification{ + OldEntry: &filer_pb.Entry{Name: "a.txt"}, + }, + }) + assert.False(t, f.isRemoteDeletionPending(context.Background(), util.FullPath(filePath), "/buckets/mybucket")) +} + +func TestRemoteDeletionTombstones_RaisedLocalTombstoneIsPendingAgain(t *testing.T) { + tombs := newRemoteDeletionTombstones() + + tombs.addFromEvent("/m/a.txt", false, 100) + _, pending := tombs.blockedSince("/m/a.txt") + assert.False(t, pending) + + // a newer local delete restamps the tombstone before its own event + // lands, so the write-back offset cannot vouch for it yet + tombs.add("/m/a.txt", false, 200) + _, pending = tombs.blockedSince("/m/a.txt") + assert.True(t, pending) + + // once that delete's event confirms the new timestamp it releases normally + tombs.addFromEvent("/m/a.txt", false, 200) + _, pending = tombs.blockedSince("/m/a.txt") + assert.False(t, pending) +} + +func TestMaybeLazyFetchFromRemote_SkipsTombstonedPath(t *testing.T) { + const storageType = "stub_tomb_fetch" + stub := &countingRemoteClient{ + stubRemoteClient: stubRemoteClient{ + statResult: &filer_pb.RemoteEntry{RemoteMtime: 1700000000, RemoteSize: 11}, + }, + } + f, _ := newMountedTestFiler(t, storageType, stub, 0) + + // a delete under the mount tombstones the path; the remote object is + // still there until the write-back daemon consumes the delete event + filePath := util.FullPath("/buckets/mybucket/dir/a.txt") + f.noteRemoteDeletion(filePath, false, time.Now().UnixNano()) + + entry, err := f.maybeLazyFetchFromRemote(context.Background(), filePath) + require.NoError(t, err) + assert.Nil(t, entry) + assert.Equal(t, 0, stub.statCalls, "the remote object must not even be consulted") +} + +func TestMaybeLazyFetchFromRemote_NewerRemoteMtimeStillBlocked(t *testing.T) { + const storageType = "stub_tomb_regen" + stub := &countingRemoteClient{ + stubRemoteClient: stubRemoteClient{ + statResult: &filer_pb.RemoteEntry{RemoteMtime: time.Now().Unix() + 60, RemoteSize: 11}, + }, + } + f, _ := newMountedTestFiler(t, storageType, stub, 0) + + // a remote object whose mtime postdates the delete is still hidden: the + // remote clock cannot distinguish a new generation from the pending + // delete's target + filePath := util.FullPath("/buckets/mybucket/dir/a.txt") + f.noteRemoteDeletion(filePath, false, time.Now().UnixNano()) + + entry, err := f.maybeLazyFetchFromRemote(context.Background(), filePath) + require.NoError(t, err) + assert.Nil(t, entry) + assert.Equal(t, 0, stub.statCalls) +} + +func TestMaybeLazyFetchFromRemote_SyncOffsetReleasesTombstone(t *testing.T) { + const storageType = "stub_tomb_release" + stub := &countingRemoteClient{ + stubRemoteClient: stubRemoteClient{ + statResult: &filer_pb.RemoteEntry{RemoteMtime: 1700000000, RemoteSize: 11}, + }, + } + f, store := newMountedTestFiler(t, storageType, stub, 0) + + filePath := util.FullPath("/buckets/mybucket/dir/a.txt") + deleteTsNs := time.Now().UnixNano() + f.onMetadataChangeEvent(&filer_pb.SubscribeMetadataResponse{ + Directory: "/buckets/mybucket/dir", + TsNs: deleteTsNs, + EventNotification: &filer_pb.EventNotification{ + OldEntry: &filer_pb.Entry{Name: "a.txt"}, + }, + }) + + // the daemon's persisted watermark is behind the delete: still blocked + putRemoteSyncOffset(t, store, "/buckets/mybucket", deleteTsNs-1) + entry, err := f.maybeLazyFetchFromRemote(context.Background(), filePath) + require.NoError(t, err) + assert.Nil(t, entry) + assert.Equal(t, 0, stub.statCalls) + + // once the watermark reaches the delete event's own timestamp, the + // remote delete has landed and the lookup may consult the remote again + putRemoteSyncOffset(t, store, "/buckets/mybucket", deleteTsNs) + entry, err = f.maybeLazyFetchFromRemote(context.Background(), filePath) + require.NoError(t, err) + require.NotNil(t, entry) + assert.Equal(t, 1, stub.statCalls) +} + +func TestDeleteEntryMetaAndData_TombstonesPath(t *testing.T) { + const storageType = "stub_tomb_delete" + stub := &countingRemoteClient{ + stubRemoteClient: stubRemoteClient{ + statResult: &filer_pb.RemoteEntry{RemoteMtime: 1700000000, RemoteSize: 11}, + }, + } + f, store := newMountedTestFiler(t, storageType, stub, 0) + + filePath := util.FullPath("/buckets/mybucket/dir/a.txt") + store.entries[string(filePath)] = &Entry{ + FullPath: filePath, + Attr: Attr{ + Mtime: time.Unix(1700000000, 0), + Crtime: time.Unix(1700000000, 0), + Mode: 0644, + FileSize: 11, + }, + } + + require.NoError(t, f.DeleteEntryMetaAndData(context.Background(), filePath, false, false, false, false, nil, 0)) + + // the deleted path must not resurrect through the lazy fetch even while + // the remote object is still present + entry, err := f.FindEntry(context.Background(), filePath) + assert.ErrorIs(t, err, filer_pb.ErrNotFound) + assert.Nil(t, entry) + assert.Equal(t, 0, stub.statCalls) + + // a peer-observed create at the path lifts the tombstone + f.onMetadataChangeEvent(&filer_pb.SubscribeMetadataResponse{ + Directory: "/buckets/mybucket/dir", + TsNs: time.Now().UnixNano(), + EventNotification: &filer_pb.EventNotification{ + NewEntry: &filer_pb.Entry{Name: "a.txt"}, + }, + }) + entry, err = f.maybeLazyFetchFromRemote(context.Background(), filePath) + require.NoError(t, err) + require.NotNil(t, entry) + assert.Equal(t, 1, stub.statCalls) +} + +func TestOnMetadataChangeEvent_PeerDeleteTombstones(t *testing.T) { + const storageType = "stub_tomb_peer" + stub := &countingRemoteClient{ + stubRemoteClient: stubRemoteClient{ + statResult: &filer_pb.RemoteEntry{RemoteMtime: 1700000000, RemoteSize: 11}, + }, + } + f, _ := newMountedTestFiler(t, storageType, stub, 0) + + f.onMetadataChangeEvent(&filer_pb.SubscribeMetadataResponse{ + Directory: "/buckets/mybucket/dir", + TsNs: time.Now().UnixNano(), + EventNotification: &filer_pb.EventNotification{ + OldEntry: &filer_pb.Entry{Name: "a.txt"}, + }, + }) + + entry, err := f.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/dir/a.txt") + require.NoError(t, err) + assert.Nil(t, entry) + assert.Equal(t, 0, stub.statCalls) +} + +func TestOnMetadataChangeEvent_PeerDirDeleteTombstonesSubtree(t *testing.T) { + const storageType = "stub_tomb_peer_dir" + stub := &countingRemoteClient{ + stubRemoteClient: stubRemoteClient{ + statResult: &filer_pb.RemoteEntry{RemoteMtime: 1700000000, RemoteSize: 11}, + }, + } + f, _ := newMountedTestFiler(t, storageType, stub, 0) + + f.onMetadataChangeEvent(&filer_pb.SubscribeMetadataResponse{ + Directory: "/buckets/mybucket", + TsNs: time.Now().UnixNano(), + EventNotification: &filer_pb.EventNotification{ + OldEntry: &filer_pb.Entry{Name: "dir", IsDirectory: true}, + }, + }) + + entry, err := f.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/dir/deep/a.txt") + require.NoError(t, err) + assert.Nil(t, entry) + assert.Equal(t, 0, stub.statCalls) +} + +func TestMaybeLazyListFromRemote_SkipsTombstonedChild(t *testing.T) { + const storageType = "stub_tomb_list" + stub := &stubRemoteClient{ + listDirFn: func(loc *remote_pb.RemoteStorageLocation, visitFn remote_storage.VisitFunc) error { + if err := visitFn("/", "deleted.txt", false, &filer_pb.RemoteEntry{RemoteMtime: 1700000000, RemoteSize: 10}); err != nil { + return err + } + return visitFn("/", "fresh.txt", false, &filer_pb.RemoteEntry{RemoteMtime: 1700000000, RemoteSize: 20}) + }, + } + f, store := newMountedTestFiler(t, storageType, stub, 300) + + f.noteRemoteDeletion("/buckets/mybucket/deleted.txt", false, time.Now().UnixNano()) + + f.maybeLazyListFromRemote(context.Background(), util.FullPath("/buckets/mybucket")) + assert.Equal(t, 1, stub.listDirCalls) + + assert.Nil(t, store.getEntry("/buckets/mybucket/deleted.txt"), "deleted child must not resurrect through a listing") + require.NotNil(t, store.getEntry("/buckets/mybucket/fresh.txt"), "other remote objects still list") +} + +func TestMaybeLazyListFromRemote_RecreatedDirStillHidesChildren(t *testing.T) { + const storageType = "stub_tomb_recreate" + deleteTs := time.Now().UnixNano() + stub := &stubRemoteClient{ + listDirFn: func(loc *remote_pb.RemoteStorageLocation, visitFn remote_storage.VisitFunc) error { + if err := visitFn("/", "stale.txt", false, &filer_pb.RemoteEntry{RemoteMtime: deleteTs/int64(time.Second) - 10, RemoteSize: 10}); err != nil { + return err + } + return visitFn("/", "fresh.txt", false, &filer_pb.RemoteEntry{RemoteMtime: deleteTs/int64(time.Second) + 10, RemoteSize: 20}) + }, + } + f, store := newMountedTestFiler(t, storageType, stub, 300) + + // the directory is deleted then recreated; its remote children stay + // hidden — old or new mtime alike — until the remote delete lands + f.onMetadataChangeEvent(&filer_pb.SubscribeMetadataResponse{ + Directory: "/buckets/mybucket", + TsNs: deleteTs, + EventNotification: &filer_pb.EventNotification{ + OldEntry: &filer_pb.Entry{Name: "dir", IsDirectory: true}, + }, + }) + f.onMetadataChangeEvent(&filer_pb.SubscribeMetadataResponse{ + Directory: "/buckets/mybucket", + TsNs: deleteTs + 1, + EventNotification: &filer_pb.EventNotification{ + NewEntry: &filer_pb.Entry{Name: "dir", IsDirectory: true}, + }, + }) + + f.maybeLazyListFromRemote(context.Background(), util.FullPath("/buckets/mybucket/dir")) + + assert.Nil(t, store.getEntry("/buckets/mybucket/dir/stale.txt")) + assert.Nil(t, store.getEntry("/buckets/mybucket/dir/fresh.txt")) + + // after the write-back daemon confirms the remote delete, listing merges again + f.remoteTombstones.releaseThrough("/buckets/mybucket/dir/stale.txt", deleteTs) + delete(store.getEntry("/buckets/mybucket/dir").Extended, xattrRemoteListingSyncedAt) + f.maybeLazyListFromRemote(context.Background(), util.FullPath("/buckets/mybucket/dir")) + assert.NotNil(t, store.getEntry("/buckets/mybucket/dir/stale.txt")) + assert.NotNil(t, store.getEntry("/buckets/mybucket/dir/fresh.txt")) +} + +func TestRemoteDeletionRebuildStartTsNs_UsesOldestMountOffset(t *testing.T) { + f, store := newMountedTestFiler(t, "stub_rebuild", nil, 0) + f.RemoteStorage.mapDirectoryToRemoteStorage("/buckets/other", &remote_pb.RemoteStorageLocation{ + Name: "tombstonestore", Bucket: "other", Path: "/", + }) + mounts := f.RemoteStorage.MountedDirectories() + require.Len(t, mounts, 2) + + now := time.Now().UnixNano() + putRemoteSyncOffset(t, store, "/buckets/mybucket", now-200) + putRemoteSyncOffset(t, store, "/buckets/other", now-300) + assert.Equal(t, now-300, f.remoteDeletionRebuildStartTsNs(context.Background(), mounts), + "rebuild must replay from the least-synced mount") + + // a mount whose offset was never written replays only within the TTL + require.NoError(t, store.KvDelete(context.Background(), remote_storage.SyncOffsetKey("/buckets/other"))) + floor := f.remoteDeletionRebuildStartTsNs(context.Background(), mounts) + assert.GreaterOrEqual(t, floor, time.Now().Add(-remoteDeletionTombstoneTTL-time.Second).UnixNano()) + + // an offset older than the TTL floor is raised to it + putRemoteSyncOffset(t, store, "/buckets/other", time.Now().Add(-remoteDeletionTombstoneTTL-time.Hour).UnixNano()) + assert.Greater(t, f.remoteDeletionRebuildStartTsNs(context.Background(), mounts), + time.Now().Add(-remoteDeletionTombstoneTTL-time.Hour).UnixNano()) +} + +func TestRebuildRemoteDeletionTombstones_EmptyLogReleasesGate(t *testing.T) { + f, _ := newMountedTestFiler(t, "stub_rebuild_empty", nil, 0) + + f.RebuildRemoteDeletionTombstones(context.Background()) + done := f.remoteTombstonesDone.Load() + require.NotNil(t, done, "the gate must be set synchronously") + + select { + case <-*done: + case <-time.After(10 * time.Second): + t.Fatal("rebuild never released the gate") + } + assert.Nil(t, f.remoteTombstonesDone.Load()) + ts, _ := f.remoteTombstones.blockedSince("/buckets/mybucket/a.txt") + assert.Zero(t, ts) +} + +func TestMaybeLazyFetchFromRemote_WaitsForRebuild(t *testing.T) { + const storageType = "stub_tomb_pending" + stub := &countingRemoteClient{ + stubRemoteClient: stubRemoteClient{ + statResult: &filer_pb.RemoteEntry{RemoteMtime: 1700000000, RemoteSize: 10}, + }, + } + f, _ := newMountedTestFiler(t, storageType, stub, 0) + + // an unfinished rebuild blocks the fetch until the context gives up + gate := make(chan struct{}) + f.remoteTombstonesDone.Store(&gate) + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + entry, err := f.maybeLazyFetchFromRemote(ctx, "/buckets/mybucket/a.txt") + assert.ErrorIs(t, err, context.DeadlineExceeded) + assert.Nil(t, entry) + assert.Equal(t, 0, stub.statCalls) + + // once the rebuild finishes, the fetch proceeds + close(gate) + entry, err = f.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/a.txt") + require.NoError(t, err) + require.NotNil(t, entry) + assert.Equal(t, 1, stub.statCalls) +} diff --git a/weed/filer/remote_storage.go b/weed/filer/remote_storage.go index 3a2043f39..e1733ce77 100644 --- a/weed/filer/remote_storage.go +++ b/weed/filer/remote_storage.go @@ -36,6 +36,7 @@ type FilerRemoteStorage struct { // whenever /etc/remote changes mu sync.RWMutex rules ptrie.Trie[*remote_pb.RemoteStorageLocation] + mountDirs []util.FullPath storageNameToConf map[string]*remote_pb.RemoteConf // confValidator, when set, is applied to every RemoteConf as it is loaded // from /etc/remote. A conf that fails is dropped from storageNameToConf so @@ -87,11 +88,13 @@ func (rs *FilerRemoteStorage) LoadRemoteStorageConfigurationsAndMapping(filer *F // build into fresh containers so an unmounted directory disappears instead // of lingering in the trie, which has no way to drop a key rules := ptrie.New[*remote_pb.RemoteStorageLocation]() + var mountDirs []util.FullPath storageNameToConf := make(map[string]*remote_pb.RemoteConf) for _, entry := range entries { if entry.Name() == REMOTE_STORAGE_MOUNT_FILE { - if err := loadRemoteStorageMountMapping(rules, entry.Content); err != nil { + mountDirs, err = loadRemoteStorageMountMapping(rules, entry.Content) + if err != nil { return err } continue @@ -117,27 +120,46 @@ func (rs *FilerRemoteStorage) LoadRemoteStorageConfigurationsAndMapping(filer *F } rs.mu.Lock() - rs.rules, rs.storageNameToConf = rules, storageNameToConf + rs.rules, rs.mountDirs, rs.storageNameToConf = rules, mountDirs, storageNameToConf rs.mu.Unlock() return nil } -func loadRemoteStorageMountMapping(rules ptrie.Trie[*remote_pb.RemoteStorageLocation], data []byte) (err error) { +func loadRemoteStorageMountMapping(rules ptrie.Trie[*remote_pb.RemoteStorageLocation], data []byte) (mountDirs []util.FullPath, err error) { mappings := &remote_pb.RemoteStorageMapping{} if err := proto.Unmarshal(data, mappings); err != nil { - return fmt.Errorf("unmarshal %s/%s: %v", DirectoryEtcRemote, REMOTE_STORAGE_MOUNT_FILE, err) + return nil, fmt.Errorf("unmarshal %s/%s: %v", DirectoryEtcRemote, REMOTE_STORAGE_MOUNT_FILE, err) } for dir, storageLocation := range mappings.Mappings { putDirectoryToRemoteStorage(rules, util.FullPath(dir), storageLocation) + mountDirs = append(mountDirs, util.FullPath(dir)) } - return nil + return mountDirs, nil } func (rs *FilerRemoteStorage) mapDirectoryToRemoteStorage(dir util.FullPath, loc *remote_pb.RemoteStorageLocation) { rs.mu.Lock() defer rs.mu.Unlock() putDirectoryToRemoteStorage(rs.rules, dir, loc) + found := false + for _, d := range rs.mountDirs { + if d == dir { + found = true + break + } + } + if !found { + rs.mountDirs = append(rs.mountDirs, dir) + } +} + +// MountedDirectories returns the directories currently mapped to remote +// storage, for callers that must reason about every mount rather than one path. +func (rs *FilerRemoteStorage) MountedDirectories() []util.FullPath { + rs.mu.RLock() + defer rs.mu.RUnlock() + return append([]util.FullPath(nil), rs.mountDirs...) } func putDirectoryToRemoteStorage(rules ptrie.Trie[*remote_pb.RemoteStorageLocation], dir util.FullPath, loc *remote_pb.RemoteStorageLocation) { diff --git a/weed/remote_storage/track_sync_offset.go b/weed/remote_storage/track_sync_offset.go index f5de002d0..7fa77a468 100644 --- a/weed/remote_storage/track_sync_offset.go +++ b/weed/remote_storage/track_sync_offset.go @@ -14,15 +14,19 @@ const ( SyncKeyPrefix = "remote.sync." ) +// SyncOffsetKey is the filer store key holding the write-back sync watermark +// for a mounted directory. +func SyncOffsetKey(dir string) []byte { + syncKey := make([]byte, len(SyncKeyPrefix)+4) + copy(syncKey, SyncKeyPrefix) + util.Uint32toBytes(syncKey[len(SyncKeyPrefix):], uint32(util.HashStringToLong(dir))) + return syncKey +} + func GetSyncOffset(grpcDialOption grpc.DialOption, filer pb.ServerAddress, dir string) (lastOffsetTsNs int64, readErr error) { - dirHash := uint32(util.HashStringToLong(dir)) - readErr = pb.WithFilerClient(false, 0, filer, grpcDialOption, func(client filer_pb.SeaweedFilerClient) error { - syncKey := []byte(SyncKeyPrefix + "____") - util.Uint32toBytes(syncKey[len(SyncKeyPrefix):len(SyncKeyPrefix)+4], dirHash) - - resp, err := client.KvGet(context.Background(), &filer_pb.KvGetRequest{Key: syncKey}) + resp, err := client.KvGet(context.Background(), &filer_pb.KvGetRequest{Key: SyncOffsetKey(dir)}) if err != nil { return err } @@ -45,18 +49,13 @@ func GetSyncOffset(grpcDialOption grpc.DialOption, filer pb.ServerAddress, dir s func SetSyncOffset(grpcDialOption grpc.DialOption, filer pb.ServerAddress, dir string, offsetTsNs int64) error { - dirHash := uint32(util.HashStringToLong(dir)) - return pb.WithFilerClient(false, 0, filer, grpcDialOption, func(client filer_pb.SeaweedFilerClient) error { - syncKey := []byte(SyncKeyPrefix + "____") - util.Uint32toBytes(syncKey[len(SyncKeyPrefix):len(SyncKeyPrefix)+4], dirHash) - valueBuf := make([]byte, 8) util.Uint64toBytes(valueBuf, uint64(offsetTsNs)) resp, err := client.KvPut(context.Background(), &filer_pb.KvPutRequest{ - Key: syncKey, + Key: SyncOffsetKey(dir), Value: valueBuf, }) if err != nil { diff --git a/weed/server/filer_server.go b/weed/server/filer_server.go index 7d998ad8f..959e1ac8c 100644 --- a/weed/server/filer_server.go +++ b/weed/server/filer_server.go @@ -311,6 +311,8 @@ func NewFilerServer(defaultMux, readonlyMux *http.ServeMux, option *FilerOption) fs.filer.LoadRemoteStorageConfAndMapping() + fs.filer.RebuildRemoteDeletionTombstones(context.Background()) + grace.OnReload(fs.Reload) fs.SetupDlmReplication()