Skip to main content

recon_protocols/
logged_uniform_reliable_broadcast.rs

1//! Logged uniform reliable broadcast.
2//!
3//! Cachin, Guerraoui & Rodrigues, Module 3.6 and Algorithm 3.8 ("Logged Majority-Ack Uniform
4//! Reliable Broadcast").
5//!
6//! **Status: transcription. Space: unbounded — and on disk. Write cost: a fixed number of appends
7//! per message.** `pending` and `delivered` grow with every message handled, as in
8//! [`crate::majority_ack_uniform_reliable_broadcast`], except that here they are written down.
9//! Each is recorded by appending one entry, so the cost of recording a message does not depend on
10//! how many preceded it; the record itself is still unbounded. See `docs/bounded-space.md`.
11//!
12//! **Assumption: a correct majority, `N > 2f`** — where in this model a *correct* process is one
13//! that always recovers from its crashes, and what it knows after recovering is what it wrote.
14//!
15//! ```text
16//! upon event ⟨ lurb, Init ⟩ do
17//!     delivered := ∅; pending := ∅;
18//!     forall m do ack[m] := ∅;
19//!     store(pending, delivered);
20//!
21//! upon event ⟨ lurb, Recovery ⟩ do
22//!     retrieve(pending, delivered);
23//!     trigger ⟨ lurb, Deliver | delivered ⟩;
24//!     forall (s, m) ∈ pending do
25//!         trigger ⟨ sbeb, Broadcast | [DATA, s, m] ⟩;
26//!
27//! upon event ⟨ lurb, Broadcast | m ⟩ do
28//!     pending := pending ∪ {(self, m)};
29//!     store(pending);
30//!     trigger ⟨ sbeb, Broadcast | [DATA, self, m] ⟩;
31//!
32//! upon event ⟨ sbeb, Deliver | p, [DATA, s, m] ⟩ do
33//!     if (s, m) ∉ pending then
34//!         pending := pending ∪ {(s, m)};
35//!         store(pending);
36//!         trigger ⟨ sbeb, Broadcast | [DATA, s, m] ⟩;
37//!     if p ∉ ack[m] then
38//!         ack[m] := ack[m] ∪ {p};
39//!         if #(ack[m]) > N/2 ∧ (s, m) ∉ delivered then
40//!             delivered := delivered ∪ {(s, m)};
41//!             store(delivered);
42//!             trigger ⟨ lurb, Deliver | delivered ⟩;
43//! ```
44//!
45//! # `ack` is deliberately not durable
46//!
47//! The book: *"Variable `ack` is not logged because it will be reconstructed upon recovery."*
48//! Getting this wrong in the direction of storing more looks safer and is worse — it would cost a
49//! write per acknowledgement to save work that retransmission does anyway, and would make the
50//! durable state grow with *traffic* rather than with messages.
51//!
52//! What rebuilds it is the recovery clause: a recovered process re-broadcasts everything still
53//! pending, and acknowledgements accumulate as the answers arrive. The stubborn broadcast beneath
54//! never stops retransmitting, so a process that was down when something was sent gets it anyway.
55//!
56//! # Why the child is stubborn broadcast and not the logged link
57//!
58//! The book's logged abstractions do not stack: Algorithm 2.3 is over stubborn links, and this one is
59//! over stubborn *broadcast*. Each keeps its own log. A perfect link's deduplication is volatile,
60//! so after a restart it would re-deliver anyway — the deduplication buys nothing a logged layer
61//! above does not already do for itself, and it is the **retransmission** a recovered process
62//! needs. Deduplicating beneath would suppress exactly that.
63//!
64//! # Departures from the page
65//!
66//! - `⟨ Init ⟩` is here and performs the book's initial store, as [`crate::logged_link`] does.
67//! - `pending` and `delivered` share one appended sequence, distinguished by a tag on each entry,
68//!   rather than the book's two stores. Recovery replays the sequence and rebuilds both.
69//! - Messages are keyed by an identifier carrying the originator and a sequence number rather
70//!   than by content, so identical content broadcast twice is delivered twice.
71//! - The sequence number is therefore as durable as the log it keys, and recovery recomputes it
72//!   rather than resuming from zero. See the note below; the obligation is part of the departure.
73//! - No `Stop`: retransmission for ever is what reaches a recovered process.
74//!
75//! # The sequence number survives a crash, without being written
76//!
77//! Content-keyed identity, as the book has it, needs nothing restored: a payload names itself. An
78//! id-keyed departure owes the counter the same durability as the set it keys, or a recovered
79//! process re-mints `(me, 1)` for something *new* and two distinct payloads collide under one
80//! identifier — last-write-wins in `pending`, acks for the old counting toward the new, at most
81//! one of the two ever delivered locally, and different processes log-delivering different
82//! payloads under the same id. No-creation, validity and uniform agreement all fail, in the one
83//! module whose model expects a recovered process to keep working.
84//!
85//! Recovery recomputes it as the greatest `seq` over replayed records originating here, which is
86//! sound because every own broadcast appends its `Pending` record in the handler that emits it:
87//! a torn write discards that handler's sends, so no broadcast escapes without a record. The
88//! alternative — the counter in `Meta` — is also correct and costs a metadata write per
89//! broadcast, which is what buys nothing here.
90
91use core::time::Duration;
92use recon_core::{Child, NodeId, Position, ProtoCx, Protocol, TimerId};
93use serde::{Deserialize, Serialize};
94use std::collections::{BTreeMap, BTreeSet};
95
96use crate::stubborn_broadcast::{self as sbeb, StubbornBroadcast};
97use crate::uniform_reliable_broadcast::{BroadcastId, Data};
98
99/// What survives a crash: what has been seen, and what has been log-delivered.
100///
101/// `ack` is absent, deliberately. See the module note.
102#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
103#[serde(bound(deserialize = "P: Ord + Deserialize<'de>"))]
104pub struct Logged<P: Ord> {
105    pending: BTreeMap<BroadcastId, P>,
106    delivered: BTreeSet<(BroadcastId, P)>,
107}
108
109impl<P: Ord> Default for Logged<P> {
110    fn default() -> Self {
111        Logged { pending: BTreeMap::new(), delivered: BTreeSet::new() }
112    }
113}
114
115impl<P: Ord> Logged<P> {
116    /// Everything log-delivered, with the identifier naming its originator.
117    pub fn delivered(&self) -> impl Iterator<Item = &(BroadcastId, P)> + '_ {
118        self.delivered.iter()
119    }
120
121    /// How many messages have been log-delivered.
122    pub fn delivered_count(&self) -> usize {
123        self.delivered.len()
124    }
125
126    /// How many messages have been seen and are still being re-broadcast.
127    pub fn pending_count(&self) -> usize {
128        self.pending.len()
129    }
130}
131
132/// One thing written down. Replaying these in order rebuilds `pending` and `delivered`.
133#[derive(Debug, Clone, PartialEq, Eq)]
134pub enum Record<P> {
135    /// Seen, and being re-broadcast until a majority has it.
136    Pending(BroadcastId, P),
137    /// Log-delivered.
138    Delivered(BroadcastId, P),
139}
140
141/// Requests from the layer above.
142#[derive(Debug, Clone, PartialEq, Eq)]
143pub enum Cmd<P> {
144    Broadcast(P),
145}
146
147/// Indications to the layer above: the durable log, not a message.
148#[derive(Debug, Clone, PartialEq, Eq)]
149pub enum Ind<P: Ord> {
150    Delivered(Logged<P>),
151}
152
153/// Uniform agreement over log-delivery, in the fail-recovery model.
154#[derive(Debug)]
155pub struct LoggedUniformReliableBroadcast<P: Clone + Ord> {
156    me: NodeId,
157    seq: u64,
158    members: usize,
159    /// Durable. Written on every change and retrieved on recovery.
160    log: Logged<P>,
161    /// Volatile, and rebuilt by re-broadcasting on recovery. Not written down.
162    ack: BTreeMap<BroadcastId, BTreeSet<NodeId>>,
163    /// Names each fan-out to the stubborn broadcast beneath. Volatile, and matched to the child's
164    /// own volatile state — a crash takes both, so nothing outlives the counter that keys it.
165    /// Not the durable [`BroadcastId`]: nothing here is ever stopped, so the two never meet.
166    beb_seq: u64,
167    beb: Child<StubbornBroadcast<Data<P>>>,
168}
169
170impl<P: Clone + Ord> LoggedUniformReliableBroadcast<P> {
171    /// Broadcast among `members`, which must include `me`, retransmitting every `interval`.
172    pub fn new(me: NodeId, members: impl IntoIterator<Item = NodeId>, interval: Duration) -> Self {
173        let mut members: BTreeSet<NodeId> = members.into_iter().collect();
174        members.insert(me);
175        let n = members.len();
176        LoggedUniformReliableBroadcast {
177            me,
178            seq: 0,
179            members: n,
180            log: Logged::default(),
181            ack: BTreeMap::new(),
182            beb_seq: 0,
183            beb: Child::new(StubbornBroadcast::new(me, members, interval)),
184        }
185    }
186
187    /// The durable log. The same value the layer above is handed.
188    pub fn log(&self) -> &Logged<P> {
189        &self.log
190    }
191
192    /// How many processes must have re-broadcast a message before it is log-delivered.
193    pub fn majority(&self) -> usize {
194        self.members / 2 + 1
195    }
196
197    /// Which processes have been seen to re-broadcast `id`. Volatile, and rebuilt on recovery.
198    pub fn acknowledged_by(&self, id: BroadcastId) -> impl Iterator<Item = NodeId> + '_ {
199        self.ack.get(&id).into_iter().flatten().copied()
200    }
201}
202
203impl<P: Clone + Ord> LoggedUniformReliableBroadcast<P> {
204    fn with_beb(
205        &mut self,
206        cx: &mut ProtoCx<'_, Self>,
207        f: impl FnOnce(&mut StubbornBroadcast<Data<P>>, &mut ProtoCx<'_, StubbornBroadcast<Data<P>>>),
208    ) {
209        let mut inds = self.beb.run(cx, core::convert::identity, f);
210        for sbeb::Ind::Deliver { from, msg } in inds.drain(..) {
211            self.on_arrival(from, msg, cx);
212        }
213        self.beb.reclaim(inds);
214    }
215
216    fn rebroadcast(&mut self, data: Data<P>, cx: &mut ProtoCx<'_, Self>) {
217        self.beb_seq += 1;
218        let id = sbeb::BroadcastId(self.beb_seq);
219        // Re-enters the child while its inbox is out on loan, so `run` hands back a fresh one.
220        let inds = self.beb.run(cx, core::convert::identity, |beb, ccx| {
221            beb.on_cmd(sbeb::Cmd::Broadcast { id, msg: data }, ccx)
222        });
223        debug_assert!(inds.is_empty(), "broadcasting must not deliver synchronously");
224        self.beb.reclaim(inds);
225    }
226
227    /// `upon event ⟨ sbeb, Deliver | p, [DATA, s, m] ⟩`.
228    ///
229    /// Arrives many times for one broadcast — the child never stops retransmitting — so every
230    /// branch here is guarded and idempotent.
231    fn on_arrival(&mut self, from: NodeId, data: Data<P>, cx: &mut ProtoCx<'_, Self>) {
232        let id = data.id;
233
234        // Re-broadcast only on first sight. An identifier determines its payload, so re-inserting
235        // cannot change what is pending — the returned Option is read only to learn whether this
236        // was the first time.
237        if self.log.pending.insert(id, data.payload.clone()).is_none() {
238            // Durable at the point of insertion, in this handler's own text: the re-broadcast is
239            // this process's acknowledgement, and under an eager sink it escapes before the
240            // handler returns. Everything in `pending` has a durable `Pending` record by
241            // construction, not by where control happens to leave.
242            cx.storage().append(Record::Pending(id, data.payload.clone()));
243            self.rebroadcast(data.clone(), cx);
244        }
245
246        if self.ack.entry(id).or_default().insert(from) {
247            let acked = self.ack.get(&id).map(|a| a.len()).unwrap_or(0);
248            let already = self.log.delivered.iter().any(|(i, _)| *i == id);
249            if 2 * acked > self.members && !already {
250                self.log.delivered.insert((id, data.payload.clone()));
251                // Durable before announced: the layer above must not learn of a delivery a crash
252                // could erase.
253                cx.storage().append(Record::Delivered(id, data.payload));
254                cx.indicate(Ind::Delivered(self.log.clone()));
255            }
256        }
257    }
258}
259
260impl<P: Clone + Ord> Protocol for LoggedUniformReliableBroadcast<P> {
261    type Cmd = Cmd<P>;
262    type Ind = Ind<P>;
263    type Msg = Data<P>;
264    type Scope = core::convert::Infallible;
265    type Note = crate::Note;
266    /// Nothing is rewritten; the metadata is written once so a restart finds something.
267    type Meta = ();
268    /// One record per message seen or log-delivered. `ack` is not among them, by design.
269    type Entry = Record<P>;
270
271    /// `⟨ lurb, Init ⟩ do delivered := ∅; pending := ∅; ...; store(pending, delivered)`.
272    fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
273        cx.storage().set(());
274    }
275
276    fn on_cmd(&mut self, Cmd::Broadcast(msg): Cmd<P>, cx: &mut ProtoCx<'_, Self>) {
277        self.seq += 1;
278        let id = BroadcastId { origin: self.me, seq: self.seq };
279        self.log.pending.insert(id, msg.clone());
280        cx.storage().append(Record::Pending(id, msg.clone()));
281        self.rebroadcast(Data { id, payload: msg }, cx);
282    }
283
284    fn on_msg(&mut self, from: NodeId, msg: Data<P>, cx: &mut ProtoCx<'_, Self>) {
285        self.with_beb(cx, |beb, ccx| beb.on_msg(from, msg, ccx));
286    }
287
288    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
289        self.with_beb(cx, |beb, ccx| beb.on_timer(id, ccx));
290    }
291
292    /// `upon event ⟨ lurb, Recovery ⟩`.
293    ///
294    /// Re-announce the log, then re-broadcast everything pending — which is what rebuilds `ack`,
295    /// and why it need not be durable. The replay also restores the send counter, so a process
296    /// that goes on to broadcast something new does not reuse an identifier.
297    fn on_recovery(&mut self, cx: &mut ProtoCx<'_, Self>) {
298        // Nothing else is dispatched until this returns, so an empty index a moment ago is safe.
299        let records: Vec<Record<P>> =
300            cx.storage().read_from(Position::START).into_iter().cloned().collect();
301        for r in records {
302            let id = match r {
303                Record::Pending(id, p) => {
304                    self.log.pending.insert(id, p);
305                    id
306                }
307                Record::Delivered(id, p) => {
308                    self.log.delivered.insert((id, p));
309                    id
310                }
311            };
312            // The counter is as durable as the set it keys. Resuming from zero would re-mint an
313            // identifier already in use for a different payload; see the module note.
314            if id.origin == self.me {
315                self.seq = self.seq.max(id.seq);
316            }
317        }
318        cx.indicate(Ind::Delivered(self.log.clone()));
319        let outstanding: Vec<Data<P>> = self
320            .log
321            .pending
322            .iter()
323            .map(|(id, payload)| Data { id: *id, payload: payload.clone() })
324            .collect();
325        for data in outstanding {
326            self.rebroadcast(data, cx);
327        }
328    }
329}