Skip to main content

recon_protocols/
epoch_change.rs

1//! Epoch-change — a sequence of epochs, each with a timestamp and a leader.
2//!
3//! **Status: implementation. Space: bounded by membership.**
4//!
5//! Cachin, Guerraoui & Rodrigues, Module 5.3 and Algorithm 5.5 ("Leader-Based Epoch-Change"),
6//! quoted from the book:
7//!
8//! ```text
9//! Algorithm 5.5: Leader-Based Epoch-Change
10//! Implements: EpochChange, instance ec.
11//! Uses:
12//!     PerfectPointToPointLinks, instance pl;
13//!     BestEffortBroadcast, instance beb;
14//!     EventualLeaderDetector, instance Ω.
15//!
16//! upon event ⟨ ec, Init ⟩ do
17//!     trusted := ℓ0;
18//!     lastts := 0;
19//!     ts := rank(self);
20//!
21//! upon event ⟨ Ω, Trust | p ⟩ do
22//!     trusted := p;
23//!     if p = self then
24//!         ts := ts + N;
25//!         trigger ⟨ beb, Broadcast | [NEWEPOCH, ts] ⟩;
26//!
27//! upon event ⟨ beb, Deliver | ℓ, [NEWEPOCH, newts] ⟩ do
28//!     if ℓ = trusted ∧ newts > lastts then
29//!         lastts := newts;
30//!         trigger ⟨ ec, StartEpoch | newts, ℓ ⟩;
31//!     else
32//!         trigger ⟨ pl, Send | ℓ, [NACK] ⟩;
33//!
34//! upon event ⟨ pl, Deliver | p, [NACK] ⟩ do
35//!     if trusted = self then
36//!         ts := ts + N;
37//!         trigger ⟨ beb, Broadcast | [NEWEPOCH, ts] ⟩;
38//! ```
39//!
40//! # Why timestamps are unique without anyone coordinating
41//!
42//! `ts := rank(self)` and `ts := ts + N`. Each process therefore draws from its own residue class
43//! modulo `N`, and two processes cannot mint the same timestamp however far apart they drift. A
44//! plain counter would not have that property, and the layer above uses the timestamp to order
45//! writes — so two epochs sharing one would make its safety argument meaningless.
46//!
47//! # The NACK, and what it is for
48//!
49//! A process that receives a `NEWEPOCH` it will not act on — because the sender is not who it
50//! trusts, or because the timestamp is not newer than one it has already started — answers `NACK`.
51//! A leader that is nacked bumps its timestamp and tries again.
52//!
53//! Without it a leader whose timestamp has fallen behind another's would broadcast for ever and
54//! never be started by anyone. The NACK is what lets it discover that and climb past.
55//!
56//! # Departure: the NACK travels by directed broadcast, not by a separate link
57//!
58//! The book names two message children, `pl` for the NACK and `beb` for the NEWEPOCH. Here there is
59//! one: [`crate::best_effort_broadcast`] gained a directed [`beb::Cmd::SendTo`] in the
60//! `link-parameterisation` change — same wire message, same link, strictly fewer recipients, no
61//! new communication step — and a NACK sent that way is a perfect-link send with an extra layer's
62//! name on it.
63//!
64//! What this buys is one child fewer and one wire variant fewer. What it costs is that the module
65//! no longer mirrors the book's `Uses:` line exactly, so it is recorded here rather than left for a
66//! reader to notice. Nothing about the guarantee changes: `beb::Cmd::SendTo` reaches exactly the
67//! one addressed process, which is what `pl, Send` does.
68//!
69//! # Departure: a leader is told where the processes trusting it have reached
70//!
71//! `⟨ Ω, Trust | p ⟩` is raised when the trusted process **changes**, and Algorithm 5.5 announces an
72//! epoch only on that edge. So a process that trusted *itself* all along is never told it has become
73//! everyone's leader — and if the others ran their epochs ahead under other leaders while it did
74//! not, nothing it would announce is high enough for them to accept, and nothing prompts it to climb.
75//!
76//! Measured, five processes partitioned `[A,B] [C] [D,E]` and then healed: Ω converges correctly on
77//! `E`, and afterwards `trusted = [E,E,E,E,E]` with `lastts = [32, 27, 43, 10, 10]`. `E` is trusted
78//! by everyone, sits in epoch 10, and announces nothing — for ever. Retransmission does not rescue
79//! it either: what `E` re-sends is `NEWEPOCH(10)`, which its recipients' links deduplicate, so it
80//! draws no refusal. `E` received **zero** in thirty timeouts.
81//!
82//! This is a gap in Algorithm 5.5 composed with **Algorithm 2.8**, rather than in either module's
83//! specification: Module 2.9 says only that Ω eventually agrees, and an Ω that re-raised `Trust`
84//! would leave 5.5 correct as written. It was unreachable while the detector beneath was
85//! [`crate::perfect_failure_detector`], whose accusations are permanent, because a partition never
86//! healed for it.
87//!
88//! So a process whose trusted leader changes to one that is **not itself**, while its current epoch
89//! was started by somebody else, tells that leader the timestamp it has reached. The leader then
90//! chooses its next candidate above what it was told. Nothing is sent while nothing has changed: the
91//! report rides the same edge the announcement does.
92//!
93//! # Departure: a refused leader climbs past the refusal in one step
94//!
95//! The same message carries the refuser's own `lastts`, and the leader jumps its candidate above it
96//! rather than adding `N`. Algorithm 5.5 steps once per refusal, which costs a round trip per step —
97//! the gap above is 33, or seven round trips. Boundedness is unaffected, and for the same reason as
98//! before: after the jump the leader's candidate is strictly above what it was told, so a repeated
99//! report names a timestamp it has already passed and moves nothing.
100//!
101//! # Departure: the NACK names the timestamp it refuses
102//!
103//! Algorithm 5.5 sends a bare `[NACK]` and bumps `ts` on every one that arrives. Algorithm 5.8 —
104//! the same abstraction in the fail-recovery model — sends `[NACK, nts]` and guards the handler
105//! with `such that nts = ts`. This module takes 5.8's form, because over a link that retransmits,
106//! 5.5's does not terminate.
107//!
108//! The loop: the leader broadcasts `NEWEPOCH(t)`; every process that does not yet trust it answers
109//! `NACK`; each NACK bumps `ts` and broadcasts again, so one announcement to `N` processes produces
110//! `N − 1` further announcements, each of which produces its own. The stubborn link beneath resends
111//! everything it has ever sent, so nothing decays. Measured before the guard: a five-process run
112//! with one crash reached epoch **647,309** and 2.3 million sends inside a second of virtual time,
113//! and no epoch ever lasted long enough for the consensus above it to finish a write. Measured
114//! after: single figures.
115//!
116//! With the guard, an announcement is answered at most once — the first NACK moves `ts`, and every
117//! later NACK naming the old timestamp is for an announcement already superseded. That is the whole
118//! of the fix, and the book states it one algorithm later.
119//!
120//! ```text
121//! EC1 [always]     Monotonicity — timestamps strictly increase, and one timestamp names one leader
122//! EC2 [eventual]   Consistency — eventually every correct process starts the same last epoch
123//! ```
124
125use recon_core::{Child, NodeId, ProtoCx, Protocol, TimerId};
126use serde::{Deserialize, Serialize};
127use std::collections::BTreeSet;
128
129use crate::best_effort_broadcast::{self as beb, BestEffortBroadcast};
130use crate::eventual_leader_detector::{self as eld, EventualLeaderDetector};
131use crate::perfect_failure_detector::Heartbeat;
132use crate::perfect_link as pl;
133use crate::{Note, Refusal, Timing};
134
135/// What this layer puts on the wire, beneath the broadcast.
136#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
137pub enum EpochMsg {
138    /// `[NEWEPOCH, ts]` — the trusted leader announcing the epoch it wants to start.
139    NewEpoch { ts: u64 },
140    /// `[NACK, nts]` — "I will not start that one", sent back to the would-be leader, naming the
141    /// timestamp it is refusing. Algorithm 5.5 writes a bare `[NACK]`; the timestamp is taken from
142    /// Algorithm 5.8, and the reason is in the module documentation.
143    Nack { nts: u64 },
144}
145
146/// The wire, multiplexing the two children.
147#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
148pub enum Wire<B> {
149    /// The leader detector's heartbeats.
150    Detector(Heartbeat),
151    /// The broadcast's traffic, carrying [`EpochMsg`].
152    Epoch(B),
153}
154
155/// Requests from the layer above.
156///
157/// Uninhabited: epochs begin at initialisation and change when leadership does. There is nothing
158/// for the layer above to ask for.
159pub type Cmd = core::convert::Infallible;
160
161/// Indications to the layer above.
162#[derive(Debug, Clone, Copy, PartialEq, Eq)]
163pub enum Ind {
164    /// `⟨ ec, StartEpoch | newts, ℓ ⟩` — begin the epoch numbered `ts`, led by `leader`.
165    StartEpoch { ts: u64, leader: NodeId },
166}
167
168/// What the broadcast beneath puts on the wire for this layer's messages.
169///
170/// Written concretely rather than as a projection, for the reason
171/// [`crate::uniform_reliable_broadcast::BebMsg`] gives: the projection is only well-formed where
172/// `EpochMsg: Clone`, which would push that bound onto every use. The assertion below keeps the two
173/// from drifting apart.
174pub type BebMsg = pl::Wire<EpochMsg>;
175
176const _: () = {
177    /// Fails to compile if best-effort broadcast ever puts something else on the wire.
178    fn _beb_msg_is_what_we_say_it_is(
179        m: BebMsg,
180    ) -> <BestEffortBroadcast<EpochMsg> as Protocol>::Msg {
181        m
182    }
183};
184
185/// A sequence of epochs, driven by who is trusted.
186#[derive(Debug)]
187pub struct EpochChange {
188    me: NodeId,
189    /// Π. Its size is the book's `N`, and a process's position in it is `rank`.
190    peers: BTreeSet<NodeId>,
191    /// `trusted`.
192    trusted: NodeId,
193    /// `lastts` — the timestamp of the last epoch this process started.
194    lastts: u64,
195    /// Who started it. Not the book's; needed to tell whether a newly trusted leader is already the
196    /// one this process is following. See the departure on telling a leader where others have got to.
197    started_by: NodeId,
198    /// `ts` — this process's own next candidate, always in its own residue class mod `N`.
199    ts: u64,
200    omega: Child<EventualLeaderDetector>,
201    beb: Child<BestEffortBroadcast<EpochMsg>>,
202}
203
204impl EpochChange {
205    /// Epoch-change among `peers`, over a leader detector with the given heartbeat and timeout and
206    /// a best-effort broadcast whose links retransmit every `retransmit`.
207    pub fn new(me: NodeId, peers: impl IntoIterator<Item = NodeId>, timing: Timing) -> Self {
208        let Timing { retransmit, heartbeat, detect_after } = timing;
209        let mut peers: BTreeSet<NodeId> = peers.into_iter().collect();
210        peers.insert(me);
211        // `ℓ0` — the book's initial leader, fixed and known to all. `maxrank(Π)` is what Ω will
212        // trust first with nobody suspected, so starting there means the first `Trust` that agrees
213        // with it changes nothing.
214        let l0 = peers.iter().next_back().copied().expect("Π contains at least this process");
215        EpochChange {
216            me,
217            peers: peers.clone(),
218            trusted: l0,
219            lastts: 0,
220            started_by: l0,
221            ts: rank(&peers, me),
222            omega: Child::new(EventualLeaderDetector::new(
223                me,
224                peers.clone(),
225                heartbeat,
226                detect_after,
227            )),
228            beb: Child::new(BestEffortBroadcast::new(me, peers, retransmit)),
229        }
230    }
231
232    /// The epoch this process has most recently started, if any.
233    pub fn last_timestamp(&self) -> u64 {
234        self.lastts
235    }
236
237    /// Who this process currently trusts.
238    pub fn trusted(&self) -> NodeId {
239        self.trusted
240    }
241
242    /// `ts := ts + N; trigger ⟨ beb, Broadcast | [NEWEPOCH, ts] ⟩`.
243    fn announce(&mut self, cx: &mut ProtoCx<'_, Self>) {
244        self.ts += self.peers.len() as u64;
245        let ts = self.ts;
246        self.through_beb(cx, |b, ccx| {
247            b.on_cmd(beb::Cmd::Broadcast(EpochMsg::NewEpoch { ts }), ccx)
248        });
249    }
250
251    /// Announce the next candidate strictly above `floor`, staying in this process's residue class.
252    ///
253    /// The class is what keeps timestamps unique across processes — see the note on that above — so
254    /// jumping means rounding up to the next value congruent to `rank(self)` modulo `N`, not simply
255    /// taking `floor + 1`.
256    fn announce_above(&mut self, floor: u64, cx: &mut ProtoCx<'_, Self>) {
257        let n = self.peers.len() as u64;
258        let residue = self.ts % n;
259        let mut next = floor.max(self.ts) + 1;
260        next += (n + residue - next % n) % n;
261        debug_assert!(next > floor && next % n == residue);
262        self.ts = next;
263        let ts = self.ts;
264        self.through_beb(cx, |b, ccx| {
265            b.on_cmd(beb::Cmd::Broadcast(EpochMsg::NewEpoch { ts }), ccx)
266        });
267    }
268
269    /// Tell `leader` the timestamp this process has reached, so a leader that was never told it
270    /// became one can climb above it. Same message as a refusal, and for the same purpose.
271    fn report_to(&mut self, leader: NodeId, cx: &mut ProtoCx<'_, Self>) {
272        let nts = self.lastts;
273        // The same `NACK` goes on the wire as when an announcement is refused, so the trace cannot
274        // tell the two decisions apart. This is the half the trace cannot say.
275        cx.note(Note::ReachReported { leader, nts });
276        self.through_beb(cx, |b, ccx| {
277            b.on_cmd(beb::Cmd::SendTo { to: leader, msg: EpochMsg::Nack { nts } }, ccx)
278        });
279    }
280
281    /// `upon event ⟨ Ω, Trust | p ⟩`.
282    fn on_trust(&mut self, leader: NodeId, cx: &mut ProtoCx<'_, Self>) {
283        self.trusted = leader;
284        if leader == self.me {
285            self.announce(cx);
286        } else if self.started_by != leader {
287            // Departure: this process is following an epoch its new leader did not start, and that
288            // leader may never have been told anything changed. Tell it where we have reached.
289            self.report_to(leader, cx);
290        }
291    }
292
293    /// `upon event ⟨ beb, Deliver | ℓ, … ⟩` and `⟨ pl, Deliver | p, [NACK] ⟩`.
294    fn on_epoch_msg(&mut self, from: NodeId, msg: EpochMsg, cx: &mut ProtoCx<'_, Self>) {
295        match msg {
296            EpochMsg::NewEpoch { ts } => {
297                if from == self.trusted && ts > self.lastts {
298                    self.lastts = ts;
299                    self.started_by = from;
300                    cx.indicate(Ind::StartEpoch { ts, leader: from });
301                } else {
302                    // The would-be leader is not who this process trusts, or its timestamp is not
303                    // newer than one already started. Either way: tell it, so it can climb past —
304                    // and name the higher of the two, so it climbs past in one step rather than one
305                    // step per refusal.
306                    let why = if from != self.trusted {
307                        Refusal::NotTrusted { trusted: self.trusted }
308                    } else {
309                        Refusal::NotAhead { reached: self.lastts }
310                    };
311                    cx.note(Note::EpochRefused { from, ts, why });
312                    let nts = ts.max(self.lastts);
313                    self.through_beb(cx, |b, ccx| {
314                        b.on_cmd(beb::Cmd::SendTo { to: from, msg: EpochMsg::Nack { nts } }, ccx)
315                    });
316                }
317            }
318            // `upon event ⟨ sl, Deliver | p, [NACK, nts] ⟩ such that nts = ts` — Algorithm 5.8's
319            // guard, relaxed to `nts ≥ ts` because the report now carries how far the sender has
320            // reached rather than only what it refused. Boundedness is unchanged: the jump leaves
321            // the candidate strictly above `nts`, so a repeat names a timestamp already passed.
322            EpochMsg::Nack { nts } => {
323                if self.trusted == self.me && nts >= self.ts {
324                    self.announce_above(nts, cx);
325                } else {
326                    // **Nothing whatever reaches the trace from here.** This is the shape of
327                    // silence that cost the most to diagnose: a leader that everyone trusts and
328                    // which announces nothing looks, in a record of effects, exactly like a leader
329                    // nobody told anything.
330                    let why = if self.trusted != self.me {
331                        Refusal::NotLeader { trusted: self.trusted }
332                    } else {
333                        Refusal::NotAhead { reached: self.ts }
334                    };
335                    cx.note(Note::ReportIgnored { from, nts, why });
336                }
337            }
338        }
339    }
340
341    fn through_omega(
342        &mut self,
343        cx: &mut ProtoCx<'_, Self>,
344        f: impl FnOnce(&mut EventualLeaderDetector, &mut ProtoCx<'_, EventualLeaderDetector>),
345    ) {
346        let mut inds = self.omega.run(cx, Wire::Detector, f);
347        for eld::Ind::Trust { leader } in inds.drain(..) {
348            self.on_trust(leader, cx);
349        }
350        self.omega.reclaim(inds);
351    }
352
353    fn through_beb(
354        &mut self,
355        cx: &mut ProtoCx<'_, Self>,
356        f: impl FnOnce(
357            &mut BestEffortBroadcast<EpochMsg>,
358            &mut ProtoCx<'_, BestEffortBroadcast<EpochMsg>>,
359        ),
360    ) {
361        let mut inds = self.beb.run(cx, Wire::Epoch, f);
362        for ind in inds.drain(..) {
363            match ind {
364                beb::Ind::Deliver { from, msg } => self.on_epoch_msg(from, msg, cx),
365                // The broadcast is over a perfect link, which reports no scope boundary.
366                beb::Ind::SessionEnded { .. } | beb::Ind::SessionEstablished { .. } => {}
367            }
368        }
369        self.beb.reclaim(inds);
370    }
371}
372
373/// `rank(p)` — a process's position in `Π`, counting from one.
374///
375/// Any fixed injective map would do; this is the one that makes `ts := rank(self)` put each process
376/// in its own residue class modulo `N`.
377fn rank(peers: &BTreeSet<NodeId>, p: NodeId) -> u64 {
378    peers.iter().position(|q| *q == p).expect("p ∈ Π") as u64 + 1
379}
380
381impl Protocol for EpochChange {
382    type Cmd = Cmd;
383    type Ind = Ind;
384    type Msg = Wire<BebMsg>;
385    type Scope = core::convert::Infallible;
386    type Note = crate::Note;
387    /// Keeps nothing durably. `logged_epoch_change` is the variant that does.
388    type Meta = core::convert::Infallible;
389    type Entry = core::convert::Infallible;
390
391    fn on_cmd(&mut self, cmd: Cmd, _: &mut ProtoCx<'_, Self>) {
392        match cmd {}
393    }
394
395    fn on_msg(&mut self, from: NodeId, msg: Self::Msg, cx: &mut ProtoCx<'_, Self>) {
396        match msg {
397            Wire::Detector(h) => self.through_omega(cx, |o, ccx| o.on_msg(from, h, ccx)),
398            Wire::Epoch(m) => self.through_beb(cx, |b, ccx| b.on_msg(from, m, ccx)),
399        }
400    }
401
402    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
403        self.through_omega(cx, |o, ccx| o.on_timer(id, ccx));
404        self.through_beb(cx, |b, ccx| b.on_timer(id, ccx));
405    }
406
407    /// `upon event ⟨ ec, Init ⟩` — the state is set in `new`; this starts the detector, whose first
408    /// `Trust` may immediately make this process announce an epoch.
409    fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
410        self.through_omega(cx, |o, ccx| o.on_init(ccx));
411    }
412}