Skip to main content

recon_protocols/
uniform_reliable_broadcast.rs

1//! Uniform reliable broadcast.
2//!
3//! Cachin, Guerraoui & Rodrigues, Module 3.3 and Algorithm 3.4 ("All-Ack Uniform Reliable
4//! Broadcast").
5//!
6//! **Status: transcription. Space: unbounded.** `pending` holds payloads, `ack` holds a process
7//! set per message, and `delivered` grows without limit. Deployable once collected — and this
8//! layer already computes the predicate it would need, since `correct ⊆ ack[m]` is a stability
9//! test: a message every correct process has seen can be dropped from `pending` and `ack` the
10//! moment it is delivered. See `docs/bounded-space.md`.
11//!
12//! Reliable broadcast guarantees agreement only among *correct* processes. A process that
13//! delivers a message and then crashes may leave the survivors never delivering it — and if that
14//! delivery had any external effect, the divergence cannot be repaired from above. Uniform
15//! agreement quantifies over *any* process: delivered by anyone at all, correct or not, means
16//! eventually delivered by everyone correct.
17//!
18//! It is bought by waiting. A message is delivered only once every process still believed correct
19//! has been seen to acknowledge it, so nobody can deliver something the others have not yet seen.
20//!
21//! ```text
22//! upon event ⟨ urb, Broadcast | m ⟩ do
23//!     pending := pending ∪ {(self, m)};
24//!     trigger ⟨ beb, Broadcast | [DATA, self, m] ⟩;
25//!
26//! upon event ⟨ beb, Deliver | p, [DATA, s, m] ⟩ do
27//!     ack[m] := ack[m] ∪ {p};
28//!     if (s, m) ∉ pending then
29//!         pending := pending ∪ {(s, m)};
30//!         trigger ⟨ beb, Broadcast | [DATA, s, m] ⟩;
31//!
32//! upon event ⟨ P, Crash | p ⟩ do
33//!     correct := correct \ {p};
34//!
35//! function candeliver(m) is  return correct ⊆ ack[m];
36//!
37//! upon exists (s, m) ∈ pending such that candeliver(m) ∧ m ∉ delivered do
38//!     delivered := delivered ∪ {m};
39//!     trigger ⟨ urb, Deliver | s, m ⟩;
40//! ```
41//!
42//! # This layer depends on a timing assumption, and cannot detect its failure
43//!
44//! Uniform agreement holds only while the failure detector is *accurate*, which holds only while
45//! the network delivers within a known bound. A wrongly accused process is removed from `correct`,
46//! `candeliver` is satisfied too early, and a message can be delivered by some processes and not
47//! others.
48//!
49//! That dependency is stated here rather than expressed with the scope annotation of
50//! `docs/scope-annotated-modules.md`, and deliberately: a scope must have a boundary the module
51//! can observe. This one has none. Synchrony failing arrives here as the detector reporting a
52//! crash, indistinguishable from the detector being right. An assumption a layer rests on but
53//! cannot detect is not a scope — tagging it would create an obligation no implementation could
54//! discharge and no test could exercise.
55//!
56//! # Departures from the page
57//!
58//! - `ack` and `delivered` are keyed by an identifier carrying the originator and a per-sender
59//!   sequence number, not by message content. The book's `ack[m]` assumes messages are unique
60//!   across senders; identical content broadcast twice must be delivered twice.
61//! - Two children both send, so this layer's wire type is an enum distinguishing a broadcast
62//!   payload from a heartbeat. It is the first multiplexing in the stack, and it is typed: a
63//!   mis-wiring is a compile error rather than a silently undelivered message.
64//! - `⟨urb, Init⟩` **is** a separate event: `new` establishes the state, and
65//!   [`Protocol::on_init`] starts the detector beneath. It was a `Cmd::Start` before the trait had
66//!   an init event, which is why the only request now is [`Cmd::Broadcast`].
67//! - Neither `ack` nor `pending` is garbage collected, as in the book. Long runs grow.
68//!
69//! # Over a link that reports scope boundaries
70//!
71//! `L` is a parameter, so this one module is also what `session_uniform_reliable_broadcast` used
72//! to be. Algorithm 3.4 gains **one** clause and nothing else changes:
73//!
74//! ```text
75//! upon event ⟨ SessionEstablished | q ⟩ do
76//!     forall (s, m) ∈ pending do
77//!         trigger ⟨ beb, SendTo | q, [DATA, s, m] ⟩;
78//! ```
79//!
80//! It reuses `pending`, which the algorithm already maintains, and `SendTo`, which is a narrowing
81//! of an existing action rather than a new communication step.
82//!
83//! Two things about that clause are deliberate and neither is what one would write first. It is
84//! **unconditional**, not filtered by `q ∉ ack[m]`: the filtered version deadlocks, for the reason
85//! set out at `resend_to`, which is where a test found it. And it is **directed** at the peer
86//! whose scope returned rather than broadcast to everyone, because the ending was per peer and so
87//! is the repair. Nothing is attempted on the *ending* itself: the peer is unreachable at that
88//! moment and anything sent would be discarded.
89//!
90//! # This layer keeps the **perfect** detector, and that is not an oversight
91//!
92//! [`crate::eventually_perfect_failure_detector`] exists, and Ω moved onto it. This layer must not.
93//! Its agreement rests on *strong* accuracy by name — the book's own proof invokes it — and one
94//! false suspicion breaks it, permanently and silently. A detector allowed to be wrong would turn a
95//! correct module into a broken one.
96//!
97//! That is not a gap to be closed later: it is the whole subject of
98//! [`crate::majority_ack_uniform_reliable_broadcast`], which asks *has more than half relayed it?*
99//! instead of *has everyone still believed correct relayed it?*, and drops the detector entirely.
100//! The pair is the demonstration of what strong accuracy is worth, and replacing the detector here
101//! would delete the demonstration rather than improve it.
102//!
103//! ```text
104//! URB1 [always]       Validity — conditional on the two mechanisms below
105//! URB2 [incarnation]  No duplication — `delivered` is volatile, so a restart forgets it
106//! URB3 [always]       No creation
107//! URB4 [always]       Uniform agreement — conditional on the two mechanisms below
108//! ```
109//!
110//! `URB2` is `[incarnation]` by `docs/scope-annotated-modules.md` Corollary 7.2: the set that
111//! would have to survive is `delivered`, it is held in memory, and the boundary it cannot cross is
112//! this process's own `⟨Init⟩`.
113//!
114//! `URB1` and `URB4` are `[always]` only because between two mechanisms no third outcome is left —
115//! the scope comes back and the resend repairs it, or the peer never returns and the detector's
116//! timeout drops it from `correct`. Each carries a condition that is an assumption rather than a
117//! property of this code. The reconnection path needs the peer to be reachable again; the
118//! accusation path needs the detector's synchrony assumption, and `perfect_failure_detector` is
119//! explicit that outside a synchronous system it accuses correct processes. **Both failing at once
120//! is a permanent split**: each side of a partition accuses the other, each has `correct ⊆ ack[m]`
121//! satisfied among itself, and both deliver — not uniform agreement failing on a technicality but
122//! two disjoint sets of processes proceeding as though the other did not exist.
123//! [`crate::majority_ack_uniform_reliable_broadcast`] cannot suffer that and blocks instead; that
124//! difference is what a quorum buys and a detector costs.
125//!
126//! Read against reliable broadcast, which has neither mechanism and whose agreement is therefore
127//! scoped, this is the clearest statement of what a failure detector buys.
128
129use core::time::Duration;
130use recon_core::{Child, NodeId, ProtoCx, Protocol, TimerId};
131use serde::{Deserialize, Serialize};
132use std::collections::{BTreeMap, BTreeSet};
133
134use crate::best_effort_broadcast::{self as beb, BestEffortBroadcast};
135use crate::link::VolatileLink;
136use crate::perfect_failure_detector::{self as pfd, Heartbeat, PerfectFailureDetector};
137use crate::perfect_link::{self as pl, PerfectLink};
138
139/// Names one broadcast uniquely: who originated it, and their sequence number for it.
140#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
141pub struct BroadcastId {
142    pub origin: NodeId,
143    pub seq: u64,
144}
145
146/// What this layer adds to a broadcast payload.
147#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
148pub struct Data<P> {
149    pub id: BroadcastId,
150    pub payload: P,
151}
152
153/// What best-effort broadcast puts on the wire for this layer's payloads.
154///
155/// Written concretely rather than as `<BestEffortBroadcast<Data<P>> as Protocol>::Msg`: the
156/// projection is only well-formed where `Data<P>: Clone`, which would push that bound onto every
157/// use of [`Wire`]. The assertion below keeps the two from drifting apart.
158pub type BebMsg<P> = pl::Wire<Data<P>>;
159
160const _: () = {
161    /// Fails to compile if best-effort broadcast ever puts something else on the wire.
162    fn _beb_msg_is_what_we_say_it_is<P: Clone>(
163        m: BebMsg<P>,
164    ) -> <BestEffortBroadcast<Data<P>> as Protocol>::Msg {
165        m
166    }
167};
168
169/// What a link beneath this layer must carry.
170///
171/// A caller supplying their own link needs this and should not have to read the source to find it:
172/// the payload is wrapped as `Data` — the payload with the identifier Algorithm 3.4 acknowledges against, so the link carries `Carried<P>` rather than `P`.
173pub type Carried<P> = Data<P>;
174
175/// The wire type, multiplexing the two children.
176///
177/// The first multiplexing in this stack. It is an enum rather than a keyed registry, so a message
178/// routed to the wrong child cannot compile.
179#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
180pub enum Wire<M> {
181    Broadcast(M),
182    Detector(Heartbeat),
183}
184
185/// Requests from the layer above.
186#[derive(Debug, Clone, PartialEq, Eq)]
187pub enum Cmd<P> {
188    Broadcast(P),
189}
190
191/// Indications to the layer above.
192#[derive(Debug, Clone, PartialEq, Eq)]
193pub enum Ind<P> {
194    /// `from` is the process that originated the message, never a relayer.
195    Deliver { from: NodeId, msg: P },
196    /// The scope with `peer` ended at `epoch`.
197    ///
198    /// Raised only over a link that reports boundaries. Unlike reliable broadcast, this layer
199    /// *can* bridge one — it holds `pending` until every correct process has acknowledged, so a
200    /// re-establishment lets it resend. It reports the ending anyway, because the layer above may
201    /// have its own reason to know, and because a guarantee that lapsed and was repaired is not
202    /// the same as one that never lapsed.
203    SessionEnded { peer: NodeId, epoch: u64 },
204    /// A scope with `peer` is in force at `epoch`.
205    SessionEstablished { peer: NodeId, epoch: u64 },
206}
207
208/// Broadcast with uniform agreement, over best-effort broadcast and a failure detector.
209#[derive(Debug)]
210pub struct UniformReliableBroadcast<P: Clone, L: VolatileLink<Data<P>> = PerfectLink<Data<P>>> {
211    me: NodeId,
212    seq: u64,
213    /// Every process believed correct. Shrinks on a crash indication and never grows.
214    correct: BTreeSet<NodeId>,
215    /// Seen and not yet delivered, with the payload kept for delivery.
216    pending: BTreeMap<BroadcastId, P>,
217    /// Which processes have been seen to acknowledge each message.
218    ack: BTreeMap<BroadcastId, BTreeSet<NodeId>>,
219    delivered: BTreeSet<BroadcastId>,
220    beb: Child<BestEffortBroadcast<Data<P>, L>>,
221    detector: Child<PerfectFailureDetector>,
222}
223
224impl<P: Clone> UniformReliableBroadcast<P, PerfectLink<Data<P>>> {
225    /// Broadcast among `members`, which must include `me`.
226    ///
227    /// `heartbeat` and `detect_after` configure the failure detector; `detect_after` must exceed
228    /// `heartbeat` plus the network's delivery bound, or the detector will accuse correct
229    /// processes and uniform agreement can break.
230    pub fn new(
231        me: NodeId,
232        members: impl IntoIterator<Item = NodeId>,
233        retransmit: Duration,
234        heartbeat: Duration,
235        detect_after: Duration,
236    ) -> Self {
237        let link = PerfectLink::new(me, retransmit);
238        Self::with_link(me, members, link, heartbeat, detect_after)
239    }
240}
241
242impl<P: Clone, L: VolatileLink<Data<P>>> UniformReliableBroadcast<P, L> {
243    /// Broadcast among `members`, over the link supplied.
244    pub fn with_link(
245        me: NodeId,
246        members: impl IntoIterator<Item = NodeId>,
247        link: L,
248        heartbeat: Duration,
249        detect_after: Duration,
250    ) -> Self {
251        let mut members: BTreeSet<NodeId> = members.into_iter().collect();
252        members.insert(me);
253        UniformReliableBroadcast {
254            me,
255            seq: 0,
256            correct: members.clone(),
257            pending: BTreeMap::new(),
258            ack: BTreeMap::new(),
259            delivered: BTreeSet::new(),
260            beb: Child::new(BestEffortBroadcast::with_link(me, members.clone(), link)),
261            detector: Child::new(PerfectFailureDetector::new(me, members, heartbeat, detect_after)),
262        }
263    }
264
265    /// The processes still believed correct, in a stable order.
266    pub fn correct(&self) -> impl Iterator<Item = NodeId> + '_ {
267        self.correct.iter().copied()
268    }
269
270    /// How many distinct messages have been delivered upward.
271    pub fn delivered_count(&self) -> usize {
272        self.delivered.len()
273    }
274
275    /// Messages seen but not yet deliverable.
276    pub fn pending_count(&self) -> usize {
277        self.pending.len()
278    }
279
280    /// Which processes have acknowledged `id`, for tests that need to see the condition forming.
281    pub fn acknowledged_by(&self, id: BroadcastId) -> impl Iterator<Item = NodeId> + '_ {
282        self.ack.get(&id).into_iter().flatten().copied()
283    }
284}
285
286impl<P: Clone, L> UniformReliableBroadcast<P, L>
287where
288    L: VolatileLink<Data<P>>,
289{
290    /// Run the broadcast child, then act on what it reported.
291    fn with_beb(
292        &mut self,
293        cx: &mut ProtoCx<'_, Self>,
294        f: impl FnOnce(
295            &mut BestEffortBroadcast<Data<P>, L>,
296            &mut ProtoCx<'_, BestEffortBroadcast<Data<P>, L>>,
297        ),
298    ) {
299        let mut inds = self.beb.run(cx, Wire::Broadcast, f);
300        for ind in inds.drain(..) {
301            match ind {
302                beb::Ind::Deliver { from, msg: Data { id, payload } } => {
303                    self.on_beb_deliver(from, id, payload, cx)
304                }
305                beb::Ind::SessionEnded { peer, epoch } => {
306                    // Informative only: the peer is unreachable, so nothing can be resent yet.
307                    // Attempting the repair here would send into a session already gone.
308                    cx.indicate(Ind::SessionEnded { peer, epoch });
309                }
310                beb::Ind::SessionEstablished { peer, epoch } => {
311                    // The moment the repair becomes possible. This layer holds `pending` until
312                    // every correct process has acknowledged, so its redundancy outlives the
313                    // scope and it *can* bridge — which is why it resends rather than merely
314                    // propagating, as reliable broadcast beneath it must.
315                    cx.indicate(Ind::SessionEstablished { peer, epoch });
316                    self.resend_to(peer, cx);
317                }
318            }
319        }
320        self.beb.reclaim(inds);
321        self.check_deliverable(cx);
322    }
323
324    /// Run the detector child, then act on what it reported.
325    /// Resend everything still outstanding to the peer whose scope has just come back.
326    ///
327    /// Only ever reached over a link that reports boundaries: `Link::classify` yields one only for
328    /// a link that can observe one, so over a perfect link this is unreachable rather than merely
329    /// unused. A `ScopedLink` bound was the original plan and does not work — the arm that calls
330    /// this lives in the `Link` impl, so the tighter bound would fall on every link. What keeps it
331    /// honest is the port's own guarantee. See `crate::link` and this change's `design.md`.
332    fn resend_to(&mut self, peer: NodeId, cx: &mut ProtoCx<'_, Self>) {
333        let outstanding: Vec<Data<P>> = self
334            .pending
335            .iter()
336            .map(|(id, payload)| Data { id: *id, payload: payload.clone() })
337            .collect();
338        for data in outstanding {
339            // Same wire message as a relay, strictly fewer recipients, no new communication step.
340            self.with_beb(cx, |beb, ccx| beb.on_cmd(beb::Cmd::SendTo { to: peer, msg: data }, ccx));
341        }
342    }
343
344    fn with_detector(
345        &mut self,
346        cx: &mut ProtoCx<'_, Self>,
347        f: impl FnOnce(&mut PerfectFailureDetector, &mut ProtoCx<'_, PerfectFailureDetector>),
348    ) {
349        let mut inds = self.detector.run(cx, Wire::Detector, f);
350        for pfd::Ind::Crash { node } in inds.drain(..) {
351            // A process never returns to `correct`; the detector's reports are permanent.
352            self.correct.remove(&node);
353        }
354        self.detector.reclaim(inds);
355        self.check_deliverable(cx);
356    }
357
358    /// `upon event ⟨ beb, Deliver | p, [DATA, s, m] ⟩`.
359    fn on_beb_deliver(
360        &mut self,
361        from: NodeId,
362        id: BroadcastId,
363        payload: P,
364        cx: &mut ProtoCx<'_, Self>,
365    ) {
366        self.ack.entry(id).or_default().insert(from);
367        // Relay only on first sight. An identifier determines its payload, so re-inserting the
368        // same id cannot change what is pending — the returned Option is only being read to
369        // learn whether this was the first time.
370        if self.pending.insert(id, payload.clone()).is_none() {
371            self.relay(Data { id, payload }, cx);
372        }
373    }
374
375    /// Re-broadcast, so the message survives its originator's crash.
376    fn relay(&mut self, data: Data<P>, cx: &mut ProtoCx<'_, Self>) {
377        // Re-enters the child while its inbox is out on loan, so `run` hands back a fresh one.
378        let inds = self
379            .beb
380            .run(cx, Wire::Broadcast, |beb, ccx| beb.on_cmd(beb::Cmd::Broadcast(data), ccx));
381        debug_assert!(
382            inds.is_empty(),
383            "relaying must not deliver synchronously; if it does, on_beb_deliver can recurse"
384        );
385        self.beb.reclaim(inds);
386    }
387
388    /// `upon exists (s, m) ∈ pending such that candeliver(m) ∧ m ∉ delivered`.
389    ///
390    /// The book's last clause is a predicate over state rather than an event, so it is evaluated
391    /// wherever its inputs change: `ack` growing on a delivery from below, `correct` shrinking on
392    /// a crash. Called from both child helpers for that reason.
393    fn check_deliverable(&mut self, cx: &mut ProtoCx<'_, Self>) {
394        let ready: Vec<BroadcastId> = self
395            .pending
396            .keys()
397            .copied()
398            .filter(|id| !self.delivered.contains(id))
399            .filter(|id| self.can_deliver(*id))
400            .collect();
401
402        for id in ready {
403            self.delivered.insert(id);
404            let payload = self.pending.get(&id).expect("pending by construction").clone();
405            cx.indicate(Ind::Deliver { from: id.origin, msg: payload });
406        }
407    }
408
409    /// `correct ⊆ ack[m]` — every process still believed correct has been seen to acknowledge it.
410    fn can_deliver(&self, id: BroadcastId) -> bool {
411        match self.ack.get(&id) {
412            None => false,
413            Some(acked) => self.correct.iter().all(|p| acked.contains(p)),
414        }
415    }
416}
417
418impl<P: Clone, L> Protocol for UniformReliableBroadcast<P, L>
419where
420    L: VolatileLink<Data<P>>,
421{
422    type Cmd = Cmd<P>;
423    type Ind = Ind<P>;
424    type Msg = Wire<L::Msg>;
425    /// No scope conditions: this protocol's guarantees do not lapse.
426    /// Whatever the link's guarantees are conditional on. This layer bridges an ending rather
427    /// than absorbing it, but bridging is not the same as never having lapsed.
428    type Scope = L::Scope;
429    type Note = crate::Note;
430    /// Keeps nothing durably: a crash loses everything this protocol knows.
431    type Meta = core::convert::Infallible;
432    type Entry = core::convert::Infallible;
433
434    /// Failure detection begins here, as Module 2.6 has it. It used to need a `Start` command
435    /// because there was no init event to hang the detector's first timer on.
436    fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
437        self.with_detector(cx, |d, ccx| d.on_init(ccx));
438    }
439
440    fn on_cmd(&mut self, cmd: Cmd<P>, cx: &mut ProtoCx<'_, Self>) {
441        match cmd {
442            Cmd::Broadcast(msg) => {
443                self.seq += 1;
444                let id = BroadcastId { origin: self.me, seq: self.seq };
445                self.pending.insert(id, msg.clone());
446                self.ack.entry(id).or_default();
447                let data = Data { id, payload: msg };
448                self.with_beb(cx, |beb, ccx| beb.on_cmd(beb::Cmd::Broadcast(data), ccx));
449            }
450        }
451    }
452
453    fn on_msg(&mut self, from: NodeId, msg: Wire<L::Msg>, cx: &mut ProtoCx<'_, Self>) {
454        match msg {
455            Wire::Broadcast(m) => self.with_beb(cx, |beb, ccx| beb.on_msg(from, m, ccx)),
456            Wire::Detector(h) => self.with_detector(cx, |d, ccx| d.on_msg(from, h, ccx)),
457        }
458    }
459
460    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
461        // Handed to both children: the identity does not say which registered it, and the one that
462        // did not will recognise that and do nothing.
463        self.with_beb(cx, |beb, ccx| beb.on_timer(id, ccx));
464        self.with_detector(cx, |d, ccx| d.on_timer(id, ccx));
465    }
466
467    /// Hand the scope ending down to the children. The trait's default would drop it.
468    fn on_scope_event(&mut self, scope: L::Scope, cx: &mut ProtoCx<'_, Self>) {
469        self.with_beb(cx, |beb, ccx| beb.on_scope_event(scope, ccx));
470    }
471}