Skip to main content

recon_protocols/
flooding_consensus.rs

1//! Flooding consensus.
2//!
3//! Cachin, Guerraoui & Rodrigues, Module 5.1 and Algorithm 5.1 ("Flooding Consensus").
4//!
5//! **Status: academic, fail-stop. Space: bounded by membership and rounds.** The state is one
6//! proposal set and one heard-from set per round entered, each holding at most one entry per
7//! process, and a run enters at most `N` rounds — so it is `O(N²)` and satisfies the rule in
8//! `docs/bounded-space.md` without any collection being added. That is not what makes it
9//! academic. What makes it academic is the assumption underneath it.
10//!
11//! Processes flood their accumulated proposal sets in rounds. A process leaves a round when it
12//! has heard, in that round, from every process it has not been told has crashed. If a round ends
13//! having heard from exactly the same processes as the one before it, nobody new crashed, so
14//! everyone holds the same proposal set and it is safe to decide the minimum of it.
15//!
16//! ```text
17//! upon event ⟨ c, Init ⟩ do
18//!     correct := Π;
19//!     round := 1;
20//!     decision := ⊥;
21//!     receivedfrom := [∅]^N;
22//!     proposals := [∅]^N;
23//!     receivedfrom[0] := Π;
24//!
25//! upon event ⟨ P, Crash | p ⟩ do
26//!     correct := correct \ {p};
27//!
28//! upon event ⟨ c, Propose | v ⟩ do
29//!     proposals[1] := proposals[1] ∪ {v};
30//!     trigger ⟨ beb, Broadcast | [PROPOSAL, 1, proposals[1]] ⟩;
31//!
32//! upon event ⟨ beb, Deliver | p, [PROPOSAL, r, ps] ⟩ do
33//!     receivedfrom[r] := receivedfrom[r] ∪ {p};
34//!     proposals[r] := proposals[r] ∪ ps;
35//!
36//! upon correct ⊆ receivedfrom[round] ∧ decision = ⊥ do
37//!     if receivedfrom[round] = receivedfrom[round − 1] then
38//!         decision := min(proposals[round]);
39//!         trigger ⟨ beb, Broadcast | [DECIDED, decision] ⟩;
40//!         trigger ⟨ c, Decide | decision ⟩;
41//!     else
42//!         round := round + 1;
43//!         trigger ⟨ beb, Broadcast | [PROPOSAL, round, proposals[round − 1]] ⟩;
44//!
45//! upon event ⟨ beb, Deliver | p, [DECIDED, v] ⟩ such that p ∈ correct ∧ decision = ⊥ do
46//!     decision := v;
47//!     trigger ⟨ beb, Broadcast | [DECIDED, decision] ⟩;
48//!     trigger ⟨ c, Decide | decision ⟩;
49//! ```
50//!
51//! # Agreement rests entirely on strong accuracy
52//!
53//! `correct` appears in exactly two places, and a false suspicion corrupts both. The round guard
54//! is `correct ⊆ receivedfrom[round]`, so wrongly shrinking `correct` lets a process finish a
55//! round *without having heard from a correct process*, and two processes can then take `min`
56//! over different sets. The decision-adoption rule is guarded by `p ∈ correct`, so a process that
57//! wrongly suspects the decider discards its `DECIDED` message. The book's own proof names the
58//! dependency: "Because of the *strong accuracy* property of the failure detector, no process
59//! that reaches the end of round r receives a proposal containing a smaller value than v."
60//!
61//! Losing the detector's **accuracy** costs *safety* — two correct processes decide differently,
62//! permanently. Losing its **completeness** costs only *liveness* — everyone blocks, but nobody
63//! is wrong. The asymmetry is the reason this protocol is worth writing.
64//!
65//! # This layer keeps the **perfect** detector, and that is not an oversight
66//!
67//! [`crate::eventually_perfect_failure_detector`] exists, and Ω moved onto it. This layer must not:
68//! it is the fail-**stop** algorithm, and its agreement rests on strong accuracy so completely that
69//! this suite exists largely to demonstrate what one false suspicion costs it. Giving it a detector
70//! allowed to be wrong would not make it deployable; it would make the demonstration untrue.
71//! [`crate::leader_driven_consensus`] is the algorithm that survives an inaccurate detector, and it
72//! is a different algorithm rather than this one reconfigured.
73//!
74//! # Why stabilising later does not help
75//!
76//! The model these algorithms are written against is not one in which the set of correct
77//! processes decays. It is one of eventual stability: bounds come to hold, and after that point
78//! the correct set is agreed and stays agreed. `docs/scope-annotated-modules.md` names this
79//! Assumption F and observes that it is the partial-synchrony global stabilisation time and the ◇
80//! of an eventually-accurate detector in the same clothes.
81//!
82//! An eventually perfect detector would therefore *withdraw* a false suspicion, and every process
83//! would again be held correct by every other — and the split would still be there, because a
84//! decision is irrevocable and was taken while the system was unstable. This is what separates
85//! this protocol from the leader-driven family: flooding consensus commits during instability and so
86//! has nothing left for stabilisation to rescue, whereas a quorum-based algorithm declines to
87//! commit until no conflicting decision is possible. Stated this way the limitation survives
88//! replacing `P` with `◇P`; "the detector never withdraws an accusation" would not.
89//!
90//! # Departures from the page
91//!
92//! - `receivedfrom` and `proposals` are maps keyed by round rather than arrays of size `N`. Only
93//!   rounds actually entered hold an entry; the bound is the same and for the same reason.
94//! - The total order the book assumes on proposals ("we implicitly assume here that the set of
95//!   all possible proposals is totally ordered and the order is known by all processes") is a
96//!   `P: Ord` bound. A value that cannot be totally ordered cannot be proposed.
97//! - The standing condition is re-evaluated in a loop rather than once, because a message for a
98//!   later round may arrive before this process enters it. It terminates: the guard requires this
99//!   process to appear in `receivedfrom[round]`, which happens only when its own broadcast for
100//!   that round returns to it.
101//! - A second `Propose` from the same process is ignored. The book's model has one proposal per
102//!   process, and Module 5.1 provides one decision per instance.
103//! - `⟨c, Init⟩` **is** a separate event: `new` establishes the state, and [`Protocol::on_init`]
104//!   starts the detector beneath, without which no round can ever complete after a crash. It was a
105//!   `Cmd::Start` before the trait had an init event, which is why the only request now is
106//!   [`Cmd::Propose`].
107
108use core::time::Duration;
109use recon_core::{Child, NodeId, ProtoCx, Protocol, TimerId};
110use serde::{Deserialize, Serialize};
111use std::collections::{BTreeMap, BTreeSet};
112
113use crate::best_effort_broadcast::{self as beb, BestEffortBroadcast};
114use crate::link::VolatileLink;
115use crate::perfect_failure_detector::{self as pfd, Heartbeat, PerfectFailureDetector};
116use crate::perfect_link::{self as pl, PerfectLink};
117
118/// What this layer puts on the wire: the two messages Algorithm 5.1 sends, and nothing else.
119///
120/// The round number is the one field this layer adds, and it adds it because it is the one thing
121/// this layer keeps that its children do not.
122#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
123// The derived bound would be `P: Deserialize`; rebuilding the set on the way in also needs the
124// total order the algorithm assumes anyway.
125#[serde(bound(deserialize = "P: Ord + Deserialize<'de>"))]
126pub enum Flood<P> {
127    Proposal { round: u64, proposals: BTreeSet<P> },
128    Decided(P),
129}
130
131/// What best-effort broadcast puts on the wire for this layer's payloads.
132///
133/// Written concretely rather than as a projection, for the reason given in
134/// [`crate::uniform_reliable_broadcast::BebMsg`]. The assertion below keeps the two in step.
135pub type BebMsg<P> = pl::Wire<Flood<P>>;
136
137const _: () = {
138    /// Fails to compile if best-effort broadcast ever puts something else on the wire.
139    fn _beb_msg_is_what_we_say_it_is<P: Clone>(
140        m: BebMsg<P>,
141    ) -> <BestEffortBroadcast<Flood<P>> as Protocol>::Msg {
142        m
143    }
144};
145
146/// What a link beneath this layer must carry.
147///
148/// A caller supplying their own link needs this and should not have to read the source to find it:
149/// the payload is wrapped as `Flood` — a round's proposal or a decision, so the link carries `Carried<P>` rather than `P`.
150pub type Carried<P> = Flood<P>;
151
152/// The wire type, multiplexing the two children. Typed, so a mis-route cannot compile.
153#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
154#[serde(bound(deserialize = "M: Deserialize<'de>"))]
155pub enum Wire<M> {
156    Broadcast(M),
157    Detector(Heartbeat),
158}
159
160/// Requests from the layer above.
161#[derive(Debug, Clone, PartialEq, Eq)]
162pub enum Cmd<P> {
163    Propose(P),
164}
165
166/// Indications to the layer above.
167#[derive(Debug, Clone, PartialEq, Eq)]
168pub enum Ind<P> {
169    Decide(P),
170}
171
172/// Regular consensus in the fail-stop model, over best-effort broadcast and a failure detector.
173///
174/// `L` is the link under the broadcast, and is a parameter rather than a fixed type: what this
175/// stack needs is a link speaking [`pl::Cmd`] and [`pl::Ind`], not one particular implementation.
176/// An application with its own driver-backed link runs this consensus over it unedited. It
177/// defaults to [`PerfectLink`], so the ordinary stack is still written `FloodingConsensus<P>`.
178#[derive(Debug)]
179pub struct FloodingConsensus<P: Clone + Ord, L: VolatileLink<Flood<P>> = PerfectLink<Flood<P>>> {
180    /// Every process believed correct. Shrinks on a crash indication and never grows.
181    correct: BTreeSet<NodeId>,
182    round: u64,
183    decision: Option<P>,
184    proposed: bool,
185    /// Who this process has heard from in each round. `receivedfrom[0]` is the full membership.
186    receivedfrom: BTreeMap<u64, BTreeSet<NodeId>>,
187    /// The proposals accumulated in each round.
188    proposals: BTreeMap<u64, BTreeSet<P>>,
189    beb: Child<BestEffortBroadcast<Flood<P>, L>>,
190    detector: Child<PerfectFailureDetector>,
191}
192
193impl<P: Clone + Ord, L: VolatileLink<Flood<P>>> FloodingConsensus<P, L> {
194    /// Consensus among `members`, over a link the caller supplies.
195    ///
196    /// The link is anything satisfying [`crate::link::Link`]; this layer never names an implementation.
197    pub fn with_link(
198        me: NodeId,
199        members: impl IntoIterator<Item = NodeId>,
200        link: L,
201        heartbeat: Duration,
202        detect_after: Duration,
203    ) -> Self {
204        let mut members: BTreeSet<NodeId> = members.into_iter().collect();
205        members.insert(me);
206        let receivedfrom = BTreeMap::from([(0, members.clone())]);
207        FloodingConsensus {
208            correct: members.clone(),
209            round: 1,
210            decision: None,
211            proposed: false,
212            receivedfrom,
213            proposals: BTreeMap::new(),
214            beb: Child::new(BestEffortBroadcast::with_link(me, members.clone(), link)),
215            detector: Child::new(PerfectFailureDetector::new(me, members, heartbeat, detect_after)),
216        }
217    }
218}
219
220impl<P: Clone + Ord> FloodingConsensus<P, PerfectLink<Flood<P>>> {
221    /// Consensus among `members`, which must include `me`.
222    ///
223    /// `detect_after` must exceed `heartbeat` plus the network's delivery bound, or the detector
224    /// will accuse correct processes and agreement can break — which is the whole subject of this
225    /// module's documentation.
226    pub fn new(
227        me: NodeId,
228        members: impl IntoIterator<Item = NodeId>,
229        retransmit: Duration,
230        heartbeat: Duration,
231        detect_after: Duration,
232    ) -> Self {
233        let mut members: BTreeSet<NodeId> = members.into_iter().collect();
234        members.insert(me);
235        // `receivedfrom[0] := Π` — which is what makes a first-round decision require having
236        // heard from every process, not merely from everyone still believed correct.
237        let receivedfrom = BTreeMap::from([(0, members.clone())]);
238        FloodingConsensus {
239            correct: members.clone(),
240            round: 1,
241            decision: None,
242            proposed: false,
243            receivedfrom,
244            proposals: BTreeMap::new(),
245            beb: Child::new(BestEffortBroadcast::new(me, members.clone(), retransmit)),
246            detector: Child::new(PerfectFailureDetector::new(me, members, heartbeat, detect_after)),
247        }
248    }
249
250    /// The processes still believed correct, in a stable order.
251    pub fn correct(&self) -> impl Iterator<Item = NodeId> + '_ {
252        self.correct.iter().copied()
253    }
254
255    /// The round this process is currently in.
256    pub fn round(&self) -> u64 {
257        self.round
258    }
259
260    /// What this process decided, if it has.
261    pub fn decision(&self) -> Option<&P> {
262        self.decision.as_ref()
263    }
264
265    /// Who this process heard from in `round`, for tests watching the guard form.
266    pub fn heard_from(&self, round: u64) -> impl Iterator<Item = NodeId> + '_ {
267        self.receivedfrom.get(&round).into_iter().flatten().copied()
268    }
269
270    /// How many rounds hold state. Bounded by the membership; see the space note above.
271    pub fn rounds_recorded(&self) -> usize {
272        self.receivedfrom.len().max(self.proposals.len())
273    }
274
275    /// Every entry held across every round — the measure a bounded-space test asserts on.
276    pub fn state_entries(&self) -> usize {
277        let heard: usize = self.receivedfrom.values().map(|s| s.len()).sum();
278        let props: usize = self.proposals.values().map(|s| s.len()).sum();
279        heard + props + self.correct.len()
280    }
281}
282
283impl<P: Clone + Ord, L: VolatileLink<Flood<P>>> FloodingConsensus<P, L> {
284    /// Run the broadcast child, then act on what it reported.
285    fn with_beb(
286        &mut self,
287        cx: &mut ProtoCx<'_, Self>,
288        f: impl FnOnce(
289            &mut BestEffortBroadcast<Flood<P>, L>,
290            &mut ProtoCx<'_, BestEffortBroadcast<Flood<P>, L>>,
291        ),
292    ) {
293        let mut inds = self.beb.run(cx, Wire::Broadcast, f);
294        for ind in inds.drain(..) {
295            // The broadcast beneath reports scope boundaries only over a link that raises them.
296            // This layer is not yet parameterised over its link, so its child is the perfect link
297            // and a boundary cannot arrive. Named rather than dropped: silently absorbing a scope
298            // end is the failure `docs/conditional-guarantees.md` calls cardinal, and when this
299            // layer gains its link parameter the arm becomes real handling.
300            let beb::Ind::Deliver { from, msg } = ind else {
301                unreachable!("this link raises no scope boundary")
302            };
303            self.on_beb_deliver(from, msg, cx);
304        }
305        self.beb.reclaim(inds);
306        self.check_round(cx);
307    }
308
309    /// Run the detector child, then act on what it reported.
310    fn with_detector(
311        &mut self,
312        cx: &mut ProtoCx<'_, Self>,
313        f: impl FnOnce(&mut PerfectFailureDetector, &mut ProtoCx<'_, PerfectFailureDetector>),
314    ) {
315        let mut inds = self.detector.run(cx, Wire::Detector, f);
316        for pfd::Ind::Crash { node } in inds.drain(..) {
317            // `upon event ⟨ P, Crash | p ⟩ do correct := correct \ {p}`. Permanent: this
318            // detector has no Restore, and nothing here would act on one if it did.
319            self.correct.remove(&node);
320        }
321        self.detector.reclaim(inds);
322        // A crash alone can complete a round, by shrinking `correct` to a set already heard
323        // from. That is why the guard is checked here and not only on the message path.
324        self.check_round(cx);
325    }
326
327    /// `upon event ⟨ beb, Deliver | p, ... ⟩`.
328    fn on_beb_deliver(&mut self, from: NodeId, msg: Flood<P>, cx: &mut ProtoCx<'_, Self>) {
329        match msg {
330            Flood::Proposal { round, proposals } => {
331                self.receivedfrom.entry(round).or_default().insert(from);
332                self.proposals.entry(round).or_default().extend(proposals);
333            }
334            // `such that p ∈ correct ∧ decision = ⊥`. A process that wrongly suspects the
335            // decider discards this, which is half of why a false suspicion is unrecoverable.
336            Flood::Decided(v) if self.correct.contains(&from) && self.decision.is_none() => {
337                self.decide(v, cx);
338            }
339            Flood::Decided(_) => {}
340        }
341    }
342
343    /// `upon correct ⊆ receivedfrom[round] ∧ decision = ⊥`.
344    ///
345    /// A standing condition over state, not an event handler: its inputs change both when a
346    /// message arrives and when `correct` shrinks, so it is evaluated from both paths.
347    ///
348    /// Looped, because advancing a round can immediately satisfy the guard again — messages for
349    /// a later round may already have arrived. It terminates because the guard requires this
350    /// process to be in `receivedfrom[round]`, and this process only appears there once its own
351    /// broadcast for that round has come back to it.
352    ///
353    /// Calling it from the detector path is load-bearing and easy to think redundant. It is not:
354    /// after a crash, no further consensus message need ever arrive, so the message path may
355    /// never run again. Under this stack it *looks* redundant, because the stubborn link
356    /// retransmits for ever and its timer re-enters the broadcast child often enough to
357    /// re-evaluate the guard by accident. That is a property of the link, not of this layer, and
358    /// a link that does not retransmit — the session link, for one — would remove it. The test
359    /// `a_round_completes_on_a_crash_indication_alone` drives the protocol directly rather than
360    /// through a run, for exactly this reason.
361    fn check_round(&mut self, cx: &mut ProtoCx<'_, Self>) {
362        while self.decision.is_none() && self.round_complete() {
363            if self.heard(self.round) == self.heard(self.round - 1) {
364                // Nobody new crashed during the round, so every process that reaches its end
365                // holds the same proposal set, and `min` needs no further communication.
366                let Some(v) = self.proposals.get(&self.round).and_then(|p| p.first()).cloned()
367                else {
368                    // Vacuous: the guard is only satisfiable once this process has heard its own
369                    // proposal for the round, which carries at least one value.
370                    debug_assert!(false, "a completed round must hold at least one proposal");
371                    return;
372                };
373                self.decide(v, cx);
374            } else {
375                self.round += 1;
376                // `[PROPOSAL, round, proposals[round − 1]]` — the *previous* round's set, under
377                // the new round's number. Rendering this as `proposals[round]` compiles and is
378                // silently wrong until a crash cascade exposes it.
379                let carried = self.heard_proposals(self.round - 1);
380                self.send(Flood::Proposal { round: self.round, proposals: carried }, cx);
381            }
382        }
383    }
384
385    fn round_complete(&self) -> bool {
386        let heard = self.receivedfrom.get(&self.round);
387        self.correct.iter().all(|p| heard.is_some_and(|h| h.contains(p)))
388    }
389
390    fn heard(&self, round: u64) -> BTreeSet<NodeId> {
391        self.receivedfrom.get(&round).cloned().unwrap_or_default()
392    }
393
394    fn heard_proposals(&self, round: u64) -> BTreeSet<P> {
395        self.proposals.get(&round).cloned().unwrap_or_default()
396    }
397
398    fn decide(&mut self, v: P, cx: &mut ProtoCx<'_, Self>) {
399        self.decision = Some(v.clone());
400        self.send(Flood::Decided(v.clone()), cx);
401        cx.indicate(Ind::Decide(v));
402    }
403
404    fn send(&mut self, msg: Flood<P>, cx: &mut ProtoCx<'_, Self>) {
405        // Re-enters the child while its inbox is out on loan, so `run` hands back a fresh one.
406        let inds =
407            self.beb.run(cx, Wire::Broadcast, |beb, ccx| beb.on_cmd(beb::Cmd::Broadcast(msg), ccx));
408        debug_assert!(
409            inds.is_empty(),
410            "sending must not deliver synchronously; if it does, check_round can recurse"
411        );
412        self.beb.reclaim(inds);
413    }
414}
415
416impl<P: Clone + Ord, L: VolatileLink<Flood<P>>> Protocol for FloodingConsensus<P, L> {
417    type Cmd = Cmd<P>;
418    type Ind = Ind<P>;
419    type Msg = Wire<L::Msg>;
420    /// No session beneath, so no scope end can be constructed — as for both children.
421    type Scope = core::convert::Infallible;
422    type Note = crate::Note;
423    /// Keeps nothing durably: a crash loses everything this protocol knows.
424    type Meta = core::convert::Infallible;
425    type Entry = core::convert::Infallible;
426
427    /// Failure detection begins here, as Module 2.6 has it. It used to need a `Start` command
428    /// because there was no init event to hang the detector's first timer on.
429    fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
430        self.with_detector(cx, |d, ccx| d.on_init(ccx));
431    }
432
433    fn on_cmd(&mut self, cmd: Cmd<P>, cx: &mut ProtoCx<'_, Self>) {
434        match cmd {
435            // One proposal per process, one decision per instance; a second is not a second
436            // consensus and is ignored.
437            Cmd::Propose(_) if self.proposed => {}
438            Cmd::Propose(v) => {
439                self.proposed = true;
440                self.proposals.entry(1).or_default().insert(v);
441                let ps = self.heard_proposals(1);
442                self.send(Flood::Proposal { round: 1, proposals: ps }, cx);
443            }
444        }
445    }
446
447    fn on_msg(&mut self, from: NodeId, msg: Self::Msg, cx: &mut ProtoCx<'_, Self>) {
448        match msg {
449            Wire::Broadcast(m) => self.with_beb(cx, |beb, ccx| beb.on_msg(from, m, ccx)),
450            Wire::Detector(h) => self.with_detector(cx, |d, ccx| d.on_msg(from, h, ccx)),
451        }
452    }
453
454    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
455        // Handed to both children: the identity does not say which registered it, and the one that
456        // did not will recognise that and do nothing.
457        self.with_beb(cx, |beb, ccx| beb.on_timer(id, ccx));
458        self.with_detector(cx, |d, ccx| d.on_timer(id, ccx));
459    }
460}