Skip to content

fix(cluster): preserve file updates during in-flight pulls - #907

Merged
xe-nvdk merged 12 commits into
Basekick-Labs:mainfrom
efegokdemir:fix-798-version-aware-pulls
Oct 5, 2026
Merged

xe-nvdk merged 12 commits into
Basekick-Labs:mainfrom
efegokdemir:fix-798-version-aware-pulls

Conversation

@efegokdemir

@efegokdemir efegokdemir commented Sep 19, 2026 •

Copy link
Copy Markdown
Contributor

Refs #798
Fixes #972

Problem

When the manifest advances a path while its pull is in flight, the puller could finish or retry stale work, miss the forced callback paired with ordinary registration, or discard a verified local generation before its replacement was ready.

Changes

  • Hand off the FSM registration/content-change callback pair even when both callbacks describe the same version.
  • Keep forced refresh retries keyed by path, and clear them on success, deletion, or stale-manifest pruning.
  • Ensure a stale catch-up page hands off to the current manifest version instead of falsely opening the readiness gate.
  • Retain verified bytes from a superseded pull until the successor is validated and installed.
  • Remove the legacy ManifestHas hook and update fixtures to use ManifestEntry.

This covers the in-flight replica update shape. Same-size rewrites on the origin node and delete-then-same-size re-registration during the delete grace window remain outside this PR.

Tests and validation

  • go test ./internal/cluster/filereplication ./internal/cluster/raft -count=1
  • go test -race ./internal/cluster/filereplication -count=1
  • go vet ./internal/cluster/filereplication ./internal/cluster/raft
  • gofmt and git diff --check

efegokdemir and others added 2 commits September 19, 2026 05:12
Resolves puller.go: keep the Basekick-Labs#961 self-origin fast-path condition
(RepullMissingSelfOrigin) and the Basekick-Labs#965 statLocal/presentAtSize pre-pull
check, adding this branch's !request.force to the latter.
@xe-nvdk

xe-nvdk commented Oct 1, 2026

Copy link
Copy Markdown
Member

Thanks for this one. The in-flight supersession is the right shape for #798, and the tests are precise: all four that compile on main fail there, each on the path it names (the fifth references the new field). Your quarantine-test change is exactly the fix for #972, which I filed this morning before seeing your branch; please add Fixes #972 to the description.

