Chapter 02 · Distributed Systems Foundations: Nodes, Networks, Failure, and State
Latency, Timeouts, Retries, Duplicate Work, Backpressure, and the Fallacy of a Reliable Network
Trace ambiguous commits, duplicate side effects, retry amplification, backoff, jitter, retry budgets, idempotency, and backpressure with deterministic AtlasMart request models.
Learning outcomes
AtlasMart's payment client sends a charge request. The server commits the charge at 80 ms, but the response is delayed in the network. At 100 ms the client times out. From the client's perspective, “no response” is true; “the charge failed” is not. If it immediately retries a non-idempotent operation, the retry can create a second charge. If thousands of clients do the same during overload, retries can become more load than the original traffic.
Distinguish latency, deadline/timeout, failure, and commit ambiguity in a remote request timeline.
Demonstrate how retries can duplicate side effects and amplify load even when the retry policy was intended to improve reliability.
Explain idempotency keys as an application/API semantic mechanism rather than a transport guarantee of exactly-once delivery.
Apply bounded retry budgets, exponential backoff, deterministic jitter, and backpressure/admission control to reduce positive feedback.
Identify metrics that reveal retry storms: attempt/request ratio, timeout rate, queue depth, saturation, rejected load, and duplicate/idempotency hits.
Remote requests incur variable latency and can lose/delay messages. A deadline limits how long the caller waits; it does not rewind the server. When a request can change durable state, retry policy must be designed together with operation semantics.
1. A timeout creates uncertainty, not necessarily failure
Latency is elapsed time for an operation from a chosen observation point. A deadline is the latest time the caller is willing to wait; a timeout is the resulting event when that deadline expires. In a distributed system, the request and response each traverse queues and network paths. The server can finish after the caller stops waiting, or it can finish before the deadline while the response is delayed.
This produces commit ambiguity: the caller cannot tell whether a state-changing operation took effect. The safe response depends on the operation's semantics. Retrying a pure read may waste capacity but normally does not duplicate a business side effect. Retrying “charge card,” “send shipment,” or “decrement one-time coupon” can be disastrous unless the API has a way to recognize semantic duplicates.
2. Idempotency makes a repeated intent recognizable
An operation is idempotent when applying the same intended operation repeatedly has the same externally relevant effect as applying it once. HTTP methods have protocol-level idempotency conventions, but business operations still need careful semantics. A common API technique is an idempotency key: the client creates a stable identifier for one business intent, and the server records the result associated with that key.
On retry, the server detects the same key and returns/reuses the original result instead of creating a second effect. This requires a scope, retention window, authentication/tenant binding, request-shape validation, and durable deduplication state appropriate to the business risk. A key must not let tenant A replay or learn tenant B's result. “Same key, different payload” must be defined as an error or a product-specific rule, not silently accepted.
Idempotency does not create magical end-to-end exactly-once delivery. Messages can still be delivered multiple times, clients can lose responses, deduplication records can expire, and downstream side effects may need their own idempotency/transaction boundaries. The claim is narrower: repeated requests with the same semantic key can be made safe under the server's documented deduplication contract.
3. Deliberately wrong approach: immediate retries everywhere
Suppose a frontend retries a database call three times, an API gateway retries the frontend three times, and a mobile client retries the gateway three times. A single user action can multiply into many backend attempts. When the backend is already slow because it is saturated, retries add more queueing, extend latency, cause more timeouts, and trigger still more retries: a positive-feedback loop.
Exponential backoff increases delay between successive attempts. Jitter randomizes those delays so many clients do not synchronize on the same retry boundary. A retry budget caps the retry volume a process/service is allowed to create. Backpressure communicates or enforces that a downstream component cannot accept unlimited work; bounded queues, admission control, load shedding, concurrency limits, or explicit overload responses can fail some work quickly rather than letting queues grow without bound.
| Mechanism | Problem addressed | Important limit |
|---|---|---|
| Timeout/deadline | Bounds caller waiting/resource occupancy | Does not reveal whether remote work committed |
| Retry | Can recover from some transient faults | Duplicates work and increases load; unsafe for ambiguous side effects without semantics |
| Exponential backoff | Reduces retry frequency as failures persist | Still synchronized without jitter; can increase user latency |
| Jitter | Spreads retry attempts across time | Does not reduce total attempts unless paired with limits/budgets |
| Retry budget | Caps retry amplification | Some transient requests will fail rather than retry |
| Backpressure/load shedding | Protects finite capacity and queues | Requires prioritization and explicit degraded/failure behavior |
| Idempotency key | Makes duplicate business intent recognizable | Needs durable scoped deduplication semantics; not exactly-once transport |
4. Run the ambiguous-commit and retry-load simulator
The first half of the lab hard-codes a timeline: attempt 1
commits at 80 ms, the client times out at 100 ms, and the first
response arrives at 150 ms. Without an idempotency key, attempt
2 creates charge-2. With key
checkout-9001, the simulated service returns the
first result on the retry.
The second half is a discrete load model, not a benchmark. Sixty original requests produce 20 immediate retries and 10 second retries under a naive policy: 90 attempts total. The safer example allows only 10 retries and schedules them with deterministic seeded jitter, giving 70 attempts and a peak retry bucket of four. Backpressure admits eight operations from a burst of 20 and explicitly rejects 12 rather than queueing them without a bound.
from collections import Counterfrom random import RandomTIMEOUT_MS = 100idempotency_key = "checkout-9001"print("ambiguous_commit_without_idempotency:")# First request commits before the client timeout, but its response is delayed.events = [ (0, "client", "send attempt=1"), (80, "server", "attempt=1 COMMIT charge=charge-1"), (100, "client", "timeout; outcome unknown; send attempt=2"), (150, "network", "attempt=1 response finally arrives"), (160, "server", "attempt=2 COMMIT charge=charge-2"), (170, "client", "attempt=2 response arrives"),]for t, actor, event in events: print(f" t+{t:03d}ms {actor:7s} {event}")print("charges_created=2")print("ambiguous_commit_with_idempotency:")seen = {}def charge(key, charge_id): if key in seen: return seen[key], "deduplicated" seen[key] = charge_id return charge_id, "committed"first = charge(idempotency_key, "charge-1")second = charge(idempotency_key, "charge-2")print(f" attempt=1 key={idempotency_key} result={first}")print(f" attempt=2 key={idempotency_key} result={second}")print(f"unique_charges={len(set(seen.values()))}")# Deterministic retry-load model: 60 initial requests, 20 time out once, 10 of those would time out again.initial = 60first_timeouts = list(range(20))second_timeouts = list(range(10))naive_attempts = initial + len(first_timeouts) + len(second_timeouts)print("retry_load_model:")print(f" naive attempts={naive_attempts} immediate_retry_batch={len(first_timeouts)}")# Safer policy: retry budget of 10, exponential base delay with deterministic jitter spread.rng = Random(42)retry_budget = 10buckets = Counter()for req_id in first_timeouts[:retry_budget]: base = 100 jitter = rng.randrange(0, 101) # 0..100 ms, deterministic for the lab scheduled = base + jitter bucket = (scheduled // 25) * 25 buckets[bucket] += 1safer_attempts = initial + retry_budgetprint(f" budgeted attempts={safer_attempts} retry_budget={retry_budget}")print(" retry_buckets_ms=" + ", ".join(f"{k}:{buckets[k]}" for k in sorted(buckets)))print(f" peak_retry_bucket={max(buckets.values())}")# Backpressure/admission example.queue_capacity = 8burst = 20accepted = min(queue_capacity, burst)rejected = burst - acceptedprint(f"backpressure burst={burst} queue_capacity={queue_capacity} accepted={accepted} rejected_fast={rejected}")print("note=this is a discrete deterministic model, not a throughput or latency benchmark")
Verified deterministic output
ambiguous_commit_without_idempotency: t+000ms client send attempt=1 t+080ms server attempt=1 COMMIT charge=charge-1 t+100ms client timeout; outcome unknown; send attempt=2 t+150ms network attempt=1 response finally arrives t+160ms server attempt=2 COMMIT charge=charge-2 t+170ms client attempt=2 response arrivescharges_created=2ambiguous_commit_with_idempotency: attempt=1 key=checkout-9001 result=('charge-1', 'committed') attempt=2 key=checkout-9001 result=('charge-1', 'deduplicated')unique_charges=1retry_load_model: naive attempts=90 immediate_retry_batch=20 budgeted attempts=70 retry_budget=10 retry_buckets_ms=100:4, 125:3, 175:3 peak_retry_bucket=4backpressure burst=20 queue_capacity=8 accepted=8 rejected_fast=12note=this is a discrete deterministic model, not a throughput or latency benchmark
5. Choosing timeouts requires a latency model
A universal timeout such as “100 ms for every database call” is not defensible. The timeout must consider the request path, connection establishment, network distance, expected tail latency, cold starts, storage behavior, and the caller's end-to-end deadline. A timeout below normal tail latency manufactures failures and retries. A timeout so large that requests occupy concurrency indefinitely can spread overload.
Measure latency distributions, not only averages. Separate connection/setup latency from request latency where the client/library does so. Track successful and failed-request latency separately; a fast error is not a fast successful service. When an operation fans out to several nodes, the slowest required response can dominate the user's tail latency.
6. Backpressure is a correctness tool for finite systems
Every queue is a promise to do work later. An unbounded queue under sustained overload is an unbounded promise that cannot be kept; latency rises until callers time out, while work they have already abandoned may still consume resources. Bounded queues and concurrency limits make overload visible sooner. Load shedding can protect critical checkout traffic by rejecting optional recommendation/search-enrichment work before it consumes the same constrained pool.
Backpressure must be propagated deliberately. If a database client has reached a concurrency limit, upstream code should not spawn unlimited tasks that merely wait for a permit. If an overload response is retriable, clients need a bounded and jittered policy. If the request is not safe to retry, the API should surface an ambiguous/failed status that the business workflow can reconcile rather than blindly replaying the side effect.
7. Production judgment and bridge to failure domains
For every retryable AtlasMart operation, document: the end-to-end deadline, per-attempt timeout, maximum attempts, which failure codes are retryable, backoff/jitter algorithm, retry budget, concurrency/queue limit, idempotency semantics, deduplication retention, and observability labels. Test the behavior under overload and delayed responses, not only clean server crashes.
Lesson 4 changes the question from “what if one request is slow?” to “what if failures are correlated?” Three replicas do not buy three independent chances if they share a host, rack, zone, region, provider control plane, or one dangerous deployment/configuration action.
Verification checklist
- The first request commits before the client timeout in the deterministic timeline.
- Without idempotency, the retry creates a second charge.
- With the stable key, the retry returns the original charge result.
- The naive model creates 90 attempts from 60 logical requests.
- The budgeted model caps retries and jitter spreads them across time buckets.
- The lesson labels all attempt counts as simulation evidence rather than real throughput/latency measurements.
Check your understanding
- Why is a timeout not proof that a write failed?
- What problem does an idempotency key solve?
- Why can retries create a cascading failure?
- What distinct purposes do backoff and jitter serve?
- Why can a bounded queue be safer than an unbounded queue during overload?
Review the answers
The server may have received and committed the request while the response was delayed or lost; the caller only knows its deadline expired.
It lets the server recognize repeated requests that represent the same business intent and reuse the original effect/result under a defined deduplication contract.
Retries add load exactly when a dependency may already be slow or overloaded, increasing queues/timeouts and causing positive feedback.
Backoff reduces how frequently a client retries as failures persist; jitter prevents many clients from retrying at the same synchronized moments.
It bounds memory/concurrency/queued work and exposes overload quickly, allowing load shedding/degraded behavior instead of ever-growing latency and abandoned work.
Authoritative references
- AWS Builders Library: Timeouts, retries, and backoff with jitter — Operational guidance on timeouts, retry amplification, backoff, and jitter.
- AWS Builders Library: Making retries safe with idempotent APIs — Detailed treatment of retryable side effects, client request IDs, and idempotent API contracts.
- Google SRE: Addressing Cascading Failures — Guidance showing how retries amplify overload and recommending randomized backoff and retry budgets.
- Google SRE: Handling Overload — Guidance on graceful degradation, load shedding, and protecting finite serving capacity.