Skip to main content

recon_protocols/
stubborn_broadcast.rs

1//! Stubborn best-effort broadcast.
2//!
3//! Cachin, Guerraoui & Rodrigues, §3.5.
4//!
5//! **Status: deployable in the fail-recovery model. Space: bounded by membership and by what is
6//! outstanding.** It holds the process set and the messages it is still transmitting, and nothing
7//! per delivery.
8//!
9//! ```text
10//! upon event ⟨ sbeb, Broadcast | m ⟩ do
11//!     forall q ∈ Π do trigger ⟨ sl, Send | q, m ⟩;
12//!
13//! upon event ⟨ sl, Deliver | p, m ⟩ do
14//!     trigger ⟨ sbeb, Deliver | p, m ⟩;
15//! ```
16//!
17//! # Why the repeats are the point
18//!
19//! [`crate::best_effort_broadcast`] fans out over perfect links, which deduplicate and stop
20//! retransmitting once a message has arrived. That is right in the crash-stop model and wrong in
21//! the fail-recovery one: a process that was **down** when a message was sent has no record of it
22//! and no way to ask, so the only thing that reaches it is a sender that never stopped trying.
23//!
24//! So this layer does not deduplicate, and must not. The repeats are what the layer above is for:
25//! [`crate::logged_uniform_reliable_broadcast`] checks its own durable log before acting, and is
26//! idempotent by construction. Deduplicating here would suppress exactly the retransmission a
27//! recovered process depends on.
28//!
29//! # Departures from the page
30//!
31//! - The message carries no identifier, because nothing here deduplicates. That leaves the layer
32//!   above to name its own messages, which is what it already does.
33//! - `Stop` is offered so a caller that knows a message is everywhere can retire it. The book has
34//!   no such request and never lets go; what is outstanding is the whole of the space claim
35//!   above, so a `Stop` nobody can name would make that claim rest on nothing. The caller names
36//!   the broadcast, and one name retires the fan-out of `N` link transmissions it became.
37//!
38//!   Nothing in this repository calls it yet: [`crate::logged_uniform_reliable_broadcast`] never
39//!   stops, because retransmission for ever is what reaches a recovered process, and its space is
40//!   unbounded for that reason and says so.
41
42use core::time::Duration;
43use recon_core::{Child, NodeId, ProtoCx, Protocol, TimerId};
44use std::collections::{BTreeMap, BTreeSet};
45
46use crate::stubborn_link::{self as sl, SendId, StubbornLink};
47
48/// Names one broadcast, so the caller can later retire it.
49///
50/// Distinct from the [`SendId`]s beneath: one broadcast becomes `N` stubborn transmissions, and
51/// the caller has no way to name those — they are minted here. The layer above allocates this,
52/// and the same rule applies as to a `SendId`: it must not name a broadcast that is still live.
53#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
54pub struct BroadcastId(pub u64);
55
56/// Requests from the layer above.
57#[derive(Debug, Clone, PartialEq, Eq)]
58pub enum Cmd<P> {
59    /// Broadcast `msg` as `id`, and keep transmitting it until stopped.
60    Broadcast { id: BroadcastId, msg: P },
61    /// Stop retransmitting the broadcast named `id`, on every link it went out over.
62    Stop { id: BroadcastId },
63}
64
65/// Indications to the layer above.
66#[derive(Debug, Clone, PartialEq, Eq)]
67pub enum Ind<P> {
68    /// A message arrived. Raised many times for one broadcast, by design.
69    Deliver { from: NodeId, msg: P },
70}
71
72/// Fan-out that never gives up.
73#[derive(Debug)]
74pub struct StubbornBroadcast<P: Clone> {
75    peers: BTreeSet<NodeId>,
76    seq: u64,
77    /// What each live broadcast became: one stubborn transmission per peer. Bounded by
78    /// membership times what is outstanding, and this is what `Stop` needs in order to work.
79    outstanding: BTreeMap<BroadcastId, Vec<SendId>>,
80    link: Child<StubbornLink<P>>,
81}
82
83impl<P: Clone> StubbornBroadcast<P> {
84    /// Broadcast among `members`, which must include `me`, retransmitting every `interval`.
85    pub fn new(me: NodeId, members: impl IntoIterator<Item = NodeId>, interval: Duration) -> Self {
86        let mut peers: BTreeSet<NodeId> = members.into_iter().collect();
87        peers.insert(me);
88        StubbornBroadcast {
89            peers,
90            seq: 0,
91            outstanding: BTreeMap::new(),
92            link: Child::new(StubbornLink::new(interval)),
93        }
94    }
95
96    /// The process set. This layer's whole state, besides what is still being transmitted.
97    pub fn peers(&self) -> impl Iterator<Item = NodeId> + '_ {
98        self.peers.iter().copied()
99    }
100
101    /// How many broadcasts are still being retransmitted. Nothing retires one but [`Cmd::Stop`].
102    pub fn outstanding_count(&self) -> usize {
103        self.outstanding.len()
104    }
105}
106
107impl<P: Clone> StubbornBroadcast<P> {
108    fn with_link(
109        &mut self,
110        cx: &mut ProtoCx<'_, Self>,
111        f: impl FnOnce(&mut StubbornLink<P>, &mut ProtoCx<'_, StubbornLink<P>>),
112    ) {
113        let mut inds = self.link.run(cx, core::convert::identity, f);
114        for sl::Ind::Deliver { from, msg } in inds.drain(..) {
115            // Straight through. No deduplication, deliberately.
116            cx.indicate(Ind::Deliver { from, msg });
117        }
118        self.link.reclaim(inds);
119    }
120}
121
122impl<P: Clone> Protocol for StubbornBroadcast<P> {
123    type Cmd = Cmd<P>;
124    type Ind = Ind<P>;
125    type Msg = P;
126    type Scope = core::convert::Infallible;
127    type Note = crate::Note;
128    /// Keeps nothing durably: what it is transmitting is rebuilt by the layer above on recovery.
129    type Meta = core::convert::Infallible;
130    type Entry = core::convert::Infallible;
131
132    fn on_cmd(&mut self, cmd: Cmd<P>, cx: &mut ProtoCx<'_, Self>) {
133        match cmd {
134            Cmd::Broadcast { id, msg } => {
135                debug_assert!(!self.outstanding.contains_key(&id), "BroadcastId {id:?} is live");
136                let peers: Vec<NodeId> = self.peers.iter().copied().collect();
137                let mut sends = Vec::with_capacity(peers.len());
138                for q in peers {
139                    self.seq += 1;
140                    let send = SendId(self.seq);
141                    sends.push(send);
142                    let msg = msg.clone();
143                    self.with_link(cx, |link, ccx| {
144                        link.on_cmd(sl::Cmd::Send { id: send, to: q, msg }, ccx)
145                    });
146                }
147                self.outstanding.insert(id, sends);
148            }
149            Cmd::Stop { id } => {
150                // One name, `N` transmissions. The link cannot do this itself: it never learns
151                // that the fan-out was one broadcast.
152                for send in self.outstanding.remove(&id).unwrap_or_default() {
153                    self.with_link(cx, |link, ccx| link.on_cmd(sl::Cmd::Stop { id: send }, ccx));
154                }
155            }
156        }
157    }
158
159    fn on_msg(&mut self, from: NodeId, msg: P, cx: &mut ProtoCx<'_, Self>) {
160        self.with_link(cx, |link, ccx| link.on_msg(from, msg, ccx));
161    }
162
163    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
164        self.with_link(cx, |link, ccx| link.on_timer(id, ccx));
165    }
166}