filer: persist pending chunk deletions across restarts (durable deletion ledger) (#11550)

* fix(filer): persist pending chunk deletions across restarts

The in-memory FileIdDeletionQueue and DeletionRetryQueue lose every
queued-but-unconfirmed deletion when the filer process restarts. Because
deletions only enter the pipeline through that queue, a crash between
enqueue and the volume confirming the delete leaks the chunk permanently:
nothing remembers it. In a multi-filer deployment this was observed as
growing collections of orphaned chunks after filer restarts, and — via
meta-replay from a peer that still had the entry — orphans being
"resurrected" as live references on the recovered filer.

This implements the "periodic snapshot with recovery on startup" option
noted in the existing DeletionRetryQueue TODO, using the store's KV layer
(no new iterator API required across the 15+ store backends):

- queueDeletions() is the single entry point that keeps the hot in-memory
  queue and the durable ledger in sync.
- Only terminal outcomes (success / not-found / permanent) remove an id
  from the ledger; retryable failures keep it, which is the point.
- A timer and Shutdown() snapshot the pending set to a single KV key.
- On startup, reloadDeletionLedger() re-queues recovered ids after a
  grace window so the initial peer meta-aggregation settles first. This
  avoids a new hazard: purging a chunk that a lagging peer is about to
  re-reference as live data (stale replay turns a stale read into a
  dangling read otherwise).
- Volume deletes are idempotent (not-found == success), so re-deleting
  after a crash never double-frees.
- Kill switch via viper: filer.deleteQueue.persist=false opts out entirely
  (reload also refuses to recover so a stale ledger never comes back).
  Tunables: filer.deleteQueue.persistInterval, .recoveryGrace.

Adds unit tests covering snapshot+recover, retry-keeps-entry, disabled
switch, and zero-value Filer safety (run green under -race).

Co-Authored-By: Athena 🏛️ <hermes-agent@local> (custom / Qwen3.8-Flash-Next-ROCmFP4)

* filer: harden the deletion ledger

- Scope the ledger key by filer address so filers sharing one store do
  not overwrite each other's pending sets; ledgers written under the
  old unscoped key are claimed once on startup.
- Serialize snapshots on deletionSnapshotLock so an in-flight timer
  snapshot cannot overwrite a newer shutdown snapshot, and wake the
  snapshotter on every queue/forget so a queued id persists within
  milliseconds instead of a full interval.
- Merge recovered ids into the pending set immediately on reload; only
  the queue push waits out the grace window, so an early snapshot
  rewrites the recovered ids rather than dropping them.
- A failed or unparseable ledger read blocks persistence for the run
  instead of letting snapshots overwrite the unread ledger.
- Split the ledger into part keys when it exceeds one 64KB value so
  stores with a size cap (FoundationDB) do not strand the backlog.
- GetReadyItems reports retry-exhausted ids so they are forgotten in
  the ledger instead of replaying after every restart.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: close the remaining deletion-ledger durability gaps

- A manifest referencing a missing part is corruption: surface a wrapped
  error and block persistence instead of treating the ledger as absent.
- Multipart snapshots write generation-scoped part keys and publish the
  manifest last, so a crash never mixes old and new part contents.
- Orphaned parts are tracked in a persisted .stale sidecar and retried.
- Legacy/index ledgers are republished under the scoped key before the
  old keys are removed.
- A ledger index lets a filer restart under a new address claim the
  ledger its previous incarnation left behind.
- Expired and permanently-failed retry items only forget the ledger
  epoch they recorded, so they cannot erase a re-queued id.
- A failed startup read no longer disables persistence: every snapshot
  retries the reload until the store reads again.

* filer: tighten ledger claiming, index updates, and retry epochs

- touchLedgerIndex verifies its write and retries so a concurrent
  filer's merge cannot silently drop this key from the index.
- Foreign-ledger claims abort on any unreadable source instead of
  leaving it stranded once the new scoped key exists.
- A source that republished during the claim is left in place and its
  newer ids merge into the claimant's pending set.
- AddOrUpdate no longer overwrites the ledger epoch of an in-flight
  retry item, so its expiry or permanent outcome cannot forget a record
  that was re-queued after the attempt began.
- The recovery grace wait exits on shutdown instead of re-queueing
  after the filer has stopped.

* filer: requeue surviving records, persist claim deltas, guard index writes

- A dropped retry item (expired or permanent) whose ledger record was
  re-enqueued now pushes the id back through the hot queue instead of
  leaving it pending with nothing scheduled.
- Ids merged from a claim source that republished mid-claim are
  rewritten under our ledger immediately, so they are durable even if
  the claimant crashes before the next snapshot.
- touchLedgerIndex aborts when the index read fails for a real error;
  only ErrKvNotFound means the index is empty, so a transient failure
  can no longer wipe peer entries with a one-key write.

---------

Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com>
Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
authored and GitHub committed 2026-10-02 22:37:31 +08:00
1 parent 43abe21ffa
commit 1445960f8c
6 files changed
+1481 -25

No files matched your search

+39 -1
View File
@@ -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()
}
+55 -15
View File
@@ -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())
}
}
+758
View File
@@ -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
}
+584
View File
@@ -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
}
+33 -9
View File
@@ -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)
}
}
+12
View File
@@ -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
}