Skip to content
Course syllabus

Lab 02: One node is not a system

Add replicated storage using an existing consensus implementation. Stop the leader and partition the network.

Lab source, setup, and report template

Lab 2. A majority, a live minority, and a returning replica

Continue Lab1 using the same three-member etcd store. Separate safety of acknowledged operations from availability, and leader replacement from restoration of redundancy. Prerequisites: a corrected Lab1, quorums, and the Raft lecture. About 90 minutes is an author’s planning estimate.

Run

After README setup, run from the repository root:

PYTHONDONTWRITEBYTECODE=1 /tmp/ds-course-venv/bin/python labs/distributed-systems/runner.py --lab 2 --solution starter --report /tmp/lab02.json

A unique Compose project runs three real etcd members using Raft. We neither implement consensus nor substitute an in-memory model. Python uses the JSON gateway. Select --solution reference for comparison; failure assertions are shared.

Assignment

  1. Draw the peer replication network and the control HTTP network. Identify the DNS alias available only on the peer network. Use evidence to show that the client can still reach the isolated process.
  2. Inspect endpoint selection and timeouts in common.py. Explain why endpoint failover does not authorize silently replacing linearizable reads with stale reads. Retry uncertain writes with the original Idempotency-Key.
  3. Distinguish three intervals: observation of a new leader, the first successful client request, and return of an up-to-date third copy. Define each before calling it “recovery time.”

Three scenarios

  • Leader stop. The runner identifies the actual leader with etcdctl endpoint status, stops that container, retries a write through remaining endpoints, and reads a previously acknowledged task. It then restarts the old member.
  • Live minority. Only a follower’s peer network is disconnected. /version and the control network remain reachable. A write and a linearizable read pinned strictly to that endpoint must not return success; an explicitly serializable read of an earlier value may succeed. The majority continues accepting tasks. The peer network is restored with its original alias and IP.
  • Lagging replica. A follower stops while the majority accepts eight tasks. After restart, checks read those tasks through that member and verify an increased applied index. This short experiment exercises log catch-up; it does not force snapshot installation.

Invariant: previously acknowledged tasks survive these faults, and a live minority cannot acknowledge a new write without a majority. A missing acknowledgment does not prove the operation can never apply after recovery: the report records the probe’s later outcome separately. Availability depends on a majority, networking, disks and client deadlines.

Deliverables

  • JSON report, network diagram, and a short invoke/response/unknown history for keys before, during and after partition.
  • Proof of a live minority: /version, the pinned endpoint, failed consensus operation, and a successful majority write.
  • Applied indexes before/after and an explanation of why an index alone does not prove the correct user-visible result.
  • An ADR for read semantics and a report with separate intervals. Include attempts and errors, including timeouts.

Reading and defense

Read Raft on leader election, safety and replication; Dynamo for a different availability choice; and etcd API guarantees. The Raft, Dynamo and Spanner research decks support the discussion, but their local addresses are not public course pages.

Explain why a network fault differs from stopping a process; why etcd’s serializable range is not the same term as serializable transaction isolation; why a leader alone is insufficient for successful operations; and which guarantee is lost when a client silently switches to a stale replica after timeout.

Recovery and scope

The runner restores peer networking in finally and removes only its own project. After SIGKILL, use the saved compose.env and README instructions; do not alter unrelated networks. All containers share one physical host. We test crash/partition and short catch-up, not Byzantine faults, a datacenter outage, or destruction of majority disks.