1mod config;
26pub use config::Config;
27mod engine;
28pub use engine::Engine;
29mod ingress;
30pub use ingress::Mailbox;
31pub(crate) use ingress::Message;
32mod metrics;
33
34#[cfg(test)]
35pub mod mocks;
36
37#[cfg(test)]
38mod tests {
39 use super::{mocks::TestMessage, *};
40 use crate::Broadcaster;
41 use commonware_actor::{
42 Feedback,
43 mailbox::{Overflow, Policy},
44 };
45 use commonware_codec::RangeCfg;
46 use commonware_cryptography::{
47 Digestible, Hasher, Sha256, Signer as _,
48 ed25519::{PrivateKey, PublicKey},
49 };
50 use commonware_macros::test_traced;
51 use commonware_p2p::{
52 Manager as _, Recipients, Sender as _, TrackedPeers,
53 simulated::{Link, Network, Oracle, Receiver, Sender},
54 };
55 use commonware_runtime::{
56 Clock, Error, IoBuf, Metrics as _, Quota, Runner, Supervisor as _, deterministic,
57 telemetry::metrics::count_running_tasks,
58 };
59 use commonware_utils::{NZUsize, Probability, probability};
60 use std::{
61 collections::{BTreeMap, VecDeque},
62 num::NonZeroU32,
63 sync::Arc,
64 time::Duration,
65 };
66
67 const CACHE_SIZE: usize = 10;
69
70 const A_JIFFY: Duration = Duration::from_millis(10);
73
74 const NETWORK_SPEED: Duration = Duration::from_millis(100);
76
77 const NETWORK_SPEED_WITH_BUFFER: Duration = Duration::from_millis(200);
79
80 const TEST_QUOTA: Quota = Quota::per_second(NonZeroU32::MAX);
82
83 type Registrations = BTreeMap<
84 PublicKey,
85 (
86 Sender<PublicKey, deterministic::Context>,
87 Receiver<PublicKey>,
88 ),
89 >;
90
91 async fn initialize_simulation(
92 context: deterministic::Context,
93 num_peers: usize,
94 success_rate: Probability,
95 ) -> (
96 Vec<PublicKey>,
97 Registrations,
98 Oracle<PublicKey, deterministic::Context>,
99 ) {
100 let (network, oracle) = Network::<deterministic::Context, PublicKey>::new(
101 context,
102 commonware_p2p::simulated::Config {
103 max_size: 1024 * 1024,
104 max_peers_per_set: NZUsize!(num_peers),
105 disconnect_on_block: true,
106 tracked_peer_sets: NZUsize!(1),
107 },
108 );
109 network.start();
110
111 let mut schemes = (0..num_peers)
112 .map(|i| PrivateKey::from_seed(i as u64))
113 .collect::<Vec<_>>();
114 schemes.sort_by_key(|s| s.public_key());
115 let peers: Vec<PublicKey> = schemes.iter().map(|c| c.public_key()).collect();
116
117 let mut registrations: Registrations = BTreeMap::new();
118 for peer in peers.iter() {
119 let (sender, receiver) = oracle
120 .control(peer.clone())
121 .register(0, TEST_QUOTA)
122 .await
123 .unwrap();
124 registrations.insert(peer.clone(), (sender, receiver));
125 }
126
127 let link = Link {
129 latency: NETWORK_SPEED,
130 jitter: Duration::ZERO,
131 success_rate,
132 };
133 for p1 in peers.iter() {
134 for p2 in peers.iter() {
135 if p2 == p1 {
136 continue;
137 }
138 oracle
139 .add_link(p1.clone(), p2.clone(), link.clone())
140 .await
141 .unwrap();
142 }
143 }
144
145 let all_peers = commonware_utils::ordered::Set::from_iter_dedup(peers.clone());
147 oracle.manager().track(0, all_peers);
148
149 (peers, registrations, oracle)
150 }
151
152 #[test]
153 fn policy_handles_closed_responders() {
154 let mut overflow = <Message<PublicKey, TestMessage> as Policy>::Overflow::default();
155 let pending_subscribe = TestMessage::shared(b"pending_subscribe");
156 let pending_get = TestMessage::shared(b"pending_get");
157 let open_subscribe = TestMessage::shared(b"open_subscribe");
158 let open_get = TestMessage::shared(b"open_get");
159 let current_get = TestMessage::shared(b"current_get");
160
161 let (closed_responder, closed_receiver) = commonware_utils::channel::oneshot::channel();
162 <Message<PublicKey, TestMessage> as Policy>::handle(
163 &mut overflow,
164 Message::Subscribe {
165 digest: pending_subscribe.digest(),
166 responder: closed_responder,
167 },
168 );
169 drop(closed_receiver);
170
171 let (open_responder, _open_receiver) = commonware_utils::channel::oneshot::channel();
172 <Message<PublicKey, TestMessage> as Policy>::handle(
173 &mut overflow,
174 Message::Subscribe {
175 digest: open_subscribe.digest(),
176 responder: open_responder,
177 },
178 );
179
180 let (closed_responder, closed_receiver) = commonware_utils::channel::oneshot::channel();
181 <Message<PublicKey, TestMessage> as Policy>::handle(
182 &mut overflow,
183 Message::Get {
184 digest: pending_get.digest(),
185 responder: closed_responder,
186 },
187 );
188 drop(closed_receiver);
189
190 let (open_responder, _open_receiver) = commonware_utils::channel::oneshot::channel();
191 <Message<PublicKey, TestMessage> as Policy>::handle(
192 &mut overflow,
193 Message::Get {
194 digest: open_get.digest(),
195 responder: open_responder,
196 },
197 );
198
199 let (current_responder, current_receiver) = commonware_utils::channel::oneshot::channel();
200 drop(current_receiver);
201 <Message<PublicKey, TestMessage> as Policy>::handle(
202 &mut overflow,
203 Message::Get {
204 digest: current_get.digest(),
205 responder: current_responder,
206 },
207 );
208
209 let mut drained = VecDeque::new();
210 overflow.drain(|message| {
211 drained.push_back(message);
212 None
213 });
214
215 assert_eq!(drained.len(), 2);
216 assert!(drained.iter().any(|message| matches!(
217 message,
218 Message::Subscribe { digest, responder }
219 if *digest == open_subscribe.digest() && !responder.is_closed()
220 )));
221 assert!(drained.iter().any(|message| matches!(
222 message,
223 Message::Get { digest, responder }
224 if *digest == open_get.digest() && !responder.is_closed()
225 )));
226 }
227
228 #[test]
229 fn policy_drain_continues_until_rejected_message() {
230 let mut overflow = <Message<PublicKey, TestMessage> as Policy>::Overflow::default();
231 let first = TestMessage::shared(b"first");
232 let second = TestMessage::shared(b"second");
233 let third = TestMessage::shared(b"third");
234
235 let (closed_responder, closed_receiver) = commonware_utils::channel::oneshot::channel();
236 <Message<PublicKey, TestMessage> as Policy>::handle(
237 &mut overflow,
238 Message::Subscribe {
239 digest: TestMessage::shared(b"closed").digest(),
240 responder: closed_responder,
241 },
242 );
243 drop(closed_receiver);
244
245 let (first_responder, _first_receiver) = commonware_utils::channel::oneshot::channel();
246 <Message<PublicKey, TestMessage> as Policy>::handle(
247 &mut overflow,
248 Message::Get {
249 digest: first.digest(),
250 responder: first_responder,
251 },
252 );
253 let (second_responder, _second_receiver) = commonware_utils::channel::oneshot::channel();
254 <Message<PublicKey, TestMessage> as Policy>::handle(
255 &mut overflow,
256 Message::Get {
257 digest: second.digest(),
258 responder: second_responder,
259 },
260 );
261 let (third_responder, _third_receiver) = commonware_utils::channel::oneshot::channel();
262 <Message<PublicKey, TestMessage> as Policy>::handle(
263 &mut overflow,
264 Message::Get {
265 digest: third.digest(),
266 responder: third_responder,
267 },
268 );
269
270 let mut drained = VecDeque::new();
271 overflow.drain(|message| {
272 drained.push_back(message);
273 if drained.len() == 3 {
274 drained.pop_back()
275 } else {
276 None
277 }
278 });
279
280 assert_eq!(drained.len(), 2);
281 assert!(matches!(
282 &drained[0],
283 Message::Get { digest, responder }
284 if *digest == first.digest() && !responder.is_closed()
285 ));
286 assert!(matches!(
287 &drained[1],
288 Message::Get { digest, responder }
289 if *digest == second.digest() && !responder.is_closed()
290 ));
291
292 overflow.drain(|message| {
293 drained.push_back(message);
294 None
295 });
296 assert_eq!(drained.len(), 3);
297 assert!(matches!(
298 &drained[2],
299 Message::Get { digest, responder }
300 if *digest == third.digest() && !responder.is_closed()
301 ));
302 }
303
304 #[test]
305 fn policy_drain_stops_after_returned_message_closes() {
306 let mut overflow = <Message<PublicKey, TestMessage> as Policy>::Overflow::default();
307 let first = TestMessage::shared(b"first");
308 let second = TestMessage::shared(b"second");
309
310 let (first_responder, first_receiver) = commonware_utils::channel::oneshot::channel();
311 <Message<PublicKey, TestMessage> as Policy>::handle(
312 &mut overflow,
313 Message::Get {
314 digest: first.digest(),
315 responder: first_responder,
316 },
317 );
318 let (second_responder, _second_receiver) = commonware_utils::channel::oneshot::channel();
319 <Message<PublicKey, TestMessage> as Policy>::handle(
320 &mut overflow,
321 Message::Get {
322 digest: second.digest(),
323 responder: second_responder,
324 },
325 );
326
327 let mut first_receiver = Some(first_receiver);
328 let mut attempts = 0;
329 overflow.drain(|message| {
330 attempts += 1;
331 drop(first_receiver.take());
332 Some(message)
333 });
334 assert_eq!(attempts, 1);
335
336 let mut drained = VecDeque::new();
337 overflow.drain(|message| {
338 drained.push_back(message);
339 None
340 });
341 assert_eq!(drained.len(), 1);
342 assert!(matches!(
343 &drained[0],
344 Message::Get { digest, responder }
345 if *digest == second.digest() && !responder.is_closed()
346 ));
347 }
348
349 async fn spawn_peer_engines(
350 context: deterministic::Context,
351 oracle: &Oracle<PublicKey, deterministic::Context>,
352 registrations: &mut Registrations,
353 ) -> BTreeMap<PublicKey, Mailbox<PublicKey, TestMessage>> {
354 let mut mailboxes = BTreeMap::new();
355 while let Some((peer, network)) = registrations.pop_first() {
356 let context = context.child("peer").with_attribute("public_key", &peer);
357 let config = Config {
358 public_key: peer.clone(),
359 mailbox_size: NZUsize!(1024),
360 deque_size: CACHE_SIZE,
361 priority: false,
362 codec_config: RangeCfg::from(..),
363 peer_provider: oracle.manager(),
364 };
365 let (engine, engine_mailbox) =
366 Engine::<_, PublicKey, TestMessage, _>::new(context, config);
367 mailboxes.insert(peer.clone(), engine_mailbox);
368 engine.start(network);
369 }
370
371 context.sleep(A_JIFFY).await;
374 mailboxes
375 }
376
377 #[test_traced]
378 fn test_broadcast() {
379 let runner = deterministic::Runner::timed(Duration::from_secs(5));
380 runner.start(|context| async move {
381 let (peers, mut registrations, oracle) =
382 initialize_simulation(context.child("network"), 4, probability!(1.0)).await;
383 let mailboxes =
384 spawn_peer_engines(context.child("peers"), &oracle, &mut registrations).await;
385
386 let message = TestMessage::shared(b"hello world test message");
388 let first_mailbox = mailboxes.get(peers.first().unwrap()).unwrap().clone();
389 assert!(
390 first_mailbox
391 .broadcast(Recipients::All, message.clone())
392 .accepted()
393 );
394
395 context.sleep(Duration::from_secs(1)).await;
397
398 for peer in peers.iter() {
400 let mailbox = mailboxes.get(peer).unwrap().clone();
401 let digest = message.digest();
402 let receiver = mailbox.subscribe(digest);
403 let received_message = receiver.await.ok();
404 assert_eq!(received_message.unwrap().as_ref(), &message);
405 }
406
407 let message = TestMessage::shared(b"hello world again");
409 assert!(
410 first_mailbox
411 .broadcast(Recipients::All, message.clone())
412 .accepted()
413 );
414
415 context.sleep(Duration::from_secs(1)).await;
417
418 let mut found = 0;
420 for peer in peers.iter() {
421 let mailbox = mailboxes.get(peer).unwrap().clone();
422 let digest = message.digest();
423 let receiver = mailbox.get(digest).await;
424 if let Some(receiver) = receiver {
425 assert_eq!(receiver.as_ref(), &message);
426 found += 1;
427 }
428 }
429 assert!(found > 0, "No peers received the message");
430 });
431 }
432
433 #[test_traced]
434 fn test_self_retrieval() {
435 let runner = deterministic::Runner::timed(Duration::from_secs(5));
436 runner.start(|context| async move {
437 let (peers, mut registrations, oracle) =
439 initialize_simulation(context.child("network"), 1, probability!(1.0)).await;
440 let mailboxes =
441 spawn_peer_engines(context.child("peers"), &oracle, &mut registrations).await;
442
443 let mailbox_a = mailboxes.get(&peers[0]).unwrap().clone();
445
446 let m1 = TestMessage::shared(b"hello world");
448 let digest_m1 = m1.digest();
449
450 let receiver_before = mailbox_a.get(digest_m1).await;
452 assert!(receiver_before.is_none());
453
454 let receiver_before = mailbox_a.subscribe(digest_m1);
456
457 assert!(mailbox_a.broadcast(Recipients::All, m1.clone()).accepted());
459
460 let msg_before = receiver_before
462 .await
463 .expect("Pre-broadcast retrieval failed");
464 assert_eq!(msg_before.as_ref(), &m1);
465
466 let receiver_after = mailbox_a.get(digest_m1).await;
468 assert_eq!(receiver_after.as_deref(), Some(&m1));
469
470 let receiver_after = mailbox_a.subscribe(digest_m1);
472
473 let start = context.current();
475 let msg_after = receiver_after
476 .await
477 .expect("Post-broadcast retrieval failed");
478 let duration = context.current().duration_since(start).unwrap();
479
480 assert_eq!(msg_after.as_ref(), &m1);
482
483 assert!(duration < A_JIFFY, "get not instant");
485 });
486 }
487
488 #[test_traced]
489 fn test_shared_broadcast_reuses_message() {
490 let runner = deterministic::Runner::timed(Duration::from_secs(5));
491 runner.start(|context| async move {
492 let (peers, mut registrations, oracle) =
493 initialize_simulation(context.child("network"), 1, probability!(1.0)).await;
494 let mailboxes =
495 spawn_peer_engines(context.child("peers"), &oracle, &mut registrations).await;
496 let mailbox = mailboxes.get(&peers[0]).unwrap();
497
498 let message = Arc::new(TestMessage::shared(b"shared broadcast"));
499 let digest = message.digest();
500 assert!(
501 mailbox
502 .broadcast_shared(Recipients::All, Arc::clone(&message))
503 .accepted()
504 );
505
506 let cached = mailbox.get(digest).await.expect("message should be cached");
507 assert!(Arc::ptr_eq(&message, &cached));
508 });
509 }
510
511 #[test_traced]
512 fn test_packet_loss() {
513 let runner = deterministic::Runner::timed(Duration::from_secs(30));
514 runner.start(|context| async move {
515 let (peers, mut registrations, oracle) =
516 initialize_simulation(context.child("network"), 10, probability!(0.1)).await;
517 let mailboxes =
518 spawn_peer_engines(context.child("peers"), &oracle, &mut registrations).await;
519
520 let message = TestMessage::shared(b"hello world test message");
522 let first_mailbox = mailboxes.get(peers.first().unwrap()).unwrap().clone();
523
524 let digest = message.digest();
526 for i in 0..100 {
527 assert!(
529 first_mailbox
530 .broadcast(Recipients::All, message.clone())
531 .accepted()
532 );
533 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
534
535 let mut all_received = true;
537 for peer in peers.iter() {
538 let mailbox = mailboxes.get(peer).unwrap().clone();
539 let receiver = mailbox.subscribe(digest);
540 let has = match context.timeout(A_JIFFY, receiver).await {
541 Ok(r) => r.is_ok(),
542 Err(Error::Timeout) => false,
543 Err(e) => panic!("unexpected error: {e:?}"),
544 };
545 all_received &= has;
546 }
547 if all_received {
549 assert!(i > 0, "Message received on first try");
550 return;
551 }
552 }
553 panic!("Not all peers received the message after retries");
554 });
555 }
556
557 #[test_traced]
558 fn test_get_cached() {
559 let runner = deterministic::Runner::timed(Duration::from_secs(5));
560 runner.start(|context| async move {
561 let (peers, mut registrations, oracle) =
562 initialize_simulation(context.child("network"), 2, probability!(1.0)).await;
563 let mailboxes =
564 spawn_peer_engines(context.child("peers"), &oracle, &mut registrations).await;
565
566 let message = TestMessage::shared(b"cached message");
568 let first_mailbox = mailboxes.get(peers.first().unwrap()).unwrap().clone();
569 assert!(
570 first_mailbox
571 .broadcast(Recipients::All, message.clone())
572 .accepted()
573 );
574
575 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
577
578 let digest = message.digest();
580 let mailbox = mailboxes.get(peers.last().unwrap()).unwrap().clone();
581 let receiver = mailbox.subscribe(digest);
582 let start = context.current();
583 let received = receiver.await.expect("failed to get cached message");
584 let duration = context.current().duration_since(start).unwrap();
585 assert_eq!(received.as_ref(), &message);
586 assert!(duration < A_JIFFY, "get not instant",);
587 });
588 }
589
590 #[test_traced]
591 fn test_get_nonexistent() {
592 let runner = deterministic::Runner::timed(Duration::from_secs(5));
593 runner.start(|context| async move {
594 let (peers, mut registrations, oracle) =
595 initialize_simulation(context.child("network"), 2, probability!(1.0)).await;
596 let mailboxes =
597 spawn_peer_engines(context.child("peers"), &oracle, &mut registrations).await;
598
599 let message = TestMessage::shared(b"future message");
601 let digest = message.digest();
602 let mailbox1 = mailboxes.get(&peers[0]).unwrap().clone();
603 let mailbox2 = mailboxes.get(&peers[1]).unwrap().clone();
604 let receiver = mailbox1.subscribe(digest);
605
606 let dummy1 = mailbox1.subscribe(digest);
608 let dummy2 = mailbox2.subscribe(digest);
609 drop(dummy1);
610 drop(dummy2);
611
612 assert!(
614 mailbox1
615 .broadcast(Recipients::All, message.clone())
616 .accepted()
617 );
618
619 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
621
622 let received = receiver.await.expect("receiver1 should get message");
624 assert_eq!(received.as_ref(), &message);
625 });
626 }
627
628 #[test_traced]
629 fn test_cache_eviction_single_peer() {
630 let runner = deterministic::Runner::timed(Duration::from_secs(5));
631 runner.start(|context| async move {
632 let (peers, mut registrations, oracle) =
633 initialize_simulation(context.child("network"), 2, probability!(1.0)).await;
634 let mailboxes =
635 spawn_peer_engines(context.child("peers"), &oracle, &mut registrations).await;
636
637 let mailbox = mailboxes.get(&peers[0]).unwrap().clone();
639 let mut messages = vec![];
640 for i in 0..CACHE_SIZE + 1 {
641 messages.push(TestMessage::shared(format!("message {i}").as_bytes()));
642 }
643 for message in messages.iter() {
644 assert!(
645 mailbox
646 .broadcast(Recipients::All, message.clone())
647 .accepted()
648 );
649 }
650
651 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
653
654 let peer_mailbox = mailboxes.get(&peers[1]).unwrap().clone();
656 for msg in messages.iter().skip(1) {
657 let result = peer_mailbox.subscribe(msg.digest()).await.unwrap();
658 assert_eq!(result.as_ref(), msg);
659 }
660
661 let receiver = peer_mailbox.subscribe(messages[0].digest());
663 match context.timeout(A_JIFFY, receiver).await {
664 Ok(_) => panic!("receiver should have failed"),
665 Err(Error::Timeout) => {} Err(e) => panic!("unexpected error: {e:?}"),
667 }
668 });
669 }
670
671 #[test_traced]
672 fn test_cache_eviction_multi_peer() {
673 let runner = deterministic::Runner::timed(Duration::from_secs(10));
674 runner.start(|context| async move {
675 let (peers, mut registrations, oracle) =
677 initialize_simulation(context.child("network"), 3, probability!(1.0)).await;
678 let mailboxes =
679 spawn_peer_engines(context.child("peers"), &oracle, &mut registrations).await;
680
681 let mailbox_a = mailboxes.get(&peers[0]).unwrap().clone();
683 let mailbox_b = mailboxes.get(&peers[1]).unwrap().clone();
684 let mailbox_c = mailboxes.get(&peers[2]).unwrap().clone();
685
686 let m1 = TestMessage::shared(b"message M1");
688 let digest_m1 = m1.digest();
689 assert!(mailbox_a.broadcast(Recipients::All, m1.clone()).accepted());
690 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
691
692 assert!(mailbox_c.broadcast(Recipients::All, m1.clone()).accepted());
694 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
695
696 let mut new_messages_a = Vec::with_capacity(CACHE_SIZE);
700 for i in 0..CACHE_SIZE {
701 new_messages_a.push(TestMessage::shared(format!("A{i}").as_bytes()));
702 }
703 for msg in &new_messages_a {
704 assert!(mailbox_a.broadcast(Recipients::All, msg.clone()).accepted());
705 }
706 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
707
708 let receiver = mailbox_b.subscribe(digest_m1);
710 let received = receiver.await.expect("M1 should be retrievable");
711 assert_eq!(received.as_ref(), &m1);
712
713 let mut new_messages_c = Vec::with_capacity(CACHE_SIZE);
715 for i in 0..CACHE_SIZE {
716 new_messages_c.push(TestMessage::shared(format!("C{i}").as_bytes()));
717 }
718 for msg in &new_messages_c {
719 assert!(mailbox_c.broadcast(Recipients::All, msg.clone()).accepted());
720 }
721 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
722
723 let receiver = mailbox_b.subscribe(digest_m1);
725 match context.timeout(A_JIFFY, receiver).await {
726 Ok(_) => panic!("M1 should not be retrievable"),
727 Err(Error::Timeout) => {} Err(e) => panic!("unexpected error: {e:?}"),
729 }
730 });
731 }
732
733 #[test_traced]
734 fn test_selective_recipients() {
735 let runner = deterministic::Runner::timed(Duration::from_secs(5));
736 runner.start(|context| async move {
737 let (peers, mut registrations, oracle) =
738 initialize_simulation(context.child("network"), 4, probability!(1.0)).await;
739
740 let sender_pk = peers[0].clone();
741 let target_peer = peers[1].clone();
742 let non_target_peer = peers[2].clone();
743
744 let mailboxes =
745 spawn_peer_engines(context.child("peers"), &oracle, &mut registrations).await;
746 let sender_mb = mailboxes.get(&sender_pk).unwrap().clone();
747
748 let msg = TestMessage::shared(b"selective-broadcast");
749 assert!(
750 sender_mb
751 .broadcast(Recipients::One(target_peer.clone()), msg.clone())
752 .accepted()
753 );
754
755 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
756
757 let got_target = mailboxes
759 .get(&target_peer)
760 .unwrap()
761 .clone()
762 .get(msg.digest())
763 .await;
764 assert_eq!(got_target.as_deref(), Some(&msg));
765
766 let got_other = mailboxes
768 .get(&non_target_peer)
769 .unwrap()
770 .clone()
771 .get(msg.digest())
772 .await;
773 assert!(got_other.is_none());
774 });
775 }
776
777 #[test_traced]
778 fn test_ref_count_across_peers() {
779 let runner = deterministic::Runner::timed(Duration::from_secs(10));
780 runner.start(|context| async move {
781 let (peers, mut registrations, oracle) =
783 initialize_simulation(context.child("network"), 3, probability!(1.0)).await;
784 let mailboxes =
785 spawn_peer_engines(context.child("peers"), &oracle, &mut registrations).await;
786
787 let p0 = peers[0].clone();
788 let p1 = peers[1].clone();
789 let observer = peers[2].clone();
790
791 let mb0 = mailboxes.get(&p0).unwrap().clone();
792 let mb1 = mailboxes.get(&p1).unwrap().clone();
793 let obs = mailboxes.get(&observer).unwrap().clone();
794
795 let dup = TestMessage::shared(b"dup");
797 let digest = dup.digest();
798
799 assert!(mb0.broadcast(Recipients::All, dup.clone()).accepted());
801 assert!(mb1.broadcast(Recipients::All, dup.clone()).accepted());
802 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
803
804 assert_eq!(obs.get(digest).await.as_deref(), Some(&dup));
806
807 for i in 0..CACHE_SIZE {
809 let spam = TestMessage::shared(format!("p0-{i}").into_bytes());
810 assert!(mb0.broadcast(Recipients::All, spam).accepted());
811 }
812 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
813 assert_eq!(obs.get(digest).await.as_deref(), Some(&dup));
814
815 for i in 0..CACHE_SIZE {
817 let spam = TestMessage::shared(format!("p1-{i}").into_bytes());
818 assert!(mb1.broadcast(Recipients::All, spam).accepted());
819 }
820 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
821 assert!(obs.get(digest).await.is_none());
822 });
823 }
824
825 #[test_traced]
826 fn test_deterministic_retrieval() {
827 let run = |seed: u64| {
828 let config = deterministic::Config::new()
829 .with_seed(seed)
830 .with_timeout(Some(Duration::from_secs(5)));
831 let runner = deterministic::Runner::new(config);
832 runner.start(|context| async move {
833 let (peers, mut registrations, oracle) =
834 initialize_simulation(context.child("network"), 1, probability!(1.0)).await;
835 let mailboxes =
836 spawn_peer_engines(context.child("peers"), &oracle, &mut registrations).await;
837
838 let sender1 = peers[0].clone();
839 let mb1 = mailboxes.get(&sender1).unwrap().clone();
840
841 let m1 = TestMessage::shared(b"content-1");
843 let m2 = TestMessage::shared(b"content-2");
844 let m3 = TestMessage::shared(b"content-3");
845 assert!(mb1.broadcast(Recipients::All, m1.clone()).accepted());
846 assert!(mb1.broadcast(Recipients::All, m2.clone()).accepted());
847 assert!(mb1.broadcast(Recipients::All, m3.clone()).accepted());
848
849 let mut hasher = Sha256::default();
850 for msg in [&m1, &m2, &m3] {
851 if let Some(value) = mb1.get(msg.digest()).await {
852 hasher.update(&value.content);
853 }
854 }
855 hasher.finalize().1
856 })
857 };
858
859 for seed in 0..10 {
860 let h1 = run(seed);
861 let h2 = run(seed);
862
863 assert_eq!(h1, h2, "Messages returned in different order for {seed}");
864 }
865 }
866
867 #[test_traced]
868 fn test_malformed_network_payload_does_not_break_valid_traffic() {
869 let runner = deterministic::Runner::timed(Duration::from_secs(10));
870 runner.start(|context| async move {
871 let (peers, mut registrations, oracle) =
872 initialize_simulation(context.child("network"), 3, probability!(1.0)).await;
873
874 let attacker = peers[0].clone();
875 let honest = peers[1].clone();
876 let victim = peers[2].clone();
877
878 let (mut attacker_sender, _) = registrations.remove(&attacker).unwrap();
879 let mailboxes =
880 spawn_peer_engines(context.child("peers"), &oracle, &mut registrations).await;
881 let honest_mailbox = mailboxes.get(&honest).unwrap().clone();
882 let victim_mailbox = mailboxes.get(&victim).unwrap().clone();
883
884 let sent = attacker_sender.send(
886 Recipients::One(victim.clone()),
887 IoBuf::from(vec![0xFF]),
888 false,
889 );
890 assert_eq!(sent, vec![victim.clone()]);
891 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
892
893 let message = TestMessage::shared(b"valid-after-malformed");
895 assert!(
896 honest_mailbox
897 .broadcast(Recipients::One(victim.clone()), message.clone())
898 .accepted()
899 );
900 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
901
902 let received = victim_mailbox
903 .subscribe(message.digest())
904 .await
905 .expect("victim should receive valid message after malformed payload");
906 assert_eq!(received.as_ref(), &message);
907 });
908 }
909
910 #[test_traced]
911 fn test_dropped_waiters_for_missing_digest_are_cleaned_up() {
912 let runner = deterministic::Runner::timed(Duration::from_secs(10));
913 runner.start(|context| async move {
914 let (peers, mut registrations, oracle) =
915 initialize_simulation(context.child("network"), 1, probability!(1.0)).await;
916 let peer = peers[0].clone();
917 let (sender, receiver) = registrations.remove(&peer).unwrap();
918
919 let engine_context = context.child("waiter_cleanup");
920 let config = Config {
921 public_key: peer,
922 mailbox_size: NZUsize!(1024),
923 deque_size: CACHE_SIZE,
924 priority: false,
925 codec_config: RangeCfg::from(..),
926 peer_provider: oracle.manager(),
927 };
928 let (engine, mailbox) =
929 Engine::<_, PublicKey, TestMessage, _>::new(engine_context, config);
930 engine.start((sender, receiver));
931
932 let missing = TestMessage::shared(b"never-arrives");
933 let missing_digest = missing.digest();
934 let rx1 = mailbox.subscribe(missing_digest);
935 let rx2 = mailbox.subscribe(missing_digest);
936
937 let _ = mailbox
939 .get(TestMessage::shared(b"before-cleanup").digest())
940 .await;
941 context.sleep(A_JIFFY).await;
942 let metrics_before = context.encode();
943 let waiter_values_before: Vec<f64> = metrics_before
944 .lines()
945 .filter(|line| {
946 line.starts_with("waiters")
947 || (line.contains("_waiters")
948 && !line.starts_with("# HELP")
949 && !line.starts_with("# TYPE"))
950 })
951 .filter_map(|line| line.split_whitespace().last())
952 .filter_map(|value| value.parse::<f64>().ok())
953 .collect();
954 assert!(
955 !waiter_values_before.is_empty(),
956 "waiters metric not found in output:\n{metrics_before}"
957 );
958 assert!(
959 waiter_values_before.iter().any(|value| *value > 0.0),
960 "expected positive waiters before cleanup, got:\n{metrics_before}"
961 );
962
963 drop(rx1);
964 drop(rx2);
965
966 let _ = mailbox
968 .get(TestMessage::shared(b"after-cleanup").digest())
969 .await;
970 context.sleep(A_JIFFY).await;
971
972 let metrics_after = context.encode();
973 let waiter_values_after: Vec<f64> = metrics_after
974 .lines()
975 .filter(|line| {
976 line.starts_with("waiters")
977 || (line.contains("_waiters")
978 && !line.starts_with("# HELP")
979 && !line.starts_with("# TYPE"))
980 })
981 .filter_map(|line| line.split_whitespace().last())
982 .filter_map(|value| value.parse::<f64>().ok())
983 .collect();
984 assert!(
985 !waiter_values_after.is_empty(),
986 "waiters metric not found in output:\n{metrics_after}"
987 );
988 assert!(
989 waiter_values_after.iter().all(|value| *value == 0.0),
990 "expected zero retained waiters, got:\n{metrics_after}"
991 );
992 });
993 }
994
995 #[allow(clippy::type_complexity)]
996 async fn spawn_peer_engines_with_handles(
997 context: deterministic::Context,
998 oracle: &Oracle<PublicKey, deterministic::Context>,
999 registrations: &mut Registrations,
1000 ) -> (
1001 BTreeMap<PublicKey, Mailbox<PublicKey, TestMessage>>,
1002 Vec<commonware_runtime::Handle<()>>,
1003 ) {
1004 let mut mailboxes = BTreeMap::new();
1005 let mut handles = Vec::new();
1006 while let Some((peer, network)) = registrations.pop_first() {
1007 let ctx = context.child("peer").with_attribute("public_key", &peer);
1008 let config = Config {
1009 public_key: peer.clone(),
1010 mailbox_size: NZUsize!(1024),
1011 deque_size: CACHE_SIZE,
1012 priority: false,
1013 codec_config: RangeCfg::from(..),
1014 peer_provider: oracle.manager(),
1015 };
1016 let (engine, engine_mailbox) = Engine::<_, PublicKey, TestMessage, _>::new(ctx, config);
1017 mailboxes.insert(peer.clone(), engine_mailbox);
1018 handles.push(engine.start(network));
1019 }
1020
1021 context.sleep(A_JIFFY).await;
1022 (mailboxes, handles)
1023 }
1024
1025 #[test_traced]
1026 fn test_operations_after_shutdown_do_not_panic() {
1027 let runner = deterministic::Runner::timed(Duration::from_secs(5));
1028 runner.start(|context| async move {
1029 let (peers, mut registrations, oracle) =
1030 initialize_simulation(context.child("network"), 2, probability!(1.0)).await;
1031 let (mut mailboxes, handles) = spawn_peer_engines_with_handles(
1032 context.child("peers"),
1033 &oracle,
1034 &mut registrations,
1035 )
1036 .await;
1037
1038 let message = TestMessage::shared(b"test message");
1040 let mailbox = mailboxes.remove(&peers[0]).unwrap();
1041 assert!(
1042 mailbox
1043 .broadcast(Recipients::All, message.clone())
1044 .accepted(),
1045 "broadcast should succeed before shutdown"
1046 );
1047
1048 for handle in handles {
1050 handle.abort();
1051 }
1052 context.sleep(Duration::from_millis(100)).await;
1053
1054 assert_eq!(
1058 mailbox.broadcast(Recipients::All, message.clone()),
1059 Feedback::Closed,
1060 "broadcast after shutdown should return Closed"
1061 );
1062
1063 let digest = message.digest();
1065 let receiver = mailbox.subscribe(digest);
1066 let result = receiver.await;
1067 assert!(
1068 result.is_err(),
1069 "subscribe after shutdown should return Canceled"
1070 );
1071
1072 let result = mailbox.get(digest).await;
1074 assert!(result.is_none(), "get after shutdown should return None");
1075 });
1076 }
1077
1078 fn clean_shutdown(seed: u64) {
1079 let cfg = deterministic::Config::new()
1080 .with_seed(seed)
1081 .with_timeout(Some(Duration::from_secs(30)));
1082 let runner = deterministic::Runner::new(cfg);
1083 runner.start(|context| async move {
1084 let (peers, mut registrations, oracle) =
1085 initialize_simulation(context.child("network"), 2, probability!(1.0)).await;
1086
1087 let (mailboxes, handles) = spawn_peer_engines_with_handles(
1088 context.child("peers"),
1089 &oracle,
1090 &mut registrations,
1091 )
1092 .await;
1093
1094 context.sleep(Duration::from_millis(100)).await;
1096
1097 let running_before = count_running_tasks(&context, "peers");
1099 assert!(
1100 running_before > 0,
1101 "at least one peer engine task should be running"
1102 );
1103
1104 let message = TestMessage::shared(b"test message");
1106 let mailbox = mailboxes.get(&peers[0]).unwrap().clone();
1107 assert!(
1108 mailbox
1109 .broadcast(Recipients::All, message.clone())
1110 .accepted(),
1111 "broadcast should succeed"
1112 );
1113
1114 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
1116
1117 let peer_mailbox = mailboxes.get(&peers[1]).unwrap().clone();
1119 let received = peer_mailbox.get(message.digest()).await;
1120 assert_eq!(received.as_deref(), Some(&message));
1121
1122 for handle in handles {
1124 handle.abort();
1125 }
1126 context.sleep(Duration::from_millis(100)).await;
1127
1128 let running_after = count_running_tasks(&context, "peers");
1130 assert_eq!(
1131 running_after, 0,
1132 "all peer engine tasks should be stopped, but {running_after} still running"
1133 );
1134 });
1135 }
1136
1137 #[test]
1138 fn test_clean_shutdown() {
1139 for seed in 0..25 {
1140 clean_shutdown(seed);
1141 }
1142 }
1143
1144 #[test_traced]
1145 fn test_peer_set_update_evicts_disconnected_peer_buffers() {
1146 let runner = deterministic::Runner::timed(Duration::from_secs(5));
1147 runner.start(|context| async move {
1148 let (peers, mut registrations, oracle) =
1149 initialize_simulation(context.child("network"), 3, probability!(1.0)).await;
1150
1151 let peer_a = peers[0].clone();
1152 let peer_b = peers[1].clone();
1153 let peer_c = peers[2].clone();
1154
1155 let network_b = registrations.remove(&peer_b).unwrap();
1157 let config_b = Config {
1158 public_key: peer_b.clone(),
1159 mailbox_size: NZUsize!(1024),
1160 deque_size: CACHE_SIZE,
1161 priority: false,
1162 codec_config: RangeCfg::from(..),
1163 peer_provider: oracle.manager(),
1164 };
1165 let (engine_b, mailbox_b) =
1166 Engine::<_, PublicKey, TestMessage, _>::new(context.child("peer_b"), config_b);
1167 engine_b.start(network_b);
1168
1169 let mut mailboxes = BTreeMap::new();
1171 mailboxes.insert(peer_b.clone(), mailbox_b);
1172 for (peer, network) in registrations {
1173 let ctx = context.child("peer").with_attribute("public_key", &peer);
1174 let config = Config {
1175 public_key: peer.clone(),
1176 mailbox_size: NZUsize!(1024),
1177 deque_size: CACHE_SIZE,
1178 priority: false,
1179 codec_config: RangeCfg::from(..),
1180 peer_provider: oracle.manager(),
1181 };
1182 let (engine, mailbox) = Engine::<_, PublicKey, TestMessage, _>::new(ctx, config);
1183 mailboxes.insert(peer, mailbox);
1184 engine.start(network);
1185 }
1186 context.sleep(A_JIFFY).await;
1187
1188 let msg = TestMessage::shared(b"eviction-test");
1190 let mailbox_a = mailboxes.get(&peer_a).unwrap().clone();
1191 assert!(mailbox_a.broadcast(Recipients::All, msg.clone()).accepted());
1192 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
1193
1194 let mailbox_b = mailboxes.get(&peer_b).unwrap().clone();
1196 assert_eq!(
1197 mailbox_b.get(msg.digest()).await.as_deref(),
1198 Some(&msg),
1199 "peer B should have the message before eviction"
1200 );
1201
1202 let remaining = commonware_utils::ordered::Set::from_iter_dedup(vec![peer_b, peer_c]);
1204 oracle.manager().track(1, remaining);
1205 context.sleep(A_JIFFY).await;
1206
1207 assert!(
1209 mailbox_b.get(msg.digest()).await.is_none(),
1210 "message should be evicted after peer A left the peer set"
1211 );
1212 });
1213 }
1214
1215 #[test_traced]
1216 fn test_peer_set_update_evicts_peers_not_in_latest_set_even_if_still_in_overlap() {
1217 let runner = deterministic::Runner::timed(Duration::from_secs(5));
1218 runner.start(|context| async move {
1219 let (network, oracle) = Network::<deterministic::Context, PublicKey>::new(
1221 context.child("network"),
1222 commonware_p2p::simulated::Config {
1223 max_size: 1024 * 1024,
1224 max_peers_per_set: NZUsize!(3),
1225 disconnect_on_block: true,
1226 tracked_peer_sets: NZUsize!(2),
1227 },
1228 );
1229 network.start();
1230
1231 let mut schemes = (0..3)
1232 .map(|i| PrivateKey::from_seed(i as u64))
1233 .collect::<Vec<_>>();
1234 schemes.sort_by_key(|s| s.public_key());
1235 let peers: Vec<PublicKey> = schemes.iter().map(|c| c.public_key()).collect();
1236 let peer_a = peers[0].clone();
1237 let peer_b = peers[1].clone();
1238 let peer_c = peers[2].clone();
1239
1240 let mut registrations: Registrations = BTreeMap::new();
1241 for peer in peers.iter() {
1242 let (sender, receiver) = oracle
1243 .control(peer.clone())
1244 .register(0, TEST_QUOTA)
1245 .await
1246 .unwrap();
1247 registrations.insert(peer.clone(), (sender, receiver));
1248 }
1249 let link = Link {
1250 latency: NETWORK_SPEED,
1251 jitter: Duration::ZERO,
1252 success_rate: probability!(1.0),
1253 };
1254 for p1 in peers.iter() {
1255 for p2 in peers.iter() {
1256 if p2 != p1 {
1257 oracle
1258 .add_link(p1.clone(), p2.clone(), link.clone())
1259 .await
1260 .unwrap();
1261 }
1262 }
1263 }
1264
1265 let all = commonware_utils::ordered::Set::from_iter_dedup(peers.clone());
1267 oracle.manager().track(0, all);
1268
1269 let network_b = registrations.remove(&peer_b).unwrap();
1271 let config_b = Config {
1272 public_key: peer_b.clone(),
1273 mailbox_size: NZUsize!(1024),
1274 deque_size: CACHE_SIZE,
1275 priority: false,
1276 codec_config: RangeCfg::from(..),
1277 peer_provider: oracle.manager(),
1278 };
1279 let (engine_b, mailbox_b) =
1280 Engine::<_, PublicKey, TestMessage, _>::new(context.child("peer_b"), config_b);
1281 engine_b.start(network_b);
1282
1283 let mut mailboxes = BTreeMap::new();
1284 mailboxes.insert(peer_b.clone(), mailbox_b);
1285 for (peer, network) in registrations {
1286 let ctx = context.child("peer").with_attribute("public_key", &peer);
1287 let config = Config {
1288 public_key: peer.clone(),
1289 mailbox_size: NZUsize!(1024),
1290 deque_size: CACHE_SIZE,
1291 priority: false,
1292 codec_config: RangeCfg::from(..),
1293 peer_provider: oracle.manager(),
1294 };
1295 let (engine, mailbox) = Engine::<_, PublicKey, TestMessage, _>::new(ctx, config);
1296 mailboxes.insert(peer, mailbox);
1297 engine.start(network);
1298 }
1299 context.sleep(A_JIFFY).await;
1300
1301 let msg = TestMessage::shared(b"eviction-latest-test");
1303 let mailbox_a = mailboxes.get(&peer_a).unwrap().clone();
1304 assert!(mailbox_a.broadcast(Recipients::All, msg.clone()).accepted());
1305 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
1306
1307 let mailbox_b = mailboxes.get(&peer_b).unwrap().clone();
1308 assert_eq!(
1309 mailbox_b.get(msg.digest()).await.as_deref(),
1310 Some(&msg),
1311 "peer B should have the message before eviction"
1312 );
1313
1314 let remaining = commonware_utils::ordered::Set::from_iter_dedup(vec![
1318 peer_b.clone(),
1319 peer_c.clone(),
1320 ]);
1321 oracle.manager().track(1, remaining);
1322 context.sleep(A_JIFFY).await;
1323
1324 assert!(
1325 mailbox_b.get(msg.digest()).await.is_none(),
1326 "message should be evicted: peer A is not in the latest peer set"
1327 );
1328
1329 let fresh = TestMessage::shared(b"post-eviction-latest-test");
1331 assert!(
1332 mailbox_a
1333 .broadcast(Recipients::All, fresh.clone())
1334 .accepted()
1335 );
1336 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
1337
1338 assert!(
1339 mailbox_b.get(fresh.digest()).await.is_none(),
1340 "message should not be rebuffered after peer A left latest.primary"
1341 );
1342 });
1343 }
1344
1345 #[test_traced]
1346 fn test_initial_latest_peer_set_blocks_sender_not_in_latest_primary() {
1347 let runner = deterministic::Runner::timed(Duration::from_secs(5));
1348 runner.start(|context| async move {
1349 let (network, oracle) = Network::<deterministic::Context, PublicKey>::new(
1350 context.child("network"),
1351 commonware_p2p::simulated::Config {
1352 max_size: 1024 * 1024,
1353 max_peers_per_set: NZUsize!(3),
1354 disconnect_on_block: true,
1355 tracked_peer_sets: NZUsize!(1),
1356 },
1357 );
1358 network.start();
1359
1360 let mut schemes = (0..3)
1361 .map(|i| PrivateKey::from_seed(i as u64))
1362 .collect::<Vec<_>>();
1363 schemes.sort_by_key(|s| s.public_key());
1364 let peers: Vec<PublicKey> = schemes.iter().map(|c| c.public_key()).collect();
1365 let peer_a = peers[0].clone();
1366 let peer_b = peers[1].clone();
1367 let peer_c = peers[2].clone();
1368
1369 let mut registrations: Registrations = BTreeMap::new();
1370 for peer in &peers {
1371 let (sender, receiver) = oracle
1372 .control(peer.clone())
1373 .register(0, TEST_QUOTA)
1374 .await
1375 .unwrap();
1376 registrations.insert(peer.clone(), (sender, receiver));
1377 }
1378 let link = Link {
1379 latency: NETWORK_SPEED,
1380 jitter: Duration::ZERO,
1381 success_rate: probability!(1.0),
1382 };
1383 for p1 in &peers {
1384 for p2 in &peers {
1385 if p1 != p2 {
1386 oracle
1387 .add_link(p1.clone(), p2.clone(), link.clone())
1388 .await
1389 .unwrap();
1390 }
1391 }
1392 }
1393
1394 let latest_primary = commonware_utils::ordered::Set::from_iter_dedup(vec![
1395 peer_b.clone(),
1396 peer_c.clone(),
1397 ]);
1398 let latest_secondary =
1399 commonware_utils::ordered::Set::from_iter_dedup(vec![peer_a.clone()]);
1400 oracle
1401 .manager()
1402 .track(0, TrackedPeers::new(latest_primary, latest_secondary));
1403
1404 let mut mailboxes = BTreeMap::new();
1405 for (peer, network) in registrations {
1406 let ctx = context.child("peer").with_attribute("public_key", &peer);
1407 let config = Config {
1408 public_key: peer.clone(),
1409 mailbox_size: NZUsize!(1024),
1410 deque_size: CACHE_SIZE,
1411 priority: false,
1412 codec_config: RangeCfg::from(..),
1413 peer_provider: oracle.manager(),
1414 };
1415 let (engine, mailbox) = Engine::<_, PublicKey, TestMessage, _>::new(ctx, config);
1416 mailboxes.insert(peer, mailbox);
1417 engine.start(network);
1418 }
1419 context.sleep(A_JIFFY).await;
1420
1421 let mailbox_a = mailboxes.get(&peer_a).unwrap().clone();
1422 let mailbox_b = mailboxes.get(&peer_b).unwrap().clone();
1423 let msg = TestMessage::shared(b"startup-latest-primary-only");
1424 assert!(
1425 mailbox_a
1426 .broadcast(Recipients::All, msg.clone())
1427 .accepted(),
1428 "Recipients::All is accepted locally; cache policy is separate"
1429 );
1430 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
1431
1432 assert_eq!(
1433 mailbox_a.get(msg.digest()).await,
1434 None,
1435 "sender not in latest.primary should not buffer, including own broadcasts"
1436 );
1437 assert!(
1438 mailbox_b.get(msg.digest()).await.is_none(),
1439 "peer B should not cache messages from a sender excluded by the initial latest.primary set"
1440 );
1441 });
1442 }
1443
1444 #[test_traced]
1448 fn test_broadcast_queued_before_start_respects_initial_latest_primary() {
1449 let runner = deterministic::Runner::timed(Duration::from_secs(5));
1450 runner.start(|context| async move {
1451 let (peers, mut registrations, oracle) =
1453 initialize_simulation(context.child("network"), 1, probability!(1.0)).await;
1454 let peer = peers[0].clone();
1455 let network = registrations.remove(&peer).unwrap();
1456 let config = Config {
1457 public_key: peer.clone(),
1458 mailbox_size: NZUsize!(1024),
1459 deque_size: CACHE_SIZE,
1460 priority: false,
1461 codec_config: RangeCfg::from(..),
1462 peer_provider: oracle.manager(),
1463 };
1464 let (engine, mailbox) =
1465 Engine::<_, PublicKey, TestMessage, _>::new(context.child("peer"), config);
1466
1467 let msg = TestMessage::shared(b"queued-before-start");
1469 assert!(mailbox
1470 .broadcast(Recipients::All, msg.clone())
1471 .accepted());
1472
1473 engine.start(network);
1475
1476 assert_eq!(
1477 mailbox.get(msg.digest()).await.as_deref(),
1478 Some(&msg),
1479 "sender is already in the initial latest.primary set, so its local broadcast should be cached"
1480 );
1481 });
1482 }
1483
1484 #[test_traced]
1485 fn test_engine_starts_before_initial_peer_set_and_delivers_after_tracking() {
1486 let runner = deterministic::Runner::timed(Duration::from_secs(5));
1487 runner.start(|context| async move {
1488 let (network, oracle) = Network::<deterministic::Context, PublicKey>::new(
1489 context.child("network"),
1490 commonware_p2p::simulated::Config {
1491 max_size: 1024 * 1024,
1492 max_peers_per_set: NZUsize!(2),
1493 disconnect_on_block: true,
1494 tracked_peer_sets: NZUsize!(1),
1495 },
1496 );
1497 network.start();
1498
1499 let mut schemes = (0..2)
1500 .map(|i| PrivateKey::from_seed(i as u64))
1501 .collect::<Vec<_>>();
1502 schemes.sort_by_key(|s| s.public_key());
1503 let peers: Vec<PublicKey> = schemes.iter().map(|c| c.public_key()).collect();
1504 let peer_a = peers[0].clone();
1505 let peer_b = peers[1].clone();
1506
1507 let mut registrations: Registrations = BTreeMap::new();
1508 for peer in &peers {
1509 let (sender, receiver) = oracle
1510 .control(peer.clone())
1511 .register(0, TEST_QUOTA)
1512 .await
1513 .unwrap();
1514 registrations.insert(peer.clone(), (sender, receiver));
1515 }
1516
1517 let link = Link {
1518 latency: NETWORK_SPEED,
1519 jitter: Duration::ZERO,
1520 success_rate: probability!(1.0),
1521 };
1522 for p1 in &peers {
1523 for p2 in &peers {
1524 if p1 != p2 {
1525 oracle
1526 .add_link(p1.clone(), p2.clone(), link.clone())
1527 .await
1528 .unwrap();
1529 }
1530 }
1531 }
1532
1533 let mut mailboxes = BTreeMap::new();
1534 for (peer, network) in registrations {
1535 let ctx = context.child("peer").with_attribute("public_key", &peer);
1536 let config = Config {
1537 public_key: peer.clone(),
1538 mailbox_size: NZUsize!(1024),
1539 deque_size: CACHE_SIZE,
1540 priority: false,
1541 codec_config: RangeCfg::from(..),
1542 peer_provider: oracle.manager(),
1543 };
1544 let (engine, mailbox) = Engine::<_, PublicKey, TestMessage, _>::new(ctx, config);
1545 mailboxes.insert(peer, mailbox);
1546 engine.start(network);
1547 }
1548
1549 let mailbox_a = mailboxes.get(&peer_a).unwrap().clone();
1550 let mailbox_b = mailboxes.get(&peer_b).unwrap().clone();
1551
1552 let before = TestMessage::shared(b"before-tracking");
1553 assert!(
1554 mailbox_a
1555 .broadcast(Recipients::All, before.clone())
1556 .accepted(),
1557 "broadcast request should be accepted before a peer set is tracked"
1558 );
1559 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
1560
1561 assert_eq!(
1562 mailbox_a.get(before.digest()).await,
1563 None,
1564 "without latest.primary, local broadcasts are not buffered"
1565 );
1566 assert!(
1567 mailbox_b.get(before.digest()).await.is_none(),
1568 "without latest.primary, remote peers do not cache inbound messages"
1569 );
1570
1571 oracle.manager().track(
1572 0,
1573 commonware_utils::ordered::Set::from_iter_dedup(peers.clone()),
1574 );
1575 context.sleep(A_JIFFY).await;
1576
1577 let after = TestMessage::shared(b"after-tracking");
1578 assert!(
1579 mailbox_a
1580 .broadcast(Recipients::All, after.clone())
1581 .accepted()
1582 );
1583 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
1584
1585 assert_eq!(mailbox_b.get(after.digest()).await.as_deref(), Some(&after));
1586 });
1587 }
1588
1589 #[test_traced]
1590 fn test_peer_set_update_preserves_shared_messages() {
1591 let runner = deterministic::Runner::timed(Duration::from_secs(5));
1592 runner.start(|context| async move {
1593 let (peers, mut registrations, oracle) =
1594 initialize_simulation(context.child("network"), 3, probability!(1.0)).await;
1595
1596 let peer_a = peers[0].clone();
1597 let peer_b = peers[1].clone();
1598 let peer_c = peers[2].clone();
1599
1600 let network_b = registrations.remove(&peer_b).unwrap();
1602 let config_b = Config {
1603 public_key: peer_b.clone(),
1604 mailbox_size: NZUsize!(1024),
1605 deque_size: CACHE_SIZE,
1606 priority: false,
1607 codec_config: RangeCfg::from(..),
1608 peer_provider: oracle.manager(),
1609 };
1610 let (engine_b, mailbox_b) =
1611 Engine::<_, PublicKey, TestMessage, _>::new(context.child("peer_b"), config_b);
1612 engine_b.start(network_b);
1613
1614 let mut mailboxes = BTreeMap::new();
1616 mailboxes.insert(peer_b.clone(), mailbox_b);
1617 for (peer, network) in registrations {
1618 let ctx = context.child("peer").with_attribute("public_key", &peer);
1619 let config = Config {
1620 public_key: peer.clone(),
1621 mailbox_size: NZUsize!(1024),
1622 deque_size: CACHE_SIZE,
1623 priority: false,
1624 codec_config: RangeCfg::from(..),
1625 peer_provider: oracle.manager(),
1626 };
1627 let (engine, mailbox) = Engine::<_, PublicKey, TestMessage, _>::new(ctx, config);
1628 mailboxes.insert(peer, mailbox);
1629 engine.start(network);
1630 }
1631 context.sleep(A_JIFFY).await;
1632
1633 let msg = TestMessage::shared(b"shared-msg");
1635 let mailbox_a = mailboxes.get(&peer_a).unwrap().clone();
1636 let mailbox_c = mailboxes.get(&peer_c).unwrap().clone();
1637 assert!(mailbox_a.broadcast(Recipients::All, msg.clone()).accepted());
1638 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
1639 assert!(mailbox_c.broadcast(Recipients::All, msg.clone()).accepted());
1640 context.sleep(NETWORK_SPEED_WITH_BUFFER).await;
1641
1642 let mailbox_b = mailboxes.get(&peer_b).unwrap().clone();
1644 assert_eq!(mailbox_b.get(msg.digest()).await.as_deref(), Some(&msg));
1645
1646 let remaining = commonware_utils::ordered::Set::from_iter_dedup(vec![peer_b, peer_c]);
1648 oracle.manager().track(1, remaining);
1649 context.sleep(A_JIFFY).await;
1650
1651 assert_eq!(
1653 mailbox_b.get(msg.digest()).await.as_deref(),
1654 Some(&msg),
1655 "message should survive when another peer in the primary set still references it"
1656 );
1657 });
1658 }
1659}