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}