Skip to main content

recon_protocols/
logged_epoch_change.rs

1//! Epoch-change that survives a restart.
2//!
3//! **Status: implementation. Space: bounded by membership, plus what the stubborn children hold
4//! outstanding — which nothing here retires, so see the departure on `Stop` below.**
5//!
6//! Cachin, Guerraoui & Rodrigues, Module 5.6 (`LoggedEpochChange`) and Algorithm 5.8 ("Logged
7//! Leader-Based Epoch-Change"), quoted from the book:
8//!
9//! ```text
10//! Algorithm 5.8: Logged Leader-Based Epoch-Change
11//! Implements: LoggedEpochChange, instance lec.
12//! Uses:
13//!     StubbornPointToPointLinks, instance sl;
14//!     StubbornBestEffortBroadcast, instance sbeb;
15//!     EventualLeaderDetector, instance Ω.
16//!
17//! upon event ⟨ lec, Init ⟩ do
18//!     trusted := ℓ0;
19//!     (startts, start) := (0, ℓ0);
20//!     ts := rank(self) − N;
21//!
22//! upon event ⟨ lec, Recovery ⟩ do
23//!     retrieve(startts, start);
24//!
25//! upon event ⟨ Ω, Trust | p ⟩ do
26//!     trusted := p;
27//!     if p = self then
28//!         ts := ts + N;
29//!         trigger ⟨ sbeb, Broadcast | [NEWEPOCH, ts] ⟩;
30//!
31//! upon event ⟨ sbeb, Deliver | ℓ, [NEWEPOCH, newts] ⟩ do
32//!     if ℓ = trusted ∧ newts > startts then
33//!         (startts, start) := (newts, ℓ);
34//!         store(startts, start);
35//!         trigger ⟨ lec, StartEpoch | startts, start ⟩;
36//!     else
37//!         trigger ⟨ sl, Send | ℓ, [NACK, newts] ⟩;
38//!
39//! upon event ⟨ sl, Deliver | p, [NACK, nts] ⟩ such that nts = ts do
40//!     if trusted = self then
41//!         ts := ts + N;
42//!         trigger ⟨ sbeb, Broadcast | [NEWEPOCH, ts] ⟩;
43//! ```
44//!
45//! # What is durable, and what is not
46//!
47//! `(startts, start)` — the epoch this process has actually entered, and who leads it. Written
48//! before `StartEpoch` is raised, in the handler's own text, because that indication is what the
49//! consensus above acts on: a process that told its consensus to enter epoch 20 and then came back
50//! believing it had entered nothing would read an empty state where an accepted value should be.
51//!
52//! `ts` — this process's own next candidate — is **not** durable, and the book does not store it.
53//! A recovered leader therefore starts climbing again from `rank(self)`, and has to walk back up in
54//! steps of `N` before it can announce a timestamp anybody will accept. That is slow but it is not
55//! wrong, and the reason it is not wrong is that `startts` *is* durable: every process refuses a
56//! timestamp no greater than the epoch it has already entered, so a reused candidate is refused
57//! rather than confused with the epoch that first used it. The NACK carries the timestamp it
58//! refuses, so each refusal moves the leader up exactly once.
59//!
60//! Reusing a candidate is safe for a second reason too: `ts ≡ rank(self) (mod N)` holds across
61//! incarnations, because `rank` is a function of the membership rather than of anything this
62//! process remembers. Two processes still cannot mint the same timestamp, so `EC2` — one timestamp
63//! names one leader — survives a restart even though `ts` does not.
64//!
65//! # Identity, and how durable it has to be
66//!
67//! `CLAUDE.md`: an identifier that crosses the wire or lands in storage outlives the handler that
68//! minted it. The [`sl::SendId`] and [`sbeb::BroadcastId`] counters here mint identifiers that do
69//! neither — they name entries in the stubborn children's own volatile tables, which a crash
70//! empties. A restarted process therefore restarts its counters at zero and names nothing that is
71//! still live, because nothing is. Their scope is the incarnation, and that is the whole of it.
72//!
73//! The timestamp is the identifier that does cross the wire, and it is the one that is durable in
74//! the sense that matters: not stored, but re-derived from `rank`, which does not change.
75//!
76//! # Departure: a repeat of the epoch already entered is not refused
77//!
78//! Algorithm 5.8 answers every NEWEPOCH it does not act on with a NACK. Over the stubborn broadcast
79//! the same algorithm's `Uses:` line names, that does not terminate, and the loop is tighter than
80//! the one [`crate::epoch_change`] describes: the leader announces `t`, every process enters it and
81//! writes it down, and then the broadcast — which retransmits until retired, and nothing here
82//! retires it — delivers `t` again. The second delivery fails `newts > startts`, because `startts`
83//! is now `t`. So every process refuses the announcement it has just accepted, the leader climbs to
84//! `t + N`, and the cycle restarts one retransmission interval later, for ever. Measured: epoch 380
85//! and still climbing after eight timeouts, with leadership settled and nothing faulty.
86//!
87//! [`crate::epoch_change`] does not have this, and the reason is the child rather than the
88//! algorithm: a best-effort broadcast over *perfect* links delivers each announcement exactly once,
89//! so the repeat never reaches the handler. Moving to a stubborn broadcast — which must not
90//! deduplicate, because repeating for ever is what reaches a recovered process — brings it back.
91//!
92//! The guard is that a repeat is not a refusal. `newts = startts` from the leader of the epoch
93//! already entered is silence: there is nothing for the leader to climb past, because its
94//! announcement was accepted. A NACK is still sent when the sender is not trusted, and when the
95//! timestamp is genuinely stale — those are the two cases the book's `else` is for.
96//!
97//! # Departure: each distinct announcement is refused once per peer
98//!
99//! The same shape one step earlier, for the announcements that *are* stale. A NEWEPOCH below
100//! `startts` from a process that is not the leader of the epoch entered is refused; the stubborn
101//! broadcast delivers it again next interval; Algorithm 5.8 refuses it again, on a fresh stubborn
102//! transmission, and so on for ever. `nacked` remembers the highest timestamp refused per peer and
103//! refuses nothing at or below it. Bounded by membership. The one NACK sent is itself stubborn, so
104//! it reaches the leader; a second carries no information the first did not.
105//!
106//! # Departure: nothing calls `Stop`
107//!
108//! [`crate::stubborn_broadcast`] and [`crate::stubborn_link`] retransmit until retired, and this
109//! module retires nothing, so its space grows with the number of *distinct* announcements and
110//! refusals rather than with the membership — though not, after the two guards above, with time. That is the same unbounded transcription
111//! [`crate::logged_uniform_reliable_broadcast`] has and for the same reason: retransmitting for
112//! ever is what reaches a process that was down when the message was sent, and a recovered process
113//! has no way to ask for what it missed.
114//!
115//! It is bounded in practice by the thing that bounds the announcements themselves — leadership
116//! settling — and the NACK's timestamp guard is what makes that a finite number rather than a
117//! feedback loop. See [`crate::epoch_change`], whose module documentation records what the
118//! unguarded form cost.
119//!
120//! ```text
121//! EC1 [always]     Monotonicity — the timestamps a process starts strictly increase, across
122//!                  restarts as well as within one incarnation, and one timestamp names one leader
123//! EC2 [conditional] Consistency — every correct process eventually starts the same last epoch,
124//!                  provided the leader detector settles
125//! ```
126
127use recon_core::{Child, NodeId, ProtoCx, Protocol, TimerId};
128use serde::{Deserialize, Serialize};
129use std::collections::{BTreeMap, BTreeSet};
130
131use crate::Timing;
132
133use crate::eventual_leader_detector::{self as eld, EventualLeaderDetector};
134use crate::perfect_failure_detector::Heartbeat;
135use crate::stubborn_broadcast::{self as sbeb, BroadcastId, StubbornBroadcast};
136use crate::stubborn_link::{self as sl, SendId, StubbornLink};
137
138/// `[NEWEPOCH, ts]` — the trusted leader announcing the epoch it wants to start.
139#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
140pub struct NewEpoch {
141    pub ts: u64,
142}
143
144/// `[NACK, nts]` — "I will not start that one", naming the timestamp refused.
145#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
146pub struct Nack {
147    pub nts: u64,
148}
149
150/// The wire, multiplexing the three children the book names.
151#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
152pub enum Wire {
153    /// The leader detector's heartbeats.
154    Detector(Heartbeat),
155    /// `sbeb` — the announcements.
156    Announce(NewEpoch),
157    /// `sl` — the refusals, which go to one process rather than all.
158    Refuse(Nack),
159}
160
161/// Requests from the layer above.
162///
163/// Uninhabited, as in [`crate::epoch_change`]: epochs begin at initialisation and change when
164/// leadership does.
165pub type Cmd = core::convert::Infallible;
166
167/// Indications to the layer above.
168#[derive(Debug, Clone, Copy, PartialEq, Eq)]
169pub enum Ind {
170    /// `⟨ lec, StartEpoch | startts, start ⟩`. Raised only after the pair is durable.
171    StartEpoch { ts: u64, leader: NodeId },
172}
173
174/// `(startts, start)` — the one metadata value this layer rewrites.
175#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
176pub struct Started {
177    pub ts: u64,
178    pub leader: NodeId,
179}
180
181/// A sequence of epochs whose current position survives a restart.
182#[derive(Debug)]
183pub struct LoggedEpochChange {
184    me: NodeId,
185    /// Π. Its size is the book's `N`, and a process's position in it is `rank`.
186    peers: BTreeSet<NodeId>,
187    /// `trusted`.
188    trusted: NodeId,
189    /// `(startts, start)` — durable.
190    started: Started,
191    /// `ts` — volatile, and re-derived from `rank` on a restart. See the module documentation.
192    ts: u64,
193    /// Names the next stubborn transmission. Volatile, and so is what it keys.
194    next_send: u64,
195    /// Names the next stubborn broadcast. Volatile, and so is what it keys.
196    next_broadcast: u64,
197    /// The highest timestamp refused, per peer. Bounded by membership; see the departure.
198    nacked: BTreeMap<NodeId, u64>,
199    omega: Child<EventualLeaderDetector>,
200    sbeb: Child<StubbornBroadcast<NewEpoch>>,
201    sl: Child<StubbornLink<Nack>>,
202}
203
204impl LoggedEpochChange {
205    /// Epoch-change among `peers`, over a leader detector with the given heartbeat and timeout and
206    /// stubborn children retransmitting 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        let l0 = peers.iter().next_back().copied().expect("Π contains at least this process");
212        // `ts := rank(self) − N`, so that the `ts := ts + N` in the first `Trust` makes the first
213        // announcement `rank(self)` rather than `rank(self) + N`. Held as a signed value only here;
214        // `rank` counts from one and `N ≥ 1`, so the first increment lands on `rank`.
215        let ts = rank(&peers, me);
216        let n = peers.len() as u64;
217        LoggedEpochChange {
218            me,
219            peers: peers.clone(),
220            trusted: l0,
221            started: Started { ts: 0, leader: l0 },
222            ts: ts.wrapping_sub(n),
223            next_send: 0,
224            next_broadcast: 0,
225            nacked: BTreeMap::new(),
226            omega: Child::new(EventualLeaderDetector::new(
227                me,
228                peers.clone(),
229                heartbeat,
230                detect_after,
231            )),
232            sbeb: Child::new(StubbornBroadcast::new(me, peers.clone(), retransmit)),
233            sl: Child::new(StubbornLink::new(retransmit)),
234        }
235    }
236
237    /// `startts` — the epoch this process has entered, as its durable record has it.
238    pub fn last_timestamp(&self) -> u64 {
239        self.started.ts
240    }
241
242    /// `start` — who leads the epoch this process has entered.
243    pub fn last_leader(&self) -> NodeId {
244        self.started.leader
245    }
246
247    /// Who this process currently trusts.
248    pub fn trusted(&self) -> NodeId {
249        self.trusted
250    }
251
252    /// `ts` — this process's own next candidate. Volatile; see the module documentation.
253    pub fn candidate(&self) -> u64 {
254        self.ts
255    }
256
257    /// `ts := ts + N; trigger ⟨ sbeb, Broadcast | [NEWEPOCH, ts] ⟩`.
258    fn announce(&mut self, cx: &mut ProtoCx<'_, Self>) {
259        self.ts = self.ts.wrapping_add(self.peers.len() as u64);
260        let msg = NewEpoch { ts: self.ts };
261        let id = BroadcastId(self.next_broadcast);
262        self.next_broadcast += 1;
263        self.through_sbeb(cx, |b, ccx| b.on_cmd(sbeb::Cmd::Broadcast { id, msg }, ccx));
264    }
265
266    /// `upon event ⟨ Ω, Trust | p ⟩`.
267    fn on_trust(&mut self, leader: NodeId, cx: &mut ProtoCx<'_, Self>) {
268        self.trusted = leader;
269        if leader == self.me {
270            self.announce(cx);
271        }
272    }
273
274    /// `upon event ⟨ sbeb, Deliver | ℓ, [NEWEPOCH, newts] ⟩`.
275    ///
276    /// **The order of the two statements in the `then` branch is the obligation.** `store` comes
277    /// before `trigger`, here in the handler's own text, because `StartEpoch` is what makes the
278    /// epoch visible to the consensus above.
279    fn on_new_epoch(&mut self, from: NodeId, newts: u64, cx: &mut ProtoCx<'_, Self>) {
280        if from == self.trusted && newts > self.started.ts {
281            self.started = Started { ts: newts, leader: from };
282            cx.storage().set(self.started);
283            cx.indicate(Ind::StartEpoch { ts: newts, leader: from });
284        } else if from == self.started.leader && newts == self.started.ts {
285            // A repeat of the announcement this process has already accepted. Silence, not a
286            // refusal — see the departure in the module documentation.
287        } else {
288            // Once per distinct announcement per peer — see the departure. A peer's candidates
289            // strictly increase, so anything at or below the last one refused is a repeat.
290            if self.nacked.get(&from).is_some_and(|last| newts <= *last) {
291                return;
292            }
293            self.nacked.insert(from, newts);
294            let id = SendId(self.next_send);
295            self.next_send += 1;
296            self.through_sl(cx, |l, ccx| {
297                l.on_cmd(sl::Cmd::Send { id, to: from, msg: Nack { nts: newts } }, ccx)
298            });
299        }
300    }
301
302    /// `upon event ⟨ sl, Deliver | p, [NACK, nts] ⟩ such that nts = ts`.
303    ///
304    /// The stubborn link beneath repeats a refusal until it is retired, and nothing retires one, so
305    /// this handler sees the same NACK many times. `nts = ts` makes that idempotent: the first one
306    /// moves `ts`, and every repeat then names a candidate already superseded.
307    fn on_nack(&mut self, nts: u64, cx: &mut ProtoCx<'_, Self>) {
308        if nts == self.ts && self.trusted == self.me {
309            self.announce(cx);
310        }
311    }
312
313    fn through_omega(
314        &mut self,
315        cx: &mut ProtoCx<'_, Self>,
316        f: impl FnOnce(&mut EventualLeaderDetector, &mut ProtoCx<'_, EventualLeaderDetector>),
317    ) {
318        let mut inds = self.omega.run(cx, Wire::Detector, f);
319        for eld::Ind::Trust { leader } in inds.drain(..) {
320            self.on_trust(leader, cx);
321        }
322        self.omega.reclaim(inds);
323    }
324
325    fn through_sbeb(
326        &mut self,
327        cx: &mut ProtoCx<'_, Self>,
328        f: impl FnOnce(&mut StubbornBroadcast<NewEpoch>, &mut ProtoCx<'_, StubbornBroadcast<NewEpoch>>),
329    ) {
330        let mut inds = self.sbeb.run(cx, Wire::Announce, f);
331        for sbeb::Ind::Deliver { from, msg } in inds.drain(..) {
332            self.on_new_epoch(from, msg.ts, cx);
333        }
334        self.sbeb.reclaim(inds);
335    }
336
337    fn through_sl(
338        &mut self,
339        cx: &mut ProtoCx<'_, Self>,
340        f: impl FnOnce(&mut StubbornLink<Nack>, &mut ProtoCx<'_, StubbornLink<Nack>>),
341    ) {
342        let mut inds = self.sl.run(cx, Wire::Refuse, f);
343        for sl::Ind::Deliver { msg, .. } in inds.drain(..) {
344            self.on_nack(msg.nts, cx);
345        }
346        self.sl.reclaim(inds);
347    }
348}
349
350/// `rank(p)` — a process's position in `Π`, counting from one.
351///
352/// A function of the membership alone, which is why it survives a restart without being stored.
353fn rank(peers: &BTreeSet<NodeId>, p: NodeId) -> u64 {
354    peers.iter().position(|q| *q == p).expect("p ∈ Π") as u64 + 1
355}
356
357impl Protocol for LoggedEpochChange {
358    type Cmd = Cmd;
359    type Ind = Ind;
360    type Msg = Wire;
361    type Scope = core::convert::Infallible;
362    type Note = crate::Note;
363    type Meta = Started;
364    /// Nothing accumulates: the epoch entered is one value, rewritten.
365    type Entry = core::convert::Infallible;
366
367    fn on_cmd(&mut self, cmd: Cmd, _: &mut ProtoCx<'_, Self>) {
368        match cmd {}
369    }
370
371    fn on_msg(&mut self, from: NodeId, msg: Wire, cx: &mut ProtoCx<'_, Self>) {
372        match msg {
373            Wire::Detector(h) => self.through_omega(cx, |o, ccx| o.on_msg(from, h, ccx)),
374            Wire::Announce(m) => self.through_sbeb(cx, |b, ccx| b.on_msg(from, m, ccx)),
375            Wire::Refuse(m) => self.through_sl(cx, |l, ccx| l.on_msg(from, m, ccx)),
376        }
377    }
378
379    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
380        self.through_omega(cx, |o, ccx| o.on_timer(id, ccx));
381        self.through_sbeb(cx, |b, ccx| b.on_timer(id, ccx));
382        self.through_sl(cx, |l, ccx| l.on_timer(id, ccx));
383    }
384
385    /// `upon event ⟨ lec, Init ⟩` — the state is set in `new`; this starts the detector.
386    fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
387        self.through_omega(cx, |o, ccx| o.on_init(ccx));
388    }
389
390    /// `upon event ⟨ lec, Recovery ⟩ do retrieve(startts, start)`.
391    ///
392    /// No `StartEpoch` is raised. The epoch is not new — this process entered it before it went
393    /// down, and told the layer above so at the time. Re-raising it would announce as fresh an
394    /// epoch whose consensus instance already exists, which is the layer above's business to
395    /// reconstruct from its own record and not this layer's to invent.
396    fn on_recovery(&mut self, cx: &mut ProtoCx<'_, Self>) {
397        if let Some(started) = cx.storage().get().copied() {
398            self.started = started;
399        }
400        self.through_omega(cx, |o, ccx| o.on_init(ccx));
401    }
402}