Skip to main content

Retries You Don't Write: Surviving Ownership Churn in the Zaris Client

· 10 min read
Clustron Team
Distributed Systems Engineering

Transparent retry and reroute under ownership churn

Here is a fact about any partitioned key-value store that the glossy diagrams tend to skip: the node that owns your key will move. A node fails and its partitions fail over to their replicas. You scale out and partitions migrate to the new member. A preferred primary comes back and ownership fails back to it. Every one of those events means that the node your client was happily talking to a millisecond ago is, right now, the wrong node for some set of keys.

The naive client experience of that moment is ugly: a connection reset, a generic "server error," or — worse — a write that lands on a node that is no longer the active owner and quietly goes nowhere. The Zaris .NET client is built so that none of that reaches your code. An ownership change becomes a typed status, the client reroutes to the new owner, and the retry is bounded by both an attempt budget and a wall-clock deadline. In the common case your PutAsync just takes a few milliseconds longer and returns Ok. This post is about the machinery that makes that true, and — just as importantly — about the cases it deliberately does not paper over.

The server tells the truth with a status, not an exception​

When a request reaches a node that cannot serve it, the node does not hang up or throw a vague error. It answers with a specific KvStatus that says exactly why it is the wrong place:

  • Moved — this node does not own the segment at all. Your routing table is stale; the key lives elsewhere.
  • OwnershipChanging — this node isn't the active owner right now, but ownership for that segment is actively in transition (a preferred-primary failback, say). It's not a settled routing bug; it's a handoff happening under you.
  • ReplicaWriteRejected — this node owns the segment but is the replica, not the active primary. Your cached active-owner is stale after a failover, and a write here must not be accepted.
  • Migrating — the segment is mid-migration to a new owner; writes are briefly held.
  • Unavailable — the partition has no reachable owner at this instant.

These are not invented for the client's benefit — they are produced on the data path. The owner-vs-non-owner decision, for example, is a single branch on the server: a non-owner request during an active handoff gets OwnershipChanging, while an ordinary stale-routing hit gets Moved.

// server-side, Clustron.Zaris.Core.ExternalStoreService
private KvStatus OwnershipFailureStatus(ushort segmentId)
=> _topology is not null && _topology.IsSegmentOwnershipTransitioning(segmentId)
? KvStatus.OwnershipChanging // handoff in progress: refresh and retry
: KvStatus.Moved; // settled, stale routing entry

Keeping those two distinct matters: a transient handoff should not look like a routing bug in your logs, and the client can react to each appropriately. The point of the whole status vocabulary is that "wrong node" is a first-class, machine-readable answer, not an error you have to pattern-match out of an exception message.

The retry brain: what is transient, and what is a real answer​

All of that routing intelligence lives in one place on the client — ClientRetryHelper — which the docs in the code call, fairly, "the retry/backoff brain of the client." Its most important decision is not how to retry but whether to at all. Two predicates draw that line:

// Retry only these — they mean "transient; the same request can succeed shortly."
private static bool ShouldRetry(KvStatus status) => status switch
{
KvStatus.Moved => true,
KvStatus.ReplicaWriteRejected => true,
KvStatus.Unavailable => true,
KvStatus.Migrating => true,
KvStatus.OwnershipChanging => true,
_ => false
};

// Of those, these three also mean "you reached the wrong node — refresh routing first."
private static bool RequiresMapRefresh(KvStatus status)
=> status == KvStatus.Moved
|| status == KvStatus.ReplicaWriteRejected
|| status == KvStatus.OwnershipChanging;

Look at what is absent from ShouldRetry. NotFound is not retryable — the key genuinely isn't there, and asking again won't conjure it. Conflict is not retryable either, and that one is the crucial distinction. As we covered in the CAS post, a Conflict from an IfMatchVersion write is a real answer: someone else wrote first, and your value was stale. The client must hand that straight back to you so your compare-and-swap loop can re-read and recompute. If the retry brain "helpfully" retried Conflict, it would silently defeat optimistic concurrency.

So the rule the client enforces is sharp: retry transient routing and availability failures; never retry a meaningful outcome. A failover is the store's problem to hide. A version conflict or a missing key is your application's business to decide.

Retry is not enough — you have to reroute​

Retrying a Moved against the same stale node would just produce another Moved forever. That's why the three "wrong node" statuses also trigger a cluster-map refresh through an onMoved callback before the next attempt. The retry loop, stripped to its spine, reads:

result = await operation().ConfigureAwait(false);

if (result.IsSuccess || !ShouldRetry(result.Status))
return result; // success, or a real answer — hand it back

if (RequiresMapRefresh(result.Status) && onMoved != null)
await onMoved(cancellationToken); // re-fetch ownership, THEN retry elsewhere

The same exception path exists for the transport-level failures that ownership churn also throws off — PartitionUnavailableException, NodeUnavailableException, TransportException, SocketException, TimeoutException. Each of those also fires the map refresh before retrying, because a socket error to a node that just died is itself a signal that routing has changed. A genuinely non-retryable exception, by contrast, is rethrown immediately — the loop never swallows what it doesn't understand.

