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
24struct Waiter<M> {
26 responder: oneshot::Sender<Arc<M>>,
28}
29
30enum InsertMessageResult {
32 Inserted,
33 Duplicate,
34 Ineligible,
35}
36
37pub 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 context: ContextCell<E>,
55
56 public_key: P,
61
62 priority: bool,
64
65 deque_size: usize,
67
68 codec_config: M::Cfg,
70
71 mailbox_receiver: mailbox::Receiver<Message<P, M>>,
76
77 waiters: BTreeMap<M::Digest, Vec<Waiter<M>>>,
79
80 peer_provider: D,
82
83 items: BTreeMap<M::Digest, Arc<M>>,
88
89 deques: BTreeMap<P, VecDeque<M::Digest>>,
95
96 counts: BTreeMap<M::Digest, usize>,
101
102 latest_primary_peers: Set<P>,
104
105 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 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 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 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 self.cleanup_waiters();
170 let _ = self.metrics.waiters.try_set(self.waiters.len());
171 },
172 on_stopped => {
173 debug!("shutdown");
174 },
175 Some(update) = peer_set_subscription.recv() else {
177 debug!("peer set subscription closed");
178 break;
179 } => {
180 self.update_latest_primary_peers(update.latest.primary);
182 },
183 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 msg = receiver.recv() => {
206 let (peer, msg) = match msg {
208 Ok(r) => r,
209 Err(err) => {
210 error!(?err, "receiver failed");
211 break;
212 }
213 };
214
215 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 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 let digest = msg.digest();
245 let _ = self.insert_shared_message(self.public_key.clone(), digest, &msg);
246
247 sender.send_ref(recipients, msg.as_ref(), self.priority);
249 }
250
251 fn handle_subscribe(&mut self, digest: M::Digest, responder: oneshot::Sender<Arc<M>>) {
256 if let Some(item) = self.items.get(&digest).cloned() {
258 self.respond_subscribe(responder, item);
259 return;
260 }
261
262 self.waiters
264 .entry(digest)
265 .or_default()
266 .push(Waiter { responder });
267 }
268
269 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 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 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 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 fn insert_cache_entry(
328 &mut self,
329 peer: P,
330 digest: M::Digest,
331 make_shared: impl FnOnce() -> Arc<M>,
332 ) -> InsertMessageResult {
333 if self.latest_primary_peers.position(&peer).is_none() {
335 return InsertMessageResult::Ineligible;
336 }
337
338 let deque = self
340 .deques
341 .entry(peer)
342 .or_insert_with(|| VecDeque::with_capacity(self.deque_size + 1));
343
344 if let Some(i) = deque.iter().position(|d| *d == digest) {
346 if i != 0 {
347 let v = deque.remove(i).unwrap(); deque.push_front(v);
349 }
350 return InsertMessageResult::Duplicate;
351 };
352
353 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 deque.len() > self.deque_size {
369 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 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 for _ in 0..dropped_count {
405 self.metrics.get.inc(Status::Dropped);
406 }
407
408 !waiters.is_empty()
409 });
410 }
411
412 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 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
444fn 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}