I pushed a merge commit to your branch resolving the two conflicts with main (#961's self-origin fast-path condition and #965's statLocal/presentAtSize pre-pull check, with your && !request.force added). CI is green on it (08ca7ede), and build, vet, gofmt, the package under -race, and the Issue798 tests 5× under -race all pass locally. Line numbers below refer to that merged head.

Verdict: the supersession core is correct. I traced the interleavings of a finishing worker against a newer enqueue, an older one, OnManifestDelete, a queue-full drop, and Stop(), and found no way to lose a newer version, leak a slot, leak or prematurely clear a catch-up tag, or land old bytes after new ones. What makes it hold: pending implies a slot, the handoff is the first action under the lock, same-worker continuation, the tag survives the handoff with the one inheritance rule that's needed, and the slot is registered only on a successful send. One design point has to change before this can merge, and a few smaller things.

High: observed is a second copy of the manifest, on every node, forever

observed keeps a FileEntry per path for the life of the process and is written on every accepted enqueue, so after one startup catch-up walk it holds the entire manifest. I measured it: 376 MB per 1M paths as written, 164 MB even slimmed to {lsn, refresh} (the key strings dominate). The FSM already holds the manifest at roughly 450 B/path, so this nearly doubles reader-side manifest memory, paid for three things that don't need per-path state. Suggested replacement, as one package:

  1. Let the FSM say "content changed". applyUpdateFile exists to make readers re-fetch (fsm.go:1358), and applyRegisterFileStruct has the previous entry in hand (old, existed := f.files[entry.Path]). Fire a forced enqueue (a second hook next to onRegister in coordinator.go, or a flag on the existing callback) when, and only when, SHA256 or SizeBytes differ from the previous entry. That covers the sequential same-size update without memory, and it avoids a trap in the current rule (see M2).
  2. Ask the manifest at pull time instead of remembering. Replace the ManifestHas hook with ManifestEntry(path) (FileEntry, bool) and, at the top of each attempt and on a checksum mismatch, abort the request when the manifest's LSN is newer than the request's (skipped_superseded, neither failed nor succeeded). This replaces observed's stale skip, and it also fixes M1 below and the walker-page staleness case (a 1000-entry page is a snapshot, and the walker can sleep minutes mid-page at high water, catchup.go:258-271).
  3. Remember only what failed. A small refreshPending map[string]uint64 (path → LSN) for forced requests that failed or were dropped on queue-full, consulted by an unforced enqueue of the same LSN to re-force it, cleared when that LSN succeeds and on OnManifestDelete. Bounded by failures and drops, not by the manifest. The drop case matters: your version survives it only because observed is written in the send branch; without that write, the next walk would enqueue v2 unforced, find the old copy at the same size, and skip it.

All five of your scenarios still hold under this shape except the skipped_dup == 1 assertion in TestPullerSequentialSameSizeUpdateIssue798:227-230, which becomes an unforced enqueue that the presence check skips.

Medium

M1. A superseded version burns its full retry budget, and deletes the previous copy, before yielding (puller.go:1289-1455). The production shape of #798 is the delete API: it rewrites the bytes first (delete.go:914) and proposes UpdateFile after (:979). A reader whose v1 request is queued or between attempts pulls from a peer whose manifest is already at v2; the ack carries the peer's manifest SHA (coordinator.go:2344) and the client rejects it immediately (fetch_client.go:249-253). So: three attempts plus 500 ms + 1 s of backoff before the handoff, totalFailed and totalChecksumMismatch bumped for something that isn't a failure, and deleteFile on each mismatch, so the reader serves nothing for that path during the spin. Item 2 above fixes it. Untested: the only in-flight test serves old bytes for LSN 1, i.e. a peer that still has the old version.

M2. "Newer LSN ⇒ force" re-downloads identical content (newerFileVersion, puller.go:114-120, used at :978-981 and :1006-1007). A same-path register or update with identical SHA and size but a new LSN (an idempotent re-register, a future metadata-only update such as a tier change) forces a full re-pull on every reader. Nothing proposes that today, but the rule invites it. Force only when SHA or size differ; use LSN only for ordering. Keep the equal-LSN SHA/size tie-break: applyBatchFileOps passes one logIndex to every op in a batch (fsm.go:1500, 1504), so equal LSNs with different content are possible in principle.

M3. The issue's second shape isn't fixed, so "Closes #798" overstates. After a manifest delete, the delete worker waits deleteGrace (500 ms, coordinator.go:57) and skips the unlink when the manifest lists the path again (deleteManifestHas, :950). A re-register inside that window at the same size: no slot, observed was cleared by OnManifestDelete (:711), force stays false, the worker finds the old copy present at the new size and skips it. #798 names this case. Either handle it (a puller-side marker set in OnManifestDelete and consumed by the next enqueue of that path, local backend only) or narrow the claim to the in-flight case and I'll track the remainder separately. Narrowing is fine by me.

M4. Docs and release hygiene. The Enqueue doc (puller.go:884-889) still says a same-path in-flight entry "is counted as a skip and never reaches a worker"; false for a newer version now. RELEASE_NOTES_2026.09.3.md exists on main (first entry landed today): please add an entry under Bug fixes in the house style, with what an operator saw (a reader kept serving the old bytes after an in-place rewrite until the next periodic reconciliation pass), what changed, and the limitation that a same-size stale copy across a restart that installs a Raft snapshot is not detected (FSM.Restore fires no callbacks). Credit line: Contributed by [@efegokdemir](https://github.com/efegokdemir) in [#907](https://github.com/Basekick-Labs/arc/pull/907).

Tests worth adding

  • The production shape: the in-flight version fails by ack mismatch, the pending version then succeeds, and no catch-up failure is recorded.
  • A pending replaced by a newer pending; Stop() with a pending; OnManifestDelete while a pending exists; an older-LSN enqueue while a slot exists (skipped_dup, :991).
  • TestPullerUpdatedVersionIsNotLostIssue798:119-123 waits on p.ctx.Done() with no timeout, so a fetch that never starts hangs the test instead of failing it.

Style

Pre-existing, out of scope, I'll file separately

Found while tracing, not introduced here: tryResumeFromPartial sizes via StatFile, so a shorter final file of the previous version reads as a resumable partial; after a writer failover the delete API's rewrite keeps the old OriginNodeID, so the origin node skips its own update and other readers try the stale origin first; and a same-size stale copy survives a restart through a snapshot install. They share one root cause, presence decided by size, which is worth a separate conversation.

Once the rework is in, I'll run it live on the per-node-storage rig (an in-place rewrite through the delete API while a reader is pulling) before merging.

@efegokdemir

Copy link
Copy Markdown
Contributor Author

Resolved the current merge conflict against upstream main in 1dca43c, preserving the #798 release-note entry alongside upstream's newer #978 notes. Existing Build & Test, CLA, and Enterprise chart checks were green before the branch refresh; GitHub is recalculating mergeability now.

…-pulls

# Conflicts:
#	RELEASE_NOTES_2026.09.3.md
#	internal/cluster/filereplication/puller.go
@efegokdemir

Copy link
Copy Markdown
Contributor Author

I resolved the new conflict against current main at bd8fbe9 and pushed 58827cc. The regression now covers the in-flight request receiving a checksum mismatch after the manifest advances, then pulling the pending version without recording a catch-up failure. The full filereplication package, go vet, and five focused -race runs pass locally; GitHub CI is running. The per-node-storage rewrite run you offered remains the final live validation.

@xe-nvdk

xe-nvdk commented Oct 5, 2026

Copy link
Copy Markdown
Member

Thanks for the rework — the structural asks from the last round are genuinely done, and I checked them in the code rather than taking the description's word for it:

  • observed is gone. refreshPending replaces it and is fed by both forced failures (puller.go:576) and queue drops (:1148). That second one was the hole in the redesign I proposed last round, and you closed it.
  • The FSM signals content changes on content, not on version: fsm.go:1209 and :1348 gate on SHA256 != || SizeBytes !=, it fires for the batch path through the *Struct variants, and an identical re-register correctly fires nothing.
  • The equal-LSN SHA/size tie-break is preserved in both newerFileVersion and manifestSupersedes.
  • test(filereplication): TestPullerQuarantineIsIdempotentAcrossReenqueues flakes when a re-enqueue lands before the inflight slot is released #972 is genuinely fixed — the assertion matches what the issue prescribes, and it holds at -race -count=20.
  • Release note is in the right file under ## Bug fixes, with credit.

Build, vet, gofmt, ./internal/cluster/filereplication and ./internal/cluster/raft are all green under -race on current main merged in.

Unfortunately there is a blocker, and it is the headline fix itself.

Blocker: force never survives the production call order, so the same-size rewrite is still skipped

fsm.go:1381-1388 fires two callbacks for one content-changing apply, in this fixed order on the apply goroutine:

if callback != nil { callback(&entryCopy) }                        // → Enqueue(entry)              force=false
if contentChanged && contentCallback != nil { contentCallback(…) }  // → EnqueueContentChanged(…)    force=true

The forced arrival lands second and finds inflight[path] already holding the identical version — equal LSN, equal SHA, equal size. So newerFileVersion is false at puller.go:1117, it is counted totalSkippedDup, and the forced request is discarded. The request that is queued has force=false, so at :1475 the else if !request.force branch takes skipped_local, sets succeeded = true, and the reader keeps serving the old bytes. Which is #798.

The reason this passed your tests is worth knowing, because it is the more useful lesson here: TestPullerSequentialSameSizeUpdateIssue798 and TestPullerFailedRefreshRetriesSameVersionIssue798 both call EnqueueContentChanged on its own, with no preceding Enqueue. That is a sequence coordinator.go cannot produce — every content change arrives as a pair. Both tests are green against the live bug. TestPullerUpdatedVersionIsNotLostIssue798 does pin the pending handoff, but there the first fetch failed so no local file exists and force is never consulted.

When a fix depends on a callback ordering, the test has to reproduce the caller's sequence, not the one the unit under test finds convenient. Modelling the FSM's two calls is enough:

backend.Write(ctx, path, oldBody)   // previous generation, same length as the new one
current = newEntry                  // manifest now at LSN 2 / new SHA
p.Enqueue(&newEntry)                // onFileRegistered
p.EnqueueContentChanged(&newEntry)  // onFileContentChanged
p.Start(ctx)

Today that yields skipped_dup:1 skipped_local:1 fetches=0 with the old bytes still on disk.

Fix, puller-side and order-independent (preferred over changing the FSM, so correctness does not depend on call order again): in the exists branch at :1111, treat a forced arrival of an equal version as a successor rather than a duplicate — roughly if request.force && (current == nil || !current.force) { p.pending[entry.Path] = request; … }. Please do not mutate active.force in place: the worker reads that field without holding inflightMu.

High

1. refreshPending is keyed by LSN, and that misses the refresh in the dangerous direction. The case I worried about last round — equal LSN with a different SHA from applyBatchFileOps' single logIndex — turns out to be the safe one: the equality matches, you over-force, one wasted pull. The hole is the opposite. Every re-register and every metadata-only UpdateFile re-stamps the LSN (fsm.go:1204, :1343) while contentChanged is correctly false. So: a forced v2 pull fails → refreshPending[p] = L2; a later identical-content re-register arrives at L3 → no content callback (right) → pendingLSN(L2) != L3 → not forced → presentAtSize → skipped_local. The local bytes stay at v1 permanently, every reconciliation pass repeats the identical skip because the walk always carries the manifest's current LSN, and the clear at :582 requires an exact LSN match so the entry is immortal. There is also no cap and no sweep — pruneStaleCatchUpState does not touch it, unlike staleKeptPaths, so the Config comment at :313-316 only holds if every forced failure is re-enqueued at exactly its own LSN.

Make it a set: map[string]struct{}, where an unresolved forced failure forces the next enqueue of any version of that path, cleared on any succeeded. That is sound because a forced request cannot reach skipped_local, so succeeded under force always means bytes were fetched and verified.

2. A superseded abort with no pending successor opens the catch-up query gate on a stale path. :1452-1458 sets neither failed nor succeeded, so finishEntry drops the catch-up tag and records nothing — correct when pending holds the successor, wrong when it does not. Reachable with the queue full during catch-up: both of v2's callbacks drop, the walker later enqueues its stale page entry v1 (not forced, per High 1), the worker aborts, the tag is dropped with no failure and no drop recorded, and FullyCaughtUp() flips true with the path still stale. Pre-PR that path burned three checksum mismatches and kept the gate red. On a superseded abort with no successor, account it like the drop branch.

3. The post-pull superseded branch unlinks without telling tiering. :1554-1557 calls deleteFile but not recordAbandonedFile, unlike its sibling at :1546-1547. That hook is what keeps a node's tier rows in step with its own disk (#1062), so this leaves a row for a file it just removed. Separately, the branch discards bytes that passed verification for the version the request named — if the successor's pull then fails, the node has lost a copy it had, which silently reverses the post-#999 "keep the stale copy rather than serve nothing" decision. It also currently masks the blocker above, since deleting the file is what defeats presentAtSize on the handoff. Once the blocker is fixed, I'd drop the deleteFile (the forced successor overwrites in place) or justify it and add recordAbandonedFile.

4. On the node that originally wrote the file, nothing here can fix a same-size rewrite. An in-place rewrite by node B keeps OriginNodeID = A (internal/api/delete.go:991). On node A both callbacks hit the self-origin fast path at :1099 → skipped_self, and the walks skip it because selfOriginPresent is presentAtSize, true for a same-size rewrite. This makes the release note's "until periodic reconciliation" wrong rather than narrow: for the same-size shape, reconciliation never heals it on any node (the walk enqueues, then processEntry takes skipped_local), and on the origin node it is not even enqueued. Either consult ManifestEntry's SHA before the self-skip for a content-changed arrival, or state the exclusion.

Medium

  • The reactive drop path now takes the FSM read lock it used to short-circuit. :1159-1162 calls manifestEntry unconditionally and then throws the result away for non-catch-up sources. That branch runs on the Raft apply goroutine, which must not block, and the hook takes a lock a manifest page fetch can hold for a full key sort — so the regression lands exactly in the queue-full row where drops are most frequent. Same shape in the processEntryOnce defer: !failed no longer short-circuits, so every successful pull pays an extra lookup. Please restore both.
  • ManifestHas kept in the production Config is a shape where the fix is off. All three supersede sites are gated on ManifestEntry != nil, so a puller wired the old way behaves exactly like main — and that is what puller_manifest_delete_test.go (7 sites) and tier_register_test.go:57 still use, so the existing manifest-delete suite proves nothing about the new paths. Migrate the fixtures and remove the field.
  • Fixes #798 will auto-close the issue while the delete→re-register-inside-deleteGrace-same-size shape it names is still open. Use Refs #798, and in the release note swap the snapshot-restore exclusion for that one (plus High 4).
  • manifestSupersedes' zero-LSN rationale cites a path that does not exist — every production entry is stamped, and the only LSN: 0 in the tree is your test table. Name the real case or drop it.
  • newerFileVersion provides no ordering for equal LSNs; it means "different", so arrival order decides. Safe today because the two callbacks run in order on one goroutine and :1452 re-reads the manifest — worth saying out loud, since the name claims an ordering it does not provide.
  • One worker can be pinned indefinitely by the handoff loop at :1400-1406 with no yield back to the queue; with Workers: 1 the queue stalls behind one hot path. A cap or a comment.
  • reconciliation_test.go:520-527 adds a 1 s wall-clock escape inside the fetch of a cancellation test, unrelated to cluster: a manifest UpdateFile (or re-register) for a path whose pull is in flight is deduped away; the reader keeps the old version #798/test(filereplication): TestPullerQuarantineIsIdempotentAcrossReenqueues flakes when a re-enqueue lands before the inflight slot is released #972. It does not mask a broken Stop, but if you saw a hang there it has a cause worth writing down. Otherwise drop the hunk.

Style

enqueue(entry, source, forced ...bool) — a variadic bool as an optional flag; make it explicit. enqueued still over-reports (:1121 counts a pending insert that never enters the channel). fsm.go:1208/1213 and 1347/1350 look up the same key twice. puller.go:300-301 has the same sentence as both a block and a trailing comment, and the #999 and "one stale peer is routine" rationales were deleted from metrics this PR does not touch — last round already asked for deleted comments back. Config.ManifestEntry's doc dropped the #759, #795 attribution and the "never invoked with inflightMu held" caller contract. Missing blank line before ### Audit events … in the release note.

Two corrections to my own last review

  • My M1 premise was stale post-Checksum mismatch stops the peer loop, so one stale origin can block a path permanently #999: the checksum-mismatch path no longer deleteFiles on each attempt — it keeps the local copy and marks it stale, so the reader served stale bytes during the spin, not nothing. The abort is still worth having (it saves two attempts plus backoffs), just for a less dramatic reason than I gave.
  • The equal-LSN/different-SHA case I flagged as needing the tie-break is the benign direction, as above. The tie-break is still right to keep; it just is not what protects you.

Tests

TestPullerUpdatedVersionIsNotLostIssue798 and TestPullerSupersededCatchUpClearsTagIssue798 earn their keep. Beyond the two discussed above, TestPullerManifestDeleteClearsRefreshPendingIssue798 is vacuous — its pull is unforced, so nothing is ever inserted into refreshPending, and both assertions still hold with the delete removed from forgetCatchUpPathLocked. And TestManifestSupersedesUsesContentAndVersionTogether has no equal-nonzero-LSN/different-SHA row, so the tie-break you were asked to preserve is itself untested.

No flakiness or goroutine leaks found.

The core handoff design is still right, and the memory problem from last round is properly solved — this is a blocker in the wiring, not in the idea. The one thing I would take from it beyond the fixes: when behaviour depends on how a caller sequences two callbacks, let the test drive that sequence.

@efegokdemir

Copy link
Copy Markdown
Contributor Author

Fixed in f36dc429d46e30fc173c3df58990f64488bacc0f (with current main merged at e83bb1e). Equal-version forced callbacks now become pending work; failed forced refreshes persist by path and are cleared by success, deletion, or stale-manifest pruning. Superseded catch-up work forces the current version, and verified old bytes remain until the successor is installed. Removed ManifestHas, corrected the #798 scope/closing reference, and added the callback-pair, catch-up, stale-copy, retry, delete-cleanup, and equal-LSN regressions.

go test ./internal/cluster/filereplication ./internal/cluster/raft -count=1, go test -race ./internal/cluster/filereplication -count=1, go vet ./internal/cluster/filereplication ./internal/cluster/raft, gofmt, and git diff --check pass locally. GitHub Build & Test is running; the per-node-storage live validation has not been run.

…callback pair; move the note to 27.01.1

Round-3 review fold-in on top of the contributor's change (Basekick-Labs#907):

- A forced content refresh on a shared backend would download the writer's
  own object from a peer and upload it back over the same key: a wasted
  transfer, and a regression of the object if a second rewrite raced the
  upload. Config.ForceContentRefresh, set by the coordinator for local
  storage only (as RepullMissingSelfOrigin is), turns a content-change
  signal into a plain registration elsewhere. Test: the FSM pair on a
  shared backend leaves a same-size copy alone and fetches nothing.
- The puller tests drove the FSM's callback pair by hand; nothing pinned
  the producer. TestFSMContentChangeCallbackFiresAfterRegisterAndOnlyOnContent
  asserts register-then-content order, no content callback for an
  identical re-register or an unchanged update, the pair through the batch
  path, and no panic with the callback unwired.
- The coordinator wires the content callback before the registration
  callback, so no apply can fire with only its non-forced half delivered.
- Release note moved from the shipped 26.09.3 file to 27.01.1 and its
  exclusions corrected: the origin-node same-size case is gone since Basekick-Labs#984
  re-stamps the rewriting node as origin; a rewrite learned from a snapshot
  restore is excluded instead (Basekick-Labs#1071); the delete-grace shape stays.
- Stale comments (snapshot restores "fire no callbacks", "the real
  ManifestHas hook", a duplicated field comment) and a leftover single-case
  select in reconciliation_test.go.
@xe-nvdk

xe-nvdk commented Oct 5, 2026

Copy link
Copy Markdown
Member

Round 3 reviewed against the same matrix (Pattern 1 reader, old origin after a #984 rewrite, the rewriting node, queue full, stale catch-up page, snapshot restore, Pattern 2 with replication on). All four round-2 findings and the blocker are genuinely fixed, and each is now proven by a test that fails against the round-2 code: the callback-pair test (result 3, the forced arrival discarded as a duplicate), the keep-copy test (deletes=1), the failed-refresh retry and the superseded catch-up re-enqueue. The drop path no longer touches the FSM lock for reactive arrivals, refreshPending is a path set, ManifestHas is gone, Refs #798 and the honest exclusions are right. Thank you for driving the caller's sequence in the test this time; that is what made the proof possible.

Three Mediums remained, folded in on your branch (one commit on top of yours):

  • Shared storage. With replication_enabled on a shared backend (off-recipe but reachable), a forced refresh would download the writer's own object from a peer and upload it back over the same key, and two rewrites racing the upload could regress the object. Config.ForceContentRefresh, set by the coordinator for local storage only (like RepullMissingSelfOrigin), turns the content signal into a plain registration there; a test pins it.
  • FSM side untested. The puller tests drove the callback pair by hand; nothing asserted that the FSM fires content after register, only on a SHA/size change, through the batch path too. One table-driven test in internal/cluster/raft does now. The coordinator also wires the content callback first, so no apply can fire with only its non-forced half delivered.
  • Note exclusions. "A same-size rewrite on the origin node" is stale since fix(cluster): refresh file origin after failover rewrite #984 re-stamped the rewriting node as origin; a rewrite learned from a snapshot restore is the exclusion that remains ([medium] cluster: after a restart, the snapshot-restore diff misses files pulled after the node's last local Raft snapshot #1071). Moved to RELEASE_NOTES_2027.01.1.md, since 26.09.3 shipped this morning.

Merging on green CI. #972 closes with it; #1061 carried the same test fix and will be closed as superseded.

@xe-nvdk xe-nvdk left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Round 3 clears the bar: the blocker and all four Highs are fixed and proven by tests that fail on the round-2 code; the three remaining Mediums were folded in on the branch (see the comment above). CI green on bfc7b79.

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.

test(filereplication): TestPullerQuarantineIsIdempotentAcrossReenqueues flakes when a re-enqueue lands before the inflight slot is released

2 participants