Skip to main content

Migrating a Partition With Zero Data Loss: The Cut-Over Protocol

· 13 min read
Clustron Team
Distributed Systems Engineering

The Zaris partition cut-over protocol

A failover is what happens when an owner dies — the cluster has seconds to pick a survivor and the partition map never moves. Migration is the calmer, more dangerous cousin: the owner is perfectly healthy, and you want to move a partition off it on purpose — to grow the cluster, to rebalance, to drain a node you're about to retire. Nothing has failed, which is exactly why the bar is higher. A failover is allowed a few honestly-lost tail writes in the worst case; a planned migration is not. If moving a partition from node A to node B can lose a single acknowledged write, migration isn't a feature, it's a liability.

This post walks the cut-over protocol Zaris actually runs, in order, and spends most of its time on the one step everyone gets wrong: the moment the source is finally "free" to drop its copy. That step is a loaded gun, and the interesting engineering is the safety on it.

The shape of the problem​

A partition is a set of segments, and a segment is a bag of key-value records, each carrying a monotonic revision. Migration moves a segment's data from a source node to a target node and then flips ownership so clients route to the target. The whole thing happens while the segment is live: reads and writes keep landing the entire time.

That liveness is the problem. If you snapshot the source, ship the bytes, and flip ownership, every write that landed during the shipping is on the source and not on the target — and the instant you flip, those writes vanish from the routed view. You acknowledged them; the client was told "committed"; now they're gone. The naive migration loses exactly the writes that happened while it was busy being careful about the ones that happened before.

So the protocol has to do two things that are in tension: keep serving writes on the source right up until the flip, and guarantee the target holds every one of those writes before the flip takes effect. The way Zaris reconciles that tension is a sequence of five steps, each of which refuses to advance until the previous one is provably complete.

Step 1 — a revision-anchored snapshot, with the trim floor pinned​

The target kicks off by asking the source for the segment's data (SegmentMigrationRequest). The source doesn't just serialize its current dictionary; it takes a revision-anchored snapshot. It records the segment's current revision — call it R — and captures the live records as of R.

Two subtleties live in this one step. First, the source arms its op-log before capturing the anchor. Zaris can skip the op-log for a sole-primary (RF1) segment as a write-path optimization; arming it first means every write from this point forward is logged, so a write that was previously skip-logged still has a revision at or below R and is carried by the snapshot itself. The delta we compute later — "ops after R" — is then provably hole-free.

Second, the source pins its op-log trim floor at R:

var (snapshotRevision, snapshotEntries) = _store.CreateSegmentSnapshot(segmentId);
// Pin BEFORE streaming so any op the source acks between the snapshot and the
// new owner's final-delta request is still replayable — the floor cannot advance past R.
_store.PinSegmentTrim(segmentId, snapshotRevision.Value);

Normally the op-log trims old entries to reclaim memory. During a migration that would be a disaster: if the floor advanced past R while the bulk transfer was in flight, the writes between R and the flip would be un-replayable, and the final delta would silently be partial. Pinning the floor at R guarantees those writes survive until the target has them. The source then streams the snapshot entries — and keeps serving reads and writes the whole time. It is still the routed owner. Nothing has moved yet.

Step 2 — the target waits, deliberately​

The target applies the bulk snapshot and then does something that looks like inaction but is the heart of the design: it moves the segment into an awaiting-cut-over state and stops. It does not acknowledge the source. It does not claim ownership. It does not tell the leader it's done. The comment in the code is blunt about why:

NO Ack to the source (it is still the routed owner), NO ownership claim, and NO completion notify yet — the completion is withheld until the fenced final delta confirms this node holds everything the source acked.

This is the "old-owner-authoritative" contract. Until the flip, exactly one node answers for the segment, and it's the source. A target that jumped the gun — started serving off a snapshot that's already stale by a few hundred writes — would be a split brain. So the target holds, and asks for the one thing that will let it prove it's caught up: the final delta.

Step 3 — the write-fenced final delta and the digest proof​

The target sends a SegmentMigrationFinalDeltaRequest anchored at R. This is where the source does the careful part.

// Fence FIRST, then capture a revision after the fence — so no acked write can exist beyond it.
var fenceEpoch = _topologyManager.EngageMigrationWriteFence(seg, request.CreationMapVersion);
var fenceRevision = _localStore.GetSegmentCurrentRevision(seg);

It engages a write fence on the segment and then reads the current revision. Ordering matters: fence first, measure second, so there is no window in which a write slips past the measurement. Then it reads every op-log entry after R — the delta of everything that landed during the bulk transfer — and computes a content digest of the segment's live state. That digest is the crux of the whole protocol, so the source records it as the cut-over proof:

response.Operations = finalOps.ToList();
response.SourceDigest = _localStore.ComputeSegmentLiveDigest(seg);
// Record this served digest as the cut-over PROOF: the post-cut-over orphan drop will refuse to
// destroy this segment's data unless its live digest still equals this exact value.
_localStore.RecordCutoverProofDigest(seg, response.SourceDigest);

There's an honest failure mode wired in right here. If the trim floor did somehow advance past R — a lost or regressed pin — then "ops after R" would be incomplete, and shipping that delta would silently drop acked writes. The source refuses:

if (!_localStore.TryGetSegmentOperationsAfter(seg, request.SnapshotRevision, out var finalOps))
{
response.Aborted = true;
response.AbortReason = "FinalDeltaBelowTrimFloor";
// release the fence, re-snapshot, retry — a loud abort, never a silent partial cut-over.
}

A below-floor abort is not a loss; it tears down and re-snapshots from a fresh anchor. The rule the code holds to is loud abort over silent partial — the migration would rather restart than flip on an incomplete delta.

Step 4 — the gate: the target proves it dominates​

