Skip to main content

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}