Part 0 · 2 chapters · ~15 min

Why Distributed Is Different

Partial failure: four different failures that look identical to the caller, and why a timeout means "I do not know"; then the eight fallacies of distributed computing, the failure each causes, and the defences, including retries done properly and retry storms.

1

Partial failure

On one machine a call either returns or the program dies with it. Across a network, a call can fail in part: the request lost, the server dead before or after acting, the reply lost, or everything simply slow. The caller sees the same thing in every case, which is nothing.

four failures that look identical
  1. Request lost: no effect.
  2. Crash before commit: no effect.
  3. Reply lost: the effect happened.
  4. Slow past the timeout: the effect happens later.
  5. A timeout means "I do not know". Blind retries can double the effect, and not retrying can lose it.
  6. The way out: make the operation safe to repeat, then ask what happened. Never guess.
code
// the wrong way: a timeout treated as a failure
try { await ledger.debit(account, 50_000_00); }
catch (e) { if (isTimeout(e)) return markFailed(transfer); }     // ✗ the debit may have happened

// the right way: a key that makes repeats safe, then resolve the unknown
const key = transfer.id;                                          // stable across retries
try { await ledger.debit(account, 50_000_00, { idempotencyKey: key }); }
catch (e) {
  if (!isTimeout(e)) throw e;
  const status = await ledger.status(key);                        // ask, do not guess
  if (status === 'unknown') await retryLater(transfer);           // same key: at most one debit
}
the user-facing version
The Trust course part 7 shows the screen for this state: "We're confirming your transfer", not "Failed". Telling a user their money did not move when it did is how support tickets and double payments begin.
PARTIAL FAILURE: DID IT WORK?
one request, a timeout, and the six things that might have happened, none of which the caller can tell apart
swipe the figure sideways, or tap expand for full screen
1/6
happy path
The happy path: the payments service sends "debit ₦50,000" to the ledger, the ledger debits and replies "ok", the reply arrives in 40 ms. One round trip, one effect, one answer. Everything below is what happens when one of those three things goes wrong.
2

The fallacies, and the network as the enemy

eight false assumptions
  1. The network is reliable. Use timeouts, plus retries with backoff and jitter on idempotent calls.
  2. Latency is zero. Batch, parallelise, cache, and give every request path a budget.
  3. Bandwidth is infinite. Send deltas, compress, paginate.
  4. The network is secure. Use TLS and authentication on every call.
  5. Topology, administrators and transport cost never change. Use discovery, name an owner for each hop, and put cost in the design review.
  6. The network is homogeneous. Version your APIs, read tolerantly, and test on the worst network your users have.
code
// retries done properly: only idempotent calls, capped, exponential, with full jitter
async function withRetry<T>(fn: () => Promise<T>, { tries = 4, base = 100, cap = 2_000 } = {}) {
  for (let attempt = 0; ; attempt++) {
    try { return await fn(); }
    catch (e) {
      if (attempt + 1 >= tries || !isRetryable(e)) throw e;     // 4xx are not retryable
      const backoff = Math.min(cap, base * 2 ** attempt);
      await sleep(Math.random() * backoff);                      // full jitter: no thundering herd
    }
  }
}
retry storms
If every layer retries three times, five layers turn one failure into 35 = 243 requests against the dependency that is already struggling. Retry at one layer only (usually the edge or the caller nearest the user), use retry budgets and circuit breakers, and let deeper layers fail fast.
THE EIGHT FALLACIES
the assumptions every first distributed design makes, and the failure each one causes in production
swipe the figure sideways, or tap expand for full screen
1/6
reliable
The network is reliable. It drops, duplicates and reorders packets; connections reset; load balancers time out idle connections; a cable is unplugged. Defence: timeouts on every call, retries with backoff and jitter on idempotent operations, and designs that tolerate any single message vanishing.