Skip to main content

recon_protocols/
lazy_probabilistic_broadcast.rs

1//! Lazy probabilistic broadcast — gossip, then pull back what it missed.
2//!
3//! **Status: implementation. Space: bounded by a retention window.**
4//!
5//! Cachin, Guerraoui & Rodrigues, Module 3.7 and **Algorithms 3.10 and 3.11**. The book splits it
6//! in two — a data half and a recovery half — and this module quotes both, verbatim from the source
7//! rather than from any recollection of it.
8//!
9//! ```text
10//! Algorithm 3.10: Lazy Probabilistic Broadcast (part 1, data dissemination)
11//! Implements: ProbabilisticBroadcast, instance pb.
12//! Uses:
13//!     FairLossPointToPointLinks, instance fll;
14//!     ProbabilisticBroadcast, instance upb.        // an unreliable implementation
15//!
16//! upon event ⟨ pb, Init ⟩ do
17//!     next := [1]^N; lsn := 0; pending := ∅; stored := ∅;
18//!
19//! procedure gossip(msg) is
20//!     forall t ∈ picktargets(k) do trigger ⟨ fll, Send | t, msg ⟩;
21//!
22//! upon event ⟨ pb, Broadcast | m ⟩ do
23//!     lsn := lsn + 1;
24//!     trigger ⟨ upb, Broadcast | [DATA, self, m, lsn] ⟩;
25//!
26//! upon event ⟨ upb, Deliver | p, [DATA, s, m, sn] ⟩ do
27//!     if random([0, 1]) > α then
28//!         stored := stored ∪ {[DATA, s, m, sn]};
29//!     if sn = next[s] then
30//!         next[s] := next[s] + 1;
31//!         trigger ⟨ pb, Deliver | s, m ⟩;
32//!     else if sn > next[s] then
33//!         pending := pending ∪ {[DATA, s, m, sn]};
34//!         forall missing ∈ [next[s], . . . , sn − 1] do
35//!             if no m′ exists such that [DATA, s, m′, missing] ∈ pending then
36//!                 gossip([REQUEST, self, s, missing, R − 1]);
37//!         starttimer(Δ, s, sn);
38//! ```
39//!
40//! ```text
41//! Algorithm 3.11: Lazy Probabilistic Broadcast (part 2, recovery)
42//!
43//! upon event ⟨ fll, Deliver | p, [REQUEST, q, s, sn, r] ⟩ do
44//!     if exists m such that [DATA, s, m, sn] ∈ stored then
45//!         trigger ⟨ fll, Send | q, [DATA, s, m, sn] ⟩;
46//!     else if r > 0 then
47//!         gossip([REQUEST, q, s, sn, r − 1]);
48//!
49//! upon event ⟨ fll, Deliver | p, [DATA, s, m, sn] ⟩ do
50//!     pending := pending ∪ {[DATA, s, m, sn]};
51//!
52//! upon exists [DATA, s, x, sn] ∈ pending such that sn = next[s] do
53//!     next[s] := next[s] + 1;
54//!     pending := pending \ {[DATA, s, x, sn]};
55//!     trigger ⟨ pb, Deliver | s, x ⟩;
56//!
57//! upon event ⟨ Timeout | s, sn ⟩ do
58//!     if sn > next[s] then
59//!         next[s] := sn + 1;
60//! ```
61//!
62//! # Two children, and why the second one matters
63//!
64//! `Uses:` names both `fll` and `upb`. Data is gossiped by the unreliable broadcast beneath;
65//! **requests and their answers travel directly over the link**, bypassing the gossip. That is what
66//! makes the second phase a *pull* and the algorithm lazy. Routing a request through `upb` would
67//! flood the membership to repair one process's gap, which is exactly the cost this phase exists to
68//! avoid. So this layer multiplexes two children onto one wire, as
69//! [`crate::uniform_reliable_broadcast`] does for its broadcast and its detector.
70//!
71//! # Three readings the page settled, each of which was about to go the other way
72//!
73//! - **`next := [1]^N`.** Sequence numbers start at one. A zero-based `next` leaves every process
74//!   waiting for a message no sender ever sends.
75//! - **The timeout skips *past* the gap.** `if sn > next[s] then next[s] := sn + 1` abandons the
76//!   message at `sn` too, not just those before it. Setting `next[s] := sn` would deliver a message
77//!   the process has already given up on.
78//! - **Draining `pending` is a standing condition.** `upon exists … such that sn = next[s]` is
79//!   re-evaluated whenever `next` or `pending` changes, so closing one gap can release a long run at
80//!   once. Written here as a loop after every mutation of either, which is the same thing.
81//!
82//! # The α, which the book states twice and inconsistently
83//!
84//! The pseudocode stores when `random([0,1]) > α`, so α is the probability of *not* storing. Page 99
85//! says in prose that a process stores "with probability α", which is the opposite. Page 100 breaks
86//! the tie: it describes setting `α = 0` as every process storing, which only holds under the
87//! pseudocode's reading.
88//!
89//! `docs/postmortem.md` disagrees with itself on this too — its re-examination reaches the same
90//! conclusion, and its own worked sketch writes `gen_bool(alpha)`. This module ends the question by
91//! not using α at all: [`Config::store_probability`] is the probability of storing, named for what
92//! it does, and the book's α is one minus it.
93//!
94//! # What this buys, and what it costs
95//!
96//! ```text
97//! PB1 [probabilistic]  Probabilistic validity — strictly better than the eager algorithm's under
98//!                      loss, because a gap is repaired rather than lost
99//! PB2 [window]         No duplication — within the retention window
100//! PB3 [always]         No creation
101//! ```
102//!
103//! Recovery depends on some reachable process having stored the message, so `PB1` here is
104//! conditional on `store_probability` and on that process being reachable — not absolute. A gap
105//! nobody stored is skipped by the timeout, which converts a permanent stall into a lost message,
106//! and a stall would be the worse outcome.
107//!
108//! # The retention window, which is this project's and not the book's
109//!
110//! Page 100: "garbage collection of the stored message copies is omitted in the pseudo code for
111//! simplicity." Both `stored` and `pending` are bounded here by a per-sender window and evicted on
112//! insert, for the reasons [`crate::probabilistic_broadcast`] gives at length. A request for
113//! something evicted is answered as unavailable, and the requester's timeout moves it past the gap.
114//!
115//! # Identity is scoped to the originator's incarnation — departure
116//!
117//! The book's `s` is a process, and `next[s]`, `pending` and `stored` are keyed by it. `lsn` is
118//! volatile, so a process that crashes and comes back numbers its messages from one again — and
119//! every receiver, holding `next[s] = 4`, would drop its first three as already delivered, silently,
120//! through the `sn < next[s]` case the pseudocode does not even write. Under the book's crash-stop
121//! model that case never arises; in the real-world set it is the first thing a restart does.
122//!
123//! So the sender of a [`Data`] is a [`Sender`] — the originator **and its incarnation**, a value
124//! drawn from the seeded generator at `Init` exactly as [`crate::probabilistic_broadcast`] draws
125//! its own — and every per-sender structure is keyed by that. A restarted originator is a new
126//! sender with `next = 1`, and its messages are delivered.
127//!
128//! What bounds it: a receiver remembers the **two most recent incarnations** of each originator, and
129//! admitting a third retires the oldest — its `next`, its pending and stored messages, its timers.
130//! Two rather than one because relayed copies from the incarnation just retired can still be
131//! arriving while the new one's begin, and a one-deep memory would flip between them, losing both.
132//! Two rather than more because a process has one live incarnation and at most one being retired;
133//! a message from an incarnation older than that is a straggler this abstraction may lose. State is
134//! therefore bounded by `2 × membership × window`, and a restart costs one purge, not a leak.
135
136use core::time::Duration;
137use recon_core::{Child, NodeId, ProtoCx, Protocol, TimerId};
138use serde::{Deserialize, Serialize};
139use std::collections::{BTreeMap, BTreeSet, VecDeque};
140
141use crate::fair_loss_link::FairLossLink;
142use crate::link::{Boundary, LinkInd, VolatileLink};
143use crate::probabilistic_broadcast::{self as pb, ProbabilisticBroadcast};
144
145/// `[DATA, s, m, sn]` — this layer's header, carried as the gossip's payload.
146#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
147pub struct Data<P> {
148    /// `s` — who originated it.
149    pub origin: NodeId,
150    /// Which incarnation of `origin`. See the module note on identity.
151    pub incarnation: u64,
152    /// `sn` — that sender's sequence number for it.
153    pub seq: u64,
154    pub payload: P,
155}
156
157impl<P> Data<P> {
158    /// The book's `s`, as this module keys on it: the originator in a particular incarnation.
159    pub fn sender(&self) -> Sender {
160        Sender { origin: self.origin, incarnation: self.incarnation }
161    }
162}
163
164/// An originator in one incarnation — what `next`, `pending` and `stored` are keyed by.
165#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
166pub struct Sender {
167    pub origin: NodeId,
168    pub incarnation: u64,
169}
170
171/// How many incarnations of one originator a receiver keeps state for. See the module note.
172const INCARNATIONS_REMEMBERED: usize = 2;
173
174/// What travels over the link directly, outside the gossip: `[REQUEST, …]` and its answer.
175#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
176pub enum Recovery<P> {
177    /// `[REQUEST, q, s, sn, r]` — `q` is the process that wants it, carried so that whoever holds
178    /// it answers the requester rather than the relayer.
179    Request { requester: NodeId, origin: NodeId, incarnation: u64, seq: u64, ttl: u32 },
180    /// `[DATA, s, m, sn]` sent back to a requester.
181    Data(Data<P>),
182}
183
184/// The wire, multiplexing the two children Algorithm 3.10 names.
185#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
186pub enum Wire<G, R> {
187    /// The gossip child's traffic.
188    Gossip(G),
189    /// Recovery traffic, which deliberately does not go through the gossip.
190    Recovery(R),
191}
192
193/// Requests from the layer above.
194#[derive(Debug, Clone, PartialEq, Eq)]
195pub enum Cmd<P> {
196    Broadcast(P),
197}
198
199/// Indications to the layer above.
200#[derive(Debug, Clone, PartialEq, Eq)]
201pub enum Ind<P> {
202    /// `from` is the originator, never a relayer, and deliveries from one sender are in sequence.
203    Deliver { from: NodeId, msg: P },
204    /// The scope with `peer` ended at `epoch`. Raised only over a link that reports boundaries.
205    SessionEnded { peer: NodeId, epoch: u64 },
206    /// A scope with `peer` is in force at `epoch`.
207    SessionEstablished { peer: NodeId, epoch: u64 },
208}
209
210/// How this instance recovers.
211#[derive(Debug, Clone, Copy, PartialEq)]
212pub struct Config {
213    /// The gossip beneath: its fanout, rounds and window.
214    pub gossip: pb::Config,
215    /// How likely a process is to keep a copy for answering requests.
216    ///
217    /// **The book's α is one minus this.** Named for what it does; see the module note. At `1.0`
218    /// every process stores everything, which is the certain and expensive case the book describes
219    /// as `α = 0`.
220    pub store_probability: f64,
221    /// `R` for a request — how far a request is relayed before it is abandoned.
222    pub request_rounds: u32,
223    /// `Δ` — how long to wait for a gap before giving up on it.
224    pub gap_timeout: Duration,
225    /// How many messages to keep in `stored` and `pending`, per sender.
226    pub window: usize,
227}
228
229/// The gossip this layer rides on: Algorithm 3.9 carrying [`Data`], over `G` — the fair-loss link
230/// the book names unless the caller says otherwise.
231pub type Gossiper<P, G = FairLossLink<pb::Carried<Data<P>>>> = ProbabilisticBroadcast<Data<P>, G>;
232
233/// Gossip with recovery.
234///
235/// Two links, both parameters: `L` carries the recovery traffic and `G` carries the gossip. Over
236/// sessions both are session links — two instances each holding one epoch per peer, both handed
237/// every scope event, on one wire — which is how [`crate::uniform_reliable_broadcast`] already puts
238/// a broadcast and a detector together. Their scopes must agree (`G::Scope = L::Scope`), because a
239/// scope event reaching this layer is one event about one session and goes to both.
240pub struct LazyProbabilisticBroadcast<
241    P,
242    L = FairLossLink<Recovery<P>>,
243    G = FairLossLink<pb::Carried<Data<P>>>,
244> where
245    P: Clone + serde::Serialize + serde::de::DeserializeOwned,
246    L: VolatileLink<Recovery<P>>,
247    L::Scope: Clone,
248    G: VolatileLink<pb::Carried<Data<P>>, Scope = L::Scope>,
249{
250    me: NodeId,
251    peers: BTreeSet<NodeId>,
252    config: Config,
253    /// This incarnation's name, drawn at `Init`. Zero until then.
254    incarnation: u64,
255    /// `lsn` — this process's own sequence counter.
256    lsn: u64,
257    /// The incarnations of each originator this process keeps state for, oldest first. At most
258    /// [`INCARNATIONS_REMEMBERED`]; admitting another retires the oldest.
259    incarnations: BTreeMap<NodeId, VecDeque<u64>>,
260    /// `next[s]` — the next sequence number expected from `s`. Absent means one, per `[1]^N`.
261    next: BTreeMap<Sender, u64>,
262    /// `pending` — received, but ahead of a gap.
263    pending: BTreeMap<(Sender, u64), P>,
264    /// `stored` — kept so this process can answer a request.
265    stored: BTreeMap<(Sender, u64), P>,
266    /// Insertion order per sender, so both collections evict in constant time.
267    pending_order: BTreeMap<Sender, VecDeque<u64>>,
268    stored_order: BTreeMap<Sender, VecDeque<u64>>,
269    /// Which gap each outstanding timer is waiting on.
270    timers: BTreeMap<TimerId, (Sender, u64)>,
271    upb: Child<Gossiper<P, G>>,
272    link: Child<L>,
273}
274
275impl<P> LazyProbabilisticBroadcast<P, FairLossLink<Recovery<P>>>
276where
277    P: Clone + serde::Serialize + serde::de::DeserializeOwned,
278{
279    /// Lazy probabilistic broadcast among `peers`, over the fair-loss links the book names.
280    pub fn new(me: NodeId, peers: impl IntoIterator<Item = NodeId>, config: Config) -> Self {
281        Self::with_link(me, peers, FairLossLink::new(), config)
282    }
283}
284
285impl<P, L> LazyProbabilisticBroadcast<P, L>
286where
287    P: Clone + serde::Serialize + serde::de::DeserializeOwned,
288    L: VolatileLink<Recovery<P>, Scope = core::convert::Infallible>,
289{
290    /// Lazy probabilistic broadcast, over the link supplied for its recovery traffic and the
291    /// book's fair-loss link for the gossip.
292    ///
293    /// Only for a recovery link that reports no boundary: the gossip beneath runs over a fair-loss
294    /// link here, and the two links' scopes have to agree. Over sessions use
295    /// [`LazyProbabilisticBroadcast::with_links`] with a session link for both.
296    pub fn with_link(
297        me: NodeId,
298        peers: impl IntoIterator<Item = NodeId>,
299        link: L,
300        config: Config,
301    ) -> Self {
302        Self::with_links(me, peers, FairLossLink::new(), link, config)
303    }
304}
305
306impl<P, L, G> LazyProbabilisticBroadcast<P, L, G>
307where
308    P: Clone + serde::Serialize + serde::de::DeserializeOwned,
309    L: VolatileLink<Recovery<P>>,
310    L::Scope: Clone,
311    G: VolatileLink<pb::Carried<Data<P>>, Scope = L::Scope>,
312{
313    /// Lazy probabilistic broadcast over two links: `gossip` beneath the eager broadcast that
314    /// disseminates data, `recovery` for the requests and answers that repair gaps.
315    pub fn with_links(
316        me: NodeId,
317        peers: impl IntoIterator<Item = NodeId>,
318        gossip: G,
319        recovery: L,
320        config: Config,
321    ) -> Self {
322        let mut peers: BTreeSet<NodeId> = peers.into_iter().collect();
323        peers.insert(me);
324        LazyProbabilisticBroadcast {
325            me,
326            peers: peers.clone(),
327            config,
328            incarnation: 0,
329            lsn: 0,
330            incarnations: BTreeMap::new(),
331            next: BTreeMap::new(),
332            pending: BTreeMap::new(),
333            stored: BTreeMap::new(),
334            pending_order: BTreeMap::new(),
335            stored_order: BTreeMap::new(),
336            timers: BTreeMap::new(),
337            upb: Child::new(ProbabilisticBroadcast::with_link(me, peers, gossip, config.gossip)),
338            link: Child::new(recovery),
339        }
340    }
341
342    /// The link the gossip travels over.
343    pub fn gossip_link(&self) -> &G {
344        self.upb.link()
345    }
346
347    /// The link the recovery traffic travels over.
348    pub fn recovery_link(&self) -> &L {
349        &self.link
350    }
351
352    /// `next[s]`, which is one until this process has delivered anything from `s`.
353    pub fn next_expected_of(&self, sender: Sender) -> u64 {
354        self.next.get(&sender).copied().unwrap_or(1)
355    }
356
357    /// `next[s]` for the incarnation of `from` most recently heard from, or one if none has been.
358    pub fn next_expected(&self, from: NodeId) -> u64 {
359        self.latest(from).map(|s| self.next_expected_of(s)).unwrap_or(1)
360    }
361
362    /// The incarnation of `from` most recently admitted, if any.
363    pub fn latest(&self, from: NodeId) -> Option<Sender> {
364        self.incarnations
365            .get(&from)
366            .and_then(|v| v.back())
367            .map(|incarnation| Sender { origin: from, incarnation: *incarnation })
368    }
369
370    /// How many incarnations of `from` this process keeps state for.
371    pub fn incarnations_of(&self, from: NodeId) -> usize {
372        self.incarnations.get(&from).map(|v| v.len()).unwrap_or(0)
373    }
374
375    /// How many messages this process is holding ahead of a gap.
376    pub fn pending_count(&self) -> usize {
377        self.pending.len()
378    }
379
380    /// How many copies this process is holding for answering requests.
381    pub fn stored_count(&self) -> usize {
382        self.stored.len()
383    }
384
385    /// Whether this process could answer a request for `seq` from the incarnation of `origin`
386    /// most recently heard from.
387    pub fn has_stored(&self, origin: NodeId, seq: u64) -> bool {
388        self.latest(origin).is_some_and(|s| self.stored.contains_key(&(s, seq)))
389    }
390}
391
392impl<P: Clone, L, G> LazyProbabilisticBroadcast<P, L, G>
393where
394    L: VolatileLink<Recovery<P>>,
395    L::Scope: Clone,
396    G: VolatileLink<pb::Carried<Data<P>>, Scope = L::Scope>,
397    P: serde::Serialize + serde::de::DeserializeOwned,
398{
399    /// Run the gossip child, then act on what it reported.
400    fn through_upb(
401        &mut self,
402        cx: &mut ProtoCx<'_, Self>,
403        f: impl FnOnce(&mut Gossiper<P, G>, &mut ProtoCx<'_, Gossiper<P, G>>),
404    ) {
405        let mut inds = self.upb.run(cx, Wire::Gossip, f);
406        for ind in inds.drain(..) {
407            match ind {
408                pb::Ind::Deliver { msg, .. } => self.on_upb_deliver(msg, cx),
409                // A boundary the gossip's link observed. The recovery link observed the same one —
410                // both links are handed every scope event, and their scopes are one type by the
411                // bound on `G` — so it is reported upward from `through_link` and once. Reporting
412                // it here too would tell the layer above one session ended twice.
413                pb::Ind::SessionEnded { .. } | pb::Ind::SessionEstablished { .. } => {}
414            }
415        }
416        self.upb.reclaim(inds);
417    }
418
419    /// Run the link, then act on the recovery traffic it reported.
420    fn through_link(
421        &mut self,
422        cx: &mut ProtoCx<'_, Self>,
423        f: impl FnOnce(&mut L, &mut ProtoCx<'_, L>),
424    ) {
425        let mut inds = self.link.run(cx, Wire::Recovery, f);
426        for ind in inds.drain(..) {
427            match L::classify(ind) {
428                LinkInd::Deliver { msg, .. } => self.on_recovery(msg, cx),
429                LinkInd::Boundary(Boundary::Ended { peer, epoch }) => {
430                    cx.indicate(Ind::SessionEnded { peer, epoch })
431                }
432                LinkInd::Boundary(Boundary::Established { peer, epoch }) => {
433                    cx.indicate(Ind::SessionEstablished { peer, epoch })
434                }
435            }
436        }
437        self.link.reclaim(inds);
438    }
439
440    /// `upon event ⟨ upb, Deliver | p, [DATA, s, m, sn] ⟩` — Algorithm 3.10.
441    fn on_upb_deliver(&mut self, data: Data<P>, cx: &mut ProtoCx<'_, Self>) {
442        use rand::Rng;
443
444        let sender = data.sender();
445        self.admit(sender);
446
447        // `if random([0, 1]) > α then stored := stored ∪ {…}`, with α restated as its complement.
448        // Before the sequence checks, so a message this process cannot deliver yet is still one it
449        // can answer a request for.
450        if cx.rng().random_bool(self.config.store_probability) {
451            self.store(data.clone());
452        }
453
454        let next = self.next_expected_of(sender);
455        if data.seq == next {
456            self.next.insert(sender, next + 1);
457            cx.indicate(Ind::Deliver { from: data.origin, msg: data.payload });
458            self.drain_pending(sender, cx);
459        } else if data.seq > next {
460            // `forall missing ∈ [next[s], …, sn − 1] do if no m′ … ∈ pending then gossip(REQUEST)`.
461            // The guard is `pending`, and nothing else: the book re-requests a gap on every
462            // out-of-order arrival that does not already have it pending. That looks like a defect
463            // and was once reported as one; it is the page.
464            for missing in next..data.seq {
465                if !self.pending.contains_key(&(sender, missing)) {
466                    self.request(sender, missing, cx);
467                }
468            }
469            let seq = data.seq;
470            self.hold(data);
471            // `starttimer(Δ, s, sn)` — one timer per gap, remembered by its handle so the expiry
472            // can be matched back to the gap it was waiting on.
473            let id = cx.set_timer(self.config.gap_timeout);
474            self.timers.insert(id, (sender, seq));
475        }
476    }
477
478    /// Note an incarnation of an originator, retiring the oldest if this makes one too many. See
479    /// the module note on identity for why two, and what retiring costs.
480    fn admit(&mut self, sender: Sender) {
481        let known = self.incarnations.entry(sender.origin).or_default();
482        if known.contains(&sender.incarnation) {
483            return;
484        }
485        known.push_back(sender.incarnation);
486        if known.len() > INCARNATIONS_REMEMBERED
487            && let Some(retired) = known.pop_front()
488        {
489            let retired = Sender { origin: sender.origin, incarnation: retired };
490            self.next.remove(&retired);
491            self.pending.retain(|(s, _), _| *s != retired);
492            self.stored.retain(|(s, _), _| *s != retired);
493            self.pending_order.remove(&retired);
494            self.stored_order.remove(&retired);
495            self.timers.retain(|_, (s, _)| *s != retired);
496        }
497    }
498
499    /// `upon event ⟨ fll, Deliver | p, [REQUEST | DATA, …] ⟩` — Algorithm 3.11.
500    fn on_recovery(&mut self, r: Recovery<P>, cx: &mut ProtoCx<'_, Self>) {
501        match r {
502            Recovery::Request { requester, origin, incarnation, seq, ttl } => {
503                let sender = Sender { origin, incarnation };
504                if let Some(payload) = self.stored.get(&(sender, seq)).cloned() {
505                    // `trigger ⟨ fll, Send | q, [DATA, s, m, sn] ⟩` — to the requester, not the
506                    // relayer, which is why `q` travels in the request.
507                    let answer = Recovery::Data(Data { origin, incarnation, seq, payload });
508                    self.send_to(requester, answer, cx);
509                } else if ttl > 0 {
510                    // `else if r > 0 then gossip([REQUEST, q, s, sn, r − 1])` — `q` is preserved.
511                    self.gossip_request(
512                        Recovery::Request { requester, origin, incarnation, seq, ttl: ttl - 1 },
513                        cx,
514                    );
515                }
516            }
517            // `upon event ⟨ fll, Deliver | p, [DATA, s, m, sn] ⟩ do pending := pending ∪ {…}`.
518            // A recovered message joins `pending` and is released by the standing condition, which
519            // is what lets it close a gap without a second code path for delivery.
520            Recovery::Data(data) => {
521                let sender = data.sender();
522                self.admit(sender);
523                self.hold(data);
524                self.drain_pending(sender, cx);
525            }
526        }
527    }
528
529    /// `upon exists [DATA, s, x, sn] ∈ pending such that sn = next[s]` — the standing condition.
530    ///
531    /// A loop rather than a single step: closing one gap can release an arbitrarily long run, and
532    /// the book's `upon` is re-evaluated after every change.
533    fn drain_pending(&mut self, from: Sender, cx: &mut ProtoCx<'_, Self>) {
534        while let Some(payload) = self.pending.remove(&(from, self.next_expected_of(from))) {
535            let seq = self.next_expected_of(from);
536            self.next.insert(from, seq + 1);
537            if let Some(order) = self.pending_order.get_mut(&from) {
538                order.retain(|s| *s != seq);
539            }
540            cx.indicate(Ind::Deliver { from: from.origin, msg: payload });
541        }
542    }
543
544    /// `gossip([REQUEST, self, s, missing, R − 1])` — over the link, not through the gossip child.
545    fn request(&mut self, from: Sender, seq: u64, cx: &mut ProtoCx<'_, Self>) {
546        let r = Recovery::Request {
547            requester: self.me,
548            origin: from.origin,
549            incarnation: from.incarnation,
550            seq,
551            ttl: self.config.request_rounds.saturating_sub(1),
552        };
553        self.gossip_request(r, cx);
554    }
555
556    /// `procedure gossip(msg)` for recovery traffic: `picktargets(k)` over the link.
557    fn gossip_request(&mut self, r: Recovery<P>, cx: &mut ProtoCx<'_, Self>) {
558        for target in self.picktargets(cx) {
559            self.send_to(target, r.clone(), cx);
560        }
561    }
562
563    /// `trigger ⟨ fll, Send | t, msg ⟩`.
564    fn send_to(&mut self, to: NodeId, r: Recovery<P>, cx: &mut ProtoCx<'_, Self>) {
565        let inds = self.link.run(cx, Wire::Recovery, |link, ccx| link.on_cmd(L::send(to, r), ccx));
566        debug_assert!(inds.is_empty(), "a send must not deliver synchronously");
567        self.link.reclaim(inds);
568    }
569
570    /// `picktargets(k)` — the same uniform draw without replacement the gossip uses.
571    fn picktargets(&self, cx: &mut ProtoCx<'_, Self>) -> Vec<NodeId> {
572        use rand::Rng;
573        let mut candidates: Vec<NodeId> =
574            self.peers.iter().copied().filter(|p| *p != self.me).collect();
575        let take = self.config.gossip.fanout.min(candidates.len());
576        for i in 0..take {
577            let j = i + cx.rng().random_range(0..candidates.len() - i);
578            candidates.swap(i, j);
579        }
580        candidates.truncate(take);
581        candidates
582    }
583
584    /// `pending := pending ∪ {[DATA, s, m, sn]}`, bounded by the window.
585    fn hold(&mut self, data: Data<P>) {
586        let sender = data.sender();
587        if self.pending.insert((sender, data.seq), data.payload).is_none() {
588            let order = self.pending_order.entry(sender).or_default();
589            order.push_back(data.seq);
590            if order.len() > self.config.window
591                && let Some(evicted) = order.pop_front()
592            {
593                self.pending.remove(&(sender, evicted));
594            }
595        }
596    }
597
598    /// `stored := stored ∪ {[DATA, s, m, sn]}`, bounded by the window.
599    fn store(&mut self, data: Data<P>) {
600        let sender = data.sender();
601        if self.stored.insert((sender, data.seq), data.payload).is_none() {
602            let order = self.stored_order.entry(sender).or_default();
603            order.push_back(data.seq);
604            if order.len() > self.config.window
605                && let Some(evicted) = order.pop_front()
606            {
607                self.stored.remove(&(sender, evicted));
608            }
609        }
610    }
611}
612
613impl<P, L, G> Protocol for LazyProbabilisticBroadcast<P, L, G>
614where
615    P: Clone + serde::Serialize + serde::de::DeserializeOwned,
616    L: VolatileLink<Recovery<P>>,
617    L::Scope: Clone,
618    G: VolatileLink<pb::Carried<Data<P>>, Scope = L::Scope>,
619{
620    type Cmd = Cmd<P>;
621    type Ind = Ind<P>;
622    type Msg = Wire<G::Msg, L::Msg>;
623    type Scope = L::Scope;
624    type Note = crate::Note;
625    /// Keeps nothing durably.
626    type Meta = core::convert::Infallible;
627    type Entry = core::convert::Infallible;
628
629    /// `upon event ⟨ pb, Broadcast | m ⟩ do lsn := lsn + 1; trigger ⟨ upb, Broadcast | [DATA, …] ⟩`.
630    ///
631    /// Note what is *not* here: no delivery to self. The eager child beneath delivers a broadcast
632    /// to its own process, and that arrives back through `on_upb_deliver` like any other, which is
633    /// what puts this process's own messages through the same sequence check as everyone else's.
634    fn on_cmd(&mut self, Cmd::Broadcast(payload): Cmd<P>, cx: &mut ProtoCx<'_, Self>) {
635        self.lsn += 1;
636        let data = Data { origin: self.me, incarnation: self.incarnation, seq: self.lsn, payload };
637        self.through_upb(cx, |upb, ccx| upb.on_cmd(pb::Cmd::Broadcast(data), ccx));
638    }
639
640    fn on_msg(&mut self, from: NodeId, msg: Self::Msg, cx: &mut ProtoCx<'_, Self>) {
641        match msg {
642            Wire::Gossip(m) => self.through_upb(cx, |upb, ccx| upb.on_msg(from, m, ccx)),
643            Wire::Recovery(m) => self.through_link(cx, |link, ccx| link.on_msg(from, m, ccx)),
644        }
645    }
646
647    /// `upon event ⟨ Timeout | s, sn ⟩ do if sn > next[s] then next[s] := sn + 1`.
648    ///
649    /// The gap is abandoned, and `sn` with it — the book skips *past* the message the timer was
650    /// waiting on, not to it. Whatever is now deliverable is released by the standing condition,
651    /// which is why the drain follows.
652    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
653        if let Some((sender, seq)) = self.timers.remove(&id) {
654            if seq > self.next_expected_of(sender) {
655                self.next.insert(sender, seq + 1);
656                self.drain_pending(sender, cx);
657            }
658            return;
659        }
660        // Not this layer's. Hand it to both children, since neither's expiry is distinguishable
661        // from the other's by its handle alone.
662        self.through_upb(cx, |upb, ccx| upb.on_timer(id, ccx));
663        self.through_link(cx, |link, ccx| link.on_timer(id, ccx));
664    }
665
666    /// Both children run over the same session, so both are told when it ends or begins.
667    fn on_scope_event(&mut self, scope: L::Scope, cx: &mut ProtoCx<'_, Self>) {
668        let for_gossip = scope.clone();
669        self.through_upb(cx, |upb, ccx| upb.on_scope_event(for_gossip, ccx));
670        self.through_link(cx, |link, ccx| link.on_scope_event(scope, ccx));
671    }
672
673    /// Name this incarnation, then start the children. Runs on every restart, which is the point.
674    fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
675        self.incarnation = cx.rng().next_u64();
676        self.through_upb(cx, |upb, ccx| upb.on_init(ccx));
677        self.through_link(cx, |link, ccx| link.on_init(ccx));
678    }
679}