Skip to main content

recon_protocols/
leader_driven_consensus.rs

1//! Leader-driven consensus — Paxos.
2//!
3//! **Status: implementation. Space: bounded by membership.**
4//!
5//! Cachin, Guerraoui & Rodrigues, Module 5.2 (`UniformConsensus`) and Algorithm 5.7
6//! ("Leader-Driven Consensus"), quoted from the book:
7//!
8//! ```text
9//! Algorithm 5.7: Leader-Driven Consensus
10//! Implements: UniformConsensus, instance uc.
11//! Uses:
12//!     EpochChange, instance ec;
13//!     EpochConsensus (multiple instances).
14//!
15//! upon event ⟨ uc, Init ⟩ do
16//!     val := ⊥; proposed := FALSE; decided := FALSE;
17//!     Obtain the leader ℓ0 of the initial epoch with timestamp 0 from epoch-change instance ec;
18//!     Initialize a new instance ep.0 of epoch consensus with timestamp 0, leader ℓ0, state (0, ⊥);
19//!     (ets, ℓ) := (0, ℓ0);
20//!     (newts, newℓ) := (0, ⊥);
21//!
22//! upon event ⟨ uc, Propose | v ⟩ do
23//!     val := v;
24//!
25//! upon event ⟨ ec, StartEpoch | newts′, newℓ′ ⟩ do
26//!     (newts, newℓ) := (newts′, newℓ′);
27//!     trigger ⟨ ep.ets, Abort ⟩;
28//!
29//! upon event ⟨ ep.ts, Aborted | state ⟩ such that ts = ets do
30//!     (ets, ℓ) := (newts, newℓ);
31//!     proposed := FALSE;
32//!     Initialize a new instance ep.ets of epoch consensus with timestamp ets, leader ℓ, and
33//!         state state;
34//!
35//! upon ℓ = self ∧ val ≠ ⊥ ∧ proposed = FALSE do
36//!     proposed := TRUE;
37//!     trigger ⟨ ep.ets, Propose | val ⟩;
38//!
39//! upon event ⟨ ep.ts, Decide | v ⟩ such that ts = ets do
40//!     if decided = FALSE then
41//!         decided := TRUE;
42//!         trigger ⟨ uc, Decide | v ⟩;
43//! ```
44//!
45//! # The child is replaced while running, and the state is what carries across
46//!
47//! Every other layer in this repository constructs its children once. This one does not: each epoch
48//! gets a **new** epoch-consensus instance, seeded with the state the previous one returned when it
49//! was aborted. `CLAUDE.md`'s rule still holds — the field is one concrete type, replaced, not a map
50//! from timestamp to instance resolved while running — but it is the first layer here whose child is
51//! rebuilt at all.
52//!
53//! **The abort handshake is asynchronous and the wait is load-bearing.** `Abort` is a request;
54//! `Aborted` is its answer, carrying `(valts, val)`. The replacement is not constructed until that
55//! answer arrives, because the state it carries is what stops the new epoch contradicting the old
56//! one. Aborting and immediately replacing — which is the obvious implementation — loses it, and
57//! with it the property the whole algorithm exists to have.
58//!
59//! # Addition: epoch-consensus traffic is tagged with its epoch
60//!
61//! The book writes `ep.ts` and `such that ts = ets`, so instances are addressed by timestamp and a
62//! message for one never reaches another. Nothing in this codebase's wire does that for free, and
63//! the consequence of omitting it is a **safety** failure rather than a lost message: a `WRITE` from
64//! epoch 7 arriving after epoch 11 has begun would be accepted and recorded at timestamp 11,
65//! inventing an acceptance that never happened.
66//!
67//! The stamp lives in [`ep::Tagged`], inside the child, rather than in this layer's wire. The epoch
68//! is the instance's own identity, so the instance stamps what it sends and drops what is not
69//! addressed to it — and this layer cannot forget to. An earlier draft put it here, and the type
70//! system objected for an unrelated reason, which turned out to be pointing at the better place.
71//!
72//! ```text
73//! UC1 [always]      Validity — a decided value was proposed by some process
74//! UC2 [always]      Uniform agreement — no two processes decide differently, **including while the
75//!                   leader detector is wrong and two processes each believe they lead**
76//! UC3 [always]      Integrity — a process decides at most once
77//! UC4 [conditional] Termination — every correct process decides, provided a majority is correct
78//!                   and the leader detector eventually settles
79//! ```
80//!
81//! `UC4` is conditional on both, and saying so is the point. A majority that never forms, or a
82//! detector that never settles, leaves this waiting — which is what FLP requires of it and what
83//! `flooding_consensus` pretends away by assuming a perfect detector.
84
85use recon_core::{Child, NodeId, ProtoCx, Protocol, TimerId};
86use serde::{Deserialize, Serialize};
87use std::collections::BTreeSet;
88
89use crate::Timing;
90use crate::epoch_change::{self as ec, EpochChange};
91use crate::epoch_consensus::{self as ep, EpochConsensus, State};
92
93/// The wire, multiplexing the epoch-change child and whichever epoch-consensus instance is live.
94#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
95pub enum Wire<E, C> {
96    /// The epoch-change child's traffic.
97    Change(E),
98    /// The live epoch-consensus instance's traffic.
99    ///
100    /// The epoch tag that makes `ep.ts` addressable lives inside `C` rather than here — see
101    /// [`ep::Tagged`]. It belongs to the instance, so the instance stamps it and drops what is not
102    /// its own, and a parent cannot forget to.
103    Consensus(C),
104}
105
106/// Requests from the layer above.
107#[derive(Debug, Clone, PartialEq, Eq)]
108pub enum Cmd<V> {
109    /// `⟨ uc, Propose | v ⟩`.
110    Propose(V),
111}
112
113/// Indications to the layer above.
114#[derive(Debug, Clone, PartialEq, Eq)]
115pub enum Ind<V> {
116    /// `⟨ uc, Decide | v ⟩`. Raised at most once.
117    Decide(V),
118}
119
120/// Paxos: uniform consensus over an epoch-change and a sequence of abortable epoch consensuses.
121pub struct LeaderDrivenConsensus<V: Clone> {
122    me: NodeId,
123    peers: BTreeSet<NodeId>,
124    timing: Timing,
125    /// `val` — what this process wants decided, once the layer above has said.
126    val: Option<V>,
127    /// `proposed` — whether this process has proposed in the current epoch.
128    proposed: bool,
129    /// `decided` — whether it has reported a decision. Reported at most once.
130    decided: bool,
131    /// `(ets, ℓ)` — the epoch now live, and who leads it.
132    ets: u64,
133    leader: NodeId,
134    /// `(newts, newℓ)` — the epoch waiting for the current one to finish aborting.
135    pending: Option<(u64, NodeId)>,
136    ec: Child<EpochChange>,
137    /// `ep.ets`. One instance, replaced on each epoch change — never a map.
138    ep: Child<EpochConsensus<V>>,
139}
140
141impl<V: Clone> LeaderDrivenConsensus<V> {
142    /// Paxos among `peers`.
143    pub fn new(me: NodeId, peers: impl IntoIterator<Item = NodeId>, timing: Timing) -> Self {
144        let mut peers: BTreeSet<NodeId> = peers.into_iter().collect();
145        peers.insert(me);
146        // "Obtain the leader ℓ0 of the initial epoch with timestamp 0 from epoch-change" — the same
147        // `maxrank(Π)` the leader detector will trust first, so the two agree before anything moves.
148        let l0 = peers.iter().next_back().copied().expect("Π contains at least this process");
149        LeaderDrivenConsensus {
150            me,
151            peers: peers.clone(),
152            timing,
153            val: None,
154            proposed: false,
155            decided: false,
156            ets: 0,
157            leader: l0,
158            pending: None,
159            ec: Child::new(EpochChange::new(me, peers.clone(), timing)),
160            ep: Child::new(EpochConsensus::new(
161                me,
162                peers,
163                0,
164                l0,
165                State::default(),
166                timing.retransmit,
167            )),
168        }
169    }
170
171    /// The epoch now live at this process.
172    pub fn epoch(&self) -> u64 {
173        self.ets
174    }
175
176    /// Who leads the epoch now live.
177    pub fn leader(&self) -> NodeId {
178        self.leader
179    }
180
181    /// Whether this process has reported a decision.
182    pub fn has_decided(&self) -> bool {
183        self.decided
184    }
185
186    /// `(valts, val)` as the epoch now live holds it.
187    ///
188    /// This is what the abort handshake carries forward: the state a new instance is constructed
189    /// from is the state the aborted one returned, so a value accepted in an early epoch is still
190    /// here after several leadership changes.
191    pub fn state(&self) -> &State<V> {
192        self.ep.state()
193    }
194
195    /// `upon ℓ = self ∧ val ≠ ⊥ ∧ proposed = FALSE` — a standing condition, re-evaluated whenever
196    /// any of the three could have changed.
197    fn maybe_propose(&mut self, cx: &mut ProtoCx<'_, Self>) {
198        if self.leader == self.me
199            && !self.proposed
200            && let Some(v) = self.val.clone()
201        {
202            self.proposed = true;
203            self.through_ep(cx, |e, ccx| e.on_cmd(ep::Cmd::Propose(v), ccx));
204        }
205    }
206
207    /// `upon event ⟨ ec, StartEpoch | newts, newℓ ⟩ do … trigger ⟨ ep.ets, Abort ⟩`.
208    fn on_start_epoch(&mut self, ts: u64, leader: NodeId, cx: &mut ProtoCx<'_, Self>) {
209        self.pending = Some((ts, leader));
210        self.through_ep(cx, |e, ccx| e.on_cmd(ep::Cmd::Abort, ccx));
211    }
212
213    /// `upon event ⟨ ep.ts, Aborted | state ⟩ such that ts = ets`.
214    ///
215    /// The replacement is built here and nowhere else, because `state` is only available here.
216    fn on_aborted(&mut self, state: State<V>, cx: &mut ProtoCx<'_, Self>) {
217        let Some((ts, leader)) = self.pending.take() else {
218            // An `Aborted` with no epoch waiting: the book's `such that ts = ets` guard, which
219            // discards an answer from an instance already superseded.
220            return;
221        };
222        self.ets = ts;
223        self.leader = leader;
224        self.proposed = false;
225        self.ep.replace(EpochConsensus::new(
226            self.me,
227            self.peers.clone(),
228            ts,
229            leader,
230            state,
231            self.timing.retransmit,
232        ));
233        self.maybe_propose(cx);
234    }
235
236    fn through_ec(
237        &mut self,
238        cx: &mut ProtoCx<'_, Self>,
239        f: impl FnOnce(&mut EpochChange, &mut ProtoCx<'_, EpochChange>),
240    ) {
241        let mut inds = self.ec.run(cx, Wire::Change, f);
242        for ec::Ind::StartEpoch { ts, leader } in inds.drain(..) {
243            self.on_start_epoch(ts, leader, cx);
244        }
245        self.ec.reclaim(inds);
246    }
247
248    fn through_ep(
249        &mut self,
250        cx: &mut ProtoCx<'_, Self>,
251        f: impl FnOnce(&mut EpochConsensus<V>, &mut ProtoCx<'_, EpochConsensus<V>>),
252    ) {
253        let mut inds = self.ep.run(cx, Wire::Consensus, f);
254        for ind in inds.drain(..) {
255            match ind {
256                // `upon event ⟨ ep.ts, Decide | v ⟩ such that ts = ets do if decided = FALSE …`
257                ep::Ind::Decide(v) => {
258                    if !self.decided {
259                        self.decided = true;
260                        cx.indicate(Ind::Decide(v));
261                    }
262                }
263                ep::Ind::Aborted(state) => self.on_aborted(state, cx),
264            }
265        }
266        self.ep.reclaim(inds);
267    }
268}
269
270impl<V: Clone> Protocol for LeaderDrivenConsensus<V> {
271    type Cmd = Cmd<V>;
272    type Ind = Ind<V>;
273    type Msg = Wire<<EpochChange as Protocol>::Msg, ep::BebMsg<V>>;
274    type Scope = core::convert::Infallible;
275    type Note = crate::Note;
276    /// Keeps nothing durably. `logged_leader_driven_consensus` is the variant that does.
277    type Meta = core::convert::Infallible;
278    type Entry = core::convert::Infallible;
279
280    /// `upon event ⟨ uc, Propose | v ⟩ do val := v`.
281    fn on_cmd(&mut self, Cmd::Propose(v): Cmd<V>, cx: &mut ProtoCx<'_, Self>) {
282        self.val = Some(v);
283        self.maybe_propose(cx);
284    }
285
286    fn on_msg(&mut self, from: NodeId, msg: Self::Msg, cx: &mut ProtoCx<'_, Self>) {
287        match msg {
288            Wire::Change(m) => self.through_ec(cx, |ec, ccx| ec.on_msg(from, m, ccx)),
289            // The instance guard is inside the child, which drops anything not stamped with its
290            // own epoch. See `ep::Tagged`.
291            Wire::Consensus(m) => self.through_ep(cx, |ep, ccx| ep.on_msg(from, m, ccx)),
292        }
293    }
294
295    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
296        self.through_ec(cx, |ec, ccx| ec.on_timer(id, ccx));
297        self.through_ep(cx, |ep, ccx| ep.on_timer(id, ccx));
298    }
299
300    fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
301        self.through_ec(cx, |ec, ccx| ec.on_init(ccx));
302    }
303}