Skip to main content

commonware_broadcast/buffered/
engine.rs

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