Skip to content

fix(cluster): never drop a local delete, and drain the pending ones on stop - #958

Merged
xe-nvdk merged 2 commits into
mainfrom
fix/local-delete-drain
Sep 28, 2026
Merged

xe-nvdk merged 2 commits into
mainfrom
fix/local-delete-drain

Conversation

@xe-nvdk

@xe-nvdk xe-nvdk commented Sep 28, 2026 •

Copy link
Copy Markdown
Member

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:

  • Drop on a full channel. The FSM delete callback handed each manifest delete to a 1024-slot channel with a non-blocking send; when full, the delete was dropped with "will reconcile on restart" — and nothing reconciles it: startup catch-up and periodic reconciliation only pull what the manifest lists. Full was easy to reach: each of the two workers takes one item and then waits out a 500 ms grace, while retention proposes deletes in chunks of 1000 with no pause (api/retention.go) and a compaction backlog applies every pending completion manifest in one watcher poll (compaction/watcher.go).
  • Loss at shutdown. Stop cancels 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):

  • An unbounded pending list behind its own mutex replaces the channel. The callback appends and nudges a worker: O(1), no I/O, no 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.
  • Workers exit only on an explicit stop signal, never on the coordinator's context. Stop closes 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.
  • Before each unlink a worker asks the manifest whether it lists the path again (the same lookup the puller already uses) and leaves such a file alone. Arc's own file names never repeat, so this is a belt for imports and restores, not a known race.
  • Gauge 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 drops delete_queue_size.
  • Release notes: new ## Bug fixes entry; 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_size key — 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

Configuration Reaches new code? Preconditions established?
OSS standalone (no cluster) no — clusterCoordinator == nil in cmd/arc/main.go n/a
Cluster, replication_enabled=false (Pattern 2 shared bucket, or Pattern 1 without replication) no — the workers are created in startFilePullerLocked, reached only from Start with replication on; Stop's drain guards on deleteStop != nil nothing constructed, nothing to drain
Cluster + replication, shared backend (S3/Azure) callback returns before the enqueue (backend.Type() != "local") → list stays empty; Stop drains nothing c.storage set before Start (SetStorageBackend)
Cluster + replication, local backend (Pattern 1) — the target yes: enqueue → list; workers drain; Stop drains completely puller != nil by construction; the manifest closure treats a nil FSM as "not listed" (unlink), the opposite of the puller's closure, on purpose
Compactor node with a local backend yes — its own source deletes fire the callback on itself; the subprocess already removed the files; LocalBackend.Delete returns nil on ENOENT harmless no-op
Retention run on Pattern 1 (the realistic producer) yes — 1000-per-chunk bursts go to the list, no drop (TestLocalDeleteNeverDrops, 3000) same as the target row
Process restart with un-snapshotted Raft log yes — post-snapshot deletes re-fire on re-apply → ENOENT no-ops harmless; the gauge is live, not cumulative
Graceful stop with a burst in flight drain completes before Stop returns (TestStopDrainsPendingDeletes, ctx cancelled first) callbacks unregistered before deleteStop is closed → no late arrivals
Stop while Delete hangs Stop returns at the bound, Error with the count (TestStopBoundsTheDrain) the only loss path left on a graceful stop
Two Stop/Start cycles in-process (tests) second Start creates fresh list/channels/WaitGroup startDeleteWorkers no-op guard; Stop nils deleteStop/deleteWg

Test plan

  • go build, gofmt -l empty, go vet ./internal/cluster/ ./internal/metrics/
  • New: 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). TestDeleteWorkerClassifiesAnUnusableKey and the four TestStop_* shutdown tests keep passing.
  • Pre-fix proofs (revert-run-restore): with the old 1024 drop reintroduced in the enqueue, TestLocalDeleteNeverDrops fails; with the old ctx.Done() exit reintroduced in the worker, TestStopDrainsPendingDeletes fails. Results: 1024 of 3000 deletes issued (workers issued 1024 deletes, want 3000), and stop drained 0 deletes, want 50 three 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 way Stop does; that test passes on main in 1.3 s and here in 1.5 s, and the whole package then passes
  • Review: one deep reviewer with the matrix (every row confirmed by trace; no data-loss path found) plus a cluster-operations reviewer on the FSM-callback contract, lock ordering with the FSM, Stop and Start ordering, and the checklist (all pass). Addressed: H1 the bound-expiry Error reported only the untaken count, which is zero exactly when a stuck worker holds the batch — an in-flight counter now makes the Error and the gauge report what is still on disk (TestStopBoundsTheDrain asserts 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 inside server.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, release closed via t.Cleanup, comments softened. Pre-existing gap the cluster-ops reviewer surfaced, filed separately: FSM.Restore never fires the delete callback, so a node that catches up through a snapshot install keeps replicas deleted in the gap.
  • Live, Pattern 1 rig built from this branch (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 (six removed local copy lines, arc_cluster_local_delete_pending back 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 cpu with 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, hot file_count=6) → compaction stranded the six hot rows → reader1's scheduled cycle five minutes later logged hot_retired=6, and the no-range read returned 60 with cold file_count=2. A graceful docker stop of 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 by TestStopDrainsPendingDeletes, not by the rig. Reader's startup log shows delete_workers=2 and no queue size.

…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.
@xe-nvdk
xe-nvdk merged commit 1625480 into main Sep 28, 2026
7 checks passed
@xe-nvdk
xe-nvdk deleted the fix/local-delete-drain branch September 28, 2026 23:15
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
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant