Skip to main content

recon_protocols/
logged_uniform_total_order_broadcast.rs

1//! Logged uniform 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, and `delivered` and `proposals` grow in
5//! stable storage as well.** That is the page. See `docs/bounded-space.md`.
6//!
7//! Cachin, Guerraoui & Rodrigues, Module `LoggedUniformTotalOrderBroadcast` and Algorithm 6.12,
8//! p. 327, quoted from the book:
9//!
10//! ```text
11//! Algorithm 6.12: Logged Uniform Total-Order Broadcast
12//! Implements: LoggedUniformTotalOrderBroadcast, instance lutob.
13//! Uses:
14//!     LoggedUniformReliableBroadcast, instance lurb;
15//!     LoggedUniformConsensus (multiple instances).
16//!
17//! upon event ⟨ lutob, Init ⟩ do
18//!     unordered := ∅;
19//!     delivered := [];
20//!     round := 1;
21//!     recovering := FALSE;
22//!     wait := FALSE;
23//!     forall r > 0 do proposals[r] := ⊥;
24//!
25//! upon event ⟨ Recovery ⟩ do
26//!     unordered := ∅;
27//!     delivered := [];
28//!     round := 1;
29//!     recovering := TRUE;
30//!     wait := FALSE;
31//!     retrieve(proposals);
32//!     if proposals[1] ≠ ⊥ then
33//!         trigger ⟨ luc.1, Propose | proposals[1] ⟩;
34//!
35//! upon event ⟨ lutob, Broadcast | m ⟩ do
36//!     trigger ⟨ lurb, Broadcast | m ⟩;
37//!
38//! upon event ⟨ lurb, Deliver | lurbdelivered ⟩ do
39//!     unordered := unordered ∪ lurbdelivered;
40//!
41//! upon unordered \ delivered ≠ ∅ ∧ wait = FALSE ∧ recovering = FALSE do
42//!     wait := TRUE;
43//!     Initialize a new instance luc.round of logged uniform consensus;
44//!     proposals[round] := unordered \ delivered;
45//!     store(proposals);
46//!     trigger ⟨ luc.round, Propose | proposals[round] ⟩;
47//!
48//! upon event ⟨ luc.r, Decide | decided ⟩ such that r = round do
49//!     forall (s, m) ∈ sort(decided) do     // by the order in the resulting sorted list
50//!         append(delivered, (s, m));
51//!     store(delivered);
52//!     round := round + 1;
53//!     if recovering = TRUE then
54//!         if proposals[round] ≠ ⊥ then
55//!             trigger ⟨ luc.round, Propose | proposals[round] ⟩;
56//!         else
57//!             recovering := FALSE;
58//!     else
59//!         wait := FALSE;
60//!     trigger ⟨ lutob, Deliver | delivered ⟩;
61//! ```
62//!
63//! The pair to [`crate::consensus_based_total_order_broadcast`], and held to the same suite. What
64//! differs is exactly one thing: the ordered sequence survives a restart. Everything else — one
65//! consensus instance per round, propose the unordered set, sort what is decided — is the same
66//! shape, which is what makes the comparison worth having.
67//!
68//! **Why the recovery is not simply "read it back".** A process that proposed for a round and then
69//! forgot would, on recovering, propose something *different* for the same round — and a uniform
70//! consensus that has already decided cannot accommodate it. So `proposals[r]` is durable before the
71//! proposal is visible to anyone, and recovery re-proposes what was recorded, round by round, until
72//! it reaches one it never proposed for. That is what `recovering` is counting through.
73//!
74//! # Departures from the page
75//!
76//! - **A read**, as the port requires. See [`crate::total_order_log`].
77//!
78//! - **Consensus instances are a family, and the conditional event handler is discharged here.**
79//!   Both as in the crash-stop member, and for the same reasons; that module's header states them.
80//!   Unlike that member, this one runs over a consensus assuming no synchrony, so processes can
81//!   genuinely drift and the family is doing work rather than standing on faithfulness alone.
82//!
83//! - **`delivered` and `proposals` are appended, not rewritten.** The page writes
84//!   `append(delivered, (s, m)); store(delivered)` and `store(proposals)` — rewriting a whole
85//!   growing structure on every change, which costs `O(n²)` bytes over a run. This repository's
86//!   storage interface splits the two cases so the choice is visible in the types, and
87//!   `docs/bounded-space.md` records both logged modules having had exactly this defect and losing
88//!   it. So both go into the appended sequence, one entry each, and recovery replays them.
89//!
90//! - **Consensus instances are re-created here, not by a runtime.** The page says so outright:
91//!   "During the recovery operation after a crash, the total-order algorithm runs again through all
92//!   rounds executed before the crash and executes the same consensus instances once more. *(We
93//!   assume that the runtime environment re-instantiates all instances of consensus that had been
94//!   dynamically initialized before the crash.)*" There is no such runtime here — a crash rebuilds
95//!   a process from its constructor and nothing else survives but storage — so `on_recovery`
96//!   re-creates every instance the durable record names and runs each one's own recovery, all of
97//!   them before any decision is acted on: a decided instance announces its decision again from its
98//!   record, and an undecided one must have read its state back before recovery re-proposes into
99//!   it. The same shape of departure as the conditional event handler: a facility the book assumes,
100//!   discharged in the module.
101//!
102//! - **The decided prefix is replayed from the record, not re-decided.** The page rebuilds
103//!   `delivered` by running every round again, which is what its runtime's re-instantiated
104//!   instances are for. The appended `Record::Ordered` entries already hold the sequence in order,
105//!   so recovery replays them directly and the walk over re-announced decisions advances `round`
106//!   without appending or announcing anything twice — the same guard that makes duplicate decisions
107//!   harmless in a live run makes the replay idempotent. Each replayed entry *is* announced again,
108//!   once, for the reason [`crate::logged_leader_driven_consensus`] gives for re-announcing a
109//!   decision: the layer above may have crashed with this process and never seen the first
110//!   indication. Positions make the re-announcement idempotent for a client.
111//!
112//! - **The child that appends is composed through a sequence slot.** `logged_uniform_reliable_
113//!   broadcast` keeps an appended record of its own, and until this module nothing composed over a
114//!   child that appends — `store.rs` said so, and said what the missing half would be. This is its
115//!   second consumer, and [`recon_core::SeqSlot`] is what that paragraph described. Parent and child
116//!   append into **one** sequence, so the order between their entries is real rather than invented
117//!   at recovery.
118
119use recon_core::{Child, KeyedSlot, NodeId, Position, ProtoCx, Protocol, SeqSlot, Slot, TimerId};
120use serde::{Deserialize, Serialize};
121use std::collections::{BTreeMap, BTreeSet};
122
123use crate::Timing;
124use crate::consensus_based_total_order_broadcast::{Batch, Slot as OrderedSlot};
125use crate::logged_leader_driven_consensus::{self as luc, LoggedLeaderDrivenConsensus};
126use crate::logged_uniform_reliable_broadcast::{self as lurb, LoggedUniformReliableBroadcast};
127use crate::total_order_log::{LogInd, TotalOrderLog};
128
129/// The consensus one round runs. Not a type parameter, for the reason the crash-stop member gives.
130pub type Consensus<V> = LoggedLeaderDrivenConsensus<Batch<V>>;
131
132/// This layer's messages: the broadcast's, and a consensus instance's stamped with its round.
133#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
134pub enum Wire<B, C> {
135    Broadcast(B),
136    /// Stamped with the round whose instance it belongs to — the round is this layer's concept, not
137    /// the consensus's, so this layer stamps.
138    Consensus {
139        round: u64,
140        msg: C,
141    },
142}
143
144/// What this protocol appends. One sequence, carrying its own entries and its child's.
145///
146/// The page's `store(delivered)` and `store(proposals)` rewrite whole growing structures; these are
147/// appended instead, which is the departure the module header records.
148#[derive(Debug, Clone, PartialEq, Eq)]
149pub enum Record<V: Clone + Ord> {
150    /// One entry taking its place in the agreed sequence.
151    Ordered(OrderedSlot<V>),
152    /// What this process proposed for a round, durable before the proposal was visible.
153    Proposed { round: u64, batch: Batch<V> },
154    /// The reliable broadcast's own record, in this sequence rather than a second one.
155    Broadcast(lurb::Record<V>),
156}
157
158/// This protocol's rewritten metadata, and its children's inside it.
159///
160/// The consensus half is a **family**: one record per round, since one instance per round keeps its
161/// own. That is what [`recon_core::KeyedSlot`] is for — a place named as a function of a key, with
162/// the key supplied as data so the slot is still one fixed function.
163#[derive(Debug, Clone, PartialEq, Eq)]
164pub struct Durable<V: Clone + Ord> {
165    /// The broadcast's record. `()` — it writes once so a restart finds something.
166    broadcast: Option<()>,
167    /// Each round's consensus instance's record.
168    rounds: BTreeMap<u64, luc::Durable<Batch<V>>>,
169}
170
171impl<V: Clone + Ord> Default for Durable<V> {
172    fn default() -> Self {
173        Durable { broadcast: None, rounds: BTreeMap::new() }
174    }
175}
176
177/// Requests from the layer above.
178#[derive(Debug, Clone, PartialEq, Eq)]
179pub enum Cmd<V> {
180    /// `⟨ lutob, Broadcast | m ⟩`, in the port's terms.
181    Append(V),
182    Read {
183        from: Position,
184    },
185}
186
187/// Indications to the layer above.
188#[derive(Debug, Clone, PartialEq, Eq)]
189pub enum Ind<V> {
190    /// One entry of `⟨ lutob, Deliver | delivered ⟩`, with the position it took.
191    ///
192    /// The page hands up the whole list each round; the port's vocabulary is one entry at a time,
193    /// and the list is [`LoggedUniformTotalOrderBroadcast::entries`].
194    Ordered {
195        position: Position,
196        from: NodeId,
197        value: V,
198    },
199    Contents {
200        from: Position,
201        entries: Vec<V>,
202    },
203}
204
205/// A totally ordered log whose sequence survives a restart.
206pub struct LoggedUniformTotalOrderBroadcast<V: Clone + Ord> {
207    me: NodeId,
208    peers: BTreeSet<NodeId>,
209    timing: Timing,
210    /// `unordered`.
211    unordered: Batch<V>,
212    /// `delivered`, as the agreed sequence.
213    delivered: Vec<OrderedSlot<V>>,
214    ordered: BTreeSet<OrderedSlot<V>>,
215    /// `proposals[r]`, durable.
216    proposals: BTreeMap<u64, Batch<V>>,
217    /// `round`.
218    round: u64,
219    /// `wait`.
220    wait: bool,
221    /// `recovering`.
222    recovering: bool,
223    lurb: Child<LoggedUniformReliableBroadcast<V>>,
224    consensus: BTreeMap<u64, Child<Consensus<V>>>,
225    decisions: BTreeMap<u64, Batch<V>>,
226}
227
228fn broadcast_slot<V: Clone + Ord>() -> Slot<Durable<V>, ()> {
229    Slot {
230        read: |d| d.broadcast.as_ref(),
231        write: |d, c| {
232            let mut whole = d.cloned().unwrap_or_default();
233            whole.broadcast = Some(c);
234            whole
235        },
236    }
237}
238
239/// One round's consensus record, inside this protocol's. The round is the key.
240fn round_slot<V: Clone + Ord>() -> KeyedSlot<Durable<V>, luc::Durable<Batch<V>>, u64> {
241    KeyedSlot {
242        read: |d, r| d.rounds.get(r),
243        write: |d, r, c| {
244            let mut whole = d.cloned().unwrap_or_default();
245            whole.rounds.insert(*r, c);
246            whole
247        },
248    }
249}
250
251fn broadcast_entries<V: Clone + Ord>() -> SeqSlot<Record<V>, lurb::Record<V>> {
252    SeqSlot {
253        wrap: Record::Broadcast,
254        project: |r| match r {
255            Record::Broadcast(inner) => Some(inner),
256            _ => None,
257        },
258    }
259}
260
261impl<V: Clone + Ord> LoggedUniformTotalOrderBroadcast<V> {
262    /// A totally ordered log among `peers` whose sequence survives a restart.
263    pub fn new(me: NodeId, peers: impl IntoIterator<Item = NodeId>, timing: Timing) -> Self {
264        let mut peers: BTreeSet<NodeId> = peers.into_iter().collect();
265        peers.insert(me);
266        LoggedUniformTotalOrderBroadcast {
267            me,
268            peers: peers.clone(),
269            timing,
270            unordered: BTreeSet::new(),
271            delivered: Vec::new(),
272            ordered: BTreeSet::new(),
273            proposals: BTreeMap::new(),
274            round: 1,
275            wait: false,
276            recovering: false,
277            lurb: Child::new(LoggedUniformReliableBroadcast::new(me, peers, timing.retransmit)),
278            consensus: BTreeMap::new(),
279            decisions: BTreeMap::new(),
280        }
281    }
282
283    /// The agreed sequence as this process holds it.
284    pub fn entries(&self) -> &[OrderedSlot<V>] {
285        &self.delivered
286    }
287
288    pub fn len(&self) -> usize {
289        self.delivered.len()
290    }
291
292    pub fn is_empty(&self) -> bool {
293        self.delivered.is_empty()
294    }
295
296    pub fn round(&self) -> u64 {
297        self.round
298    }
299
300    pub fn instances(&self) -> usize {
301        self.consensus.len()
302    }
303
304    /// `upon unordered \ delivered ≠ ∅ ∧ wait = FALSE ∧ recovering = FALSE do`.
305    fn maybe_propose(&mut self, cx: &mut ProtoCx<'_, Self>) {
306        if self.wait || self.recovering {
307            return;
308        }
309        let batch: Batch<V> =
310            self.unordered.difference(&self.ordered).cloned().collect::<BTreeSet<_>>();
311        if batch.is_empty() {
312            return;
313        }
314        self.wait = true;
315        let round = self.round;
316        self.proposals.insert(round, batch.clone());
317        // `store(proposals)` — durable **before** the proposal is visible to anyone. A process that
318        // proposed and then forgot would propose something different for the same round on
319        // recovering, and a decided uniform consensus cannot accommodate that.
320        cx.storage().append(Record::Proposed { round, batch: batch.clone() });
321        self.through_consensus(round, cx, |c, ccx| c.on_cmd(luc::Cmd::Propose(batch), ccx));
322    }
323
324    /// The decide handler, with the buffering the book's run-time system would have done.
325    fn drain_decisions(&mut self, cx: &mut ProtoCx<'_, Self>) {
326        while let Some(decided) = self.decisions.remove(&self.round) {
327            for (from, value) in &decided {
328                if self.ordered.contains(&(*from, value.clone())) {
329                    continue;
330                }
331                let position = Position(self.delivered.len() as u64);
332                self.delivered.push((*from, value.clone()));
333                self.ordered.insert((*from, value.clone()));
334                // `append(delivered, (s, m)); store(delivered)` — appended rather than rewritten.
335                cx.storage().append(Record::Ordered((*from, value.clone())));
336                cx.indicate(Ind::Ordered { position, from: *from, value: value.clone() });
337            }
338            self.round += 1;
339            if self.recovering {
340                // Re-propose what was recorded for the next round, or stop recovering when this
341                // process never got that far.
342                match self.proposals.get(&self.round).cloned() {
343                    Some(batch) => {
344                        let round = self.round;
345                        self.through_consensus(round, cx, |c, ccx| {
346                            c.on_cmd(luc::Cmd::Propose(batch), ccx)
347                        });
348                    }
349                    None => {
350                        self.recovering = false;
351                        self.wait = false;
352                    }
353                }
354            } else {
355                self.wait = false;
356            }
357        }
358        self.maybe_propose(cx);
359    }
360
361    fn through_lurb(
362        &mut self,
363        cx: &mut ProtoCx<'_, Self>,
364        f: impl FnOnce(
365            &mut LoggedUniformReliableBroadcast<V>,
366            &mut ProtoCx<'_, LoggedUniformReliableBroadcast<V>>,
367        ),
368    ) {
369        let mut inds = self.lurb.run_appending(
370            cx,
371            Wire::Broadcast,
372            broadcast_slot::<V>(),
373            broadcast_entries::<V>(),
374            f,
375        );
376        for lurb::Ind::Delivered(log) in inds.drain(..) {
377            // `upon event ⟨ lurb, Deliver | lurbdelivered ⟩ do unordered := unordered ∪ lurbdelivered`
378            for (id, msg) in log.delivered() {
379                self.unordered.insert((id.origin, msg.clone()));
380            }
381        }
382        self.lurb.reclaim(inds);
383        self.maybe_propose(cx);
384    }
385
386    fn through_consensus(
387        &mut self,
388        r: u64,
389        cx: &mut ProtoCx<'_, Self>,
390        f: impl FnOnce(&mut Consensus<V>, &mut ProtoCx<'_, Consensus<V>>),
391    ) {
392        // "Initialize a new instance luc.round" — an event, run before whatever provoked the
393        // creation; the crash-stop member's departures say what skipping it cost. An instance
394        // re-created by `on_recovery` takes the recovery branch there instead, and this path then
395        // finds it present.
396        let created = !self.consensus.contains_key(&r);
397        self.consensus.entry(r).or_insert_with(|| {
398            Child::new(LoggedLeaderDrivenConsensus::new(self.me, self.peers.clone(), self.timing))
399        });
400        let child = self.consensus.get_mut(&r).expect("just inserted");
401        let mut inds = child.run_keyed(
402            cx,
403            move |m| Wire::Consensus { round: r, msg: m },
404            round_slot::<V>(),
405            r,
406            |c, ccx| {
407                if created {
408                    c.on_init(ccx);
409                }
410                f(c, ccx)
411            },
412        );
413        for luc::Ind::Decide(decided) in inds.drain(..) {
414            self.decisions.entry(r).or_insert(decided);
415        }
416        if let Some(child) = self.consensus.get_mut(&r) {
417            child.reclaim(inds);
418        }
419        self.drain_decisions(cx);
420    }
421}
422
423impl<V: Clone + Ord> Protocol for LoggedUniformTotalOrderBroadcast<V> {
424    type Cmd = Cmd<V>;
425    type Ind = Ind<V>;
426    type Msg = Wire<<LoggedUniformReliableBroadcast<V> as Protocol>::Msg, luc::Wire<Batch<V>>>;
427    type Scope = core::convert::Infallible;
428    type Note = crate::Note;
429    type Meta = Durable<V>;
430    type Entry = Record<V>;
431
432    fn on_cmd(&mut self, cmd: Cmd<V>, cx: &mut ProtoCx<'_, Self>) {
433        match cmd {
434            Cmd::Append(v) => {
435                self.through_lurb(cx, |b, ccx| b.on_cmd(lurb::Cmd::Broadcast(v), ccx));
436            }
437            Cmd::Read { from } => {
438                let entries: Vec<V> =
439                    self.delivered.iter().skip(from.0 as usize).map(|(_, v)| v.clone()).collect();
440                cx.indicate(Ind::Contents { from, entries });
441            }
442        }
443    }
444
445    fn on_msg(&mut self, from: NodeId, msg: Self::Msg, cx: &mut ProtoCx<'_, Self>) {
446        match msg {
447            Wire::Broadcast(m) => self.through_lurb(cx, |b, ccx| b.on_msg(from, m, ccx)),
448            Wire::Consensus { round, msg } => {
449                self.through_consensus(round, cx, |c, ccx| c.on_msg(from, msg, ccx));
450            }
451        }
452    }
453
454    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
455        self.through_lurb(cx, |b, ccx| b.on_timer(id, ccx));
456        let rounds: Vec<u64> = self.consensus.keys().copied().collect();
457        for r in rounds {
458            self.through_consensus(r, cx, |c, ccx| c.on_timer(id, ccx));
459        }
460    }
461
462    fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
463        self.through_lurb(cx, |b, ccx| b.on_init(ccx));
464    }
465
466    /// `upon event ⟨ Recovery ⟩`, with the two facilities the page assumes discharged here: the
467    /// runtime that re-instantiates consensus instances, and the buffering behind `such that`.
468    fn on_recovery(&mut self, cx: &mut ProtoCx<'_, Self>) {
469        self.recovering = true;
470
471        // `retrieve(proposals)` — and the ordered entries, which the page rebuilds by re-running
472        // every round and this module replays from its appended record instead; the departures say
473        // why, and why each replayed entry is announced again.
474        let records: Vec<Record<V>> =
475            cx.storage().read_from(Position::START).into_iter().cloned().collect();
476        for record in records {
477            match record {
478                Record::Ordered((from, value)) => {
479                    if self.ordered.contains(&(from, value.clone())) {
480                        continue;
481                    }
482                    let position = Position(self.delivered.len() as u64);
483                    self.delivered.push((from, value.clone()));
484                    self.ordered.insert((from, value.clone()));
485                    cx.indicate(Ind::Ordered { position, from, value });
486                }
487                Record::Proposed { round, batch } => {
488                    self.proposals.insert(round, batch);
489                }
490                // The child's, replayed by the child below through its own filtered view.
491                Record::Broadcast(_) => {}
492            }
493        }
494
495        // The broadcast re-announces its log — which rebuilds `unordered` — and re-sends what was
496        // still pending. `recovering` holds the proposal this would otherwise trigger.
497        self.through_lurb(cx, |b, ccx| b.on_recovery(ccx));
498
499        // Re-instantiate every consensus instance the durable record names, and recover all of
500        // them before acting on any decision. Not `through_consensus`: that drains after each
501        // instance, and the walk's re-proposal must never reach an instance that has not read its
502        // record back — a process that proposed and then forgot is the exact failure the durable
503        // proposal exists to prevent.
504        let rounds: Vec<u64> =
505            cx.storage().get().map(|d| d.rounds.keys().copied().collect()).unwrap_or_default();
506        for r in rounds {
507            self.consensus.entry(r).or_insert_with(|| {
508                Child::new(LoggedLeaderDrivenConsensus::new(
509                    self.me,
510                    self.peers.clone(),
511                    self.timing,
512                ))
513            });
514            let child = self.consensus.get_mut(&r).expect("just inserted");
515            let mut inds = child.run_keyed(
516                cx,
517                move |m| Wire::Consensus { round: r, msg: m },
518                round_slot::<V>(),
519                r,
520                |c, ccx| c.on_recovery(ccx),
521            );
522            for luc::Ind::Decide(decided) in inds.drain(..) {
523                self.decisions.entry(r).or_insert(decided);
524            }
525            if let Some(child) = self.consensus.get_mut(&r) {
526                child.reclaim(inds);
527            }
528        }
529
530        // Walk the re-announced decisions forward. Replayed entries are already in `ordered`, so
531        // the walk advances `round` without appending or announcing anything twice, and its
532        // recovering branch re-proposes for a round that was proposed and never decided.
533        let before = self.round;
534        self.drain_decisions(cx);
535
536        // The page's own Recovery handler: `if proposals[1] ≠ ⊥ then trigger ⟨ luc.1, Propose ⟩`.
537        // Needed only when the walk did not run — no round had decided — since the walk's
538        // recovering branch otherwise made this same choice at the round it stopped at. No recorded
539        // proposal means the crash landed before this process proposed anything still undecided,
540        // and recovery is over.
541        if self.recovering && self.round == before {
542            match self.proposals.get(&self.round).cloned() {
543                Some(batch) => {
544                    let round = self.round;
545                    self.through_consensus(round, cx, |c, ccx| {
546                        c.on_cmd(luc::Cmd::Propose(batch), ccx)
547                    });
548                }
549                None => {
550                    self.recovering = false;
551                    self.wait = false;
552                    self.maybe_propose(cx);
553                }
554            }
555        }
556    }
557}
558
559impl<V: Clone + Ord> TotalOrderLog<V> for LoggedUniformTotalOrderBroadcast<V> {
560    fn append(value: V) -> Cmd<V> {
561        Cmd::Append(value)
562    }
563
564    fn read(from: Position) -> Cmd<V> {
565        Cmd::Read { from }
566    }
567
568    fn classify(ind: Ind<V>) -> LogInd<V> {
569        match ind {
570            Ind::Ordered { position, from, value } => LogInd::Ordered { position, from, value },
571            Ind::Contents { from, entries } => LogInd::Contents { from, entries },
572        }
573    }
574}