Chapter 08 · Membership, Failure Detection, Gossip, and Cluster Coordination
Split Brain, Fencing, Leases, Epochs/Terms, and Avoiding Two Active Owners
Prevent a paused or partitioned old owner from corrupting state by combining ownership generations with fencing enforced at the downstream resource.
Learning outcomes
AtlasMart's inventory allocator pauses after obtaining ownership token 41. During the pause, its lease expires and another allocator receives token 42. If the first process later resumes and the storage layer accepts its old write, the system can corrupt inventory despite having used a lock.
Define split brain, lease, epoch/term, quorum ownership, and monotonically increasing fencing token.
Explain why a lease or lock without downstream fencing can be unsafe after pauses/partitions.
Trace two competing owners through a partition and show which writes the storage resource must reject.
Distinguish clock-bounded lease assumptions from logical ownership generations.
Design failover so stale owners lose authority even if they remain alive and later reconnect.
The dangerous question is not merely “who currently holds the lock?” It is “can an old holder still reach the resource and issue a side effect after a newer holder took over?” Fencing makes the resource answer that question.
1. Split brain is conflicting authority
Split brain is a state in which two partitions/processes both act as if they hold an exclusive role or ownership. A lease grants authority for a bounded interval under timing assumptions. An epoch or term is a monotonically increasing logical generation for leadership/ownership. A fencing token is a monotonically increasing value attached to operations so the downstream resource can reject stale generations.
A quorum can prevent two disjoint minorities from both obtaining a new authoritative term when quorum intersection and membership rules are correct. But the old owner may still be paused with outstanding I/O. The resource must distinguish its stale requests from the new owner's requests.
| Mechanism | Helps with | Residual risk |
|---|---|---|
| lease | bounds how long owner should act | clock/timing assumptions; paused owner can resume late |
| term/epoch | orders ownership generations | resource must enforce generation where side effects occur |
| quorum ownership | prevents two independent new owners under assumptions | old in-flight/stale actor may still send requests |
| fencing token | lets resource reject stale actor | all protected side-effect paths must carry/check token |
2. Why a lock alone can fail
t0 allocator A acquires lease, token=41
t1 A reads inventory=9
t2 A pauses before write (scheduler/GC/network stall)
t3 lease expires; quorum grants allocator B token=42
t4 B writes inventory=8 with token=42
t5 A resumes and sends delayed inventory=7 with token=41
without fencing at storage: A's late write may be accepted
with fencing: storage has max_token=42 and rejects token=41
The lease tells A when it should stop, but a paused process cannot execute “stop.” A clock correction or long pause can also make local timing misleading. Fencing moves the safety check to the resource that sees both generations.
3. Epochs, terms, leases, and quorum are complementary
Many coordination systems assign a monotonically increasing generation each time exclusive ownership changes. Consensus protocols often call this a term/ballot/epoch; lock services may expose a sequencer; application designs can persist a version. The exact source of the number matters: it must not be reusable by two conflicting owners as if both were newest.
Leases can reduce coordination for repeated operations while their assumptions hold. A quorum can authorize a new generation. Fencing then ensures storage, queue, payment gateway wrapper, or other side-effect boundary rejects requests from an older generation. If the downstream API cannot enforce fencing, the application may need an alternate single-writer design, conditional update, idempotent command protocol, or stronger transaction boundary.
4. Wrong approach: “we use a distributed lock, so split brain is impossible”
A lock service can correctly transfer ownership while the protected resource remains vulnerable to delayed old requests. Another error is to enforce fencing only in one code path while a maintenance script or retry worker can write without a token. Safety then exists only on paper.
The repair is end-to-end: every mutating path carries the generation, the storage boundary atomically compares it with the highest accepted generation, stale generations are rejected and observable, and failover tests explicitly pause the old owner across handoff.
5. AtlasMart lab — stale owner rejected by fencing
This local simulation first appends writes to an unfenced history, then uses a tiny store that remembers the largest accepted token. No real locks, threads, clocks, or network partitions are involved.
class FencedStore:
def __init__(self):
self.max_token = 0
self.value = None
def write(self, token, value):
if token < self.max_token:
return f"REJECT token={token} < max={self.max_token}"
self.max_token = token
self.value = value
return f"ACCEPT token={token} value={value}"
print("UNFENCED history:")
unfenced=[]
unfenced.append(("A", 41, "inventory=9")) # old owner pauses afterward
unfenced.append(("B", 42, "inventory=8")) # new owner after lease/term change
unfenced.append(("A", 41, "inventory=7")) # stale A resumes late
for row in unfenced:
print(row)
print("FAILURE: stale owner A overwrites newer owner B when resource ignores epochs")
print("\nFENCED history:")
store=FencedStore()
print("A:", store.write(41, "inventory=9"))
print("B:", store.write(42, "inventory=8"))
print("A resumes:", store.write(41, "inventory=7"))
print("final:", store.value, "max_token=", store.max_token)
print("SAFETY CONDITION: downstream resource must actually compare fencing tokens")
Expected output
UNFENCED history:
('A', 41, 'inventory=9')
('B', 42, 'inventory=8')
('A', 41, 'inventory=7')
FAILURE: stale owner A overwrites newer owner B when resource ignores epochs
FENCED history:
A: ACCEPT token=41 value=inventory=9
B: ACCEPT token=42 value=inventory=8
A resumes: REJECT token=41 < max=42
final: inventory=8 max_token= 42
SAFETY CONDITION: downstream resource must actually compare fencing tokens
The important verification is
REJECT token=41 < max=42. Fencing does not prove
that B's business decision is correct; it proves only that an
older ownership generation cannot overwrite a newer one through
this protected resource.
Check your understanding
- Why can a paused lock holder remain dangerous?
- What makes a fencing token useful?
- Does a lease remove the need for timing assumptions?
- Where must fencing be enforced?
- What if a downstream system cannot compare tokens?
Review the answers
1. It may resume after ownership changed and issue an already-prepared or delayed side effect.
2. It monotonically orders ownership generations and is checked by the downstream resource.
3. No. Leases explicitly depend on bounded timing/clock behavior and expiry interpretation.
4. At every side-effect boundary that stale owners could reach, not merely in the lock-acquisition code.
5. Use another enforceable concurrency boundary, conditional update, single-owner architecture, or redesign the side effect rather than claiming fencing exists.
6. Production judgment
Track term/epoch changes, lease-expiry handoffs, stale-token rejections, double-owner alerts, quorum health, and the identities performing administrative failover. Treat forced promotion as a high-risk operation: document which owner is fenced, which data may be stale, and how rollback avoids reactivating an older generation.
Security and correctness align here. If an attacker can mint a larger fencing token, change membership, or bypass the fenced write path, the protection collapses. Apply least privilege to coordination and storage APIs, audit generation changes, and test delayed-owner recovery. The final lesson turns these primitives into a safe rolling-maintenance procedure.
Authoritative references
- The Chubby lock service for loosely-coupled distributed systems — lock service design including sequencer-style protection against stale lock holders
- In Search of an Understandable Consensus Algorithm (Raft) — terms and quorum-based leadership; consensus is covered in depth later in the course
- SWIM membership protocol — membership/failure-detection context; dissemination alone does not fence owners