Skip to main content

commonware_resolver/p2p/
mod.rs

1//! Resolve data identified by a fixed-length key by using the P2P network.
2//!
3//! # Overview
4//!
5//! The `p2p` module enables resolving data by fixed-length keys in a P2P network. Central to the
6//! module is the `peer` actor which manages the fetch lifecycle. Its mailbox allows
7//! initiation and pruning of fetches via the `Resolver` interface.
8//!
9//! The peer handles an arbitrarily large number of concurrent fetches by sending requests
10//! to other peers and processing their responses. It selects peers based on performance, retrying
11//! with another peer if one fails or provides invalid data. Fetches persist until pruned or
12//! fulfilled, delivering data to the `Consumer` for verification.
13//!
14//! The `Consumer` checks data integrity and authenticity (critical in an adversarial environment)
15//! and returns `true` if valid, completing the fetch, or `false` to retry. Pruning a fetch with
16//! in-progress response validation aborts that validation. If the aborted validation would have
17//! returned `false`, the peer is not blocked for that response.
18//!
19//! The peer also serves data to other peers, forwarding network requests to the `Producer`. The
20//! `Producer` provides data asynchronously (e.g., from storage). If it fails, the peer sends an
21//! empty response, prompting the requester to retry elsewhere. Each message between peers contains
22//! an ID. Each request is sent with a unique ID, and each response includes the ID of the request
23//! it responds to.
24//!
25//! # Targeting
26//!
27//! Callers can restrict fetches to specific target peers using
28//! [`TargetedResolver::fetch_targeted`](crate::TargetedResolver::fetch_targeted).
29//! Only target peers are tried, there is no automatic fallback to other peers. Targets persist through
30//! transient failures (timeout, "no data" response, send failure) since the peer might be slow or
31//! receive the data later.
32//!
33//! While a fetch is in progress, callers can modify targeting:
34//! - [`TargetedResolver::fetch_targeted`](crate::TargetedResolver::fetch_targeted) adds peers to the existing target set
35//!   (only if the fetch already has targets, an "all" fetch remains unrestricted)
36//! - [`Resolver::fetch`](crate::Resolver::fetch) clears all targets, allowing fallback to any peer
37//!
38//! These modifications only apply to in-progress fetches. Once a fetch completes (success, pruning,
39//! or blocked peer), the targets for that key are cleared automatically.
40//!
41//! # Subscribers
42//!
43//! [`Resolver::fetch`](crate::Resolver::fetch) accepts a peer-visible key and a
44//! subscriber. This is useful when several subscribers can share the same peer-visible
45//! fetch. A fetch remains active while at least one attached subscriber satisfies the latest
46//! [`Resolver::retain`](crate::Resolver::retain) predicate. When the fetch resolves, the
47//! key and currently retained subscribers are supplied to
48//! [`Consumer::deliver`](crate::Consumer::deliver). Subscribers added while response validation
49//! is in progress are delivered the same accepted response locally.
50//!
51//! # Peer Selection
52//!
53//! Outbound fetches are only sent to peers in `latest.primary` (see [commonware_p2p::Provider]) but inbound
54//! requests are handled for all connected peers. Thus, callers that still expect a key to be fetchable after
55//! a peer set update must ensure the latest primary set can serve it.
56//!
57//! [`TargetedResolver::fetch_targeted`](crate::TargetedResolver::fetch_targeted) can narrow the current primary set
58//! further, but it does not bypass that latest-primary filter. Explicit targets that are no longer
59//! in the latest primary set are ignored until they become primary again.
60//!
61//! # Performance Considerations
62//!
63//! The peer supports arbitrarily many concurrent fetches, but resource usage generally
64//! depends on the rate-limiting configuration of the underlying P2P network.
65
66use bytes::Bytes;
67use commonware_utils::{channel::oneshot, Span};
68
69mod config;
70pub use config::Config;
71mod engine;
72pub use engine::Engine;
73mod fetcher;
74mod inflight;
75mod ingress;
76pub use ingress::Mailbox;
77mod metrics;
78mod wire;
79
80#[cfg(feature = "mocks")]
81pub mod mocks;
82
83/// Serves data requested by the network.
84pub trait Producer: Clone + Send + 'static {
85    /// Type used to key data requested from peers.
86    type Key: Span;
87
88    /// Serve a request received from the network.
89    fn produce(&mut self, key: Self::Key) -> oneshot::Receiver<Bytes>;
90}
91
92#[cfg(test)]
93mod tests {
94    use super::{
95        mocks::{Consumer, Key, Producer},
96        Config, Engine, Mailbox,
97    };
98    use crate::{Delivery, Fetch, Resolver, TargetedResolver};
99    use bytes::Bytes;
100    use commonware_cryptography::{
101        ed25519::{PrivateKey, PublicKey},
102        Signer,
103    };
104    use commonware_macros::{select, test_traced};
105    use commonware_p2p::{
106        simulated::{Link, Network, Oracle, Receiver, Sender},
107        Blocker, Manager as _, Provider, TrackedPeers,
108    };
109    use commonware_runtime::{
110        deterministic, telemetry::metrics::count_running_tasks, Clock, Metrics as _, Quota, Runner,
111        Spawner as _, Supervisor as _,
112    };
113    use commonware_utils::{
114        channel::{fallible::FallibleExt, mpsc, oneshot},
115        non_empty_vec,
116        ordered::Set,
117        sync::Mutex,
118        NZUsize, NZU32,
119    };
120    use std::{
121        collections::{HashMap, VecDeque},
122        num::{NonZeroU32, NonZeroUsize},
123        sync::Arc,
124        time::Duration,
125    };
126
127    const MAILBOX_SIZE: NonZeroUsize = NZUsize!(1024);
128    const RATE_LIMIT: NonZeroU32 = NZU32!(10);
129    const INITIAL_DURATION: Duration = Duration::from_millis(100);
130    const TIMEOUT: Duration = Duration::from_millis(400);
131    const FETCH_RETRY_TIMEOUT: Duration = Duration::from_millis(100);
132    const LINK: Link = Link {
133        latency: Duration::from_millis(10),
134        jitter: Duration::from_millis(1),
135        success_rate: 1.0,
136    };
137    const LINK_UNRELIABLE: Link = Link {
138        latency: Duration::from_millis(10),
139        jitter: Duration::from_millis(1),
140        success_rate: 0.5,
141    };
142
143    fn status_metric_total(metrics: &str, name: &str, status: &str) -> u64 {
144        let prefix = format!("{name}{{");
145        let status_label = format!("status=\"{status}\"");
146        metrics
147            .lines()
148            .filter(|line| line.starts_with(&prefix) && line.contains(&status_label))
149            .map(|line| {
150                line.split_whitespace()
151                    .next_back()
152                    .expect("metric line must have a value")
153                    .parse::<u64>()
154                    .expect("status metric value must be an integer")
155            })
156            .sum()
157    }
158
159    async fn setup_network_and_peers(
160        context: &deterministic::Context,
161        peer_seeds: &[u64],
162    ) -> (
163        Oracle<PublicKey, deterministic::Context>,
164        Vec<PrivateKey>,
165        Vec<PublicKey>,
166        Vec<(
167            Sender<PublicKey, deterministic::Context>,
168            Receiver<PublicKey>,
169        )>,
170    ) {
171        setup_network_and_peers_with_rate_limit(context, peer_seeds, Quota::per_second(RATE_LIMIT))
172            .await
173    }
174
175    async fn setup_network_and_peers_with_rate_limit(
176        context: &deterministic::Context,
177        peer_seeds: &[u64],
178        rate_limit: Quota,
179    ) -> (
180        Oracle<PublicKey, deterministic::Context>,
181        Vec<PrivateKey>,
182        Vec<PublicKey>,
183        Vec<(
184            Sender<PublicKey, deterministic::Context>,
185            Receiver<PublicKey>,
186        )>,
187    ) {
188        let (network, oracle) = Network::new(
189            context.child("network"),
190            commonware_p2p::simulated::Config {
191                max_size: 1024 * 1024,
192                disconnect_on_block: true,
193                tracked_peer_sets: NZUsize!(3),
194            },
195        );
196        network.start();
197
198        let schemes: Vec<PrivateKey> = peer_seeds
199            .iter()
200            .map(|seed| PrivateKey::from_seed(*seed))
201            .collect();
202        let peers: Vec<PublicKey> = schemes.iter().map(|s| s.public_key()).collect();
203        let mut manager = oracle.manager();
204        manager.track(0, Set::try_from(peers.clone()).unwrap());
205
206        let mut connections = Vec::new();
207        for peer in &peers {
208            let (sender, receiver) = oracle
209                .control(peer.clone())
210                .register(0, rate_limit)
211                .await
212                .unwrap();
213            connections.push((sender, receiver));
214        }
215
216        (oracle, schemes, peers, connections)
217    }
218
219    async fn add_link(
220        oracle: &mut Oracle<PublicKey, deterministic::Context>,
221        link: Link,
222        peers: &[PublicKey],
223        from: usize,
224        to: usize,
225    ) {
226        oracle
227            .add_link(peers[from].clone(), peers[to].clone(), link.clone())
228            .await
229            .unwrap();
230        oracle
231            .add_link(peers[to].clone(), peers[from].clone(), link)
232            .await
233            .unwrap();
234    }
235
236    #[derive(Clone, Default)]
237    struct SequencedProducer {
238        data: Arc<Mutex<HashMap<Key, VecDeque<Bytes>>>>,
239    }
240
241    impl SequencedProducer {
242        fn insert(&mut self, key: Key, values: impl IntoIterator<Item = Bytes>) {
243            self.data.lock().insert(key, values.into_iter().collect());
244        }
245
246        fn remaining(&self, key: &Key) -> Vec<Bytes> {
247            self.data
248                .lock()
249                .get(key)
250                .map(|values| values.iter().cloned().collect())
251                .unwrap_or_default()
252        }
253    }
254
255    impl crate::p2p::Producer for SequencedProducer {
256        type Key = Key;
257
258        fn produce(&mut self, key: Self::Key) -> oneshot::Receiver<Bytes> {
259            let (sender, receiver) = oneshot::channel();
260            if let Some(value) = self.data.lock().get_mut(&key).and_then(VecDeque::pop_front) {
261                let _ = sender.send(value);
262            }
263            receiver
264        }
265    }
266
267    fn setup_and_spawn_actor<C, R>(
268        context: &deterministic::Context,
269        provider: impl Provider<PublicKey = PublicKey>,
270        blocker: impl Blocker<PublicKey = PublicKey>,
271        signer: impl Signer<PublicKey = PublicKey>,
272        connection: (
273            Sender<PublicKey, deterministic::Context>,
274            Receiver<PublicKey>,
275        ),
276        consumer: C,
277        producer: Producer<Key, Bytes>,
278    ) -> Mailbox<Key, PublicKey, R>
279    where
280        C: crate::Consumer<Key = Key, Subscriber = R, Value = Bytes>,
281        R: Clone + Ord + Send + 'static,
282    {
283        setup_and_spawn_actor_with_producer(
284            context, provider, blocker, signer, connection, consumer, producer,
285        )
286    }
287
288    fn setup_and_spawn_actor_with_producer<C, R, Pro>(
289        context: &deterministic::Context,
290        provider: impl Provider<PublicKey = PublicKey>,
291        blocker: impl Blocker<PublicKey = PublicKey>,
292        signer: impl Signer<PublicKey = PublicKey>,
293        connection: (
294            Sender<PublicKey, deterministic::Context>,
295            Receiver<PublicKey>,
296        ),
297        consumer: C,
298        producer: Pro,
299    ) -> Mailbox<Key, PublicKey, R>
300    where
301        C: crate::Consumer<Key = Key, Subscriber = R, Value = Bytes>,
302        Pro: crate::p2p::Producer<Key = Key>,
303        R: Clone + Ord + Send + 'static,
304    {
305        let public_key = signer.public_key();
306        let (engine, mailbox) = Engine::new(
307            context.child("actor").with_attribute("peer", &public_key),
308            Config {
309                peer_provider: provider,
310                blocker,
311                consumer,
312                producer,
313                mailbox_size: MAILBOX_SIZE,
314                me: Some(public_key),
315                initial: INITIAL_DURATION,
316                timeout: TIMEOUT,
317                fetch_retry_timeout: FETCH_RETRY_TIMEOUT,
318                priority_requests: false,
319                priority_responses: false,
320            },
321        );
322        engine.start(connection);
323
324        mailbox
325    }
326
327    type DeliveryGate = (oneshot::Receiver<()>, bool);
328    type DeliveryGates = Arc<Mutex<VecDeque<DeliveryGate>>>;
329
330    #[derive(Clone)]
331    struct BlockingConsumer {
332        context: Arc<deterministic::Context>,
333        sender: mpsc::UnboundedSender<(Key, Bytes)>,
334        started: mpsc::UnboundedSender<Key>,
335        gates: DeliveryGates,
336    }
337
338    impl BlockingConsumer {
339        fn new(
340            context: deterministic::Context,
341            gates: Vec<DeliveryGate>,
342        ) -> (
343            Self,
344            mpsc::UnboundedReceiver<(Key, Bytes)>,
345            mpsc::UnboundedReceiver<Key>,
346        ) {
347            let (sender, receiver) = mpsc::unbounded_channel();
348            let (started, started_receiver) = mpsc::unbounded_channel();
349            (
350                Self {
351                    context: Arc::new(context),
352                    sender,
353                    started,
354                    gates: Arc::new(Mutex::new(gates.into())),
355                },
356                receiver,
357                started_receiver,
358            )
359        }
360    }
361
362    impl crate::Consumer for BlockingConsumer {
363        type Key = Key;
364        type Value = Bytes;
365        type Subscriber = ();
366
367        fn deliver(
368            &mut self,
369            delivery: Delivery<Self::Key, Self::Subscriber>,
370            value: Self::Value,
371        ) -> oneshot::Receiver<bool> {
372            let key = delivery.key;
373            self.started.send_lossy(key.clone());
374            let (gate, valid) = self
375                .gates
376                .lock()
377                .pop_front()
378                .map_or((None, true), |(gate, valid)| (Some(gate), valid));
379            let (mut response, receiver) = oneshot::channel();
380            let sender = self.sender.clone();
381            self.context.child("delivery").spawn(move |_| async move {
382                if let Some(gate) = gate {
383                    select! {
384                        _ = response.closed() => return,
385                        result = gate => {
386                            if result.is_err() {
387                                let _ = response.send(false);
388                                return;
389                            }
390                        },
391                    }
392                }
393                if valid {
394                    sender.send_lossy((key, value));
395                }
396                let _ = response.send(valid);
397            });
398            receiver
399        }
400    }
401
402    #[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)]
403    struct SubscriberTag(u16);
404
405    type RecordedDelivery = (Delivery<Key, SubscriberTag>, Bytes);
406
407    #[derive(Clone)]
408    struct BlockingSubscriberRecordingConsumer {
409        context: Arc<deterministic::Context>,
410        sender: mpsc::UnboundedSender<RecordedDelivery>,
411        started: mpsc::UnboundedSender<Delivery<Key, SubscriberTag>>,
412        gates: DeliveryGates,
413    }
414
415    impl BlockingSubscriberRecordingConsumer {
416        fn new(
417            context: deterministic::Context,
418            gates: Vec<DeliveryGate>,
419        ) -> (
420            Self,
421            mpsc::UnboundedReceiver<RecordedDelivery>,
422            mpsc::UnboundedReceiver<Delivery<Key, SubscriberTag>>,
423        ) {
424            let (sender, receiver) = mpsc::unbounded_channel();
425            let (started, started_receiver) = mpsc::unbounded_channel();
426            (
427                Self {
428                    context: Arc::new(context),
429                    sender,
430                    started,
431                    gates: Arc::new(Mutex::new(gates.into())),
432                },
433                receiver,
434                started_receiver,
435            )
436        }
437    }
438
439    impl crate::Consumer for BlockingSubscriberRecordingConsumer {
440        type Key = Key;
441        type Value = Bytes;
442        type Subscriber = SubscriberTag;
443
444        fn deliver(
445            &mut self,
446            delivery: Delivery<Self::Key, Self::Subscriber>,
447            value: Self::Value,
448        ) -> oneshot::Receiver<bool> {
449            self.started.send_lossy(delivery.clone());
450            let (gate, valid) = self
451                .gates
452                .lock()
453                .pop_front()
454                .map_or((None, true), |(gate, valid)| (Some(gate), valid));
455            let (mut response, receiver) = oneshot::channel();
456            let sender = self.sender.clone();
457            self.context.child("delivery").spawn(move |_| async move {
458                if let Some(gate) = gate {
459                    select! {
460                        _ = response.closed() => return,
461                        result = gate => {
462                            if result.is_err() {
463                                let _ = response.send(false);
464                                return;
465                            }
466                        },
467                    }
468                }
469                if valid {
470                    sender.send_lossy((delivery, value));
471                }
472                let _ = response.send(valid);
473            });
474            receiver
475        }
476    }
477
478    #[derive(Clone)]
479    struct SubscriberRecordingConsumer {
480        sender: mpsc::UnboundedSender<RecordedDelivery>,
481    }
482
483    impl SubscriberRecordingConsumer {
484        fn new() -> (Self, mpsc::UnboundedReceiver<RecordedDelivery>) {
485            let (sender, receiver) = mpsc::unbounded_channel();
486            (Self { sender }, receiver)
487        }
488    }
489
490    impl crate::Consumer for SubscriberRecordingConsumer {
491        type Key = Key;
492        type Value = Bytes;
493        type Subscriber = SubscriberTag;
494
495        fn deliver(
496            &mut self,
497            delivery: Delivery<Self::Key, Self::Subscriber>,
498            value: Self::Value,
499        ) -> oneshot::Receiver<bool> {
500            let (sender, receiver) = oneshot::channel();
501            self.sender.send_lossy((delivery, value));
502            let _ = sender.send(true);
503            receiver
504        }
505    }
506
507    fn dummy_consumer() -> Consumer<Key, Bytes> {
508        Consumer::dummy()
509    }
510
511    fn consumer() -> (Consumer<Key, Bytes>, mpsc::UnboundedReceiver<(Key, Bytes)>) {
512        Consumer::new()
513    }
514
515    async fn wait_for_blocked(
516        context: &deterministic::Context,
517        oracle: &Oracle<PublicKey, deterministic::Context>,
518        blocker: &PublicKey,
519        blocked: &PublicKey,
520    ) {
521        loop {
522            let blocked_peers = oracle.blocked().await.unwrap();
523            if blocked_peers
524                .iter()
525                .any(|(a, b)| a == blocker && b == blocked)
526            {
527                return;
528            }
529            context.sleep(Duration::from_millis(10)).await;
530        }
531    }
532
533    /// Tests that fetching a key from another peer succeeds when data is available.
534    /// This test sets up two peers, where Peer 1 requests data that Peer 2 has,
535    /// and verifies that the data is correctly delivered to Peer 1's consumer.
536    #[test_traced]
537    fn test_fetch_success() {
538        let executor = deterministic::Runner::timed(Duration::from_secs(10));
539        executor.start(|context| async move {
540            let (mut oracle, mut schemes, peers, mut connections) =
541                setup_network_and_peers(&context, &[1, 2]).await;
542
543            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
544
545            let key = Key(2);
546            let mut prod2 = Producer::default();
547            prod2.insert(key.clone(), Bytes::from("data for key 2"));
548
549            let (cons1, mut cons_out1) = consumer();
550
551            let scheme = schemes.remove(0);
552            let mut mailbox1 = setup_and_spawn_actor(
553                &context,
554                oracle.manager(),
555                oracle.control(scheme.public_key()),
556                scheme,
557                connections.remove(0),
558                cons1,
559                Producer::default(),
560            );
561
562            let scheme = schemes.remove(0);
563            let _mailbox2 = setup_and_spawn_actor(
564                &context,
565                oracle.manager(),
566                oracle.control(scheme.public_key()),
567                scheme,
568                connections.remove(0),
569                dummy_consumer(),
570                prod2,
571            );
572
573            mailbox1.fetch(key.clone());
574
575            let (key_actual, value) = cons_out1.recv().await.unwrap();
576            assert_eq!(key_actual, key);
577            assert_eq!(value, Bytes::from("data for key 2"));
578        });
579    }
580
581    #[test_traced]
582    fn test_pending_delivery_does_not_block_engine() {
583        let executor = deterministic::Runner::timed(Duration::from_secs(10));
584        executor.start(|context| async move {
585            let (mut oracle, mut schemes, peers, mut connections) =
586                setup_network_and_peers(&context, &[1, 2, 3]).await;
587
588            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
589            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
590
591            let key1 = Key(1);
592            let key2 = Key(2);
593            let data1 = Bytes::from("data for key 1");
594            let data2 = Bytes::from("data for key 2");
595
596            let mut prod2 = Producer::default();
597            prod2.insert(key1.clone(), data1.clone());
598
599            let mut prod3 = Producer::default();
600            prod3.insert(key2.clone(), data2.clone());
601
602            let (gate_sender1, gate_receiver1) = oneshot::channel();
603            let (gate_sender2, gate_receiver2) = oneshot::channel();
604            let (cons1, mut cons_out1, mut started) = BlockingConsumer::new(
605                context.child("consumer"),
606                vec![(gate_receiver1, true), (gate_receiver2, true)],
607            );
608
609            let scheme = schemes.remove(0);
610            let mut mailbox1 = setup_and_spawn_actor(
611                &context,
612                oracle.manager(),
613                oracle.control(scheme.public_key()),
614                scheme,
615                connections.remove(0),
616                cons1,
617                Producer::default(),
618            );
619
620            let scheme = schemes.remove(0);
621            let _mailbox2 = setup_and_spawn_actor(
622                &context,
623                oracle.manager(),
624                oracle.control(scheme.public_key()),
625                scheme,
626                connections.remove(0),
627                dummy_consumer(),
628                prod2,
629            );
630
631            let scheme = schemes.remove(0);
632            let _mailbox3 = setup_and_spawn_actor(
633                &context,
634                oracle.manager(),
635                oracle.control(scheme.public_key()),
636                scheme,
637                connections.remove(0),
638                dummy_consumer(),
639                prod3,
640            );
641
642            mailbox1.fetch(key1.clone());
643            let started_key = started.recv().await.expect("delivery did not start");
644            assert_eq!(started_key, key1);
645
646            mailbox1.fetch(key2.clone());
647            select! {
648                started_key = started.recv() => {
649                    assert_eq!(started_key.expect("delivery did not start"), key2);
650                },
651                _ = context.sleep(Duration::from_secs(2)) => {
652                    panic!("resolver engine blocked on pending delivery");
653                },
654            };
655
656            gate_sender2.send(()).unwrap();
657            let (key_actual, value) = cons_out1.recv().await.expect("consumer channel closed");
658            assert_eq!(key_actual, key2);
659            assert_eq!(value, data2);
660
661            gate_sender1.send(()).unwrap();
662            let (key_actual, value) = cons_out1.recv().await.expect("consumer channel closed");
663            assert_eq!(key_actual, key1);
664            assert_eq!(value, data1);
665        });
666    }
667
668    #[test_traced]
669    fn test_retain_drops_pending_delivery() {
670        let executor = deterministic::Runner::timed(Duration::from_secs(10));
671        executor.start(|context| async move {
672            let (mut oracle, mut schemes, peers, mut connections) =
673                setup_network_and_peers(&context, &[1, 2]).await;
674
675            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
676
677            let key = Key(1);
678            let data = Bytes::from("data for key 1");
679            let mut prod2 = Producer::default();
680            prod2.insert(key.clone(), data.clone());
681
682            let (mut first_gate_sender, first_gate_receiver) = oneshot::channel();
683            let (second_gate_sender, second_gate_receiver) = oneshot::channel();
684            let (cons1, mut cons_out1, mut started) = BlockingConsumer::new(
685                context.child("consumer"),
686                vec![(first_gate_receiver, true), (second_gate_receiver, true)],
687            );
688
689            let scheme = schemes.remove(0);
690            let mut mailbox1 = setup_and_spawn_actor(
691                &context,
692                oracle.manager(),
693                oracle.control(scheme.public_key()),
694                scheme,
695                connections.remove(0),
696                cons1,
697                Producer::default(),
698            );
699
700            let scheme = schemes.remove(0);
701            let _mailbox2 = setup_and_spawn_actor(
702                &context,
703                oracle.manager(),
704                oracle.control(scheme.public_key()),
705                scheme,
706                connections.remove(0),
707                dummy_consumer(),
708                prod2,
709            );
710
711            mailbox1.fetch(key.clone());
712            let started_key = started.recv().await.expect("delivery did not start");
713            assert_eq!(started_key, key);
714
715            let canceled = key.clone();
716            mailbox1.retain(move |key, _| key != &canceled);
717            mailbox1.fetch(key.clone());
718
719            first_gate_sender.closed().await;
720            let started_key = started.recv().await.expect("second delivery did not start");
721            assert_eq!(started_key, key);
722
723            second_gate_sender.send(()).unwrap();
724            let (key_actual, value) = cons_out1.recv().await.expect("consumer channel closed");
725            assert_eq!(key_actual, key);
726            assert_eq!(value, data);
727
728            select! {
729                _ = cons_out1.recv() => panic!("unexpected extra event"),
730                _ = context.sleep(Duration::from_millis(100)) => {},
731            };
732        });
733    }
734
735    #[test_traced]
736    fn test_invalid_delivery_retries_and_rearms_slot() {
737        let executor = deterministic::Runner::timed(Duration::from_secs(10));
738        executor.start(|context| async move {
739            let (mut oracle, mut schemes, peers, mut connections) =
740                setup_network_and_peers(&context, &[1, 2, 3]).await;
741
742            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
743
744            let key = Key(1);
745            let data = Bytes::from("data for key 1");
746
747            let mut prod2 = Producer::default();
748            prod2.insert(key.clone(), data.clone());
749
750            let mut prod3 = Producer::default();
751            prod3.insert(key.clone(), data.clone());
752
753            let (first_gate_sender, first_gate_receiver) = oneshot::channel();
754            let (second_gate_sender, second_gate_receiver) = oneshot::channel();
755            let (cons1, mut cons_out1, mut started) = BlockingConsumer::new(
756                context.child("consumer"),
757                vec![(first_gate_receiver, false), (second_gate_receiver, true)],
758            );
759
760            let scheme = schemes.remove(0);
761            let mut mailbox1 = setup_and_spawn_actor(
762                &context,
763                oracle.manager(),
764                oracle.control(scheme.public_key()),
765                scheme,
766                connections.remove(0),
767                cons1,
768                Producer::default(),
769            );
770
771            let scheme = schemes.remove(0);
772            let _mailbox2 = setup_and_spawn_actor(
773                &context,
774                oracle.manager(),
775                oracle.control(scheme.public_key()),
776                scheme,
777                connections.remove(0),
778                dummy_consumer(),
779                prod2,
780            );
781
782            let scheme = schemes.remove(0);
783            let _mailbox3 = setup_and_spawn_actor(
784                &context,
785                oracle.manager(),
786                oracle.control(scheme.public_key()),
787                scheme,
788                connections.remove(0),
789                dummy_consumer(),
790                prod3,
791            );
792
793            mailbox1.fetch_targeted(
794                key.clone(),
795                non_empty_vec![peers[1].clone(), peers[2].clone()],
796            );
797            let started_key = started.recv().await.expect("delivery did not start");
798            assert_eq!(started_key, key);
799
800            first_gate_sender.send(()).unwrap();
801            wait_for_blocked(&context, &oracle, &peers[0], &peers[1]).await;
802
803            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
804            oracle.manager().track(
805                1,
806                Set::try_from([peers[0].clone(), peers[2].clone()]).unwrap(),
807            );
808
809            let started_key = started.recv().await.expect("retry delivery did not start");
810            assert_eq!(started_key, key);
811
812            second_gate_sender.send(()).unwrap();
813            let (key_actual, value) = cons_out1.recv().await.expect("consumer channel closed");
814            assert_eq!(key_actual, key);
815            assert_eq!(value, data);
816        });
817    }
818
819    async fn run_pending_invalid_delivery_race(
820        context: &deterministic::Context,
821        validation_first: bool,
822    ) {
823        let (mut oracle, mut schemes, peers, mut connections) =
824            setup_network_and_peers(context, &[1, 2]).await;
825
826        add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
827
828        let key = Key(1);
829        let mut prod2 = Producer::default();
830        prod2.insert(key.clone(), Bytes::from("data for key 1"));
831
832        let (mut gate_sender, gate_receiver) = oneshot::channel();
833        let (cons1, mut cons_out1, mut started) =
834            BlockingConsumer::new(context.child("consumer"), vec![(gate_receiver, false)]);
835
836        let scheme = schemes.remove(0);
837        let mut mailbox1 = setup_and_spawn_actor(
838            context,
839            oracle.manager(),
840            oracle.control(scheme.public_key()),
841            scheme,
842            connections.remove(0),
843            cons1,
844            Producer::default(),
845        );
846
847        let scheme = schemes.remove(0);
848        let _mailbox2 = setup_and_spawn_actor(
849            context,
850            oracle.manager(),
851            oracle.control(scheme.public_key()),
852            scheme,
853            connections.remove(0),
854            dummy_consumer(),
855            prod2,
856        );
857
858        mailbox1.fetch(key.clone());
859        let started_key = started.recv().await.expect("delivery did not start");
860        assert_eq!(started_key, key);
861
862        if validation_first {
863            gate_sender.send(()).unwrap();
864            wait_for_blocked(context, &oracle, &peers[0], &peers[1]).await;
865            mailbox1.retain(|_, _| false);
866            let blocked = oracle.blocked().await.unwrap();
867            assert_eq!(blocked.len(), 1);
868            assert_eq!(blocked[0].0, peers[0]);
869            assert_eq!(blocked[0].1, peers[1]);
870        } else {
871            mailbox1.retain(|_, _| false);
872            gate_sender.closed().await;
873            assert!(oracle.blocked().await.unwrap().is_empty());
874        }
875
876        select! {
877            _ = cons_out1.recv() => panic!("unexpected event"),
878            _ = context.sleep(Duration::from_millis(100)) => {},
879        };
880    }
881
882    #[test_traced]
883    fn test_retain_pending_invalid_delivery_race() {
884        let executor = deterministic::Runner::timed(Duration::from_secs(10));
885        executor.start(|context| async move {
886            run_pending_invalid_delivery_race(&context, false).await;
887            run_pending_invalid_delivery_race(&context, true).await;
888        });
889    }
890
891    /// Tests that pruning a fetch leaves the consumer untouched.
892    /// This test initiates a fetch and immediately prunes it, verifying
893    /// that the consumer does not receive any event.
894    #[test_traced]
895    fn test_retain_drops_fetch() {
896        let executor = deterministic::Runner::timed(Duration::from_secs(10));
897        executor.start(|context| async move {
898            let (oracle, mut schemes, _peers, mut connections) =
899                setup_network_and_peers(&context, &[1]).await;
900
901            let (cons1, mut cons_out1) = consumer();
902            let prod1 = Producer::default();
903
904            let scheme = schemes.remove(0);
905            let mut mailbox1 = setup_and_spawn_actor(
906                &context,
907                oracle.manager(),
908                oracle.control(scheme.public_key()),
909                scheme,
910                connections.remove(0),
911                cons1,
912                prod1,
913            );
914
915            let key = Key(3);
916            mailbox1.fetch(key.clone());
917            let canceled = key.clone();
918            mailbox1.retain(move |key, _| key != &canceled);
919
920            select! {
921                _ = cons_out1.recv() => panic!("unexpected event"),
922                _ = context.sleep(Duration::from_millis(100)) => {},
923            };
924        });
925    }
926
927    /// Tests fetching data from a peer when some peers lack the data.
928    /// This test sets up three peers, where Peer 1 requests data that only Peer 3 has.
929    /// It verifies that the resolver retries with another peer and successfully
930    /// delivers the data to Peer 1's consumer.
931    #[test_traced]
932    fn test_peer_no_data() {
933        let executor = deterministic::Runner::timed(Duration::from_secs(10));
934        executor.start(|context| async move {
935            let (mut oracle, mut schemes, peers, mut connections) =
936                setup_network_and_peers(&context, &[1, 2, 3]).await;
937
938            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
939            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
940
941            let prod1 = Producer::default();
942            let prod2 = Producer::default();
943            let mut prod3 = Producer::default();
944            let key = Key(3);
945            prod3.insert(key.clone(), Bytes::from("data for key 3"));
946
947            let (cons1, mut cons_out1) = consumer();
948
949            let scheme = schemes.remove(0);
950            let mut mailbox1 = setup_and_spawn_actor(
951                &context,
952                oracle.manager(),
953                oracle.control(scheme.public_key()),
954                scheme,
955                connections.remove(0),
956                cons1,
957                prod1,
958            );
959
960            let scheme = schemes.remove(0);
961            let _mailbox2 = setup_and_spawn_actor(
962                &context,
963                oracle.manager(),
964                oracle.control(scheme.public_key()),
965                scheme,
966                connections.remove(0),
967                dummy_consumer(),
968                prod2,
969            );
970
971            let scheme = schemes.remove(0);
972            let _mailbox3 = setup_and_spawn_actor(
973                &context,
974                oracle.manager(),
975                oracle.control(scheme.public_key()),
976                scheme,
977                connections.remove(0),
978                dummy_consumer(),
979                prod3,
980            );
981
982            mailbox1.fetch(key.clone());
983
984            let (key_actual, value) = cons_out1.recv().await.unwrap();
985            assert_eq!(key_actual, key);
986            assert_eq!(value, Bytes::from("data for key 3"));
987        });
988    }
989
990    /// Tests fetching when no peers are available.
991    /// This test sets up a single peer with an empty peer provider (no peers).
992    /// It initiates a fetch, waits beyond the retry timeout, prunes the fetch,
993    /// and verifies that the consumer receives a failure notification.
994    #[test_traced]
995    fn test_no_peers_available() {
996        let executor = deterministic::Runner::timed(Duration::from_secs(10));
997        executor.start(|context| async move {
998            let (oracle, mut schemes, _peers, mut connections) =
999                setup_network_and_peers(&context, &[1]).await;
1000
1001            let (cons1, mut cons_out1) = consumer();
1002            let prod1 = Producer::default();
1003
1004            let scheme = schemes.remove(0);
1005            let mut mailbox1 = setup_and_spawn_actor(
1006                &context,
1007                oracle.manager(),
1008                oracle.control(scheme.public_key()),
1009                scheme,
1010                connections.remove(0),
1011                cons1,
1012                prod1,
1013            );
1014
1015            mailbox1.fetch(Key(4));
1016            context.sleep(Duration::from_secs(5)).await;
1017
1018            // With no peers, no event should arrive
1019            select! {
1020                _ = cons_out1.recv() => panic!("Fetch should have failed due to no peers"),
1021                _ = context.sleep(Duration::from_millis(100)) => {},
1022            };
1023        });
1024    }
1025
1026    /// Tests that fetches issued before the first peer set arrives stay pending and complete once
1027    /// the initial update is tracked.
1028    #[test_traced]
1029    fn test_fetch_before_initial_peer_set_waits_for_update() {
1030        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1031        executor.start(|context| async move {
1032            let (network, mut oracle) = Network::new(
1033                context.child("network"),
1034                commonware_p2p::simulated::Config {
1035                    max_size: 1024 * 1024,
1036                    disconnect_on_block: true,
1037                    tracked_peer_sets: NZUsize!(1),
1038                },
1039            );
1040            network.start();
1041
1042            let mut schemes = [1_u64, 2]
1043                .into_iter()
1044                .map(PrivateKey::from_seed)
1045                .collect::<Vec<_>>();
1046            schemes.sort_by_key(|s| s.public_key());
1047            let peers: Vec<PublicKey> = schemes.iter().map(|s| s.public_key()).collect();
1048
1049            let mut connections = Vec::new();
1050            for peer in &peers {
1051                let (sender, receiver) = oracle
1052                    .control(peer.clone())
1053                    .register(0, Quota::per_second(RATE_LIMIT))
1054                    .await
1055                    .unwrap();
1056                connections.push((sender, receiver));
1057            }
1058
1059            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1060
1061            let key = Key(2);
1062            let mut prod2 = Producer::default();
1063            prod2.insert(key.clone(), Bytes::from("data for key 2"));
1064
1065            let (cons1, mut cons_out1) = consumer();
1066
1067            let scheme = schemes.remove(0);
1068            let mut mailbox1 = setup_and_spawn_actor(
1069                &context,
1070                oracle.manager(),
1071                oracle.control(scheme.public_key()),
1072                scheme,
1073                connections.remove(0),
1074                cons1,
1075                Producer::default(),
1076            );
1077
1078            let scheme = schemes.remove(0);
1079            let _mailbox2 = setup_and_spawn_actor(
1080                &context,
1081                oracle.manager(),
1082                oracle.control(scheme.public_key()),
1083                scheme,
1084                connections.remove(0),
1085                dummy_consumer(),
1086                prod2,
1087            );
1088
1089            mailbox1.fetch(key.clone());
1090
1091            select! {
1092                event = cons_out1.recv() => {
1093                    panic!("fetch should wait for the initial peer set, got {event:?}");
1094                },
1095                _ = context.sleep(Duration::from_millis(200)) => {},
1096            };
1097
1098            oracle
1099                .manager()
1100                .track(0, Set::try_from(peers.clone()).unwrap());
1101
1102            let (key_actual, value) = cons_out1.recv().await.unwrap();
1103            assert_eq!(key_actual, key);
1104            assert_eq!(value, Bytes::from("data for key 2"));
1105        });
1106    }
1107
1108    /// Tests that concurrent fetches are handled correctly.
1109    /// Also tests that the peer can recover from having no peers available.
1110    /// Also tests that the peer can get data from multiple peers that have different sets of data.
1111    #[test_traced]
1112    fn test_concurrent_fetch_requests() {
1113        let executor = deterministic::Runner::default();
1114        executor.start(|context| async move {
1115            let (mut oracle, mut schemes, peers, mut connections) =
1116                setup_network_and_peers(&context, &[1, 2, 3]).await;
1117
1118            let key2 = Key(2);
1119            let key3 = Key(3);
1120            let mut prod2 = Producer::default();
1121            prod2.insert(key2.clone(), Bytes::from("data for key 2"));
1122            let mut prod3 = Producer::default();
1123            prod3.insert(key3.clone(), Bytes::from("data for key 3"));
1124
1125            let (cons1, mut cons_out1) = consumer();
1126
1127            let scheme = schemes.remove(0);
1128            let mut mailbox1 = setup_and_spawn_actor(
1129                &context,
1130                oracle.manager(),
1131                oracle.control(scheme.public_key()),
1132                scheme,
1133                connections.remove(0),
1134                cons1,
1135                Producer::default(),
1136            );
1137
1138            let scheme = schemes.remove(0);
1139            let _mailbox2 = setup_and_spawn_actor(
1140                &context,
1141                oracle.manager(),
1142                oracle.control(scheme.public_key()),
1143                scheme,
1144                connections.remove(0),
1145                dummy_consumer(),
1146                prod2,
1147            );
1148
1149            let scheme = schemes.remove(0);
1150            let _mailbox3 = setup_and_spawn_actor(
1151                &context,
1152                oracle.manager(),
1153                oracle.control(scheme.public_key()),
1154                scheme,
1155                connections.remove(0),
1156                dummy_consumer(),
1157                prod3,
1158            );
1159
1160            // Add choppy links between the requester and the two producers
1161            add_link(&mut oracle, LINK_UNRELIABLE.clone(), &peers, 0, 1).await;
1162            add_link(&mut oracle, LINK_UNRELIABLE.clone(), &peers, 0, 2).await;
1163
1164            // Run the fetches multiple times to ensure that the peer tries both of its peers
1165            for _ in 0..10 {
1166                // Initiate concurrent fetches.
1167                mailbox1.fetch(key2.clone());
1168                mailbox1.fetch(key3.clone());
1169
1170                // Collect both events without assuming order
1171                let mut events = Vec::new();
1172                events.push(cons_out1.recv().await.expect("Consumer channel closed"));
1173                events.push(cons_out1.recv().await.expect("Consumer channel closed"));
1174
1175                // Check that both keys were successfully fetched
1176                let mut found_key2 = false;
1177                let mut found_key3 = false;
1178                for (key_actual, value) in events {
1179                    if key_actual == key2 {
1180                        assert_eq!(value, Bytes::from("data for key 2"));
1181                        found_key2 = true;
1182                    } else if key_actual == key3 {
1183                        assert_eq!(value, Bytes::from("data for key 3"));
1184                        found_key3 = true;
1185                    } else {
1186                        panic!("Unexpected key received");
1187                    }
1188                }
1189                assert!(found_key2 && found_key3,);
1190            }
1191        });
1192    }
1193
1194    /// Tests that pruning an inactive fetch has no effect.
1195    /// Prunes a key before, after, and during the fetch process.
1196    #[test_traced]
1197    fn test_retain_drops_key() {
1198        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1199        executor.start(|context| async move {
1200            let (mut oracle, mut schemes, peers, mut connections) =
1201                setup_network_and_peers(&context, &[1, 2]).await;
1202
1203            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1204
1205            let key = Key(6);
1206            let mut prod2 = Producer::default();
1207            prod2.insert(key.clone(), Bytes::from("data for key 6"));
1208
1209            let (cons1, mut cons_out1) = consumer();
1210
1211            let scheme = schemes.remove(0);
1212            let mut mailbox1 = setup_and_spawn_actor(
1213                &context,
1214                oracle.manager(),
1215                oracle.control(scheme.public_key()),
1216                scheme,
1217                connections.remove(0),
1218                cons1,
1219                Producer::default(),
1220            );
1221
1222            let scheme = schemes.remove(0);
1223            let _mailbox2 = setup_and_spawn_actor(
1224                &context,
1225                oracle.manager(),
1226                oracle.control(scheme.public_key()),
1227                scheme,
1228                connections.remove(0),
1229                dummy_consumer(),
1230                prod2,
1231            );
1232
1233            // Prune before sending the fetch, expecting no effect.
1234            let canceled = key.clone();
1235            mailbox1.retain(move |key, _| key != &canceled);
1236            select! {
1237                _ = cons_out1.recv() => {
1238                    panic!("unexpected event");
1239                },
1240                _ = context.sleep(Duration::from_millis(100)) => {},
1241            };
1242
1243            // Initiate fetch and wait for data to be delivered
1244            mailbox1.fetch(key.clone());
1245            let (key_actual, value) = cons_out1.recv().await.unwrap();
1246            assert_eq!(key_actual, key);
1247            assert_eq!(value, Bytes::from("data for key 6"));
1248
1249            // Attempt to prune after data has been delivered, expecting no effect
1250            let canceled = key.clone();
1251            mailbox1.retain(move |key, _| key != &canceled);
1252            select! {
1253                _ = cons_out1.recv() => {
1254                    panic!("unexpected event");
1255                },
1256                _ = context.sleep(Duration::from_millis(100)) => {},
1257            };
1258
1259            // Initiate and prune another fetch.
1260            let key = Key(7);
1261            mailbox1.fetch(key.clone());
1262            let canceled = key.clone();
1263            mailbox1.retain(move |key, _| key != &canceled);
1264
1265            // No event should arrive after pruning.
1266            select! {
1267                _ = cons_out1.recv() => panic!("unexpected event"),
1268                _ = context.sleep(Duration::from_millis(100)) => {},
1269            };
1270        });
1271    }
1272
1273    /// Tests that a peer is blocked after delivering invalid data,
1274    /// preventing further fetches from that peer.
1275    #[test_traced]
1276    fn test_blocking_peer() {
1277        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1278        executor.start(|context| async move {
1279            let (mut oracle, mut schemes, peers, mut connections) =
1280                setup_network_and_peers(&context, &[1, 2, 3]).await;
1281
1282            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1283            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
1284            add_link(&mut oracle, LINK.clone(), &peers, 1, 2).await;
1285
1286            let key_a = Key(1);
1287            let key_b = Key(2);
1288            let invalid_data_a = Bytes::from("invalid for A");
1289            let valid_data_a = Bytes::from("valid for A");
1290            let valid_data_b = Bytes::from("valid for B");
1291
1292            // Set up producers
1293            let mut prod2 = Producer::default();
1294            prod2.insert(key_a.clone(), invalid_data_a.clone());
1295            prod2.insert(key_b.clone(), valid_data_b.clone());
1296
1297            let mut prod3 = Producer::default();
1298            prod3.insert(key_a.clone(), valid_data_a.clone());
1299
1300            // Set up consumer for Peer1 with expected values
1301            let (mut cons1, mut cons_out1) = consumer();
1302            cons1.add_expected(key_a.clone(), valid_data_a.clone());
1303            cons1.add_expected(key_b.clone(), valid_data_b.clone());
1304
1305            // Spawn actors
1306            let scheme = schemes.remove(0);
1307            let mut mailbox1 = setup_and_spawn_actor(
1308                &context,
1309                oracle.manager(),
1310                oracle.control(scheme.public_key()),
1311                scheme,
1312                connections.remove(0),
1313                cons1,
1314                Producer::default(),
1315            );
1316
1317            let scheme = schemes.remove(0);
1318            let _mailbox2 = setup_and_spawn_actor(
1319                &context,
1320                oracle.manager(),
1321                oracle.control(scheme.public_key()),
1322                scheme,
1323                connections.remove(0),
1324                dummy_consumer(),
1325                prod2,
1326            );
1327
1328            let scheme = schemes.remove(0);
1329            let _mailbox3 = setup_and_spawn_actor(
1330                &context,
1331                oracle.manager(),
1332                oracle.control(scheme.public_key()),
1333                scheme,
1334                connections.remove(0),
1335                dummy_consumer(),
1336                prod3,
1337            );
1338
1339            // Fetch keyA multiple times to ensure that Peer2 is blocked.
1340            for _ in 0..20 {
1341                // Fetch keyA
1342                mailbox1.fetch(key_a.clone());
1343
1344                // Wait for success event for keyA
1345                let (key_actual, value) = cons_out1.recv().await.unwrap();
1346                assert_eq!(key_actual, key_a);
1347                assert_eq!(value, valid_data_a);
1348            }
1349
1350            // Fetch keyB
1351            mailbox1.fetch(key_b.clone());
1352
1353            // Wait for some time (longer than retry timeout)
1354            context.sleep(Duration::from_secs(5)).await;
1355
1356            // No success event should be received for keyB since the only peer with valid data is blocked
1357            select! {
1358                _ = cons_out1.recv() => panic!("unexpected event"),
1359                _ = context.sleep(Duration::from_millis(100)) => {},
1360            };
1361
1362            // Prune the fetch for keyB.
1363            let canceled = key_b.clone();
1364            mailbox1.retain(move |key, _| key != &canceled);
1365
1366            // Check oracle
1367            let blocked = oracle.blocked().await.unwrap();
1368            assert_eq!(blocked.len(), 1);
1369            assert_eq!(blocked[0].0, peers[0]);
1370            assert_eq!(blocked[0].1, peers[1]);
1371        });
1372    }
1373
1374    /// Tests that duplicate fetches for the same key are handled properly.
1375    /// The test verifies that when the same key is fetched multiple times,
1376    /// the data is correctly delivered once without errors.
1377    #[test_traced]
1378    fn test_duplicate_fetch_key() {
1379        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1380        executor.start(|context| async move {
1381            let (mut oracle, mut schemes, peers, mut connections) =
1382                setup_network_and_peers(&context, &[1, 2]).await;
1383
1384            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1385
1386            let key = Key(5);
1387            let mut prod2 = Producer::default();
1388            prod2.insert(key.clone(), Bytes::from("data for key 5"));
1389
1390            let (cons1, mut cons_out1) = consumer();
1391
1392            let scheme = schemes.remove(0);
1393            let mut mailbox1 = setup_and_spawn_actor(
1394                &context,
1395                oracle.manager(),
1396                oracle.control(scheme.public_key()),
1397                scheme,
1398                connections.remove(0),
1399                cons1,
1400                Producer::default(),
1401            );
1402
1403            let scheme = schemes.remove(0);
1404            let _mailbox2 = setup_and_spawn_actor(
1405                &context,
1406                oracle.manager(),
1407                oracle.control(scheme.public_key()),
1408                scheme,
1409                connections.remove(0),
1410                dummy_consumer(),
1411                prod2,
1412            );
1413
1414            // Send duplicate fetches for the same key.
1415            mailbox1.fetch(key.clone());
1416            mailbox1.fetch(key.clone());
1417
1418            // Should receive the data only once
1419            let (key_actual, value) = cons_out1.recv().await.unwrap();
1420            assert_eq!(key_actual, key);
1421            assert_eq!(value, Bytes::from("data for key 5"));
1422
1423            // Make sure we don't receive a second event for the duplicate fetch
1424            select! {
1425                _ = cons_out1.recv() => {
1426                    panic!("Unexpected second event received for duplicate fetch");
1427                },
1428                _ = context.sleep(Duration::from_millis(500)) => {
1429                    // This is expected - no additional events should be produced
1430                },
1431            };
1432        });
1433    }
1434
1435    /// Tests that changing peer sets is handled correctly using the update channel.
1436    /// This test verifies that when the peer set changes from peer A to peer B,
1437    /// the resolver correctly adapts and fetches from the new peer.
1438    #[test_traced]
1439    fn test_changing_peer_sets() {
1440        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1441        executor.start(|context| async move {
1442            let (mut oracle, mut schemes, peers, mut connections) =
1443                setup_network_and_peers(&context, &[1, 2, 3]).await;
1444
1445            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1446            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
1447
1448            let key1 = Key(1);
1449            let key2 = Key(2);
1450
1451            let mut prod2 = Producer::default();
1452            prod2.insert(key1.clone(), Bytes::from("data from peer 2"));
1453
1454            let mut prod3 = Producer::default();
1455            prod3.insert(key2.clone(), Bytes::from("data from peer 3"));
1456
1457            let (cons1, mut cons_out1) = consumer();
1458
1459            let scheme = schemes.remove(0);
1460            let mut mailbox1 = setup_and_spawn_actor(
1461                &context,
1462                oracle.manager(),
1463                oracle.control(scheme.public_key()),
1464                scheme,
1465                connections.remove(0),
1466                cons1,
1467                Producer::default(),
1468            );
1469
1470            let scheme = schemes.remove(0);
1471            let _mailbox2 = setup_and_spawn_actor(
1472                &context,
1473                oracle.manager(),
1474                oracle.control(scheme.public_key()),
1475                scheme,
1476                connections.remove(0),
1477                dummy_consumer(),
1478                prod2,
1479            );
1480
1481            // Fetch key1 from peer 2
1482            mailbox1.fetch(key1.clone());
1483
1484            // Wait for successful fetch
1485            let (key_actual, value) = cons_out1.recv().await.unwrap();
1486            assert_eq!(key_actual, key1);
1487            assert_eq!(value, Bytes::from("data from peer 2"));
1488
1489            // Change peer set to include peer 3
1490            let scheme = schemes.remove(0);
1491            let _mailbox3 = setup_and_spawn_actor(
1492                &context,
1493                oracle.manager(),
1494                oracle.control(scheme.public_key()),
1495                scheme,
1496                connections.remove(0),
1497                dummy_consumer(),
1498                prod3,
1499            );
1500
1501            // Need to wait for the peer set change to propagate
1502            context.sleep(Duration::from_millis(200)).await;
1503
1504            // Fetch key2 from peer 3
1505            mailbox1.fetch(key2.clone());
1506
1507            // Wait for successful fetch
1508            let (key_actual, value) = cons_out1.recv().await.unwrap();
1509            assert_eq!(key_actual, key2);
1510            assert_eq!(value, Bytes::from("data from peer 3"));
1511        });
1512    }
1513
1514    #[test_traced]
1515    fn test_fetch_targeted() {
1516        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1517        executor.start(|context| async move {
1518            let (mut oracle, mut schemes, peers, mut connections) =
1519                setup_network_and_peers(&context, &[1, 2, 3]).await;
1520
1521            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1522            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
1523
1524            let key = Key(1);
1525            let invalid_data = Bytes::from("invalid data");
1526            let valid_data = Bytes::from("valid data");
1527
1528            // Peer 2 has invalid data, peer 3 has valid data
1529            let mut prod2 = Producer::default();
1530            prod2.insert(key.clone(), invalid_data.clone());
1531
1532            let mut prod3 = Producer::default();
1533            prod3.insert(key.clone(), valid_data.clone());
1534
1535            // Consumer expects only valid_data
1536            let (mut cons1, mut cons_out1) = consumer();
1537            cons1.add_expected(key.clone(), valid_data.clone());
1538
1539            let scheme = schemes.remove(0);
1540            let mut mailbox1 = setup_and_spawn_actor(
1541                &context,
1542                oracle.manager(),
1543                oracle.control(scheme.public_key()),
1544                scheme,
1545                connections.remove(0),
1546                cons1,
1547                Producer::default(),
1548            );
1549
1550            let scheme = schemes.remove(0);
1551            let _mailbox2 = setup_and_spawn_actor(
1552                &context,
1553                oracle.manager(),
1554                oracle.control(scheme.public_key()),
1555                scheme,
1556                connections.remove(0),
1557                dummy_consumer(),
1558                prod2,
1559            );
1560
1561            let scheme = schemes.remove(0);
1562            let _mailbox3 = setup_and_spawn_actor(
1563                &context,
1564                oracle.manager(),
1565                oracle.control(scheme.public_key()),
1566                scheme,
1567                connections.remove(0),
1568                dummy_consumer(),
1569                prod3,
1570            );
1571
1572            // Wait for peer set to be established
1573            context.sleep(Duration::from_millis(100)).await;
1574
1575            // Start fetch with targets for both peer 2 (invalid data) and peer 3 (valid data)
1576            // When peer 2 returns invalid data, only peer 2 should be removed from targets
1577            // Peer 3 should still be tried as a target and succeed
1578            mailbox1.fetch_targeted(
1579                key.clone(),
1580                non_empty_vec![peers[1].clone(), peers[2].clone()],
1581            );
1582
1583            // Should eventually succeed from peer 3
1584            let (key_actual, value) = cons_out1.recv().await.unwrap();
1585            assert_eq!(key_actual, key);
1586            assert_eq!(value, valid_data);
1587
1588            // Verify peer 2 was blocked (sent invalid data)
1589            let blocked = oracle.blocked().await.unwrap();
1590            assert_eq!(blocked.len(), 1);
1591            assert_eq!(blocked[0].0, peers[0]);
1592            assert_eq!(blocked[0].1, peers[1]);
1593
1594            // Verify metrics: 1 successful fetch (from peer 3 after peer 2 was blocked)
1595            let metrics = context.encode();
1596            assert_eq!(
1597                status_metric_total(&metrics, "actor_fetch_total", "Success"),
1598                1
1599            );
1600        });
1601    }
1602
1603    #[test_traced]
1604    fn test_fetch_targeted_no_fallback() {
1605        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1606        executor.start(|context| async move {
1607            let (mut oracle, mut schemes, peers, mut connections) =
1608                setup_network_and_peers(&context, &[1, 2, 3, 4]).await;
1609
1610            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1611            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
1612            add_link(&mut oracle, LINK.clone(), &peers, 0, 3).await;
1613
1614            let key = Key(1);
1615
1616            // Only peer 4 has the data, peers 2 and 3 don't
1617            let mut prod4 = Producer::default();
1618            prod4.insert(key.clone(), Bytes::from("data from peer 4"));
1619
1620            let (cons1, mut cons_out1) = consumer();
1621
1622            let scheme = schemes.remove(0);
1623            let mut mailbox1 = setup_and_spawn_actor(
1624                &context,
1625                oracle.manager(),
1626                oracle.control(scheme.public_key()),
1627                scheme,
1628                connections.remove(0),
1629                cons1,
1630                Producer::default(),
1631            );
1632
1633            let scheme = schemes.remove(0);
1634            let _mailbox2 = setup_and_spawn_actor(
1635                &context,
1636                oracle.manager(),
1637                oracle.control(scheme.public_key()),
1638                scheme,
1639                connections.remove(0),
1640                dummy_consumer(),
1641                Producer::default(), // no data
1642            );
1643
1644            let scheme = schemes.remove(0);
1645            let _mailbox3 = setup_and_spawn_actor(
1646                &context,
1647                oracle.manager(),
1648                oracle.control(scheme.public_key()),
1649                scheme,
1650                connections.remove(0),
1651                dummy_consumer(),
1652                Producer::default(), // no data
1653            );
1654
1655            let scheme = schemes.remove(0);
1656            let _mailbox4 = setup_and_spawn_actor(
1657                &context,
1658                oracle.manager(),
1659                oracle.control(scheme.public_key()),
1660                scheme,
1661                connections.remove(0),
1662                dummy_consumer(),
1663                prod4,
1664            );
1665
1666            // Wait for peer set to be established
1667            context.sleep(Duration::from_millis(100)).await;
1668
1669            // Start fetch with targets for peers 2 and 3 (both don't have data)
1670            // Peer 4 has data but is NOT a target - it should NEVER be tried
1671            mailbox1.fetch_targeted(
1672                key.clone(),
1673                non_empty_vec![peers[1].clone(), peers[2].clone()],
1674            );
1675
1676            // Wait enough time for targets to fail and retry multiple times
1677            // The fetch should not succeed because peer 4 (which has data) is not targeted
1678            select! {
1679                event = cons_out1.recv() => {
1680                    panic!("Fetch should not succeed, but got: {event:?}");
1681                },
1682                _ = context.sleep(Duration::from_secs(3)) => {
1683                    // Expected: no success event because peer 4 is not targeted
1684                },
1685            };
1686        });
1687    }
1688
1689    #[test_traced]
1690    fn test_fetch_all_targeted() {
1691        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1692        executor.start(|context| async move {
1693            let (mut oracle, mut schemes, peers, mut connections) =
1694                setup_network_and_peers(&context, &[1, 2, 3, 4]).await;
1695
1696            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1697            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
1698            add_link(&mut oracle, LINK.clone(), &peers, 0, 3).await;
1699
1700            let key1 = Key(1);
1701            let key2 = Key(2);
1702            let key3 = Key(3);
1703
1704            // Peer 2 has key1
1705            let mut prod2 = Producer::default();
1706            prod2.insert(key1.clone(), Bytes::from("data for key 1"));
1707
1708            // Peer 3 has key3
1709            let mut prod3 = Producer::default();
1710            prod3.insert(key3.clone(), Bytes::from("data for key 3"));
1711
1712            // Peer 4 has key2
1713            let mut prod4 = Producer::default();
1714            prod4.insert(key2.clone(), Bytes::from("data for key 2"));
1715
1716            // Consumer expects all three keys
1717            let (mut cons1, mut cons_out1) = consumer();
1718            cons1.add_expected(key1.clone(), Bytes::from("data for key 1"));
1719            cons1.add_expected(key2.clone(), Bytes::from("data for key 2"));
1720            cons1.add_expected(key3.clone(), Bytes::from("data for key 3"));
1721
1722            let scheme = schemes.remove(0);
1723            let mut mailbox1 = setup_and_spawn_actor(
1724                &context,
1725                oracle.manager(),
1726                oracle.control(scheme.public_key()),
1727                scheme,
1728                connections.remove(0),
1729                cons1,
1730                Producer::default(),
1731            );
1732
1733            let scheme = schemes.remove(0);
1734            let _mailbox2 = setup_and_spawn_actor(
1735                &context,
1736                oracle.manager(),
1737                oracle.control(scheme.public_key()),
1738                scheme,
1739                connections.remove(0),
1740                dummy_consumer(),
1741                prod2,
1742            );
1743
1744            let scheme = schemes.remove(0);
1745            let _mailbox3 = setup_and_spawn_actor(
1746                &context,
1747                oracle.manager(),
1748                oracle.control(scheme.public_key()),
1749                scheme,
1750                connections.remove(0),
1751                dummy_consumer(),
1752                prod3,
1753            );
1754
1755            let scheme = schemes.remove(0);
1756            let _mailbox4 = setup_and_spawn_actor(
1757                &context,
1758                oracle.manager(),
1759                oracle.control(scheme.public_key()),
1760                scheme,
1761                connections.remove(0),
1762                dummy_consumer(),
1763                prod4,
1764            );
1765
1766            // Wait for peer set to be established
1767            context.sleep(Duration::from_millis(100)).await;
1768
1769            // Fetch keys with mixed targeting:
1770            // - key1 targeted to peer 2 (has data) -> should succeed from target
1771            // - key2 targeted to peer 4 (has data) -> should succeed from target
1772            // - key3 no targeting -> fetched from any peer (peer 3 has it)
1773            mailbox1.fetch_all_targeted(vec![
1774                (key1.clone(), non_empty_vec![peers[1].clone()]), // peer 2 has key1
1775                (key2.clone(), non_empty_vec![peers[3].clone()]), // peer 4 has key2
1776            ]);
1777            mailbox1.fetch(key3.clone()); // no targeting for key3
1778
1779            // Collect all three events
1780            let mut results = HashMap::new();
1781            for _ in 0..3 {
1782                let (key, value) = cons_out1.recv().await.unwrap();
1783                results.insert(key, value);
1784            }
1785
1786            // Verify all keys received correct data
1787            assert_eq!(results.len(), 3);
1788            assert_eq!(results.get(&key1).unwrap(), &Bytes::from("data for key 1"));
1789            assert_eq!(results.get(&key2).unwrap(), &Bytes::from("data for key 2"));
1790            assert_eq!(results.get(&key3).unwrap(), &Bytes::from("data for key 3"));
1791
1792            // Verify metrics: 3 successful fetches
1793            let metrics = context.encode();
1794            assert_eq!(
1795                status_metric_total(&metrics, "actor_fetch_total", "Success"),
1796                3
1797            );
1798        });
1799    }
1800
1801    /// Tests that calling fetch() on an in-progress targeted fetch clears the targets,
1802    /// allowing the fetch to succeed from any available peer.
1803    #[test_traced]
1804    fn test_fetch_clears_targets() {
1805        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1806        executor.start(|context| async move {
1807            let (mut oracle, mut schemes, peers, mut connections) =
1808                setup_network_and_peers(&context, &[1, 2, 3]).await;
1809
1810            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1811            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
1812
1813            let key = Key(1);
1814            let valid_data = Bytes::from("valid data");
1815
1816            // Peer 2 has no data, peer 3 has the data
1817            let mut prod3 = Producer::default();
1818            prod3.insert(key.clone(), valid_data.clone());
1819
1820            let (cons1, mut cons_out1) = consumer();
1821
1822            let scheme = schemes.remove(0);
1823            let mut mailbox1 = setup_and_spawn_actor(
1824                &context,
1825                oracle.manager(),
1826                oracle.control(scheme.public_key()),
1827                scheme,
1828                connections.remove(0),
1829                cons1,
1830                Producer::default(),
1831            );
1832
1833            let scheme = schemes.remove(0);
1834            let _mailbox2 = setup_and_spawn_actor(
1835                &context,
1836                oracle.manager(),
1837                oracle.control(scheme.public_key()),
1838                scheme,
1839                connections.remove(0),
1840                dummy_consumer(),
1841                Producer::default(), // no data
1842            );
1843
1844            let scheme = schemes.remove(0);
1845            let _mailbox3 = setup_and_spawn_actor(
1846                &context,
1847                oracle.manager(),
1848                oracle.control(scheme.public_key()),
1849                scheme,
1850                connections.remove(0),
1851                dummy_consumer(),
1852                prod3,
1853            );
1854
1855            // Wait for peer set to be established
1856            context.sleep(Duration::from_millis(100)).await;
1857
1858            // Start fetch with target for peer 2 only (who doesn't have data)
1859            mailbox1.fetch_targeted(key.clone(), non_empty_vec![peers[1].clone()]);
1860
1861            // Wait for the targeted fetch to fail a few times
1862            context.sleep(Duration::from_millis(500)).await;
1863
1864            // Call fetch() which should clear the targets and allow fallback to any peer
1865            mailbox1.fetch(key.clone());
1866
1867            // Should now succeed from peer 3 (who has data but wasn't originally targeted)
1868            let (key_actual, value) = cons_out1.recv().await.unwrap();
1869            assert_eq!(key_actual, key);
1870            assert_eq!(value, valid_data);
1871        });
1872    }
1873
1874    #[test_traced]
1875    fn test_fetch_targeted_does_not_restrict_all() {
1876        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1877        executor.start(|context| async move {
1878            let (mut oracle, mut schemes, peers, mut connections) =
1879                setup_network_and_peers(&context, &[1, 2, 3]).await;
1880
1881            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1882            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
1883
1884            let key = Key(1);
1885            let valid_data = Bytes::from("valid data");
1886
1887            // Peer 2 has no data, peer 3 has the data
1888            let mut prod3 = Producer::default();
1889            prod3.insert(key.clone(), valid_data.clone());
1890
1891            let (cons1, mut cons_out1) = consumer();
1892
1893            let scheme = schemes.remove(0);
1894            let mut mailbox1 = setup_and_spawn_actor(
1895                &context,
1896                oracle.manager(),
1897                oracle.control(scheme.public_key()),
1898                scheme,
1899                connections.remove(0),
1900                cons1,
1901                Producer::default(),
1902            );
1903
1904            let scheme = schemes.remove(0);
1905            let _mailbox2 = setup_and_spawn_actor(
1906                &context,
1907                oracle.manager(),
1908                oracle.control(scheme.public_key()),
1909                scheme,
1910                connections.remove(0),
1911                dummy_consumer(),
1912                Producer::default(), // no data
1913            );
1914
1915            let scheme = schemes.remove(0);
1916            let _mailbox3 = setup_and_spawn_actor(
1917                &context,
1918                oracle.manager(),
1919                oracle.control(scheme.public_key()),
1920                scheme,
1921                connections.remove(0),
1922                dummy_consumer(),
1923                prod3,
1924            );
1925
1926            // Wait for peer set to be established
1927            context.sleep(Duration::from_millis(100)).await;
1928
1929            // Start fetch without targets (can try any peer)
1930            mailbox1.fetch(key.clone());
1931
1932            // Wait a bit for the fetch to start
1933            context.sleep(Duration::from_millis(50)).await;
1934
1935            // Call fetch_targeted with peer 2 only (who doesn't have data)
1936            // This should NOT restrict the existing "all" fetch
1937            mailbox1.fetch_targeted(key.clone(), non_empty_vec![peers[1].clone()]);
1938
1939            // Should still succeed from peer 3 (who has data but wasn't in the targeted call)
1940            // because the original fetch was "all" and shouldn't be restricted
1941            let (key_actual, value) = cons_out1.recv().await.unwrap();
1942            assert_eq!(key_actual, key);
1943            assert_eq!(value, valid_data);
1944        });
1945    }
1946
1947    #[test_traced]
1948    fn test_retain() {
1949        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1950        executor.start(|context| async move {
1951            let (mut oracle, mut schemes, peers, mut connections) =
1952                setup_network_and_peers(&context, &[1, 2]).await;
1953
1954            let key = Key(5);
1955            let mut prod2 = Producer::default();
1956            prod2.insert(key.clone(), Bytes::from("data for key 5"));
1957
1958            let (cons1, mut cons_out1) = consumer();
1959
1960            let scheme = schemes.remove(0);
1961            let mut mailbox1 = setup_and_spawn_actor(
1962                &context,
1963                oracle.manager(),
1964                oracle.control(scheme.public_key()),
1965                scheme,
1966                connections.remove(0),
1967                cons1,
1968                Producer::default(),
1969            );
1970
1971            let scheme = schemes.remove(0);
1972            let _mailbox2 = setup_and_spawn_actor(
1973                &context,
1974                oracle.manager(),
1975                oracle.control(scheme.public_key()),
1976                scheme,
1977                connections.remove(0),
1978                dummy_consumer(),
1979                prod2,
1980            );
1981
1982            // Retain before fetching should have no effect
1983            mailbox1.retain(|_, _| true);
1984            select! {
1985                _ = cons_out1.recv() => {
1986                    panic!("unexpected event");
1987                },
1988                _ = context.sleep(Duration::from_millis(100)) => {},
1989            };
1990
1991            // Start a fetch (no link, so fetch stays in-flight)
1992            mailbox1.fetch(key.clone());
1993
1994            // Retain with predicate that excludes the key. This must clean up
1995            // the in-flight entry for the key.
1996            let key_clone = key.clone();
1997            mailbox1.retain(move |key, _| key != &key_clone);
1998
1999            // Now add link so fetches can complete
2000            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2001
2002            // Fetch same key again, if the in-flight entry wasn't cleaned up, this would
2003            // be treated as a duplicate and silently ignored
2004            mailbox1.fetch(key.clone());
2005
2006            // Should succeed
2007            let (key_actual, value) = cons_out1.recv().await.unwrap();
2008            assert_eq!(key_actual, key);
2009            assert_eq!(value, Bytes::from("data for key 5"));
2010        });
2011    }
2012
2013    #[test_traced]
2014    fn test_retain_uses_subscribers() {
2015        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2016        executor.start(|context| async move {
2017            let (mut oracle, mut schemes, peers, mut connections) =
2018                setup_network_and_peers(&context, &[1, 2]).await;
2019
2020            let key = Key(5);
2021            let mut prod2 = Producer::default();
2022            prod2.insert(key.clone(), Bytes::from("data for key 5"));
2023
2024            let (cons1, mut cons_out1): (Consumer<Key, Bytes, SubscriberTag>, _) = Consumer::new();
2025
2026            let scheme = schemes.remove(0);
2027            let mut mailbox1 = setup_and_spawn_actor(
2028                &context,
2029                oracle.manager(),
2030                oracle.control(scheme.public_key()),
2031                scheme,
2032                connections.remove(0),
2033                cons1,
2034                Producer::default(),
2035            );
2036
2037            let scheme = schemes.remove(0);
2038            let _mailbox2 = setup_and_spawn_actor(
2039                &context,
2040                oracle.manager(),
2041                oracle.control(scheme.public_key()),
2042                scheme,
2043                connections.remove(0),
2044                dummy_consumer(),
2045                prod2,
2046            );
2047
2048            let dropped_subscriber = SubscriberTag(50);
2049            let kept_subscriber = SubscriberTag(51);
2050            mailbox1.fetch(Fetch {
2051                key: key.clone(),
2052                subscriber: dropped_subscriber,
2053                span: tracing::Span::none(),
2054            });
2055            mailbox1.fetch(Fetch {
2056                key: key.clone(),
2057                subscriber: kept_subscriber.clone(),
2058                span: tracing::Span::none(),
2059            });
2060
2061            context.sleep(Duration::from_millis(100)).await;
2062            mailbox1.retain(move |_, subscriber| subscriber == &kept_subscriber);
2063            context.sleep(Duration::from_millis(100)).await;
2064
2065            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2066
2067            let (key_actual, value) = cons_out1.recv().await.unwrap();
2068            assert_eq!(key_actual, key);
2069            assert_eq!(value, Bytes::from("data for key 5"));
2070        });
2071    }
2072
2073    #[test_traced]
2074    fn test_deliver_receives_subscribers() {
2075        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2076        executor.start(|context| async move {
2077            let (mut oracle, mut schemes, peers, mut connections) =
2078                setup_network_and_peers(&context, &[1, 2]).await;
2079
2080            let key = Key(5);
2081            let mut prod2 = Producer::default();
2082            prod2.insert(key.clone(), Bytes::from("data for key 5"));
2083
2084            let (cons1, mut cons_out1) = SubscriberRecordingConsumer::new();
2085
2086            let scheme = schemes.remove(0);
2087            let mut mailbox1 = setup_and_spawn_actor(
2088                &context,
2089                oracle.manager(),
2090                oracle.control(scheme.public_key()),
2091                scheme,
2092                connections.remove(0),
2093                cons1,
2094                Producer::default(),
2095            );
2096
2097            let scheme = schemes.remove(0);
2098            let _mailbox2 = setup_and_spawn_actor(
2099                &context,
2100                oracle.manager(),
2101                oracle.control(scheme.public_key()),
2102                scheme,
2103                connections.remove(0),
2104                dummy_consumer(),
2105                prod2,
2106            );
2107
2108            let first_subscriber = SubscriberTag(50);
2109            let second_subscriber = SubscriberTag(51);
2110            mailbox1.fetch(Fetch {
2111                key: key.clone(),
2112                subscriber: second_subscriber.clone(),
2113                span: tracing::Span::none(),
2114            });
2115            mailbox1.fetch(Fetch {
2116                key: key.clone(),
2117                subscriber: first_subscriber.clone(),
2118                span: tracing::Span::none(),
2119            });
2120
2121            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2122
2123            let (delivery, value) = cons_out1.recv().await.unwrap();
2124            assert_eq!(
2125                delivery,
2126                Delivery {
2127                    key,
2128                    subscribers: non_empty_vec![
2129                        (first_subscriber, tracing::Span::none()),
2130                        (second_subscriber, tracing::Span::none())
2131                    ],
2132                }
2133            );
2134            assert_eq!(value, Bytes::from("data for key 5"));
2135        });
2136    }
2137
2138    #[test_traced]
2139    fn test_deliver_receives_multiple_subscribers() {
2140        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2141        executor.start(|context| async move {
2142            let (mut oracle, mut schemes, peers, mut connections) =
2143                setup_network_and_peers(&context, &[1, 2]).await;
2144
2145            let key = Key(5);
2146            let mut prod2 = Producer::default();
2147            prod2.insert(key.clone(), Bytes::from("data for key 5"));
2148
2149            let (cons1, mut cons_out1) = SubscriberRecordingConsumer::new();
2150
2151            let scheme = schemes.remove(0);
2152            let mut mailbox1 = setup_and_spawn_actor(
2153                &context,
2154                oracle.manager(),
2155                oracle.control(scheme.public_key()),
2156                scheme,
2157                connections.remove(0),
2158                cons1,
2159                Producer::default(),
2160            );
2161
2162            let scheme = schemes.remove(0);
2163            let _mailbox2 = setup_and_spawn_actor(
2164                &context,
2165                oracle.manager(),
2166                oracle.control(scheme.public_key()),
2167                scheme,
2168                connections.remove(0),
2169                dummy_consumer(),
2170                prod2,
2171            );
2172
2173            let first_subscriber = SubscriberTag(49);
2174            let second_subscriber = SubscriberTag(50);
2175            mailbox1.fetch(Fetch {
2176                key: key.clone(),
2177                subscriber: first_subscriber.clone(),
2178                span: tracing::Span::none(),
2179            });
2180            mailbox1.fetch(Fetch {
2181                key: key.clone(),
2182                subscriber: second_subscriber.clone(),
2183                span: tracing::Span::none(),
2184            });
2185
2186            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2187
2188            let (delivery, value) = cons_out1.recv().await.unwrap();
2189            assert_eq!(
2190                delivery,
2191                Delivery {
2192                    key: key.clone(),
2193                    subscribers: non_empty_vec![
2194                        (first_subscriber, tracing::Span::none()),
2195                        (second_subscriber, tracing::Span::none())
2196                    ],
2197                }
2198            );
2199            assert_eq!(value, Bytes::from("data for key 5"));
2200        });
2201    }
2202
2203    #[test_traced]
2204    fn test_fetch_during_validation_reuses_response_after_success() {
2205        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2206        executor.start(|context| async move {
2207            let (mut oracle, mut schemes, peers, mut connections) =
2208                setup_network_and_peers(&context, &[1, 2]).await;
2209
2210            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2211
2212            let key = Key(5);
2213            let first_response = Bytes::from("data for key 5");
2214            let second_response = Bytes::from("refetched data for key 5");
2215            let mut prod2 = SequencedProducer::default();
2216            prod2.insert(
2217                key.clone(),
2218                [first_response.clone(), second_response.clone()],
2219            );
2220            let prod2_observer = prod2.clone();
2221
2222            let (first_gate_sender, first_gate_receiver) = oneshot::channel();
2223            let (second_gate_sender, second_gate_receiver) = oneshot::channel();
2224            let (cons1, mut deliveries, mut started) = BlockingSubscriberRecordingConsumer::new(
2225                context.child("consumer"),
2226                vec![(first_gate_receiver, true), (second_gate_receiver, true)],
2227            );
2228
2229            let scheme = schemes.remove(0);
2230            let mut mailbox1 = setup_and_spawn_actor(
2231                &context,
2232                oracle.manager(),
2233                oracle.control(scheme.public_key()),
2234                scheme,
2235                connections.remove(0),
2236                cons1,
2237                Producer::default(),
2238            );
2239
2240            let scheme = schemes.remove(0);
2241            let _mailbox2 = setup_and_spawn_actor_with_producer(
2242                &context,
2243                oracle.manager(),
2244                oracle.control(scheme.public_key()),
2245                scheme,
2246                connections.remove(0),
2247                dummy_consumer(),
2248                prod2,
2249            );
2250
2251            let first_subscriber = SubscriberTag(49);
2252            let second_subscriber = SubscriberTag(50);
2253            mailbox1.fetch(Fetch {
2254                key: key.clone(),
2255                subscriber: first_subscriber.clone(),
2256                span: tracing::Span::none(),
2257            });
2258
2259            let delivery = started.recv().await.expect("delivery did not start");
2260            assert_eq!(
2261                delivery,
2262                Delivery {
2263                    key: key.clone(),
2264                    subscribers: non_empty_vec![(first_subscriber.clone(), tracing::Span::none())],
2265                }
2266            );
2267
2268            mailbox1.fetch(Fetch {
2269                key: key.clone(),
2270                subscriber: second_subscriber.clone(),
2271                span: tracing::Span::none(),
2272            });
2273            context.sleep(Duration::from_millis(100)).await;
2274            assert_eq!(
2275                prod2_observer.remaining(&key),
2276                vec![second_response.clone()]
2277            );
2278
2279            first_gate_sender.send(()).unwrap();
2280            let (delivery, value) = deliveries.recv().await.expect("consumer channel closed");
2281            assert_eq!(
2282                delivery,
2283                Delivery {
2284                    key: key.clone(),
2285                    subscribers: non_empty_vec![(first_subscriber, tracing::Span::none())],
2286                }
2287            );
2288            assert_eq!(value, first_response);
2289
2290            let delivery = select! {
2291                delivery = started.recv() => delivery.expect("second delivery did not start"),
2292                _ = context.sleep(Duration::from_secs(2)) => {
2293                    panic!("late subscriber was not delivered");
2294                },
2295            };
2296            assert_eq!(
2297                delivery,
2298                Delivery {
2299                    key: key.clone(),
2300                    subscribers: non_empty_vec![(second_subscriber.clone(), tracing::Span::none())],
2301                }
2302            );
2303
2304            second_gate_sender.send(()).unwrap();
2305            let (delivery, value) = deliveries.recv().await.expect("consumer channel closed");
2306            assert_eq!(
2307                delivery,
2308                Delivery {
2309                    key: key.clone(),
2310                    subscribers: non_empty_vec![(second_subscriber, tracing::Span::none())],
2311                }
2312            );
2313            assert_eq!(value, first_response);
2314            assert_eq!(prod2_observer.remaining(&key), vec![second_response]);
2315        });
2316    }
2317
2318    #[test_traced]
2319    fn test_late_subscriber_delivery_ignores_unrelated_waiter() {
2320        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2321        executor.start(|context| async move {
2322            let (mut oracle, mut schemes, peers, mut connections) =
2323                setup_network_and_peers(&context, &[1, 2, 3]).await;
2324
2325            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2326            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
2327
2328            let blocked_key = Key(4);
2329            let waiting_key = Key(5);
2330            let main_key = Key(6);
2331            let data = Bytes::from("data for key 6");
2332
2333            let mut prod2 = Producer::default();
2334            prod2.insert(blocked_key.clone(), Bytes::from("bad data"));
2335
2336            let mut prod3 = Producer::default();
2337            prod3.insert(main_key.clone(), data.clone());
2338
2339            let (first_gate_sender, first_gate_receiver) = oneshot::channel();
2340            let (second_gate_sender, second_gate_receiver) = oneshot::channel();
2341            let (cons1, mut deliveries, mut started) = BlockingSubscriberRecordingConsumer::new(
2342                context.child("consumer"),
2343                vec![(first_gate_receiver, false), (second_gate_receiver, true)],
2344            );
2345
2346            let scheme = schemes.remove(0);
2347            let mut mailbox1 = setup_and_spawn_actor(
2348                &context,
2349                oracle.manager(),
2350                oracle.control(scheme.public_key()),
2351                scheme,
2352                connections.remove(0),
2353                cons1,
2354                Producer::default(),
2355            );
2356
2357            let scheme = schemes.remove(0);
2358            let _mailbox2 = setup_and_spawn_actor(
2359                &context,
2360                oracle.manager(),
2361                oracle.control(scheme.public_key()),
2362                scheme,
2363                connections.remove(0),
2364                dummy_consumer(),
2365                prod2,
2366            );
2367
2368            let scheme = schemes.remove(0);
2369            let _mailbox3 = setup_and_spawn_actor(
2370                &context,
2371                oracle.manager(),
2372                oracle.control(scheme.public_key()),
2373                scheme,
2374                connections.remove(0),
2375                dummy_consumer(),
2376                prod3,
2377            );
2378
2379            mailbox1.fetch(Fetch {
2380                key: blocked_key.clone(),
2381                subscriber: SubscriberTag(1),
2382                span: tracing::Span::none(),
2383            });
2384            started
2385                .recv()
2386                .await
2387                .expect("blocking delivery did not start");
2388            first_gate_sender.send(()).unwrap();
2389            wait_for_blocked(&context, &oracle, &peers[0], &peers[1]).await;
2390
2391            mailbox1.fetch_targeted(
2392                Fetch {
2393                    key: waiting_key,
2394                    subscriber: SubscriberTag(2),
2395                    span: tracing::Span::none(),
2396                },
2397                non_empty_vec![peers[1].clone()],
2398            );
2399            context.sleep(Duration::from_millis(100)).await;
2400
2401            let first_subscriber = SubscriberTag(3);
2402            let second_subscriber = SubscriberTag(4);
2403            mailbox1.fetch(Fetch {
2404                key: main_key.clone(),
2405                subscriber: first_subscriber.clone(),
2406                span: tracing::Span::none(),
2407            });
2408
2409            let delivery = started.recv().await.expect("delivery did not start");
2410            assert_eq!(
2411                delivery,
2412                Delivery {
2413                    key: main_key.clone(),
2414                    subscribers: non_empty_vec![(first_subscriber.clone(), tracing::Span::none())],
2415                }
2416            );
2417
2418            mailbox1.fetch(Fetch {
2419                key: main_key.clone(),
2420                subscriber: second_subscriber.clone(),
2421                span: tracing::Span::none(),
2422            });
2423            context.sleep(Duration::from_millis(100)).await;
2424
2425            second_gate_sender.send(()).unwrap();
2426            let (delivery, value) = deliveries.recv().await.expect("consumer channel closed");
2427            assert_eq!(
2428                delivery,
2429                Delivery {
2430                    key: main_key.clone(),
2431                    subscribers: non_empty_vec![(first_subscriber, tracing::Span::none())],
2432                }
2433            );
2434            assert_eq!(value, data);
2435
2436            let delivery = select! {
2437                delivery = started.recv() => delivery.expect("second delivery did not start"),
2438                _ = context.sleep(Duration::from_secs(2)) => {
2439                    panic!("late subscriber was not delivered while an unrelated waiter was armed");
2440                },
2441            };
2442            assert_eq!(
2443                delivery,
2444                Delivery {
2445                    key: main_key,
2446                    subscribers: non_empty_vec![(second_subscriber, tracing::Span::none())],
2447                }
2448            );
2449        });
2450    }
2451
2452    #[test_traced]
2453    fn test_deliver_receives_distinct_subscriber_type() {
2454        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2455        executor.start(|context| async move {
2456            let (mut oracle, mut schemes, peers, mut connections) =
2457                setup_network_and_peers(&context, &[1, 2]).await;
2458
2459            let key = Key(5);
2460            let mut prod2 = Producer::default();
2461            prod2.insert(key.clone(), Bytes::from("data for key 5"));
2462
2463            let (cons1, mut cons_out1) = SubscriberRecordingConsumer::new();
2464
2465            let scheme = schemes.remove(0);
2466            let mut mailbox1 = setup_and_spawn_actor(
2467                &context,
2468                oracle.manager(),
2469                oracle.control(scheme.public_key()),
2470                scheme,
2471                connections.remove(0),
2472                cons1,
2473                Producer::default(),
2474            );
2475
2476            let scheme = schemes.remove(0);
2477            let _mailbox2 = setup_and_spawn_actor(
2478                &context,
2479                oracle.manager(),
2480                oracle.control(scheme.public_key()),
2481                scheme,
2482                connections.remove(0),
2483                dummy_consumer(),
2484                prod2,
2485            );
2486
2487            let subscriber = SubscriberTag(50);
2488            let retained = subscriber.clone();
2489            mailbox1.fetch(Fetch {
2490                key: key.clone(),
2491                subscriber: subscriber.clone(),
2492                span: tracing::Span::none(),
2493            });
2494
2495            context.sleep(Duration::from_millis(100)).await;
2496            mailbox1.retain(move |_, subscriber| subscriber == &retained);
2497            context.sleep(Duration::from_millis(100)).await;
2498
2499            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2500
2501            let (delivery, value) = cons_out1.recv().await.unwrap();
2502            assert_eq!(
2503                delivery,
2504                Delivery {
2505                    key,
2506                    subscribers: non_empty_vec![(subscriber, tracing::Span::none())],
2507                }
2508            );
2509            assert_eq!(value, Bytes::from("data for key 5"));
2510        });
2511    }
2512
2513    #[test_traced]
2514    fn test_deliver_receives_single_subscriber() {
2515        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2516        executor.start(|context| async move {
2517            let (mut oracle, mut schemes, peers, mut connections) =
2518                setup_network_and_peers(&context, &[1, 2]).await;
2519
2520            let key = Key(5);
2521            let mut prod2 = Producer::default();
2522            prod2.insert(key.clone(), Bytes::from("data for key 5"));
2523
2524            let (cons1, mut cons_out1) = SubscriberRecordingConsumer::new();
2525
2526            let scheme = schemes.remove(0);
2527            let mut mailbox1 = setup_and_spawn_actor(
2528                &context,
2529                oracle.manager(),
2530                oracle.control(scheme.public_key()),
2531                scheme,
2532                connections.remove(0),
2533                cons1,
2534                Producer::default(),
2535            );
2536
2537            let scheme = schemes.remove(0);
2538            let _mailbox2 = setup_and_spawn_actor(
2539                &context,
2540                oracle.manager(),
2541                oracle.control(scheme.public_key()),
2542                scheme,
2543                connections.remove(0),
2544                dummy_consumer(),
2545                prod2,
2546            );
2547
2548            let subscriber = SubscriberTag(50);
2549            mailbox1.fetch(Fetch {
2550                key: key.clone(),
2551                subscriber: subscriber.clone(),
2552                span: tracing::Span::none(),
2553            });
2554            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2555
2556            let (delivery, value) = cons_out1.recv().await.unwrap();
2557            assert_eq!(
2558                delivery,
2559                Delivery {
2560                    key,
2561                    subscribers: non_empty_vec![(subscriber, tracing::Span::none())],
2562                }
2563            );
2564            assert_eq!(value, Bytes::from("data for key 5"));
2565        });
2566    }
2567
2568    #[test_traced]
2569    fn test_retain_drops_all() {
2570        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2571        executor.start(|context| async move {
2572            let (mut oracle, mut schemes, peers, mut connections) =
2573                setup_network_and_peers(&context, &[1, 2]).await;
2574
2575            // No link yet - fetch will stay in-flight
2576            let key = Key(6);
2577            let mut prod2 = Producer::default();
2578            prod2.insert(key.clone(), Bytes::from("data for key 6"));
2579
2580            let (cons1, mut cons_out1) = consumer();
2581
2582            let scheme = schemes.remove(0);
2583            let mut mailbox1 = setup_and_spawn_actor(
2584                &context,
2585                oracle.manager(),
2586                oracle.control(scheme.public_key()),
2587                scheme,
2588                connections.remove(0),
2589                cons1,
2590                Producer::default(),
2591            );
2592
2593            let scheme = schemes.remove(0);
2594            let _mailbox2 = setup_and_spawn_actor(
2595                &context,
2596                oracle.manager(),
2597                oracle.control(scheme.public_key()),
2598                scheme,
2599                connections.remove(0),
2600                dummy_consumer(),
2601                prod2,
2602            );
2603
2604            // Pruning before fetching should have no effect.
2605            mailbox1.retain(|_, _| false);
2606            select! {
2607                _ = cons_out1.recv() => {
2608                    panic!("unexpected event");
2609                },
2610                _ = context.sleep(Duration::from_millis(100)) => {},
2611            };
2612
2613            // Start a fetch (no link, so fetch stays in-flight)
2614            mailbox1.fetch(key.clone());
2615
2616            // Prune all fetches.
2617            mailbox1.retain(|_, _| false);
2618
2619            // Now add link so fetches can complete
2620            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2621
2622            // Fetch same key again, if the in-flight entry wasn't cleaned up, this would
2623            // be treated as a duplicate and silently ignored
2624            mailbox1.fetch(key.clone());
2625
2626            // Should succeed
2627            let (key_actual, value) = cons_out1.recv().await.unwrap();
2628            assert_eq!(key_actual, key);
2629            assert_eq!(value, Bytes::from("data for key 6"));
2630        });
2631    }
2632
2633    /// Tests that when a peer is rate-limited, the fetcher spills over to another peer.
2634    /// With 2 peers and rate limit of 1/sec each, 2 requests issued simultaneously should
2635    /// both complete immediately (one to each peer) without waiting for rate limit reset.
2636    #[test_traced]
2637    fn test_rate_limit_spillover() {
2638        let executor = deterministic::Runner::timed(Duration::from_secs(30));
2639        executor.start(|context| async move {
2640            // Use a very restrictive rate limit: 1 request per second per peer
2641            let (mut oracle, mut schemes, peers, mut connections) =
2642                setup_network_and_peers_with_rate_limit(
2643                    &context,
2644                    &[1, 2, 3],
2645                    Quota::per_second(NZU32!(1)),
2646                )
2647                .await;
2648
2649            // Add links between peer 1 and both peer 2 and peer 3
2650            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2651            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
2652
2653            // Both peer 2 and peer 3 have the same data
2654            let mut prod2 = Producer::default();
2655            let mut prod3 = Producer::default();
2656            prod2.insert(Key(0), Bytes::from("data for key 0"));
2657            prod2.insert(Key(1), Bytes::from("data for key 1"));
2658            prod3.insert(Key(0), Bytes::from("data for key 0"));
2659            prod3.insert(Key(1), Bytes::from("data for key 1"));
2660
2661            let (cons1, mut cons_out1) = consumer();
2662
2663            // Set up peer 1 (the requester)
2664            let scheme = schemes.remove(0);
2665            let mut mailbox1 = setup_and_spawn_actor(
2666                &context,
2667                oracle.manager(),
2668                oracle.control(scheme.public_key()),
2669                scheme,
2670                connections.remove(0),
2671                cons1,
2672                Producer::default(),
2673            );
2674
2675            // Set up peer 2 (has data)
2676            let scheme = schemes.remove(0);
2677            let _mailbox2 = setup_and_spawn_actor(
2678                &context,
2679                oracle.manager(),
2680                oracle.control(scheme.public_key()),
2681                scheme,
2682                connections.remove(0),
2683                dummy_consumer(),
2684                prod2,
2685            );
2686
2687            // Set up peer 3 (also has data)
2688            let scheme = schemes.remove(0);
2689            let _mailbox3 = setup_and_spawn_actor(
2690                &context,
2691                oracle.manager(),
2692                oracle.control(scheme.public_key()),
2693                scheme,
2694                connections.remove(0),
2695                dummy_consumer(),
2696                prod3,
2697            );
2698
2699            // Wait for peer set to be established
2700            context.sleep(Duration::from_millis(100)).await;
2701            let start = context.current();
2702
2703            // Issue 2 fetches rapidly.
2704            // With rate limit of 1/sec per peer and 2 peers, both should complete
2705            // immediately via spill-over (one request to each peer)
2706            mailbox1.fetch(Key(0));
2707            mailbox1.fetch(Key(1));
2708
2709            // Collect results
2710            let mut results = HashMap::new();
2711            for _ in 0..2 {
2712                let (key, value) = cons_out1.recv().await.unwrap();
2713                results.insert(key.clone(), value);
2714            }
2715
2716            // Verify both keys were fetched successfully
2717            assert_eq!(results.len(), 2);
2718            assert_eq!(
2719                results.get(&Key(0)).unwrap(),
2720                &Bytes::from("data for key 0")
2721            );
2722            assert_eq!(
2723                results.get(&Key(1)).unwrap(),
2724                &Bytes::from("data for key 1")
2725            );
2726
2727            // Verify it completed quickly (well under 1 second) - proves spill-over worked
2728            // Without spill-over, the second request would wait ~1 second for rate limit reset
2729            let elapsed = context.current().duration_since(start).unwrap();
2730            assert!(
2731                elapsed < Duration::from_millis(500),
2732                "Expected quick completion via spill-over, but took {elapsed:?}"
2733            );
2734        });
2735    }
2736
2737    /// Tests that rate limiting causes retries to eventually succeed after the rate limit resets.
2738    /// This test uses a single peer with a restrictive rate limit and verifies that
2739    /// fetches eventually complete after waiting for the rate limit to reset.
2740    #[test_traced]
2741    fn test_rate_limit_retry_after_reset() {
2742        let executor = deterministic::Runner::timed(Duration::from_secs(30));
2743        executor.start(|context| async move {
2744            // Use a restrictive rate limit: 1 request per second
2745            let (mut oracle, mut schemes, peers, mut connections) =
2746                setup_network_and_peers_with_rate_limit(
2747                    &context,
2748                    &[1, 2],
2749                    Quota::per_second(NZU32!(1)),
2750                )
2751                .await;
2752
2753            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2754
2755            // Peer 2 has data for multiple keys
2756            let mut prod2 = Producer::default();
2757            prod2.insert(Key(1), Bytes::from("data for key 1"));
2758            prod2.insert(Key(2), Bytes::from("data for key 2"));
2759            prod2.insert(Key(3), Bytes::from("data for key 3"));
2760
2761            let (cons1, mut cons_out1) = consumer();
2762
2763            let scheme = schemes.remove(0);
2764            let mut mailbox1 = setup_and_spawn_actor(
2765                &context,
2766                oracle.manager(),
2767                oracle.control(scheme.public_key()),
2768                scheme,
2769                connections.remove(0),
2770                cons1,
2771                Producer::default(),
2772            );
2773
2774            let scheme = schemes.remove(0);
2775            let _mailbox2 = setup_and_spawn_actor(
2776                &context,
2777                oracle.manager(),
2778                oracle.control(scheme.public_key()),
2779                scheme,
2780                connections.remove(0),
2781                dummy_consumer(),
2782                prod2,
2783            );
2784
2785            // Wait for peer set to be established
2786            context.sleep(Duration::from_millis(100)).await;
2787            let start = context.current();
2788
2789            // Issue 3 fetches to a single peer with rate limit of 1/sec.
2790            // Only 1 can be sent immediately, the others must wait for rate limit reset
2791            mailbox1.fetch(Key(1));
2792            mailbox1.fetch(Key(2));
2793            mailbox1.fetch(Key(3));
2794
2795            // All 3 should eventually succeed (after rate limit resets)
2796            let mut results = HashMap::new();
2797            for _ in 0..3 {
2798                let (key, value) = cons_out1.recv().await.unwrap();
2799                results.insert(key.clone(), value);
2800            }
2801
2802            assert_eq!(results.len(), 3);
2803            for i in 1..=3 {
2804                assert_eq!(
2805                    results.get(&Key(i)).unwrap(),
2806                    &Bytes::from(format!("data for key {}", i))
2807                );
2808            }
2809
2810            // Verify it took significant time due to rate limiting
2811            // With 3 requests at 1/sec to a single peer, requests 2 and 3 must wait
2812            // for rate limit resets (~1 second each), so total should be > 2 seconds
2813            let elapsed = context.current().duration_since(start).unwrap();
2814            assert!(
2815                elapsed > Duration::from_secs(2),
2816                "Expected rate limiting to cause delay > 2s, but took {elapsed:?}"
2817            );
2818        });
2819    }
2820
2821    /// Tests that the resolver never sends fetches to itself (me exclusion).
2822    /// Even when the local peer has the data in its producer, it should fetch from
2823    /// another peer instead.
2824    #[test_traced]
2825    fn test_self_exclusion() {
2826        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2827        executor.start(|context| async move {
2828            let (mut oracle, mut schemes, peers, mut connections) =
2829                setup_network_and_peers(&context, &[1, 2]).await;
2830
2831            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2832
2833            let key = Key(1);
2834            let data = Bytes::from("shared data");
2835
2836            // Both peers have the data - peer 1 (requester) and peer 2
2837            let mut prod1 = Producer::default();
2838            prod1.insert(key.clone(), data.clone());
2839            let mut prod2 = Producer::default();
2840            prod2.insert(key.clone(), data.clone());
2841
2842            let (cons1, mut cons_out1) = consumer();
2843
2844            // Set up peer 1 with `me` set - it has the data but should NOT fetch from itself
2845            let scheme = schemes.remove(0);
2846            let mut mailbox1 = setup_and_spawn_actor(
2847                &context,
2848                oracle.manager(),
2849                oracle.control(scheme.public_key()),
2850                scheme,
2851                connections.remove(0),
2852                cons1,
2853                prod1, // peer 1 has the data
2854            );
2855
2856            // Set up peer 2 - also has the data
2857            let scheme = schemes.remove(0);
2858            let _mailbox2 = setup_and_spawn_actor(
2859                &context,
2860                oracle.manager(),
2861                oracle.control(scheme.public_key()),
2862                scheme,
2863                connections.remove(0),
2864                dummy_consumer(),
2865                prod2,
2866            );
2867
2868            // Wait for peer set to be established
2869            context.sleep(Duration::from_millis(100)).await;
2870
2871            // Fetch the key - should get it from peer 2, not from self
2872            mailbox1.fetch(key.clone());
2873
2874            // Should succeed (from peer 2)
2875            let (key_actual, value) = cons_out1.recv().await.unwrap();
2876            assert_eq!(key_actual, key);
2877            assert_eq!(value, data);
2878        });
2879    }
2880
2881    #[test_traced]
2882    fn test_fetch_uses_primary_peers_only() {
2883        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2884        executor.start(|context| async move {
2885            let (network, oracle) = Network::new(
2886                context.child("network"),
2887                commonware_p2p::simulated::Config {
2888                    max_size: 1024 * 1024,
2889                    disconnect_on_block: true,
2890                    tracked_peer_sets: NZUsize!(1),
2891                },
2892            );
2893            network.start();
2894
2895            let schemes: Vec<PrivateKey> = [1u64, 2, 3]
2896                .into_iter()
2897                .map(PrivateKey::from_seed)
2898                .collect();
2899            let peers: Vec<PublicKey> = schemes.iter().map(|s| s.public_key()).collect();
2900            let mut schemes = schemes;
2901
2902            let mut connections = Vec::new();
2903            for peer in &peers {
2904                let (sender, receiver) = oracle
2905                    .control(peer.clone())
2906                    .register(0, Quota::per_second(RATE_LIMIT))
2907                    .await
2908                    .unwrap();
2909                connections.push((sender, receiver));
2910            }
2911
2912            // Topology: peer 1 (requester) linked to peers 2 and 3.
2913            // Peer 2 is primary (no data), peer 3 is secondary (has data).
2914            // Fetch should only query primary peers, so the request must time out.
2915            let mut oracle = oracle;
2916            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2917            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
2918
2919            oracle.manager().track(
2920                1,
2921                TrackedPeers::new(
2922                    Set::try_from([peers[1].clone()]).unwrap(),
2923                    Set::try_from([peers[2].clone()]).unwrap(),
2924                ),
2925            );
2926            context.sleep(Duration::from_millis(100)).await;
2927
2928            let key = Key(1);
2929            let data = Bytes::from("secondary only data");
2930
2931            let (cons1, mut cons_out1) = consumer();
2932
2933            // Peer 1: the requester, has no data.
2934            let scheme = schemes.remove(0);
2935            let mut mailbox1 = setup_and_spawn_actor(
2936                &context,
2937                oracle.manager(),
2938                oracle.control(scheme.public_key()),
2939                scheme,
2940                connections.remove(0),
2941                cons1,
2942                Producer::default(),
2943            );
2944
2945            // Peer 2: primary, has no data.
2946            let scheme = schemes.remove(0);
2947            let _mailbox2 = setup_and_spawn_actor(
2948                &context,
2949                oracle.manager(),
2950                oracle.control(scheme.public_key()),
2951                scheme,
2952                connections.remove(0),
2953                dummy_consumer(),
2954                Producer::default(),
2955            );
2956
2957            // Peer 3: secondary, has the data. Should not be queried.
2958            let mut prod3 = Producer::default();
2959            prod3.insert(key.clone(), data);
2960            let scheme = schemes.remove(0);
2961            let _mailbox3 = setup_and_spawn_actor(
2962                &context,
2963                oracle.manager(),
2964                oracle.control(scheme.public_key()),
2965                scheme,
2966                connections.remove(0),
2967                dummy_consumer(),
2968                prod3,
2969            );
2970
2971            // Fetch should time out because the only peer with data (peer 3)
2972            // is secondary and won't be queried.
2973            mailbox1.fetch(key.clone());
2974
2975            select! {
2976                event = cons_out1.recv() => {
2977                    panic!("fetch should not succeed from a secondary peer, got: {event:?}");
2978                },
2979                _ = context.sleep(Duration::from_secs(2)) => {},
2980            }
2981        });
2982    }
2983
2984    #[test_traced]
2985    fn test_fetch_uses_latest_primary_set_only() {
2986        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2987        executor.start(|context| async move {
2988            let (network, oracle) = Network::new(
2989                context.child("network"),
2990                commonware_p2p::simulated::Config {
2991                    max_size: 1024 * 1024,
2992                    disconnect_on_block: true,
2993                    tracked_peer_sets: NZUsize!(2),
2994                },
2995            );
2996            network.start();
2997
2998            let schemes: Vec<PrivateKey> = [1u64, 2, 3]
2999                .into_iter()
3000                .map(PrivateKey::from_seed)
3001                .collect();
3002            let peers: Vec<PublicKey> = schemes.iter().map(|s| s.public_key()).collect();
3003            let mut schemes = schemes;
3004
3005            let mut connections = Vec::new();
3006            for peer in &peers {
3007                let (sender, receiver) = oracle
3008                    .control(peer.clone())
3009                    .register(0, Quota::per_second(RATE_LIMIT))
3010                    .await
3011                    .unwrap();
3012                connections.push((sender, receiver));
3013            }
3014
3015            let mut oracle = oracle;
3016            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3017            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
3018
3019            // Keep the requester tracked across the cutover so the fetch path itself remains
3020            // active, while peer 2 is retained only through the overlap window after peer 3
3021            // becomes the newest primary set.
3022            oracle
3023                .manager()
3024                .track(
3025                    0,
3026                    Set::try_from([peers[0].clone(), peers[1].clone()]).unwrap(),
3027                );
3028            context.sleep(Duration::from_millis(100)).await;
3029
3030            let key = Key(7);
3031            let targeted_key = Key(8);
3032            let data = Bytes::from("old primary data");
3033
3034            let (cons1, mut cons_out1) = consumer();
3035
3036            // Peer 1: requester.
3037            let scheme = schemes.remove(0);
3038            let mut mailbox1 = setup_and_spawn_actor(
3039                &context,
3040                oracle.manager(),
3041                oracle.control(scheme.public_key()),
3042                scheme,
3043                connections.remove(0),
3044                cons1,
3045                Producer::default(),
3046            );
3047
3048            // Peer 2: old primary, still retained in `all.primary`, has the data.
3049            let mut prod2 = Producer::default();
3050            prod2.insert(key.clone(), data.clone());
3051            prod2.insert(targeted_key.clone(), data);
3052            let scheme = schemes.remove(0);
3053            let _mailbox2 = setup_and_spawn_actor(
3054                &context,
3055                oracle.manager(),
3056                oracle.control(scheme.public_key()),
3057                scheme,
3058                connections.remove(0),
3059                dummy_consumer(),
3060                prod2,
3061            );
3062
3063            // Peer 3: latest primary, has no data.
3064            let scheme = schemes.remove(0);
3065            let _mailbox3 = setup_and_spawn_actor(
3066                &context,
3067                oracle.manager(),
3068                oracle.control(scheme.public_key()),
3069                scheme,
3070                connections.remove(0),
3071                dummy_consumer(),
3072                Producer::default(),
3073            );
3074
3075            context.sleep(Duration::from_millis(100)).await;
3076
3077            // Track peer 3 as the latest primary while keeping the requester in the peer set.
3078            // Peer 2 remains in the provider's overlap window (`all.primary`), but new resolver traffic
3079            // should use only `latest.primary`.
3080            oracle
3081                .manager()
3082                .track(
3083                    1,
3084                    Set::try_from([peers[0].clone(), peers[2].clone()]).unwrap(),
3085                );
3086            context.sleep(Duration::from_millis(100)).await;
3087
3088            mailbox1.fetch(key);
3089
3090            select! {
3091                event = cons_out1.recv() => {
3092                    panic!(
3093                        "fetch should not succeed from an old primary retained only in the overlap window, got: {event:?}"
3094                    );
3095                },
3096                _ = context.sleep(Duration::from_secs(1)) => {},
3097            }
3098
3099            // Explicit targets still respect the latest-primary filter.
3100            mailbox1
3101                .fetch_targeted(targeted_key, non_empty_vec![peers[1].clone()]);
3102
3103            select! {
3104                event = cons_out1.recv() => {
3105                    panic!(
3106                        "targeted fetch should not bypass the latest-primary filter, got: {event:?}"
3107                    );
3108                },
3109                _ = context.sleep(Duration::from_secs(1)) => {},
3110            }
3111        });
3112    }
3113
3114    #[test_traced]
3115    fn test_fetch_after_cutover_relies_on_latest_primary_history() {
3116        let executor = deterministic::Runner::timed(Duration::from_secs(10));
3117        executor.start(|context| async move {
3118            let (network, oracle) = Network::new(
3119                context.child("network"),
3120                commonware_p2p::simulated::Config {
3121                    max_size: 1024 * 1024,
3122                    disconnect_on_block: true,
3123                    tracked_peer_sets: NZUsize!(2),
3124                },
3125            );
3126            network.start();
3127
3128            let schemes: Vec<PrivateKey> = [1u64, 2, 3]
3129                .into_iter()
3130                .map(PrivateKey::from_seed)
3131                .collect();
3132            let peers: Vec<PublicKey> = schemes.iter().map(|s| s.public_key()).collect();
3133            let mut schemes = schemes;
3134
3135            let mut connections = Vec::new();
3136            for peer in &peers {
3137                let (sender, receiver) = oracle
3138                    .control(peer.clone())
3139                    .register(0, Quota::per_second(RATE_LIMIT))
3140                    .await
3141                    .unwrap();
3142                connections.push((sender, receiver));
3143            }
3144
3145            let mut oracle = oracle;
3146            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3147            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
3148
3149            // Keep the requester in the peer set across the cutover while peer 2 remains connected
3150            // only through the overlap window after the latest primary advances to peer 3.
3151            oracle.manager().track(
3152                0,
3153                Set::try_from([peers[0].clone(), peers[1].clone()]).unwrap(),
3154            );
3155            context.sleep(Duration::from_millis(100)).await;
3156
3157            let key = Key(9);
3158            let invalid_history = Bytes::from("stale overlap history");
3159            let valid_history = Bytes::from("latest primary history");
3160
3161            let (mut cons1, mut cons_out1) = consumer();
3162            cons1.add_expected(key.clone(), valid_history.clone());
3163
3164            // Peer 1: requester.
3165            let scheme = schemes.remove(0);
3166            let mut mailbox1 = setup_and_spawn_actor(
3167                &context,
3168                oracle.manager(),
3169                oracle.control(scheme.public_key()),
3170                scheme,
3171                connections.remove(0),
3172                cons1,
3173                Producer::default(),
3174            );
3175
3176            // Peer 2: old primary retained only via overlap. If queried, it would be blocked for
3177            // serving invalid history.
3178            let mut prod2 = Producer::default();
3179            prod2.insert(key.clone(), invalid_history);
3180            let scheme = schemes.remove(0);
3181            let _mailbox2 = setup_and_spawn_actor(
3182                &context,
3183                oracle.manager(),
3184                oracle.control(scheme.public_key()),
3185                scheme,
3186                connections.remove(0),
3187                dummy_consumer(),
3188                prod2,
3189            );
3190
3191            // Peer 3: latest primary and the only peer that should satisfy the fetch.
3192            let mut prod3 = Producer::default();
3193            prod3.insert(key.clone(), valid_history.clone());
3194            let scheme = schemes.remove(0);
3195            let _mailbox3 = setup_and_spawn_actor(
3196                &context,
3197                oracle.manager(),
3198                oracle.control(scheme.public_key()),
3199                scheme,
3200                connections.remove(0),
3201                dummy_consumer(),
3202                prod3,
3203            );
3204
3205            context.sleep(Duration::from_millis(100)).await;
3206
3207            oracle.manager().track(
3208                1,
3209                Set::try_from([peers[0].clone(), peers[2].clone()]).unwrap(),
3210            );
3211            context.sleep(Duration::from_millis(100)).await;
3212
3213            mailbox1.fetch(key.clone());
3214
3215            let (key_actual, value) = cons_out1.recv().await.unwrap();
3216            assert_eq!(key_actual, key);
3217            assert_eq!(value, valid_history);
3218
3219            assert!(
3220                oracle.blocked().await.unwrap().is_empty(),
3221                "overlap-only peers should not be queried for post-cutover history"
3222            );
3223        });
3224    }
3225
3226    #[test_traced]
3227    fn test_secondary_peer_requests_are_served() {
3228        let executor = deterministic::Runner::timed(Duration::from_secs(10));
3229        executor.start(|context| async move {
3230            let (mut oracle, mut schemes, peers, mut connections) =
3231                setup_network_and_peers(&context, &[1, 2]).await;
3232
3233            // Topology: peer 1 is primary (has data), peer 2 is secondary (requester).
3234            // Verifies that a primary peer serves requests from secondary peers
3235            // (i.e. secondary peers can't fetch via broadcast, but their direct
3236            // requests are still answered).
3237            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3238
3239            oracle.manager().track(
3240                1,
3241                TrackedPeers::new(
3242                    Set::try_from([peers[0].clone()]).unwrap(),
3243                    Set::try_from([peers[1].clone()]).unwrap(),
3244                ),
3245            );
3246            context.sleep(Duration::from_millis(100)).await;
3247
3248            let key = Key(9);
3249            let data = Bytes::from("served to secondary");
3250
3251            // Peer 1: primary, has the data.
3252            let mut prod1 = Producer::default();
3253            prod1.insert(key.clone(), data.clone());
3254
3255            let scheme = schemes.remove(0);
3256            let _mailbox1 = setup_and_spawn_actor(
3257                &context,
3258                oracle.manager(),
3259                oracle.control(scheme.public_key()),
3260                scheme,
3261                connections.remove(0),
3262                dummy_consumer(),
3263                prod1,
3264            );
3265
3266            // Peer 2: secondary, uses fetch_targeted to explicitly request from peer 1.
3267            let (mut cons2, mut cons_out2) = consumer();
3268            cons2.add_expected(key.clone(), data.clone());
3269            let scheme = schemes.remove(0);
3270            let mut mailbox2 = setup_and_spawn_actor(
3271                &context,
3272                oracle.manager(),
3273                oracle.control(scheme.public_key()),
3274                scheme,
3275                connections.remove(0),
3276                cons2,
3277                Producer::default(),
3278            );
3279
3280            mailbox2.fetch_targeted(key.clone(), non_empty_vec![peers[0].clone()]);
3281
3282            let (key_actual, value) = cons_out2.recv().await.unwrap();
3283            assert_eq!(key_actual, key);
3284            assert_eq!(value, data);
3285        });
3286    }
3287
3288    #[test_traced]
3289    fn test_shutdown_aborts_pending_delivery_without_leaked_tasks() {
3290        let executor = deterministic::Runner::timed(Duration::from_secs(10));
3291        executor.start(|context| async move {
3292            let (mut oracle, mut schemes, peers, mut connections) =
3293                setup_network_and_peers(&context, &[1, 2]).await;
3294
3295            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3296
3297            let key = Key(1);
3298            let data = Bytes::from("data for key 1");
3299            let mut prod2 = Producer::default();
3300            prod2.insert(key.clone(), data);
3301
3302            let (mut gate_sender, gate_receiver) = oneshot::channel();
3303            let (cons1, mut cons_out1, mut started) =
3304                BlockingConsumer::new(context.child("consumer"), vec![(gate_receiver, true)]);
3305
3306            let actor_context = context.child("actor");
3307
3308            let scheme = schemes.remove(0);
3309            let public_key = scheme.public_key();
3310            let (engine, mut mailbox1): (_, Mailbox<Key, PublicKey>) = Engine::new(
3311                actor_context.child("peer").with_attribute("index", 0),
3312                Config {
3313                    peer_provider: oracle.manager(),
3314                    blocker: oracle.control(public_key.clone()),
3315                    consumer: cons1,
3316                    producer: Producer::<Key, Bytes>::default(),
3317                    mailbox_size: MAILBOX_SIZE,
3318                    me: Some(public_key),
3319                    initial: INITIAL_DURATION,
3320                    timeout: TIMEOUT,
3321                    fetch_retry_timeout: FETCH_RETRY_TIMEOUT,
3322                    priority_requests: false,
3323                    priority_responses: false,
3324                },
3325            );
3326            let handle1 = engine.start(connections.remove(0));
3327
3328            let scheme = schemes.remove(0);
3329            let public_key = scheme.public_key();
3330            let (engine, _mailbox2): (_, Mailbox<Key, PublicKey>) = Engine::new(
3331                actor_context.child("peer").with_attribute("index", 1),
3332                Config {
3333                    peer_provider: oracle.manager(),
3334                    blocker: oracle.control(public_key.clone()),
3335                    consumer: dummy_consumer(),
3336                    producer: prod2,
3337                    mailbox_size: MAILBOX_SIZE,
3338                    me: Some(public_key),
3339                    initial: INITIAL_DURATION,
3340                    timeout: TIMEOUT,
3341                    fetch_retry_timeout: FETCH_RETRY_TIMEOUT,
3342                    priority_requests: false,
3343                    priority_responses: false,
3344                },
3345            );
3346            let handle2 = engine.start(connections.remove(0));
3347
3348            mailbox1.fetch(key.clone());
3349            let started_key = started.recv().await.expect("delivery did not start");
3350            assert_eq!(started_key, key);
3351
3352            assert!(count_running_tasks(&context, "actor") > 0);
3353
3354            handle1.abort();
3355            handle2.abort();
3356
3357            context.sleep(Duration::from_millis(100)).await;
3358
3359            select! {
3360                _ = gate_sender.closed() => {},
3361                _ = context.sleep(Duration::from_secs(2)) => {
3362                    panic!("pending delivery was not aborted");
3363                },
3364            };
3365
3366            select! {
3367                event = cons_out1.recv() => assert!(event.is_none(), "unexpected event"),
3368                _ = context.sleep(Duration::from_millis(100)) => {},
3369            };
3370
3371            let running_after = count_running_tasks(&context, "actor");
3372            assert_eq!(
3373                running_after, 0,
3374                "all actor tasks should be stopped, but {running_after} still running"
3375            );
3376        });
3377    }
3378
3379    #[allow(clippy::type_complexity)]
3380    fn spawn_actors_with_handles(
3381        context: &deterministic::Context,
3382        oracle: &Oracle<PublicKey, deterministic::Context>,
3383        schemes: Vec<PrivateKey>,
3384        connections: Vec<(
3385            Sender<PublicKey, deterministic::Context>,
3386            Receiver<PublicKey>,
3387        )>,
3388        consumers: Vec<Consumer<Key, Bytes>>,
3389        producers: Vec<Producer<Key, Bytes>>,
3390    ) -> (
3391        Vec<Mailbox<Key, PublicKey>>,
3392        Vec<commonware_runtime::Handle<()>>,
3393    ) {
3394        let actor_context = context.child("actor");
3395        let mut mailboxes = Vec::new();
3396        let mut handles = Vec::new();
3397
3398        for (idx, ((scheme, conn), (consumer, producer))) in schemes
3399            .into_iter()
3400            .zip(connections)
3401            .zip(consumers.into_iter().zip(producers))
3402            .enumerate()
3403        {
3404            let ctx = actor_context.child("peer").with_attribute("index", idx);
3405            let public_key = scheme.public_key();
3406            let (engine, mailbox) = Engine::new(
3407                ctx,
3408                Config {
3409                    peer_provider: oracle.manager(),
3410                    blocker: oracle.control(public_key.clone()),
3411                    consumer,
3412                    producer,
3413                    mailbox_size: MAILBOX_SIZE,
3414                    me: Some(public_key),
3415                    initial: INITIAL_DURATION,
3416                    timeout: TIMEOUT,
3417                    fetch_retry_timeout: FETCH_RETRY_TIMEOUT,
3418                    priority_requests: false,
3419                    priority_responses: false,
3420                },
3421            );
3422            handles.push(engine.start(conn));
3423            mailboxes.push(mailbox);
3424        }
3425
3426        (mailboxes, handles)
3427    }
3428
3429    #[test_traced]
3430    fn test_operations_after_shutdown_do_not_panic() {
3431        let executor = deterministic::Runner::timed(Duration::from_secs(10));
3432        executor.start(|context| async move {
3433            let (mut oracle, schemes, peers, connections) =
3434                setup_network_and_peers(&context, &[1, 2]).await;
3435
3436            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3437
3438            let key = Key(1);
3439            let mut prod2 = Producer::default();
3440            prod2.insert(key.clone(), Bytes::from("data for key 1"));
3441
3442            let (cons1, mut cons_out1) = consumer();
3443
3444            let (mut mailboxes, handles) = spawn_actors_with_handles(
3445                &context,
3446                &oracle,
3447                schemes,
3448                connections,
3449                vec![cons1, dummy_consumer()],
3450                vec![Producer::default(), prod2],
3451            );
3452
3453            // Fetch to verify network is functional
3454            mailboxes[0].fetch(key.clone());
3455            let (_, value) = cons_out1.recv().await.unwrap();
3456            assert_eq!(value, Bytes::from("data for key 1"));
3457
3458            // Abort all actors
3459            for handle in handles {
3460                handle.abort();
3461            }
3462            context.sleep(Duration::from_millis(100)).await;
3463
3464            // All operations should not panic after shutdown
3465
3466            // Fetch should not panic
3467            let key2 = Key(2);
3468            mailboxes[0].fetch(key2.clone());
3469
3470            // Retain can prune a single key after shutdown without panicking.
3471            let canceled = key2;
3472            mailboxes[0].retain(move |key, _| key != &canceled);
3473
3474            // Retain should not panic
3475            mailboxes[0].retain(|_, _| true);
3476
3477            // Fetch targeted should not panic
3478            mailboxes[0].fetch_targeted(Key(3), non_empty_vec![peers[1].clone()]);
3479        });
3480    }
3481
3482    fn clean_shutdown(seed: u64) {
3483        let cfg = deterministic::Config::default()
3484            .with_seed(seed)
3485            .with_timeout(Some(Duration::from_secs(30)));
3486        let executor = deterministic::Runner::new(cfg);
3487        executor.start(|context| async move {
3488            let (mut oracle, schemes, peers, connections) =
3489                setup_network_and_peers(&context, &[1, 2]).await;
3490
3491            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3492
3493            let key = Key(1);
3494            let mut prod2 = Producer::default();
3495            prod2.insert(key.clone(), Bytes::from("data for key 1"));
3496
3497            let (cons1, mut cons_out1) = consumer();
3498
3499            let (mut mailboxes, handles) = spawn_actors_with_handles(
3500                &context,
3501                &oracle,
3502                schemes,
3503                connections,
3504                vec![cons1, dummy_consumer()],
3505                vec![Producer::default(), prod2],
3506            );
3507
3508            // Allow tasks to start
3509            context.sleep(Duration::from_millis(100)).await;
3510
3511            // Count running tasks under the actor prefix
3512            let running_before = count_running_tasks(&context, "actor");
3513            assert!(
3514                running_before > 0,
3515                "at least one actor task should be running"
3516            );
3517
3518            // Verify network is functional
3519            mailboxes[0].fetch(key.clone());
3520            let (_, value) = cons_out1.recv().await.unwrap();
3521            assert_eq!(value, Bytes::from("data for key 1"));
3522
3523            // Abort all actors
3524            for handle in handles {
3525                handle.abort();
3526            }
3527            context.sleep(Duration::from_millis(100)).await;
3528
3529            // Verify all actor tasks are stopped
3530            let running_after = count_running_tasks(&context, "actor");
3531            assert_eq!(
3532                running_after, 0,
3533                "all actor tasks should be stopped, but {running_after} still running"
3534            );
3535        });
3536    }
3537
3538    #[test]
3539    fn test_clean_shutdown() {
3540        for seed in 0..25 {
3541            clean_shutdown(seed);
3542        }
3543    }
3544}