Skip to main content

commonware_broadcast/buffered/
engine.rs

1use super::{metrics, Config, Mailbox, Message};
2use commonware_actor::mailbox;
3use commonware_codec::Codec;
4use commonware_cryptography::{Digestible, PublicKey};
5use commonware_macros::select_loop;
6use commonware_p2p::{
7    utils::codec::{wrap, WrappedSender},
8    Provider, Receiver, Recipients, Sender,
9};
10use commonware_runtime::{
11    spawn_cell,
12    telemetry::metrics::{status::Status, GaugeExt},
13    BufferPooler, Clock, ContextCell, Handle, Metrics, Spawner,
14};
15use commonware_utils::{
16    channel::{fallible::OneshotExt, oneshot},
17    ordered::Set,
18};
19use std::{
20    collections::{BTreeMap, VecDeque},
21    sync::Arc,
22};
23use tracing::{debug, error, trace, warn};
24
25/// A responder waiting for a message.
26struct Waiter<M> {
27    /// The responder to send the message to.
28    responder: oneshot::Sender<Arc<M>>,
29}
30
31/// Result of buffering an incoming or locally sent digest (inserted, duplicate, or ineligible).
32enum InsertMessageResult {
33    Inserted,
34    Duplicate,
35    Ineligible,
36}
37
38/// Instance of the main engine for the module.
39///
40/// It is responsible for:
41/// - Broadcasting messages to the network
42/// - Receiving messages from the network
43/// - Storing messages in the cache
44/// - Responding to requests from the application
45pub struct Engine<E, P, M, D>
46where
47    E: BufferPooler + Clock + Spawner + Metrics,
48    P: PublicKey,
49    M: Digestible + Codec,
50    D: Provider<PublicKey = P>,
51{
52    ////////////////////////////////////////
53    // Interfaces
54    ////////////////////////////////////////
55    context: ContextCell<E>,
56
57    ////////////////////////////////////////
58    // Configuration
59    ////////////////////////////////////////
60    /// My public key
61    public_key: P,
62
63    /// Whether messages are sent as priority
64    priority: bool,
65
66    /// Number of messages to cache per peer
67    deque_size: usize,
68
69    /// Configuration for decoding messages
70    codec_config: M::Cfg,
71
72    ////////////////////////////////////////
73    // Messaging
74    ////////////////////////////////////////
75    /// The mailbox for receiving messages.
76    mailbox_receiver: mailbox::Receiver<Message<P, M>>,
77
78    /// Pending requests from the application.
79    waiters: BTreeMap<M::Digest, Vec<Waiter<M>>>,
80
81    /// Provider for peer set changes.
82    peer_provider: D,
83
84    ////////////////////////////////////////
85    // Cache
86    ////////////////////////////////////////
87    /// All cached messages by digest.
88    items: BTreeMap<M::Digest, Arc<M>>,
89
90    /// A LRU cache of the latest received digests from each peer.
91    ///
92    /// This is used to limit the number of digests stored per peer.
93    /// At most `deque_size` digests are stored per peer. This value is expected to be small, so
94    /// membership checks are done in linear time.
95    deques: BTreeMap<P, VecDeque<M::Digest>>,
96
97    /// The number of times each digest (globally unique) exists in one of the deques.
98    ///
99    /// Multiple peers can send the same message and we only want to store
100    /// the message once.
101    counts: BTreeMap<M::Digest, usize>,
102
103    /// Latest primary peer set allowed to keep buffered messages resident.
104    latest_primary_peers: Set<P>,
105
106    ////////////////////////////////////////
107    // Metrics
108    ////////////////////////////////////////
109    /// Metrics
110    metrics: metrics::Metrics<P>,
111}
112
113impl<E, P, M, D> Engine<E, P, M, D>
114where
115    E: BufferPooler + Clock + Spawner + Metrics,
116    P: PublicKey,
117    M: Digestible + Codec,
118    D: Provider<PublicKey = P>,
119{
120    /// Creates a new engine with the given context and configuration.
121    /// Returns the engine and a mailbox for sending messages to the engine.
122    pub fn new(context: E, cfg: Config<P, M::Cfg, D>) -> (Self, Mailbox<P, M>) {
123        let (mailbox_sender, mailbox_receiver) =
124            mailbox::new(context.child("mailbox"), cfg.mailbox_size);
125        let mailbox = Mailbox::<P, M>::new(mailbox_sender);
126
127        let metrics = metrics::Metrics::init(&context);
128
129        let result = Self {
130            context: ContextCell::new(context),
131            public_key: cfg.public_key,
132            priority: cfg.priority,
133            deque_size: cfg.deque_size,
134            codec_config: cfg.codec_config,
135            mailbox_receiver,
136            waiters: BTreeMap::new(),
137            deques: BTreeMap::new(),
138            items: BTreeMap::new(),
139            counts: BTreeMap::new(),
140            latest_primary_peers: Set::default(),
141            peer_provider: cfg.peer_provider,
142            metrics,
143        };
144
145        (result, mailbox)
146    }
147
148    /// Starts the engine with the given network.
149    pub fn start(
150        mut self,
151        network: (impl Sender<PublicKey = P>, impl Receiver<PublicKey = P>),
152    ) -> Handle<()> {
153        spawn_cell!(self.context, self.run(network))
154    }
155
156    /// Inner run loop called by `start`.
157    async fn run(mut self, network: (impl Sender<PublicKey = P>, impl Receiver<PublicKey = P>)) {
158        let (mut sender, mut receiver) = wrap(
159            self.codec_config.clone(),
160            self.context.network_buffer_pool().clone(),
161            network.0,
162            network.1,
163        );
164        let mut peer_set_subscription = self.peer_provider.subscribe().await;
165
166        select_loop! {
167            self.context,
168            on_start => {
169                // Cleanup waiters
170                self.cleanup_waiters();
171                let _ = self.metrics.waiters.try_set(self.waiters.len());
172            },
173            on_stopped => {
174                debug!("shutdown");
175            },
176            // Handle peer set subscription messages
177            Some(update) = peer_set_subscription.recv() else {
178                debug!("peer set subscription closed");
179                break;
180            } => {
181                // Evict by latest primary only; see buffered module docs.
182                self.update_latest_primary_peers(update.latest.primary);
183            },
184            // Handle mailbox messages
185            Some(msg) = self.mailbox_receiver.recv() else {
186                error!("mailbox receiver failed");
187                break;
188            } => match msg {
189                Message::Broadcast {
190                    recipients,
191                    message,
192                } => {
193                    trace!("mailbox: broadcast");
194                    self.handle_broadcast(&mut sender, recipients, message);
195                }
196                Message::Subscribe { digest, responder } => {
197                    trace!("mailbox: subscribe");
198                    self.handle_subscribe(digest, responder);
199                }
200                Message::Get { digest, responder } => {
201                    trace!("mailbox: get");
202                    self.handle_get(digest, responder);
203                }
204            },
205            // Handle incoming messages
206            msg = receiver.recv() => {
207                // Error handling
208                let (peer, msg) = match msg {
209                    Ok(r) => r,
210                    Err(err) => {
211                        error!(?err, "receiver failed");
212                        break;
213                    }
214                };
215
216                // Decode the message
217                let msg = match msg {
218                    Ok(msg) => msg,
219                    Err(err) => {
220                        warn!(?err, ?peer, "failed to decode message");
221                        self.metrics.receive.inc(Status::Invalid);
222                        continue;
223                    }
224                };
225
226                trace!(?peer, "network");
227                self.metrics.peer.get_or_create_by(&peer).inc();
228                self.handle_network(peer, msg);
229            },
230        }
231    }
232
233    ////////////////////////////////////////
234    // Handling
235    ////////////////////////////////////////
236
237    /// Handles a `broadcast` request from the application.
238    fn handle_broadcast<Sr: Sender<PublicKey = P>>(
239        &mut self,
240        sender: &mut WrappedSender<Sr, M>,
241        recipients: Recipients<P>,
242        msg: Arc<M>,
243    ) {
244        // Store the message, continue even if it was already stored
245        let digest = msg.digest();
246        let _ = self.insert_shared_message(self.public_key.clone(), digest, &msg);
247
248        // Broadcast the message to the network
249        sender.send_ref(recipients, msg.as_ref(), self.priority);
250    }
251
252    /// Handles a `subscribe` request from the application.
253    ///
254    /// If the message is already in the cache, the responder is immediately sent the message.
255    /// Otherwise, the responder is stored in the waiters list.
256    fn handle_subscribe(&mut self, digest: M::Digest, responder: oneshot::Sender<Arc<M>>) {
257        // Check if the message is already in the cache
258        if let Some(item) = self.items.get(&digest).cloned() {
259            self.respond_subscribe(responder, item);
260            return;
261        }
262
263        // Store the responder
264        self.waiters
265            .entry(digest)
266            .or_default()
267            .push(Waiter { responder });
268    }
269
270    /// Handles a `get` request from the application.
271    fn handle_get(&mut self, digest: M::Digest, responder: oneshot::Sender<Option<Arc<M>>>) {
272        let item = self.items.get(&digest).cloned();
273        self.respond_get(responder, item);
274    }
275
276    /// Handles a message that was received from a peer.
277    fn handle_network(&mut self, peer: P, msg: M) {
278        let digest = msg.digest();
279        match self.insert_message(peer.clone(), digest, msg) {
280            InsertMessageResult::Inserted => {
281                self.metrics.receive.inc(Status::Success);
282            }
283            InsertMessageResult::Duplicate => {
284                debug!(?peer, "message already stored");
285                self.metrics.receive.inc(Status::Dropped);
286            }
287            InsertMessageResult::Ineligible => {
288                debug!(?peer, "message from peer outside latest.primary not cached");
289                self.metrics.receive.inc(Status::Dropped);
290            }
291        }
292    }
293
294    ////////////////////////////////////////
295    // Cache Management
296    ////////////////////////////////////////
297
298    /// Inserts a message into the cache.
299    ///
300    /// Waiters are notified even when a sender is not eligible to keep a
301    /// buffered cache entry resident.
302    fn insert_message(&mut self, peer: P, digest: M::Digest, msg: M) -> InsertMessageResult {
303        if let Some(waiters) = self.waiters.remove(&digest) {
304            let msg = Arc::new(msg);
305            self.respond_waiters(waiters, &msg);
306            return self.insert_cache_entry(peer, digest, || Arc::clone(&msg));
307        }
308
309        self.insert_cache_entry(peer, digest, || Arc::new(msg))
310    }
311
312    /// Inserts a shared message into the cache.
313    fn insert_shared_message(
314        &mut self,
315        peer: P,
316        digest: M::Digest,
317        msg: &Arc<M>,
318    ) -> InsertMessageResult {
319        if let Some(waiters) = self.waiters.remove(&digest) {
320            self.respond_waiters(waiters, msg);
321        }
322
323        self.insert_cache_entry(peer, digest, || Arc::clone(msg))
324    }
325
326    /// Records a peer's reference to a message, acquiring an `Arc` only when
327    /// the cache needs to store the message.
328    fn insert_cache_entry(
329        &mut self,
330        peer: P,
331        digest: M::Digest,
332        make_shared: impl FnOnce() -> Arc<M>,
333    ) -> InsertMessageResult {
334        // Only peers listed in `latest.primary` may buffer
335        if self.latest_primary_peers.position(&peer).is_none() {
336            return InsertMessageResult::Ineligible;
337        }
338
339        // Get the relevant deque for the peer
340        let deque = self
341            .deques
342            .entry(peer)
343            .or_insert_with(|| VecDeque::with_capacity(self.deque_size + 1));
344
345        // If the message is already in the deque, move it to the front and return early
346        if let Some(i) = deque.iter().position(|d| *d == digest) {
347            if i != 0 {
348                let v = deque.remove(i).unwrap(); // Must exist
349                deque.push_front(v);
350            }
351            return InsertMessageResult::Duplicate;
352        };
353
354        // - Insert the digest into the peer cache
355        // - Increment the item count
356        // - Insert the message if-and-only-if the new item count is 1
357        deque.push_front(digest);
358        let count = self
359            .counts
360            .entry(digest)
361            .and_modify(|c| *c = c.checked_add(1).unwrap())
362            .or_insert(1);
363        if *count == 1 {
364            let existing = self.items.insert(digest, make_shared());
365            assert!(existing.is_none());
366        }
367
368        // If the cache is full...
369        if deque.len() > self.deque_size {
370            // Remove the oldest item from the peer cache
371            // Decrement the item count
372            // Remove the message if-and-only-if the new item count is 0
373            let stale = deque.pop_back().unwrap();
374            decrement_digest_refcount(&mut self.counts, &mut self.items, &stale);
375        }
376
377        InsertMessageResult::Inserted
378    }
379
380    fn update_latest_primary_peers(&mut self, peers: Set<P>) {
381        for (peer, deque) in self
382            .deques
383            .extract_if(.., |peer, _| peers.position(peer).is_none())
384        {
385            debug!(?peer, digests = deque.len(), "evicting disconnected peer");
386            for digest in deque {
387                decrement_digest_refcount(&mut self.counts, &mut self.items, &digest);
388            }
389        }
390        self.latest_primary_peers = peers;
391    }
392
393    ////////////////////////////////////////
394    // Utilities
395    ////////////////////////////////////////
396
397    /// Remove all waiters that have dropped receivers.
398    fn cleanup_waiters(&mut self) {
399        self.waiters.retain(|_, waiters| {
400            let initial_len = waiters.len();
401            waiters.retain(|waiter| !waiter.responder.is_closed());
402            let dropped_count = initial_len - waiters.len();
403
404            // Increment metrics for each dropped waiter
405            for _ in 0..dropped_count {
406                self.metrics.get.inc(Status::Dropped);
407            }
408
409            !waiters.is_empty()
410        });
411    }
412
413    /// Respond to a waiter with a message.
414    /// Increments the appropriate metric based on the result.
415    fn respond_subscribe(&mut self, responder: oneshot::Sender<Arc<M>>, msg: Arc<M>) {
416        self.metrics.subscribe.inc(if responder.send_lossy(msg) {
417            Status::Success
418        } else {
419            Status::Dropped
420        });
421    }
422
423    fn respond_waiters(&mut self, waiters: Vec<Waiter<M>>, msg: &Arc<M>) {
424        for waiter in waiters {
425            self.respond_subscribe(waiter.responder, Arc::clone(msg));
426        }
427    }
428
429    /// Respond to a get request.
430    /// Increments the appropriate metric based on the result.
431    fn respond_get(&mut self, responder: oneshot::Sender<Option<Arc<M>>>, msg: Option<Arc<M>>) {
432        let found = msg.is_some();
433        self.metrics.get.inc(if responder.send_lossy(msg) {
434            if found {
435                Status::Success
436            } else {
437                Status::Failure
438            }
439        } else {
440            Status::Dropped
441        });
442    }
443}
444
445/// Decrement a digest refcount and evict it from cache when no references remain.
446fn decrement_digest_refcount<D: Ord, M>(
447    counts: &mut BTreeMap<D, usize>,
448    items: &mut BTreeMap<D, M>,
449    digest: &D,
450) {
451    let should_remove = {
452        let count = counts.get_mut(digest).expect("count must exist");
453        *count = count.checked_sub(1).expect("count must be > 0");
454        *count == 0
455    };
456    if should_remove {
457        let existing = counts.remove(digest);
458        assert!(existing == Some(0));
459        items.remove(digest);
460    }
461}