Expand description
Epoch-change — a sequence of epochs, each with a timestamp and a leader.
Status: implementation. Space: bounded by membership.
Cachin, Guerraoui & Rodrigues, Module 5.3 and Algorithm 5.5 (“Leader-Based Epoch-Change”), quoted from the book:
Algorithm 5.5: Leader-Based Epoch-Change
Implements: EpochChange, instance ec.
Uses:
PerfectPointToPointLinks, instance pl;
BestEffortBroadcast, instance beb;
EventualLeaderDetector, instance Ω.
upon event ⟨ ec, Init ⟩ do
trusted := ℓ0;
lastts := 0;
ts := rank(self);
upon event ⟨ Ω, Trust | p ⟩ do
trusted := p;
if p = self then
ts := ts + N;
trigger ⟨ beb, Broadcast | [NEWEPOCH, ts] ⟩;
upon event ⟨ beb, Deliver | ℓ, [NEWEPOCH, newts] ⟩ do
if ℓ = trusted ∧ newts > lastts then
lastts := newts;
trigger ⟨ ec, StartEpoch | newts, ℓ ⟩;
else
trigger ⟨ pl, Send | ℓ, [NACK] ⟩;
upon event ⟨ pl, Deliver | p, [NACK] ⟩ do
if trusted = self then
ts := ts + N;
trigger ⟨ beb, Broadcast | [NEWEPOCH, ts] ⟩;§Why timestamps are unique without anyone coordinating
ts := rank(self) and ts := ts + N. Each process therefore draws from its own residue class
modulo N, and two processes cannot mint the same timestamp however far apart they drift. A
plain counter would not have that property, and the layer above uses the timestamp to order
writes — so two epochs sharing one would make its safety argument meaningless.
§The NACK, and what it is for
A process that receives a NEWEPOCH it will not act on — because the sender is not who it
trusts, or because the timestamp is not newer than one it has already started — answers NACK.
A leader that is nacked bumps its timestamp and tries again.
Without it a leader whose timestamp has fallen behind another’s would broadcast for ever and never be started by anyone. The NACK is what lets it discover that and climb past.
§Departure: the NACK travels by directed broadcast, not by a separate link
The book names two message children, pl for the NACK and beb for the NEWEPOCH. Here there is
one: crate::best_effort_broadcast gained a directed beb::Cmd::SendTo in the
link-parameterisation change — same wire message, same link, strictly fewer recipients, no
new communication step — and a NACK sent that way is a perfect-link send with an extra layer’s
name on it.
What this buys is one child fewer and one wire variant fewer. What it costs is that the module
no longer mirrors the book’s Uses: line exactly, so it is recorded here rather than left for a
reader to notice. Nothing about the guarantee changes: beb::Cmd::SendTo reaches exactly the
one addressed process, which is what pl, Send does.
§Departure: a leader is told where the processes trusting it have reached
⟨ Ω, Trust | p ⟩ is raised when the trusted process changes, and Algorithm 5.5 announces an
epoch only on that edge. So a process that trusted itself all along is never told it has become
everyone’s leader — and if the others ran their epochs ahead under other leaders while it did
not, nothing it would announce is high enough for them to accept, and nothing prompts it to climb.
Measured, five processes partitioned [A,B] [C] [D,E] and then healed: Ω converges correctly on
E, and afterwards trusted = [E,E,E,E,E] with lastts = [32, 27, 43, 10, 10]. E is trusted
by everyone, sits in epoch 10, and announces nothing — for ever. Retransmission does not rescue
it either: what E re-sends is NEWEPOCH(10), which its recipients’ links deduplicate, so it
draws no refusal. E received zero in thirty timeouts.
This is a gap in Algorithm 5.5 composed with Algorithm 2.8, rather than in either module’s
specification: Module 2.9 says only that Ω eventually agrees, and an Ω that re-raised Trust
would leave 5.5 correct as written. It was unreachable while the detector beneath was
crate::perfect_failure_detector, whose accusations are permanent, because a partition never
healed for it.
So a process whose trusted leader changes to one that is not itself, while its current epoch was started by somebody else, tells that leader the timestamp it has reached. The leader then chooses its next candidate above what it was told. Nothing is sent while nothing has changed: the report rides the same edge the announcement does.
§Departure: a refused leader climbs past the refusal in one step
The same message carries the refuser’s own lastts, and the leader jumps its candidate above it
rather than adding N. Algorithm 5.5 steps once per refusal, which costs a round trip per step —
the gap above is 33, or seven round trips. Boundedness is unaffected, and for the same reason as
before: after the jump the leader’s candidate is strictly above what it was told, so a repeated
report names a timestamp it has already passed and moves nothing.
§Departure: the NACK names the timestamp it refuses
Algorithm 5.5 sends a bare [NACK] and bumps ts on every one that arrives. Algorithm 5.8 —
the same abstraction in the fail-recovery model — sends [NACK, nts] and guards the handler
with such that nts = ts. This module takes 5.8’s form, because over a link that retransmits,
5.5’s does not terminate.
The loop: the leader broadcasts NEWEPOCH(t); every process that does not yet trust it answers
NACK; each NACK bumps ts and broadcasts again, so one announcement to N processes produces
N − 1 further announcements, each of which produces its own. The stubborn link beneath resends
everything it has ever sent, so nothing decays. Measured before the guard: a five-process run
with one crash reached epoch 647,309 and 2.3 million sends inside a second of virtual time,
and no epoch ever lasted long enough for the consensus above it to finish a write. Measured
after: single figures.
With the guard, an announcement is answered at most once — the first NACK moves ts, and every
later NACK naming the old timestamp is for an announcement already superseded. That is the whole
of the fix, and the book states it one algorithm later.
EC1 [always] Monotonicity — timestamps strictly increase, and one timestamp names one leader
EC2 [eventual] Consistency — eventually every correct process starts the same last epochStructs§
- Epoch
Change - A sequence of epochs, driven by who is trusted.
Enums§
- Epoch
Msg - What this layer puts on the wire, beneath the broadcast.
- Ind
- Indications to the layer above.
- Wire
- The wire, multiplexing the two children.