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}