mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-29 11:15:34 +00:00
filer: keep lazy remote reads from resurrecting deleted paths (#11452)
* filer: keep lazy remote reads from resurrecting deleted paths Under a remote mount with filer.remote.sync as write-back, a path that was deleted or renamed away could come back as a chunkless remote-only entry: between the local delete and the daemon's remote delete, a store miss made maybeLazyFetchFromRemote trust a bucket that was behind the filer. The ghost then outlived the remote object -- HEAD answered 200, GET failed, and nothing cleaned it up. The filer now tombstones paths it deletes under a remote mount, learned both synchronously from its own delete path and from peer metadata events. The lazy fetch and the lazy listing skip a tombstoned path until the path is written again, until the mount's persisted write-back sync offset has passed the delete event (the remote delete has landed), or until a generous TTL covers a mount without a daemon. Fixes #11440 Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: cover recursive remote deletes with an ancestor tombstone A recursive delete now records the directory tombstone before walking children, so a partial traversal or a store that drops the subtree without listing it still leaves every descendant covered. Directory tombstones also subsume older descendant entries on add, descendant adds covered by a standing ancestor are skipped, and an existing tombstone can be refreshed even at capacity. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: scope remote tombstones to the deleted object's generation A remote object whose own mtime postdates the local delete is a new generation, not the one the tombstone hides, so a recreated directory can surface remote writes made after its delete while old-generation objects stay hidden. Lazy fetch now stats the remote object before deciding, listings pass each child's remote mtime, and a sync offset releases a tombstone once it reaches the delete's own timestamp. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: rebuild remote deletion tombstones after restart In-memory tombstones are lost on restart while remote write-back offsets persist, so a filer boot replays the persisted metadata log from the oldest mount offset and folds deletes back into the tombstone set through the same event handler. Lazy remote reads hold off while the replay runs so a pending delete cannot resurrect in the gap. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: release remote tombstones only after their delete event lands The write-back offset orders against event timestamps, but the synchronous delete path recorded tombstones with the local clock before its event was emitted — a later unrelated event could already have pushed the mount's watermark past that guess, releasing the tombstone before the daemon applied the delete. Tombstones recorded ahead of their event are now marked pending and can only be lifted by the event confirming them or by TTL; event-stamped tombstones release through the offset as before. The remote-mtime generation bypass is dropped: remote and filer clocks are independent, and a pending remote delete removes whatever object sits at the path, so a "newer" remote object would only resurrect as a phantom. Tombstoned lookups now skip the remote stat entirely. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: drop dir tombstone when recursive delete fails before listing The ancestor tombstone is recorded before the child listing; if that listing fails nothing was deleted, and the leftover tombstone would hide still-existing remote children for the whole TTL. Tombstones for children already deleted stay, since their remote deletes are still owed. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: block lazy remote reads on startup tombstone rebuild The rebuild gate is now a done-channel set synchronously before the replay goroutine starts, so no lazy read can slip through in between. Reads wait on it with context cancellation instead of returning an empty miss that makes remote-only objects look deleted. The replay start is floored at now-TTL: mounts without a recorded write-back offset previously replayed the whole persisted history, and events older than the TTL would only build already-expired tombstones. The gate check now runs after the mount lookup so replaying the meta log's own directory listings does not deadlock on the gate, and the replay retries with backoff until it succeeds instead of failing open. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: mark restamped tombstone pending until its delete event lands When a local delete raises an existing tombstone's timestamp, the new value is only a local clock guess ahead of that delete's event. Leaving the tombstone un-pending lets a write-back offset release it before the event is actually consumed, reopening the resurrection window. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: bound tombstone replay to the tombstone TTL Persisted-log replay retried forever, keeping lazy remote reads gated indefinitely when the log cannot be read. Cap retries at the tombstone TTL measured from replay start: past that point every tombstone would have expired anyway, so opening the gate loses no protection. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: re-check deletion tombstone before persisting lazy fetch A delete landing while StatFile is in flight passed the earlier tombstone check but still persisted the fetched entry, resurrecting a path whose remote delete is pending. Re-check right before CreateEntry. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: retract a lazily persisted entry when a delete raced the insert The pre-insert tombstone check still leaves a window between the check and the store insert. Since deletes always record the tombstone before removing the entry, a tombstone visible right after a successful insert means the delete already ran: delete the entry back out so the tombstoned path stays deleted. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: note why the replay deadline can safely open the gate Deletes made after startup are captured by the live delete and event paths, so a stalled replay can only be missing pre-restart deletes, all of which are past the tombstone TTL by the deadline. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: retract only the entry a lazy remote read materialized Deleting by path after a raced delete could remove a legitimate rewrite that replaced the fetched entry. Verify the stored entry still matches the remote object (or the just-created directory shape) before deleting, and apply the same post-insert check to lazy listing children. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: require full-entry equality before retracting a lazy entry Remote-only matching still removed a write that had updated the fetched entry, e.g. appended chunks. Compare the persisted entry against what this read materialized; any change means a real update owns the path. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
co-authored by
Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
parent
df4995b894
commit
2f6c237238
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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() {}),
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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) {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user