recon_protocols/consensus_based_total_order_broadcast.rs
1//! Consensus-based total-order broadcast.
2//!
3//! **Status: transcription. Space: unbounded — `unordered`, `delivered` and the family of consensus
4//! instances all grow with the number of entries handled.** That is the page, and
5//! `docs/bounded-space.md` is explicit that inheriting the book's omissions is correct of a
6//! transcription and disqualifying of an implementation. Bounding any of them weakens a guarantee to
7//! a scope and belongs to a change with a proposal.
8//!
9//! Cachin, Guerraoui & Rodrigues, Module 6.1 (`TotalOrderBroadcast`) and Algorithm 6.1, quoted from
10//! the book:
11//!
12//! ```text
13//! Algorithm 6.1: Consensus-Based Total-Order Broadcast
14//! Implements: TotalOrderBroadcast, instance tob.
15//! Uses:
16//! ReliableBroadcast, instance rb;
17//! Consensus (multiple instances).
18//!
19//! upon event ⟨ tob, Init ⟩ do
20//! unordered := ∅;
21//! delivered := ∅;
22//! round := 1;
23//! wait := FALSE;
24//!
25//! upon event ⟨ tob, Broadcast | m ⟩ do
26//! trigger ⟨ rb, Broadcast | m ⟩;
27//!
28//! upon event ⟨ rb, Deliver | p, m ⟩ do
29//! if m ∉ delivered then
30//! unordered := unordered ∪ {(p, m)};
31//!
32//! upon unordered ≠ ∅ ∧ wait = FALSE do
33//! wait := TRUE;
34//! Initialize a new instance c.round of consensus;
35//! trigger ⟨ c.round, Propose | unordered ⟩;
36//!
37//! upon event ⟨ c.r, Decide | decided ⟩ such that r = round do
38//! forall (s, m) ∈ sort(decided) do // by the order in the resulting sorted list
39//! trigger ⟨ tob, Deliver | s, m ⟩;
40//! delivered := delivered ∪ decided;
41//! unordered := unordered \ decided;
42//! round := round + 1;
43//! wait := FALSE;
44//! ```
45//!
46//! The shape is one round at a time: everything reliable broadcast has delivered and this process
47//! has not yet ordered is proposed as a *set*, consensus agrees on a set, and every process turns
48//! that same set into the same sequence by sorting it. Ordering is therefore agreed without anyone
49//! communicating about order at all — the sort does that work, which is why it must be deterministic.
50//!
51//! # Departures from the page
52//!
53//! - **A read.** The port this satisfies offers [`crate::total_order_log::TotalOrderLog::read`],
54//! which the book's abstraction does not: its clients observe deliveries. The algorithm already
55//! maintains `delivered`, so the read exposes what the page keeps and does not offer, served
56//! locally. See the port's own documentation.
57//!
58//! - **Consensus instances are held explicitly, keyed by round.** The page writes `c.round` and
59//! `⟨ c.r, Decide ⟩`, so instances are a family addressed by round; the book's runtime routes to
60//! them and this one does not. They are created on demand — including for a round this process has
61//! not reached, which is what lets a peer that is ahead make progress — and never pruned, as the
62//! page has them. Creation runs the instance's `⟨ Init ⟩` before the event that provoked it:
63//! "Initialize a new instance c.round" is an event the book's runtime delivers, and skipping it
64//! leaves the instance's failure detector without its timers — which no fault-free run notices,
65//! because deciding under the initial epoch never consults the detector. A crash is then never
66//! detected and the survivors stall, which is what the suite's crash property caught.
67//!
68//! - **The conditional event handler is discharged here.** `such that r = round` is not a guard that
69//! discards. The book states its meaning: "An algorithm that uses conditional event handlers
70//! relies on the run-time system to buffer external events until the condition on internal
71//! variables becomes satisfied." `Cx` has no such facility, so a decision for a round this process
72//! has not reached is **held** in `decisions` and acted on when `round` catches up.
73//! [`crate::leader_driven_consensus`]'s `pending` is the same pattern, for the same reason.
74//!
75//! - **A consensus message carries its round.** The page addresses instances; nothing on this wire
76//! does. Unlike [`crate::epoch_consensus`], whose instance stamps its own messages because the
77//! epoch is its identity, the round is *this* layer's concept and not the consensus's — so this
78//! layer stamps, and the stamp is the one thing it cannot delegate.
79//!
80//! - **`unordered` is deduplicated on the pair `(p, m)`**, not on `m` alone. The page's
81//! `if m ∉ delivered` reads against a `delivered` that holds pairs, so one of the two is loose;
82//! the pair is what makes `unordered \ decided` well defined, and it is what is used here. The
83//! consequence, which the page shares: one process appending the same value twice contributes one
84//! entry.
85//!
86//! - **The consensus beneath is not a type parameter, and its link is fixed.** Every other composing
87//! layer here takes its child as a parameter with a default. This one cannot for the consensus:
88//! instances are created *at run time*, one per round, so the layer would need a link **factory**
89//! rather than a link — a runtime indirection the static composition model exists to avoid. The
90//! reliable broadcast keeps its parameter, because there is exactly one of it and the caller
91//! supplies it once. Revisit if a stack ever wants rounds agreed over something else.
92//!
93//! - **The sort is a `BTreeSet`'s own order.** Proposing an ordered set means `sort(decided)` is
94//! iteration, and every process computes the same sequence because `Ord` is a function of the
95//! values rather than of anything local. The guard forbidding hash-keyed maps in these three
96//! crates exists for the converse reason: an iteration order that varies per process is exactly
97//! what would break agreement here.
98
99use recon_core::{Child, NodeId, Position, ProtoCx, Protocol, TimerId};
100use serde::{Deserialize, Serialize};
101use std::collections::{BTreeMap, BTreeSet};
102
103use crate::Timing;
104use crate::flooding_consensus::{self as fc, FloodingConsensus};
105use crate::link::{Boundary, VolatileLink};
106use crate::perfect_link::PerfectLink;
107use crate::reliable_broadcast::{self as rb, ReliableBroadcast};
108use crate::total_order_log::{LogInd, TotalOrderLog};
109
110/// One entry with the process that appended it — the page's `(s, m)`.
111///
112/// `Ord` is what makes the ordering agreed: consensus decides a *set*, and every process must turn
113/// that set into the same sequence with no further communication.
114pub type Slot<V> = (NodeId, V);
115
116/// What a round proposes and decides: the set of entries not yet ordered.
117pub type Batch<V> = BTreeSet<Slot<V>>;
118
119/// What reliable broadcast carries for this layer.
120pub type Carried<V> = rb::Carried<V>;
121
122/// What the consensus beneath carries for this layer.
123pub type ConsensusCarried<V> = fc::Flood<Batch<V>>;
124
125/// The consensus one round runs. Not a type parameter — see the module's departures.
126pub type Consensus<V> = FloodingConsensus<Batch<V>, PerfectLink<ConsensusCarried<V>>>;
127
128/// This layer's messages: the broadcast's, and a consensus instance's stamped with its round.
129#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
130pub enum Wire<R, C> {
131 /// A reliable broadcast message.
132 Broadcast(R),
133 /// A consensus message, stamped with the round whose instance it belongs to.
134 ///
135 /// The page addresses instances as `c.round`; nothing on this wire does. The stamp is this
136 /// layer's rather than the consensus's, because the round is this layer's concept — which is why
137 /// it cannot live in the child as [`crate::epoch_consensus::Tagged`]'s does.
138 Consensus { round: u64, msg: C },
139}
140
141/// Requests from the layer above.
142#[derive(Debug, Clone, PartialEq, Eq)]
143pub enum Cmd<V> {
144 /// `⟨ tob, Broadcast | m ⟩`, in the port's terms: append `v` to the log.
145 Append(V),
146 /// Read the ordered sequence from `from` onwards. The page has no such request; see the
147 /// module's departures.
148 Read { from: Position },
149}
150
151/// Indications to the layer above.
152#[derive(Debug, Clone, PartialEq, Eq)]
153pub enum Ind<V> {
154 /// `⟨ tob, Deliver | s, m ⟩`, with the position it took in the agreed sequence.
155 Ordered { position: Position, from: NodeId, value: V },
156 /// The answer to a [`Cmd::Read`].
157 Contents { from: Position, entries: Vec<V> },
158 /// The scope with `peer` ended at `epoch`, as reliable broadcast beneath reported it.
159 ///
160 /// Propagated rather than absorbed: this layer's only redundancy is the broadcast's, which
161 /// relays once and cannot resend across an ending, so it cannot bridge either. Raised only over
162 /// a link that reports boundaries, which the default stack's perfect link does not.
163 SessionEnded { peer: NodeId, epoch: u64 },
164 /// A scope with `peer` is in force at `epoch`.
165 SessionEstablished { peer: NodeId, epoch: u64 },
166}
167
168/// A totally ordered log, agreed by one consensus instance per round.
169pub struct ConsensusBasedTotalOrderBroadcast<
170 V: Clone + Ord,
171 L: VolatileLink<Carried<V>> = PerfectLink<Carried<V>>,
172> {
173 me: NodeId,
174 peers: BTreeSet<NodeId>,
175 timing: Timing,
176 /// `unordered` — delivered by the broadcast beneath, not yet ordered.
177 unordered: Batch<V>,
178 /// `delivered`, as the agreed sequence. The page writes a set; the order is the whole point, so
179 /// it is kept as one and `ordered` is the membership test the page's `∉` needs.
180 delivered: Vec<Slot<V>>,
181 ordered: BTreeSet<Slot<V>>,
182 /// `round`.
183 round: u64,
184 /// `wait`.
185 wait: bool,
186 rb: Child<ReliableBroadcast<V, L>>,
187 /// `c.r` for every `r` this process has started or been sent a message for. A family, as the
188 /// page has it — never one instance replaced, because no round supersedes another.
189 consensus: BTreeMap<u64, Child<Consensus<V>>>,
190 /// Decisions for rounds not yet reached, held rather than discarded — the conditional event
191 /// handler `such that r = round`, discharged here because `Cx` cannot buffer on a condition.
192 decisions: BTreeMap<u64, Batch<V>>,
193}
194
195impl<V: Clone + Ord> ConsensusBasedTotalOrderBroadcast<V> {
196 /// A totally ordered log among `peers`, over perfect links.
197 ///
198 /// `timing` is the consensus beneath's: its `detect_after` must exceed `heartbeat` plus the
199 /// network's delivery bound, or the perfect failure detector under flooding consensus accuses a
200 /// correct process and agreement can break.
201 pub fn new(me: NodeId, peers: impl IntoIterator<Item = NodeId>, timing: Timing) -> Self {
202 let mut peers: BTreeSet<NodeId> = peers.into_iter().collect();
203 peers.insert(me);
204 ConsensusBasedTotalOrderBroadcast {
205 me,
206 peers: peers.clone(),
207 timing,
208 unordered: BTreeSet::new(),
209 delivered: Vec::new(),
210 ordered: BTreeSet::new(),
211 round: 1,
212 wait: false,
213 rb: Child::new(ReliableBroadcast::new(me, peers, timing.retransmit)),
214 consensus: BTreeMap::new(),
215 decisions: BTreeMap::new(),
216 }
217 }
218}
219
220impl<V: Clone + Ord, L: VolatileLink<Carried<V>>> ConsensusBasedTotalOrderBroadcast<V, L> {
221 /// The agreed sequence as this process holds it.
222 pub fn entries(&self) -> &[Slot<V>] {
223 &self.delivered
224 }
225
226 /// How many entries this process has ordered.
227 pub fn len(&self) -> usize {
228 self.delivered.len()
229 }
230
231 pub fn is_empty(&self) -> bool {
232 self.delivered.is_empty()
233 }
234
235 /// The round this process is running.
236 pub fn round(&self) -> u64 {
237 self.round
238 }
239
240 /// How many consensus instances are held. Grows with rounds, and is never pruned — see the
241 /// module's space statement.
242 pub fn instances(&self) -> usize {
243 self.consensus.len()
244 }
245
246 /// `upon unordered ≠ ∅ ∧ wait = FALSE do` — a standing condition, re-evaluated whenever either
247 /// could have changed.
248 fn maybe_propose(&mut self, cx: &mut ProtoCx<'_, Self>) {
249 if self.unordered.is_empty() || self.wait {
250 return;
251 }
252 self.wait = true;
253 let round = self.round;
254 let batch = self.unordered.clone();
255 self.through_consensus(round, cx, |c, ccx| c.on_cmd(fc::Cmd::Propose(batch), ccx));
256 }
257
258 /// `upon event ⟨ c.r, Decide | decided ⟩ such that r = round`, with the buffering the book's
259 /// run-time system would have done.
260 fn drain_decisions(&mut self, cx: &mut ProtoCx<'_, Self>) {
261 while let Some(decided) = self.decisions.remove(&self.round) {
262 // `forall (s, m) ∈ sort(decided)` — iteration of an ordered set *is* the sort, and every
263 // process computes the same one because `Ord` reads only the values.
264 for (from, value) in &decided {
265 if self.ordered.contains(&(*from, value.clone())) {
266 continue;
267 }
268 let position = Position(self.delivered.len() as u64);
269 self.delivered.push((*from, value.clone()));
270 self.ordered.insert((*from, value.clone()));
271 cx.indicate(Ind::Ordered { position, from: *from, value: value.clone() });
272 }
273 for slot in &decided {
274 self.unordered.remove(slot);
275 }
276 self.round += 1;
277 self.wait = false;
278 }
279 self.maybe_propose(cx);
280 }
281
282 fn through_rb(
283 &mut self,
284 cx: &mut ProtoCx<'_, Self>,
285 f: impl FnOnce(&mut ReliableBroadcast<V, L>, &mut ProtoCx<'_, ReliableBroadcast<V, L>>),
286 ) {
287 let mut inds = self.rb.run(cx, Wire::Broadcast, f);
288 for ind in inds.drain(..) {
289 match ind {
290 // `upon event ⟨ rb, Deliver | p, m ⟩ do if (p, m) ∉ delivered then …`
291 rb::Ind::Deliver { from, msg } => {
292 let slot = (from, msg);
293 if !self.ordered.contains(&slot) {
294 self.unordered.insert(slot);
295 }
296 }
297 // Reliable broadcast cannot bridge a scope ending and propagates; neither can this
298 // layer, whose only redundancy is the broadcast's. Over a perfect link the arm is
299 // unreachable — see `link.rs`.
300 rb::Ind::SessionEnded { peer, epoch } => {
301 cx.indicate(Ind::SessionEnded { peer, epoch });
302 }
303 rb::Ind::SessionEstablished { peer, epoch } => {
304 cx.indicate(Ind::SessionEstablished { peer, epoch });
305 }
306 }
307 }
308 self.rb.reclaim(inds);
309 self.maybe_propose(cx);
310 }
311
312 /// Run `f` against round `r`'s instance, creating it if this process has not started it.
313 fn through_consensus(
314 &mut self,
315 r: u64,
316 cx: &mut ProtoCx<'_, Self>,
317 f: impl FnOnce(&mut Consensus<V>, &mut ProtoCx<'_, Consensus<V>>),
318 ) {
319 // "Initialize a new instance c.round of consensus" — an event, which runs before whatever
320 // provoked the creation. See the module's departures for what skipping it cost.
321 let created = !self.consensus.contains_key(&r);
322 let entry = self.consensus.entry(r).or_insert_with(|| {
323 Child::new(FloodingConsensus::new(
324 self.me,
325 self.peers.clone(),
326 self.timing.retransmit,
327 self.timing.heartbeat,
328 self.timing.detect_after,
329 ))
330 });
331 let mut inds = entry.run(
332 cx,
333 |m| Wire::Consensus { round: r, msg: m },
334 |c, ccx| {
335 if created {
336 c.on_init(ccx);
337 }
338 f(c, ccx)
339 },
340 );
341 for fc::Ind::Decide(decided) in inds.drain(..) {
342 self.decisions.entry(r).or_insert(decided);
343 }
344 if let Some(child) = self.consensus.get_mut(&r) {
345 child.reclaim(inds);
346 }
347 self.drain_decisions(cx);
348 }
349}
350
351impl<V: Clone + Ord, L: VolatileLink<Carried<V>>> Protocol
352 for ConsensusBasedTotalOrderBroadcast<V, L>
353{
354 type Cmd = Cmd<V>;
355 type Ind = Ind<V>;
356 type Msg = Wire<rb::Wire<V, L>, <Consensus<V> as Protocol>::Msg>;
357 type Scope = core::convert::Infallible;
358 type Note = crate::Note;
359 /// Keeps nothing durably. The fail-recovery variant is the one that does.
360 type Meta = core::convert::Infallible;
361 type Entry = core::convert::Infallible;
362
363 fn on_cmd(&mut self, cmd: Cmd<V>, cx: &mut ProtoCx<'_, Self>) {
364 match cmd {
365 // `upon event ⟨ tob, Broadcast | m ⟩ do trigger ⟨ rb, Broadcast | m ⟩`
366 Cmd::Append(v) => self.through_rb(cx, |r, ccx| r.on_cmd(rb::Cmd::Broadcast(v), ccx)),
367 // The departure. Served from this process's own sequence, so it may lag.
368 Cmd::Read { from } => {
369 let entries: Vec<V> =
370 self.delivered.iter().skip(from.0 as usize).map(|(_, v)| v.clone()).collect();
371 cx.indicate(Ind::Contents { from, entries });
372 }
373 }
374 }
375
376 fn on_msg(&mut self, from: NodeId, msg: Self::Msg, cx: &mut ProtoCx<'_, Self>) {
377 match msg {
378 Wire::Broadcast(m) => self.through_rb(cx, |r, ccx| r.on_msg(from, m, ccx)),
379 // Routed to the round's own instance, creating it if this process has not started that
380 // round — which is what lets a peer that is ahead make progress.
381 Wire::Consensus { round, msg } => {
382 self.through_consensus(round, cx, |c, ccx| c.on_msg(from, msg, ccx));
383 }
384 }
385 }
386
387 fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
388 self.through_rb(cx, |r, ccx| r.on_timer(id, ccx));
389 let rounds: Vec<u64> = self.consensus.keys().copied().collect();
390 for r in rounds {
391 self.through_consensus(r, cx, |c, ccx| c.on_timer(id, ccx));
392 }
393 }
394
395 fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
396 self.through_rb(cx, |r, ccx| r.on_init(ccx));
397 }
398}
399
400impl<V: Clone + Ord, L: VolatileLink<Carried<V>>> TotalOrderLog<V>
401 for ConsensusBasedTotalOrderBroadcast<V, L>
402{
403 fn append(value: V) -> Cmd<V> {
404 Cmd::Append(value)
405 }
406
407 fn read(from: Position) -> Cmd<V> {
408 Cmd::Read { from }
409 }
410
411 fn classify(ind: Ind<V>) -> LogInd<V> {
412 match ind {
413 Ind::Ordered { position, from, value } => LogInd::Ordered { position, from, value },
414 Ind::Contents { from, entries } => LogInd::Contents { from, entries },
415 Ind::SessionEnded { peer, epoch } => LogInd::Boundary(Boundary::Ended { peer, epoch }),
416 Ind::SessionEstablished { peer, epoch } => {
417 LogInd::Boundary(Boundary::Established { peer, epoch })
418 }
419 }
420 }
421}