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}