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.

Beginner95–115 minutesPartial-failure + observer-view labVendor-neutral · deterministic Python modelLast reviewed: August 2026

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.

01

Define partial failure and asymmetric reachability without equating a timeout with proof that a remote node has crashed.

02

Show how a slow node, failed disk, stale replica, and network path failure create different client-visible outcomes while other work continues.

03

Explain why “the cluster is up” and “the process answers ping” are insufficient correctness statements.

04

Build path-specific health evidence using reachability, latency/deadline, durable-storage state, and freshness/version requirements.

05

Design degraded behavior and observability around user operations rather than one global green/red health indicator.

Terminology

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.

Do not collapse evidence.

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.”

python · model asymmetric reachability, slowness, failed storage, and stale state
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

text · observer-specific health and request-path evidence
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

  1. Why can a process health check succeed while a durable write path is unsafe?
  2. What does asymmetric reachability mean?
  3. Why is a timeout weaker evidence than a local function exception?
  4. What is wrong with describing the entire cluster as simply up or down?
  5. 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

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.