recon_protocols/probabilistic_broadcast.rs
1//! Eager probabilistic broadcast — gossip.
2//!
3//! **Status: implementation. Space: bounded by a retention window.** The first module above the
4//! failure detector that is not a transcription; see the guarantee table below for what the window
5//! costs.
6//!
7//! Cachin, Guerraoui & Rodrigues, Module 3.7 and Algorithm 3.9 ("Eager Probabilistic Broadcast").
8//! Quoted from the book rather than from memory, which matters here more than anywhere else in this
9//! repository: `docs/postmortem.md` records four bugs once reported in the previous implementation
10//! of this algorithm, of which **three were false positives** produced by reading code against
11//! remembered pseudocode.
12//!
13//! ```text
14//! upon event ⟨ pb, Init ⟩ do
15//! delivered := ∅;
16//!
17//! procedure gossip(msg) is
18//! forall t ∈ picktargets(k) do
19//! trigger ⟨ fll, Send | t, msg ⟩;
20//!
21//! upon event ⟨ pb, Broadcast | m ⟩ do
22//! delivered := delivered ∪ {m};
23//! trigger ⟨ pb, Deliver | self, m ⟩;
24//! gossip([GOSSIP, self, m, R]);
25//!
26//! upon event ⟨ fll, Deliver | p, [GOSSIP, s, m, r] ⟩ do
27//! if m ∉ delivered then
28//! delivered := delivered ∪ {m};
29//! trigger ⟨ pb, Deliver | s, m ⟩;
30//! if r > 1 then gossip([GOSSIP, s, m, r − 1]);
31//! ```
32//!
33//! # The relay is outside the delivery guard, and that is the book's
34//!
35//! Read the indentation. `if r > 1 then gossip(...)` sits at the same level as `if m ∉ delivered`,
36//! not inside it, so a process relays a message it has **already delivered**. This looks like a
37//! defect. It has been reported as one. It is not: the book names the consequence on the same page
38//! — "the algorithm induces a significant amount of redundancy in the message exchanges: any given
39//! process may receive the same message many times" — and the redundancy is what makes the
40//! probability work out. Relaying only on first receipt would cut the fan-out of every message that
41//! reaches a process twice, which is most of them.
42//!
43//! # What is probabilistic, and what is not
44//!
45//! ```text
46//! PB1 [probabilistic] Probabilistic validity — a correct sender's message reaches every correct
47//! process with high probability, and on some runs it does not
48//! PB2 [window] No duplication — within the retention window; see below
49//! PB3 [always] No creation
50//! ```
51//!
52//! `PB1` is the whole point and the whole cost. Best-effort broadcast reaches everyone whenever the
53//! sender is correct; this does not, and buys in exchange that no process ever sends to all of `Π`.
54//! A run in which some correct process never delivers is **not a violation** — the suite counts
55//! such runs rather than failing on them, and asserts against a stated threshold.
56//!
57//! # Identity is an identifier, not the message
58//!
59//! **Departure.** The book deduplicates on the message itself — `m ∉ delivered` — which assumes
60//! messages are unique across senders. Every broadcast here instead carries a [`BroadcastId`]: its
61//! originator and a per-sender sequence number, exactly as `reliable_broadcast` does, so identical
62//! content broadcast twice is delivered twice. The consequence for this module is that `delivered`
63//! holds identifiers rather than payloads, which is also what makes the window below affordable.
64//!
65//! The sequence counter is volatile and so is the set it keys, which is the pairing
66//! `CLAUDE.md` requires: a durable set keyed by a volatile counter is the bug, and neither half is
67//! durable here.
68//!
69//! **The identifier also names the originator's incarnation.** A volatile counter restarts at one
70//! when its process does, and every *other* process's window is keyed on it — so without this, an
71//! originator that crashed and came back would have its first `window` broadcasts discarded
72//! everywhere as duplicates of ones it sent before. The incarnation is a value drawn from the seeded
73//! generator at `Init`: distinct across restarts with probability `1 − 2⁻⁶⁴` per pair, decided by
74//! the one process that knows it restarted, and needing no storage. A session boundary could not
75//! have done this job — a receiver's link is to a relayer, not to the originator, and a link cannot
76//! tell a reconnect from a restart in any case (`docs/conditional-guarantees.md`). Under the book's
77//! crash-stop model nobody restarts and the field is inert; it exists for the real-world set.
78//!
79//! # The retention window, which is this project's and not the book's
80//!
81//! Page 100: "garbage collection of the stored message copies is omitted in the pseudo code for
82//! simplicity." So there is no page to follow, and the mechanism is a design decision with its own
83//! cost. `delivered` keeps the most recent `window` identifiers **per sender** and evicts the
84//! oldest when that is exceeded, on insert.
85//!
86//! Two things follow, and both are deliberate:
87//!
88//! - **Reclaiming is constant work.** One eviction per insert, never a pass over the set. The
89//! previous implementation expired by wall-clock age and rebuilt the whole set on every event, so
90//! receiving one message cost time linear in everything ever received. That is the one defect in
91//! that code which survived scrutiny, and this is the shape that avoids it rather than a smaller
92//! version of it.
93//! - **`PB2` is scoped to the window.** A message re-arriving after its identifier has been evicted
94//! is delivered again. That is the stated guarantee, not a violation of it, and it is why the
95//! table above says `[window]` where the book says nothing.
96
97use recon_core::{Child, NodeId, ProtoCx, Protocol, TimerId};
98use serde::{Deserialize, Serialize};
99use std::collections::{BTreeMap, BTreeSet, VecDeque};
100
101use crate::fair_loss_link::FairLossLink;
102use crate::link::{Boundary, LinkInd, VolatileLink};
103
104/// Which broadcast this is: who originated it, and their sequence number for it.
105///
106/// See the departure note above — the book keys on the message, this keys on an identifier.
107#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
108pub struct BroadcastId {
109 pub origin: NodeId,
110 /// Which incarnation of `origin` sent it. See the module note on identity.
111 pub incarnation: u64,
112 pub seq: u64,
113}
114
115/// What this layer puts on the wire: the identifier, the rounds still to live, and the payload.
116///
117/// `ttl` is the book's `r`. It is decremented at each hop and a message is not relayed once it
118/// reaches one, which is what makes a broadcast generate finitely many transmissions.
119#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
120pub struct Gossip<P> {
121 pub id: BroadcastId,
122 pub ttl: u32,
123 pub payload: P,
124}
125
126/// What a link beneath this layer must carry.
127pub type Carried<P> = Gossip<P>;
128
129/// Requests from the layer above.
130#[derive(Debug, Clone, PartialEq, Eq)]
131pub enum Cmd<P> {
132 Broadcast(P),
133}
134
135/// Indications to the layer above.
136#[derive(Debug, Clone, PartialEq, Eq)]
137pub enum Ind<P> {
138 /// `from` is the originator, never a relayer.
139 Deliver { from: NodeId, msg: P },
140 /// The scope with `peer` ended at `epoch`. Raised only over a link that reports boundaries.
141 ///
142 /// This layer cannot bridge one: it keeps identifiers rather than payloads, so it has nothing
143 /// to resend. It propagates, as `docs/conditional-guarantees.md` requires of a layer in that
144 /// position. A scope ending is in any case only one more way for a gossip to be lost, which is
145 /// a case this abstraction already tolerates by construction.
146 SessionEnded { peer: NodeId, epoch: u64 },
147 /// A scope with `peer` is in force at `epoch`.
148 SessionEstablished { peer: NodeId, epoch: u64 },
149}
150
151/// How this instance gossips.
152#[derive(Debug, Clone, Copy, PartialEq, Eq)]
153pub struct Config {
154 /// The book's `k` — how many peers a relay addresses. Must be smaller than the membership for
155 /// the abstraction to be doing anything; see [`Config::fanout`].
156 pub fanout: usize,
157 /// The book's `R` — how many hops a message travels before it stops being relayed.
158 pub rounds: u32,
159 /// How many identifiers to remember per sender. See the module note on the retention window.
160 pub window: usize,
161}
162
163impl Config {
164 /// A configuration, with the fanout and rounds the caller wants and a window large enough that
165 /// deduplication is not the thing under test.
166 pub fn new(fanout: usize, rounds: u32, window: usize) -> Self {
167 Config { fanout, rounds, window }
168 }
169}
170
171/// Gossip: relay to a random few, for a bounded number of rounds.
172///
173/// `L` is the link beneath and it is a parameter, so this composes over a perfect link, a session
174/// link, or an application's own. It bounds on [`crate::link::Link`] rather than anything narrower
175/// because gossip needs nothing of a scope boundary beyond passing it upward.
176#[derive(Debug)]
177pub struct ProbabilisticBroadcast<P: Clone, L: VolatileLink<Carried<P>> = FairLossLink<Gossip<P>>> {
178 me: NodeId,
179 /// Π — every process, including this one. The sender delivers to itself directly, as Algorithm
180 /// 3.9 has it, so `picktargets` draws from everyone else.
181 peers: BTreeSet<NodeId>,
182 /// This incarnation's name, drawn at `Init`. Zero until then, which a driver that never sends
183 /// `Init` — a test stepping the protocol by hand — will see.
184 incarnation: u64,
185 seq: u64,
186 config: Config,
187 /// Identifiers already delivered, and the order they arrived in per sender, so the oldest can
188 /// be evicted in constant time. Bounded by `config.window` per sender, across incarnations.
189 delivered: BTreeSet<BroadcastId>,
190 order: BTreeMap<NodeId, VecDeque<(u64, u64)>>,
191 link: Child<L>,
192 _payload: core::marker::PhantomData<fn() -> P>,
193}
194
195impl<P: Clone> ProbabilisticBroadcast<P, FairLossLink<Gossip<P>>> {
196 /// Gossip among `peers`, over the fair-loss link Algorithm 3.9 names.
197 ///
198 /// The book says `Uses: FairLossPointToPointLinks` and this default honours it. A perfect link
199 /// would retransmit until delivery, which masks the probabilistic guarantee this abstraction
200 /// exists to provide — and would never fall silent, because the stubborn link beneath it
201 /// re-sends everything it has ever sent. Gossip over a link that does not lose is gossip with
202 /// nothing to do.
203 pub fn new(me: NodeId, peers: impl IntoIterator<Item = NodeId>, config: Config) -> Self {
204 Self::with_link(me, peers, FairLossLink::new(), config)
205 }
206}
207
208impl<P: Clone, L: VolatileLink<Carried<P>>> ProbabilisticBroadcast<P, L> {
209 /// Gossip among `peers`, over the link supplied.
210 pub fn with_link(
211 me: NodeId,
212 peers: impl IntoIterator<Item = NodeId>,
213 link: L,
214 config: Config,
215 ) -> Self {
216 let mut peers: BTreeSet<NodeId> = peers.into_iter().collect();
217 peers.insert(me);
218 ProbabilisticBroadcast {
219 me,
220 peers,
221 incarnation: 0,
222 seq: 0,
223 config,
224 delivered: BTreeSet::new(),
225 order: BTreeMap::new(),
226 link: Child::new(link),
227 _payload: core::marker::PhantomData,
228 }
229 }
230
231 /// How many identifiers this process is currently remembering. For the bound test.
232 pub fn remembered(&self) -> usize {
233 self.delivered.len()
234 }
235
236 /// Whether this process has delivered `id` and still remembers doing so.
237 pub fn has_delivered(&self, id: BroadcastId) -> bool {
238 self.delivered.contains(&id)
239 }
240
241 /// The processes this instance gossips among, in a stable order.
242 pub fn peers(&self) -> impl Iterator<Item = NodeId> + '_ {
243 self.peers.iter().copied()
244 }
245}
246
247impl<P: Clone, L> ProbabilisticBroadcast<P, L>
248where
249 L: VolatileLink<Carried<P>>,
250{
251 /// The link beneath.
252 pub fn link(&self) -> &L {
253 &self.link
254 }
255
256 /// Run the link, then act on whatever it reported.
257 fn through_link(
258 &mut self,
259 cx: &mut ProtoCx<'_, Self>,
260 f: impl FnOnce(&mut L, &mut ProtoCx<'_, L>),
261 ) {
262 let mut inds = self.link.run(cx, core::convert::identity, f);
263 for ind in inds.drain(..) {
264 match L::classify(ind) {
265 LinkInd::Deliver { msg, .. } => self.on_arrival(msg, cx),
266 LinkInd::Boundary(Boundary::Ended { peer, epoch }) => {
267 cx.indicate(Ind::SessionEnded { peer, epoch })
268 }
269 LinkInd::Boundary(Boundary::Established { peer, epoch }) => {
270 cx.indicate(Ind::SessionEstablished { peer, epoch })
271 }
272 }
273 }
274 self.link.reclaim(inds);
275 }
276
277 /// `upon event ⟨ fll, Deliver | p, [GOSSIP, s, m, r] ⟩`.
278 ///
279 /// The two clauses are siblings, not nested. See the module note: relaying a message already
280 /// delivered is the book's, and removing the redundancy would remove the guarantee with it.
281 fn on_arrival(&mut self, msg: Gossip<P>, cx: &mut ProtoCx<'_, Self>) {
282 let Gossip { id, ttl, payload } = msg;
283
284 if !self.delivered.contains(&id) {
285 self.record(id);
286 cx.indicate(Ind::Deliver { from: id.origin, msg: payload.clone() });
287 }
288
289 if ttl > 1 {
290 self.gossip(Gossip { id, ttl: ttl - 1, payload }, cx);
291 }
292 }
293
294 /// `procedure gossip(msg) is forall t ∈ picktargets(k) do trigger ⟨ fll, Send | t, msg ⟩`.
295 fn gossip(&mut self, msg: Gossip<P>, cx: &mut ProtoCx<'_, Self>) {
296 for target in self.picktargets(cx) {
297 let out = msg.clone();
298 let inds = self.link.run(cx, core::convert::identity, |link, ccx| {
299 link.on_cmd(L::send(target, out), ccx)
300 });
301 debug_assert!(inds.is_empty(), "a send must not deliver synchronously");
302 self.link.reclaim(inds);
303 }
304 }
305
306 /// `picktargets(k)` — `k` peers other than this one, drawn without replacement.
307 ///
308 /// Never the whole membership: fanning out to all of `Π` is best-effort broadcast, which this
309 /// repository already has, and a probabilistic broadcast that does it has paid for uncertainty
310 /// and bought nothing. When the fanout exceeds the number of peers the draw returns them all,
311 /// which is a configuration mistake rather than a special case worth encoding.
312 fn picktargets(&self, cx: &mut ProtoCx<'_, Self>) -> Vec<NodeId> {
313 use rand::Rng;
314 let mut candidates: Vec<NodeId> =
315 self.peers.iter().copied().filter(|p| *p != self.me).collect();
316 let take = self.config.fanout.min(candidates.len());
317 // Partial Fisher-Yates: `take` draws, each from what is left, so no peer is chosen twice
318 // and the cost is the fanout rather than the membership.
319 for i in 0..take {
320 let j = i + cx.rng().random_range(0..candidates.len() - i);
321 candidates.swap(i, j);
322 }
323 candidates.truncate(take);
324 candidates
325 }
326
327 /// Remember `id`, evicting this sender's oldest if the window is full.
328 ///
329 /// Constant work per insert. The alternative — expiring by age with a pass over the set — is
330 /// the previous implementation's one surviving defect, and is what this shape exists to avoid.
331 fn record(&mut self, id: BroadcastId) {
332 self.delivered.insert(id);
333 let seen = self.order.entry(id.origin).or_default();
334 seen.push_back((id.incarnation, id.seq));
335 if seen.len() > self.config.window
336 && let Some((incarnation, seq)) = seen.pop_front()
337 {
338 self.delivered.remove(&BroadcastId { origin: id.origin, incarnation, seq });
339 }
340 }
341}
342
343impl<P: Clone, L> Protocol for ProbabilisticBroadcast<P, L>
344where
345 L: VolatileLink<Carried<P>>,
346{
347 type Cmd = Cmd<P>;
348 type Ind = Ind<P>;
349 type Msg = L::Msg;
350 /// Whatever the link's guarantees are conditional on. This layer adds no condition of its own
351 /// and cannot bridge the link's.
352 type Scope = L::Scope;
353 type Note = crate::Note;
354 /// Keeps nothing durably: a crash loses everything this protocol knows, which is why `PB2` is
355 /// scoped to the window *within* an incarnation and says nothing across one.
356 type Meta = core::convert::Infallible;
357 type Entry = core::convert::Infallible;
358
359 /// `upon event ⟨ pb, Broadcast | m ⟩` — deliver to self, then gossip.
360 fn on_cmd(&mut self, Cmd::Broadcast(payload): Cmd<P>, cx: &mut ProtoCx<'_, Self>) {
361 self.seq += 1;
362 let id = BroadcastId { origin: self.me, incarnation: self.incarnation, seq: self.seq };
363 self.record(id);
364 cx.indicate(Ind::Deliver { from: self.me, msg: payload.clone() });
365 self.gossip(Gossip { id, ttl: self.config.rounds, payload }, cx);
366 }
367
368 fn on_msg(&mut self, from: NodeId, msg: L::Msg, cx: &mut ProtoCx<'_, Self>) {
369 self.through_link(cx, |link, ccx| link.on_msg(from, msg, ccx));
370 }
371
372 fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
373 self.through_link(cx, |link, ccx| link.on_timer(id, ccx));
374 }
375
376 /// Name this incarnation. Runs on a first start and on every restart, which is the point: a
377 /// restarted process is a new incarnation and its identifiers must say so.
378 fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
379 self.incarnation = cx.rng().next_u64();
380 self.through_link(cx, |link, ccx| link.on_init(ccx));
381 }
382
383 /// Hand the scope ending down to the link. The trait's default would drop it.
384 fn on_scope_event(&mut self, scope: L::Scope, cx: &mut ProtoCx<'_, Self>) {
385 self.through_link(cx, |link, ccx| link.on_scope_event(scope, ccx));
386 }
387}