recon_protocols/logged_leader_driven_consensus.rs
1//! Paxos that survives a restart.
2//!
3//! **Status: implementation. Space: bounded by membership, plus what the stubborn children hold
4//! outstanding — inherited from [`crate::logged_epoch_change`] and
5//! [`crate::logged_epoch_consensus`], and retired for the consensus by each epoch's `Abort`.**
6//!
7//! Cachin, Guerraoui & Rodrigues, Module 5.5 (`LoggedUniformConsensus`) and Algorithms 5.10–5.11
8//! ("Logged Leader-Driven Consensus"), quoted from the book:
9//!
10//! ```text
11//! Algorithm 5.10: Logged Leader-Driven Consensus (part 1)
12//! Implements: LoggedUniformConsensus, instance luc.
13//! Uses:
14//! LoggedEpochChange, instance lec;
15//! LoggedEpochConsensus (multiple instances).
16//!
17//! upon event ⟨ luc, Init ⟩ do
18//! val := ⊥; decision := ⊥; aborted := FALSE; proposed := FALSE;
19//! Obtain the initial leader ℓ0 from the logged epoch-change instance lec;
20//! Initialize a new instance lep.0 of logged epoch consensus with timestamp 0, leader ℓ0,
21//! and state (0, ⊥);
22//! (ets, ℓ) := (0, ℓ0);
23//! store(ets, ℓ, decision);
24//!
25//! upon event ⟨ luc, Recovery ⟩ do
26//! retrieve(ets, ℓ, decision);
27//! retrieve(startts, start) of instance lec;
28//! (newts, newℓ) := (startts, start);
29//! retrieve(epochdecision) of instance lep.ets;
30//! if epochdecision ≠ ⊥ ∧ decision = ⊥ then
31//! decision := epochdecision;
32//! store(decision);
33//! trigger ⟨ luc, Decide | decision ⟩;
34//! aborted := FALSE;
35//!
36//! upon event ⟨ luc, Propose | v ⟩ do
37//! val := v;
38//!
39//! Algorithm 5.11: Logged Leader-Driven Consensus (part 2)
40//!
41//! upon event ⟨ lec, StartEpoch | startts, start ⟩ do
42//! retrieve(startts, start) of instance lec;
43//! (newts, newℓ) := (startts, start);
44//!
45//! upon (ets, ℓ) ≠ (newts, newℓ) ∧ aborted = FALSE do
46//! aborted := TRUE;
47//! trigger ⟨ lep.ets, Abort ⟩;
48//!
49//! upon event ⟨ lep.ts, Aborted | state ⟩ such that ts = ets do
50//! (ets, ℓ) := (newts, newℓ);
51//! store(ets, ℓ);
52//! aborted := FALSE;
53//! proposed := FALSE;
54//! Initialize a new instance lep.ets of logged epoch consensus with timestamp ets, leader ℓ,
55//! and state state;
56//!
57//! upon ℓ = self ∧ val ≠ ⊥ ∧ proposed = FALSE do
58//! proposed := TRUE;
59//! trigger ⟨ lep.ets, Propose | val ⟩;
60//!
61//! upon event ⟨ lep.ts, Decide | epochdecision ⟩ such that ts = ets do
62//! retrieve(epochdecision) of instance lep.ets;
63//! if decision = ⊥ then
64//! decision := epochdecision;
65//! store(decision);
66//! trigger ⟨ luc, Decide | decision ⟩;
67//! ```
68//!
69//! # Two durable children under one durable parent
70//!
71//! This is the first protocol here that keeps a record of its own *and* composes children that keep
72//! records of theirs, and it is the reason [`recon_core::Slot`] exists. Every other composition
73//! hands its children a `NoStore`, because a parent and a child sharing one store would each
74//! overwrite the other's metadata — and would do it silently, with nothing failing until a recovery
75//! read back half of what it wrote.
76//!
77//! A slot names the part of [`Durable`] that belongs to a child. The child's `set` becomes a
78//! read-modify-write of this record: **one write, not two**, so a crash cannot land between the
79//! parent's record and its child's.
80//!
81//! The book writes `retrieve(startts, start) of instance lec` and `retrieve(epochdecision) of
82//! instance lep.ets` — a parent reading its children's records by name. Here it does that by
83//! handing each child its slot and letting the child's own `Recovery` read it, which is the same
84//! thing said in the direction the composition already runs.
85//!
86//! # One slot for a child that is replaced every epoch
87//!
88//! The book has one logged epoch consensus instance per timestamp, each with its own record. There
89//! is one slot here, holding whichever instance is live. That is not a loss: the only instance ever
90//! read is `lep.ets`, and `ets` is in this record too, so a slot holding the current instance is
91//! exactly what `retrieve(...) of instance lep.ets` asks for.
92//!
93//! A crash can land between `store(ets, ℓ)` and the new instance's own `Init` write, and both
94//! outcomes are safe. Land before it, and recovery reads the *previous* epoch's `epochdecision`
95//! against the *new* `ets` — but a value that epoch decided is, by lock-in, the value every later
96//! epoch decides, so deciding it is right. Land after it, and recovery reads a fresh record against
97//! the old `ets` and has simply not decided yet.
98//!
99//! # Reading two lines the page does not quite give
100//!
101//! `upon (ets, ℓ) = (newts, newℓ) ∧ aborted = FALSE do … trigger ⟨ lep.ets, Abort ⟩` is printed
102//! with `=`, and must be `≠`: aborting the epoch you are in because it *is* the one you want would
103//! abort every epoch immediately and decide nothing. The same OCR class as `if v ≠ ⊥ then tmpval :=
104//! v` in Algorithm 5.6, which [`crate::epoch_consensus`] records for the same reason.
105//!
106//! `Init` does not print an assignment to `(newts, newℓ)`. It must be `(0, ℓ0)` — the same pair as
107//! `(ets, ℓ)` — or the standing condition above is true from the first event and the initial epoch
108//! is aborted before it does anything. Algorithm 5.7 sets `(newts, newℓ) := (0, ⊥)`, which works
109//! there because its condition is written on the `Aborted` handler rather than as a standing one.
110//!
111//! # The decision is announced again after a recovery, and Module 5.5 has no integrity property
112//!
113//! `⟨ luc, Decide | decision ⟩` is specified as "notifies the upper layer that variable `decision`
114//! in stable storage contains the decided value of consensus" — a pointer to a record, not a
115//! one-shot event. Module 5.5 lists **three** properties where the fail-noisy Module 5.2 lists
116//! four: termination, validity and uniform agreement, with *no* integrity clause. The book dropped
117//! it, and the reason is exactly this: a logged indication may be raised again, and the layer above
118//! reads storage and must be idempotent. `logged_link` and `logged_uniform_reliable_broadcast`
119//! already work that way, and `README.md` states it as the rule for this whole model.
120//!
121//! So a process that decided, crashed, and came back announces its decision once more. Algorithm
122//! 5.10's `Recovery` handler as printed does not — it announces only when `epochdecision ≠ ⊥ ∧
123//! decision = ⊥`, the case where the child's record survived and this layer's did not. Re-announcing
124//! the other case is the departure, and it is the module's own indication wording taken at face
125//! value: a layer above that crashed with this one never saw the first indication, and there is no
126//! other event that would tell it.
127//!
128//! ```text
129//! LUC1 [conditional] Termination — every correct process that never crashes eventually
130//! log-decides, provided a majority is correct and the leader detector settles.
131//! "Correct" here means eventually up and staying up, so a process that keeps
132//! crashing is not owed a decision
133//! LUC2 [always] Validity — a log-decided value was proposed by some process
134//! LUC3 [always] Uniform agreement — no two processes log-decide differently, **including
135//! across crashes and recoveries, and while the leader detector is wrong**
136//! ```
137//!
138//! There is deliberately no integrity clause, for the reason above. What replaces it is that the
139//! *value* never changes: a process announces the same decision every time, which
140//! [`LoggedLeaderDrivenConsensus::decision`] is the durable statement of.
141
142use recon_core::{Child, NodeId, ProtoCx, Protocol, Slot, TimerId, slot};
143use serde::{Deserialize, Serialize};
144use std::collections::BTreeSet;
145
146use crate::Timing;
147use crate::logged_epoch_change::{self as lec, LoggedEpochChange};
148use crate::logged_epoch_consensus::{self as lep, LoggedEpochConsensus, State};
149
150/// The wire, multiplexing the epoch-change child and whichever epoch-consensus instance is live.
151#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
152pub enum Wire<V> {
153 /// The logged epoch-change child's traffic.
154 Change(lec::Wire),
155 /// The live logged epoch-consensus instance's traffic.
156 ///
157 /// The epoch tag that makes `lep.ts` addressable lives inside the child — see
158 /// [`lep::Tagged`] — because the epoch is the instance's own identity.
159 Consensus(lep::Wire<V>),
160}
161
162/// Requests from the layer above.
163#[derive(Debug, Clone, PartialEq, Eq)]
164pub enum Cmd<V> {
165 /// `⟨ luc, Propose | v ⟩`.
166 Propose(V),
167}
168
169/// Indications to the layer above.
170#[derive(Debug, Clone, PartialEq, Eq)]
171pub enum Ind<V> {
172 /// `⟨ luc, Decide | v ⟩`. Raised at most once, and only after the decision is durable.
173 Decide(V),
174}
175
176/// Everything this stack keeps durably, as one rewritten value.
177///
178/// The two `Option` fields are the children's slots. They are `Option` because a child may not have
179/// written yet, and because [`Slot::write`] has to be able to build this record from nothing.
180#[derive(Debug, Clone, PartialEq, Eq)]
181pub struct Durable<V> {
182 /// `ets`.
183 pub ets: u64,
184 /// `ℓ`. `None` before this process's own `Init` has written, which is when it is `ℓ0`.
185 pub leader: Option<NodeId>,
186 /// `decision`.
187 pub decision: Option<V>,
188 /// `lec`'s record: `(startts, start)`.
189 pub lec: Option<lec::Started>,
190 /// `lep.ets`'s record: `(valts, val)` and `epochdecision`.
191 pub lep: Option<lep::Durable<V>>,
192}
193
194impl<V> Default for Durable<V> {
195 fn default() -> Self {
196 Durable { ets: 0, leader: None, decision: None, lec: None, lep: None }
197 }
198}
199
200/// Paxos in the fail-recovery model.
201pub struct LoggedLeaderDrivenConsensus<V: Clone> {
202 me: NodeId,
203 peers: BTreeSet<NodeId>,
204 timing: Timing,
205 /// `ℓ0` — the initial leader, re-derived from the membership rather than stored.
206 l0: NodeId,
207 /// `val`.
208 val: Option<V>,
209 /// `proposed`.
210 proposed: bool,
211 /// `decision` — durable.
212 decision: Option<V>,
213 /// `aborted` — whether an abort is outstanding.
214 aborted: bool,
215 /// `(ets, ℓ)` — durable.
216 ets: u64,
217 leader: NodeId,
218 /// `(newts, newℓ)` — the epoch the change child has told this process to enter.
219 newts: u64,
220 newleader: NodeId,
221 lec: Child<LoggedEpochChange>,
222 /// `lep.ets`. One instance, replaced on each epoch change, sharing one slot.
223 lep: Child<LoggedEpochConsensus<V>>,
224}
225
226/// Where `lec`'s record sits inside this one.
227fn lec_slot<V: Clone>() -> Slot<Durable<V>, lec::Started> {
228 slot!(Durable<V>, lec)
229}
230
231/// Where the live `lep` instance's record sits inside this one.
232fn lep_slot<V: Clone>() -> Slot<Durable<V>, lep::Durable<V>> {
233 slot!(Durable<V>, lep)
234}
235
236impl<V: Clone + PartialEq> LoggedLeaderDrivenConsensus<V> {
237 /// Paxos among `peers`, over stable storage.
238 pub fn new(me: NodeId, peers: impl IntoIterator<Item = NodeId>, timing: Timing) -> Self {
239 let mut peers: BTreeSet<NodeId> = peers.into_iter().collect();
240 peers.insert(me);
241 // "Obtain the initial leader ℓ0 from the logged epoch-change instance lec" — `maxrank(Π)`,
242 // which is what Ω trusts first with nobody suspected, so the two agree before anything
243 // moves. A function of the membership, so it survives a restart without being stored.
244 let l0 = peers.iter().next_back().copied().expect("Π contains at least this process");
245 LoggedLeaderDrivenConsensus {
246 me,
247 peers: peers.clone(),
248 timing,
249 l0,
250 val: None,
251 proposed: false,
252 decision: None,
253 aborted: false,
254 ets: 0,
255 leader: l0,
256 // `(newts, newℓ) := (0, ℓ0)`. See the module's note on the two lines the page does not
257 // quite give: `(0, ⊥)` would abort the initial epoch before it did anything.
258 newts: 0,
259 newleader: l0,
260 lec: Child::new(LoggedEpochChange::new(me, peers.clone(), timing)),
261 lep: Child::new(LoggedEpochConsensus::new(
262 me,
263 peers,
264 0,
265 l0,
266 State::default(),
267 timing.retransmit,
268 )),
269 }
270 }
271
272 /// The epoch now live at this process.
273 pub fn epoch(&self) -> u64 {
274 self.ets
275 }
276
277 /// Who leads the epoch now live.
278 pub fn leader(&self) -> NodeId {
279 self.leader
280 }
281
282 /// What this process decided, if it has.
283 pub fn decision(&self) -> Option<&V> {
284 self.decision.as_ref()
285 }
286
287 /// `(valts, val)` as the epoch now live holds it — what the abort handshake carries forward.
288 pub fn state(&self) -> &State<V> {
289 self.lep.state()
290 }
291
292 /// This process's own part of the durable record, with the children's slots carried across
293 /// unchanged.
294 ///
295 /// The parent writes the whole record, so it has to preserve what the children put in it —
296 /// the mirror of what [`Slot`] does for a child's write.
297 fn record(&self, cx: &mut ProtoCx<'_, Self>) -> Durable<V> {
298 let held = cx.storage().get().cloned();
299 Durable {
300 ets: self.ets,
301 leader: Some(self.leader),
302 decision: self.decision.clone(),
303 lec: held.as_ref().and_then(|d| d.lec),
304 lep: held.and_then(|d| d.lep),
305 }
306 }
307
308 fn store(&mut self, cx: &mut ProtoCx<'_, Self>) {
309 let record = self.record(cx);
310 cx.storage().set(record);
311 }
312
313 /// `upon ℓ = self ∧ val ≠ ⊥ ∧ proposed = FALSE` — a standing condition, re-evaluated whenever
314 /// any of the three could have changed.
315 fn maybe_propose(&mut self, cx: &mut ProtoCx<'_, Self>) {
316 if self.leader == self.me
317 && !self.proposed
318 && let Some(v) = self.val.clone()
319 {
320 self.proposed = true;
321 self.through_lep(cx, |e, ccx| e.on_cmd(lep::Cmd::Propose(v), ccx));
322 }
323 }
324
325 /// `upon (ets, ℓ) ≠ (newts, newℓ) ∧ aborted = FALSE do aborted := TRUE; trigger ⟨ Abort ⟩`.
326 fn maybe_abort(&mut self, cx: &mut ProtoCx<'_, Self>) {
327 if (self.ets, self.leader) != (self.newts, self.newleader) && !self.aborted {
328 self.aborted = true;
329 self.through_lep(cx, |e, ccx| e.on_cmd(lep::Cmd::Abort, ccx));
330 }
331 }
332
333 /// `upon event ⟨ lep.ts, Aborted | state ⟩ such that ts = ets`.
334 ///
335 /// **`store(ets, ℓ)` precedes the new instance**, and the new instance's own `Init` write
336 /// follows it. Neither order is safe on its own; what makes both safe is that the pair is
337 /// read back together — see the module's note on one slot for a replaced child.
338 fn on_aborted(&mut self, state: State<V>, cx: &mut ProtoCx<'_, Self>) {
339 if !self.aborted {
340 // The book's `such that ts = ets`: an answer from an instance already superseded.
341 return;
342 }
343 self.ets = self.newts;
344 self.leader = self.newleader;
345 self.aborted = false;
346 self.proposed = false;
347 self.store(cx);
348 self.lep.replace(LoggedEpochConsensus::new(
349 self.me,
350 self.peers.clone(),
351 self.ets,
352 self.leader,
353 state,
354 self.timing.retransmit,
355 ));
356 self.through_lep(cx, |e, ccx| e.on_init(ccx));
357 self.maybe_propose(cx);
358 }
359
360 /// `upon event ⟨ lep.ts, Decide | epochdecision ⟩ such that ts = ets`.
361 ///
362 /// The child has already made `epochdecision` durable before raising this — that is
363 /// [`crate::logged_epoch_consensus`]'s own obligation — and this handler makes `decision`
364 /// durable before reporting it, which is this one's.
365 fn on_epoch_decision(&mut self, v: V, cx: &mut ProtoCx<'_, Self>) {
366 if self.decision.is_some() {
367 return;
368 }
369 self.decision = Some(v.clone());
370 self.store(cx);
371 cx.indicate(Ind::Decide(v));
372 }
373
374 fn through_lec(
375 &mut self,
376 cx: &mut ProtoCx<'_, Self>,
377 f: impl FnOnce(&mut LoggedEpochChange, &mut ProtoCx<'_, LoggedEpochChange>),
378 ) {
379 let mut inds = self.lec.run_durable(cx, Wire::Change, lec_slot(), f);
380 for lec::Ind::StartEpoch { ts, leader } in inds.drain(..) {
381 self.newts = ts;
382 self.newleader = leader;
383 self.maybe_abort(cx);
384 }
385 self.lec.reclaim(inds);
386 }
387
388 fn through_lep(
389 &mut self,
390 cx: &mut ProtoCx<'_, Self>,
391 f: impl FnOnce(&mut LoggedEpochConsensus<V>, &mut ProtoCx<'_, LoggedEpochConsensus<V>>),
392 ) {
393 let mut inds = self.lep.run_durable(cx, Wire::Consensus, lep_slot(), f);
394 for ind in inds.drain(..) {
395 match ind {
396 lep::Ind::Decide(v) => self.on_epoch_decision(v, cx),
397 lep::Ind::Aborted(state) => self.on_aborted(state, cx),
398 }
399 }
400 self.lep.reclaim(inds);
401 }
402}
403
404impl<V: Clone + PartialEq> Protocol for LoggedLeaderDrivenConsensus<V> {
405 type Cmd = Cmd<V>;
406 type Ind = Ind<V>;
407 type Msg = Wire<V>;
408 type Scope = core::convert::Infallible;
409 type Note = crate::Note;
410 type Meta = Durable<V>;
411 /// Nothing accumulates: one epoch, one leader, one decision, and one record per child.
412 type Entry = core::convert::Infallible;
413
414 /// `upon event ⟨ luc, Propose | v ⟩ do val := v`.
415 fn on_cmd(&mut self, Cmd::Propose(v): Cmd<V>, cx: &mut ProtoCx<'_, Self>) {
416 self.val = Some(v);
417 self.maybe_propose(cx);
418 }
419
420 fn on_msg(&mut self, from: NodeId, msg: Wire<V>, cx: &mut ProtoCx<'_, Self>) {
421 match msg {
422 Wire::Change(m) => self.through_lec(cx, |lec, ccx| lec.on_msg(from, m, ccx)),
423 // The instance guard is inside the child, which drops anything not stamped with its
424 // own epoch. See `lep::Tagged`.
425 Wire::Consensus(m) => self.through_lep(cx, |lep, ccx| lep.on_msg(from, m, ccx)),
426 }
427 }
428
429 fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
430 self.through_lec(cx, |lec, ccx| lec.on_timer(id, ccx));
431 self.through_lep(cx, |lep, ccx| lep.on_timer(id, ccx));
432 }
433
434 /// `upon event ⟨ luc, Init ⟩ … store(ets, ℓ, decision)`.
435 ///
436 /// This process's own record goes down first, then each child's, so the record exists before
437 /// anything writes into a slot of it.
438 fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
439 self.store(cx);
440 self.through_lec(cx, |lec, ccx| lec.on_init(ccx));
441 self.through_lep(cx, |lep, ccx| lep.on_init(ccx));
442 }
443
444 /// `upon event ⟨ luc, Recovery ⟩`.
445 ///
446 /// The book reads its children's records by name — `retrieve(startts, start) of instance lec`,
447 /// `retrieve(epochdecision) of instance lep.ets`. Here each child reads its own slot in its own
448 /// `Recovery`, which is the same statement in the direction the composition runs.
449 fn on_recovery(&mut self, cx: &mut ProtoCx<'_, Self>) {
450 // `retrieve(ets, ℓ, decision)`
451 if let Some(held) = cx.storage().get().cloned() {
452 self.ets = held.ets;
453 self.leader = held.leader.unwrap_or(self.l0);
454 self.decision = held.decision;
455 }
456
457 // `retrieve(startts, start) of instance lec; (newts, newℓ) := (startts, start)`
458 self.through_lec(cx, |lec, ccx| lec.on_recovery(ccx));
459 self.newts = self.lec.last_timestamp();
460 self.newleader = self.lec.last_leader();
461
462 // `retrieve(epochdecision) of instance lep.ets`. The instance is rebuilt at the epoch just
463 // read back, and reads its own slot.
464 self.lep.replace(LoggedEpochConsensus::new(
465 self.me,
466 self.peers.clone(),
467 self.ets,
468 self.leader,
469 State::default(),
470 self.timing.retransmit,
471 ));
472 self.through_lep(cx, |lep, ccx| lep.on_recovery(ccx));
473
474 // `if epochdecision ≠ ⊥ ∧ decision = ⊥ then decision := epochdecision; store(decision);
475 // trigger ⟨ luc, Decide | decision ⟩`
476 //
477 // This is what makes `UC3` hold across a restart from the *other* side: a process that
478 // decided and crashed before anything above it saw the indication is told again.
479 let epoch_decision = self.lep.epoch_decision().cloned();
480 if let Some(v) = epoch_decision
481 && self.decision.is_none()
482 {
483 self.decision = Some(v.clone());
484 self.store(cx);
485 cx.indicate(Ind::Decide(v));
486 } else if self.decision.is_some() {
487 // Decided before the crash, and the record says so. Re-announced for the same reason:
488 // the layer above may never have seen the first indication.
489 let v = self.decision.clone().expect("just checked");
490 cx.indicate(Ind::Decide(v));
491 }
492
493 // `aborted := FALSE`
494 self.aborted = false;
495 self.proposed = false;
496 self.maybe_abort(cx);
497 }
498}