Skip to main content

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}