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}