The effect is that the second attempt doesn't go back to the node that said "not me" — it goes to whoever the refreshed map now names as the active owner. Retry plus reroute is what actually makes the operation succeed; retry alone would just spin.

Two bounds, not one — the refresh-storm we had to fix​

A retry loop needs a stopping condition, and the subtle bug is choosing only one. The Zaris client bounds every operation by both an attempt budget and a wall-clock deadline, and the loop condition fails whichever trips first:

for (int attempt = 1;
attempt <= maxAttempts && (!deadline.HasValue || DateTime.UtcNow < deadline.Value);
attempt++)

That && is there because of a real incident. An earlier version used the deadline instead of maxAttempts whenever a timeout was set. So a caller asking for maxAttempts: 3 with the default 60-second operation timeout didn't retry three times — it retried for the entire 60 seconds, firing a cluster-map refresh on every attempt. Under node churn, where lots of operations are failing and refreshing at once, that turned a single failing op into dozens of refreshes, all serialized behind the refresh lock. The send path starved and throughput collapsed toward zero. The fix was to honor both bounds: a few attempts and a hard deadline, whichever comes first.

The defaults are deliberately modest — maxAttempts: 3, backoffMillis: 100 — with the operation deadline resolved from ZARIS_REQUEST_TIMEOUT (60 seconds by default; set it to infinite to disable). The attempt budget keeps a transient blip cheap; the deadline keeps a genuine outage from hanging your call forever.

Backoff that respects the deadline — and knows migrations are slow​

Between attempts the client backs off linearly (backoffMillis * attempt), but two refinements make that production-grade:

private static TimeSpan GetDelay(int attempt, int backoffMillis, KvStatus? status, DateTime? deadline)
{
var delayMs = status == KvStatus.Migrating
? backoffMillis * attempt * 8 // a migration takes real time; don't hammer it
: backoffMillis * attempt;

var delay = TimeSpan.FromMilliseconds(Math.Min(delayMs, 2_000)); // hard cap
if (deadline.HasValue)
{
var remaining = deadline.Value - DateTime.UtcNow;
if (remaining <= TimeSpan.Zero) return TimeSpan.Zero;
if (delay > remaining) delay = remaining; // never overshoot
}
return delay;
}

A Migrating status gets an 8× longer delay, because a segment migration is a seconds-scale event and retrying every 100 ms just adds noise to a node that is already busy moving data. Every delay is capped at two seconds so backoff can't balloon, and — the detail that keeps the two bounds honest together — the delay is clamped to the time remaining before the deadline, so the client never sleeps past its own cutoff.

One refresh, not a thundering herd​

Here is the part that only matters at scale. Under churn, many operations fail at the same instant and every one of them wants to refresh the cluster map. If each fired its own refresh, you'd get a thundering herd pounding the surviving nodes exactly when they're most fragile. The client coalesces them into a single shared refresh:

internal async Task TryRefreshClusterMapSafe(CancellationToken ct = default)
{
Task? existing = _activeRefreshTask;
if (existing != null && !existing.IsCompleted)
{
// A refresh is already in flight — ride it instead of starting another.
await WaitForSharedTaskAsync(existing, ct);
return;
}
// ... under a lock, double-check and start exactly one, storing it in _activeRefreshTask
}

A handful of further guards make this robust: the refresh itself debounces (it no-ops if the map was refreshed in the last 500 ms), a caller waits at most a bounded budget for the shared refresh (15 seconds by default, ZARIS_CLIENT_MAP_REFRESH_WAIT_MS) before giving up on waiting without cancelling the refresh others still need, and the shared slot is cleared with an Interlocked.CompareExchange rather than a lock — because an earlier lock-based clear running inside a synchronous continuation could convoy the entire refresh path and wedge healthy reads. The lesson that keeps recurring: the resilience code has to be at least as careful about its own contention as the system it's protecting.

Where it stops — and why that's correct​

Retries are not magic, and the honest edge is idempotency. The retry loop gives you at-least-once send semantics: if a write's acknowledgement is lost on the wire after the store committed it, a retry will send it again. For a plain last-writer-wins PutAsync of a fixed value, replaying is harmless. For anything that isn't naturally idempotent, the fix isn't to disable retries — it's to make the operation safe to repeat, using the very primitives Zaris already gives you:

  • Guard a mutation with IfMatchVersion so a replayed write against a now-moved version comes back as Conflict — which, as established above, the client hands straight back to you rather than retrying.
  • Use IfAbsent to make a create-once genuinely once; a replay sees Conflict and you treat it as "already done."

In other words, the transport-level retry and the application-level compare-and-swap are designed to compose: the client hides the churn, and the preconditions make the hidden retries safe. If you need to turn retries off entirely — in a test that wants to observe raw statuses, say — ZARIS_DISABLE_RETRIES=true collapses the budget to a single attempt.

The deeper point is a division of labour. A distributed store that moves ownership will hand your client transient failures; that's not a defect, it's the cost of being able to fail over and rebalance at all. What a good client owes you is to make those failures the store's problem — typed, rerouted, bounded — while being scrupulously careful never to hide the failures that are actually your answer. Zaris draws that line at exactly one place: retry the churn, surface the truth.