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}