recon_protocols/majority_ack_uniform_reliable_broadcast.rs
1//! Majority-ack uniform reliable broadcast.
2//!
3//! Cachin, Guerraoui & Rodrigues, Module 3.3 and Algorithm 3.5 ("Majority-Ack Uniform Reliable
4//! Broadcast").
5//!
6//! **Status: transcription. Space: unbounded.** `pending`, `ack` and `delivered` grow exactly as
7//! in [`crate::uniform_reliable_broadcast`]; removing the detector removes a timing assumption,
8//! not the collection debt. See `docs/bounded-space.md`.
9//!
10//! **Assumption: a correct majority, `N > 2f`.** That is the whole of what this layer rests on. It
11//! is a standing property of the deployment rather than a moment-to-moment property of the
12//! network, and it is the same trade the leader-driven consensus algorithms make.
13//!
14//! # What changed, and what it bought
15//!
16//! Algorithm 3.4 delivers when every process still *believed correct* has relayed a message. That
17//! belief comes from a perfect failure detector, and
18//! `uniform_agreement_breaks_when_the_timing_assumption_is_withdrawn` shows what one wrong belief
19//! costs: a live process is dropped from `correct`, the condition is satisfied too early, and a
20//! message is delivered by some processes and not others.
21//!
22//! Algorithm 3.5 asks a different question of the same record, and the book states the change
23//! exactly:
24//!
25//! ```text
26//! // Except for the function candeliver(·) below and for the absence of ⟨ Crash ⟩ events
27//! // triggered by the perfect failure detector, it is the same as Algorithm 3.4.
28//!
29//! function candeliver(m) returns Boolean is
30//! return #(ack[m]) > N/2;
31//! ```
32//!
33//! There is no set of believed-correct processes, so no process is ever excluded, so no wrong
34//! judgement about who has crashed can be made. What is left is arithmetic over a record this
35//! layer already kept. The rest of the algorithm — `pending`, the relay on first sight, the
36//! identifier carrying the originator — is [`crate::uniform_reliable_broadcast`] unchanged:
37//!
38//! ```text
39//! upon event ⟨ urb, Broadcast | m ⟩ do
40//! pending := pending ∪ {(self, m)};
41//! trigger ⟨ beb, Broadcast | [DATA, self, m] ⟩;
42//!
43//! upon event ⟨ beb, Deliver | p, [DATA, s, m] ⟩ do
44//! ack[m] := ack[m] ∪ {p};
45//! if (s, m) ∉ pending then
46//! pending := pending ∪ {(s, m)};
47//! trigger ⟨ beb, Broadcast | [DATA, s, m] ⟩;
48//!
49//! upon exists (s, m) ∈ pending such that candeliver(m) ∧ m ∉ delivered do
50//! delivered := delivered ∪ {m};
51//! trigger ⟨ urb, Deliver | s, m ⟩;
52//! ```
53//!
54//! # When the assumption fails, this layer blocks rather than diverges
55//!
56//! With `N ≤ 2f` — half or more of the processes crashed, or a partition leaving no majority
57//! anywhere — no message reaches a majority and nothing further is delivered. That is a *worse
58//! liveness* failure and no safety failure at all, which is the opposite of what happens to
59//! Algorithm 3.4 when its detector is wrong. A blocked cluster can be repaired by restoring
60//! processes; a split delivery cannot be repaired by anything.
61//!
62//! # Departures from the page
63//!
64//! - The predicate is written `2 · #(ack[m]) > N` rather than `#(ack[m]) > N/2`. The book means
65//! real division; integer division gives the same answer for every `N`, but only by an argument
66//! the reader has to reconstruct.
67//! - `N` is the full membership including this process, and a process's own relay counts like any
68//! other, because best-effort broadcast sends to the sender too.
69//! - `ack` and `delivered` are keyed by an identifier carrying the originator and a per-sender
70//! sequence number, not by message content — as in the all-ack version, so that identical
71//! content broadcast twice is delivered twice.
72//! - This layer has one child, so it has no wire type of its own: the message *is* the broadcast
73//! child's. It is the first place in this stack where a wire type gets simpler going up, and
74//! removing an assumption is what did it.
75//! - There is no `Init` event and no `Start` command. `new` establishes the state and there is
76//! nothing to start, failure detection having gone.
77//! - Neither `ack` nor `pending` is garbage collected, as in the book. Long runs grow.
78//!
79//! # Over a link that reports scope boundaries
80//!
81//! `L` is a parameter, so this one module is also what
82//! `session_majority_ack_uniform_reliable_broadcast` used to be. Algorithm 3.5 has no scope
83//! events; the establishment clause is this layer's, and it is the same one the all-ack version
84//! carries — resend everything pending, unconditionally, directed at the peer whose scope
85//! returned.
86//!
87//! Unconditional matters here for an extra reason. `ack[m]` records who relayed `m` **to this
88//! process**. It says nothing about whether *this* process's relay reached them, and that relay is
89//! the token they are waiting for. Filtering the resend by `q ∉ ack[m]` deadlocks; the argument is
90//! recorded at `resend_to`, where a test found it. The delivery predicate changing does not change
91//! that argument, so the clause is carried over unchanged, including its cost: a re-establishment
92//! sends every pending message to that peer.
93//!
94//! # When the assumption fails, this layer blocks rather than diverges
95//!
96//! A partition leaving one side with fewer than half the processes delivers nothing on that side,
97//! rather than delivering something the majority will never deliver. When the sides rejoin, the
98//! minority catches up through the same resend clause. Compare the all-ack version, where each
99//! side accuses the other and both proceed — which is a split, and permanent.
100//!
101//! ```text
102//! URB1 [always] Validity — conditional on the reachability below
103//! URB2 [incarnation] No duplication — `delivered` is volatile, so a restart forgets it
104//! URB3 [always] No creation
105//! URB4 [always] Uniform agreement — conditional on the reachability below
106//! ```
107//!
108//! `URB2` is `[incarnation]` by `docs/scope-annotated-modules.md` Corollary 7.2: the redundancy
109//! that would have to survive is the `delivered` set, it is held in memory, and the boundary it
110//! cannot cross is this process's own `⟨Init⟩`.
111//!
112//! `URB1` and `URB4` are `[always]` **only while a majority remains mutually reachable**, which is
113//! an assumption and not a property of this code. A partition leaving no side with more than `N/2`
114//! blocks both sides rather than splitting them, and both properties survive the block — nothing
115//! is delivered that should not be. What does not survive a *permanent* split is liveness: the
116//! layer waits for ever, which is the honest outcome and the one the all-ack version cannot offer,
117//! because its detector's accusations let both sides proceed. See
118//! `without_a_majority_the_layer_blocks_rather_than_diverges`.
119//!
120//! Unlike [`crate::uniform_reliable_broadcast`], no timing assumption is among them: removing the
121//! detector removed the synchrony it needed, not just a dependency.
122
123use recon_core::{Child, NodeId, ProtoCx, Protocol, TimerId};
124use std::collections::{BTreeMap, BTreeSet};
125
126use crate::best_effort_broadcast::{self as beb, BestEffortBroadcast};
127use crate::link::VolatileLink;
128use crate::perfect_link::PerfectLink;
129use crate::uniform_reliable_broadcast::{BroadcastId, Data};
130
131/// The message type: the broadcast child's, unwrapped.
132///
133/// With one child there is nothing to multiplex and no discriminant to add. Compare
134/// [`crate::uniform_reliable_broadcast::Wire`], which needs an enum because a detector also sends.
135pub type Msg<P, L = PerfectLink<Data<P>>> = <BestEffortBroadcast<Data<P>, L> as Protocol>::Msg;
136
137/// What a link beneath this layer must carry.
138///
139/// A caller supplying their own link needs this and should not have to read the source to find it:
140/// the payload is wrapped as the same `Data` as [`crate::uniform_reliable_broadcast`], which this layer shares, so the link carries `Carried<P>` rather than `P`.
141pub type Carried<P> = Data<P>;
142
143/// Requests from the layer above. Broadcasting is the only one.
144#[derive(Debug, Clone, PartialEq, Eq)]
145pub enum Cmd<P> {
146 Broadcast(P),
147}
148
149/// Indications to the layer above.
150#[derive(Debug, Clone, PartialEq, Eq)]
151pub enum Ind<P> {
152 /// `from` is the process that originated the message, never a relayer.
153 Deliver { from: NodeId, msg: P },
154 /// The scope with `peer` ended at `epoch`. Raised only over a link that reports boundaries.
155 SessionEnded { peer: NodeId, epoch: u64 },
156 /// A scope with `peer` is in force at `epoch`. The moment the resend becomes possible.
157 SessionEstablished { peer: NodeId, epoch: u64 },
158}
159
160/// Broadcast with uniform agreement, resting on a correct majority and on nothing else.
161#[derive(Debug)]
162pub struct MajorityAckUniformReliableBroadcast<
163 P: Clone,
164 L: VolatileLink<Data<P>> = PerfectLink<Data<P>>,
165> {
166 me: NodeId,
167 seq: u64,
168 /// How many processes there are. The denominator of the majority, and fixed.
169 members: usize,
170 /// Seen and not yet delivered, with the payload kept for delivery.
171 pending: BTreeMap<BroadcastId, P>,
172 /// Which processes have been seen to relay each message.
173 ack: BTreeMap<BroadcastId, BTreeSet<NodeId>>,
174 delivered: BTreeSet<BroadcastId>,
175 beb: Child<BestEffortBroadcast<Data<P>, L>>,
176}
177
178impl<P: Clone> MajorityAckUniformReliableBroadcast<P, PerfectLink<Data<P>>> {
179 /// Broadcast among `members`, which must include `me`.
180 ///
181 /// The guarantees hold while more than half of `members` are correct. There is no timing
182 /// parameter, because there is no timeout: nothing here waits on a clock.
183 pub fn new(
184 me: NodeId,
185 members: impl IntoIterator<Item = NodeId>,
186 retransmit: core::time::Duration,
187 ) -> Self {
188 Self::with_link(me, members, PerfectLink::new(me, retransmit))
189 }
190}
191
192impl<P: Clone, L: VolatileLink<Data<P>>> MajorityAckUniformReliableBroadcast<P, L> {
193 /// Broadcast among `members`, over the link supplied.
194 pub fn with_link(me: NodeId, members: impl IntoIterator<Item = NodeId>, link: L) -> Self {
195 let mut members: BTreeSet<NodeId> = members.into_iter().collect();
196 members.insert(me);
197 let n = members.len();
198 MajorityAckUniformReliableBroadcast {
199 me,
200 seq: 0,
201 members: n,
202 pending: BTreeMap::new(),
203 ack: BTreeMap::new(),
204 delivered: BTreeSet::new(),
205 beb: Child::new(BestEffortBroadcast::with_link(me, members, link)),
206 }
207 }
208
209 /// How many distinct messages have been delivered upward.
210 pub fn delivered_count(&self) -> usize {
211 self.delivered.len()
212 }
213
214 /// Messages seen but not yet deliverable.
215 pub fn pending_count(&self) -> usize {
216 self.pending.len()
217 }
218
219 /// Which processes have relayed `id`, for tests watching the majority form.
220 pub fn acknowledged_by(&self, id: BroadcastId) -> impl Iterator<Item = NodeId> + '_ {
221 self.ack.get(&id).into_iter().flatten().copied()
222 }
223
224 /// How many processes a message must be relayed by before it can be delivered.
225 pub fn majority(&self) -> usize {
226 self.members / 2 + 1
227 }
228}
229
230impl<P: Clone, L> MajorityAckUniformReliableBroadcast<P, L>
231where
232 L: VolatileLink<Data<P>>,
233{
234 /// Run the broadcast child, then act on what it reported.
235 fn with_beb(
236 &mut self,
237 cx: &mut ProtoCx<'_, Self>,
238 f: impl FnOnce(
239 &mut BestEffortBroadcast<Data<P>, L>,
240 &mut ProtoCx<'_, BestEffortBroadcast<Data<P>, L>>,
241 ),
242 ) {
243 let mut inds = self.beb.run(cx, core::convert::identity, f);
244 for ind in inds.drain(..) {
245 match ind {
246 beb::Ind::Deliver { from, msg: Data { id, payload } } => {
247 self.on_beb_deliver(from, id, payload, cx)
248 }
249 beb::Ind::SessionEnded { peer, epoch } => {
250 // Informative only: the peer is unreachable, so nothing can be resent yet.
251 cx.indicate(Ind::SessionEnded { peer, epoch });
252 }
253 beb::Ind::SessionEstablished { peer, epoch } => {
254 cx.indicate(Ind::SessionEstablished { peer, epoch });
255 self.resend_to(peer, cx);
256 }
257 }
258 }
259 self.beb.reclaim(inds);
260 self.check_deliverable(cx);
261 }
262
263 /// `upon event ⟨ beb, Deliver | p, [DATA, s, m] ⟩`.
264 /// On a scope becoming available again, send that peer everything still pending.
265 ///
266 /// Unconditional, and directed at that peer alone. Unconditional matters: this layer has no
267 /// failure detector, so it has no notion of a peer being *excluded*, and a resend filtered on
268 /// some belief about who is still correct would deadlock the very case the majority rule
269 /// exists to survive. A peer absent for longer than any timeout the all-ack version would have
270 /// used is not a stranger when it returns, because nothing ever excluded it.
271 fn resend_to(&mut self, peer: NodeId, cx: &mut ProtoCx<'_, Self>) {
272 let outstanding: Vec<Data<P>> = self
273 .pending
274 .iter()
275 .map(|(id, payload)| Data { id: *id, payload: payload.clone() })
276 .collect();
277 for data in outstanding {
278 self.with_beb(cx, |beb, ccx| beb.on_cmd(beb::Cmd::SendTo { to: peer, msg: data }, ccx));
279 }
280 }
281
282 fn on_beb_deliver(
283 &mut self,
284 from: NodeId,
285 id: BroadcastId,
286 payload: P,
287 cx: &mut ProtoCx<'_, Self>,
288 ) {
289 self.ack.entry(id).or_default().insert(from);
290 if self.pending.insert(id, payload.clone()).is_none() {
291 self.relay(Data { id, payload }, cx);
292 }
293 }
294
295 /// Re-broadcast, so the message survives its originator's crash.
296 fn relay(&mut self, data: Data<P>, cx: &mut ProtoCx<'_, Self>) {
297 // Re-enters the child while its inbox is out on loan, so `run` hands back a fresh one.
298 let inds = self.beb.run(cx, core::convert::identity, |beb, ccx| {
299 beb.on_cmd(beb::Cmd::Broadcast(data), ccx)
300 });
301 debug_assert!(
302 inds.is_empty(),
303 "relaying must not deliver synchronously; if it does, on_beb_deliver can recurse"
304 );
305 self.beb.reclaim(inds);
306 }
307
308 /// `upon exists (s, m) ∈ pending such that candeliver(m) ∧ m ∉ delivered`.
309 ///
310 /// A predicate over state rather than an event. Its only input is `ack`, which grows on a
311 /// delivery from below — there is no second path, the detector having gone, so unlike the
312 /// all-ack version this is called from one place.
313 fn check_deliverable(&mut self, cx: &mut ProtoCx<'_, Self>) {
314 let ready: Vec<BroadcastId> = self
315 .pending
316 .keys()
317 .copied()
318 .filter(|id| !self.delivered.contains(id))
319 .filter(|id| self.can_deliver(*id))
320 .collect();
321
322 for id in ready {
323 self.delivered.insert(id);
324 let payload = self.pending.get(&id).expect("pending by construction").clone();
325 cx.indicate(Ind::Deliver { from: id.origin, msg: payload });
326 }
327 }
328
329 /// `#(ack[m]) > N/2` — more than half the processes have relayed it.
330 fn can_deliver(&self, id: BroadcastId) -> bool {
331 match self.ack.get(&id) {
332 None => false,
333 Some(acked) => 2 * acked.len() > self.members,
334 }
335 }
336}
337
338impl<P: Clone, L> Protocol for MajorityAckUniformReliableBroadcast<P, L>
339where
340 L: VolatileLink<Data<P>>,
341{
342 type Cmd = Cmd<P>;
343 type Ind = Ind<P>;
344 type Msg = Msg<P, L>;
345 /// No scope conditions: this protocol's guarantees do not lapse.
346 /// Whatever the link's guarantees are conditional on. This layer bridges an ending by
347 /// resending on the establishment that follows.
348 type Scope = L::Scope;
349 type Note = crate::Note;
350 /// Keeps nothing durably: a crash loses everything this protocol knows.
351 type Meta = core::convert::Infallible;
352 type Entry = core::convert::Infallible;
353
354 fn on_cmd(&mut self, Cmd::Broadcast(msg): Cmd<P>, cx: &mut ProtoCx<'_, Self>) {
355 self.seq += 1;
356 let id = BroadcastId { origin: self.me, seq: self.seq };
357 self.pending.insert(id, msg.clone());
358 self.ack.entry(id).or_default();
359 let data = Data { id, payload: msg };
360 self.with_beb(cx, |beb, ccx| beb.on_cmd(beb::Cmd::Broadcast(data), ccx));
361 }
362
363 fn on_msg(&mut self, from: NodeId, msg: Msg<P, L>, cx: &mut ProtoCx<'_, Self>) {
364 self.with_beb(cx, |beb, ccx| beb.on_msg(from, msg, ccx));
365 }
366
367 fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
368 self.with_beb(cx, |beb, ccx| beb.on_timer(id, ccx));
369 }
370
371 /// Hand the scope ending down to the link. The trait's default would drop it.
372 fn on_scope_event(&mut self, scope: L::Scope, cx: &mut ProtoCx<'_, Self>) {
373 self.with_beb(cx, |beb, ccx| beb.on_scope_event(scope, ccx));
374 }
375}