Skip to main content

recon_protocols/
reliable_broadcast.rs

1//! Regular reliable broadcast.
2//!
3//! Cachin, Guerraoui & Rodrigues, Module 3.2 and Algorithm 3.3 ("Eager Reliable Broadcast").
4//!
5//! **Status: transcription. Space: unbounded.** `delivered` grows with every message delivered,
6//! and the book omits its collection deliberately. Deployable once it is windowed — which weakens
7//! no duplication to hold within the retention window. See `docs/bounded-space.md`.
8//!
9//! Best-effort broadcast promises nothing when the sender crashes partway through: some processes
10//! deliver, others do not, and they disagree for ever. This layer adds **agreement** — if any
11//! correct process delivers a message, every correct process eventually does — by having every
12//! process relay each message the first time it delivers it. The redundancy therefore lives at
13//! the other processes, which is why this guarantee survives the sender's crash where best-effort
14//! broadcast's does not.
15//!
16//! ```text
17//! upon event ⟨ rb, Broadcast | m ⟩ do
18//!     trigger ⟨ beb, Broadcast | [DATA, self, m] ⟩;
19//!
20//! upon event ⟨ beb, Deliver | p, [DATA, s, m] ⟩ do
21//!     if m ∉ delivered then
22//!         delivered := delivered ∪ {m};
23//!         trigger ⟨ rb, Deliver | s, m ⟩;
24//!         trigger ⟨ beb, Broadcast | [DATA, s, m] ⟩;
25//! ```
26//!
27//! The relay is unconditional on first delivery — the book's *eager* scheme. Algorithm 3.2, the
28//! lazy variant, relays only when a perfect failure detector reports the sender crashed; that
29//! abstraction below does not exist here, and eager needs no failure detector at all. It pays
30//! for that in messages.
31//!
32//! Scope tags, in the notation of `docs/scope-annotated-modules.md`:
33//!
34//! ```text
35//! RB1 [always]       Validity
36//! RB2 [incarnation]  No duplication  — the delivered set is volatile
37//! RB3 [always]       No creation
38//! RB4 [always]       Agreement       — bridged by redundancy at the other processes
39//! ```
40//!
41//! **Two departures from the page**, both for reasons already met lower in the stack:
42//!
43//! - The book deduplicates on message content, assuming messages are unique across senders. Here
44//!   each broadcast carries an identifier — its originator and a per-sender sequence number — and
45//!   deduplication is on that, so identical content broadcast twice is delivered twice.
46//! - `⟨rb, Init⟩` is not a separate event; `new` establishes the same state.
47//!
48//! # Over a link that reports scope boundaries
49//!
50//! `L` is a parameter, so this one module is also what `session_reliable_broadcast` used to be.
51//! Algorithm 3.3 is unchanged; what changes is the scope its agreement holds within.
52//!
53//! Over a perfect link a relay always arrives, because the link retransmits until it does. Over a
54//! session link it may not, and this layer has nothing with which to retry:
55//!
56//! - It relays **once**, on first delivery. That is what makes eager reliable broadcast eager.
57//! - It keeps `delivered` as a set of **identifiers**, not payloads — so even knowing a relay was
58//!   lost, it has no copy to send again. Retaining payloads would be state growing with messages,
59//!   which `docs/bounded-space.md` forbids without a window.
60//! - It is **fail-silent**. Algorithm 3.3 uses no failure detector, so it cannot conclude that a
61//!   process is gone and stop expecting to reach it.
62//!
63//! So when a relay is lost to a scope ending, nothing retries and nothing gives up. This layer
64//! cannot bridge, so it propagates: the boundary is reported upward in its own `Ind` rather than
65//! absorbed, which is what `docs/conditional-guarantees.md` requires of a layer in that position.
66//!
67//! ```text
68//! RB1 [session]       Validity
69//! RB2 [incarnation]   No duplication — `delivered` is volatile, so a restart forgets it
70//! RB3 [always]        No creation
71//! RB4 [session]       Agreement — within the scopes carrying the relay, and not across one
72//! ```
73//!
74//! `RB2` is `[incarnation]` for the reason `docs/scope-annotated-modules.md` gives as Corollary
75//! 7.2, and by the same argument: the redundancy that would have to survive is `delivered`, that
76//! set is held in memory, and the boundary it cannot cross is this process's own `⟨Init⟩`. A
77//! recipient that restarts and is then relayed a message it had already delivered — which the
78//! eager relay of a peer that did *not* restart will happily do — delivers it a second time.
79//! `[always]` would be the claim that a volatile set survives a crash.
80//!
81//! This is not a defect to be fixed here. It is the honest reading of Algorithm 3.3 on a link that
82//! can lose a suffix, and it is exactly what uniform reliable broadcast does not share — that one
83//! has a failure detector, and between reconnection and accusation it has no third outcome.
84//! Reading the two together is the sharpest available argument for why uniform reliable broadcast
85//! needs a detector at all.
86
87use core::time::Duration;
88use recon_core::{Child, NodeId, ProtoCx, Protocol, TimerId};
89use serde::{Deserialize, Serialize};
90use std::collections::BTreeSet;
91
92use crate::best_effort_broadcast::{self as beb, BestEffortBroadcast};
93use crate::link::VolatileLink;
94use crate::perfect_link::PerfectLink;
95
96/// Names one broadcast uniquely: who originated it, and their sequence number for it.
97///
98/// A relayed message must still be attributed to its originator, so the identifier travels with
99/// the message rather than being derived from whoever most recently sent it.
100#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
101pub struct BroadcastId {
102    pub origin: NodeId,
103    pub seq: u64,
104}
105
106/// What this layer adds to the wire: the originator, and the payload.
107///
108/// The first header contributed above the perfect link's. It exists because a relayer is not the
109/// sender, and without it a recipient could not tell who originated what it received.
110#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
111pub struct Data<P> {
112    pub id: BroadcastId,
113    pub payload: P,
114}
115
116/// What a link beneath this layer must carry.
117///
118/// A caller supplying their own link needs this and should not have to read the source to find it:
119/// the payload is wrapped as `Data` — the payload with the originator and sequence number Algorithm 3.3 needs, so the link carries `Carried<P>` rather than `P`.
120pub type Carried<P> = Data<P>;
121
122/// Requests from the layer above.
123#[derive(Debug, Clone, PartialEq, Eq)]
124pub enum Cmd<P> {
125    Broadcast(P),
126}
127
128/// Indications to the layer above.
129#[derive(Debug, Clone, PartialEq, Eq)]
130pub enum Ind<P> {
131    /// `from` is the process that *originated* the message, never the one that relayed it.
132    Deliver { from: NodeId, msg: P },
133    /// The scope with `peer` ended at `epoch`.
134    ///
135    /// RB4's agreement is scoped to the sessions that carried the relay: this layer relays once,
136    /// on first receipt, and never again, so a relay lost to an ending is lost for good. It holds
137    /// no redundancy that outlives the scope, so it cannot bridge and must propagate.
138    /// Raised only over a link that reports boundaries.
139    SessionEnded { peer: NodeId, epoch: u64 },
140    /// A scope with `peer` is in force at `epoch`.
141    SessionEstablished { peer: NodeId, epoch: u64 },
142}
143
144/// The wire type: this layer's data, carried by best-effort broadcast.
145pub type Wire<P, L = PerfectLink<Data<P>>> = <BestEffortBroadcast<Data<P>, L> as Protocol>::Msg;
146
147/// Broadcast with agreement, over best-effort broadcast.
148#[derive(Debug)]
149pub struct ReliableBroadcast<P: Clone, L: VolatileLink<Data<P>> = PerfectLink<Data<P>>> {
150    me: NodeId,
151    seq: u64,
152    delivered: BTreeSet<BroadcastId>,
153    beb: Child<BestEffortBroadcast<Data<P>, L>>,
154}
155
156impl<P: Clone> ReliableBroadcast<P, PerfectLink<Data<P>>> {
157    /// Reliable broadcast for process `me` among `peers`, over links retransmitting every
158    /// `interval`.
159    pub fn new(me: NodeId, peers: impl IntoIterator<Item = NodeId>, interval: Duration) -> Self {
160        Self::with_link(me, peers, PerfectLink::new(me, interval))
161    }
162}
163
164impl<P: Clone, L: VolatileLink<Data<P>>> ReliableBroadcast<P, L> {
165    /// Reliable broadcast for process `me` among `peers`, over the link supplied.
166    pub fn with_link(me: NodeId, peers: impl IntoIterator<Item = NodeId>, link: L) -> Self {
167        ReliableBroadcast {
168            me,
169            seq: 0,
170            delivered: BTreeSet::new(),
171            beb: Child::new(BestEffortBroadcast::with_link(me, peers, link)),
172        }
173    }
174
175    /// How many distinct broadcasts this process has delivered upward.
176    pub fn delivered_count(&self) -> usize {
177        self.delivered.len()
178    }
179
180    /// Whether this process has already delivered `id`.
181    pub fn has_delivered(&self, id: BroadcastId) -> bool {
182        self.delivered.contains(&id)
183    }
184}
185
186impl<P: Clone, L> ReliableBroadcast<P, L>
187where
188    L: VolatileLink<Data<P>>,
189{
190    /// Run the child, then act on whatever it reported.
191    fn with_beb(
192        &mut self,
193        cx: &mut ProtoCx<'_, Self>,
194        f: impl FnOnce(
195            &mut BestEffortBroadcast<Data<P>, L>,
196            &mut ProtoCx<'_, BestEffortBroadcast<Data<P>, L>>,
197        ),
198    ) {
199        let mut inds = self.beb.run(cx, core::convert::identity, f);
200        for ind in inds.drain(..) {
201            match ind {
202                beb::Ind::Deliver { msg: Data { id, payload }, .. } => {
203                    self.on_beb_deliver(id, payload, cx)
204                }
205                // This layer relays once, on first receipt, and never again, so it holds no
206                // redundancy outliving the scope and cannot repair what an ending lost. It
207                // propagates instead — which is the whole of what
208                // `docs/conditional-guarantees.md` requires of a layer that cannot bridge.
209                beb::Ind::SessionEnded { peer, epoch } => {
210                    cx.indicate(Ind::SessionEnded { peer, epoch })
211                }
212                beb::Ind::SessionEstablished { peer, epoch } => {
213                    cx.indicate(Ind::SessionEstablished { peer, epoch })
214                }
215            }
216        }
217        self.beb.reclaim(inds);
218    }
219
220    /// Algorithm 3.3's second handler: deliver once, then relay once.
221    fn on_beb_deliver(&mut self, id: BroadcastId, payload: P, cx: &mut ProtoCx<'_, Self>) {
222        if !self.delivered.insert(id) {
223            return; // already seen — neither delivered again nor relayed again
224        }
225        // Attributed to the originator, not to whoever relayed it here.
226        cx.indicate(Ind::Deliver { from: id.origin, msg: payload.clone() });
227        self.relay(Data { id, payload }, cx);
228    }
229
230    /// Re-broadcast a message so that it survives the originator's crash.
231    ///
232    /// This re-enters the child while its inbox is out on loan, so `run` hands back a fresh one.
233    /// Best-effort broadcast turns a request into sends and timers only — a message to this process
234    /// travels through the links like any other and arrives later — so no indication can be raised
235    /// here. The assertion records that reasoning rather than trusting it.
236    fn relay(&mut self, data: Data<P>, cx: &mut ProtoCx<'_, Self>) {
237        let inds = self.beb.run(cx, core::convert::identity, |beb, ccx| {
238            beb.on_cmd(beb::Cmd::Broadcast(data), ccx)
239        });
240        debug_assert!(
241            inds.is_empty(),
242            "relaying must not deliver synchronously; if it does, on_beb_deliver can recurse"
243        );
244        self.beb.reclaim(inds);
245    }
246}
247
248impl<P: Clone, L> Protocol for ReliableBroadcast<P, L>
249where
250    L: VolatileLink<Data<P>>,
251{
252    type Cmd = Cmd<P>;
253    type Ind = Ind<P>;
254    type Msg = Wire<P, L>;
255    /// Whatever the link's guarantees are conditional on. RB4 is scoped to the sessions that
256    /// carried the relay, and this layer cannot bridge one ending.
257    type Scope = L::Scope;
258    type Note = crate::Note;
259    /// Keeps nothing durably: a crash loses everything this protocol knows.
260    type Meta = core::convert::Infallible;
261    type Entry = core::convert::Infallible;
262
263    fn on_cmd(&mut self, Cmd::Broadcast(msg): Cmd<P>, cx: &mut ProtoCx<'_, Self>) {
264        self.seq += 1;
265        let data = Data { id: BroadcastId { origin: self.me, seq: self.seq }, payload: msg };
266        self.with_beb(cx, |beb, ccx| beb.on_cmd(beb::Cmd::Broadcast(data), ccx));
267    }
268
269    fn on_msg(&mut self, from: NodeId, msg: Wire<P, L>, cx: &mut ProtoCx<'_, Self>) {
270        self.with_beb(cx, |beb, ccx| beb.on_msg(from, msg, ccx));
271    }
272
273    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
274        self.with_beb(cx, |beb, ccx| beb.on_timer(id, ccx));
275    }
276
277    /// Hand the scope ending down. This layer cannot bridge one, and the default would drop it.
278    fn on_scope_event(&mut self, scope: L::Scope, cx: &mut ProtoCx<'_, Self>) {
279        self.with_beb(cx, |beb, ccx| beb.on_scope_event(scope, ccx));
280    }
281}