Skip to main content

Module logged_epoch_consensus

Module logged_epoch_consensus 

Source
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 aborted

The 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 + 1 counts messages, and one process’s ACCEPT arrives for ever. The count would pass N/2 on 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/2 and upon accepted > N/2 are standing conditions the book re-arms by clearing what they count. Clearing is not enough when the messages come back: written and announced make each fire once, as they already do in crate::epoch_consensus.
  • store is a rewrite of the same value on a repeat, which is idempotent, but it is still a write. WRITE is 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 READ once and WRITE once. The book answers every delivery. Over a stubborn link the answer is itself retransmitted until this instance ends, so a second answer to a redelivered READ is 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 silent

Structs§

Durable
Everything this instance keeps durably, as one rewritten value.
LoggedEpochConsensus
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.

Enums§

Announce
What travels by sbeb — the leader speaking to everyone.
Cmd
Requests from the layer above.
Ind
Indications to the layer above.
Reply
What travels by sl — a follower answering the leader.
Wire
The wire, multiplexing the two children the book names.