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
25struct Waiter<M> {
27 responder: oneshot::Sender<Arc<M>>,
29}
30
31enum InsertMessageResult {
33 Inserted,
34 Duplicate,
35 Ineligible,
36}
37
38pub 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 context: ContextCell<E>,
56
57 public_key: P,
62
63 priority: bool,
65
66 deque_size: usize,
68
69 codec_config: M::Cfg,
71
72 mailbox_receiver: mailbox::Receiver<Message<P, M>>,
77
78 waiters: BTreeMap<M::Digest, Vec<Waiter<M>>>,
80
81 peer_provider: D,
83
84 items: BTreeMap<M::Digest, Arc<M>>,
89
90 deques: BTreeMap<P, VecDeque<M::Digest>>,
96
97 counts: BTreeMap<M::Digest, usize>,
102
103 latest_primary_peers: Set<P>,
105
106 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 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 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 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 self.cleanup_waiters();
171 let _ = self.metrics.waiters.try_set(self.waiters.len());
172 },
173 on_stopped => {
174 debug!("shutdown");
175 },
176 Some(update) = peer_set_subscription.recv() else {
178 debug!("peer set subscription closed");
179 break;
180 } => {
181 self.update_latest_primary_peers(update.latest.primary);
183 },
184 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 msg = receiver.recv() => {
207 let (peer, msg) = match msg {
209 Ok(r) => r,
210 Err(err) => {
211 error!(?err, "receiver failed");
212 break;
213 }
214 };
215
216 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 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 let digest = msg.digest();
246 let _ = self.insert_shared_message(self.public_key.clone(), digest, &msg);
247
248 sender.send_ref(recipients, msg.as_ref(), self.priority);
250 }
251
252 fn handle_subscribe(&mut self, digest: M::Digest, responder: oneshot::Sender<Arc<M>>) {
257 if let Some(item) = self.items.get(&digest).cloned() {
259 self.respond_subscribe(responder, item);
260 return;
261 }
262
263 self.waiters
265 .entry(digest)
266 .or_default()
267 .push(Waiter { responder });
268 }
269
270 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 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 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 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 fn insert_cache_entry(
329 &mut self,
330 peer: P,
331 digest: M::Digest,
332 make_shared: impl FnOnce() -> Arc<M>,
333 ) -> InsertMessageResult {
334 if self.latest_primary_peers.position(&peer).is_none() {
336 return InsertMessageResult::Ineligible;
337 }
338
339 let deque = self
341 .deques
342 .entry(peer)
343 .or_insert_with(|| VecDeque::with_capacity(self.deque_size + 1));
344
345 if let Some(i) = deque.iter().position(|d| *d == digest) {
347 if i != 0 {
348 let v = deque.remove(i).unwrap(); deque.push_front(v);
350 }
351 return InsertMessageResult::Duplicate;
352 };
353
354 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 deque.len() > self.deque_size {
370 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 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 for _ in 0..dropped_count {
406 self.metrics.get.inc(Status::Dropped);
407 }
408
409 !waiters.is_empty()
410 });
411 }
412
413 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 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
445fn 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}