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}