Chapter 02 · Distributed Systems Foundations: Nodes, Networks, Failure, and State
Partial Failure: Why a Distributed System Can Be Both Working and Broken at the Same Time
Make partial failure concrete with asymmetric reachability, a slow node, failed durable storage, and a stale replica so “cluster up” is replaced by operation-scoped evidence.
Learning outcomes
An AtlasMart status dashboard says the database cluster is “green.” Checkout still times out for one customer segment, an operator in another location can reach a different subset of nodes, and a database process answers health checks even though its local disk has failed. Nothing about this situation is paradoxical. It is the defining operational property of distributed systems: components and communication paths can fail independently, and observers can have different evidence at the same moment.
Define partial failure and asymmetric reachability without equating a timeout with proof that a remote node has crashed.
Show how a slow node, failed disk, stale replica, and network path failure create different client-visible outcomes while other work continues.
Explain why “the cluster is up” and “the process answers ping” are insufficient correctness statements.
Build path-specific health evidence using reachability, latency/deadline, durable-storage state, and freshness/version requirements.
Design degraded behavior and observability around user operations rather than one global green/red health indicator.
Partial failure means some participants or communication paths fail while others continue. Asymmetric reachability means A can reach B while C cannot, or one direction works while another does not. Stale replica means a copy is behind the version/freshness required by the operation. A timeout says a result did not arrive before a deadline; it does not reveal the remote component's complete state or whether it performed the operation.
1. Local failure is observable; remote failure is inferred
If an in-process function returns an exception, the caller directly observes the failure in the same address space. A remote call crosses queues, sockets, network devices, proxies, schedulers, and another process. When the caller's deadline expires, many histories are possible: the request never left, arrived late, executed and the response was lost, reached a slow server, waited behind queued work, or reached a server whose process was alive but storage was unhealthy.
This uncertainty is why distributed systems use failure detectors, health probes, heartbeats, leases, deadlines, and protocol-specific membership state rather than pretending they have perfect knowledge. Chandra and Toueg's classic failure-detector work formalizes the idea that failure detectors can make mistakes in asynchronous systems. This lesson stays operational: treat “suspected/unreachable/late” as evidence, not metaphysical proof that a machine is dead.
process_up=true, disk_ok=false, and
version=41 can all be true simultaneously. A
liveness probe that checks only process responsiveness may
therefore classify a node as healthy for traffic that requires
durable writes or fresh reads.
2. Four partial failures produce four different risks
| Failure | What still works | What can break | Useful evidence |
|---|---|---|---|
| Asymmetric network reachability | Some clients/nodes communicate normally | Routing, replication, leader access, repair, or one region/client path | Per-source reachability matrix, connection errors, packet/RPC traces |
| Slow node / queue saturation | Process and disk may be correct eventually | Deadlines, tail latency, retries, coordinator waiting, overload amplification | Latency percentiles, queue depth, in-flight requests, deadline-exceeded counts |
| Disk/storage failure | Process may answer health/API calls | Durable writes, WAL/fsync, reads of damaged/unavailable files, recovery | I/O errors, disk health, failed fsync/write, volume metrics |
| Stale replica | Reads may return validly encoded data | Freshness/session/invariant expectations | Replica version/offset/lag, last-applied index/timestamp, read source |
The correct response depends on the operation. A stale search index may be acceptable for 30 seconds; a stale inventory reservation may violate a checkout invariant. A slow analytics replica can be removed from interactive traffic while remaining useful for offline work. A process with failed durable storage should not acknowledge writes merely because its HTTP health endpoint returns 200.
3. Different observers can report different realities
The lab has three nodes. n1 is healthy and fast.
n2 is healthy but takes 280 ms, beyond the 120 ms
client deadline. n3 has an answering process,
failed durable storage, and stale version 41.
client-east can reach n1 and n2 but not n3.
operator-west can reach n2 and n3 but not n1.
A naive dashboard computes
cluster_up = any(process_up) and returns true.
Meanwhile a checkout write requires two durable acknowledgements
within its deadline. From client-east, only n1 satisfies
reachability + durable storage + deadline, so that write path
cannot satisfy its policy. From operator-west, the only node
responding inside the deadline is n3, which is stale and has
failed storage. Thus “cluster up” coexists with “checkout
unavailable” and “operator sees only unsafe/stale candidate.”
from dataclasses import dataclass@dataclassclass Node: process_up: bool disk_ok: bool version: int latency_ms: intnodes = { "n1": Node(True, True, 42, 18), "n2": Node(True, True, 42, 280), # alive but slow "n3": Node(True, False, 41, 25), # process answers; durable storage failed and data is stale}reachability = { "client-east": {"n1": True, "n2": True, "n3": False}, "operator-west": {"n1": False, "n2": True, "n3": True},}TIMEOUT_MS = 120REQUIRED_VERSION = 42WRITE_ACKS = 2print("reachability:")for client, row in reachability.items(): visible = ", ".join(f"{n}={'yes' if ok else 'no'}" for n, ok in row.items()) print(f" {client:13s} {visible}")print("node_state:")for name, n in nodes.items(): print(f" {name} process_up={n.process_up} disk_ok={n.disk_ok} version={n.version} latency_ms={n.latency_ms}")# A deliberately weak cluster-level health check.cluster_up = any(n.process_up for n in nodes.values())print(f"naive_cluster_up={cluster_up}")# Evaluate the checkout write from client-east.durable_acks = []for name, n in nodes.items(): if not reachability["client-east"].get(name, False): continue if n.process_up and n.disk_ok and n.latency_ms <= TIMEOUT_MS: durable_acks.append(name)write_ok = len(durable_acks) >= WRITE_ACKSprint(f"checkout_write durable_acks={durable_acks} required={WRITE_ACKS} result={'success' if write_ok else 'timeout/insufficient-acks'}")# Evaluate a freshness-sensitive read from operator-west.read_candidates = []for name, n in nodes.items(): if reachability["operator-west"].get(name, False) and n.process_up and n.latency_ms <= TIMEOUT_MS: read_candidates.append((name, n.version, n.disk_ok))print(f"operator_read_candidates={read_candidates}")unsafe_nearest = min(read_candidates, key=lambda x: nodes[x[0]].latency_ms)print(f"unsafe_nearest_read={unsafe_nearest[0]} version={unsafe_nearest[1]} disk_ok={unsafe_nearest[2]}")fresh_candidates = [x for x in read_candidates if x[1] >= REQUIRED_VERSION and x[2]]print(f"fresh_durable_candidates={fresh_candidates}")print("conclusion=the cluster can answer pings while a specific write path is unavailable and another observer sees different usable nodes")
Verified deterministic output
reachability: client-east n1=yes, n2=yes, n3=no operator-west n1=no, n2=yes, n3=yesnode_state: n1 process_up=True disk_ok=True version=42 latency_ms=18 n2 process_up=True disk_ok=True version=42 latency_ms=280 n3 process_up=True disk_ok=False version=41 latency_ms=25naive_cluster_up=Truecheckout_write durable_acks=['n1'] required=2 result=timeout/insufficient-acksoperator_read_candidates=[('n3', 41, False)]unsafe_nearest_read=n3 version=41 disk_ok=Falsefresh_durable_candidates=[]conclusion=the cluster can answer pings while a specific write path is unavailable and another observer sees different usable nodes
4. Deliberately wrong approach: one global health bit
Suppose AtlasMart's load balancer routes traffic whenever
/health returns 200 from any database process. This
check is useful for a narrow question—“can this process execute
the health handler?”—but it does not test that a specific
partition is reachable, a write can be durably acknowledged, a
read meets its freshness requirement, or a critical dependency
can satisfy the deadline.
The repair is layered health. Keep cheap local liveness checks for process supervision. Add readiness checks that test the local dependencies needed before accepting traffic. More importantly, monitor black-box service-level indicators for actual user paths: checkout write success, order read latency, search freshness, replication lag, durable-write errors, and region/tenant breakdowns. Health should answer a scoped question instead of flattening all service behavior into one boolean.
Do not make readiness probes so expensive that they become the incident. A probe that performs full distributed consensus, writes production records, or scans every partition on every request can create load or dependencies of its own. Use targeted synthetic transactions and service-level monitoring with known blast radius.
5. A timeout is not the same as a known failure
In the lab, n2's 280 ms response exceeds a 120 ms deadline. The caller may abandon the request, but n2 can still be working. If the operation is a write, it may commit after the caller stops waiting. This creates commit ambiguity: the client does not know whether retrying repeats a side effect. Lesson 3 turns that ambiguity into a duplicate-charge failure and then repairs it with idempotency, bounded retries, backoff, jitter, and backpressure.
The distinction also matters for failover. Reassigning ownership because one node is merely slow can create two actors that both believe they own the same work unless the protocol uses terms/epochs/leases/fencing correctly. Chapter 08 will treat failure detectors as suspicion and show why split-brain prevention needs stronger ownership evidence than “I could not ping the old owner.”
6. Design degraded modes before an incident
Partial failure is easier to survive when the application knows which capabilities are optional. AtlasMart can define: checkout remains authoritative and may fail closed if order/payment invariants cannot be protected; product search may fail open to a simpler catalog browse with explicit freshness limits; recommendation/fraud-enrichment can be omitted from the synchronous path if their dependency is slow; analytics can lag without blocking order placement.
A degraded mode must preserve security and invariants. “Serve stale” is not safe for authorization, tenant boundaries, revocation, or money merely because it improves availability. Every degraded path needs a scoped data source, freshness bound, user-visible behavior, metrics, recovery trigger, and test.
7. Production judgment and bridge to timeout/retry dynamics
Distributed health is multi-dimensional. Record health by operation, partition/shard, client region, node role, and dependency. Alert on symptoms users experience—errors, latency, freshness, rejected writes—while retaining white-box evidence such as disk errors, queue depth, replication lag, and reachability to diagnose causes. “Three nodes running” is inventory, not an SLO.
Lesson 3 asks what clients do when evidence is incomplete. A timeout often triggers a retry, and retries can transform a small partial failure into duplicate side effects or a system-wide overload. We will model both paths.
Verification checklist
- The two clients have different reachability to the same three nodes.
- n2 is alive but too slow for the configured request deadline.
- n3 answers as a process while durable storage is failed and its version is stale.
- The naive global health value is true while the checkout path fails its write policy.
- A freshness-sensitive read has no safe candidate for operator-west in the model.
- No failure detector or timeout is described as perfect proof of remote death.
Check your understanding
- Why can a process health check succeed while a durable write path is unsafe?
- What does asymmetric reachability mean?
- Why is a timeout weaker evidence than a local function exception?
- What is wrong with describing the entire cluster as simply up or down?
- When can serving stale data be a dangerous degraded mode?
Review the answers
The process can execute code while its durable storage, replication path, partition ownership, or required dependencies are unhealthy.
Different observers or directions have different network connectivity, so one client/node can reach a participant that another cannot.
The timeout only says no result arrived before the deadline; the request may have been delayed, executed, committed, or had its response lost.
Different partitions, operations, clients, replicas, and dependencies can have different outcomes at the same time; a single bit hides that scope.
When freshness is part of correctness or security—for example inventory reservation, authorization, revocation, or tenant isolation—stale data may violate invariants rather than merely reduce quality.
Authoritative references
- A Note on Distributed Computing — Classic paper explaining why latency, concurrency, and partial failure prevent remote calls from behaving like local calls.
- Unreliable Failure Detectors for Reliable Distributed Systems — Chandra and Toueg primary work formalizing imperfect failure-detection information in asynchronous distributed systems.
- Google SRE: Monitoring Distributed Systems — Operational guidance distinguishing symptoms from causes and using user-visible latency, traffic, errors, and saturation.
- Google SRE: Service Level Objectives — Guidance for defining measurable service behavior rather than relying on component-up indicators.