Skip to main content

commonware_broadcast/buffered/
mod.rs

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