An Owner That Restarts Has No Right to Its Own Emptiness
A failover covers the moment an owner dies: the cluster picks a survivor and keeps serving. A planned migration covers moving a healthy owner's partition on purpose. This post is about the third act, the one that runs after the sirens stop — the original owner process restarts, rejoins the cluster, and wants its partition back. Its on-disk-nothing, in-memory-everything store came back empty. The data it is responsible for is still alive, but it is alive on the replicas that kept serving while it was down.
The whole correctness problem is contained in one sentence: the restarted owner cannot be allowed to trust its own emptiness. An empty store looks exactly like a store someone deliberately cleared. If the node assumes "I'm the owner and I have no data, therefore this partition is empty," it will happily serve NotFound for millions of live keys and, worse, let its empty state flow back onto the replicas and delete them. This post walks the five small, pure policies Zaris uses to make sure that never happens.
Why emptiness is undecidable from the inside
Start with the trap, because every rule below is shaped by it. A node that boots up holding zero records for a segment it owns is in one of two states, and it cannot tell which from local information alone:
- It is empty because it just restarted and never pulled its data back. It must pull.
- It is empty because it authoritatively cleared or expired those keys, and a replica that was briefly down still holds the stale copies. If it pulls, it resurrects deleted data.
These two are locally indistinguishable — the store is empty either way. An earlier version of Zaris tried to resolve this with an EmptyActive anomaly signal and it was removed precisely because the signal is a guess. You cannot ask "am I empty for a good reason?" and get a trustworthy answer from the empty node.
So Zaris does not ask emptiness anything. It routes the decision through a different fact: a sticky, partition-level flag called certified, which means exactly "my startup sync completed and my store is now authoritative." A node sets it only after it has proven possession. A freshly restarted owner is, by definition, uncertified — and an uncertified owner has no authority to be empty. It must pull from whoever holds the data. A certified owner is authoritative and is left strictly alone, so a legitimate clear is never resurrected.
That single reframing — from emptiness (undecidable) to certification (a proven, sticky fact) — is what makes the rest of recovery safe.
Step 1 — the restarted owner pulls, instead of assuming
When the node rejoins, the map still assigns it as owner of its old partition. The reverse-sync trigger decides whether each owned segment should be pulled from a live replica rather than treated as locally authoritative:
public static bool ShouldReverseSyncOwnedSegment(
bool partitionCertified,
bool isMapPreferredPrimary,
bool liveNonSelfReplicaExists)
=> !partitionCertified && liveNonSelfReplicaExists;
Three things about this tiny predicate are load-bearing, and each is a bug that reached a certification run before the rule was written down.
It ignores isMapPreferredPrimary. That parameter is still passed in — as documented decision-context — but the decision deliberately does not read it. An earlier version only fired reverse-sync when the node was also the map's preferred primary. A node that is the runtime active but not the map primary — the exact state a re-homed or promoted-after-recovery owner lands in — had its owned segments skipped entirely. It never requested its data and stranded empty forever while a live replica held every record (one cert run: active holding 0 vs replica holding 4600 on fifteen segments, never converged over eight minutes). Ownership of the data-pull cannot be gated on being the preferred primary; it has to fire for any uncertified active.
The signal is !partitionCertified, not emptiness. This is the undecidability fix from the previous section, in code. An uncertified active has never proved possession, so it pulls; a certified one is authoritative and is skipped.
liveNonSelfReplicaExists is a guard, not a nicety. If the configured replica is down, there is nothing to sync from. Driving recovery against a dead source spins forever (RECOVERY STUCK source=self synced=0/128). So when no live replica exists, the active is the only surviving copy — it is the authority, and it certifies immediately rather than waiting for a source that will never answer. The pull itself is revision-fenced and watermark-bounded, which makes it a harmless no-op when the replica happens to hold nothing.
Step 2 — don't serve a key you haven't caught up to yet
Pulling is not instantaneous, and the window while it runs is where acknowledged reads quietly turn into NotFound. A naive gate is a single boolean, IsSynced, meaning "a sync pass finished." That boolean is not enough, and three separate cert failures proved it: a reclaiming active promoted 34 keys short; a handoff lost one acked key; a rejoining active served NotFound for 57 keys during churn. In every case a sync pass had finished — but the source (the temporary active that kept serving) had since accepted more writes the recovering node had not pulled. IsSynced was true and the node was still behind.
The fix is to gate on a position, not a flag. The op-log revision is a replicated, cross-node-comparable number: the active assigns it and ships it, so a caught-up replica converges to the same head. The serve-gate compares heads:
public static SegmentServeDecision Decide(bool isSynced, long localHeadRevision, long sourceHeadRevision)
{
if (!isSynced)
return SegmentServeDecision.HoldUnsynced;
// A watermark of 0 carries no information — do not block on it; trust the completed sync pass.
if (sourceHeadRevision > 0 && localHeadRevision < sourceHeadRevision)
return SegmentServeDecision.HoldBehindSource;
return SegmentServeDecision.Serve;
}
HoldBehindSource is the new state that matters: a sync pass completed, but the source's head (captured at that sync) is ahead of ours, so the source accepted writes we haven't pulled. Serving now would return NotFound for exactly those keys. The node holds until its local head reaches the watermark. Note the deliberate degrade: a recorded watermark of 0 means "no watermark information," and the gate falls back to the IsSynced-only behavior rather than blocking forever on missing data. The rule never blocks on the absence of a fact — only on a fact that says you are behind.
Step 3 — pull, but do not clobber
Now the dangerous direction: applying the snapshot the replica ships back. The instinct is to clear the local store and copy the snapshot in wholesale. That instinct is correct exactly half the time, and catastrophically wrong the other half. The deciding question is whose revisions are comparable to whose:
public static ResyncApplyMode Decide(
bool digestMatches,
bool isPassiveSubordinateReplica,
bool priorReconcileShortfall = false)
=> digestMatches
? ResyncApplyMode.Skip
: (isPassiveSubordinateReplica && !priorReconcileShortfall)
? ResyncApplyMode.Reconcile
: ResyncApplyMode.Wipe;
Three outcomes:
Skip— the content digests already match. Nothing to transfer, nothing to apply. (See the cut-over protocol for why a content digest, not a revision count, is the thing worth comparing.)Reconcile(non-destructive) — the local node is a pure passive subordinate of the one live active. Its per-key revisions all came from that single active, so they are directly comparable. The snapshot is merged revision-fenced: aPutRawkeeps the locally-newer record and fills in missing keys; tombstones apply viaDeleteRaw. Crucially, this preserves writes the replica accumulated over the live channel during a slow bootstrap instead of throwing them away by clearing first.Wipe(destructive clear + copy) — the local node is or was itself an active for this partition. Two independently-active nodes have incomparable per-key revisions: the same revision number on each refers to different writes. Merging them would silently interleave two histories. So the only safe move is to drop local state and adopt the authoritative snapshot wholesale.
The default is the safe destructive copy; Reconcile is unlocked only by proving the node is a purely passive subordinate. This is the "wipe-guard": it prevents the dangerous merge across two former actives, and symmetrically it prevents a wholesale wipe from discarding good live-channel writes on a plain replica.
There's one more wrinkle encoded in priorReconcileShortfall, and it comes straight out of a twelve-hour-certification abort. A revision-fenced reconcile can fence out records it should have applied and complete while still behind the authoritative snapshot — a segment stuck at 12643/15435 across 232 sync sessions, reconciling the same records out of the window every time, looping forever. The flag says: if a prior reconcile completed but left us short, stop reconciling and escalate to a wipe so the frozen snapshot is adopted whole and we actually converge. It is only ever set after a measured shortfall — an ahead-or-equal replica never trips it, so live-channel writes are never discarded to chase a phantom.
Step 4 — the other side: an uncertified owner can still be a good source
Recovery is not only about the node that is behind. Consider the mirror case: node A restarts, and node B — a replica that has been down — now wants to sync from A. Should A serve the request while A is itself still uncertified?
The naive guard says no: "uncertified means my store is empty, so serving would overwrite B's valid data with my gap." But that premise is false the moment A actually holds the authoritative copy — which happens two ways: A was a preferred primary that verified catch-up before promotion but whose promotion path never called CertifyPartition (so the sticky flag is stuck false while the data is all there), or A took ownership through a total-partition-failure recovery and holds the full data as the active without ever being stamped IsSynced. In both, A's copy is the source of truth, and refusing stranded the replica at zero (one dump: 1,792 requests refused; another: all 128 segments refused for a whole partition).
So the source-side guard is made data-aware:
public static bool ShouldRefuse(
bool selfIsRuntimeOwnerMapPrimary,
bool partitionCertified,
bool segmentSynced,
bool selfHoldsSegmentData)
=> selfIsRuntimeOwnerMapPrimary && !partitionCertified && !segmentSynced && !selfHoldsSegmentData;
Refuse to serve only when the node is the owner, is uncertified, and the segment is both unsynced and holds no data — the one genuine "mid-recovery and actually empty" case where serving a gap could let the requester certify against it. Because the receiver applies every chunk revision-fenced, serving real data is always safe; the guard narrows refusal to the single state where there is nothing safe to serve.
Step 5 — declare victory, exactly once
Finally the recovering active has pulled every owned segment and caught each one up to its source head. It is authoritative again, and it needs to say so — flip the sticky certified flag — so the reverse-sync trigger from Step 1 stops re-yielding its segments and the cluster registers the partition as converged.
public static bool ShouldCertify(
bool selfIsRuntimeActive,
bool partitionCertified,
int ownedSegmentCount,
int syncedSegmentCount)
=> selfIsRuntimeActive
&& !partitionCertified
&& ownedSegmentCount > 0
&& syncedSegmentCount >= ownedSegmentCount;
Certification is the symmetric completion of the trigger: the trigger enqueues an uncertified active's segments to pull; this flips the flag once every owned segment is synced, and the partition settles. The pre-existing certify path only ran for the map's preferred primary, so a re-homed non-primary active pulled all its data correctly and then never certified — it re-enqueued reverse-sync forever, and the cluster never registered as converged (one run: 30 checkpoints timing out ~106 seconds after start, while Missing=0 and counts never diverged — the data was right, only the settle was missing).
And note the one shortcut this policy deliberately refuses: there is no "no live replica, so certify my emptiness as authority" path for a non-primary active. A non-primary active is not the map's authority; certifying it empty while a temporarily-down replica still holds the data would let that replica later sync from this empty node and lose everything. The authority-to-be-empty shortcut belongs only to the map preferred primary, in its own reconcile path. Everywhere else, emptiness earns nothing.
The shape of the whole thing
Step back and the five policies are one idea applied in five places: possession is proven, never assumed.
- A restarted owner is uncertified, so it pulls instead of trusting its empty store (Step 1).
- It will not serve a segment until its op-log head reaches the source's — a finished sync pass is not the same as being caught up (Step 2).
- It merges rather than clobbers when revisions are comparable, and clobbers rather than interleaves when they are not (Step 3).
- As a source, it serves whenever it genuinely holds the data, and refuses only when it is truly empty mid-recovery (Step 4).
- It certifies — flips the sticky authoritative flag — only on proof that every owned segment is synced, with no emptiness-as-authority shortcut for a non-primary (Step 5).
Every one of these rules is a pure, IO-free predicate — a handful of booleans and two revision numbers in, one decision out — which is why they are unit-testable in isolation and why each could be pinned to the exact certification failure that motivated it. They are the inside view of the invariant the twelve-hour chaos certification proves from the outside: kill an owner mid-flight, let its replicas carry the load, bring it back, and watch Missing=0 hold across millions of acknowledged operations. The reason the restart does not lose a write is that the returning owner is never once allowed to believe its own emptiness. It has to earn its data back — from the replicas that kept it safe.
If you want the pieces this leans on: the replayable op-log is what makes "catch my head up to the source's head" a real, hole-free operation, and reading your own writes is the client-facing promise that the serve-gate in Step 2 exists to keep.