The target applies the delta, then recomputes its own content digest of the segment and compares it to the source's:

var localDigest = _localStore.ComputeSegmentLiveDigest(seg);
var matchesSource = string.Equals(localDigest, response.SourceDigest, StringComparison.Ordinal);

If the two digests are equal, the target holds byte-for-byte what the source held at the fence — the target's state is a superset of every write the source acknowledged. If they differ, it's almost always because the streamed snapshot is still being applied when the first delta is computed (the completion signal can overtake queued data batches), so the target stays pending and lets a spaced retry re-request the delta rather than spinning. A mismatch never flips. The flip is gated on a positive proof of equality, not on the absence of an error.

There's one more thing the target does before it declares victory, and it's a guarantee most migration implementations skip. The source isn't the only copy that might be ahead — a replica of the segment could hold a write the fenced source missed. So before completing, the target folds in every live replica's read-only tail:

The flip may fire only when the new owner's digest dominates the source and every live replica of the segment.

The fold is revision-preserving: a lagging replica's older records are ignored by the revision gate, and an ahead replica contributes exactly the write the source didn't have. Only when the target dominates the source and every live replica does it move on. This is the "durability before ownership" invariant — ownership is a reward for provable possession, not a precondition for it.

Step 5 — the flip, and the drop that almost loses your data​

Now the target tells the leader it's done (SegmentMigrationCompletionNotify). The leader validates the completion against the live migrating override and arms an atomic cut-over regeneration of the partition map — a single versioned map change that flips the segment's owner to the target and removes the migrating override. After that map applies, the target is the routed owner. Clients route to it. The move is, as far as the cluster is concerned, done.

And the source, now holding a segment it no longer owns a role for, runs a housekeeping sweep to drop the orphaned copy and reclaim the memory. This is the loaded gun.

Here is the trap, and it's subtle enough that it survived into a certification run before it was caught (Missing=8629 acknowledged keys gone). A write fence reduces acking but, on a single-copy segment under load, does not perfectly quiesce it to zero instantaneously. So this sequence is possible:

  1. The source serves the final-delta digest as the proof. The target matches it and flips.
  2. In the sliver between serving the proof and the flip landing, one more write is acked on the source.
  3. The target never received that write — the proof it matched was from before it.
  4. The source, now an orphan, drops its copy to reclaim space.
  5. That last write existed only on the source. It's gone. A client was told it committed.

The fix is to make the drop conditional on the proof still being true. The source does not blindly destroy an orphaned segment. It compares the segment's current live digest against the digest it recorded as the cut-over proof, and only destroys the data if they still match:

public static bool ShouldRetainAfterCutover(bool hasCutoverProof, string? currentLiveDigest, string? proofDigest)
{
if (!hasCutoverProof || string.IsNullOrEmpty(proofDigest))
return false; // not a proof-carrying cut-over segment -> this policy does not intervene
return !string.Equals(currentLiveDigest, proofDigest, System.StringComparison.Ordinal);
}

If a write landed after the proof, the current digest differs from the proof digest, ShouldRetainAfterCutover returns true, and the source keeps its copy:

if (CutoverDropSafetyPolicy.ShouldRetainAfterCutover(hasProof, currentDigest, proofDigest))
{
Log.CutoverDropRetainedUnprovenTail(_logger, store.Count, segmentId);
continue; // keep the store — a post-proof tail exists that the new owner has not received
}

The asymmetry is the whole point. Retention is always safe: the source is no longer the routed owner, so a kept copy is a redundant, non-served orphan — it can never cause a split brain, it just costs a little memory until a re-serve re-delivers the tail to the new owner and a later sweep, finding the digests agree again, drops it cleanly. Destruction on a stale proof is never safe: it's the one operation that can permanently delete an acknowledged write. So the policy is deliberately narrow — it only ever overrides a drop that would otherwise happen, and only for a segment carrying a recorded proof whose digest no longer matches. Segments with no proof (an eviction, an immediate flip, a never-migrated segment) keep their existing behavior untouched. The gate adds a retain; it never adds a drop.

Why the digest, and not a revision number?​

You might ask why the proof is a content digest and not just "the highest revision the target acked." The answer is that a revision number tells you how far the target got, but a digest tells you whether what it got is what you had. The two can diverge in the seams — a replica fold, a re-anchored retry, a snapshot still draining — and the only question that keeps data safe is the content question: does the new owner hold exactly this? A digest answers that directly. A watermark only answers a proxy for it, and proxies are where acknowledged writes go to die.

What this buys, honestly​

Stated plainly, the guarantee is: a planned partition migration does not lose an acknowledged write. Every write the source acked before the fence is in the snapshot or the delta; every write it acked after the proof keeps the source's copy alive until the new owner has it too. The protocol will abort and restart loudly rather than flip on an incomplete delta, and it folds live replicas so the new owner dominates every copy, not just the source.

What it is not is free. The fenced final delta is a synchronous round-trip that briefly fences writes on the segment being cut over — a short, bounded pause on one segment, not the partition and not the cluster. A migration that can't converge (a trim floor that keeps advancing, a replica tail that can't be expressed from its log) re-snapshots and retries rather than forcing progress, which under pathological churn means a migration takes longer rather than taking a shortcut through your data. Those are the right trade-offs for a planned operation: a move you chose to do should never be the reason a committed write disappears.

This is the same backbone that the twelve-hour chaos certification exercises from the outside — stop a node mid-migration, churn the map, and watch acknowledged writes survive. This post is the inside view of one of the seams that run proves closed: the moment a partition changes hands, and the safety on the trigger that keeps the handoff from eating the very writes it was built to protect. It leans on two other pieces worth reading on their own — the replayable op-log that makes "ops after revision R" a real, hole-free query, and the elastic scale-out that is the most common reason a partition ever needs to move at all.