Skip to main content

commonware_p2p/simulated/
mod.rs

1//! Simulate networking between peers with configurable link behavior (i.e. drops, latency, corruption, etc.).
2//!
3//! Both peer and link modification can be performed dynamically over the lifetime of the simulated network. This
4//! can be used to mimic transient network partitions, offline nodes (that later connect), and/or degrading link
5//! quality. Messages on a link are delivered in order, and optional per-peer bandwidth limits account for
6//! transmission delay and queueing.
7//!
8//! # Determinism
9//!
10//! `commonware-p2p::simulated` can be run deterministically when paired with `commonware-runtime::deterministic`.
11//! This makes it possible to reproduce an arbitrary order of delivered/dropped messages with a given seed.
12//!
13//! # Bandwidth Simulation
14//!
15//! The simulator provides a realistic model of bandwidth contention where network
16//! capacity is a shared, finite resource. Bandwidth is allocated via progressive
17//! filling to provide max-min fairness.
18//!
19//! _If no bandwidth constraints are provided (default behavior), progressive filling and bandwidth
20//! tracking are not performed (avoiding unnecessary overhead for minimal p2p testing common in CI)._
21//!
22//! ## Core Model
23//!
24//! Whenever a transfer starts or finishes, or a bandwidth limit is updated, we execute a scheduling tick:
25//!
26//! 1. **Collect Active Flows:** Gather every active transfer that still has
27//!    bytes to send. A flow is bound to one sender and to one receiver (if the message will be delivered).
28//! 2. **Compute Progressive Filling:** Run progressive filling to raise the rate of
29//!    every active flow in lock-step until some sender's egress or receiver's ingress
30//!    limit saturates (at which point the flow is frozen and the process repeats with what remains).
31//! 3. **Wait for the Next Event:** Using those rates, determine which flow will
32//!    finish first by computing how long it needs to transmit its remaining
33//!    bytes. Advance simulated time directly to that completion instant (advancing all other flows
34//!    by the bytes transferred over the interval).
35//! 4. **Deliver Message:** Remove the completed flow and pass the message to the receiver. Repeat from step 1
36//!    until all flows are processed.
37//!
38//! _Messages between the same pair of peers remain strictly ordered. When one
39//! message finishes, the next message on that link may begin sending at
40//! `arrival_time - new_latency` so that its first byte arrives immediately after
41//! the previous one is fully received._
42//!
43//! ## Latency vs. Transmission Delay
44//!
45//! The simulation correctly distinguishes between two key components of message delivery:
46//!
47//! - **Transmission Delay:** The time it takes to send all bytes of a message over
48//!   the link. This is determined by the message size and the available bandwidth
49//!   (e.g., a 10KB message on a 10KB/s link has a 1-second transmission delay).
50//! - **Network Latency:** The time it takes for a byte to travel from the sender
51//!   to the receiver, independent of bandwidth. This is configured via the `Link`
52//!   properties.
53//!
54//! The final delivery time of a message is the sum of when its transmission completes
55//! plus the simulated network latency. This model ensures that large messages correctly
56//! occupy the network link for longer periods, affecting other concurrent transfers,
57//! while still accounting for the physical travel time of the data.
58//!
59//! # Example
60//!
61//! ```rust
62//! use commonware_p2p::simulated::{Config, Link, Network};
63//! use commonware_cryptography::{ed25519, PrivateKey, Signer as _, PublicKey as _, };
64//! use commonware_runtime::{deterministic, Metrics, Quota, Runner, Spawner, Supervisor};
65//! use commonware_utils::{NZU32, NZUsize, probability};
66//! use std::time::Duration;
67//!
68//! // Generate peers
69//! let peers = vec![
70//!     ed25519::PrivateKey::from_seed(0).public_key(),
71//!     ed25519::PrivateKey::from_seed(1).public_key(),
72//!     ed25519::PrivateKey::from_seed(2).public_key(),
73//!     ed25519::PrivateKey::from_seed(3).public_key(),
74//! ];
75//!
76//! // Configure network
77//! let p2p_cfg = Config {
78//!     max_size: 1024 * 1024, // 1MB
79//!     max_peers_per_set: NZUsize!(4),
80//!     disconnect_on_block: true,
81//!     tracked_peer_sets: NZUsize!(3),
82//! };
83//!
84//! // Rate limit quota (1000 messages per second per peer)
85//! let quota = Quota::per_second(NZU32!(1000));
86//!
87//! // Start context
88//! let executor = deterministic::Runner::seeded(0);
89//! executor.start(|context| async move {
90//!     // Initialize the network with an initial peer set (tracked at id 0).
91//!     let (network, oracle) =
92//!         Network::new_with_peers(context.child("network"), p2p_cfg, peers.clone())
93//!             .await;
94//!
95//!     // Start network
96//!     let network_handler = network.start();
97//!
98//!     let (sender1, receiver1) = oracle.control(peers[0].clone()).register(0, quota).await.unwrap();
99//!     let (sender2, receiver2) = oracle.control(peers[1].clone()).register(0, quota).await.unwrap();
100//!
101//!     // Set bandwidth limits
102//!     // peer[0]: 10KB/s egress, unlimited ingress
103//!     // peer[1]: unlimited egress, 5KB/s ingress
104//!     oracle.limit_bandwidth(peers[0].clone(), Some(10_000), None).await.unwrap();
105//!     oracle.limit_bandwidth(peers[1].clone(), None, Some(5_000)).await.unwrap();
106//!
107//!     // Link 2 peers
108//!     oracle.add_link(
109//!         peers[0].clone(),
110//!         peers[1].clone(),
111//!         Link {
112//!             latency: Duration::from_millis(5),
113//!             jitter: Duration::from_millis(2),
114//!             success_rate: probability!(0.75),
115//!         },
116//!     ).await.unwrap();
117//!
118//!     // ... Use sender and receiver ...
119//!
120//!     // Update link
121//!     oracle.remove_link(
122//!         peers[0].clone(),
123//!         peers[1].clone(),
124//!     ).await.unwrap();
125//!     oracle.add_link(
126//!         peers[0].clone(),
127//!         peers[1].clone(),
128//!         Link {
129//!             latency: Duration::from_millis(100),
130//!             jitter: Duration::from_millis(25),
131//!             success_rate: probability!(0.8),
132//!         },
133//!     ).await.unwrap();
134//!
135//!     // ... Use sender and receiver ...
136//!
137//!     // Shutdown network
138//!     network_handler.abort();
139//! });
140//! ```
141
142mod bandwidth;
143mod ingress;
144mod metrics;
145mod network;
146mod transmitter;
147
148use thiserror::Error;
149
150/// Errors that can occur when interacting with the network.
151#[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    /// Default rate limit set high enough to not interfere with normal operation
214    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            // Create simulated network
243            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            // Start network
254            network.start();
255
256            // Register agents
257            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                    // Exiting early here tests the case where the recipient end of an agent is dropped
275                });
276            }
277            track_peers(&oracle, agents.keys().cloned()).await;
278
279            // Link all outbound-capable agents.
280            let only_inbound = PrivateKey::from_seed(0).public_key();
281            for agent in agents.keys() {
282                if agent == &only_inbound {
283                    // Leave this peer inbound-only to exercise missing-link handling.
284                    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                    // BTreeMap iteration gives deterministic sender selection.
310                    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            // Wait for all recipients
325            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        // Collect outputs
335        let mut outputs = Vec::new();
336        for seed in 0..seeds {
337            outputs.push(simulate_messages(seed, size));
338        }
339
340        // Confirm outputs are deterministic
341        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            // Create simulated network
359            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            // Start network
370            network.start();
371
372            // Register agents
373            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            // Send invalid message
385            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            // Create simulated network
400            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            // Start network
411            network.start();
412
413            // Register agents
414            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            // Attempt to link self
422            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            // Confirm error is correct
435            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            // Create simulated network
444            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            // Start network
455            network.start();
456
457            // Setup links
458            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            // Register channels
486            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            // Send messages
499            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            // Update channel
512            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            // Send message
519            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            // Listen on original
532            assert!(matches!(
533                my_receiver.recv().await,
534                Err(Error::NetworkClosed)
535            ));
536
537            // Send on original gracefully handles a closed channel.
538            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            // Create peers
547            let pk1 = PrivateKey::from_seed(0).public_key();
548            let pk2 = PrivateKey::from_seed(1).public_key();
549
550            // Create simulated network
551            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            // Add link
565            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            // Register channels
579            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            // Send message
591            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            // Create simulated network
604            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            // Start network
615            network.start();
616
617            // Register agents
618            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            // Register unused channels
632            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            // Link agents
645            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            // Send messages
671            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            // Confirm message delivery
677            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            // Create simulated network
691            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            // Start network
702            network.start();
703
704            // Register agents
705            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            // Link agents
720            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            // Send message
734            let msg = IoBuf::from(b"hello from pk1");
735            sender1.send(Recipients::One(pk2), msg, false);
736
737            // Confirm no message delivery
738            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            // Create simulated network
752            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            // Start network
763            network.start();
764
765            // Define agents
766            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            // Link agents
781            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            // Send messages
807            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            // Confirm message delivery
813            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            // Create simulated network
827            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            // Start network
838            network.start();
839
840            // Register agents
841            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            // Send messages
856            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            // Confirm no message delivery
862            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            // Link agents
873            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            // Send messages
899            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            // Confirm message delivery
905            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            // Remove links
913            oracle.remove_link(pk1.clone(), pk2.clone()).await.unwrap();
914            oracle.remove_link(pk2.clone(), pk1.clone()).await.unwrap();
915
916            // Send messages
917            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            // Confirm no message delivery
923            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            // Remove non-existent links
934            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        // Create two agents
949        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        // Set bandwidth limits
965        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        // Link the two agents
975        oracle
976            .add_link(
977                pk1.clone(),
978                pk2.clone(),
979                Link {
980                    // No latency so it doesn't interfere with bandwidth delay calculation
981                    latency: Duration::ZERO,
982                    jitter: Duration::ZERO,
983                    success_rate: probability!(1.0),
984                },
985            )
986            .await
987            .unwrap();
988
989        // Send a message from agent 1 to 2
990        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        // Measure how long it takes for agent 2 to receive the message
995        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            // Both sender and receiver have the same bandiwdth (1000 B/s)
1026            // 500 bytes at 1000 B/s = 0.5 seconds
1027            test_bandwidth_between_peers(
1028                &mut context,
1029                &oracle,
1030                0,
1031                Some(1000), // sender egress
1032                Some(1000), // receiver ingress
1033                500,        // message size
1034                500,        // expected duration in ms
1035            )
1036            .await;
1037
1038            // Sender has lower bandwidth (500 B/s) than receiver (2000 B/s)
1039            // Should be limited by sender's 500 B/s
1040            // 250 bytes at 500 B/s = 0.5 seconds
1041            test_bandwidth_between_peers(
1042                &mut context,
1043                &oracle,
1044                1,
1045                Some(500),  // sender egress
1046                Some(2000), // receiver ingress
1047                250,        // message size
1048                500,        // expected duration in ms
1049            )
1050            .await;
1051
1052            // Sender has higher bandwidth (2000 B/s) than receiver (500 B/s)
1053            // Should be limited by receiver's 500 B/s
1054            // 250 bytes at 500 B/s = 0.5 seconds
1055            test_bandwidth_between_peers(
1056                &mut context,
1057                &oracle,
1058                2,
1059                Some(2000), // sender egress
1060                Some(500),  // receiver ingress
1061                250,        // message size
1062                500,        // expected duration in ms
1063            )
1064            .await;
1065
1066            // Unlimited sender, limited receiver
1067            // Should be limited by receiver's 1000 B/s
1068            // 500 bytes at 1000 B/s = 0.5 seconds
1069            test_bandwidth_between_peers(
1070                &mut context,
1071                &oracle,
1072                3,
1073                None,       // sender egress (unlimited)
1074                Some(1000), // receiver ingress
1075                500,        // message size
1076                500,        // expected duration in ms
1077            )
1078            .await;
1079
1080            // Limited sender, unlimited receiver
1081            // Should be limited by sender's 1000 B/s
1082            // 500 bytes at 1000 B/s = 0.5 seconds
1083            test_bandwidth_between_peers(
1084                &mut context,
1085                &oracle,
1086                4,
1087                Some(1000), // sender egress
1088                None,       // receiver ingress (unlimited)
1089                500,        // message size
1090                500,        // expected duration in ms
1091            )
1092            .await;
1093
1094            // Unlimited sender, unlimited receiver
1095            // Delivery should be (almost) instant
1096            test_bandwidth_between_peers(
1097                &mut context,
1098                &oracle,
1099                5,
1100                None, // sender egress (unlimited)
1101                None, // receiver ingress (unlimited)
1102                500,  // message size
1103                0,    // expected duration in ms
1104            )
1105            .await;
1106        });
1107    }
1108
1109    #[test]
1110    fn test_bandwidth_contention() {
1111        // Test bandwidth contention with many peers (one-to-many and many-to-one scenarios)
1112        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            // Configuration
1126            const NUM_PEERS: usize = 100;
1127            const MESSAGE_SIZE: usize = 1000; // 1KB per message
1128            const EFFECTIVE_BPS: usize = 10_000; // 10KB/s egress/ingress per peer
1129
1130            // Create peers
1131            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            // Create the main peer (index 0) and 100 other peers
1136            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            // Set bandwidth limits for all peers
1150            for pk in &peers {
1151                oracle
1152                    .limit_bandwidth(pk.clone(), Some(EFFECTIVE_BPS), Some(EFFECTIVE_BPS))
1153                    .await
1154                    .unwrap();
1155            }
1156
1157            // Link all peers to the main peer (peers[0]) with zero latency
1158            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            // One-to-many (main peer sends to all others). Verifies that bandwidth limits
1186            // are properly enforced when sending to multiple recipients
1187            let start = context.current();
1188
1189            // Send message to all peers concurrently
1190            // and wait for all sends to be acknowledged
1191            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            // Verify all messages are received
1197            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            // Calculate expected time
1206            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            // Many-to-one (all peers send to the main peer)
1218            let start = context.current();
1219
1220            // Each peer sends a message to the main peer concurrently and we wait for all
1221            // sends to be acknowledged
1222            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            // Collect all messages at the main peer
1228            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            // Calculate expected time
1241            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            // Verify we received from all peers
1253            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        // Test that messages arrive in order even with variable latency
1263        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            // Register agents
1277            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            // Link agents with high jitter to create variable delays
1292            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            // Send multiple messages that should arrive in order
1306            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            // Receive messages and verify they arrive in order
1319            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            // Send slow message
1359            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            // Update link
1379            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            // Send fast message
1394            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            // Create 10 senders and 1 receiver
1448            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                // Each sender has 10KB/s egress
1461                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            // Receiver has 100KB/s ingress
1476            oracle
1477                .limit_bandwidth(receiver.clone(), None, Some(100_000))
1478                .await
1479                .unwrap();
1480
1481            // Add links with no latency
1482            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            // All senders send 10KB simultaneously
1500            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            // All 10 messages should be received at ~1s
1507            // (100KB total data at 100KB/s aggregate bandwidth)
1508            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                // Messages should all complete around 1s
1513                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        // Test that 1 fast sender (100KB/s) sending to 10 receivers (10KB/s each)
1525        // should complete all sends in ~1s and all messages received in ~1s
1526        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            // Create fast sender
1540            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            // Sender has 100KB/s egress
1548            oracle
1549                .limit_bandwidth(sender.clone(), Some(100_000), None)
1550                .await
1551                .unwrap();
1552
1553            // Create 10 receivers
1554            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                // Each receiver has 10KB/s ingress
1567                oracle
1568                    .limit_bandwidth(receiver.clone(), None, Some(10_000))
1569                    .await
1570                    .unwrap();
1571
1572                // Add link with no latency
1573                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            // Send 10KB to each receiver (100KB total)
1595            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            // Each receiver should receive their 10KB message in ~1s (10KB at 10KB/s)
1601            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                // All messages should be received around 1s
1607                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        // Test that 10 slow senders (1KB/s each) sending to a fast receiver (10KB/s)
1619        // should complete all transfers in ~1s
1620        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            // Create 10 slow senders
1634            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                // Each sender has 1KB/s egress (slow)
1647                oracle
1648                    .limit_bandwidth(sender.clone(), Some(1_000), None)
1649                    .await
1650                    .unwrap();
1651            }
1652
1653            // Create fast receiver
1654            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            // Receiver has 10KB/s ingress (can handle all 10 senders at full speed)
1663            oracle
1664                .limit_bandwidth(receiver.clone(), None, Some(10_000))
1665                .await
1666                .unwrap();
1667
1668            // Add links with no latency
1669            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            // All senders send 1KB simultaneously
1687            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            // Each sender takes 1s to transmit 1KB at 1KB/s
1694            // All transmissions happen in parallel, so total send time is ~1s
1695
1696            // All 10 messages (10KB total) should be received at ~1s
1697            // Receiver processes at 10KB/s, can handle all 10KB in 1s
1698            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                // All messages should complete around 1s
1703                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        // Test that bandwidth is dynamically allocated as
1715        // transfers start and complete at different times
1716        //
1717        // 3 senders to 1 receiver, starting at different times
1718        // Receiver has 30KB/s, senders each have 30KB/s
1719        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            // Create 3 senders
1733            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                // Each sender has 30KB/s egress
1746                oracle
1747                    .limit_bandwidth(sender.clone(), Some(30_000), None)
1748                    .await
1749                    .unwrap();
1750            }
1751
1752            // Create receiver with 30KB/s ingress
1753            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            // Add links with minimal latency
1766            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            // Sender 0: sends 30KB at t=0
1785            // Gets full 30KB/s for the first 0.5s, then shares with sender 1
1786            // at 15KB/s until completion at t=1.5s
1787            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            // Sender 1: sends 30KB at t=0.5s
1795            // Shares bandwidth with sender 0 (15KB/s each) until t=1.5s,
1796            // then gets the full 30KB/s
1797            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            // Sender 2: sends 15KB at t=1.5s and shares the receiver with
1806            // sender 1, completing at roughly t=2.5s
1807            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            // Receive and verify timing
1816            // Message 0: starts at t=0, shares bandwidth after 0.5s,
1817            // and completes at t=1.5s (plus link latency)
1818            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            // The algorithm may deliver messages in a different order based on
1827            // efficient bandwidth usage. Let's collect the next two messages and
1828            // verify their timings regardless of order.
1829            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            // Figure out which message is which based on content
1836            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            // Message 1 (30KB) started at t=0.5s
1846            // Message 2 (15KB) started at t=1.5s
1847            // With efficient scheduling, message 2 might complete first since it's smaller
1848            // Both should complete between 1.5s and 2.5s
1849            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        // Test dynamic allocation with different message sizes arriving simultaneously
1864        // This tests that smaller messages complete first when bandwidth is shared
1865        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            // Create 3 senders
1879            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                // Each sender has unlimited egress
1892                oracle
1893                    .limit_bandwidth(sender.clone(), None, None)
1894                    .await
1895                    .unwrap();
1896            }
1897
1898            // Create receiver with 30KB/s ingress
1899            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            // Add links
1912            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            // All start at the same time but with different sizes
1930            //
1931            // The scheduler reserves bandwidth in advance, the actual behavior
1932            // depends on the order tasks are processed. Since all senders
1933            // start at once, they'll compete for bandwidth
1934            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            // Receive messages. They arrive in the order they were scheduled,
1943            // not necessarily size order. Collect all messages and sort by
1944            // receive time to verify timing
1945            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            // When all start at once, they'll reserve bandwidth slots
1953            // sequentially. First gets full 30KB/s, others wait or get
1954            // remaining bandwidth. Just verify all messages arrived and total
1955            // time is reasonable
1956            assert_eq!(messages.len(), 3);
1957
1958            // Total data is 60KB at 30KB/s receiver ingress = 2s minimum
1959            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        // Test that bandwidth pipe is only reserved for transmission duration, not latency
1970        // This means new messages can start transmitting while others are still in flight
1971        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            // Create two peers
1985            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            // Set bandwidth: 1000 B/s (1 byte per millisecond)
2001            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            // Add link with significant latency (1 second)
2011            oracle
2012                .add_link(
2013                    sender.clone(),
2014                    receiver.clone(),
2015                    Link {
2016                        latency: Duration::from_secs(1), // 1 second latency
2017                        jitter: Duration::ZERO,
2018                        success_rate: probability!(1.0),
2019                    },
2020                )
2021                .await
2022                .unwrap();
2023
2024            // Send 3 messages of 500 bytes each
2025            // At 1000 B/s, each message takes 500ms to transmit
2026            // With 1s latency, if pipe was reserved for tx+latency, total would be:
2027            //   - Msg 1: 0-1500ms (500ms tx + 1000ms latency)
2028            //   - Msg 2: 1500-3000ms (starts after msg 1 fully delivered)
2029            //   - Msg 3: 3000-4500ms
2030            // But if pipe is only reserved during tx (correct behavior):
2031            //   - Msg 1: tx 0-500ms, delivered at 1500ms
2032            //   - Msg 2: tx 500-1000ms, delivered at 2000ms
2033            //   - Msg 3: tx 1000-1500ms, delivered at 2500ms
2034            let start = context.current();
2035
2036            // Send all messages in quick succession
2037            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            // Wait for all receives to complete and record their completion times
2043            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            // Messages should be received at:
2051            // - Msg 1: ~1500ms (500ms transmission + 1000ms latency)
2052            // - Msg 2: ~2000ms (500ms wait + 500ms transmission + 1000ms latency)
2053            // - Msg 3: ~2500ms (1000ms wait + 500ms transmission + 1000ms latency)
2054            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        // This test verifies that bandwidth changes affect NEW transfers,
2073        // not transfers already in progress (which have their reservations locked in)
2074        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            // Register peers and establish link
2091            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), // Small latency
2108                        jitter: Duration::ZERO,
2109                        success_rate: probability!(1.0),
2110                    },
2111                )
2112                .await
2113                .unwrap();
2114
2115            // Initial bandwidth: 10 KB/s
2116            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            // Send first message at 10 KB/s
2126            let msg1 = IoBuf::from(vec![1u8; 20_000]); // 20 KB
2127            let start_time = context.current();
2128            sender_tx.send(Recipients::One(pk_receiver.clone()), msg1.clone(), false);
2129
2130            // Receive first message (should take ~2s at 10KB/s)
2131            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            // Change bandwidth to 2 KB/s
2141            oracle
2142                .limit_bandwidth(pk_sender.clone(), Some(2_000), None)
2143                .await
2144                .unwrap();
2145
2146            // Send second message at new bandwidth
2147            let msg2 = IoBuf::from(vec![2u8; 10_000]); // 10 KB
2148            let msg2_start = context.current();
2149            sender_tx.send(Recipients::One(pk_receiver.clone()), msg2.clone(), false);
2150
2151            // Receive second message (should take ~5s at 2KB/s)
2152            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            // Register peers and establish link
2182            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            // Set sender bandwidth to 0
2207            oracle
2208                .limit_bandwidth(pk_receiver.clone(), None, Some(0))
2209                .await
2210                .unwrap();
2211
2212            // Send message to receiver
2213            let msg1 = IoBuf::from(vec![1u8; 20_000]); // 20 KB
2214            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            // Message should not be received after 10 seconds
2219            select! {
2220                _ = receiver_rx.recv() => {
2221                    panic!("unexpected message");
2222                },
2223                _ = context.sleep(Duration::from_secs(10)) => {},
2224            }
2225
2226            // Unset bandwidth
2227            oracle
2228                .limit_bandwidth(pk_receiver.clone(), None, None)
2229                .await
2230                .unwrap();
2231
2232            // Message should be immediately received
2233            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            // Register peers and establish link
2261            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            // Set sender bandwidth to 0
2286            oracle
2287                .limit_bandwidth(pk_sender.clone(), Some(0), None)
2288                .await
2289                .unwrap();
2290
2291            // Send message to receiver
2292            let msg1 = IoBuf::from(vec![1u8; 20_000]); // 20 KB
2293            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            // Message should not be received after 10 seconds
2298            select! {
2299                _ = receiver_rx.recv() => {
2300                    panic!("unexpected message");
2301                },
2302                _ = context.sleep(Duration::from_secs(10)) => {},
2303            }
2304
2305            // Unset bandwidth
2306            oracle
2307                .limit_bandwidth(pk_sender.clone(), None, None)
2308                .await
2309                .unwrap();
2310
2311            // Message should be immediately received
2312            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            // pk2 is in both primary and secondary TrackedPeers; stored secondary is pk3 only.
2452            // latest and aggregate.secondary omit pk2; aggregate.primary still includes pk2.
2453            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            // Same key in primary and secondary maps; primary address and role win (secondary ignored).
2540            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            // Use asymmetric addresses where ingress (dial) differs from egress (filter)
2593            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            // Verify peer set contains expected keys (addresses are ignored by simulated network)
2612            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            // Verify subscription works
2617            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), // Only track 2 peer sets
2637                },
2638            );
2639            network.start();
2640
2641            // Create 4 peers
2642            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            // Register first peer set with pk1 and pk2
2648            let mut manager = oracle.manager();
2649            manager.track(1, Set::try_from(vec![pk1.clone(), pk2.clone()]).unwrap());
2650
2651            // Register channels for all peers
2652            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            // Create bidirectional links between all peers
2674            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            // pk1 can broadcast to pk2, but not pk3 while pk3 is untracked.
2694            let recipients = sender1.check(Recipients::All).unwrap().recipients();
2695            assert_eq!(recipients, vec![pk2.clone()]);
2696            assert!(!recipients.contains(&pk3));
2697
2698            // Register second peer set with pk2 and pk3
2699            manager.track(2, Set::try_from(vec![pk2.clone(), pk3.clone()]).unwrap());
2700            assert!(manager.peer_set(2).await.is_some());
2701
2702            // Now pk3 is in a peer set and pk1 can broadcast to it.
2703            let recipients = sender1.check(Recipients::All).unwrap().recipients();
2704            assert!(recipients.contains(&pk3));
2705
2706            // Register third peer set with pk3 and pk4 (this will evict peer set 1)
2707            manager.track(3, Set::try_from(vec![pk3.clone(), pk4.clone()]).unwrap());
2708            assert!(manager.peer_set(3).await.is_some());
2709
2710            // pk1 should now be removed from all peer sets.
2711            let recipients = sender2.check(Recipients::All).unwrap().recipients();
2712            assert!(!recipients.contains(&pk1));
2713
2714            // pk3 should still be reachable (in sets 2 and 3).
2715            assert!(recipients.contains(&pk3));
2716
2717            // pk4 should be reachable (in set 3).
2718            let recipients = sender3.check(Recipients::All).unwrap().recipients();
2719            assert!(recipients.contains(&pk4));
2720
2721            // Verify peer set contents
2722            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            // Peer set 1 should no longer exist
2731            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            // Create a simulated network
2782            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            // Register a peer set
2796            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            // Register channels
2806            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            // Add link
2818            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            // Send and confirm message
2832            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            // Register another peer set
2844            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            // Explicit sends are accepted locally, but the network drops
2853            // messages from peers no longer in any peer set.
2854            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            // Confirm message was not delivered
2862            select! {
2863                _ = receiver.recv() => {
2864                    panic!("unexpected message");
2865                },
2866                _ = context.sleep(Duration::from_secs(10)) => {},
2867            }
2868
2869            // Add a peer back to a peer set
2870            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            // Send message from a peer now back in a peer set
2878            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            // Subscribe to peer set updates
2906            let mut manager = oracle.manager();
2907            let mut subscription = manager.subscribe().await;
2908
2909            // Create peers
2910            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            // Register first peer set
2915            manager.track(1, Set::try_from(vec![pk1.clone(), pk2.clone()]).unwrap());
2916
2917            // Verify we receive the notification
2918            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            // Register second peer set
2932            manager.track(2, Set::try_from(vec![pk2.clone(), pk3.clone()]).unwrap());
2933
2934            // Verify we receive the notification
2935            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            // Register third peer set
2951            manager.track(3, Set::try_from(vec![pk1.clone(), pk3.clone()]).unwrap());
2952
2953            // Verify we receive the notification
2954            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            // Register fourth peer set
2970            manager.track(4, Set::try_from(vec![pk1.clone(), pk3.clone()]).unwrap());
2971
2972            // Verify we receive the notification
2973            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            // Create multiple subscriptions
3004            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            // Create peers
3010            let pk1 = PrivateKey::from_seed(1).public_key();
3011            let pk2 = PrivateKey::from_seed(2).public_key();
3012
3013            // Register a peer set
3014            manager.track(1, Set::try_from(vec![pk1.clone(), pk2.clone()]).unwrap());
3015
3016            // Verify all subscriptions receive the notification
3017            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 one subscription
3026            drop(subscription2);
3027
3028            // Register another peer set
3029            manager.track(2, Set::try_from(vec![pk1.clone(), pk2.clone()]).unwrap());
3030
3031            // Verify remaining subscriptions still receive notifications
3032            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            // Create "self" and "other" peers
3056            let self_pk = PrivateKey::from_seed(0).public_key();
3057            let other_pk = PrivateKey::from_seed(1).public_key();
3058
3059            // Register a channel for self (this creates the peer in the network)
3060            let (_sender, _receiver) = oracle
3061                .control(self_pk.clone())
3062                .register(0, TEST_QUOTA)
3063                .await
3064                .unwrap();
3065
3066            // Subscribe to peer set updates
3067            let mut manager = oracle.manager();
3068            let mut subscription = manager.subscribe().await;
3069
3070            // Register a peer set that does NOT include self
3071            manager.track(1, Set::try_from(vec![other_pk.clone()]).unwrap());
3072
3073            // Receive subscription notification
3074            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            // Self should NOT be in the latest primary set
3082            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            // Self should NOT be in the peer set (not tracked)
3092            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            // Now register a peer set that DOES include self
3102            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            // Both peers should be in the latest primary set
3115            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            // Both peers should be in the peer set
3125            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            // Create two public keys
3147            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            // Register with a very restrictive quota: 1 message per second
3156            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            // Add bidirectional links
3163            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            // First message should succeed immediately
3175            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            // Verify first message is received
3180            let (_, received1) = receiver.recv().await.unwrap();
3181            assert_eq!(received1, msg1);
3182
3183            // Second message should be rate-limited (quota is 1/sec, no time has passed)
3184            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            // Advance time by 1 second to allow the rate limiter to reset
3193            context.sleep(Duration::from_secs(1)).await;
3194
3195            // Third message should succeed after waiting
3196            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            // Verify third message is received
3201            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            // A new subscription starts with the current, empty set.
3224            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            // Blocking publishes the peer to the blocker's subscribers only.
3229            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            // Neither the other peer's block nor a repeated block reaches pk1's
3244            // subscription, since pk1's set is unchanged.
3245            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            // Lifting the block publishes the empty set again.
3250            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            // Create peers
3267            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            // Register channels
3277            let control1 = oracle.control(pk1.clone());
3278            let (mut sender, _receiver) = control1.register(0, TEST_QUOTA).await.unwrap();
3279
3280            // Add link
3281            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            // Abort the network
3293            handle.abort();
3294            let _ = handle.await;
3295            wait_for_task_count(&context, "network", |count| count == 0).await;
3296
3297            // Sending messages should not panic and should return empty
3298            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 operations should not panic
3303            manager.track(1, Set::try_from([pk1.clone()]).unwrap());
3304            let _ = manager.peer_set(0).await;
3305            let _ = manager.subscribe().await;
3306
3307            // Oracle operations should not panic
3308            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            // Control operations should not panic
3315            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            // Create peers
3332            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            // Register channels
3341            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            // Add bidirectional links
3347            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 until the network has started at least one task.
3362            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            // Send and receive a message to verify network is functional
3370            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            // Abort the network
3378            handle.abort();
3379            let _ = handle.await;
3380
3381            // Wait until task shutdown is reflected in the metrics.
3382            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            // Create simulated network with peer set tracking
3403            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            // Generate keys
3415            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            // Simulated network ignores addresses, so overwrite is a no-op
3422            let addr: Address = "127.0.0.1:8000".parse::<SocketAddr>().unwrap().into();
3423
3424            // Register a peer set
3425            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            // overwrite is a no-op for simulated network (addresses not used)
3438            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            // Create peers and track them
3458            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            // Subscribe after tracking. The current peer set should be
3466            // available immediately on the subscription channel.
3467            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}