fix(cluster): never drop a local delete, and drain the pending ones on stop - #958
Merged
Merged
Conversation
…n stop On a per-node-storage cluster with file replication, every node unlinks its local copy of a file when the manifest drops it. The hand-off from the FSM delete callback to the unlink was a 1024-slot channel with a non-blocking send: when full, the delete was dropped with a log that promised a reconcile on restart — and nothing reconciles it, because startup catch-up only pulls what the manifest lists. Full was easy to reach: each of the two workers takes one item and waits out a 500 ms grace, while retention proposes deletes in chunks of 1000 with no pause and a compaction backlog applies dozens of completion manifests in one watcher poll. A second loss sat at shutdown: Stop cancels the coordinator's context before it closed the channel, so a worker parked in its grace returned with the item it held and whatever was still buffered was closed away. The replica stayed on that node forever, and every read there read the file twice (#953's own release note documented the leak). Manifest deletes now go onto an unbounded pending list behind its own mutex — the callback appends and nudges a worker, no I/O, no c.mu, within the FSM callback contract — and the workers take the whole list by swap every grace period. The list cannot outgrow the manifest the node already holds in memory, so there is no cap. The workers exit only on an explicit stop, which Stop closes after it has unregistered the FSM callbacks, so a stop drains everything that is pending; the drain runs against a fresh context and is bounded at ten seconds, past which the untaken count is logged at Error. Before each unlink a worker asks the manifest whether it lists the path again and leaves such a file alone. A gauge, arc_cluster_local_delete_pending, reports the backlog. What remains: a crash loses the deletes pending at that instant, except those Raft re-applies on restart because they came after the last snapshot. The #953 release-notes paragraph that described the leak — and told a Pattern 1 operator to reclaim the replica with the reconciler, which in local mode reports every replica of another node's file as orphan storage (a re-check stops the deletion, but the audit is unusable; filed separately) — is replaced by a pointer to the new entry.
This was referenced Sep 28, 2026
xe-nvdk
added a commit
that referenced
this pull request
Sep 29, 2026
… hold the manifest sweep until then (#961) * fix(cluster): re-pull a node's own files after an empty-disk restore; hold the manifest sweep until then On per-node storage with file replication the puller assumed a node still holds every file it once wrote: enqueue returned SkippedSelf for any entry whose origin was this node, and the startup catch-up walk fast-pathed the same entries. A node restored with an empty data disk under a stable cluster.node_id pulled every other node's files back and never its own, and its reads of those partitions returned fewer rows with no error, for good. An act-mode reconciliation on that node then found each of those entries missing locally and proposed its deletion, which since #958 every other node carries out on its replica. The walks now let the disk decide on per-node storage (a puller flag the coordinator sets for a local backend): a self-origin entry is checked inline at the manifest's size with the same rule the worker's pre-pull check uses, skipped when present as before, and otherwise catch-up-tagged and pulled from a peer that holds a replica like any other entry. The reactive register of an own file — this node just wrote it — is still never pulled, and on a shared bucket nothing changes, since a missing own object there is not on any peer either. The reconciler checks the manifest-sweep gate before step 4 and, when it is withheld, records manifest_sweep_held on the run and goes on to the storage half; going through the sweep's chunk-boundary check instead aborted the whole run as a lost lease. In main, the local-storage gate withholds the manifest sweep until ReplicationReady() when replication and its catch-up walker are both enabled, and warns at startup when the walker is off, since readiness would never be reached then. Known: an own-origin file no peer holds any more keeps the node from converging and its manifest sweep held until an operator clears it; with the walker disabled neither the re-pull nor the hold applies. Fixes #959 * docs(release-notes): link the empty-disk restore entry to #961
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
Follow-up from the #953 review, whose release note documented this leak. On a per-node-storage cluster with file replication (Pattern 1), every node unlinks its local copy when the manifest drops a file, and that unlink could be lost in two ways:
api/retention.go) and a compaction backlog applies every pending completion manifest in one watcher poll (compaction/watcher.go).Stopcancels the coordinator's context before it closed the channel; a worker parked in its grace returned on that cancellation with the item it held, and whatever was still buffered was closed away.Either way the replica stayed on that node forever and every read there read the file twice.
The fix (about 60 lines of production code):
c.mu, inside the FSM-callback contract. No cap — the list cannot outgrow the manifest this node already holds in memory (the callback fires only for a path the FSM listed), and the workers take the whole list by swap every grace period. One Warn when the backlog first crosses 10 000.Stopcloses it after it has unregistered the FSM callbacks, so a stop drains everything that is pending; the drain uses a fresh context, skips the grace, and is bounded at 10 s (past which the untaken count is logged at Error). The WaitGroup is per Start so a bounded wait that is still parked cannot be reused by a later Start.arc_cluster_local_delete_pending(four metrics sites; documented in docs(arc-enterprise): list the cluster local-delete gauge docs.basekick.net#88). The callback comment and log no longer claim a restart reconcile; the startup log dropsdelete_queue_size.## Bug fixesentry; the fix(tiering): keep the cluster manifest in step with migration; gate replicating per-node clusters #953 paragraph that documented the leak is replaced by a pointer to it and loses its "reclaim it with the reconciler" advice — in local mode the reconciler keeps only this node's own-origin manifest entries while walking all of local storage, so on a replicating node it reports every replica of another node's file as orphan storage (a per-candidate re-check stops the deletion, but the dry-run audit is unusable). Confirmed with a probe test and filed as reconciliation: local-mode origin filter reports every replicated foreign-origin file as orphan storage; only the sweep re-check prevents deletion #957.Dropped from the first plan on the adversarial pass: a shutdown journal with startup replay, a 1M cap with a drop counter, and a
cluster.local_delete_queue_sizekey — all protected only a graceful stop, which the drain now covers completely. Residual, documented: a crash loses the deletes pending at that instant, except those Raft re-applies on restart because they came after the last snapshot. Plan + matrix:docs/progress/2026-09-28-local-delete-overflow-and-tier-walk-move.md(untracked).Configuration matrix
clusterCoordinator == nilincmd/arc/main.goreplication_enabled=false(Pattern 2 shared bucket, or Pattern 1 without replication)startFilePullerLocked, reached only fromStartwith replication on;Stop's drain guards ondeleteStop != nilbackend.Type() != "local") → list stays empty; Stop drains nothingc.storageset before Start (SetStorageBackend)puller != nilby construction; the manifest closure treats a nil FSM as "not listed" (unlink), the opposite of the puller's closure, on purposeLocalBackend.Deletereturns nil on ENOENTTestLocalDeleteNeverDrops, 3000)Stopreturns (TestStopDrainsPendingDeletes, ctx cancelled first)deleteStopis closed → no late arrivalsDeletehangsStopreturns at the bound, Error with the count (TestStopBoundsTheDrain)startDeleteWorkersno-op guard;StopnilsdeleteStop/deleteWgTest plan
go build,gofmt -lempty,go vet ./internal/cluster/ ./internal/metrics/TestLocalDeleteNeverDrops(3000 in one burst, all unlinked, gauge back to 0),TestStopDrainsPendingDeletes(context cancelled first, 50 pending, all unlinked, prompt),TestDeleteWorkerSkipsPathBackInManifest,TestStopBoundsTheDrain(blocking disk, returns at the bound).TestDeleteWorkerClassifiesAnUnusableKeyand the fourTestStop_*shutdown tests keep passing.TestLocalDeleteNeverDropsfails; with the oldctx.Done()exit reintroduced in the worker,TestStopDrainsPendingDeletesfails. Results: 1024 of 3000 deletes issued (workers issued 1024 deletes, want 3000), andstop drained 0 deletes, want 50three runs out of three (the test enqueues on live workers, lets them enter the grace, cancels the context, then stops — the real order)go test -race ./internal/cluster/...(whole package): 324 tests pass in about two minutes on the final tree. The first whole-package run hung past the ten-minute default:TestManifestDelete_ReopensGateAfterFailedCatchUp's rig cancelled the coordinator's context and then waited on the delete workers' WaitGroup, which was enough for the old workers and a deadlock for the new ones, whose only exit is the stop signal (the very property this PR adds). The rig now stops them the wayStopdoes; that test passes on main in 1.3 s and here in 1.5 s, and the whole package then passesTestStopBoundsTheDrainasserts 2 with one stuck and one behind); M1/M2 the gauge is published under the mutex from pending + in-flight, so it no longer drops to zero at take time and two publishers cannot leave a stale value; M3 the manifest lookup is captured by the workers at start instead of read from the field, so a worker that outlived a timed-out drain cannot race a later in-process Start; M4/M5 the note says Raft re-applies pending deletes after the node re-registers its callbacks (usually all, since a lone node must first win an election) and that the 10 s bound sits insideserver.shutdown_timeout; M6 the docs monitoring page gains the gauge (paired docs PR); M7 the plan's nil-FSM row was wrong (that branch is dead; it reads "unlink" on purpose) — fixed; style: "below" → "above", test renamed,releaseclosed viat.Cleanup, comments softened. Pre-existing gap the cluster-ops reviewer surfaced, filed separately:FSM.Restorenever fires the delete callback, so a node that catches up through a snapshot install keeps replicas deleted in the gap.deploy/docker-compose/enterprise-local: writer1-3 + reader1, replication on, plus a SeaweedFS cold bucket, tiering on every node, per-minute schedule,compaction.daily_min_files=2). The first Pattern 1 reader live run. Six flushes into one day on the primary → six replicas on reader1 within 2 s → daily compaction on the primary (one completion manifest, six source deletes in one Raft entry) → reader1 unlinked all six in one batch within 2 s (sixremoved local copylines,arc_cluster_local_delete_pendingback to 0) → the per-minute tick migrated the daily to cold and its manifest delete removed the daily replica from reader1 → reader1's next cycle synced the cold row (cold_synced=1) → on reader1,SELECT COUNT(*) FROM cpuwith no time range returned 30 (the fix(tiering): read an all-cold measurement without a time range; retire hot rows for vanished files #954 read on a reader, never run live before). Second burst with a reader scan while the replicas existed (files_registered=6, hotfile_count=6) → compaction stranded the six hot rows → reader1's scheduled cycle five minutes later loggedhot_retired=6, and the no-range read returned 60 with coldfile_count=2. A gracefuldocker stopof reader1 right after the second compaction and a restart were clean (no drain-bound Error, no errors after restart); the stop did not catch the drain mid-flight, because reader1 had unlinked the six sources nine seconds before its shutdown hooks ran, so the drain-on-stop path is covered byTestStopDrainsPendingDeletes, not by the rig. Reader's startup log showsdelete_workers=2and no queue size.