diff --git a/weed/filer/filer.go b/weed/filer/filer.go index 2e7dae156..b2c09fcd1 100644 --- a/weed/filer/filer.go +++ b/weed/filer/filer.go @@ -7,6 +7,7 @@ import ( "os" "sort" "strings" + "sync" "sync/atomic" "time" @@ -85,6 +86,30 @@ type Filer struct { // rebuild finishes; lazy remote reads wait on it so a pending delete // cannot resurrect in the gap. remoteTombstonesDone atomic.Pointer[chan struct{}] + + // Durable deletion ledger (see filer_deletion_persist.go). The set of + // fileIds that still need deleting but are not yet confirmed gone, mirrored + // to the store so a restart does not leak chunks. Guarded by + // deletionLedgerLock; nil-safe for Filer literals in tests. + // pendingDeletions maps each id to its enqueue epoch so an expiry-forget + // cannot erase a re-queued id. + deletionLedgerLock sync.Mutex + pendingDeletions map[string]uint64 + deletionSeq uint64 + deletionLedgerDirty bool + deletionLedgerParts int + deletionLedgerGen int + // deletionLedgerStale holds orphan part keys from abandoned multipart + // writes, retried on the next snapshot. + deletionLedgerStale []string + // deletionSnapshotLock serializes ledger writes in copy order so an + // in-flight timer snapshot cannot overwrite a newer shutdown snapshot. + deletionSnapshotLock sync.Mutex + // deletionLedgerBlocked is set when the startup ledger read fails: the + // persisted set is then unknown, so snapshots retry the read instead of + // overwriting the unread ledger with a partial set. + deletionLedgerBlocked atomic.Bool + deletionLedgerFlush 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 { @@ -98,6 +123,7 @@ func NewFiler(masters pb.ServerDiscovery, grpcDialOption grpc.DialOption, filerH Dlm: lock_manager.NewDistributedLockManager(filerHost), MaxFilenameLength: maxFilenameLength, deletionQuit: make(chan struct{}), + deletionLedgerFlush: make(chan struct{}, 1), DeletionRetryQueue: NewDeletionRetryQueue(), persistedLogCache: newPersistedLogCache(persistedLogCacheMaxBytes), remoteTombstones: newRemoteDeletionTombstones(), @@ -228,7 +254,16 @@ func (f *Filer) ListExistingPeerUpdates(ctx context.Context) (existingNodes []*m func (f *Filer) SetStore(store FilerStore) (isFresh bool) { f.Store = NewFilerStoreWrapper(store) - return f.setOrLoadFilerStoreSignature(store) + isFresh = f.setOrLoadFilerStoreSignature(store) + + // Recover deletions that were pending when a previous process died, and keep + // the durable ledger snapshotted while running (see filer_deletion_persist.go). + // A failed ledger read leaves writes blocked until a snapshot retries and + // the read succeeds, so the unread ledger is never overwritten. + f.reloadDeletionLedger() + f.startDeletionLedgerSnapshotter() + + return isFresh } func (f *Filer) setOrLoadFilerStoreSignature(store FilerStore) (isFresh bool) { @@ -776,6 +811,9 @@ func (f *Filer) Shutdown() { f.LocalMetaLogBuffer.ShutdownLogBuffer() // The final metadata-log flush still needs the store to append its entry. f.LocalMetaLogBuffer.WaitForShutdown() + // Persist the deletion ledger one last time before the store closes, so a + // clean shutdown leaves the recovery set exactly consistent with reality. + f.snapshotDeletionLedger() f.Store.Shutdown() } diff --git a/weed/filer/filer_deletion.go b/weed/filer/filer_deletion.go index 5ca125a31..ba65a325f 100644 --- a/weed/filer/filer_deletion.go +++ b/weed/filer/filer_deletion.go @@ -66,8 +66,9 @@ type DeletionRetryItem struct { RetryCount int NextRetryAt time.Time LastError string - heapIndex int // index in the heap (for heap.Interface) - inFlight bool // true when item is being processed, prevents duplicate additions + ledgerEpoch uint64 // pendingDeletions epoch when the item was queued; expiry forgets only a matching epoch + heapIndex int // index in the heap (for heap.Interface) + inFlight bool // true when item is being processed, prevents duplicate additions } // retryHeap implements heap.Interface for DeletionRetryItem @@ -163,7 +164,7 @@ func calculateBackoff(retryCount int) time.Duration { // AddOrUpdate adds a new failed deletion or updates an existing one // Time complexity: O(log N) for insertion/update -func (q *DeletionRetryQueue) AddOrUpdate(fileId string, errorMsg string) { +func (q *DeletionRetryQueue) AddOrUpdate(fileId string, errorMsg string, ledgerEpoch uint64) { q.lock.Lock() defer q.lock.Unlock() @@ -173,6 +174,13 @@ func (q *DeletionRetryQueue) AddOrUpdate(fileId string, errorMsg string) { // The existing retry schedule should proceed. // RetryCount is only incremented in RequeueForRetry when an actual retry is performed. item.LastError = errorMsg + // Keep the recorded epoch while the item is in flight: an expiry or + // permanent outcome from that attempt must forget only the record it + // started with, not a newer enqueue for the same id. This also keeps + // the field immutable once the worker can read it without the lock. + if !item.inFlight { + item.ledgerEpoch = ledgerEpoch + } if item.inFlight { glog.V(2).Infof("retry for %s in-flight: attempt %d, will preserve retry state", fileId, item.RetryCount) } else { @@ -188,6 +196,7 @@ func (q *DeletionRetryQueue) AddOrUpdate(fileId string, errorMsg string) { RetryCount: 1, NextRetryAt: time.Now().Add(delay), LastError: errorMsg, + ledgerEpoch: ledgerEpoch, inFlight: false, } heap.Push(&q.heap, item) @@ -216,18 +225,19 @@ func (q *DeletionRetryQueue) RequeueForRetry(item *DeletionRetryItem, errorMsg s heap.Push(&q.heap, item) } -// GetReadyItems returns items that are ready to be retried and marks them as in-flight +// GetReadyItems returns items that are ready to be retried and marks them as in-flight. +// Expired reports fileIds dropped for exceeding MaxRetryAttempts — permanent +// dropouts the caller must also forget in the durable deletion ledger. // Time complexity: O(K log N) where K is the number of ready items // Items are processed in order of NextRetryAt (earliest first) -func (q *DeletionRetryQueue) GetReadyItems(maxItems int) []*DeletionRetryItem { +func (q *DeletionRetryQueue) GetReadyItems(maxItems int) (ready []*DeletionRetryItem, expired []*DeletionRetryItem) { q.lock.Lock() defer q.lock.Unlock() now := time.Now() - var readyItems []*DeletionRetryItem // Peek at items from the top of the heap (earliest NextRetryAt) - for len(q.heap) > 0 && len(readyItems) < maxItems { + for len(q.heap) > 0 && len(ready) < maxItems { item := q.heap[0] // If the earliest item is not ready yet, no other items are ready either @@ -240,15 +250,16 @@ func (q *DeletionRetryQueue) GetReadyItems(maxItems int) []*DeletionRetryItem { if item.RetryCount <= MaxRetryAttempts { item.inFlight = true // Mark as being processed - readyItems = append(readyItems, item) + ready = append(ready, item) } else { // Max attempts reached, log and discard completely delete(q.itemIndex, item.FileId) + expired = append(expired, item) glog.Warningf("max retry attempts (%d) reached for %s, last error: %s", MaxRetryAttempts, item.FileId, item.LastError) } } - return readyItems + return ready, expired } // Remove removes an item from the queue (called when deletion succeeds or fails permanently) @@ -343,6 +354,13 @@ func (f *Filer) processDeletionBatch(ctx context.Context, toDeleteFileIds []stri return } + // Remember each id's ledger epoch before the remote deletes run: a + // permanent outcome below must not erase a record re-queued mid-flight. + epochs := make(map[string]uint64, len(uniqueFileIdsSlice)) + for _, fileId := range uniqueFileIdsSlice { + epochs[fileId] = f.deletionEpoch(fileId) + } + // Delete files and classify outcomes outcomes := deleteFilesAndClassify(ctx, f.GrpcDialOption, uniqueFileIdsSlice, lookupFunc) @@ -356,16 +374,21 @@ func (f *Filer) processDeletionBatch(ctx context.Context, toDeleteFileIds []stri switch outcome.status { case deletionOutcomeSuccess: successCount++ + f.forgetDeletion(fileId) // confirmed gone: drop from durable ledger case deletionOutcomeNotFound: notFoundCount++ + f.forgetDeletion(fileId) // already gone: drop from durable ledger case deletionOutcomeRetryable, deletionOutcomeNoResult: retryableErrorCount++ - f.DeletionRetryQueue.AddOrUpdate(fileId, outcome.errorMsg) + // Keep in the durable ledger: not confirmed, must survive restarts. + f.DeletionRetryQueue.AddOrUpdate(fileId, outcome.errorMsg, f.deletionEpoch(fileId)) if len(errorDetails) < MaxLoggedErrorDetails { errorDetails = append(errorDetails, fileId+": "+outcome.errorMsg+" (will retry)") } case deletionOutcomePermanent: permanentErrorCount++ + // gave up: drop the record, unless the id was re-queued meanwhile + f.forgetDeletionEpoch(fileId, epochs[fileId]) if len(errorDetails) < MaxLoggedErrorDetails { errorDetails = append(errorDetails, fileId+": "+outcome.errorMsg+" (permanent)") } @@ -523,7 +546,16 @@ func (f *Filer) loopProcessingDeletionRetry(lookupFunc func([]string) (map[strin // Process all ready items in batches until queue is empty totalProcessed := 0 for { - readyItems := f.DeletionRetryQueue.GetReadyItems(DeletionRetryBatchSize) + readyItems, expired := f.DeletionRetryQueue.GetReadyItems(DeletionRetryBatchSize) + for _, item := range expired { + // Permanently discarded — drop the ledger record, but only + // if the id was not re-queued after this retry item was + // recorded. A surviving newer record still needs work, so + // it goes back through the hot queue. + if !f.forgetDeletionEpoch(item.FileId, item.ledgerEpoch) { + f.queueDeletions(item.FileId) + } + } if len(readyItems) == 0 { break } @@ -561,19 +593,27 @@ func (f *Filer) processRetryBatch(readyItems []*DeletionRetryItem, lookupFunc fu case deletionOutcomeSuccess: successCount++ f.DeletionRetryQueue.Remove(item) // Remove from queue (success) + f.forgetDeletion(item.FileId) // confirmed gone: drop from durable ledger glog.V(2).Infof("retry successful for %s after %d attempts", item.FileId, item.RetryCount) case deletionOutcomeNotFound: notFoundCount++ f.DeletionRetryQueue.Remove(item) // Remove from queue (already deleted) + f.forgetDeletion(item.FileId) // already gone: drop from durable ledger case deletionOutcomeRetryable, deletionOutcomeNoResult: retryCount++ if outcome.status == deletionOutcomeNoResult { glog.Warningf("no deletion result for retried file %s, re-queuing to avoid loss", item.FileId) } + // Keep in the durable ledger: still not confirmed. f.DeletionRetryQueue.RequeueForRetry(item, outcome.errorMsg) case deletionOutcomePermanent: permanentErrorCount++ f.DeletionRetryQueue.Remove(item) // Remove from queue (permanent failure) + // gave up: drop the record, unless the id was re-queued meanwhile + // — a surviving newer record goes back through the hot queue. + if !f.forgetDeletionEpoch(item.FileId, item.ledgerEpoch) { + f.queueDeletions(item.FileId) + } glog.Warningf("permanent error on retry for %s after %d attempts: %s", item.FileId, item.RetryCount, outcome.errorMsg) } } @@ -604,7 +644,7 @@ func (f *Filer) DeleteChunks(ctx context.Context, fullpath util.FullPath, chunks func (f *Filer) doDeleteChunks(ctx context.Context, chunks []*filer_pb.FileChunk) { for _, chunk := range chunks { if !chunk.IsChunkManifest { - f.FileIdDeletionQueue.EnQueue(chunk.GetFileIdString()) + f.queueDeletions(chunk.GetFileIdString()) continue } dataChunks, manifestResolveErr := ResolveOneChunkManifest(ctx, f.MasterClient.LookupFileId, chunk, f.MasterClient) @@ -612,15 +652,15 @@ func (f *Filer) doDeleteChunks(ctx context.Context, chunks []*filer_pb.FileChunk glog.V(0).InfofCtx(ctx, "failed to resolve manifest %s: %v", chunk.FileId, manifestResolveErr) } for _, dChunk := range dataChunks { - f.FileIdDeletionQueue.EnQueue(dChunk.GetFileIdString()) + f.queueDeletions(dChunk.GetFileIdString()) } - f.FileIdDeletionQueue.EnQueue(chunk.GetFileIdString()) + f.queueDeletions(chunk.GetFileIdString()) } } func (f *Filer) DeleteChunksNotRecursive(chunks []*filer_pb.FileChunk) { for _, chunk := range chunks { - f.FileIdDeletionQueue.EnQueue(chunk.GetFileIdString()) + f.queueDeletions(chunk.GetFileIdString()) } } diff --git a/weed/filer/filer_deletion_persist.go b/weed/filer/filer_deletion_persist.go new file mode 100644 index 000000000..33b2371cf --- /dev/null +++ b/weed/filer/filer_deletion_persist.go @@ -0,0 +1,758 @@ +package filer + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "time" + + "github.com/seaweedfs/seaweedfs/weed/glog" + "github.com/seaweedfs/seaweedfs/weed/util" +) + +// Persistent deletion ledger. +// +// Background: the in-memory FileIdDeletionQueue and DeletionRetryQueue lose every +// queued-but-unconfirmed deletion when the filer process restarts (the upstream +// TODO in filer_deletion.go notes this and proposes exactly the "periodic snapshot +// with recovery on startup" strategy implemented here). Because a delete only ever +// enters the pipeline through the queue, a crash between enqueue and the volume +// actually confirming the delete leaks the chunk: nothing remembers it, so it +// becomes an orphan the next fsck sees as 100% orphaned and a later +// meta-replay from a lagging peer can "resurrect" as if it were live data. +// +// The ledger is a KV entry per filer (KvKeyDeletionLedger suffixed with this +// filer's address, so filers sharing one store never overwrite each other's +// pending sets) holding the fileIds that still need to be deleted but have not +// yet been confirmed gone. It is: +// - additive: every EnQueue also records the fileId here, +// - subtractive: only terminal outcomes (success / not-found / permanent) +// remove it; retryable failures keep it, +// - snapshotted when it changes, on a timer, and on Shutdown, so a crash +// loses only ids queued in the last in-flight write — recovery re-enqueues +// the whole pending set and the idempotent volume delete absorbs anything +// that actually completed before the crash. +// +// Stores that cap value size (FoundationDB at 100KB) get the set split into +// part keys when one value would exceed deletionLedgerPartSize. Multipart +// snapshots write each generation under generation-scoped part keys and +// publish only the manifest last, so a crash mid-write leaves the previous +// generation intact; orphaned parts from abandoned generations are tracked in +// deletionLedgerStale and retried on the next write. +// +// Safety on a lagging peer (the resurrection case): recovered entries are replayed +// into the queue only after a grace window (DeletionRecoveryGrace), long enough for +// the initial peer meta-aggregation to settle so we don't purge a chunk that a +// peer is about to re-reference as live. New live deletes are never gated. +const ( + // KvKeyDeletionLedger is the reserved store key prefix for the persisted + // ledger, namespaced next to FilerStoreId so it never collides with user + // data. The full key is this prefix plus the filer's own address. + KvKeyDeletionLedger = "filer.deleteQueue.ledger.v1" + + // KvKeyDeletionLedgerIndex lists the scoped ledger keys in the store so a + // filer that restarts under a different advertised address can still find + // and claim the ledger it left behind. + KvKeyDeletionLedgerIndex = KvKeyDeletionLedger + ".index" + + // deletionLedgerPartSize bounds one ledger value; stores with a value cap + // reject anything larger, which would strand a big backlog in memory only. + deletionLedgerPartSize = 64 * 1024 + + // Defaults; all three overridable via viper (filer.deleteQueue.*). + defaultDeletionPersistInterval = 10 * time.Second + defaultDeletionRecoveryGrace = 30 * time.Second +) + +// deletionPersistEnabled reports whether ledger persistence is on. +// Read from viper with a default of true so a plain `weed filer` gets the +// durability guarantee without extra flags. +func deletionPersistEnabled() bool { + v := util.GetViper() + v.SetDefault("filer.deleteQueue.persist", true) + return v.GetBool("filer.deleteQueue.persist") +} + +func deletionPersistInterval() time.Duration { + v := util.GetViper() + v.SetDefault("filer.deleteQueue.persistInterval", defaultDeletionPersistInterval) + d := v.GetDuration("filer.deleteQueue.persistInterval") + if d <= 0 { + return defaultDeletionPersistInterval + } + return d +} + +func deletionRecoveryGrace() time.Duration { + v := util.GetViper() + v.SetDefault("filer.deleteQueue.recoveryGrace", defaultDeletionRecoveryGrace) + d := v.GetDuration("filer.deleteQueue.recoveryGrace") + if d < 0 { + return 0 + } + return d +} + +// deletionLedgerKey scopes the ledger to this filer so several filers sharing +// one metadata store do not overwrite each other's pending sets. A filer +// literal without an advertised address (tests) falls back to the base key. +func (f *Filer) deletionLedgerKey() string { + if f.Dlm != nil && f.Dlm.Host != "" { + return KvKeyDeletionLedger + "." + f.Dlm.Host.ToHttpAddress() + } + return KvKeyDeletionLedger +} + +func deletionLedgerPartKey(key string, part int) []byte { + return []byte(fmt.Sprintf("%s.part.%05d", key, part)) +} + +// deletionLedgerGenPartKey returns the part key for one generation. Generation +// 0 keeps the original format so ledgers written before generation publishing +// still decode. +func deletionLedgerGenPartKey(key string, gen, part int) []byte { + if gen == 0 { + return deletionLedgerPartKey(key, part) + } + return []byte(fmt.Sprintf("%s.g%06d.part.%05d", key, gen, part)) +} + +// startDeletionLedgerSnapshotter periodically flushes the pending deletion set +// to durable storage. It also wakes on deletionLedgerFlush so a queued id is +// persisted within milliseconds instead of a full interval. Started from +// SetStore only after the ledger read succeeded. +func (f *Filer) startDeletionLedgerSnapshotter() { + go func() { + ticker := time.NewTicker(deletionPersistInterval()) + defer ticker.Stop() + for { + select { + case <-f.deletionQuit: + return + case <-ticker.C: + case <-f.deletionLedgerFlush: + } + // A close of deletionQuit may be concurrent with this wake; the + // final snapshot in Shutdown covers whatever the loop left dirty. + select { + case <-f.deletionQuit: + return + default: + } + f.snapshotDeletionLedger() + } + }() +} + +// signalLedgerFlush wakes the snapshotter without blocking the caller. +func (f *Filer) signalLedgerFlush() { + f.deletionLedgerLock.Lock() + if f.deletionLedgerFlush == nil { + f.deletionLedgerFlush = make(chan struct{}, 1) + } + ch := f.deletionLedgerFlush + f.deletionLedgerLock.Unlock() + select { + case ch <- struct{}{}: + default: + } +} + +// queueDeletions is the single entry point for adding fileIds to the deletion +// pipeline. It keeps the in-memory hot queue AND the durable ledger in sync. +// Safe on a zero-value Filer (tests build struct literals without NewFiler): +// the map is lazily created and the mutex is zero-value friendly. +func (f *Filer) queueDeletions(fileIds ...string) { + if len(fileIds) == 0 { + return + } + // Hot path: existing in-memory queue, unchanged. + f.FileIdDeletionQueue.EnQueue(fileIds...) + + // Durable ledger: record what still needs deleting. Every enqueue bumps the + // id's epoch so an expiry-forget carrying an older epoch cannot erase the + // re-queued intent. + f.deletionLedgerLock.Lock() + if f.pendingDeletions == nil { + f.pendingDeletions = make(map[string]uint64, len(fileIds)) + } + for _, id := range fileIds { + if id == "" { + continue + } + _, exists := f.pendingDeletions[id] + f.deletionSeq++ + f.pendingDeletions[id] = f.deletionSeq + if !exists { + f.deletionLedgerDirty = true + } + } + f.deletionLedgerLock.Unlock() + f.signalLedgerFlush() +} + +// pendingDeletionCount returns the number of fileIds still tracked as needing +// deletion in the durable ledger. Handy for diagnostics and tests; 0 when the +// map was never initialised. +func (f *Filer) pendingDeletionCount() int { + f.deletionLedgerLock.Lock() + defer f.deletionLedgerLock.Unlock() + return len(f.pendingDeletions) +} + +// deletionEpoch returns the pending record's current epoch for fileId. +func (f *Filer) deletionEpoch(fileId string) uint64 { + f.deletionLedgerLock.Lock() + defer f.deletionLedgerLock.Unlock() + return f.pendingDeletions[fileId] +} + +// forgetDeletion removes a fileId from the ledger once its deletion is terminal +// (deleted or already absent: the chunks are gone, so any pending record — even +// a re-queued one — refers to work that no longer exists). Non-terminal +// outcomes must NOT call this so the entry survives until confirmed. +func (f *Filer) forgetDeletion(fileId string) { + if fileId == "" { + return + } + f.deletionLedgerLock.Lock() + if _, exists := f.pendingDeletions[fileId]; exists { + delete(f.pendingDeletions, fileId) + f.deletionLedgerDirty = true + } + f.deletionLedgerLock.Unlock() + f.signalLedgerFlush() +} + +// forgetDeletionEpoch removes the ledger record only when its epoch still +// matches the one captured when the deletion attempt began. A mismatch means +// the id was re-queued meanwhile, so the newer record wins. Reports whether +// the record was removed. +func (f *Filer) forgetDeletionEpoch(fileId string, epoch uint64) bool { + if fileId == "" { + return false + } + removed := false + f.deletionLedgerLock.Lock() + if cur, exists := f.pendingDeletions[fileId]; exists && cur == epoch { + delete(f.pendingDeletions, fileId) + f.deletionLedgerDirty = true + removed = true + } + f.deletionLedgerLock.Unlock() + if removed { + f.signalLedgerFlush() + } + return removed +} + +// snapshotDeletionLedger serialises the current pending set to the store. +// Only writes when something changed since the last snapshot to keep KV churn low. +// Writes serialize on deletionSnapshotLock so a snapshot in flight when Shutdown +// starts cannot overwrite the final one with an older copy. +func (f *Filer) snapshotDeletionLedger() { + if !deletionPersistEnabled() || f.Store == nil { + return + } + f.deletionSnapshotLock.Lock() + defer f.deletionSnapshotLock.Unlock() + + if f.deletionLedgerBlocked.Load() && !f.tryUnblockDeletionLedger() { + return + } + + f.deletionLedgerLock.Lock() + if !f.deletionLedgerDirty { + hasStale := len(f.deletionLedgerStale) > 0 + f.deletionLedgerLock.Unlock() + if !hasStale { + return + } + // Nothing to republish, but orphaned part keys still need cleanup. + ctx := context.Background() + f.flushStaleLedgerParts(ctx) + f.persistStaleLedgerParts(ctx, f.deletionLedgerKey()) + return + } + // Take a stable copy under lock; marshal + KV write happen outside it. + ids := make([]string, 0, len(f.pendingDeletions)) + for id := range f.pendingDeletions { + ids = append(ids, id) + } + f.deletionLedgerDirty = false + f.deletionLedgerLock.Unlock() + + if err := f.writeDeletionLedger(ids); err != nil { + glog.Warningf("failed to persist deletion ledger (%d ids): %v", len(ids), err) + f.deletionLedgerLock.Lock() + f.deletionLedgerDirty = true + f.deletionLedgerLock.Unlock() + return + } + glog.V(3).Infof("persisted deletion ledger: %d pending deletions", len(ids)) +} + +// tryUnblockDeletionLedger retries the ledger reload that originally failed. +// Once the read succeeds the persisted ids merge back into the pending set and +// snapshots resume; until then nothing is written so the unread ledger cannot +// be overwritten. +func (f *Filer) tryUnblockDeletionLedger() bool { + f.reloadDeletionLedger() + return !f.deletionLedgerBlocked.Load() +} + +// writeDeletionLedger persists the id set. A single value is published +// atomically; a multipart set writes a new generation's part keys first and +// only then the manifest, so a crash or failed write leaves the previously +// published generation readable. Orphaned parts are retried on the next write. +func (f *Filer) writeDeletionLedger(ids []string) error { + ctx := context.Background() + key := f.deletionLedgerKey() + + f.deletionLedgerLock.Lock() + prevGen, prevParts := f.deletionLedgerGen, f.deletionLedgerParts + f.deletionLedgerLock.Unlock() + + f.flushStaleLedgerParts(ctx) + + var parts [][]byte + var batch []string + size := 2 // "[]" + flush := func() error { + if batch == nil { + batch = []string{} + } + payload, err := json.Marshal(batch) + if err != nil { + return err + } + parts = append(parts, payload) + batch = nil + size = 2 + return nil + } + for _, id := range ids { + if need := len(id) + 3; len(batch) > 0 && size+need > deletionLedgerPartSize { + if err := flush(); err != nil { + return err + } + } + batch = append(batch, id) + size += len(id) + 3 + } + if len(batch) > 0 || len(parts) == 0 { + if err := flush(); err != nil { + return err + } + } + + if len(parts) == 1 { + if err := f.Store.KvPut(ctx, []byte(key), parts[0]); err != nil { + return err + } + f.deleteLedgerParts(ctx, key, prevGen, prevParts) + f.deletionLedgerLock.Lock() + f.deletionLedgerGen, f.deletionLedgerParts = 0, 0 + f.deletionLedgerLock.Unlock() + } else { + gen := prevGen + 1 + wrote := 0 + for i, payload := range parts { + if err := f.Store.KvPut(ctx, deletionLedgerGenPartKey(key, gen, i), payload); err != nil { + f.markStaleLedgerParts(key, gen, wrote) + return err + } + wrote++ + } + manifest, _ := json.Marshal(struct { + Parts int `json:"parts"` + Gen int `json:"gen,omitempty"` + }{len(parts), gen}) + if err := f.Store.KvPut(ctx, []byte(key), manifest); err != nil { + f.markStaleLedgerParts(key, gen, wrote) + return err + } + f.deleteLedgerParts(ctx, key, prevGen, prevParts) + f.deletionLedgerLock.Lock() + f.deletionLedgerGen, f.deletionLedgerParts = gen, len(parts) + f.deletionLedgerLock.Unlock() + } + f.persistStaleLedgerParts(ctx, key) + f.touchLedgerIndex(ctx, key) + return nil +} + +// deleteLedgerParts removes a published generation's part keys; failures are +// tracked as stale so the next snapshot retries them. +func (f *Filer) deleteLedgerParts(ctx context.Context, key string, gen, parts int) { + for i := 0; i < parts; i++ { + if err := f.Store.KvDelete(ctx, deletionLedgerGenPartKey(key, gen, i)); err != nil { + f.deletionLedgerLock.Lock() + f.deletionLedgerStale = append(f.deletionLedgerStale, string(deletionLedgerGenPartKey(key, gen, i))) + f.deletionLedgerLock.Unlock() + } + } +} + +// markStaleLedgerParts remembers parts of an abandoned generation so they are +// deleted by a later snapshot instead of leaking store space. +func (f *Filer) markStaleLedgerParts(key string, gen, wrote int) { + if wrote == 0 { + return + } + f.deletionLedgerLock.Lock() + for i := 0; i < wrote; i++ { + f.deletionLedgerStale = append(f.deletionLedgerStale, string(deletionLedgerGenPartKey(key, gen, i))) + } + f.deletionLedgerLock.Unlock() +} + +func (f *Filer) flushStaleLedgerParts(ctx context.Context) { + f.deletionLedgerLock.Lock() + stale := f.deletionLedgerStale + f.deletionLedgerStale = nil + f.deletionLedgerLock.Unlock() + + var keep []string + for _, k := range stale { + if err := f.Store.KvDelete(ctx, []byte(k)); err != nil { + keep = append(keep, k) + } + } + if len(keep) > 0 { + f.deletionLedgerLock.Lock() + f.deletionLedgerStale = append(keep, f.deletionLedgerStale...) + f.deletionLedgerLock.Unlock() + } +} + +// persistStaleLedgerParts mirrors the stale-part list under a sidecar key so a +// crash between manifest publication and part cleanup can retry after restart +// instead of orphaning the part keys. +func (f *Filer) persistStaleLedgerParts(ctx context.Context, key string) { + f.deletionLedgerLock.Lock() + stale := append([]string(nil), f.deletionLedgerStale...) + f.deletionLedgerLock.Unlock() + staleKey := []byte(key + ".stale") + if len(stale) == 0 { + _ = f.Store.KvDelete(ctx, staleKey) + return + } + payload, _ := json.Marshal(stale) + _ = f.Store.KvPut(ctx, staleKey, payload) +} + +// loadStaleLedgerParts seeds the stale-part list left by a previous run. +func (f *Filer) loadStaleLedgerParts(ctx context.Context, key string) { + payload, err := f.Store.KvGet(ctx, []byte(key+".stale")) + if err != nil { + return + } + var stale []string + if json.Unmarshal(payload, &stale) != nil { + return + } + f.deletionLedgerLock.Lock() + f.deletionLedgerStale = append(f.deletionLedgerStale, stale...) + f.deletionLedgerLock.Unlock() +} + +// touchLedgerIndex records this filer's scoped key in the shared index so the +// ledger remains discoverable if the filer restarts under a new address. The +// index is a read-modify-write list with no CAS, so the write is verified and +// retried: a peer's concurrent update must not drop this key. +func (f *Filer) touchLedgerIndex(ctx context.Context, key string) { + if key == KvKeyDeletionLedger { + return + } + for attempt := 0; attempt < 3; attempt++ { + keys, err := f.readLedgerIndex(ctx) + if err != nil { + // Only a genuinely absent index means "empty"; a store error must + // not let us write a one-key list that drops every peer entry. + if err != ErrKvNotFound { + glog.V(1).Infof("deletion ledger index unreadable, skipping update: %v", err) + return + } + keys = nil + } + found := false + for _, k := range keys { + if k == key { + found = true + break + } + } + if found { + return + } + payload, _ := json.Marshal(append(keys, key)) + if err := f.Store.KvPut(ctx, []byte(KvKeyDeletionLedgerIndex), payload); err != nil { + glog.V(1).Infof("failed to update deletion ledger index: %v", err) + return + } + } + glog.V(1).Infof("deletion ledger index lost %q to a concurrent update; retrying on the next snapshot", key) +} + +func (f *Filer) readLedgerIndex(ctx context.Context) ([]string, error) { + payload, err := f.Store.KvGet(ctx, []byte(KvKeyDeletionLedgerIndex)) + if err != nil { + return nil, err + } + var keys []string + if err := json.Unmarshal(payload, &keys); err != nil { + return nil, err + } + return keys, nil +} + +func (f *Filer) pruneLedgerIndex(ctx context.Context, key string) { + keys, err := f.readLedgerIndex(ctx) + if err != nil { + return + } + kept := keys[:0] + for _, k := range keys { + if k != key { + kept = append(kept, k) + } + } + payload, _ := json.Marshal(kept) + _ = f.Store.KvPut(ctx, []byte(KvKeyDeletionLedgerIndex), payload) +} + +// readDeletionLedger reads the manifest key: a JSON array is the whole set; a +// {"parts":N,"gen":G} manifest points at that generation's part keys. A part +// the manifest references but the store lacks is corruption, not absence, so +// it surfaces as a wrapped error rather than ErrKvNotFound. +func (f *Filer) readDeletionLedger(key string) (ids []string, gen, parts int, err error) { + ctx := context.Background() + + payload, err := f.Store.KvGet(ctx, []byte(key)) + if err != nil { + return nil, 0, 0, err + } + if !bytes.HasPrefix(bytes.TrimSpace(payload), []byte("{")) { + if err := json.Unmarshal(payload, &ids); err != nil { + return nil, 0, 0, err + } + return ids, 0, 0, nil + } + var manifest struct { + Parts int `json:"parts"` + Gen int `json:"gen,omitempty"` + } + if err := json.Unmarshal(payload, &manifest); err != nil { + return nil, 0, 0, err + } + for i := 0; i < manifest.Parts; i++ { + partPayload, err := f.Store.KvGet(ctx, deletionLedgerGenPartKey(key, manifest.Gen, i)) + if err != nil { + return nil, 0, 0, fmt.Errorf("deletion ledger part %d of %s unreadable: %w", i, key, err) + } + var part []string + if err := json.Unmarshal(partPayload, &part); err != nil { + return nil, 0, 0, fmt.Errorf("deletion ledger part %d of %s corrupt: %w", i, key, err) + } + ids = append(ids, part...) + } + return ids, manifest.Gen, manifest.Parts, nil +} + +// recoverForeignDeletionLedgers runs when the scoped key is absent: the ledger +// may sit under the pre-scoping base key, or under a scoped key belonging to an +// earlier incarnation of this filer whose advertised address changed. Every +// discovered set is published under this filer's key first and the source keys +// deleted only after that write succeeds. A live peer's ledger claimed here is +// rewritten by the peer's next snapshot, so the ids only get processed twice — +// deletions are not owner-specific. +func (f *Filer) recoverForeignDeletionLedgers(key string) (ids []string, err error) { + ctx := context.Background() + + claimed := map[string]ledgerClaim{} // source key -> what was read + lIds, lGen, lParts, lErr := f.readDeletionLedger(KvKeyDeletionLedger) + switch { + case lErr == nil: + ids = append(ids, lIds...) + claimed[KvKeyDeletionLedger] = ledgerClaim{lIds, lGen, lParts} + f.loadStaleLedgerParts(ctx, KvKeyDeletionLedger) + case lErr != ErrKvNotFound: + return nil, lErr + } + + // Any read failure aborts the whole claim so reload retries: skipping an + // unreadable source would strand it, since a successful claim here means + // the index is never searched again. + indexKeys, idxErr := f.readLedgerIndex(ctx) + if idxErr != nil && idxErr != ErrKvNotFound { + return nil, idxErr + } + for _, other := range indexKeys { + if other == key { + continue + } + oIds, oGen, oParts, oErr := f.readDeletionLedger(other) + if oErr == ErrKvNotFound { + f.pruneLedgerIndex(ctx, other) + continue + } + if oErr != nil { + return nil, fmt.Errorf("deletion ledger %s unreadable: %w", other, oErr) + } + ids = append(ids, oIds...) + claimed[other] = ledgerClaim{oIds, oGen, oParts} + f.loadStaleLedgerParts(ctx, other) + } + + if len(claimed) == 0 { + return nil, ErrKvNotFound + } + if err := f.writeDeletionLedger(ids); err != nil { + return nil, err + } + mergedMore := false + for src, claim := range claimed { + // A source republished since we read it belongs to a live filer + // (or a racing claimer): merge the newer ids and leave it in place. + curIds, curGen, curParts, rerr := f.readDeletionLedger(src) + if rerr != nil || curGen != claim.gen || curParts != claim.parts || !sameStringSet(curIds, claim.ids) { + if rerr == nil { + for _, id := range curIds { + if !containsId(ids, id) { + ids = append(ids, id) + mergedMore = true + } + } + } + glog.V(0).Infof("deletion ledger %s changed while claiming; leaving it for its owner", src) + continue + } + f.deleteLedgerParts(ctx, src, claim.gen, claim.parts) + _ = f.Store.KvDelete(ctx, []byte(src)) + _ = f.Store.KvDelete(ctx, []byte(src+".stale")) + f.pruneLedgerIndex(ctx, src) + } + if mergedMore { + // Ids from a changed source are durable only under that source's key; + // rewrite our ledger so they survive under ours too. + if err := f.writeDeletionLedger(ids); err != nil { + return nil, err + } + } + return ids, nil +} + +func containsId(ids []string, id string) bool { + for _, s := range ids { + if s == id { + return true + } + } + return false +} + +type ledgerClaim struct { + ids []string + gen int + parts int +} + +func sameStringSet(a, b []string) bool { + if len(a) != len(b) { + return false + } + seen := make(map[string]struct{}, len(a)) + for _, s := range a { + seen[s] = struct{}{} + } + for _, s := range b { + if _, ok := seen[s]; !ok { + return false + } + } + return true +} + +// reloadDeletionLedger re-enqueues any pending deletions found in the store after +// a restart, so a crash that killed the in-memory queues does not leak chunks. +// Recovered ids join the pending set immediately so snapshots rewrite the full +// ledger; only the re-queue waits out DeletionRecoveryGrace. +// +// It reports whether the ledger is usable. A read error or an unparseable +// payload leaves the persisted set unknown, so snapshots retry the read rather +// than overwrite the unread ledger with a partial set. +// +// Safe to call on a Filer with a nil Store (no-op). Idempotent for the volume +// side: a chunk that was actually deleted before the crash re-deletes as not-found. +func (f *Filer) reloadDeletionLedger() bool { + if !deletionPersistEnabled() || f.Store == nil { + return false + } + + key := f.deletionLedgerKey() + ids, gen, parts, err := f.readDeletionLedger(key) + stateWritten := false + if err == ErrKvNotFound && key != KvKeyDeletionLedger { + ids, err = f.recoverForeignDeletionLedgers(key) + stateWritten = err == nil + } + if err != nil { + if err == ErrKvNotFound { + f.deletionLedgerBlocked.Store(false) + return true + } + f.deletionLedgerBlocked.Store(true) + glog.Warningf("failed to read persisted deletion ledger; persistence disabled until it reads: %v", err) + return false + } + f.deletionLedgerBlocked.Store(false) + f.loadStaleLedgerParts(context.Background(), key) + + // Merge into the pending set now so an early snapshot rewrites the + // recovered ids instead of overwriting the ledger with only new ones. + f.deletionLedgerLock.Lock() + if len(ids) > 0 { + if f.pendingDeletions == nil { + f.pendingDeletions = make(map[string]uint64, len(ids)) + } + for _, id := range ids { + if id != "" { + if _, exists := f.pendingDeletions[id]; !exists { + f.deletionSeq++ + f.pendingDeletions[id] = f.deletionSeq + } + } + } + } + if !stateWritten { + f.deletionLedgerGen = gen + f.deletionLedgerParts = parts + } + f.deletionLedgerLock.Unlock() + + if len(ids) == 0 { + return true + } + + grace := deletionRecoveryGrace() + glog.V(0).Infof("recovered %d pending deletions from ledger, applying in %v", len(ids), grace) + + go func() { + timer := time.NewTimer(grace) + defer timer.Stop() + select { + case <-f.deletionQuit: + return + case <-timer.C: + } + // The ids are already in the pending set; only the queue push waits + // for peer meta-aggregation to settle. The ledger is deliberately NOT + // cleared here: it shrinks only as the delete pipeline confirms each + // id terminal, so a second crash mid-recovery replays everything. + f.queueDeletions(ids...) + glog.V(0).Infof("re-queued %d recovered pending deletions", len(ids)) + }() + return true +} diff --git a/weed/filer/filer_deletion_persist_test.go b/weed/filer/filer_deletion_persist_test.go new file mode 100644 index 000000000..d9f82f54b --- /dev/null +++ b/weed/filer/filer_deletion_persist_test.go @@ -0,0 +1,584 @@ +package filer + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/cluster/lock_manager" + "github.com/seaweedfs/seaweedfs/weed/util" +) + +// readPersistedLedger pulls the persisted ledger payload straight from the store +// the filer is wired to, the same way a restarting filer would. +func readPersistedLedger(t *testing.T, f *Filer) ([]string, bool) { + t.Helper() + raw, err := f.Store.KvGet(context.Background(), []byte(f.deletionLedgerKey())) + if err != nil { + if err == ErrKvNotFound { + return nil, false + } + t.Fatalf("unexpected error reading ledger: %v", err) + } + var ids []string + if err := json.Unmarshal(raw, &ids); err != nil { + t.Fatalf("persisted payload not valid JSON: %v", err) + } + return ids, true +} + +// withPersistedConfig pins the ledger knobs at viper's override tier so the test +// sees deterministic semantics regardless of global SetDefault ordering (the +// production getters SetDefault(true), which would otherwise win the tier race). +// The overrides are restored on cleanup so later tests see production defaults. +func withPersistedConfig(t *testing.T, enabled bool) { + t.Helper() + v := util.GetViper() + v.Set("filer.deleteQueue.persist", enabled) + v.Set("filer.deleteQueue.recoveryGrace", defaultDeletionRecoveryGrace) + t.Cleanup(func() { + v.Set("filer.deleteQueue.persist", true) + v.Set("filer.deleteQueue.recoveryGrace", defaultDeletionRecoveryGrace) + }) +} + +// newLedgerTestFiler builds a Filer with the pieces the ledger touches, backed by +// a real (stub) store, bypassing NewFiler's master/aggregator wiring. +func newLedgerTestFiler(store FilerStore) *Filer { + return &Filer{ + FileIdDeletionQueue: util.NewUnboundedQueue(), + DeletionRetryQueue: NewDeletionRetryQueue(), + deletionQuit: make(chan struct{}), + Store: NewFilerStoreWrapper(store), + } +} + +// TestDeletionLedgerSnapshotAndRecover is the core guarantee: enqueued-but- +// unconfirmed deletions survive a simulated process restart and land back in the +// hot queue (minus whatever was confirmed terminal in the meantime). +func TestDeletionLedgerSnapshotAndRecover(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + f := newLedgerTestFiler(store) + + f.queueDeletions("1,01", "1,02", "1,03") + f.snapshotDeletionLedger() + + persisted, ok := readPersistedLedger(t, f) + if !ok { + t.Fatalf("expected persisted ledger under %q", KvKeyDeletionLedger) + } + if len(persisted) != 3 { + t.Fatalf("expected 3 persisted ids, got %d: %v", len(persisted), persisted) + } + for _, want := range []string{"1,01", "1,02", "1,03"} { + if !contains(persisted, want) { + t.Errorf("expected %q in persisted ledger, got %v", want, persisted) + } + } + + // Confirmed-gone for one id, then snapshot again. + f.forgetDeletion("1,02") + f.snapshotDeletionLedger() + persisted, _ = readPersistedLedger(t, f) + if len(persisted) != 2 { + t.Fatalf("expected 2 persisted ids after forget, got %d: %v", len(persisted), persisted) + } + if contains(persisted, "1,02") { + t.Errorf("forgotten id still in persisted ledger: %v", persisted) + } + + // --- simulated restart: a fresh Filer with NO ledger in memory must + // recover the still-pending ids from the store and re-queue them. --- + f2 := newLedgerTestFiler(store) + // recoveryGrace 0 so the background re-enqueue is near-instant. + util.GetViper().Set("filer.deleteQueue.recoveryGrace", time.Duration(0)) + f2.reloadDeletionLedger() + + // Wait for the async recovery goroutine to drain into the hot queue. + seen := make(map[string]bool) + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + f2.FileIdDeletionQueue.Consume(func(ids []string) { + for _, id := range ids { + seen[id] = true + } + }) + if len(seen) >= 2 { + break + } + time.Sleep(10 * time.Millisecond) + } + + if !seen["1,01"] || !seen["1,03"] { + t.Fatalf("expected recovered ids 1,01 and 1,03 back in hot queue, got %v", seen) + } + if seen["1,02"] { + t.Errorf("forgotten id 1,02 must NOT come back: %v", seen) + } +} + +// TestDeletionLedgerRetryKeepsEntry pins the key semantic: a non-terminal +// (retryable) failure must keep the entry persisted, so a crash mid-retry does +// not orphan the chunk. This is the whole reason forgetDeletion is only called on +// success / not-found / permanent. +func TestDeletionLedgerRetryKeepsEntry(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + f := newLedgerTestFiler(store) + + f.queueDeletions("9,01") + f.snapshotDeletionLedger() + + // A retryable outcome deliberately does NOT call forgetDeletion, so the id + // must still be in the durable ledger. + persisted, ok := readPersistedLedger(t, f) + if !ok { + t.Fatalf("ledger missing — retryable ids must remain persisted") + } + if len(persisted) != 1 || persisted[0] != "9,01" { + t.Fatalf("retryable id must stay in ledger, got %v", persisted) + } +} + +// TestDeletionLedgerDisabled verifies the kill switch: with persist off, nothing +// reaches KV, and reload refuses to recover even if a ledger exists in the store +// (so an operator turning the feature off never resurrects a stale ledger). +func TestDeletionLedgerDisabled(t *testing.T) { + // First, with persistence ON, put a real ledger into the store. + withPersistedConfig(t, true) + store := newStubFilerStore() + fOn := newLedgerTestFiler(store) + fOn.queueDeletions("5,01") + fOn.snapshotDeletionLedger() + if _, ok := readPersistedLedger(t, fOn); !ok { + t.Fatalf("setup: expected a persisted ledger with persist enabled") + } + + // Now flip the switch OFF and check the contract. + withPersistedConfig(t, false) + + // 1. Snapshot does not touch KV (no new writes, no clear). + f := newLedgerTestFiler(store) + f.queueDeletions("6,01") // in-memory only + f.snapshotDeletionLedger() + persisted, ok := readPersistedLedger(t, f) + if !ok { + t.Fatalf("setup: ledger should still exist from the enabled phase") + } + // The enabled-phase ledger had "5,01"; disabled phase must not have added "6,01". + if contains(persisted, "6,01") { + t.Errorf("persist disabled but snapshot wrote 6,01: %v", persisted) + } + + // 2. Reload with the switch off must NOT recover, even though a ledger exists. + f2 := newLedgerTestFiler(store) + f2.reloadDeletionLedger() + if f2.pendingDeletionCount() != 0 { + t.Fatalf("reload must be a no-op when disabled, got %d pending", f2.pendingDeletionCount()) + } +} + +// Filers sharing one metadata store must not overwrite each other's ledgers: +// each filer keys its ledger by its own address. +func TestDeletionLedgerScopedPerFiler(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + fA := newLedgerTestFiler(store) + fA.Dlm = lock_manager.NewDistributedLockManager("filer-a:8888") + fB := newLedgerTestFiler(store) + fB.Dlm = lock_manager.NewDistributedLockManager("filer-b:8888") + + fA.queueDeletions("1,01") + fB.queueDeletions("2,02") + fA.snapshotDeletionLedger() + fB.snapshotDeletionLedger() + + idsA, okA := readPersistedLedger(t, fA) + idsB, okB := readPersistedLedger(t, fB) + if !okA || !okB { + t.Fatalf("both filers must have their own ledger, got %v %v", idsA, idsB) + } + if !contains(idsA, "1,01") || contains(idsA, "2,02") { + t.Fatalf("filer A ledger wrong: %v", idsA) + } + if !contains(idsB, "2,02") || contains(idsB, "1,01") { + t.Fatalf("filer B ledger wrong: %v", idsB) + } + + // A restarted filer on B's address recovers only B's pending deletions. + fB2 := newLedgerTestFiler(store) + fB2.Dlm = lock_manager.NewDistributedLockManager("filer-b:8888") + if !fB2.reloadDeletionLedger() { + t.Fatalf("reload should succeed") + } + if fB2.pendingDeletionCount() != 1 { + t.Fatalf("B's restart should recover exactly its own id, got %d", fB2.pendingDeletionCount()) + } +} + +// A ledger bigger than one store value must still persist: it splits into part +// keys under a manifest, and shrinking below the part size removes the parts. +func TestDeletionLedgerChunked(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + f := newLedgerTestFiler(store) + + var ids []string + for i := 0; i < 4000; i++ { + ids = append(ids, fmt.Sprintf("7,%08xbeefcafe", i)) + } + f.queueDeletions(ids...) + f.snapshotDeletionLedger() + + raw, err := f.Store.KvGet(context.Background(), []byte(f.deletionLedgerKey())) + if err != nil || len(raw) == 0 || raw[0] != '{' { + t.Fatalf("expected a chunked manifest, got %v %q", err, raw) + } + + f2 := newLedgerTestFiler(store) + if !f2.reloadDeletionLedger() { + t.Fatalf("chunked ledger should reload") + } + if got := f2.pendingDeletionCount(); got != len(ids) { + t.Fatalf("recovered %d ids, want %d", got, len(ids)) + } + + // Forget everything; the next snapshot is a single value again and the + // part keys are removed. + for _, id := range ids { + f2.forgetDeletion(id) + } + f2.snapshotDeletionLedger() + raw, err = f.Store.KvGet(context.Background(), []byte(f2.deletionLedgerKey())) + if err != nil || len(raw) == 0 || raw[0] != '[' { + t.Fatalf("expected single-value ledger after shrink, got %v %q", err, raw) + } + if _, err := f.Store.KvGet(context.Background(), deletionLedgerPartKey(f2.deletionLedgerKey(), 0)); err != ErrKvNotFound { + t.Fatalf("stale part key must be removed, got %v", err) + } +} + +// A startup ledger read that fails for a reason other than not-found must not +// let snapshots overwrite the unread ledger with a partial set. +func TestDeletionLedgerReadFailureBlocksPersistence(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + f := newLedgerTestFiler(store) + store.kvGetErr = errors.New("kv read down") + + if f.reloadDeletionLedger() { + t.Fatalf("reload should report the ledger as unusable") + } + + f.queueDeletions("1,01") + f.snapshotDeletionLedger() + if len(store.kv) != 0 { + t.Fatalf("snapshot must not overwrite a ledger that was never read, wrote %v", store.kv) + } + + // The block is not permanent: once the store reads again, persistence + // resumes — a transient failure must not disable the ledger for the life + // of the process. + store.kvGetErr = nil + f.snapshotDeletionLedger() + persisted, ok := readPersistedLedger(t, f) + if !ok || len(persisted) != 1 || persisted[0] != "1,01" { + t.Fatalf("persistence should resume once the ledger reads, got %v", persisted) + } +} + +// Recovered ids join the pending set immediately — before the grace delay — so +// an early snapshot rewrites the recovered ids instead of dropping them. +func TestDeletionLedgerRecoveryKeepsIdsPending(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + f := newLedgerTestFiler(store) + f.queueDeletions("1,01", "1,02") + f.snapshotDeletionLedger() + + util.GetViper().Set("filer.deleteQueue.recoveryGrace", time.Hour) + f2 := newLedgerTestFiler(store) + if !f2.reloadDeletionLedger() { + t.Fatalf("reload should succeed") + } + // Still inside the grace window, but the ids are already pending. + if f2.pendingDeletionCount() != 2 { + t.Fatalf("recovered ids must be pending immediately, got %d", f2.pendingDeletionCount()) + } + // And an early shutdown snapshot keeps them. + f2.snapshotDeletionLedger() + persisted, ok := readPersistedLedger(t, f2) + if !ok || len(persisted) != 2 { + t.Fatalf("snapshot during grace must carry recovered ids, got %v", persisted) + } +} + +// TestDeletionLedgerZeroValueFiler proves nil-map safety for the struct literals +// used across the test suite (they never call NewFiler). +func TestDeletionLedgerZeroValueFiler(t *testing.T) { + withPersistedConfig(t, true) + f := &Filer{ + FileIdDeletionQueue: util.NewUnboundedQueue(), + DeletionRetryQueue: NewDeletionRetryQueue(), + deletionQuit: make(chan struct{}), + } + f.queueDeletions("0,01") + if f.pendingDeletionCount() != 1 { + t.Fatalf("zero-value filer should lazily init the ledger, got %d", f.pendingDeletionCount()) + } + f.forgetDeletion("0,01") + if f.pendingDeletionCount() != 0 { + t.Fatalf("forget on zero-value filer failed, got %d", f.pendingDeletionCount()) + } +} + +// A manifest that references a part the store cannot return is corruption, not +// absence: reload must report failure and nothing may overwrite the surviving +// manifest with only the in-memory set. +func TestDeletionLedgerMissingPartBlocks(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + f := newLedgerTestFiler(store) + + key := f.deletionLedgerKey() + store.kv[key] = []byte(`{"parts":2,"gen":4}`) + store.kv[string(deletionLedgerGenPartKey(key, 4, 0))] = []byte(`["1,01"]`) + // Part 1 of generation 4 is missing. + + if f.reloadDeletionLedger() { + t.Fatalf("a manifest referencing a missing part is corrupt, not absent") + } + f.queueDeletions("9,99") + f.snapshotDeletionLedger() + if got := string(store.kv[key]); got != `{"parts":2,"gen":4}` { + t.Fatalf("unread ledger must not be overwritten, got %q", got) + } +} + +// A filer restarting under a new advertised address must still find its +// previous ledger: the index lists scoped keys, the ids are published under +// the new key first, and only then is the stranded key removed. +func TestDeletionLedgerAddressChangeClaim(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + fOld := newLedgerTestFiler(store) + fOld.Dlm = lock_manager.NewDistributedLockManager("filer-old:8888") + fOld.queueDeletions("1,01", "1,02") + fOld.snapshotDeletionLedger() + + fNew := newLedgerTestFiler(store) + fNew.Dlm = lock_manager.NewDistributedLockManager("filer-new:9999") + if !fNew.reloadDeletionLedger() { + t.Fatalf("reload should claim the stranded ledger") + } + if fNew.pendingDeletionCount() != 2 { + t.Fatalf("expected 2 claimed ids, got %d", fNew.pendingDeletionCount()) + } + if _, err := store.KvGet(context.Background(), []byte(fOld.deletionLedgerKey())); err != ErrKvNotFound { + t.Fatalf("stranded ledger must be removed after claim, got %v", err) + } + persisted, ok := readPersistedLedger(t, fNew) + if !ok || len(persisted) != 2 { + t.Fatalf("claimed ids must be durably stored under the new key, got %v", persisted) + } + indexKeys, _ := fNew.readLedgerIndex(context.Background()) + if contains(indexKeys, fOld.deletionLedgerKey()) { + t.Fatalf("claimed key must leave the index, got %v", indexKeys) + } +} + +// Every multipart snapshot writes a fresh generation's part keys and publishes +// the manifest only after all parts land, so the committed generation is never +// overwritten. The previous generation is removed after publication. +func TestDeletionLedgerGenerationPublish(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + f := newLedgerTestFiler(store) + + var ids []string + for i := 0; i < 4000; i++ { + ids = append(ids, fmt.Sprintf("7,%08xbeefcafe", i)) + } + f.queueDeletions(ids...) + f.snapshotDeletionLedger() + key := f.deletionLedgerKey() + if _, err := store.KvGet(context.Background(), deletionLedgerGenPartKey(key, 1, 0)); err != nil { + t.Fatalf("generation 1 part must exist, got %v", err) + } + + f.queueDeletions("8,08") + f.snapshotDeletionLedger() + if _, err := store.KvGet(context.Background(), deletionLedgerGenPartKey(key, 2, 0)); err != nil { + t.Fatalf("generation 2 part must exist, got %v", err) + } + if _, err := store.KvGet(context.Background(), deletionLedgerGenPartKey(key, 1, 0)); err != ErrKvNotFound { + t.Fatalf("generation 1 part must be cleaned after publication, got %v", err) + } + var manifest struct { + Parts int `json:"parts"` + Gen int `json:"gen"` + } + raw, _ := store.KvGet(context.Background(), []byte(key)) + if err := json.Unmarshal(raw, &manifest); err != nil || manifest.Gen != 2 { + t.Fatalf("manifest must publish generation 2, got %q", raw) + } +} + +// A failed cleanup of a superseded generation is retried by the next snapshot +// instead of leaking the part keys. +func TestDeletionLedgerStalePartsRetried(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + f := newLedgerTestFiler(store) + + var ids []string + for i := 0; i < 4000; i++ { + ids = append(ids, fmt.Sprintf("7,%08xbeefcafe", i)) + } + f.queueDeletions(ids...) + f.snapshotDeletionLedger() + key := f.deletionLedgerKey() + + store.kvDeleteErr = errors.New("delete down") + f.queueDeletions("8,08") + f.snapshotDeletionLedger() + if len(f.deletionLedgerStale) == 0 { + t.Fatalf("failed part deletes must be tracked for retry") + } + store.kvDeleteErr = nil + f.queueDeletions("8,09") + f.snapshotDeletionLedger() + if len(f.deletionLedgerStale) != 0 { + t.Fatalf("stale parts must flush once deletes work, got %v", f.deletionLedgerStale) + } + if _, err := store.KvGet(context.Background(), deletionLedgerGenPartKey(key, 1, 0)); err != ErrKvNotFound { + t.Fatalf("stale generation part must be deleted, got %v", err) + } +} + +// A pre-scoping ledger under the base key must migrate only after the scoped +// copy is durable: the recovered ids are written under the scoped key first, +// then the base key is removed. +func TestDeletionLedgerLegacyMigrationDurable(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + store.kv[KvKeyDeletionLedger] = []byte(`["1,01","1,02"]`) + + f := newLedgerTestFiler(store) + f.Dlm = lock_manager.NewDistributedLockManager("filer-new:9999") + if !f.reloadDeletionLedger() { + t.Fatalf("reload should migrate the legacy ledger") + } + persisted, ok := readPersistedLedger(t, f) + if !ok || len(persisted) != 2 { + t.Fatalf("legacy ids must land under the scoped key first, got %v", persisted) + } + if _, err := store.KvGet(context.Background(), []byte(KvKeyDeletionLedger)); err != ErrKvNotFound { + t.Fatalf("legacy key must be removed once the scoped write is durable, got %v", err) + } + if f.pendingDeletionCount() != 2 { + t.Fatalf("migrated ids must join the pending set, got %d", f.pendingDeletionCount()) + } +} + +// An indexed ledger that cannot be read must not be skipped: claiming the +// readable ones and leaving the unreadable one behind would strand it, since +// the index is only searched while the filer's own key is absent. The claim +// aborts so the reload is retried. +func TestDeletionLedgerUnreadableForeignBlocks(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + stranded := KvKeyDeletionLedger + ".filer-old:8888" + store.kv[stranded] = []byte(`{"parts":2,"gen":1}`) // manifest without parts = corrupt + index, _ := json.Marshal([]string{stranded}) + store.kv[KvKeyDeletionLedgerIndex] = index + + f := newLedgerTestFiler(store) + f.Dlm = lock_manager.NewDistributedLockManager("filer-new:9999") + if f.reloadDeletionLedger() { + t.Fatalf("an unreadable indexed ledger must fail the reload") + } + if _, err := store.KvGet(context.Background(), []byte(stranded)); err != nil { + t.Fatalf("unreadable ledger must be left untouched, got %v", err) + } + if _, ok := readPersistedLedger(t, f); ok { + t.Fatalf("no scoped ledger may be written while a claim source is unreadable") + } +} + +// A live filer's ledger is not deleted when it republished between the claim's +// read and cleanup: the source stays and its newer ids merge into the claim. +// The hook swaps the source content on its second read — the claim's first +// read sees the old set, the verification read sees the republished one, +// exactly like a peer snapshot landing mid-claim. +func TestDeletionLedgerClaimKeepsChangedSource(t *testing.T) { + withPersistedConfig(t, true) + store := newStubFilerStore() + srcKey := KvKeyDeletionLedger + ".filer-old:8888" + store.kv[srcKey] = []byte(`["1,01","1,02"]`) + index, _ := json.Marshal([]string{srcKey}) + store.kv[KvKeyDeletionLedgerIndex] = index + + srcReads := 0 + store.kvGetHook = func(key []byte) { + if string(key) == srcKey { + srcReads++ + if srcReads == 2 { + store.kv[srcKey] = []byte(`["1,01","1,02","1,03"]`) + } + } + } + + f := newLedgerTestFiler(store) + f.Dlm = lock_manager.NewDistributedLockManager("filer-new:9999") + if !f.reloadDeletionLedger() { + t.Fatalf("claim should succeed") + } + if got := string(store.kv[srcKey]); got != `["1,01","1,02","1,03"]` { + t.Fatalf("republished source must be left intact, got %q", got) + } + if f.pendingDeletionCount() != 3 { + t.Fatalf("republished ids must merge into the pending set, got %d", f.pendingDeletionCount()) + } + persisted, ok := readPersistedLedger(t, f) + if !ok || len(persisted) != 3 { + t.Fatalf("merged ids must be durable under our key, got %v", persisted) + } +} + +// An expired retry item must not erase a newer enqueue for the same file id: +// expiry forgets only the epoch the retry item recorded. +func TestForgetDeletionEpochSkipsNewer(t *testing.T) { + withPersistedConfig(t, true) + f := newLedgerTestFiler(newStubFilerStore()) + + f.queueDeletions("1,01") + oldEpoch := f.deletionEpoch("1,01") + f.queueDeletions("1,01") // re-enqueue bumps the epoch + newEpoch := f.deletionEpoch("1,01") + if oldEpoch == newEpoch { + t.Fatalf("re-enqueue must bump the epoch") + } + + f.forgetDeletionEpoch("1,01", oldEpoch) + if f.pendingDeletionCount() != 1 { + t.Fatalf("stale expiry must not erase a newer enqueue") + } + f.forgetDeletionEpoch("1,01", newEpoch) + if f.pendingDeletionCount() != 0 { + t.Fatalf("matching expiry must forget the id") + } +} + +func contains(haystack []string, needle string) bool { + for _, s := range haystack { + if s == needle { + return true + } + } + return false +} diff --git a/weed/filer/filer_deletion_test.go b/weed/filer/filer_deletion_test.go index 77ac2310f..f1ac3343d 100644 --- a/weed/filer/filer_deletion_test.go +++ b/weed/filer/filer_deletion_test.go @@ -10,15 +10,15 @@ func TestDeletionRetryQueue_AddAndRetrieve(t *testing.T) { queue := NewDeletionRetryQueue() // Add items - queue.AddOrUpdate("file1", "is read only") - queue.AddOrUpdate("file2", "connection reset") + queue.AddOrUpdate("file1", "is read only", 0) + queue.AddOrUpdate("file2", "connection reset", 0) if queue.Size() != 2 { t.Errorf("Expected queue size 2, got %d", queue.Size()) } // Items not ready yet (initial delay is 5 minutes) - readyItems := queue.GetReadyItems(10) + readyItems, _ := queue.GetReadyItems(10) if len(readyItems) != 0 { t.Errorf("Expected 0 ready items, got %d", len(readyItems)) } @@ -104,7 +104,7 @@ func TestDeletionRetryQueue_MaxAttemptsReached(t *testing.T) { queue := NewDeletionRetryQueue() // Add item - queue.AddOrUpdate("file1", "error") + queue.AddOrUpdate("file1", "error", 0) // Manually set retry count to max queue.lock.Lock() @@ -119,7 +119,7 @@ func TestDeletionRetryQueue_MaxAttemptsReached(t *testing.T) { queue.lock.Unlock() // Try to get ready items - should be returned for the last retry (attempt #10) - readyItems := queue.GetReadyItems(10) + readyItems, _ := queue.GetReadyItems(10) if len(readyItems) != 1 { t.Fatalf("Expected 1 item for last retry, got %d", len(readyItems)) } @@ -139,7 +139,7 @@ func TestDeletionRetryQueue_MaxAttemptsReached(t *testing.T) { queue.lock.Unlock() // Now it should be discarded (retry count is 11, exceeds max of 10) - readyItems = queue.GetReadyItems(10) + readyItems, _ = queue.GetReadyItems(10) if len(readyItems) != 0 { t.Errorf("Expected 0 items (max attempts exceeded), got %d", len(readyItems)) } @@ -248,7 +248,7 @@ func TestDeletionRetryQueue_HeapOrdering(t *testing.T) { queue.lock.Unlock() // GetReadyItems should return in NextRetryAt order - readyItems := queue.GetReadyItems(10) + readyItems, _ := queue.GetReadyItems(10) expectedOrder := []string{"file1", "file2", "file3"} if len(readyItems) != 3 { @@ -266,7 +266,7 @@ func TestDeletionRetryQueue_DuplicateFileIds(t *testing.T) { queue := NewDeletionRetryQueue() // Add same file ID twice with retryable error - simulates duplicate in batch - queue.AddOrUpdate("file1", "timeout error") + queue.AddOrUpdate("file1", "timeout error", 0) // Verify only one item exists in queue if queue.Size() != 1 { @@ -284,7 +284,7 @@ func TestDeletionRetryQueue_DuplicateFileIds(t *testing.T) { queue.lock.Unlock() // Add same file ID again - should NOT increment retry count (just update error) - queue.AddOrUpdate("file1", "timeout error again") + queue.AddOrUpdate("file1", "timeout error again", 0) // Verify still only one item exists in queue (not duplicated) if queue.Size() != 1 { @@ -306,3 +306,27 @@ func TestDeletionRetryQueue_DuplicateFileIds(t *testing.T) { t.Errorf("Expected LastError to be updated to 'timeout error again', got %q", item2.LastError) } } + +// AddOrUpdate must not overwrite the ledger epoch on an in-flight item: the +// worker that popped it reads that field without the queue lock, and its +// expiry/permanent-forget must only match the record the attempt started with. +func TestDeletionRetryQueue_InFlightKeepsEpoch(t *testing.T) { + queue := NewDeletionRetryQueue() + queue.AddOrUpdate("file1", "timeout", 7) + + queue.lock.Lock() + item := queue.itemIndex["file1"] + item.NextRetryAt = time.Now().Add(-time.Second) + heap.Init(&queue.heap) + queue.lock.Unlock() + + ready, _ := queue.GetReadyItems(1) + if len(ready) != 1 { + t.Fatalf("expected the item ready, got %d", len(ready)) + } + + queue.AddOrUpdate("file1", "newer error", 42) + if got := ready[0].ledgerEpoch; got != 7 { + t.Fatalf("in-flight epoch must stay 7, got %d", got) + } +} diff --git a/weed/filer/filer_lazy_remote_test.go b/weed/filer/filer_lazy_remote_test.go index 468142613..04798c3f2 100644 --- a/weed/filer/filer_lazy_remote_test.go +++ b/weed/filer/filer_lazy_remote_test.go @@ -35,6 +35,9 @@ type stubFilerStore struct { kv map[string][]byte insertErr error findErr error + kvGetErr error + kvDeleteErr error + kvGetHook func(key []byte) deleteErrByPath map[string]error } @@ -63,6 +66,12 @@ func (s *stubFilerStore) KvPut(_ context.Context, key []byte, value []byte) erro func (s *stubFilerStore) KvGet(_ context.Context, key []byte) ([]byte, error) { s.mu.Lock() defer s.mu.Unlock() + if s.kvGetHook != nil { + s.kvGetHook(key) + } + if s.kvGetErr != nil { + return nil, s.kvGetErr + } value, found := s.kv[string(key)] if !found { return nil, ErrKvNotFound @@ -72,6 +81,9 @@ func (s *stubFilerStore) KvGet(_ context.Context, key []byte) ([]byte, error) { func (s *stubFilerStore) KvDelete(_ context.Context, key []byte) error { s.mu.Lock() defer s.mu.Unlock() + if s.kvDeleteErr != nil { + return s.kvDeleteErr + } delete(s.kv, string(key)) return nil }