1mod bandwidth;
143mod ingress;
144mod metrics;
145mod network;
146mod transmitter;
147
148use thiserror::Error;
149
150#[derive(Debug, Error)]
152pub enum Error {
153 #[error("network closed")]
154 NetworkClosed,
155 #[error("not valid to link self")]
156 LinkingSelf,
157 #[error("link already exists")]
158 LinkExists,
159 #[error("link missing")]
160 LinkMissing,
161 #[error("send_frame failed")]
162 SendFrameFailed,
163 #[error("recv_frame failed")]
164 RecvFrameFailed,
165 #[error("bind failed")]
166 BindFailed,
167 #[error("accept failed")]
168 AcceptFailed,
169 #[error("dial failed")]
170 DialFailed,
171 #[error("peer missing")]
172 PeerMissing,
173}
174
175pub use ingress::{Control, Link, Manager, Oracle, SocketManager};
176pub use network::{
177 Config, ConnectedPeerProvider, MAX_SIZE, Network, Receiver, Sender, SplitForwarder,
178 SplitOrigin, SplitRouter, SplitSender, SplitTarget, UnlimitedSender,
179};
180
181#[cfg(test)]
182mod tests {
183 use super::*;
184 use crate::{
185 Address, AddressableManager, AddressableTrackedPeers, Blocker as _, CheckedSender as _,
186 Ingress, LimitedSender as _, Manager, Provider, Receiver, Recipients, Sender, TrackedPeers,
187 };
188 use commonware_cryptography::{
189 Signer as _,
190 ed25519::{self, PrivateKey, PublicKey},
191 };
192 use commonware_macros::{select, test_group};
193 use commonware_runtime::{
194 Clock, IoBuf, Quota, Runner, Spawner, Supervisor as _, deterministic, reschedule,
195 telemetry::metrics::count_running_tasks,
196 };
197 use commonware_utils::{
198 NZU32, NZUsize,
199 channel::mpsc,
200 hostname, ordered,
201 ordered::{Map, Set},
202 probability,
203 };
204 use futures::StreamExt;
205 use rand::RngExt as _;
206 use std::{
207 collections::{BTreeMap, HashMap, HashSet},
208 net::SocketAddr,
209 num::NonZeroU32,
210 time::Duration,
211 };
212
213 const TEST_QUOTA: Quota = Quota::per_second(NonZeroU32::MAX);
215
216 async fn track_peers<I>(oracle: &Oracle<PublicKey, deterministic::Context>, peers: I)
217 where
218 I: IntoIterator<Item = PublicKey>,
219 {
220 let mut manager = oracle.manager();
221 manager.track(0, Set::from_iter_dedup(peers));
222 assert!(manager.peer_set(0).await.is_some());
223 }
224
225 async fn wait_for_task_count(
226 context: &deterministic::Context,
227 prefix: &str,
228 expected: impl Fn(usize) -> bool,
229 ) {
230 loop {
231 let count = count_running_tasks(context, prefix);
232 if expected(count) {
233 return;
234 }
235 reschedule().await;
236 }
237 }
238
239 fn simulate_messages(seed: u64, size: usize) -> (String, Vec<usize>) {
240 let executor = deterministic::Runner::seeded(seed);
241 executor.start(|context| async move {
242 let (network, oracle) = Network::new(
244 context.child("network"),
245 Config {
246 max_size: 1024 * 1024,
247 max_peers_per_set: NZUsize!(size),
248 disconnect_on_block: true,
249 tracked_peer_sets: NZUsize!(1),
250 },
251 );
252
253 network.start();
255
256 let mut agents = BTreeMap::new();
258 let (seen_sender, mut seen_receiver) = mpsc::channel(1024);
259 for i in 0..size {
260 let pk = PrivateKey::from_seed(i as u64).public_key();
261 let (sender, mut receiver) = oracle
262 .control(pk.clone())
263 .register(0, TEST_QUOTA)
264 .await
265 .unwrap();
266 agents.insert(pk, sender);
267 let agent_sender = seen_sender.clone();
268 context.child("agent_receiver").spawn(move |_| async move {
269 for _ in 0..size {
270 receiver.recv().await.unwrap();
271 }
272 agent_sender.send(i).await.unwrap();
273
274 });
276 }
277 track_peers(&oracle, agents.keys().cloned()).await;
278
279 let only_inbound = PrivateKey::from_seed(0).public_key();
281 for agent in agents.keys() {
282 if agent == &only_inbound {
283 continue;
285 }
286 for other in agents.keys() {
287 let result = oracle
288 .add_link(
289 agent.clone(),
290 other.clone(),
291 Link {
292 latency: Duration::from_millis(5),
293 jitter: Duration::from_millis(2),
294 success_rate: probability!(0.75),
295 },
296 )
297 .await;
298 if agent == other {
299 assert!(matches!(result, Err(Error::LinkingSelf)));
300 } else {
301 assert!(result.is_ok());
302 }
303 }
304 }
305
306 context
307 .child("agent_sender")
308 .spawn(|mut context| async move {
309 let keys = agents.keys().cloned().collect::<Vec<_>>();
311
312 loop {
313 let index = context.random_range(0..keys.len());
314 let sender = &keys[index];
315 let msg = format!("hello from {sender:?}");
316 let msg = IoBuf::copy_from_slice(msg.as_bytes());
317 let message_sender = agents.get_mut(sender).unwrap();
318 let sent = message_sender.send(Recipients::All, msg, false);
319 assert_eq!(sent.len(), keys.len() - 1);
320 reschedule().await;
321 }
322 });
323
324 let mut results = Vec::new();
326 for _ in 0..size {
327 results.push(seen_receiver.recv().await.unwrap());
328 }
329 (context.auditor().state(), results)
330 })
331 }
332
333 fn compare_outputs(seeds: u64, size: usize) {
334 let mut outputs = Vec::new();
336 for seed in 0..seeds {
337 outputs.push(simulate_messages(seed, size));
338 }
339
340 for seed in 0..seeds {
342 let output = simulate_messages(seed, size);
343 assert_eq!(output, outputs[seed as usize]);
344 }
345 }
346
347 #[test_group("slow")]
348 #[test]
349 fn test_determinism() {
350 compare_outputs(25, 25);
351 }
352
353 #[test]
354 #[should_panic(expected = "message too large")]
355 fn test_message_too_big() {
356 let executor = deterministic::Runner::default();
357 executor.start(|mut context| async move {
358 let (network, oracle) = Network::new(
360 context.child("network"),
361 Config {
362 max_size: 1024 * 1024,
363 max_peers_per_set: NZUsize!(1),
364 disconnect_on_block: true,
365 tracked_peer_sets: NZUsize!(1),
366 },
367 );
368
369 network.start();
371
372 let mut agents = HashMap::new();
374 for i in 0..10 {
375 let pk = PrivateKey::from_seed(i as u64).public_key();
376 let (sender, _) = oracle
377 .control(pk.clone())
378 .register(0, TEST_QUOTA)
379 .await
380 .unwrap();
381 agents.insert(pk, sender);
382 }
383
384 let keys = agents.keys().collect::<Vec<_>>();
386 let index = context.random_range(0..keys.len());
387 let sender = keys[index];
388 let mut message_sender = agents.get(sender).unwrap().clone();
389 let mut msg = vec![0u8; 1024 * 1024 + 1];
390 context.fill(&mut msg[..]);
391 message_sender.send(Recipients::All, msg, false);
392 });
393 }
394
395 #[test]
396 fn test_linking_self() {
397 let executor = deterministic::Runner::default();
398 executor.start(|context| async move {
399 let (network, oracle) = Network::new(
401 context.child("network"),
402 Config {
403 max_size: 1024 * 1024,
404 max_peers_per_set: NZUsize!(1),
405 disconnect_on_block: true,
406 tracked_peer_sets: NZUsize!(1),
407 },
408 );
409
410 network.start();
412
413 let pk = PrivateKey::from_seed(0).public_key();
415 oracle
416 .control(pk.clone())
417 .register(0, TEST_QUOTA)
418 .await
419 .unwrap();
420
421 let result = oracle
423 .add_link(
424 pk.clone(),
425 pk,
426 Link {
427 latency: Duration::from_millis(5),
428 jitter: Duration::from_millis(2),
429 success_rate: probability!(0.75),
430 },
431 )
432 .await;
433
434 assert!(matches!(result, Err(Error::LinkingSelf)));
436 });
437 }
438
439 #[test]
440 fn test_duplicate_channel() {
441 let executor = deterministic::Runner::default();
442 executor.start(|context| async move {
443 let (network, oracle) = Network::new(
445 context.child("network"),
446 Config {
447 max_size: 1024 * 1024,
448 max_peers_per_set: NZUsize!(2),
449 disconnect_on_block: true,
450 tracked_peer_sets: NZUsize!(1),
451 },
452 );
453
454 network.start();
456
457 let my_pk = PrivateKey::from_seed(0).public_key();
459 let other_pk = PrivateKey::from_seed(1).public_key();
460 oracle
461 .add_link(
462 my_pk.clone(),
463 other_pk.clone(),
464 Link {
465 latency: Duration::from_millis(10),
466 jitter: Duration::from_millis(1),
467 success_rate: probability!(1.0),
468 },
469 )
470 .await
471 .unwrap();
472 oracle
473 .add_link(
474 other_pk.clone(),
475 my_pk.clone(),
476 Link {
477 latency: Duration::from_millis(10),
478 jitter: Duration::from_millis(1),
479 success_rate: probability!(1.0),
480 },
481 )
482 .await
483 .unwrap();
484
485 let (mut my_sender, mut my_receiver) = oracle
487 .control(my_pk.clone())
488 .register(0, TEST_QUOTA)
489 .await
490 .unwrap();
491 let (mut other_sender, mut other_receiver) = oracle
492 .control(other_pk.clone())
493 .register(0, TEST_QUOTA)
494 .await
495 .unwrap();
496 track_peers(&oracle, [my_pk.clone(), other_pk.clone()]).await;
497
498 let msg = IoBuf::from(b"hello");
500 let sent = my_sender.send(Recipients::One(other_pk.clone()), msg.clone(), false);
501 assert_eq!(sent.len(), 1);
502 let (from, message) = other_receiver.recv().await.unwrap();
503 assert_eq!(from, my_pk);
504 assert_eq!(message, msg.clone());
505 let sent = other_sender.send(Recipients::One(my_pk.clone()), msg.clone(), false);
506 assert_eq!(sent.len(), 1);
507 let (from, message) = my_receiver.recv().await.unwrap();
508 assert_eq!(from, other_pk);
509 assert_eq!(message, msg);
510
511 let (mut my_sender_2, mut my_receiver_2) = oracle
513 .control(my_pk.clone())
514 .register(0, TEST_QUOTA)
515 .await
516 .unwrap();
517
518 let msg = IoBuf::from(b"hello again");
520 let sent = my_sender_2.send(Recipients::One(other_pk.clone()), msg.clone(), false);
521 assert_eq!(sent.len(), 1);
522 let (from, message) = other_receiver.recv().await.unwrap();
523 assert_eq!(from, my_pk);
524 assert_eq!(message, msg.clone());
525 let sent = other_sender.send(Recipients::One(my_pk.clone()), msg.clone(), false);
526 assert_eq!(sent.len(), 1);
527 let (from, message) = my_receiver_2.recv().await.unwrap();
528 assert_eq!(from, other_pk);
529 assert_eq!(message, msg.clone());
530
531 assert!(matches!(
533 my_receiver.recv().await,
534 Err(Error::NetworkClosed)
535 ));
536
537 my_sender.send(Recipients::One(other_pk.clone()), msg, false);
539 });
540 }
541
542 #[test]
543 fn test_add_link_before_channel_registration() {
544 let executor = deterministic::Runner::default();
545 executor.start(|context| async move {
546 let pk1 = PrivateKey::from_seed(0).public_key();
548 let pk2 = PrivateKey::from_seed(1).public_key();
549
550 let (network, oracle) = Network::new_with_peers(
552 context.child("network"),
553 Config {
554 max_size: 1024 * 1024,
555 max_peers_per_set: NZUsize!(2),
556 disconnect_on_block: true,
557 tracked_peer_sets: NZUsize!(3),
558 },
559 [pk1.clone(), pk2.clone()],
560 )
561 .await;
562 network.start();
563
564 oracle
566 .add_link(
567 pk1.clone(),
568 pk2.clone(),
569 Link {
570 latency: Duration::ZERO,
571 jitter: Duration::ZERO,
572 success_rate: probability!(1.0),
573 },
574 )
575 .await
576 .unwrap();
577
578 let (mut sender1, _receiver1) = oracle
580 .control(pk1.clone())
581 .register(0, TEST_QUOTA)
582 .await
583 .unwrap();
584 let (_, mut receiver2) = oracle
585 .control(pk2.clone())
586 .register(0, TEST_QUOTA)
587 .await
588 .unwrap();
589
590 let msg1 = IoBuf::from(b"link-before-register-1");
592 sender1.send(Recipients::One(pk2.clone()), msg1.clone(), false);
593 let (from, received) = receiver2.recv().await.unwrap();
594 assert_eq!(from, pk1);
595 assert_eq!(received, msg1);
596 });
597 }
598
599 #[test]
600 fn test_simple_message_delivery() {
601 let executor = deterministic::Runner::default();
602 executor.start(|context| async move {
603 let (network, oracle) = Network::new(
605 context.child("network"),
606 Config {
607 max_size: 1024 * 1024,
608 max_peers_per_set: NZUsize!(2),
609 disconnect_on_block: true,
610 tracked_peer_sets: NZUsize!(1),
611 },
612 );
613
614 network.start();
616
617 let pk1 = PrivateKey::from_seed(0).public_key();
619 let pk2 = PrivateKey::from_seed(1).public_key();
620 let (mut sender1, mut receiver1) = oracle
621 .control(pk1.clone())
622 .register(0, TEST_QUOTA)
623 .await
624 .unwrap();
625 let (mut sender2, mut receiver2) = oracle
626 .control(pk2.clone())
627 .register(0, TEST_QUOTA)
628 .await
629 .unwrap();
630
631 let _ = oracle
633 .control(pk1.clone())
634 .register(1, TEST_QUOTA)
635 .await
636 .unwrap();
637 let _ = oracle
638 .control(pk2.clone())
639 .register(2, TEST_QUOTA)
640 .await
641 .unwrap();
642 track_peers(&oracle, [pk1.clone(), pk2.clone()]).await;
643
644 oracle
646 .add_link(
647 pk1.clone(),
648 pk2.clone(),
649 Link {
650 latency: Duration::from_millis(5),
651 jitter: Duration::from_millis(2),
652 success_rate: probability!(1.0),
653 },
654 )
655 .await
656 .unwrap();
657 oracle
658 .add_link(
659 pk2.clone(),
660 pk1.clone(),
661 Link {
662 latency: Duration::from_millis(5),
663 jitter: Duration::from_millis(2),
664 success_rate: probability!(1.0),
665 },
666 )
667 .await
668 .unwrap();
669
670 let msg1 = IoBuf::from(b"hello from pk1");
672 let msg2 = IoBuf::from(b"hello from pk2");
673 sender1.send(Recipients::One(pk2.clone()), msg1.clone(), false);
674 sender2.send(Recipients::One(pk1.clone()), msg2.clone(), false);
675
676 let (sender, message) = receiver1.recv().await.unwrap();
678 assert_eq!(sender, pk2);
679 assert_eq!(message, msg2);
680 let (sender, message) = receiver2.recv().await.unwrap();
681 assert_eq!(sender, pk1);
682 assert_eq!(message, msg1);
683 });
684 }
685
686 #[test]
687 fn test_send_wrong_channel() {
688 let executor = deterministic::Runner::default();
689 executor.start(|context| async move {
690 let (network, oracle) = Network::new(
692 context.child("network"),
693 Config {
694 max_size: 1024 * 1024,
695 max_peers_per_set: NZUsize!(2),
696 disconnect_on_block: true,
697 tracked_peer_sets: NZUsize!(1),
698 },
699 );
700
701 network.start();
703
704 let pk1 = PrivateKey::from_seed(0).public_key();
706 let pk2 = PrivateKey::from_seed(1).public_key();
707 let (mut sender1, _) = oracle
708 .control(pk1.clone())
709 .register(0, TEST_QUOTA)
710 .await
711 .unwrap();
712 let (_, mut receiver2) = oracle
713 .control(pk2.clone())
714 .register(1, TEST_QUOTA)
715 .await
716 .unwrap();
717 track_peers(&oracle, [pk1.clone(), pk2.clone()]).await;
718
719 oracle
721 .add_link(
722 pk1,
723 pk2.clone(),
724 Link {
725 latency: Duration::from_millis(5),
726 jitter: Duration::ZERO,
727 success_rate: probability!(1.0),
728 },
729 )
730 .await
731 .unwrap();
732
733 let msg = IoBuf::from(b"hello from pk1");
735 sender1.send(Recipients::One(pk2), msg, false);
736
737 select! {
739 _ = receiver2.recv() => {
740 panic!("unexpected message");
741 },
742 _ = context.sleep(Duration::from_secs(1)) => {},
743 }
744 });
745 }
746
747 #[test]
748 fn test_dynamic_peers() {
749 let executor = deterministic::Runner::default();
750 executor.start(|context| async move {
751 let (network, oracle) = Network::new(
753 context.child("network"),
754 Config {
755 max_size: 1024 * 1024,
756 max_peers_per_set: NZUsize!(2),
757 disconnect_on_block: true,
758 tracked_peer_sets: NZUsize!(1),
759 },
760 );
761
762 network.start();
764
765 let pk1 = PrivateKey::from_seed(0).public_key();
767 let pk2 = PrivateKey::from_seed(1).public_key();
768 let (mut sender1, mut receiver1) = oracle
769 .control(pk1.clone())
770 .register(0, TEST_QUOTA)
771 .await
772 .unwrap();
773 let (mut sender2, mut receiver2) = oracle
774 .control(pk2.clone())
775 .register(0, TEST_QUOTA)
776 .await
777 .unwrap();
778 track_peers(&oracle, [pk1.clone(), pk2.clone()]).await;
779
780 oracle
782 .add_link(
783 pk1.clone(),
784 pk2.clone(),
785 Link {
786 latency: Duration::from_millis(5),
787 jitter: Duration::from_millis(2),
788 success_rate: probability!(1.0),
789 },
790 )
791 .await
792 .unwrap();
793 oracle
794 .add_link(
795 pk2.clone(),
796 pk1.clone(),
797 Link {
798 latency: Duration::from_millis(5),
799 jitter: Duration::from_millis(2),
800 success_rate: probability!(1.0),
801 },
802 )
803 .await
804 .unwrap();
805
806 let msg1 = IoBuf::from(b"attempt 1: hello from pk1");
808 let msg2 = IoBuf::from(b"attempt 1: hello from pk2");
809 sender1.send(Recipients::One(pk2.clone()), msg1.clone(), false);
810 sender2.send(Recipients::One(pk1.clone()), msg2.clone(), false);
811
812 let (sender, message) = receiver1.recv().await.unwrap();
814 assert_eq!(sender, pk2);
815 assert_eq!(message, msg2);
816 let (sender, message) = receiver2.recv().await.unwrap();
817 assert_eq!(sender, pk1);
818 assert_eq!(message, msg1);
819 });
820 }
821
822 #[test]
823 fn test_dynamic_links() {
824 let executor = deterministic::Runner::default();
825 executor.start(|context| async move {
826 let (network, oracle) = Network::new(
828 context.child("network"),
829 Config {
830 max_size: 1024 * 1024,
831 max_peers_per_set: NZUsize!(2),
832 disconnect_on_block: true,
833 tracked_peer_sets: NZUsize!(1),
834 },
835 );
836
837 network.start();
839
840 let pk1 = PrivateKey::from_seed(0).public_key();
842 let pk2 = PrivateKey::from_seed(1).public_key();
843 let (mut sender1, mut receiver1) = oracle
844 .control(pk1.clone())
845 .register(0, TEST_QUOTA)
846 .await
847 .unwrap();
848 let (mut sender2, mut receiver2) = oracle
849 .control(pk2.clone())
850 .register(0, TEST_QUOTA)
851 .await
852 .unwrap();
853 track_peers(&oracle, [pk1.clone(), pk2.clone()]).await;
854
855 let msg1 = IoBuf::from(b"attempt 1: hello from pk1");
857 let msg2 = IoBuf::from(b"attempt 1: hello from pk2");
858 sender1.send(Recipients::One(pk2.clone()), msg1.clone(), false);
859 sender2.send(Recipients::One(pk1.clone()), msg2.clone(), false);
860
861 select! {
863 _ = receiver1.recv() => {
864 panic!("unexpected message");
865 },
866 _ = receiver2.recv() => {
867 panic!("unexpected message");
868 },
869 _ = context.sleep(Duration::from_secs(1)) => {},
870 }
871
872 oracle
874 .add_link(
875 pk1.clone(),
876 pk2.clone(),
877 Link {
878 latency: Duration::from_millis(5),
879 jitter: Duration::from_millis(2),
880 success_rate: probability!(1.0),
881 },
882 )
883 .await
884 .unwrap();
885 oracle
886 .add_link(
887 pk2.clone(),
888 pk1.clone(),
889 Link {
890 latency: Duration::from_millis(5),
891 jitter: Duration::from_millis(2),
892 success_rate: probability!(1.0),
893 },
894 )
895 .await
896 .unwrap();
897
898 let msg1 = IoBuf::from(b"attempt 2: hello from pk1");
900 let msg2 = IoBuf::from(b"attempt 2: hello from pk2");
901 sender1.send(Recipients::One(pk2.clone()), msg1.clone(), false);
902 sender2.send(Recipients::One(pk1.clone()), msg2.clone(), false);
903
904 let (sender, message) = receiver1.recv().await.unwrap();
906 assert_eq!(sender, pk2);
907 assert_eq!(message, msg2);
908 let (sender, message) = receiver2.recv().await.unwrap();
909 assert_eq!(sender, pk1);
910 assert_eq!(message, msg1);
911
912 oracle.remove_link(pk1.clone(), pk2.clone()).await.unwrap();
914 oracle.remove_link(pk2.clone(), pk1.clone()).await.unwrap();
915
916 let msg1 = IoBuf::from(b"attempt 3: hello from pk1");
918 let msg2 = IoBuf::from(b"attempt 3: hello from pk2");
919 sender1.send(Recipients::One(pk2.clone()), msg1.clone(), false);
920 sender2.send(Recipients::One(pk1.clone()), msg2.clone(), false);
921
922 select! {
924 _ = receiver1.recv() => {
925 panic!("unexpected message");
926 },
927 _ = receiver2.recv() => {
928 panic!("unexpected message");
929 },
930 _ = context.sleep(Duration::from_secs(1)) => {},
931 }
932
933 let result = oracle.remove_link(pk1, pk2).await;
935 assert!(matches!(result, Err(Error::LinkMissing)));
936 });
937 }
938
939 async fn test_bandwidth_between_peers(
940 context: &mut deterministic::Context,
941 oracle: &Oracle<PublicKey, deterministic::Context>,
942 index: u64,
943 sender_bps: Option<usize>,
944 receiver_bps: Option<usize>,
945 message_size: usize,
946 expected_duration_ms: u64,
947 ) {
948 let pk1 = PrivateKey::from_seed(context.random::<u64>()).public_key();
950 let pk2 = PrivateKey::from_seed(context.random::<u64>()).public_key();
951 let (mut sender, _) = oracle
952 .control(pk1.clone())
953 .register(0, TEST_QUOTA)
954 .await
955 .unwrap();
956 let (_, mut receiver) = oracle
957 .control(pk2.clone())
958 .register(0, TEST_QUOTA)
959 .await
960 .unwrap();
961 let mut manager = oracle.manager();
962 manager.track(index, Set::from_iter_dedup([pk1.clone(), pk2.clone()]));
963
964 oracle
966 .limit_bandwidth(pk1.clone(), sender_bps, None)
967 .await
968 .unwrap();
969 oracle
970 .limit_bandwidth(pk2.clone(), None, receiver_bps)
971 .await
972 .unwrap();
973
974 oracle
976 .add_link(
977 pk1.clone(),
978 pk2.clone(),
979 Link {
980 latency: Duration::ZERO,
982 jitter: Duration::ZERO,
983 success_rate: probability!(1.0),
984 },
985 )
986 .await
987 .unwrap();
988
989 let msg = IoBuf::from(vec![42u8; message_size]);
991 let start = context.current();
992 sender.send(Recipients::One(pk2.clone()), msg.clone(), true);
993
994 let (origin, received) = receiver.recv().await.unwrap();
996 let elapsed = context.current().duration_since(start).unwrap();
997
998 assert_eq!(origin, pk1);
999 assert_eq!(received, msg);
1000 assert!(
1001 elapsed >= Duration::from_millis(expected_duration_ms),
1002 "Message arrived too quickly: {elapsed:?} (expected >= {expected_duration_ms}ms)"
1003 );
1004 assert!(
1005 elapsed < Duration::from_millis(expected_duration_ms + 100),
1006 "Message took too long: {elapsed:?} (expected ~{expected_duration_ms}ms)"
1007 );
1008 }
1009
1010 #[test]
1011 fn test_bandwidth() {
1012 let executor = deterministic::Runner::default();
1013 executor.start(|mut context| async move {
1014 let (network, oracle) = Network::new(
1015 context.child("network"),
1016 Config {
1017 max_size: 1024 * 1024,
1018 max_peers_per_set: NZUsize!(2),
1019 disconnect_on_block: true,
1020 tracked_peer_sets: NZUsize!(1),
1021 },
1022 );
1023 network.start();
1024
1025 test_bandwidth_between_peers(
1028 &mut context,
1029 &oracle,
1030 0,
1031 Some(1000), Some(1000), 500, 500, )
1036 .await;
1037
1038 test_bandwidth_between_peers(
1042 &mut context,
1043 &oracle,
1044 1,
1045 Some(500), Some(2000), 250, 500, )
1050 .await;
1051
1052 test_bandwidth_between_peers(
1056 &mut context,
1057 &oracle,
1058 2,
1059 Some(2000), Some(500), 250, 500, )
1064 .await;
1065
1066 test_bandwidth_between_peers(
1070 &mut context,
1071 &oracle,
1072 3,
1073 None, Some(1000), 500, 500, )
1078 .await;
1079
1080 test_bandwidth_between_peers(
1084 &mut context,
1085 &oracle,
1086 4,
1087 Some(1000), None, 500, 500, )
1092 .await;
1093
1094 test_bandwidth_between_peers(
1097 &mut context,
1098 &oracle,
1099 5,
1100 None, None, 500, 0, )
1105 .await;
1106 });
1107 }
1108
1109 #[test]
1110 fn test_bandwidth_contention() {
1111 let executor = deterministic::Runner::default();
1113 executor.start(|context| async move {
1114 let (network, oracle) = Network::new(
1115 context.child("network"),
1116 Config {
1117 max_size: 1024 * 1024,
1118 max_peers_per_set: NZUsize!(101),
1119 disconnect_on_block: true,
1120 tracked_peer_sets: NZUsize!(1),
1121 },
1122 );
1123 network.start();
1124
1125 const NUM_PEERS: usize = 100;
1127 const MESSAGE_SIZE: usize = 1000; const EFFECTIVE_BPS: usize = 10_000; let mut peers = Vec::with_capacity(NUM_PEERS + 1);
1132 let mut senders = Vec::with_capacity(NUM_PEERS + 1);
1133 let mut receivers = Vec::with_capacity(NUM_PEERS + 1);
1134
1135 for i in 0..=NUM_PEERS {
1137 let pk = PrivateKey::from_seed(i as u64).public_key();
1138 let (sender, receiver) = oracle
1139 .control(pk.clone())
1140 .register(0, TEST_QUOTA)
1141 .await
1142 .unwrap();
1143 peers.push(pk);
1144 senders.push(sender);
1145 receivers.push(receiver);
1146 }
1147 track_peers(&oracle, peers.iter().cloned()).await;
1148
1149 for pk in &peers {
1151 oracle
1152 .limit_bandwidth(pk.clone(), Some(EFFECTIVE_BPS), Some(EFFECTIVE_BPS))
1153 .await
1154 .unwrap();
1155 }
1156
1157 for peer in peers.iter().skip(1) {
1159 oracle
1160 .add_link(
1161 peer.clone(),
1162 peers[0].clone(),
1163 Link {
1164 latency: Duration::ZERO,
1165 jitter: Duration::ZERO,
1166 success_rate: probability!(1.0),
1167 },
1168 )
1169 .await
1170 .unwrap();
1171 oracle
1172 .add_link(
1173 peers[0].clone(),
1174 peer.clone(),
1175 Link {
1176 latency: Duration::ZERO,
1177 jitter: Duration::ZERO,
1178 success_rate: probability!(1.0),
1179 },
1180 )
1181 .await
1182 .unwrap();
1183 }
1184
1185 let start = context.current();
1188
1189 let msg = IoBuf::from(vec![0u8; MESSAGE_SIZE]);
1192 for peer in peers.iter().skip(1) {
1193 senders[0].send(Recipients::One(peer.clone()), msg.clone(), true);
1194 }
1195
1196 for receiver in receivers.iter_mut().skip(1) {
1198 let (origin, received) = receiver.recv().await.unwrap();
1199 assert_eq!(origin, peers[0]);
1200 assert_eq!(received.len(), MESSAGE_SIZE);
1201 }
1202
1203 let elapsed = context.current().duration_since(start).unwrap();
1204
1205 let expected_ms = (NUM_PEERS * MESSAGE_SIZE * 1000) / EFFECTIVE_BPS;
1207
1208 assert!(
1209 elapsed >= Duration::from_millis(expected_ms as u64),
1210 "One-to-many completed too quickly: {elapsed:?} (expected >= {expected_ms}ms)"
1211 );
1212 assert!(
1213 elapsed < Duration::from_millis((expected_ms as u64) + 500),
1214 "One-to-many took too long: {elapsed:?} (expected ~{expected_ms}ms)"
1215 );
1216
1217 let start = context.current();
1219
1220 let msg = IoBuf::from(vec![0; MESSAGE_SIZE]);
1223 for mut sender in senders.into_iter().skip(1) {
1224 sender.send(Recipients::One(peers[0].clone()), msg.clone(), true);
1225 }
1226
1227 let mut received_from = HashSet::new();
1229 for _ in 1..=NUM_PEERS {
1230 let (origin, received) = receivers[0].recv().await.unwrap();
1231 assert_eq!(received.len(), MESSAGE_SIZE);
1232 assert!(
1233 received_from.insert(origin.clone()),
1234 "Received duplicate from {origin:?}"
1235 );
1236 }
1237
1238 let elapsed = context.current().duration_since(start).unwrap();
1239
1240 let expected_ms = (NUM_PEERS * MESSAGE_SIZE * 1000) / EFFECTIVE_BPS;
1242
1243 assert!(
1244 elapsed >= Duration::from_millis(expected_ms as u64),
1245 "Many-to-one completed too quickly: {elapsed:?} (expected >= {expected_ms}ms)"
1246 );
1247 assert!(
1248 elapsed < Duration::from_millis((expected_ms as u64) + 500),
1249 "Many-to-one took too long: {elapsed:?} (expected ~{expected_ms}ms)"
1250 );
1251
1252 assert_eq!(received_from.len(), NUM_PEERS);
1254 for peer in peers.iter().skip(1) {
1255 assert!(received_from.contains(peer));
1256 }
1257 });
1258 }
1259
1260 #[test]
1261 fn test_message_ordering() {
1262 let executor = deterministic::Runner::default();
1264 executor.start(|context| async move {
1265 let (network, oracle) = Network::new(
1266 context.child("network"),
1267 Config {
1268 max_size: 1024 * 1024,
1269 max_peers_per_set: NZUsize!(2),
1270 disconnect_on_block: true,
1271 tracked_peer_sets: NZUsize!(1),
1272 },
1273 );
1274 network.start();
1275
1276 let pk1 = PrivateKey::from_seed(1).public_key();
1278 let pk2 = PrivateKey::from_seed(2).public_key();
1279 let (mut sender, _) = oracle
1280 .control(pk1.clone())
1281 .register(0, TEST_QUOTA)
1282 .await
1283 .unwrap();
1284 let (_, mut receiver) = oracle
1285 .control(pk2.clone())
1286 .register(0, TEST_QUOTA)
1287 .await
1288 .unwrap();
1289 track_peers(&oracle, [pk1.clone(), pk2.clone()]).await;
1290
1291 oracle
1293 .add_link(
1294 pk1.clone(),
1295 pk2.clone(),
1296 Link {
1297 latency: Duration::from_millis(50),
1298 jitter: Duration::from_millis(40),
1299 success_rate: probability!(1.0),
1300 },
1301 )
1302 .await
1303 .unwrap();
1304
1305 let messages = vec![
1307 IoBuf::from(b"message 1"),
1308 IoBuf::from(b"message 2"),
1309 IoBuf::from(b"message 3"),
1310 IoBuf::from(b"message 4"),
1311 IoBuf::from(b"message 5"),
1312 ];
1313
1314 for msg in messages.clone() {
1315 sender.send(Recipients::One(pk2.clone()), msg, true);
1316 }
1317
1318 for expected_msg in messages {
1320 let (origin, received_msg) = receiver.recv().await.unwrap();
1321 assert_eq!(origin, pk1);
1322 assert_eq!(received_msg, expected_msg);
1323 }
1324 })
1325 }
1326
1327 #[test]
1328 fn test_high_latency_message_blocks_followup() {
1329 let executor = deterministic::Runner::default();
1330 executor.start(|context| async move {
1331 let (network, oracle) = Network::new(
1332 context.child("network"),
1333 Config {
1334 max_size: 1024 * 1024,
1335 max_peers_per_set: NZUsize!(2),
1336 disconnect_on_block: true,
1337 tracked_peer_sets: NZUsize!(1),
1338 },
1339 );
1340 network.start();
1341
1342 let pk1 = PrivateKey::from_seed(1).public_key();
1343 let pk2 = PrivateKey::from_seed(2).public_key();
1344 let (mut sender, _) = oracle.control(pk1.clone()).register(0, TEST_QUOTA).await.unwrap();
1345 let (_, mut receiver) = oracle.control(pk2.clone()).register(0, TEST_QUOTA).await.unwrap();
1346 track_peers(&oracle, [pk1.clone(), pk2.clone()]).await;
1347
1348 const BPS: usize = 1_000;
1349 oracle
1350 .limit_bandwidth(pk1.clone(), Some(BPS), None)
1351 .await
1352 .unwrap();
1353 oracle
1354 .limit_bandwidth(pk2.clone(), None, Some(BPS))
1355 .await
1356 .unwrap();
1357
1358 oracle
1360 .add_link(
1361 pk1.clone(),
1362 pk2.clone(),
1363 Link {
1364 latency: Duration::from_millis(5_000),
1365 jitter: Duration::ZERO,
1366 success_rate: probability!(1.0),
1367 },
1368 )
1369 .await
1370 .unwrap();
1371
1372 let egress_time = Duration::from_secs(1);
1373 let start = context.current();
1374 let slow = IoBuf::from(vec![0u8; 1_000]);
1375 sender
1376 .send(Recipients::One(pk2.clone()), slow.clone(), true);
1377
1378 oracle.remove_link(pk1.clone(), pk2.clone()).await.unwrap();
1380 oracle
1381 .add_link(
1382 pk1.clone(),
1383 pk2.clone(),
1384 Link {
1385 latency: Duration::from_millis(1),
1386 jitter: Duration::ZERO,
1387 success_rate: probability!(1.0),
1388 },
1389 )
1390 .await
1391 .unwrap();
1392
1393 let fast = IoBuf::from(vec![1u8; 1_000]);
1395 sender
1396 .send(Recipients::One(pk2.clone()), fast.clone(), true);
1397
1398 let (origin1, message1) = receiver.recv().await.unwrap();
1399 assert_eq!(origin1, pk1);
1400 assert_eq!(message1, slow);
1401 let first_elapsed = context.current().duration_since(start).unwrap();
1402
1403 let (origin2, message2) = receiver.recv().await.unwrap();
1404 let second_elapsed = context.current().duration_since(start).unwrap();
1405 assert_eq!(origin2, pk1);
1406 assert_eq!(message2, fast);
1407
1408 let slow_latency = Duration::from_millis(5_000);
1409 let expected_first = egress_time + slow_latency;
1410 let tolerance = Duration::from_millis(10);
1411 assert!(
1412 first_elapsed >= expected_first.saturating_sub(tolerance)
1413 && first_elapsed <= expected_first + tolerance,
1414 "slow message arrived outside expected window: {first_elapsed:?} (expected {expected_first:?} ± {tolerance:?})"
1415 );
1416 assert!(
1417 second_elapsed >= first_elapsed,
1418 "fast message arrived before slow transmission completed"
1419 );
1420
1421 let arrival_gap = second_elapsed
1422 .checked_sub(first_elapsed)
1423 .expect("timestamps ordered");
1424 assert!(
1425 arrival_gap >= egress_time.saturating_sub(tolerance)
1426 && arrival_gap <= egress_time + tolerance,
1427 "next arrival deviated from transmit duration (gap = {arrival_gap:?}, expected {egress_time:?} ± {tolerance:?})"
1428 );
1429 })
1430 }
1431
1432 #[test]
1433 fn test_many_to_one_bandwidth_sharing() {
1434 let executor = deterministic::Runner::default();
1435 executor.start(|context| async move {
1436 let (network, oracle) = Network::new(
1437 context.child("network"),
1438 Config {
1439 max_size: 1024 * 1024,
1440 max_peers_per_set: NZUsize!(11),
1441 disconnect_on_block: true,
1442 tracked_peer_sets: NZUsize!(1),
1443 },
1444 );
1445 network.start();
1446
1447 let mut senders = Vec::new();
1449 let mut sender_txs = Vec::new();
1450 for i in 0..10 {
1451 let sender = ed25519::PrivateKey::from_seed(i).public_key();
1452 senders.push(sender.clone());
1453 let (tx, _) = oracle
1454 .control(sender.clone())
1455 .register(0, TEST_QUOTA)
1456 .await
1457 .unwrap();
1458 sender_txs.push(tx);
1459
1460 oracle
1462 .limit_bandwidth(sender.clone(), Some(10_000), None)
1463 .await
1464 .unwrap();
1465 }
1466
1467 let receiver = ed25519::PrivateKey::from_seed(100).public_key();
1468 let (_, mut receiver_rx) = oracle
1469 .control(receiver.clone())
1470 .register(0, TEST_QUOTA)
1471 .await
1472 .unwrap();
1473 track_peers(&oracle, senders.iter().cloned().chain([receiver.clone()])).await;
1474
1475 oracle
1477 .limit_bandwidth(receiver.clone(), None, Some(100_000))
1478 .await
1479 .unwrap();
1480
1481 for sender in &senders {
1483 oracle
1484 .add_link(
1485 sender.clone(),
1486 receiver.clone(),
1487 Link {
1488 latency: Duration::ZERO,
1489 jitter: Duration::ZERO,
1490 success_rate: probability!(1.0),
1491 },
1492 )
1493 .await
1494 .unwrap();
1495 }
1496
1497 let start = context.current();
1498
1499 for (i, mut tx) in sender_txs.into_iter().enumerate() {
1501 let receiver_clone = receiver.clone();
1502 let msg = IoBuf::from(vec![i as u8; 10_000]);
1503 tx.send(Recipients::One(receiver_clone), msg, true);
1504 }
1505
1506 for i in 0..10 {
1509 let (_, _msg) = receiver_rx.recv().await.unwrap();
1510 let recv_time = context.current().duration_since(start).unwrap();
1511
1512 assert!(
1514 recv_time >= Duration::from_millis(950)
1515 && recv_time <= Duration::from_millis(1100),
1516 "Message {i} received at {recv_time:?}, expected ~1s",
1517 );
1518 }
1519 });
1520 }
1521
1522 #[test]
1523 fn test_one_to_many_fast_sender() {
1524 let executor = deterministic::Runner::default();
1527 executor.start(|context| async move {
1528 let (network, oracle) = Network::new(
1529 context.child("network"),
1530 Config {
1531 max_size: 1024 * 1024,
1532 max_peers_per_set: NZUsize!(11),
1533 disconnect_on_block: true,
1534 tracked_peer_sets: NZUsize!(1),
1535 },
1536 );
1537 network.start();
1538
1539 let sender = ed25519::PrivateKey::from_seed(0).public_key();
1541 let (mut sender_tx, _) = oracle
1542 .control(sender.clone())
1543 .register(0, TEST_QUOTA)
1544 .await
1545 .unwrap();
1546
1547 oracle
1549 .limit_bandwidth(sender.clone(), Some(100_000), None)
1550 .await
1551 .unwrap();
1552
1553 let mut receivers = Vec::new();
1555 let mut receiver_rxs = Vec::new();
1556 for i in 0..10 {
1557 let receiver = ed25519::PrivateKey::from_seed(i + 1).public_key();
1558 receivers.push(receiver.clone());
1559 let (_, rx) = oracle
1560 .control(receiver.clone())
1561 .register(0, TEST_QUOTA)
1562 .await
1563 .unwrap();
1564 receiver_rxs.push(rx);
1565
1566 oracle
1568 .limit_bandwidth(receiver.clone(), None, Some(10_000))
1569 .await
1570 .unwrap();
1571
1572 oracle
1574 .add_link(
1575 sender.clone(),
1576 receiver.clone(),
1577 Link {
1578 latency: Duration::ZERO,
1579 jitter: Duration::ZERO,
1580 success_rate: probability!(1.0),
1581 },
1582 )
1583 .await
1584 .unwrap();
1585 }
1586 track_peers(
1587 &oracle,
1588 core::iter::once(sender.clone()).chain(receivers.iter().cloned()),
1589 )
1590 .await;
1591
1592 let start = context.current();
1593
1594 for (i, receiver) in receivers.iter().enumerate() {
1596 let msg = IoBuf::from(vec![i as u8; 10_000]);
1597 sender_tx.send(Recipients::One(receiver.clone()), msg, true);
1598 }
1599
1600 for (i, mut rx) in receiver_rxs.into_iter().enumerate() {
1602 let (_, msg) = rx.recv().await.unwrap();
1603 assert_eq!(msg.as_ref()[0], i as u8);
1604 let recv_time = context.current().duration_since(start).unwrap();
1605
1606 assert!(
1608 recv_time >= Duration::from_millis(950)
1609 && recv_time <= Duration::from_millis(1100),
1610 "Receiver {i} received at {recv_time:?}, expected ~1s",
1611 );
1612 }
1613 });
1614 }
1615
1616 #[test]
1617 fn test_many_slow_senders_to_fast_receiver() {
1618 let executor = deterministic::Runner::default();
1621 executor.start(|context| async move {
1622 let (network, oracle) = Network::new(
1623 context.child("network"),
1624 Config {
1625 max_size: 1024 * 1024,
1626 max_peers_per_set: NZUsize!(11),
1627 disconnect_on_block: true,
1628 tracked_peer_sets: NZUsize!(1),
1629 },
1630 );
1631 network.start();
1632
1633 let mut senders = Vec::new();
1635 let mut sender_txs = Vec::new();
1636 for i in 0..10 {
1637 let sender = ed25519::PrivateKey::from_seed(i).public_key();
1638 senders.push(sender.clone());
1639 let (tx, _) = oracle
1640 .control(sender.clone())
1641 .register(0, TEST_QUOTA)
1642 .await
1643 .unwrap();
1644 sender_txs.push(tx);
1645
1646 oracle
1648 .limit_bandwidth(sender.clone(), Some(1_000), None)
1649 .await
1650 .unwrap();
1651 }
1652
1653 let receiver = ed25519::PrivateKey::from_seed(100).public_key();
1655 let (_, mut receiver_rx) = oracle
1656 .control(receiver.clone())
1657 .register(0, TEST_QUOTA)
1658 .await
1659 .unwrap();
1660 track_peers(&oracle, senders.iter().cloned().chain([receiver.clone()])).await;
1661
1662 oracle
1664 .limit_bandwidth(receiver.clone(), None, Some(10_000))
1665 .await
1666 .unwrap();
1667
1668 for sender in &senders {
1670 oracle
1671 .add_link(
1672 sender.clone(),
1673 receiver.clone(),
1674 Link {
1675 latency: Duration::ZERO,
1676 jitter: Duration::ZERO,
1677 success_rate: probability!(1.0),
1678 },
1679 )
1680 .await
1681 .unwrap();
1682 }
1683
1684 let start = context.current();
1685
1686 for (i, mut tx) in sender_txs.into_iter().enumerate() {
1688 let receiver_clone = receiver.clone();
1689 let msg = IoBuf::from(vec![i as u8; 1_000]);
1690 tx.send(Recipients::One(receiver_clone), msg, true);
1691 }
1692
1693 for i in 0..10 {
1699 let (_, _msg) = receiver_rx.recv().await.unwrap();
1700 let recv_time = context.current().duration_since(start).unwrap();
1701
1702 assert!(
1704 recv_time >= Duration::from_millis(950)
1705 && recv_time <= Duration::from_millis(1100),
1706 "Message {i} received at {recv_time:?}, expected ~1s",
1707 );
1708 }
1709 });
1710 }
1711
1712 #[test]
1713 fn test_dynamic_bandwidth_allocation_staggered() {
1714 let executor = deterministic::Runner::default();
1720 executor.start(|context| async move {
1721 let (network, oracle) = Network::new(
1722 context.child("network"),
1723 Config {
1724 max_size: 1024 * 1024,
1725 max_peers_per_set: NZUsize!(4),
1726 disconnect_on_block: true,
1727 tracked_peer_sets: NZUsize!(1),
1728 },
1729 );
1730 network.start();
1731
1732 let mut senders = Vec::new();
1734 let mut sender_txs = Vec::new();
1735 for i in 0..3 {
1736 let sender = ed25519::PrivateKey::from_seed(i).public_key();
1737 senders.push(sender.clone());
1738 let (tx, _) = oracle
1739 .control(sender.clone())
1740 .register(0, TEST_QUOTA)
1741 .await
1742 .unwrap();
1743 sender_txs.push(tx);
1744
1745 oracle
1747 .limit_bandwidth(sender.clone(), Some(30_000), None)
1748 .await
1749 .unwrap();
1750 }
1751
1752 let receiver = ed25519::PrivateKey::from_seed(100).public_key();
1754 let (_, mut receiver_rx) = oracle
1755 .control(receiver.clone())
1756 .register(0, TEST_QUOTA)
1757 .await
1758 .unwrap();
1759 track_peers(&oracle, senders.iter().cloned().chain([receiver.clone()])).await;
1760 oracle
1761 .limit_bandwidth(receiver.clone(), None, Some(30_000))
1762 .await
1763 .unwrap();
1764
1765 for sender in &senders {
1767 oracle
1768 .add_link(
1769 sender.clone(),
1770 receiver.clone(),
1771 Link {
1772 latency: Duration::from_millis(1),
1773 jitter: Duration::ZERO,
1774 success_rate: probability!(1.0),
1775 },
1776 )
1777 .await
1778 .unwrap();
1779 }
1780
1781 let start = context.current();
1782 let mut sender_txs = sender_txs.into_iter();
1783
1784 let mut tx0 = sender_txs.next().expect("missing sender 0");
1788 let rx_clone = receiver.clone();
1789 context.child("task").spawn(move |_| async move {
1790 let msg = IoBuf::from(vec![0u8; 30_000]);
1791 tx0.send(Recipients::One(rx_clone), msg, true);
1792 });
1793
1794 let mut tx1 = sender_txs.next().expect("missing sender 1");
1798 let rx_clone = receiver.clone();
1799 context.child("task").spawn(move |context| async move {
1800 context.sleep(Duration::from_millis(500)).await;
1801 let msg = IoBuf::from(vec![1u8; 30_000]);
1802 tx1.send(Recipients::One(rx_clone), msg, true);
1803 });
1804
1805 let mut tx2 = sender_txs.next().expect("missing sender 2");
1808 let rx_clone = receiver.clone();
1809 context.child("task").spawn(move |context| async move {
1810 context.sleep(Duration::from_millis(1500)).await;
1811 let msg = IoBuf::from(vec![2u8; 15_000]);
1812 tx2.send(Recipients::One(rx_clone), msg, true);
1813 });
1814
1815 let (_, msg0) = receiver_rx.recv().await.unwrap();
1819 assert_eq!(msg0.as_ref()[0], 0);
1820 let t0 = context.current().duration_since(start).unwrap();
1821 assert!(
1822 t0 >= Duration::from_millis(1490) && t0 <= Duration::from_millis(1600),
1823 "Message 0 received at {t0:?}, expected ~1.5s",
1824 );
1825
1826 let (_, msg_a) = receiver_rx.recv().await.unwrap();
1830 let t_a = context.current().duration_since(start).unwrap();
1831
1832 let (_, msg_b) = receiver_rx.recv().await.unwrap();
1833 let t_b = context.current().duration_since(start).unwrap();
1834
1835 let (msg1, t1, msg2, t2) = if msg_a.as_ref()[0] == 1 {
1837 (msg_a, t_a, msg_b, t_b)
1838 } else {
1839 (msg_b, t_b, msg_a, t_a)
1840 };
1841
1842 assert_eq!(msg1.as_ref()[0], 1);
1843 assert_eq!(msg2.as_ref()[0], 2);
1844
1845 assert!(
1850 t1 >= Duration::from_millis(1500) && t1 <= Duration::from_millis(2600),
1851 "Message 1 received at {t1:?}, expected between 1.5s-2.6s",
1852 );
1853
1854 assert!(
1855 t2 >= Duration::from_millis(1500) && t2 <= Duration::from_millis(2600),
1856 "Message 2 received at {t2:?}, expected between 1.5s-2.6s",
1857 );
1858 });
1859 }
1860
1861 #[test]
1862 fn test_dynamic_bandwidth_varied_sizes() {
1863 let executor = deterministic::Runner::default();
1866 executor.start(|context| async move {
1867 let (network, oracle) = Network::new(
1868 context.child("network"),
1869 Config {
1870 max_size: 1024 * 1024,
1871 max_peers_per_set: NZUsize!(4),
1872 disconnect_on_block: true,
1873 tracked_peer_sets: NZUsize!(1),
1874 },
1875 );
1876 network.start();
1877
1878 let mut senders = Vec::new();
1880 let mut sender_txs = Vec::new();
1881 for i in 0..3 {
1882 let sender = ed25519::PrivateKey::from_seed(i).public_key();
1883 senders.push(sender.clone());
1884 let (tx, _) = oracle
1885 .control(sender.clone())
1886 .register(0, TEST_QUOTA)
1887 .await
1888 .unwrap();
1889 sender_txs.push(tx);
1890
1891 oracle
1893 .limit_bandwidth(sender.clone(), None, None)
1894 .await
1895 .unwrap();
1896 }
1897
1898 let receiver = ed25519::PrivateKey::from_seed(100).public_key();
1900 let (_, mut receiver_rx) = oracle
1901 .control(receiver.clone())
1902 .register(0, TEST_QUOTA)
1903 .await
1904 .unwrap();
1905 track_peers(&oracle, senders.iter().cloned().chain([receiver.clone()])).await;
1906 oracle
1907 .limit_bandwidth(receiver.clone(), None, Some(30_000))
1908 .await
1909 .unwrap();
1910
1911 for sender in &senders {
1913 oracle
1914 .add_link(
1915 sender.clone(),
1916 receiver.clone(),
1917 Link {
1918 latency: Duration::from_millis(1),
1919 jitter: Duration::ZERO,
1920 success_rate: probability!(1.0),
1921 },
1922 )
1923 .await
1924 .unwrap();
1925 }
1926
1927 let start = context.current();
1928
1929 let sizes = [10_000, 20_000, 30_000];
1935 for (i, (mut tx, size)) in sender_txs.into_iter().zip(sizes.iter()).enumerate() {
1936 let rx_clone = receiver.clone();
1937 let msg_size = *size;
1938 let msg = IoBuf::from(vec![i as u8; msg_size]);
1939 tx.send(Recipients::One(rx_clone), msg, true);
1940 }
1941
1942 let mut messages = Vec::new();
1946 for _ in 0..3 {
1947 let (_, msg) = receiver_rx.recv().await.unwrap();
1948 let t = context.current().duration_since(start).unwrap();
1949 messages.push((msg.as_ref()[0] as usize, msg.len(), t));
1950 }
1951
1952 assert_eq!(messages.len(), 3);
1957
1958 let max_time = messages.iter().map(|&(_, _, t)| t).max().unwrap();
1960 assert!(
1961 max_time >= Duration::from_millis(2000),
1962 "Total time {max_time:?} should be at least 2s for 60KB at 30KB/s",
1963 );
1964 });
1965 }
1966
1967 #[test]
1968 fn test_bandwidth_pipe_reservation_duration() {
1969 let executor = deterministic::Runner::default();
1972 executor.start(|context| async move {
1973 let (network, oracle) = Network::new(
1974 context.child("network"),
1975 Config {
1976 max_size: 1024 * 1024,
1977 max_peers_per_set: NZUsize!(2),
1978 disconnect_on_block: true,
1979 tracked_peer_sets: NZUsize!(1),
1980 },
1981 );
1982 network.start();
1983
1984 let sender = PrivateKey::from_seed(1).public_key();
1986 let receiver = PrivateKey::from_seed(2).public_key();
1987
1988 let (mut sender_tx, _) = oracle
1989 .control(sender.clone())
1990 .register(0, TEST_QUOTA)
1991 .await
1992 .unwrap();
1993 let (_, mut receiver_rx) = oracle
1994 .control(receiver.clone())
1995 .register(0, TEST_QUOTA)
1996 .await
1997 .unwrap();
1998 track_peers(&oracle, [sender.clone(), receiver.clone()]).await;
1999
2000 oracle
2002 .limit_bandwidth(sender.clone(), Some(1000), None)
2003 .await
2004 .unwrap();
2005 oracle
2006 .limit_bandwidth(receiver.clone(), None, Some(1000))
2007 .await
2008 .unwrap();
2009
2010 oracle
2012 .add_link(
2013 sender.clone(),
2014 receiver.clone(),
2015 Link {
2016 latency: Duration::from_secs(1), jitter: Duration::ZERO,
2018 success_rate: probability!(1.0),
2019 },
2020 )
2021 .await
2022 .unwrap();
2023
2024 let start = context.current();
2035
2036 for i in 0..3 {
2038 let msg = IoBuf::from(vec![i; 500]);
2039 sender_tx.send(Recipients::One(receiver.clone()), msg, false);
2040 }
2041
2042 let mut receive_times = Vec::new();
2044 for i in 0..3 {
2045 let (_, received) = receiver_rx.recv().await.unwrap();
2046 receive_times.push(context.current().duration_since(start).unwrap());
2047 assert_eq!(received.as_ref()[0], i);
2048 }
2049
2050 for (i, time) in receive_times.iter().enumerate() {
2055 let expected_min = (i as u64 * 500) + 1500;
2056 let expected_max = expected_min + 100;
2057
2058 assert!(
2059 *time >= Duration::from_millis(expected_min)
2060 && *time < Duration::from_millis(expected_max),
2061 "Message {} should arrive at ~{}ms, got {:?}",
2062 i + 1,
2063 expected_min,
2064 time
2065 );
2066 }
2067 });
2068 }
2069
2070 #[test]
2071 fn test_dynamic_bandwidth_affects_new_transfers() {
2072 let executor = deterministic::Runner::default();
2075 executor.start(|context| async move {
2076 let (network, oracle) = Network::new(
2077 context.child("network"),
2078 Config {
2079 max_size: 1024 * 1024,
2080 max_peers_per_set: NZUsize!(2),
2081 disconnect_on_block: true,
2082 tracked_peer_sets: NZUsize!(1),
2083 },
2084 );
2085 network.start();
2086
2087 let pk_sender = PrivateKey::from_seed(1).public_key();
2088 let pk_receiver = PrivateKey::from_seed(2).public_key();
2089
2090 let (mut sender_tx, _) = oracle
2092 .control(pk_sender.clone())
2093 .register(0, TEST_QUOTA)
2094 .await
2095 .unwrap();
2096 let (_, mut receiver_rx) = oracle
2097 .control(pk_receiver.clone())
2098 .register(0, TEST_QUOTA)
2099 .await
2100 .unwrap();
2101 track_peers(&oracle, [pk_sender.clone(), pk_receiver.clone()]).await;
2102 oracle
2103 .add_link(
2104 pk_sender.clone(),
2105 pk_receiver.clone(),
2106 Link {
2107 latency: Duration::from_millis(1), jitter: Duration::ZERO,
2109 success_rate: probability!(1.0),
2110 },
2111 )
2112 .await
2113 .unwrap();
2114
2115 oracle
2117 .limit_bandwidth(pk_sender.clone(), Some(10_000), None)
2118 .await
2119 .unwrap();
2120 oracle
2121 .limit_bandwidth(pk_receiver.clone(), None, Some(10_000))
2122 .await
2123 .unwrap();
2124
2125 let msg1 = IoBuf::from(vec![1u8; 20_000]); let start_time = context.current();
2128 sender_tx.send(Recipients::One(pk_receiver.clone()), msg1.clone(), false);
2129
2130 let (_sender, received_msg1) = receiver_rx.recv().await.unwrap();
2132 let msg1_time = context.current().duration_since(start_time).unwrap();
2133 assert_eq!(received_msg1.len(), 20_000);
2134 assert!(
2135 msg1_time >= Duration::from_millis(1999)
2136 && msg1_time <= Duration::from_millis(2010),
2137 "First message should take ~2s, got {msg1_time:?}",
2138 );
2139
2140 oracle
2142 .limit_bandwidth(pk_sender.clone(), Some(2_000), None)
2143 .await
2144 .unwrap();
2145
2146 let msg2 = IoBuf::from(vec![2u8; 10_000]); let msg2_start = context.current();
2149 sender_tx.send(Recipients::One(pk_receiver.clone()), msg2.clone(), false);
2150
2151 let (_sender, received_msg2) = receiver_rx.recv().await.unwrap();
2153 let msg2_time = context.current().duration_since(msg2_start).unwrap();
2154 assert_eq!(received_msg2.len(), 10_000);
2155 assert!(
2156 msg2_time >= Duration::from_millis(4999)
2157 && msg2_time <= Duration::from_millis(5010),
2158 "Second message should take ~5s at reduced bandwidth, got {msg2_time:?}",
2159 );
2160 });
2161 }
2162
2163 #[test]
2164 fn test_zero_receiver_ingress_bandwidth() {
2165 let executor = deterministic::Runner::default();
2166 executor.start(|context| async move {
2167 let (network, oracle) = Network::new(
2168 context.child("network"),
2169 Config {
2170 max_size: 1024 * 1024,
2171 max_peers_per_set: NZUsize!(2),
2172 disconnect_on_block: true,
2173 tracked_peer_sets: NZUsize!(1),
2174 },
2175 );
2176 network.start();
2177
2178 let pk_sender = PrivateKey::from_seed(1).public_key();
2179 let pk_receiver = PrivateKey::from_seed(2).public_key();
2180
2181 let (mut sender_tx, _) = oracle
2183 .control(pk_sender.clone())
2184 .register(0, TEST_QUOTA)
2185 .await
2186 .unwrap();
2187 let (_, mut receiver_rx) = oracle
2188 .control(pk_receiver.clone())
2189 .register(0, TEST_QUOTA)
2190 .await
2191 .unwrap();
2192 track_peers(&oracle, [pk_sender.clone(), pk_receiver.clone()]).await;
2193 oracle
2194 .add_link(
2195 pk_sender.clone(),
2196 pk_receiver.clone(),
2197 Link {
2198 latency: Duration::ZERO,
2199 jitter: Duration::ZERO,
2200 success_rate: probability!(1.0),
2201 },
2202 )
2203 .await
2204 .unwrap();
2205
2206 oracle
2208 .limit_bandwidth(pk_receiver.clone(), None, Some(0))
2209 .await
2210 .unwrap();
2211
2212 let msg1 = IoBuf::from(vec![1u8; 20_000]); let sent = sender_tx.send(Recipients::One(pk_receiver.clone()), msg1.clone(), false);
2215 assert_eq!(sent.len(), 1);
2216 assert_eq!(sent[0], pk_receiver);
2217
2218 select! {
2220 _ = receiver_rx.recv() => {
2221 panic!("unexpected message");
2222 },
2223 _ = context.sleep(Duration::from_secs(10)) => {},
2224 }
2225
2226 oracle
2228 .limit_bandwidth(pk_receiver.clone(), None, None)
2229 .await
2230 .unwrap();
2231
2232 select! {
2234 _ = receiver_rx.recv() => {},
2235 _ = context.sleep(Duration::from_secs(1)) => {
2236 panic!("timeout");
2237 },
2238 }
2239 });
2240 }
2241
2242 #[test]
2243 fn test_zero_sender_egress_bandwidth() {
2244 let executor = deterministic::Runner::default();
2245 executor.start(|context| async move {
2246 let (network, oracle) = Network::new(
2247 context.child("network"),
2248 Config {
2249 max_size: 1024 * 1024,
2250 max_peers_per_set: NZUsize!(2),
2251 disconnect_on_block: true,
2252 tracked_peer_sets: NZUsize!(1),
2253 },
2254 );
2255 network.start();
2256
2257 let pk_sender = PrivateKey::from_seed(1).public_key();
2258 let pk_receiver = PrivateKey::from_seed(2).public_key();
2259
2260 let (mut sender_tx, _) = oracle
2262 .control(pk_sender.clone())
2263 .register(0, TEST_QUOTA)
2264 .await
2265 .unwrap();
2266 let (_, mut receiver_rx) = oracle
2267 .control(pk_receiver.clone())
2268 .register(0, TEST_QUOTA)
2269 .await
2270 .unwrap();
2271 track_peers(&oracle, [pk_sender.clone(), pk_receiver.clone()]).await;
2272 oracle
2273 .add_link(
2274 pk_sender.clone(),
2275 pk_receiver.clone(),
2276 Link {
2277 latency: Duration::ZERO,
2278 jitter: Duration::ZERO,
2279 success_rate: probability!(1.0),
2280 },
2281 )
2282 .await
2283 .unwrap();
2284
2285 oracle
2287 .limit_bandwidth(pk_sender.clone(), Some(0), None)
2288 .await
2289 .unwrap();
2290
2291 let msg1 = IoBuf::from(vec![1u8; 20_000]); let sent = sender_tx.send(Recipients::One(pk_receiver.clone()), msg1.clone(), false);
2294 assert_eq!(sent.len(), 1);
2295 assert_eq!(sent[0], pk_receiver);
2296
2297 select! {
2299 _ = receiver_rx.recv() => {
2300 panic!("unexpected message");
2301 },
2302 _ = context.sleep(Duration::from_secs(10)) => {},
2303 }
2304
2305 oracle
2307 .limit_bandwidth(pk_sender.clone(), None, None)
2308 .await
2309 .unwrap();
2310
2311 select! {
2313 _ = receiver_rx.recv() => {},
2314 _ = context.sleep(Duration::from_secs(1)) => {
2315 panic!("timeout");
2316 },
2317 }
2318 });
2319 }
2320
2321 #[test]
2322 fn register_peer_set() {
2323 let executor = deterministic::Runner::default();
2324 executor.start(|context| async move {
2325 let (network, oracle) = Network::new(
2326 context.child("network"),
2327 Config {
2328 max_size: 1024 * 1024,
2329 max_peers_per_set: NZUsize!(2),
2330 disconnect_on_block: true,
2331 tracked_peer_sets: NZUsize!(3),
2332 },
2333 );
2334 network.start();
2335
2336 let mut manager = oracle.manager();
2337 assert_eq!(manager.peer_set(0).await, None);
2338
2339 let pk1 = PrivateKey::from_seed(1).public_key();
2340 let pk2 = PrivateKey::from_seed(2).public_key();
2341 manager.track(0xFF, Set::try_from([pk1.clone(), pk2.clone()]).unwrap());
2342
2343 assert_eq!(
2344 manager.peer_set(0xFF).await.unwrap(),
2345 TrackedPeers::primary(Set::try_from([pk1, pk2]).unwrap())
2346 );
2347 });
2348 }
2349
2350 #[test]
2351 fn test_socket_manager() {
2352 let executor = deterministic::Runner::default();
2353 executor.start(|context| async move {
2354 let (network, oracle) = Network::new(
2355 context.child("network"),
2356 Config {
2357 max_size: 1024 * 1024,
2358 max_peers_per_set: NZUsize!(2),
2359 disconnect_on_block: true,
2360 tracked_peer_sets: NZUsize!(3),
2361 },
2362 );
2363 network.start();
2364
2365 let pk1 = PrivateKey::from_seed(1).public_key();
2366 let pk2 = PrivateKey::from_seed(2).public_key();
2367 let addr1: Address = SocketAddr::from(([127, 0, 0, 1], 4000)).into();
2368 let addr2: Address = SocketAddr::from(([127, 0, 0, 1], 4001)).into();
2369
2370 let mut manager = oracle.socket_manager();
2371 manager.track(
2372 1,
2373 Map::<_, Address>::try_from([
2374 (pk1.clone(), addr1.clone()),
2375 (pk2.clone(), addr2.clone()),
2376 ])
2377 .unwrap(),
2378 );
2379
2380 let peer_set = manager.peer_set(1).await.expect("peer set missing");
2381 let keys: Vec<_> = Vec::from(peer_set.primary.clone());
2382 assert_eq!(keys, vec![pk1.clone(), pk2.clone()]);
2383
2384 let mut subscription = manager.subscribe().await;
2385 let update = subscription.recv().await.unwrap();
2386 assert_eq!(update.index, 1);
2387 let latest_keys: Vec<_> = Vec::from(update.latest.primary.clone());
2388 assert_eq!(latest_keys, vec![pk1.clone(), pk2.clone()]);
2389 assert!(update.latest.secondary.is_empty());
2390 let all_primary_keys: Vec<_> = Vec::from(update.all.primary.clone());
2391 assert_eq!(all_primary_keys, vec![pk1.clone(), pk2.clone()]);
2392 assert!(update.all.secondary.is_empty());
2393
2394 manager.track(
2395 2,
2396 Map::<_, Address>::try_from([(pk2.clone(), addr2)]).unwrap(),
2397 );
2398
2399 let update = subscription.recv().await.unwrap();
2400 assert_eq!(update.index, 2);
2401 let latest_keys: Vec<_> = Vec::from(update.latest.primary);
2402 assert_eq!(latest_keys, vec![pk2.clone()]);
2403 assert!(update.latest.secondary.is_empty());
2404 let all_primary_keys: Vec<_> = Vec::from(update.all.primary);
2405 assert_eq!(all_primary_keys, vec![pk1, pk2]);
2406 assert!(update.all.secondary.is_empty());
2407 });
2408 }
2409
2410 #[test]
2411 fn test_manager_track_accepts_tracked_peers() {
2412 let executor = deterministic::Runner::default();
2413 executor.start(|context| async move {
2414 let (network, oracle) = Network::new(
2415 context.child("network"),
2416 Config {
2417 max_size: 1024 * 1024,
2418 max_peers_per_set: NZUsize!(2),
2419 disconnect_on_block: true,
2420 tracked_peer_sets: NZUsize!(3),
2421 },
2422 );
2423 network.start();
2424
2425 let pk1 = PrivateKey::from_seed(1).public_key();
2426 let pk2 = PrivateKey::from_seed(2).public_key();
2427 let mut manager = oracle.manager();
2428
2429 manager.track(
2430 7,
2431 TrackedPeers::new(
2432 Set::try_from([pk1.clone()]).unwrap(),
2433 Set::try_from([pk2]).unwrap(),
2434 ),
2435 );
2436
2437 assert_eq!(
2438 manager.peer_set(7).await.unwrap(),
2439 TrackedPeers::new(
2440 Set::try_from([pk1]).unwrap(),
2441 Set::try_from([PrivateKey::from_seed(2).public_key()]).unwrap(),
2442 )
2443 );
2444 });
2445 }
2446
2447 #[test]
2448 fn test_manager_track_tracked_peers_overlap_primary_wins() {
2449 let executor = deterministic::Runner::default();
2450 executor.start(|context| async move {
2451 let (network, oracle) = Network::new(
2454 context.child("network"),
2455 Config {
2456 max_size: 1024 * 1024,
2457 max_peers_per_set: NZUsize!(3),
2458 disconnect_on_block: true,
2459 tracked_peer_sets: NZUsize!(3),
2460 },
2461 );
2462 network.start();
2463
2464 let pk1 = PrivateKey::from_seed(1).public_key();
2465 let pk2 = PrivateKey::from_seed(2).public_key();
2466 let pk3 = PrivateKey::from_seed(3).public_key();
2467 let mut manager = oracle.manager();
2468
2469 manager.track(
2470 9,
2471 TrackedPeers::new(
2472 Set::try_from([pk1.clone(), pk2.clone()]).unwrap(),
2473 Set::try_from([pk2.clone(), pk3.clone()]).unwrap(),
2474 ),
2475 );
2476
2477 assert_eq!(
2478 manager.peer_set(9).await.unwrap(),
2479 TrackedPeers::new(
2480 Set::try_from([pk1.clone(), pk2.clone()]).unwrap(),
2481 Set::try_from([pk3.clone()]).unwrap(),
2482 )
2483 );
2484
2485 let mut subscription = manager.subscribe().await;
2486 let update = subscription.recv().await.unwrap();
2487 assert_eq!(update.index, 9);
2488 assert!(update.latest.primary.position(&pk2).is_some());
2489 assert!(update.latest.secondary.position(&pk2).is_none());
2490 assert!(update.latest.secondary.position(&pk3).is_some());
2491 assert!(update.all.secondary.position(&pk2).is_none());
2492 assert!(update.all.primary.position(&pk2).is_some());
2493 });
2494 }
2495
2496 #[test]
2497 fn test_socket_manager_track_accepts_addressable_tracked_peers() {
2498 let executor = deterministic::Runner::default();
2499 executor.start(|context| async move {
2500 let (network, oracle) = Network::new(
2501 context.child("network"),
2502 Config {
2503 max_size: 1024 * 1024,
2504 max_peers_per_set: NZUsize!(2),
2505 disconnect_on_block: true,
2506 tracked_peer_sets: NZUsize!(3),
2507 },
2508 );
2509 network.start();
2510
2511 let pk1 = PrivateKey::from_seed(1).public_key();
2512 let pk2 = PrivateKey::from_seed(2).public_key();
2513 let addr1: Address = SocketAddr::from(([127, 0, 0, 1], 4000)).into();
2514 let addr2: Address = SocketAddr::from(([127, 0, 0, 1], 4001)).into();
2515 let mut manager = oracle.socket_manager();
2516
2517 manager.track(
2518 7,
2519 AddressableTrackedPeers::new(
2520 Map::<_, Address>::try_from([(pk1.clone(), addr1)]).unwrap(),
2521 Map::<_, Address>::try_from([(pk2, addr2)]).unwrap(),
2522 ),
2523 );
2524
2525 assert_eq!(
2526 manager.peer_set(7).await.unwrap(),
2527 TrackedPeers::new(
2528 Set::try_from([pk1]).unwrap(),
2529 Set::try_from([PrivateKey::from_seed(2).public_key()]).unwrap(),
2530 )
2531 );
2532 });
2533 }
2534
2535 #[test]
2536 fn test_socket_manager_track_addressable_overlap_primary_wins() {
2537 let executor = deterministic::Runner::default();
2538 executor.start(|context| async move {
2539 let (network, oracle) = Network::new(
2541 context.child("network"),
2542 Config {
2543 max_size: 1024 * 1024,
2544 max_peers_per_set: NZUsize!(1),
2545 disconnect_on_block: true,
2546 tracked_peer_sets: NZUsize!(3),
2547 },
2548 );
2549 network.start();
2550
2551 let pk = PrivateKey::from_seed(1).public_key();
2552 let addr_primary: Address = SocketAddr::from(([127, 0, 0, 1], 4000)).into();
2553 let addr_secondary: Address = SocketAddr::from(([127, 0, 0, 1], 5000)).into();
2554 let mut manager = oracle.socket_manager();
2555 let mut subscription = manager.subscribe().await;
2556
2557 manager.track(
2558 11,
2559 AddressableTrackedPeers::new(
2560 Map::<_, Address>::try_from([(pk.clone(), addr_primary.clone())]).unwrap(),
2561 Map::<_, Address>::try_from([(pk.clone(), addr_secondary)]).unwrap(),
2562 ),
2563 );
2564
2565 let update = subscription.recv().await.unwrap();
2566 assert_eq!(update.index, 11);
2567 assert_eq!(update.latest.primary.len(), 1);
2568 assert!(update.latest.secondary.is_empty());
2569 assert!(update.all.secondary.is_empty());
2570 assert_eq!(update.latest.primary, Set::try_from([pk.clone()]).unwrap());
2571 });
2572 }
2573
2574 #[test]
2575 fn test_socket_manager_with_asymmetric_addresses() {
2576 let executor = deterministic::Runner::default();
2577 executor.start(|context| async move {
2578 let (network, oracle) = Network::new(
2579 context.child("network"),
2580 Config {
2581 max_size: 1024 * 1024,
2582 max_peers_per_set: NZUsize!(2),
2583 disconnect_on_block: true,
2584 tracked_peer_sets: NZUsize!(3),
2585 },
2586 );
2587 network.start();
2588
2589 let pk1 = PrivateKey::from_seed(1).public_key();
2590 let pk2 = PrivateKey::from_seed(2).public_key();
2591
2592 let addr1 = Address::Asymmetric {
2594 ingress: Ingress::Socket(SocketAddr::from(([10, 0, 0, 1], 8080))),
2595 egress: SocketAddr::from(([192, 168, 1, 1], 9090)),
2596 };
2597 let addr2 = Address::Asymmetric {
2598 ingress: Ingress::Dns {
2599 host: hostname!("node2.example.com"),
2600 port: 8080,
2601 },
2602 egress: SocketAddr::from(([192, 168, 1, 2], 9090)),
2603 };
2604
2605 let mut manager = oracle.socket_manager();
2606 manager.track(
2607 1,
2608 Map::<_, Address>::try_from([(pk1.clone(), addr1), (pk2.clone(), addr2)]).unwrap(),
2609 );
2610
2611 let peer_set = manager.peer_set(1).await.expect("peer set missing");
2613 let keys: Vec<_> = Vec::from(peer_set.primary);
2614 assert_eq!(keys, vec![pk1.clone(), pk2.clone()]);
2615
2616 let mut subscription = manager.subscribe().await;
2618 let update = subscription.recv().await.unwrap();
2619 assert_eq!(update.index, 1);
2620 let latest_keys: Vec<_> = Vec::from(update.latest.primary);
2621 assert_eq!(latest_keys, vec![pk1, pk2]);
2622 assert!(update.latest.secondary.is_empty());
2623 });
2624 }
2625
2626 #[test]
2627 fn test_peer_set_window_management() {
2628 let executor = deterministic::Runner::default();
2629 executor.start(|context| async move {
2630 let (network, oracle) = Network::new(
2631 context.child("network"),
2632 Config {
2633 max_size: 1024 * 1024,
2634 max_peers_per_set: NZUsize!(2),
2635 disconnect_on_block: true,
2636 tracked_peer_sets: NZUsize!(2), },
2638 );
2639 network.start();
2640
2641 let pk1 = PrivateKey::from_seed(1).public_key();
2643 let pk2 = PrivateKey::from_seed(2).public_key();
2644 let pk3 = PrivateKey::from_seed(3).public_key();
2645 let pk4 = PrivateKey::from_seed(4).public_key();
2646
2647 let mut manager = oracle.manager();
2649 manager.track(1, Set::try_from(vec![pk1.clone(), pk2.clone()]).unwrap());
2650
2651 let (mut sender1, _receiver1) = oracle
2653 .control(pk1.clone())
2654 .register(0, TEST_QUOTA)
2655 .await
2656 .unwrap();
2657 let (mut sender2, _receiver2) = oracle
2658 .control(pk2.clone())
2659 .register(0, TEST_QUOTA)
2660 .await
2661 .unwrap();
2662 let (mut sender3, _receiver3) = oracle
2663 .control(pk3.clone())
2664 .register(0, TEST_QUOTA)
2665 .await
2666 .unwrap();
2667 let (_mut_sender4, _receiver4) = oracle
2668 .control(pk4.clone())
2669 .register(0, TEST_QUOTA)
2670 .await
2671 .unwrap();
2672
2673 for peer_a in &[pk1.clone(), pk2.clone(), pk3.clone(), pk4.clone()] {
2675 for peer_b in &[pk1.clone(), pk2.clone(), pk3.clone(), pk4.clone()] {
2676 if peer_a != peer_b {
2677 oracle
2678 .add_link(
2679 peer_a.clone(),
2680 peer_b.clone(),
2681 Link {
2682 latency: Duration::from_millis(1),
2683 jitter: Duration::ZERO,
2684 success_rate: probability!(1.0),
2685 },
2686 )
2687 .await
2688 .unwrap();
2689 }
2690 }
2691 }
2692
2693 let recipients = sender1.check(Recipients::All).unwrap().recipients();
2695 assert_eq!(recipients, vec![pk2.clone()]);
2696 assert!(!recipients.contains(&pk3));
2697
2698 manager.track(2, Set::try_from(vec![pk2.clone(), pk3.clone()]).unwrap());
2700 assert!(manager.peer_set(2).await.is_some());
2701
2702 let recipients = sender1.check(Recipients::All).unwrap().recipients();
2704 assert!(recipients.contains(&pk3));
2705
2706 manager.track(3, Set::try_from(vec![pk3.clone(), pk4.clone()]).unwrap());
2708 assert!(manager.peer_set(3).await.is_some());
2709
2710 let recipients = sender2.check(Recipients::All).unwrap().recipients();
2712 assert!(!recipients.contains(&pk1));
2713
2714 assert!(recipients.contains(&pk3));
2716
2717 let recipients = sender3.check(Recipients::All).unwrap().recipients();
2719 assert!(recipients.contains(&pk4));
2720
2721 let peer_set_2 = manager.peer_set(2).await.unwrap();
2723 assert!(peer_set_2.primary.position(&pk2).is_some());
2724 assert!(peer_set_2.primary.position(&pk3).is_some());
2725
2726 let peer_set_3 = manager.peer_set(3).await.unwrap();
2727 assert!(peer_set_3.primary.position(&pk3).is_some());
2728 assert!(peer_set_3.primary.position(&pk4).is_some());
2729
2730 assert!(manager.peer_set(1).await.is_none());
2732 });
2733 }
2734
2735 #[test]
2736 fn test_connected_subscription_updates_after_track() {
2737 let executor = deterministic::Runner::default();
2738 executor.start(|context| async move {
2739 let (network, oracle) = Network::new(
2740 context.child("network"),
2741 Config {
2742 max_size: 1024 * 1024,
2743 max_peers_per_set: NZUsize!(2),
2744 disconnect_on_block: true,
2745 tracked_peer_sets: NZUsize!(2),
2746 },
2747 );
2748 network.start();
2749
2750 let pk1 = PrivateKey::from_seed(1).public_key();
2751 let pk2 = PrivateKey::from_seed(2).public_key();
2752 let (mut sender, _) = oracle
2753 .control(pk1.clone())
2754 .register(0, TEST_QUOTA)
2755 .await
2756 .unwrap();
2757
2758 assert!(
2759 sender
2760 .check(Recipients::All)
2761 .unwrap()
2762 .recipients()
2763 .is_empty()
2764 );
2765
2766 let mut manager = oracle.manager();
2767 manager.track(1, Set::try_from([pk1, pk2.clone()]).unwrap());
2768 assert!(manager.peer_set(1).await.is_some());
2769
2770 assert_eq!(
2771 sender.check(Recipients::All).unwrap().recipients(),
2772 vec![pk2]
2773 );
2774 });
2775 }
2776
2777 #[test]
2778 fn test_sender_removed_from_peer_set_drops_message() {
2779 let executor = deterministic::Runner::default();
2780 executor.start(|context| async move {
2781 let (network, oracle) = Network::new(
2783 context.child("network"),
2784 Config {
2785 max_size: 1024 * 1024,
2786 max_peers_per_set: NZUsize!(2),
2787 disconnect_on_block: true,
2788 tracked_peer_sets: NZUsize!(1),
2789 },
2790 );
2791 network.start();
2792 let mut manager = oracle.manager();
2793 let mut subscription = manager.subscribe().await;
2794
2795 let sender_pk = PrivateKey::from_seed(1).public_key();
2797 let recipient_pk = PrivateKey::from_seed(2).public_key();
2798 manager.track(
2799 1,
2800 Set::try_from(vec![sender_pk.clone(), recipient_pk.clone()]).unwrap(),
2801 );
2802 let update = subscription.recv().await.unwrap();
2803 assert_eq!(update.index, 1);
2804
2805 let (mut sender, _) = oracle
2807 .control(sender_pk.clone())
2808 .register(0, TEST_QUOTA)
2809 .await
2810 .unwrap();
2811 let (_sender2, mut receiver) = oracle
2812 .control(recipient_pk.clone())
2813 .register(0, TEST_QUOTA)
2814 .await
2815 .unwrap();
2816
2817 oracle
2819 .add_link(
2820 sender_pk.clone(),
2821 recipient_pk.clone(),
2822 Link {
2823 latency: Duration::from_millis(1),
2824 jitter: Duration::ZERO,
2825 success_rate: probability!(1.0),
2826 },
2827 )
2828 .await
2829 .unwrap();
2830
2831 let initial_msg = IoBuf::from(b"tracked");
2833 let sent = sender.send(
2834 Recipients::One(recipient_pk.clone()),
2835 initial_msg.clone(),
2836 false,
2837 );
2838 assert_eq!(sent.len(), 1);
2839 assert_eq!(sent[0], recipient_pk);
2840 let (_pk, received) = receiver.recv().await.unwrap();
2841 assert_eq!(received, initial_msg.clone());
2842
2843 let other_pk = PrivateKey::from_seed(3).public_key();
2845 manager.track(
2846 2,
2847 Set::try_from(vec![recipient_pk.clone(), other_pk]).unwrap(),
2848 );
2849 let update = subscription.recv().await.unwrap();
2850 assert_eq!(update.index, 2);
2851
2852 let sent = sender.send(
2855 Recipients::One(recipient_pk.clone()),
2856 IoBuf::from(b"untracked"),
2857 false,
2858 );
2859 assert_eq!(sent, vec![recipient_pk.clone()]);
2860
2861 select! {
2863 _ = receiver.recv() => {
2864 panic!("unexpected message");
2865 },
2866 _ = context.sleep(Duration::from_secs(10)) => {},
2867 }
2868
2869 manager.track(
2871 3,
2872 Set::try_from(vec![sender_pk.clone(), recipient_pk.clone()]).unwrap(),
2873 );
2874 let update = subscription.recv().await.unwrap();
2875 assert_eq!(update.index, 3);
2876
2877 let sent = sender.send(
2879 Recipients::One(recipient_pk.clone()),
2880 initial_msg.clone(),
2881 false,
2882 );
2883 assert_eq!(sent.len(), 1);
2884 assert_eq!(sent[0], recipient_pk);
2885 let (_pk, received) = receiver.recv().await.unwrap();
2886 assert_eq!(received, initial_msg);
2887 });
2888 }
2889
2890 #[test]
2891 fn test_subscribe_to_peer_sets() {
2892 let executor = deterministic::Runner::default();
2893 executor.start(|context| async move {
2894 let (network, oracle) = Network::new(
2895 context.child("network"),
2896 Config {
2897 max_size: 1024 * 1024,
2898 max_peers_per_set: NZUsize!(2),
2899 disconnect_on_block: true,
2900 tracked_peer_sets: NZUsize!(2),
2901 },
2902 );
2903 network.start();
2904
2905 let mut manager = oracle.manager();
2907 let mut subscription = manager.subscribe().await;
2908
2909 let pk1 = PrivateKey::from_seed(1).public_key();
2911 let pk2 = PrivateKey::from_seed(2).public_key();
2912 let pk3 = PrivateKey::from_seed(3).public_key();
2913
2914 manager.track(1, Set::try_from(vec![pk1.clone(), pk2.clone()]).unwrap());
2916
2917 let update = subscription.recv().await.unwrap();
2919 assert_eq!(update.index, 1);
2920 assert_eq!(
2921 update.latest.primary,
2922 Set::try_from(vec![pk1.clone(), pk2.clone()]).unwrap()
2923 );
2924 assert!(update.latest.secondary.is_empty());
2925 assert_eq!(
2926 update.all.primary,
2927 Set::try_from(vec![pk1.clone(), pk2.clone()]).unwrap()
2928 );
2929 assert!(update.all.secondary.is_empty());
2930
2931 manager.track(2, Set::try_from(vec![pk2.clone(), pk3.clone()]).unwrap());
2933
2934 let update = subscription.recv().await.unwrap();
2936 assert_eq!(update.index, 2);
2937 assert_eq!(
2938 update.latest.primary,
2939 Set::try_from(vec![pk2.clone(), pk3.clone()]).unwrap()
2940 );
2941 assert!(update.latest.secondary.is_empty());
2942 assert_eq!(
2943 update.all.primary,
2944 vec![pk1.clone(), pk2.clone(), pk3.clone()]
2945 .try_into()
2946 .unwrap()
2947 );
2948 assert!(update.all.secondary.is_empty());
2949
2950 manager.track(3, Set::try_from(vec![pk1.clone(), pk3.clone()]).unwrap());
2952
2953 let update = subscription.recv().await.unwrap();
2955 assert_eq!(update.index, 3);
2956 assert_eq!(
2957 update.latest.primary,
2958 Set::try_from(vec![pk1.clone(), pk3.clone()]).unwrap()
2959 );
2960 assert!(update.latest.secondary.is_empty());
2961 assert_eq!(
2962 update.all.primary,
2963 vec![pk1.clone(), pk2.clone(), pk3.clone()]
2964 .try_into()
2965 .unwrap()
2966 );
2967 assert!(update.all.secondary.is_empty());
2968
2969 manager.track(4, Set::try_from(vec![pk1.clone(), pk3.clone()]).unwrap());
2971
2972 let update = subscription.recv().await.unwrap();
2974 assert_eq!(update.index, 4);
2975 assert_eq!(
2976 update.latest.primary,
2977 Set::try_from(vec![pk1.clone(), pk3.clone()]).unwrap()
2978 );
2979 assert!(update.latest.secondary.is_empty());
2980 assert_eq!(
2981 update.all.primary,
2982 Set::try_from(vec![pk1.clone(), pk3.clone()]).unwrap()
2983 );
2984 assert!(update.all.secondary.is_empty());
2985 });
2986 }
2987
2988 #[test]
2989 fn test_multiple_subscriptions() {
2990 let executor = deterministic::Runner::default();
2991 executor.start(|context| async move {
2992 let (network, oracle) = Network::new(
2993 context.child("network"),
2994 Config {
2995 max_size: 1024 * 1024,
2996 max_peers_per_set: NZUsize!(2),
2997 disconnect_on_block: true,
2998 tracked_peer_sets: NZUsize!(3),
2999 },
3000 );
3001 network.start();
3002
3003 let mut manager = oracle.manager();
3005 let mut subscription1 = manager.subscribe().await;
3006 let mut subscription2 = manager.subscribe().await;
3007 let mut subscription3 = manager.subscribe().await;
3008
3009 let pk1 = PrivateKey::from_seed(1).public_key();
3011 let pk2 = PrivateKey::from_seed(2).public_key();
3012
3013 manager.track(1, Set::try_from(vec![pk1.clone(), pk2.clone()]).unwrap());
3015
3016 let update1 = subscription1.recv().await.unwrap();
3018 let update2 = subscription2.recv().await.unwrap();
3019 let update3 = subscription3.recv().await.unwrap();
3020
3021 assert_eq!(update1.index, 1);
3022 assert_eq!(update2.index, 1);
3023 assert_eq!(update3.index, 1);
3024
3025 drop(subscription2);
3027
3028 manager.track(2, Set::try_from(vec![pk1.clone(), pk2.clone()]).unwrap());
3030
3031 let update1 = subscription1.recv().await.unwrap();
3033 let update3 = subscription3.recv().await.unwrap();
3034
3035 assert_eq!(update1.index, 2);
3036 assert_eq!(update3.index, 2);
3037 });
3038 }
3039
3040 #[test]
3041 fn test_subscription_includes_self_when_registered() {
3042 let executor = deterministic::Runner::default();
3043 executor.start(|context| async move {
3044 let (network, oracle) = Network::new(
3045 context.child("network"),
3046 Config {
3047 max_size: 1024 * 1024,
3048 max_peers_per_set: NZUsize!(2),
3049 disconnect_on_block: true,
3050 tracked_peer_sets: NZUsize!(2),
3051 },
3052 );
3053 network.start();
3054
3055 let self_pk = PrivateKey::from_seed(0).public_key();
3057 let other_pk = PrivateKey::from_seed(1).public_key();
3058
3059 let (_sender, _receiver) = oracle
3061 .control(self_pk.clone())
3062 .register(0, TEST_QUOTA)
3063 .await
3064 .unwrap();
3065
3066 let mut manager = oracle.manager();
3068 let mut subscription = manager.subscribe().await;
3069
3070 manager.track(1, Set::try_from(vec![other_pk.clone()]).unwrap());
3072
3073 let update = subscription.recv().await.unwrap();
3075 assert_eq!(update.index, 1);
3076 assert_eq!(update.latest.primary.len(), 1);
3077 assert!(update.latest.secondary.is_empty());
3078 assert_eq!(update.all.primary.len(), 1);
3079 assert!(update.all.secondary.is_empty());
3080
3081 assert!(
3083 update.latest.primary.position(&self_pk).is_none(),
3084 "latest primary set should not include self"
3085 );
3086 assert!(
3087 update.latest.primary.position(&other_pk).is_some(),
3088 "latest primary set should include other"
3089 );
3090
3091 assert!(
3093 update.all.primary.position(&self_pk).is_none(),
3094 "peer set should not include self"
3095 );
3096 assert!(
3097 update.all.primary.position(&other_pk).is_some(),
3098 "peer set should include other"
3099 );
3100
3101 manager.track(
3103 2,
3104 Set::try_from(vec![self_pk.clone(), other_pk.clone()]).unwrap(),
3105 );
3106
3107 let update = subscription.recv().await.unwrap();
3108 assert_eq!(update.index, 2);
3109 assert_eq!(update.latest.primary.len(), 2);
3110 assert!(update.latest.secondary.is_empty());
3111 assert_eq!(update.all.primary.len(), 2);
3112 assert!(update.all.secondary.is_empty());
3113
3114 assert!(
3116 update.latest.primary.position(&self_pk).is_some(),
3117 "latest primary set should include self"
3118 );
3119 assert!(
3120 update.latest.primary.position(&other_pk).is_some(),
3121 "latest primary set should include other"
3122 );
3123
3124 assert!(
3126 update.all.primary.position(&self_pk).is_some(),
3127 "peer set should include self"
3128 );
3129 assert!(
3130 update.all.primary.position(&other_pk).is_some(),
3131 "peer set should include other"
3132 );
3133 });
3134 }
3135
3136 #[test]
3137 fn test_rate_limiting() {
3138 let executor = deterministic::Runner::default();
3139 executor.start(|context| async move {
3140 let cfg = Config {
3141 max_size: 1024 * 1024,
3142 max_peers_per_set: NZUsize!(2),
3143 disconnect_on_block: true,
3144 tracked_peer_sets: NZUsize!(3),
3145 };
3146 let pk1 = ed25519::PrivateKey::from_seed(1).public_key();
3148 let pk2 = ed25519::PrivateKey::from_seed(2).public_key();
3149
3150 let (network, oracle) =
3151 Network::new_with_peers(context.child("network"), cfg, [pk1.clone(), pk2.clone()])
3152 .await;
3153 network.start();
3154
3155 let restrictive_quota = Quota::per_second(NZU32!(1));
3157 let control1 = oracle.control(pk1.clone());
3158 let (mut sender, _) = control1.register(0, restrictive_quota).await.unwrap();
3159 let control2 = oracle.control(pk2.clone());
3160 let (_, mut receiver) = control2.register(0, TEST_QUOTA).await.unwrap();
3161
3162 let link = ingress::Link {
3164 latency: Duration::from_millis(0),
3165 jitter: Duration::from_millis(0),
3166 success_rate: probability!(1.0),
3167 };
3168 oracle
3169 .add_link(pk1.clone(), pk2.clone(), link.clone())
3170 .await
3171 .unwrap();
3172 oracle.add_link(pk2.clone(), pk1, link).await.unwrap();
3173
3174 let msg1 = IoBuf::from(b"message1");
3176 let result1 = sender.send(Recipients::One(pk2.clone()), msg1.clone(), false);
3177 assert_eq!(result1.len(), 1, "first message should be sent");
3178
3179 let (_, received1) = receiver.recv().await.unwrap();
3181 assert_eq!(received1, msg1);
3182
3183 let msg2 = IoBuf::from(b"message2");
3185 let result2 = sender.send(Recipients::One(pk2.clone()), msg2.clone(), false);
3186 assert_eq!(
3187 result2.len(),
3188 0,
3189 "second message should be rate-limited (skipped)"
3190 );
3191
3192 context.sleep(Duration::from_secs(1)).await;
3194
3195 let msg3 = IoBuf::from(b"message3");
3197 let result3 = sender.send(Recipients::One(pk2.clone()), msg3.clone(), false);
3198 assert_eq!(result3.len(), 1, "third message should be sent after wait");
3199
3200 let (_, received3) = receiver.recv().await.unwrap();
3202 assert_eq!(received3, msg3);
3203 });
3204 }
3205
3206 #[test]
3207 fn test_blocked_subscription_tracks_block_and_unblock() {
3208 let executor = deterministic::Runner::default();
3209 executor.start(|context| async move {
3210 let cfg = Config {
3211 max_size: 1024 * 1024,
3212 max_peers_per_set: NZUsize!(2),
3213 disconnect_on_block: true,
3214 tracked_peer_sets: NZUsize!(3),
3215 };
3216 let pk1 = ed25519::PrivateKey::from_seed(1).public_key();
3217 let pk2 = ed25519::PrivateKey::from_seed(2).public_key();
3218 let (network, oracle) =
3219 Network::new_with_peers(context.child("network"), cfg, [pk1.clone(), pk2.clone()])
3220 .await;
3221 network.start();
3222
3223 let mut control1 = oracle.control(pk1.clone());
3225 let mut blocked = control1.blocked();
3226 assert!(blocked.next().await.unwrap().iter().next().is_none());
3227
3228 crate::block_peer(&mut control1, pk2.clone());
3230 assert_eq!(
3231 blocked
3232 .next()
3233 .await
3234 .unwrap()
3235 .iter()
3236 .cloned()
3237 .collect::<Vec<_>>(),
3238 vec![pk2.clone()]
3239 );
3240 let mut control2 = oracle.control(pk2.clone());
3241 crate::block_peer(&mut control2, pk1.clone());
3242
3243 crate::block_peer(&mut control1, pk2.clone());
3246 assert_eq!(oracle.blocked().await.unwrap().len(), 2);
3247 assert!(blocked.try_recv().is_err());
3248
3249 oracle.unblock(pk1.clone(), pk2.clone()).await.unwrap();
3251 assert!(blocked.next().await.unwrap().iter().next().is_none());
3252 assert_eq!(oracle.blocked().await.unwrap(), vec![(pk2, pk1)]);
3253 });
3254 }
3255
3256 #[test]
3257 fn test_operations_after_shutdown_do_not_panic() {
3258 let executor = deterministic::Runner::default();
3259 executor.start(|context| async move {
3260 let cfg = Config {
3261 max_size: 1024 * 1024,
3262 max_peers_per_set: NZUsize!(2),
3263 disconnect_on_block: true,
3264 tracked_peer_sets: NZUsize!(3),
3265 };
3266 let pk1 = ed25519::PrivateKey::from_seed(1).public_key();
3268 let pk2 = ed25519::PrivateKey::from_seed(2).public_key();
3269
3270 let (network, oracle) =
3271 Network::new_with_peers(context.child("network"), cfg, [pk1.clone(), pk2.clone()])
3272 .await;
3273 let handle = network.start();
3274 let mut manager = oracle.manager();
3275
3276 let control1 = oracle.control(pk1.clone());
3278 let (mut sender, _receiver) = control1.register(0, TEST_QUOTA).await.unwrap();
3279
3280 let link = ingress::Link {
3282 latency: Duration::from_millis(10),
3283 jitter: Duration::from_millis(0),
3284 success_rate: probability!(1.0),
3285 };
3286 oracle
3287 .add_link(pk1.clone(), pk2.clone(), link.clone())
3288 .await
3289 .unwrap();
3290 wait_for_task_count(&context, "network", |count| count > 0).await;
3291
3292 handle.abort();
3294 let _ = handle.await;
3295 wait_for_task_count(&context, "network", |count| count == 0).await;
3296
3297 let msg = IoBuf::from(b"test");
3299 let result = sender.send(Recipients::One(pk2.clone()), msg, false);
3300 assert!(result.is_empty(), "send after shutdown should return empty");
3301
3302 manager.track(1, Set::try_from([pk1.clone()]).unwrap());
3304 let _ = manager.peer_set(0).await;
3305 let _ = manager.subscribe().await;
3306
3307 let _ = oracle
3309 .add_link(pk1.clone(), pk2.clone(), link.clone())
3310 .await;
3311 let _ = oracle.remove_link(pk1.clone(), pk2.clone()).await;
3312 let _ = oracle.blocked().await;
3313
3314 let _ = control1.register(1, TEST_QUOTA).await;
3316 });
3317 }
3318
3319 fn clean_shutdown(seed: u64) {
3320 let cfg = deterministic::Config::default()
3321 .with_seed(seed)
3322 .with_timeout(Some(Duration::from_secs(30)));
3323 let executor = deterministic::Runner::new(cfg);
3324 executor.start(|context| async move {
3325 let cfg = Config {
3326 max_size: 1024 * 1024,
3327 max_peers_per_set: NZUsize!(2),
3328 disconnect_on_block: true,
3329 tracked_peer_sets: NZUsize!(3),
3330 };
3331 let pk1 = ed25519::PrivateKey::from_seed(1).public_key();
3333 let pk2 = ed25519::PrivateKey::from_seed(2).public_key();
3334
3335 let (network, oracle) =
3336 Network::new_with_peers(context.child("network"), cfg, [pk1.clone(), pk2.clone()])
3337 .await;
3338 let handle = network.start();
3339
3340 let control1 = oracle.control(pk1.clone());
3342 let control2 = oracle.control(pk2.clone());
3343 let (mut sender, _) = control1.register(0, TEST_QUOTA).await.unwrap();
3344 let (_, mut receiver) = control2.register(0, TEST_QUOTA).await.unwrap();
3345
3346 let link = ingress::Link {
3348 latency: Duration::from_millis(10),
3349 jitter: Duration::from_millis(0),
3350 success_rate: probability!(1.0),
3351 };
3352 oracle
3353 .add_link(pk1.clone(), pk2.clone(), link.clone())
3354 .await
3355 .unwrap();
3356 oracle
3357 .add_link(pk2.clone(), pk1.clone(), link)
3358 .await
3359 .unwrap();
3360
3361 wait_for_task_count(&context, "network", |count| count > 0).await;
3363 let running_before = count_running_tasks(&context, "network");
3364 assert!(
3365 running_before > 0,
3366 "at least one network task should be running"
3367 );
3368
3369 let msg = IoBuf::from(b"test_message");
3371 let result = sender.send(Recipients::One(pk2.clone()), msg.clone(), false);
3372 assert_eq!(result.len(), 1, "message should be sent");
3373
3374 let (_, received) = receiver.recv().await.unwrap();
3375 assert_eq!(received, msg, "message should be received");
3376
3377 handle.abort();
3379 let _ = handle.await;
3380
3381 wait_for_task_count(&context, "network", |count| count == 0).await;
3383 let running_after = count_running_tasks(&context, "network");
3384 assert_eq!(
3385 running_after, 0,
3386 "all network tasks should be stopped, but {running_after} still running"
3387 );
3388 });
3389 }
3390
3391 #[test]
3392 fn test_clean_shutdown() {
3393 for seed in 0..25 {
3394 clean_shutdown(seed);
3395 }
3396 }
3397
3398 #[test]
3399 fn test_socket_manager_overwrite() {
3400 let executor = deterministic::Runner::default();
3401 executor.start(|context| async move {
3402 let (network, oracle) = Network::new(
3404 context.child("network"),
3405 Config {
3406 max_size: 1024 * 1024,
3407 max_peers_per_set: NZUsize!(2),
3408 disconnect_on_block: true,
3409 tracked_peer_sets: NZUsize!(3),
3410 },
3411 );
3412 network.start();
3413
3414 let pk1 = PrivateKey::from_seed(1).public_key();
3416 let pk2 = PrivateKey::from_seed(2).public_key();
3417 let _pk3 = PrivateKey::from_seed(3).public_key();
3418
3419 let mut socket_manager = oracle.socket_manager();
3420
3421 let addr: Address = "127.0.0.1:8000".parse::<SocketAddr>().unwrap().into();
3423
3424 socket_manager.track(
3426 0,
3427 Map::<PublicKey, Address>::try_from([
3428 (
3429 pk1.clone(),
3430 "127.0.0.1:8001".parse::<SocketAddr>().unwrap().into(),
3431 ),
3432 (pk2, "127.0.0.1:8002".parse::<SocketAddr>().unwrap().into()),
3433 ])
3434 .unwrap(),
3435 );
3436
3437 socket_manager.overwrite([(pk1, addr)].try_into().unwrap());
3439 });
3440 }
3441
3442 #[test]
3443 fn test_subscribe_returns_current_peer_set() {
3444 let executor = deterministic::Runner::default();
3445 executor.start(|context| async move {
3446 let (network, oracle) = Network::new(
3447 context.child("network"),
3448 Config {
3449 max_size: 1024 * 1024,
3450 max_peers_per_set: NZUsize!(2),
3451 disconnect_on_block: true,
3452 tracked_peer_sets: NZUsize!(3),
3453 },
3454 );
3455 network.start();
3456
3457 let pk1 = PrivateKey::from_seed(0).public_key();
3459 let pk2 = PrivateKey::from_seed(1).public_key();
3460 let peers = ordered::Set::try_from(vec![pk1.clone(), pk2.clone()]).unwrap();
3461
3462 let mut manager = oracle.manager();
3463 Manager::track(&mut manager, 0, peers.clone());
3464
3465 let mut subscription = Provider::subscribe(&mut manager).await;
3468 let update = subscription
3469 .try_recv()
3470 .expect("current peer set should be available immediately after subscribe");
3471 assert_eq!(update.index, 0);
3472 assert_eq!(update.latest.primary, peers);
3473 assert!(update.latest.secondary.is_empty());
3474 });
3475 }
3476}