Expand description
Read/write epoch consensus that survives a restart.
Status: implementation. Space: bounded by membership, plus what the stubborn children hold
outstanding — which nothing here retires. See the departure on Stop.
Cachin, Guerraoui & Rodrigues, Module 5.7 (LoggedEpochConsensus) and Algorithm 5.9 (“Logged
Read/Write Epoch Consensus”), quoted from the book:
Algorithm 5.9: Logged Read/Write Epoch Consensus
Implements: EpochConsensus, instance lep, with timestamp ets and leader ℓ.
Uses:
StubbornPointToPointLinks, instance sl;
StubbornBestEffortBroadcast, instance sbeb;
upon event ⟨ lep, Init | state ⟩ do
(valts, val) := state;
store(valts, val);
tmpval := ⊥;
states := [⊥]^N;
accepted := 0;
upon event ⟨ lep, Recovery ⟩ do
retrieve(valts, val);
upon event ⟨ lep, Propose | v ⟩ do // only leader ℓ
tmpval := v;
trigger ⟨ sbeb, Broadcast | [READ] ⟩;
upon event ⟨ sbeb, Deliver | ℓ, [READ] ⟩ do
trigger ⟨ sl, Send | ℓ, [STATE, valts, val] ⟩;
upon event ⟨ sl, Deliver | q, [STATE, ts, v] ⟩ do // only leader ℓ
states[q] := (ts, v);
upon #(states) > N/2 do // only leader ℓ
(ts, v) := highest(states);
if v ≠ ⊥ then tmpval := v;
states := [⊥]^N;
trigger ⟨ sbeb, Broadcast | [WRITE, tmpval] ⟩;
upon event ⟨ sbeb, Deliver | ℓ, [WRITE, v] ⟩ do
(valts, val) := (ets, v);
store(valts, val);
trigger ⟨ sl, Send | ℓ, [ACCEPT] ⟩;
upon event ⟨ sl, Deliver | q, [ACCEPT] ⟩ do // only leader ℓ
accepted := accepted + 1;
upon accepted > N/2 do // only leader ℓ
accepted := 0;
trigger ⟨ sbeb, Broadcast | [DECIDED, tmpval] ⟩;
upon event ⟨ sbeb, Deliver | ℓ, [DECIDED, v] ⟩ do
epochdecision := v;
store(epochdecision);
trigger ⟨ lep, Decide | epochdecision ⟩;
upon event ⟨ lep, Abort ⟩ do
trigger ⟨ lep, Aborted | (valts, val) ⟩;
halt; // stop operating when abortedThe safety argument is crate::epoch_consensus’s and is not restated here: two majorities
intersect, so a later epoch’s read reaches a process that accepted whatever an earlier epoch
decided, and if v ≠ ⊥ then tmpval := v makes the later leader adopt it. What this module adds
is that the argument still holds when the processes holding that intersection go down and come
back.
§Two store calls, one metadata value
The book writes store(valts, val) and store(epochdecision) as separate calls. Cx::storage
offers one rewritten metadata value and an appended sequence, so both land in one Durable
that is rewritten each time. Nothing accumulates: an epoch accepts at most one value and decides
at most one, so the record is a fixed size and Entry is uninhabited.
§Durable before visible, twice, and both in the handler’s own text
store(valts, val); trigger ⟨ sl, Send | ℓ, [ACCEPT] ⟩ — the acceptance is a promise to a
quorum. A process that told the leader it had accepted v at ets, and then came back with
no record of it, would answer a later epoch’s read with an empty state; the later leader would
find nothing in the intersection and be free to write something else, after v had already been
decided. That is EPC4 failing, and it fails silently.
store(epochdecision); trigger ⟨ lep, Decide | v ⟩ — same shape one step later. The layer above
reads epochdecision back on recovery (LoggedEpochConsensus::epoch_decision) and that is how
Algorithm 5.10 knows a process had decided before it went down.
Both orders are written here, in these handlers, and not left to a driver to arrange by
buffering effects until the handler returns. Cx supports eager sinks, so buffering is not
something this code may assume.
§Departure: repeats are idempotent, and the standing conditions fire once
crate::epoch_consensus runs over perfect links, which deliver each message once. This one
runs over stubborn ones, which must not deduplicate — repeating for ever is what reaches a
process that was down when the message was sent. So every handler here sees its message many
times, and the book’s counters do not survive that:
accepted := accepted + 1counts messages, and one process’s ACCEPT arrives for ever. The count would passN/2on its own with a single acceptance in the whole run. It is a set of the processes that have accepted, so a repeat adds nothing.upon #(states) > N/2andupon accepted > N/2are standing conditions the book re-arms by clearing what they count. Clearing is not enough when the messages come back:writtenandannouncedmake each fire once, as they already do incrate::epoch_consensus.storeis a rewrite of the same value on a repeat, which is idempotent, but it is still a write.WRITEis applied only when it changes something, so the write count stays one per acceptance and a test can check that rather than take it on trust.- A follower answers
READonce andWRITEonce. The book answers every delivery. Over a stubborn link the answer is itself retransmitted until this instance ends, so a second answer to a redeliveredREADis a second stubborn transmission carrying the same content — and since redeliveries never stop, neither would the transmissions. Measured before this guard: the send rate grew linearly in time, 12.6k → 76.6k per 400 ms across five windows, with nothing faulty. One answer is enough for the same reason retransmission exists at all: a leader that crashed and came back re-proposes, and what reaches its new incarnation is the follower’s original reply, still going. A follower that crashes forgets it answered and answers again, which is correct — its link forgot the transmission too.
§Departure: messages carry the epoch they belong to
As in crate::epoch_consensus: instances are addressed lep.ets in the book and by nothing at
all on a real wire, so a WRITE from epoch 7 arriving after epoch 11 began would be accepted and
recorded at timestamp 11 — an acceptance that never happened. The stamp is in Tagged, inside
the instance, because the epoch is the instance’s own identity.
§Departure: nothing calls Stop
As in crate::logged_epoch_change. The stubborn children retransmit until retired and nothing
retires them, so space grows with the number of distinct messages an epoch sends rather than
with the membership. Bounded in practice by the epoch ending, which is what Abort is for.
EPC1 [always] Validity — a decided value was proposed in this epoch, or was the highest-
timestamped value some process had already accepted
EPC2 [always] Uniform agreement — no two processes decide differently in one epoch
EPC3 [always] Integrity — a process decides at most once
EPC4 [always] Lock-in — a value decided in an earlier epoch is what a later one decides, and
**this holds across a crash**: what a process accepted is read back on recovery
EPC5 [always] Abort behaviour — an abandoned instance reports its state and then is silentStructs§
- Durable
- Everything this instance keeps durably, as one rewritten value.
- Logged
Epoch Consensus - Abortable consensus within one epoch, whose acceptances survive a restart.
- State
(valts, val)— what a process has accepted, and when.- Tagged
- A message stamped with the epoch it belongs to.