Anatomy of a Failover: What Happens to an In-Flight Write When the Owner Dies
A failover is the moment a distributed store is most honest about what it really is. In steady state every store looks the same: a write lands on an owner, gets acknowledged, you move on. The design only shows itself when the owner disappears mid-sentence — a node is holding the active copy of one of your partitions, writes are landing on it right now, and then it's gone. Hard-killed. Network-partitioned. OOM'd. The cluster has a few seconds to notice, agree on a new owner, and get traffic flowing again — and whether it does that correctly decides whether your acknowledged write survives or quietly evaporates.
This post follows that sequence end to end, in the order Zaris actually runs it, and names the cases where the answer is genuinely "the write is lost" — because a failover post that only describes the happy path isn't worth reading.
The one invariant that makes the rest tractable
Before the timeline, the single design decision everything else hangs off: Zaris separates the partition map from the partition runtime.
The partition map is the slow, deliberate assignment — which nodes are responsible for partition P: a preferred primary and its replicas. It changes only on real topology events (a node permanently joins or leaves, a scale-out, a rebalance) and it changes through a planned, versioned regeneration.
The runtime is the fast, live answer to a different question — which assigned node is the active owner serving P right now, and in what health. During a failover the map does not move at all. Only the runtime active moves. That's the invariant stated in the failover service itself: one independent failure signal opens the runtime failover gate, "PartitionMap remains stable and only PartitionRuntimeMap.Active is allowed to move during failover."
Why this matters: regenerating a partition map is expensive and dangerous — it can reshuffle healthy partitions and trigger data migration. A node dying should not trigger that. It should trigger a cheap, local runtime decision: pick a surviving replica of the same partition and make it the active. The map still says "P is owned by n1, n2, n3"; the runtime flips from "active = n1" to "active = n2". Keeping these two planes apart is what lets failover be fast without being reckless.
Step 1 — detect, but don't believe the first rumour
Failover is leader-coordinated. One node in the cluster is the elected leader, and only the leader promotes — a non-leader that hears a node-left event simply ignores it. That avoids two nodes independently electing two different new owners for the same partition.
When the leader receives a NodeLeft signal for a node, it does not immediately promote replicas for every partition that node owned. Promotion is called out in the code as "an irreversible, high-impact action," and it's guarded:
if (string.Equals(activeOwner, nodeId, StringComparison.Ordinal))
{
// A NodeLeft message can race ahead of the peer-table update.
// If the "dead" owner is in fact still reachable, do NOT fail it over.
if (IsNodeAlive(nodeId))
{
Log.NodeUnavailableRejectedReachable(/* ... */);
continue;
}
await PromoteOrMarkUnavailableAsync(partition, nodeId, runtime, ct);
}
The guard exists because the failure message and the membership layer's peer table update asynchronously. A NodeLeft can arrive a heartbeat before the peer table reflects it — or be flatly wrong. So before the high-impact action, the leader re-checks liveness directly. A still-reachable "dead" owner is left exactly where it is.
There's a second, slower path for the case the first one can't cover: an active owner that has gone silent without a clean NodeLeft. A periodic reconcile sweep (HandleUnavailableActiveOwnersAsync) walks every partition, and for an owner that isn't alive but hasn't been confirmed down, it starts a suspicion timer rather than acting:
var suspectedSince = _suspectedUnavailableSinceUtc.GetOrAdd(activeOwner, _ => DateTime.UtcNow);
if (DateTime.UtcNow - suspectedSince < _suspectEscalationTimeout)
{
Log.UnconfirmedActiveOwnerFailure(/* ... */);
continue; // not yet — give it until the escalation timeout
}
ConfirmNodeUnavailable(activeOwner, "RuntimeOwnershipSuspectEscalationTimeout");
Only after the owner has been unreachable for the whole escalation window (ten seconds by default) does the leader confirm it down and proceed. The design is deliberately asymmetric: marking a partition degraded for replica loss is cheap and reversible, so it fires immediately; promoting a new active is irreversible, so it waits for proof.
Step 2 — choose a survivor, and bump the epoch
Once the owner is confirmed down, the leader picks a replica to promote. Selection is availability-first: among the partition's assigned replicas, it prefers the one that has synced the most segments — the node most likely to be able to serve soon.
The promotion itself is a single broadcast that advances a per-partition epoch:
var nextEpoch = GenerateNextEpoch(partition.PartitionId, runtime, "ReplicaPromotion");
// ...
await _runtimeOwnership.BroadcastAsync(
partition.PartitionId, promotedNodeId, nextEpoch.Value, promotionState, ct);
The epoch is the fencing token, and it's the quiet hero of the whole mechanism. It's monotonic per partition, and the generator refuses to go backwards — a proposed epoch that isn't strictly greater than the highest already observed is rejected outright:
var nextEpoch = baseEpoch + 1;
if (nextEpoch <= highestObservedEpoch)
{
Log.StaleEpochGenerationRejected(/* ... */);
return null; // never regress ownership
}
So after the promotion, the partition's active is "n2 at epoch 8". If the old owner n1 was merely partitioned (not dead) and comes back thinking it's still the active at epoch 7, its stale epoch loses every comparison. It cannot reclaim ownership, and any write it tries to accept as "active" is fenced out. This is how Zaris avoids the two-primaries split-brain that would otherwise let a healed-but-confused old owner accept writes nobody will ever see again.
Step 3 — the possession proof: serve, or sync first?
Here is the subtle part, and the part that took real debugging to get right. Having chosen the most-available replica, can it actually start serving reads immediately?
Not necessarily — because "most segments synced" is not the same as "holds every acknowledged key." A replica can be caught up on segment count while still being behind the true data on the segments it has. Promote such a node straight to a serving state and it becomes an incomplete active: it answers NotFound for keys that were durably acknowledged on the dead owner and do exist on a more-complete sibling. That's silent data loss dressed up as a successful read.
So before the promoted node is allowed to serve, Zaris applies a possession proof. The leader keeps a fresh, term-matched table of how many keys each node reports holding for the partition. The rule is a pure comparison:
public static PartitionHealthState ResolvePromotionState(
long candidateHeldKeys, long bestOtherLiveHolderKeys)
=> bestOtherLiveHolderKeys > candidateHeldKeys
? PartitionHealthState.Recovering // a more-complete live peer exists: sync first, do NOT serve
: PartitionHealthState.Degraded; // candidate is the authority: serve now
Two outcomes:
- No strictly-more-complete live holder exists — the candidate is the most complete node, or a genuine sole survivor with nothing to sync from. It's promoted to Degraded: a serving state. It has nothing better to wait for, so wedging it non-serving would just make the partition needlessly unavailable. It serves immediately.
- A strictly-more-complete live holder does exist — the candidate is promoted to Recovering: a non-serving state. It syncs from that better source first and only begins serving once it's proven complete. Clients are told the partition is still coming up rather than being handed a node that would lie to them.
The distinction between Degraded and Recovering is exactly the distinction between "redundancy is reduced but the truth is being served" and "don't trust this owner for reads yet." Getting a sole survivor correctly into Degraded matters just as much as holding an incomplete node in Recovering — a single remaining copy must never be wedged non-serving forever just because it's the only copy left.
A quick note on honesty: this proof is count-based, drawn from per-node held-key reports with a freshness watermark (a report from a node that's gone silent is not counted as a live source). It is a strong, cheap guard against promoting an obviously-behind node to serve; it is not a cryptographic attestation that two nodes hold byte-identical data. The convergence that makes the copies actually identical is the replica-sync path, covered below.
Step 4 — the client reroutes, and your retry lands
While all of this happens on the server, what does the client that was mid-write see?
It sees the truth, as a typed status. A request that reaches the old owner (now fenced, or simply gone) doesn't hang or throw a vague error — it comes back as Moved or OwnershipChanging. The Zaris client treats those as retryable reroute signals: it refreshes its ownership map, discovers the new active, and replays the operation there. In the common case your PutAsync just takes a few milliseconds longer and returns Ok — against the newly promoted owner, at the new epoch. The application code never learns the owner moved.
That client machinery — the status taxonomy, the attempt-budget-and-deadline retry, the map refresh — is a story of its own, told in Retries You Don't Write. For this post the point is just that the server's epoch bump and the client's reroute are two halves of one handoff: the server makes the old owner un-authoritative, and the client stops talking to it.
Step 5 — converge, then heal back to Healthy
Promotion restores availability; it does not restore redundancy. Right after a failover the partition is running Degraded with one fewer copy than its target. Two things then happen in the background.
First, the promoted node (if it came up Recovering) and any lagging or returning replicas catch up by replaying the per-segment operation log from a known revision — the same replayable substrate that drives ordinary replication. A node that just inherited the active role and a replica that merely fell behind converge through the identical mechanism; that's the subject of One Log, Two Jobs.
Second, once every assigned participant is alive again and in sync, the leader's recovery sweep transitions the partition back to Healthy — but only through the same epoch-advancing broadcast, and only after a short cooldown so a flapping node can't drive the runtime through a storm of transitions. The health state machine is therefore a small, closed loop: Healthy → Degraded/Recovering → Healthy, every edge an epoch bump, every promotion fenced.
If no eligible replica survives at all — every assigned node for the partition is down — the leader does the only honest thing: it marks the partition Unavailable (no owner, no false promotion) and starts an eviction countdown that, if nobody returns in time, lets a planned map regeneration re-home the partition. An Unavailable partition fails reads and writes loudly rather than serving a lie.
Where a write is genuinely lost — stated plainly
A failover mechanism earns trust by being precise about its edges, not by implying there are none.
- Replication factor 1. A partition with a single copy and no replica has nothing to promote. When that node dies the partition goes Unavailable until it returns; the data is not lost (it's in that node's memory) but it is unavailable, and if the node never returns, it's gone. RF1 is non-durable by design — it is a deliberate choice for data you can afford to lose, not a guarantee the failover path can rescue.
- Asynchronous acknowledgement. Zaris lets a write be acknowledged on the active owner before it has reached a replica, for latency. That's the fast path — and its honest cost is a window: if the active dies in the gap between acking the client and the op reaching any surviving holder, that acknowledged write is lost. The possession proof protects you from promoting a known-behind node to serve; it cannot conjure an op that never left the dead node. If you need a write to survive the loss of its owner, it must be acknowledged only after it's durable on a replica — the synchronous path trades a little latency for exactly that guarantee. Choose per workload, with eyes open.
- The confirmation window itself. For the escalation timeout after an owner goes silent, the partition's writes are in limbo — retried by the client, not yet landing on a new owner. Correctness is preserved (nothing is wrongly promoted), but availability dips for that window. Tightening the timeout trades faster failover against more false promotions under transient blips; it's a dial, not a free lunch.
None of these are bugs the failover path can engineer away — they're the genuine trade-offs of availability versus durability versus latency, and Zaris's job is to make the safe choice the default and the fast choice explicit, never to pretend the window doesn't exist.
The shape of it
Strip the failover down and it's four moves on top of one invariant. The invariant: the map holds still; only the runtime active moves. The moves: confirm the death before acting on it, promote a survivor behind a monotonic epoch that fences the old owner, prove possession so the new owner serves the truth or syncs before it serves at all, and reroute the client so its retry lands on the new owner. Then converge the copies back and heal to Healthy.
That it holds up is not a matter of argument — it's measured. The same promotion, possession, and convergence path runs continuously under deliberate node-kill chaos for hours at a stretch, with the ledger checked for a single lost acknowledged write at the end. That certification is its own story: No Acknowledged Write Left Behind.