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}