Skip to main content

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}