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. Blocked peers are learned from the
12//! network through [`Blocker::blocked`](commonware_p2p::Blocker::blocked), so a peer becomes eligible
13//! again when its block expires. Fetches persist until pruned, fulfilled, or reported as no longer
14//! needed by the `Consumer`.
15//!
16//! The `Consumer` returns a [`crate::Outcome`], checking data integrity and authenticity unless it
17//! no longer needs the key. A complete response retires its delivered subscribers, an ambiguous
18//! response retries without penalizing the peer, an invalid response retries after blocking the
19//! peer, and an ignored response retires the key without scoring its peer. A verdict the consumer
20//! drops without answering hands the response to the remaining subscribers, or retires the key
21//! when none remain. Pruning a fetch with in-progress response validation aborts that validation;
22//! an invalid outcome produced after cancellation does not block the peer.
23//!
24//! The peer also serves data to other peers, forwarding network requests to the `Producer`. The
25//! `Producer` provides data asynchronously (e.g., from storage). If it fails, the peer sends an
26//! empty response, prompting the requester to retry elsewhere. Each message between peers contains
27//! an ID. Each request is sent with a unique ID, and each response includes the ID of the request
28//! it responds to.
29//!
30//! # Targeting
31//!
32//! Callers can restrict fetches to specific target peers using
33//! [`TargetedResolver::fetch_targeted`](crate::TargetedResolver::fetch_targeted).
34//! Only target peers are tried, there is no automatic fallback to other peers. Targets persist through
35//! transient failures (timeout, "no data" response, send failure) since the peer might be slow or
36//! receive the data later.
37//!
38//! While a fetch is in progress, callers can modify targeting:
39//! - [`TargetedResolver::fetch_targeted`](crate::TargetedResolver::fetch_targeted) adds peers to the existing target set
40//!   (only if the fetch already has targets, an "all" fetch remains unrestricted)
41//! - [`Resolver::fetch`](crate::Resolver::fetch) clears all targets, allowing fallback to any peer
42//!
43//! These modifications only apply to in-progress fetches. Once a fetch succeeds, is pruned, or is
44//! ignored by the consumer, the targets for that key are cleared automatically. A blocked peer is
45//! skipped until the network unblocks it, so a fetch whose every target is blocked stays
46//! outstanding and resumes when one of them is unblocked, new targets are added, or targeting is
47//! cleared.
48//!
49//! # Subscribers
50//!
51//! [`Resolver::fetch`](crate::Resolver::fetch) accepts a peer-visible key and a
52//! subscriber. This is useful when several subscribers can share the same peer-visible
53//! fetch. A fetch remains active while at least one attached subscriber satisfies the latest
54//! [`Resolver::retain`](crate::Resolver::retain) predicate. When the fetch resolves, the
55//! key and currently retained subscribers are supplied to
56//! [`Consumer::deliver`](crate::Consumer::deliver). Subscribers added while response validation
57//! is in progress are delivered the same response locally, once it is accepted or when the
58//! consumer drops its verdict without judging it.
59//!
60//! While a response is being validated, its key remains in flight, so no further request is sent.
61//! New fetches for the key only attach subscribers or targets. A complete outcome retires the
62//! delivered subscribers, an ambiguous outcome retries the key, an invalid outcome retries the key
63//! after blocking the serving peer, and an ignored outcome retires the entire key without scoring
64//! the serving peer. When a peer-visible key admits multiple valid responses, a consumer should
65//! return an ambiguous outcome if the delivered response does not satisfy every subscriber,
66//! allowing the resolver to try another response.
67//!
68//! # Scheduling
69//!
70//! All pending fresh keys are attempted before any pending retries. Fresh keys and retries are
71//! each ordered by their next attempt time.
72//!
73//! # Peer Selection
74//!
75//! Outbound fetches are only sent to peers in `latest.primary` (see [commonware_p2p::Provider]) but inbound
76//! requests are handled for all connected peers. Thus, callers that still expect a key to be fetchable after
77//! a peer set update must ensure the latest primary set can serve it.
78//!
79//! [`TargetedResolver::fetch_targeted`](crate::TargetedResolver::fetch_targeted) can narrow the current primary set
80//! further, but it does not bypass that latest-primary filter. Explicit targets that are no longer
81//! in the latest primary set are ignored until they become primary again.
82//!
83//! # Performance Considerations
84//!
85//! The peer supports arbitrarily many concurrent fetches, but resource usage generally
86//! depends on the rate-limiting configuration of the underlying P2P network.
87
88use bytes::Bytes;
89use commonware_utils::{Span, channel::oneshot};
90
91mod config;
92pub use config::Config;
93mod engine;
94pub use engine::Engine;
95mod fetcher;
96mod inflight;
97mod ingress;
98pub use ingress::Mailbox;
99mod metrics;
100mod wire;
101
102#[cfg(feature = "mocks")]
103pub mod mocks;
104
105/// Serves data requested by the network.
106pub trait Producer: Clone + Send + 'static {
107    /// Type used to key data requested from peers.
108    type Key: Span;
109
110    /// Serve a request received from the network.
111    fn produce(&mut self, key: Self::Key) -> oneshot::Receiver<Bytes>;
112}
113
114#[cfg(test)]
115mod tests {
116    use super::{
117        Config, Engine, Mailbox,
118        mocks::{Consumer, Key, Producer},
119    };
120    use crate::{Delivery, Fetch, Outcome, Resolver, TargetedResolver};
121    use bytes::Bytes;
122    use commonware_cryptography::{
123        Signer,
124        ed25519::{PrivateKey, PublicKey},
125    };
126    use commonware_macros::{select, test_traced};
127    use commonware_p2p::{
128        Blocker, Manager as _, Provider, TrackedPeers,
129        simulated::{Link, Network, Oracle, Receiver, Sender},
130    };
131    use commonware_runtime::{
132        Clock, Metrics as _, Quota, Runner, Spawner as _, Supervisor as _, deterministic,
133        telemetry::metrics::count_running_tasks,
134    };
135    use commonware_utils::{
136        NZU32, NZUsize,
137        channel::{
138            fallible::{FallibleExt, OneshotExt},
139            mpsc, oneshot,
140        },
141        non_empty_vec,
142        ordered::Set,
143        probability,
144        sync::Mutex,
145    };
146    use std::{
147        collections::{HashMap, VecDeque},
148        num::{NonZeroU32, NonZeroUsize},
149        sync::Arc,
150        time::Duration,
151    };
152
153    const MAILBOX_SIZE: NonZeroUsize = NZUsize!(1024);
154    const RATE_LIMIT: NonZeroU32 = NZU32!(10);
155    const TIMEOUT: Duration = Duration::from_millis(400);
156    const FETCH_RETRY_TIMEOUT: Duration = Duration::from_millis(100);
157    const LINK: Link = Link {
158        latency: Duration::from_millis(10),
159        jitter: Duration::from_millis(1),
160        success_rate: probability!(1.0),
161    };
162    const LINK_UNRELIABLE: Link = Link {
163        latency: Duration::from_millis(10),
164        jitter: Duration::from_millis(1),
165        success_rate: probability!(0.5),
166    };
167
168    fn status_metric_total(metrics: &str, name: &str, status: &str) -> u64 {
169        let prefix = format!("{name}{{");
170        let status_label = format!("status=\"{status}\"");
171        metrics
172            .lines()
173            .filter(|line| line.starts_with(&prefix) && line.contains(&status_label))
174            .map(|line| {
175                line.split_whitespace()
176                    .next_back()
177                    .expect("metric line must have a value")
178                    .parse::<u64>()
179                    .expect("status metric value must be an integer")
180            })
181            .sum()
182    }
183
184    async fn setup_network_and_peers(
185        context: &deterministic::Context,
186        peer_seeds: &[u64],
187    ) -> (
188        Oracle<PublicKey, deterministic::Context>,
189        Vec<PrivateKey>,
190        Vec<PublicKey>,
191        Vec<(
192            Sender<PublicKey, deterministic::Context>,
193            Receiver<PublicKey>,
194        )>,
195    ) {
196        setup_network_and_peers_with_rate_limit(context, peer_seeds, Quota::per_second(RATE_LIMIT))
197            .await
198    }
199
200    async fn setup_network_and_peers_with_rate_limit(
201        context: &deterministic::Context,
202        peer_seeds: &[u64],
203        rate_limit: Quota,
204    ) -> (
205        Oracle<PublicKey, deterministic::Context>,
206        Vec<PrivateKey>,
207        Vec<PublicKey>,
208        Vec<(
209            Sender<PublicKey, deterministic::Context>,
210            Receiver<PublicKey>,
211        )>,
212    ) {
213        let (network, oracle) = Network::new(
214            context.child("network"),
215            commonware_p2p::simulated::Config {
216                max_size: 1024 * 1024,
217                max_peers_per_set: NZUsize!(peer_seeds.len()),
218                disconnect_on_block: true,
219                tracked_peer_sets: NZUsize!(3),
220            },
221        );
222        network.start();
223
224        let schemes: Vec<PrivateKey> = peer_seeds
225            .iter()
226            .map(|seed| PrivateKey::from_seed(*seed))
227            .collect();
228        let peers: Vec<PublicKey> = schemes.iter().map(|s| s.public_key()).collect();
229        let mut manager = oracle.manager();
230        manager.track(0, Set::try_from(peers.clone()).unwrap());
231
232        let mut connections = Vec::new();
233        for peer in &peers {
234            let (sender, receiver) = oracle
235                .control(peer.clone())
236                .register(0, rate_limit)
237                .await
238                .unwrap();
239            connections.push((sender, receiver));
240        }
241
242        (oracle, schemes, peers, connections)
243    }
244
245    async fn add_link(
246        oracle: &mut Oracle<PublicKey, deterministic::Context>,
247        link: Link,
248        peers: &[PublicKey],
249        from: usize,
250        to: usize,
251    ) {
252        oracle
253            .add_link(peers[from].clone(), peers[to].clone(), link.clone())
254            .await
255            .unwrap();
256        oracle
257            .add_link(peers[to].clone(), peers[from].clone(), link)
258            .await
259            .unwrap();
260    }
261
262    #[derive(Clone, Default)]
263    struct SequencedProducer {
264        data: Arc<Mutex<HashMap<Key, VecDeque<Bytes>>>>,
265    }
266
267    impl SequencedProducer {
268        fn insert(&mut self, key: Key, values: impl IntoIterator<Item = Bytes>) {
269            self.data.lock().insert(key, values.into_iter().collect());
270        }
271
272        fn remaining(&self, key: &Key) -> Vec<Bytes> {
273            self.data
274                .lock()
275                .get(key)
276                .map(|values| values.iter().cloned().collect())
277                .unwrap_or_default()
278        }
279    }
280
281    impl crate::p2p::Producer for SequencedProducer {
282        type Key = Key;
283
284        fn produce(&mut self, key: Self::Key) -> oneshot::Receiver<Bytes> {
285            let (sender, receiver) = oneshot::channel();
286            if let Some(value) = self.data.lock().get_mut(&key).and_then(VecDeque::pop_front) {
287                let _ = sender.send(value);
288            }
289            receiver
290        }
291    }
292
293    fn setup_and_spawn_actor<C, R>(
294        context: &deterministic::Context,
295        provider: impl Provider<PublicKey = PublicKey>,
296        blocker: impl Blocker<PublicKey = PublicKey>,
297        signer: impl Signer<PublicKey = PublicKey>,
298        connection: (
299            Sender<PublicKey, deterministic::Context>,
300            Receiver<PublicKey>,
301        ),
302        consumer: C,
303        producer: Producer<Key, Bytes>,
304    ) -> Mailbox<Key, PublicKey, R>
305    where
306        C: crate::Consumer<Key = Key, Subscriber = R, Value = Bytes>,
307        R: Clone + Ord + Send + 'static,
308    {
309        setup_and_spawn_actor_with_producer(
310            context, provider, blocker, signer, connection, consumer, producer,
311        )
312    }
313
314    fn setup_and_spawn_actor_with_producer<C, R, Pro>(
315        context: &deterministic::Context,
316        provider: impl Provider<PublicKey = PublicKey>,
317        blocker: impl Blocker<PublicKey = PublicKey>,
318        signer: impl Signer<PublicKey = PublicKey>,
319        connection: (
320            Sender<PublicKey, deterministic::Context>,
321            Receiver<PublicKey>,
322        ),
323        consumer: C,
324        producer: Pro,
325    ) -> Mailbox<Key, PublicKey, R>
326    where
327        C: crate::Consumer<Key = Key, Subscriber = R, Value = Bytes>,
328        Pro: crate::p2p::Producer<Key = Key>,
329        R: Clone + Ord + Send + 'static,
330    {
331        let public_key = signer.public_key();
332        let (engine, mailbox) = Engine::new(
333            context.child("actor").with_attribute("peer", &public_key),
334            Config {
335                peer_provider: provider,
336                blocker,
337                consumer,
338                producer,
339                mailbox_size: MAILBOX_SIZE,
340                me: Some(public_key),
341                timeout: TIMEOUT,
342                fetch_retry_timeout: FETCH_RETRY_TIMEOUT,
343                priority_requests: false,
344                priority_responses: false,
345            },
346        );
347        engine.start(connection);
348
349        mailbox
350    }
351
352    type DeliveryGate = (oneshot::Receiver<()>, Outcome);
353    type DeliveryGates = Arc<Mutex<VecDeque<DeliveryGate>>>;
354
355    #[derive(Clone)]
356    struct BlockingConsumer {
357        context: Arc<deterministic::Context>,
358        sender: mpsc::UnboundedSender<(Key, Bytes)>,
359        started: mpsc::UnboundedSender<Key>,
360        gates: DeliveryGates,
361    }
362
363    impl BlockingConsumer {
364        fn new(
365            context: deterministic::Context,
366            gates: Vec<DeliveryGate>,
367        ) -> (
368            Self,
369            mpsc::UnboundedReceiver<(Key, Bytes)>,
370            mpsc::UnboundedReceiver<Key>,
371        ) {
372            let (sender, receiver) = mpsc::unbounded_channel();
373            let (started, started_receiver) = mpsc::unbounded_channel();
374            (
375                Self {
376                    context: Arc::new(context),
377                    sender,
378                    started,
379                    gates: Arc::new(Mutex::new(gates.into())),
380                },
381                receiver,
382                started_receiver,
383            )
384        }
385    }
386
387    impl crate::Consumer for BlockingConsumer {
388        type Key = Key;
389        type Value = Bytes;
390        type Subscriber = ();
391        type Outcome = Outcome;
392
393        fn deliver(
394            &mut self,
395            delivery: Delivery<Self::Key, Self::Subscriber>,
396            value: Self::Value,
397        ) -> oneshot::Receiver<Self::Outcome> {
398            let key = delivery.key;
399            self.started.send_lossy(key.clone());
400            let (gate, outcome) = self
401                .gates
402                .lock()
403                .pop_front()
404                .map_or((None, Outcome::Complete), |(gate, outcome)| {
405                    (Some(gate), outcome)
406                });
407            let (mut response, receiver) = oneshot::channel();
408            let sender = self.sender.clone();
409            self.context.child("delivery").spawn(move |_| async move {
410                if let Some(gate) = gate {
411                    select! {
412                        _ = response.closed() => return,
413                        result = gate => {
414                            if result.is_err() {
415                                let _ = response.send(Outcome::Invalid);
416                                return;
417                            }
418                        },
419                    }
420                }
421                if outcome == Outcome::Complete {
422                    sender.send_lossy((key, value));
423                }
424                let _ = response.send(outcome);
425            });
426            receiver
427        }
428    }
429
430    #[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)]
431    struct SubscriberTag(u16);
432
433    type RecordedDelivery = (Delivery<Key, SubscriberTag>, Bytes);
434
435    #[derive(Clone)]
436    struct BlockingSubscriberRecordingConsumer {
437        context: Arc<deterministic::Context>,
438        sender: mpsc::UnboundedSender<RecordedDelivery>,
439        started: mpsc::UnboundedSender<Delivery<Key, SubscriberTag>>,
440        gates: DeliveryGates,
441    }
442
443    impl BlockingSubscriberRecordingConsumer {
444        fn new(
445            context: deterministic::Context,
446            gates: Vec<DeliveryGate>,
447        ) -> (
448            Self,
449            mpsc::UnboundedReceiver<RecordedDelivery>,
450            mpsc::UnboundedReceiver<Delivery<Key, SubscriberTag>>,
451        ) {
452            let (sender, receiver) = mpsc::unbounded_channel();
453            let (started, started_receiver) = mpsc::unbounded_channel();
454            (
455                Self {
456                    context: Arc::new(context),
457                    sender,
458                    started,
459                    gates: Arc::new(Mutex::new(gates.into())),
460                },
461                receiver,
462                started_receiver,
463            )
464        }
465    }
466
467    impl crate::Consumer for BlockingSubscriberRecordingConsumer {
468        type Key = Key;
469        type Value = Bytes;
470        type Subscriber = SubscriberTag;
471        type Outcome = Outcome;
472
473        fn deliver(
474            &mut self,
475            delivery: Delivery<Self::Key, Self::Subscriber>,
476            value: Self::Value,
477        ) -> oneshot::Receiver<Self::Outcome> {
478            self.started.send_lossy(delivery.clone());
479            let (gate, outcome) = self
480                .gates
481                .lock()
482                .pop_front()
483                .map_or((None, Outcome::Complete), |(gate, outcome)| {
484                    (Some(gate), outcome)
485                });
486            let (mut response, receiver) = oneshot::channel();
487            let sender = self.sender.clone();
488            self.context.child("delivery").spawn(move |_| async move {
489                if let Some(gate) = gate {
490                    select! {
491                        _ = response.closed() => return,
492                        result = gate => {
493                            if result.is_err() {
494                                let _ = response.send(Outcome::Invalid);
495                                return;
496                            }
497                        },
498                    }
499                }
500                if outcome == Outcome::Complete {
501                    sender.send_lossy((delivery, value));
502                }
503                let _ = response.send(outcome);
504            });
505            receiver
506        }
507    }
508
509    #[derive(Clone)]
510    struct SubscriberRecordingConsumer {
511        sender: mpsc::UnboundedSender<RecordedDelivery>,
512    }
513
514    impl SubscriberRecordingConsumer {
515        fn new() -> (Self, mpsc::UnboundedReceiver<RecordedDelivery>) {
516            let (sender, receiver) = mpsc::unbounded_channel();
517            (Self { sender }, receiver)
518        }
519    }
520
521    impl crate::Consumer for SubscriberRecordingConsumer {
522        type Key = Key;
523        type Value = Bytes;
524        type Subscriber = SubscriberTag;
525        type Outcome = bool;
526
527        fn deliver(
528            &mut self,
529            delivery: Delivery<Self::Key, Self::Subscriber>,
530            value: Self::Value,
531        ) -> oneshot::Receiver<bool> {
532            let (sender, receiver) = oneshot::channel();
533            self.sender.send_lossy((delivery, value));
534            let _ = sender.send(true);
535            receiver
536        }
537    }
538
539    fn dummy_consumer() -> Consumer<Key, Bytes> {
540        Consumer::dummy()
541    }
542
543    fn consumer() -> (Consumer<Key, Bytes>, mpsc::UnboundedReceiver<(Key, Bytes)>) {
544        Consumer::new()
545    }
546
547    async fn wait_for_blocked(
548        context: &deterministic::Context,
549        oracle: &Oracle<PublicKey, deterministic::Context>,
550        blocker: &PublicKey,
551        blocked: &PublicKey,
552    ) {
553        loop {
554            let blocked_peers = oracle.blocked().await.unwrap();
555            if blocked_peers
556                .iter()
557                .any(|(a, b)| a == blocker && b == blocked)
558            {
559                return;
560            }
561            context.sleep(Duration::from_millis(10)).await;
562        }
563    }
564
565    /// Tests that fetching a key from another peer succeeds when data is available.
566    /// This test sets up two peers, where Peer 1 requests data that Peer 2 has,
567    /// and verifies that the data is correctly delivered to Peer 1's consumer.
568    #[test_traced]
569    fn test_fetch_success() {
570        let executor = deterministic::Runner::timed(Duration::from_secs(10));
571        executor.start(|context| async move {
572            let (mut oracle, mut schemes, peers, mut connections) =
573                setup_network_and_peers(&context, &[1, 2]).await;
574
575            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
576
577            let key = Key(2);
578            let mut prod2 = Producer::default();
579            prod2.insert(key.clone(), Bytes::from("data for key 2"));
580
581            let (cons1, mut cons_out1) = consumer();
582
583            let scheme = schemes.remove(0);
584            let mut mailbox1 = setup_and_spawn_actor(
585                &context,
586                oracle.manager(),
587                oracle.control(scheme.public_key()),
588                scheme,
589                connections.remove(0),
590                cons1,
591                Producer::default(),
592            );
593
594            let scheme = schemes.remove(0);
595            let _mailbox2 = setup_and_spawn_actor(
596                &context,
597                oracle.manager(),
598                oracle.control(scheme.public_key()),
599                scheme,
600                connections.remove(0),
601                dummy_consumer(),
602                prod2,
603            );
604
605            mailbox1.fetch(key.clone());
606
607            let (key_actual, value) = cons_out1.recv().await.unwrap();
608            assert_eq!(key_actual, key);
609            assert_eq!(value, Bytes::from("data for key 2"));
610        });
611    }
612
613    /// A consumer that hands each delivery's response sender to the test,
614    /// letting it resolve verdicts synchronously (without a task yield).
615    #[derive(Clone)]
616    struct HoldingConsumer {
617        deliveries: mpsc::UnboundedSender<(Key, oneshot::Sender<Outcome>)>,
618    }
619
620    impl HoldingConsumer {
621        fn new() -> (
622            Self,
623            mpsc::UnboundedReceiver<(Key, oneshot::Sender<Outcome>)>,
624        ) {
625            let (deliveries, receiver) = mpsc::unbounded_channel();
626            (Self { deliveries }, receiver)
627        }
628    }
629
630    impl crate::Consumer for HoldingConsumer {
631        type Key = Key;
632        type Value = Bytes;
633        type Subscriber = ();
634        type Outcome = Outcome;
635
636        fn deliver(
637            &mut self,
638            delivery: Delivery<Self::Key, Self::Subscriber>,
639            _: Self::Value,
640        ) -> oneshot::Receiver<Outcome> {
641            let (sender, receiver) = oneshot::channel();
642            self.deliveries.send_lossy((delivery.key, sender));
643            receiver
644        }
645    }
646
647    type HeldDelivery = (
648        Delivery<Key, SubscriberTag>,
649        Bytes,
650        oneshot::Sender<Outcome>,
651    );
652
653    /// A consumer that hands each delivery, its value, and its response sender
654    /// to the test, keeping the delivered subscribers visible.
655    #[derive(Clone)]
656    struct HoldingSubscriberConsumer {
657        deliveries: mpsc::UnboundedSender<HeldDelivery>,
658    }
659
660    impl HoldingSubscriberConsumer {
661        fn new() -> (Self, mpsc::UnboundedReceiver<HeldDelivery>) {
662            let (deliveries, receiver) = mpsc::unbounded_channel();
663            (Self { deliveries }, receiver)
664        }
665    }
666
667    impl crate::Consumer for HoldingSubscriberConsumer {
668        type Key = Key;
669        type Value = Bytes;
670        type Subscriber = SubscriberTag;
671        type Outcome = Outcome;
672
673        fn deliver(
674            &mut self,
675            delivery: Delivery<Self::Key, Self::Subscriber>,
676            value: Self::Value,
677        ) -> oneshot::Receiver<Outcome> {
678            let (sender, receiver) = oneshot::channel();
679            self.deliveries.send_lossy((delivery, value, sender));
680            receiver
681        }
682    }
683
684    #[test_traced]
685    fn test_fetch_after_accepted_verdict_restarts() {
686        let executor = deterministic::Runner::timed(Duration::from_secs(10));
687        executor.start(|context| async move {
688            let (mut oracle, mut schemes, peers, mut connections) =
689                setup_network_and_peers(&context, &[1, 2]).await;
690
691            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
692
693            let key = Key(2);
694            let mut prod2 = Producer::default();
695            prod2.insert(key.clone(), Bytes::from("data for key 2"));
696
697            let (cons1, mut deliveries) = HoldingConsumer::new();
698
699            let scheme = schemes.remove(0);
700            let mut mailbox1 = setup_and_spawn_actor(
701                &context,
702                oracle.manager(),
703                oracle.control(scheme.public_key()),
704                scheme,
705                connections.remove(0),
706                cons1,
707                Producer::default(),
708            );
709
710            let scheme = schemes.remove(0);
711            let _mailbox2 = setup_and_spawn_actor(
712                &context,
713                oracle.manager(),
714                oracle.control(scheme.public_key()),
715                scheme,
716                connections.remove(0),
717                dummy_consumer(),
718                prod2,
719            );
720
721            mailbox1.fetch(key.clone());
722            let (first, verdict) = deliveries.recv().await.unwrap();
723            assert_eq!(first, key);
724
725            // Accept the response and re-fetch the key before yielding: both
726            // the completion and the fetch are pending when the engine next
727            // runs, and it must settle the completion first so the re-fetch
728            // starts fresh instead of being deduplicated against the
729            // completing key and dropped with it.
730            verdict.send_lossy(Outcome::Complete);
731            mailbox1.fetch(key.clone());
732
733            select! {
734                delivered = deliveries.recv() => {
735                    let (second, verdict) = delivered.unwrap();
736                    assert_eq!(second, key);
737                    verdict.send_lossy(Outcome::Complete);
738                },
739                _ = context.sleep(Duration::from_secs(5)) => {
740                    panic!("re-fetch was dropped with the completed fetch");
741                },
742            }
743        });
744    }
745
746    #[test_traced]
747    fn test_pending_delivery_does_not_block_engine() {
748        let executor = deterministic::Runner::timed(Duration::from_secs(10));
749        executor.start(|context| async move {
750            let (mut oracle, mut schemes, peers, mut connections) =
751                setup_network_and_peers(&context, &[1, 2, 3]).await;
752
753            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
754            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
755
756            let key1 = Key(1);
757            let key2 = Key(2);
758            let data1 = Bytes::from("data for key 1");
759            let data2 = Bytes::from("data for key 2");
760
761            let mut prod2 = Producer::default();
762            prod2.insert(key1.clone(), data1.clone());
763
764            let mut prod3 = Producer::default();
765            prod3.insert(key2.clone(), data2.clone());
766
767            let (gate_sender1, gate_receiver1) = oneshot::channel();
768            let (gate_sender2, gate_receiver2) = oneshot::channel();
769            let (cons1, mut cons_out1, mut started) = BlockingConsumer::new(
770                context.child("consumer"),
771                vec![
772                    (gate_receiver1, Outcome::Complete),
773                    (gate_receiver2, Outcome::Complete),
774                ],
775            );
776
777            let scheme = schemes.remove(0);
778            let mut mailbox1 = setup_and_spawn_actor(
779                &context,
780                oracle.manager(),
781                oracle.control(scheme.public_key()),
782                scheme,
783                connections.remove(0),
784                cons1,
785                Producer::default(),
786            );
787
788            let scheme = schemes.remove(0);
789            let _mailbox2 = setup_and_spawn_actor(
790                &context,
791                oracle.manager(),
792                oracle.control(scheme.public_key()),
793                scheme,
794                connections.remove(0),
795                dummy_consumer(),
796                prod2,
797            );
798
799            let scheme = schemes.remove(0);
800            let _mailbox3 = setup_and_spawn_actor(
801                &context,
802                oracle.manager(),
803                oracle.control(scheme.public_key()),
804                scheme,
805                connections.remove(0),
806                dummy_consumer(),
807                prod3,
808            );
809
810            mailbox1.fetch(key1.clone());
811            let started_key = started.recv().await.expect("delivery did not start");
812            assert_eq!(started_key, key1);
813
814            mailbox1.fetch(key2.clone());
815            select! {
816                started_key = started.recv() => {
817                    assert_eq!(started_key.expect("delivery did not start"), key2);
818                },
819                _ = context.sleep(Duration::from_secs(2)) => {
820                    panic!("resolver engine blocked on pending delivery");
821                },
822            };
823
824            gate_sender2.send(()).unwrap();
825            let (key_actual, value) = cons_out1.recv().await.expect("consumer channel closed");
826            assert_eq!(key_actual, key2);
827            assert_eq!(value, data2);
828
829            gate_sender1.send(()).unwrap();
830            let (key_actual, value) = cons_out1.recv().await.expect("consumer channel closed");
831            assert_eq!(key_actual, key1);
832            assert_eq!(value, data1);
833        });
834    }
835
836    #[test_traced]
837    fn test_retain_drops_pending_delivery() {
838        let executor = deterministic::Runner::timed(Duration::from_secs(10));
839        executor.start(|context| async move {
840            let (mut oracle, mut schemes, peers, mut connections) =
841                setup_network_and_peers(&context, &[1, 2]).await;
842
843            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
844
845            let key = Key(1);
846            let data = Bytes::from("data for key 1");
847            let mut prod2 = Producer::default();
848            prod2.insert(key.clone(), data.clone());
849
850            let (mut first_gate_sender, first_gate_receiver) = oneshot::channel();
851            let (second_gate_sender, second_gate_receiver) = oneshot::channel();
852            let (cons1, mut cons_out1, mut started) = BlockingConsumer::new(
853                context.child("consumer"),
854                vec![
855                    (first_gate_receiver, Outcome::Complete),
856                    (second_gate_receiver, Outcome::Complete),
857                ],
858            );
859
860            let scheme = schemes.remove(0);
861            let mut mailbox1 = setup_and_spawn_actor(
862                &context,
863                oracle.manager(),
864                oracle.control(scheme.public_key()),
865                scheme,
866                connections.remove(0),
867                cons1,
868                Producer::default(),
869            );
870
871            let scheme = schemes.remove(0);
872            let _mailbox2 = setup_and_spawn_actor(
873                &context,
874                oracle.manager(),
875                oracle.control(scheme.public_key()),
876                scheme,
877                connections.remove(0),
878                dummy_consumer(),
879                prod2,
880            );
881
882            mailbox1.fetch(key.clone());
883            let started_key = started.recv().await.expect("delivery did not start");
884            assert_eq!(started_key, key);
885
886            let canceled = key.clone();
887            mailbox1.retain(move |key, _| key != &canceled);
888            mailbox1.fetch(key.clone());
889
890            first_gate_sender.closed().await;
891            let started_key = started.recv().await.expect("second delivery did not start");
892            assert_eq!(started_key, key);
893
894            second_gate_sender.send(()).unwrap();
895            let (key_actual, value) = cons_out1.recv().await.expect("consumer channel closed");
896            assert_eq!(key_actual, key);
897            assert_eq!(value, data);
898
899            select! {
900                _ = cons_out1.recv() => panic!("unexpected extra event"),
901                _ = context.sleep(Duration::from_millis(100)) => {},
902            };
903        });
904    }
905
906    #[test_traced]
907    fn test_invalid_delivery_retries_and_rearms_slot() {
908        let executor = deterministic::Runner::timed(Duration::from_secs(10));
909        executor.start(|context| async move {
910            let (mut oracle, mut schemes, peers, mut connections) =
911                setup_network_and_peers(&context, &[1, 2, 3]).await;
912
913            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
914
915            let key = Key(1);
916            let data = Bytes::from("data for key 1");
917
918            let mut prod2 = Producer::default();
919            prod2.insert(key.clone(), data.clone());
920
921            let mut prod3 = Producer::default();
922            prod3.insert(key.clone(), data.clone());
923
924            let (first_gate_sender, first_gate_receiver) = oneshot::channel();
925            let (second_gate_sender, second_gate_receiver) = oneshot::channel();
926            let (cons1, mut cons_out1, mut started) = BlockingConsumer::new(
927                context.child("consumer"),
928                vec![
929                    (first_gate_receiver, Outcome::Invalid),
930                    (second_gate_receiver, Outcome::Complete),
931                ],
932            );
933
934            let scheme = schemes.remove(0);
935            let mut mailbox1 = setup_and_spawn_actor(
936                &context,
937                oracle.manager(),
938                oracle.control(scheme.public_key()),
939                scheme,
940                connections.remove(0),
941                cons1,
942                Producer::default(),
943            );
944
945            let scheme = schemes.remove(0);
946            let _mailbox2 = setup_and_spawn_actor(
947                &context,
948                oracle.manager(),
949                oracle.control(scheme.public_key()),
950                scheme,
951                connections.remove(0),
952                dummy_consumer(),
953                prod2,
954            );
955
956            let scheme = schemes.remove(0);
957            let _mailbox3 = setup_and_spawn_actor(
958                &context,
959                oracle.manager(),
960                oracle.control(scheme.public_key()),
961                scheme,
962                connections.remove(0),
963                dummy_consumer(),
964                prod3,
965            );
966
967            mailbox1.fetch_targeted(
968                key.clone(),
969                non_empty_vec![peers[1].clone(), peers[2].clone()],
970            );
971            let started_key = started.recv().await.expect("delivery did not start");
972            assert_eq!(started_key, key);
973
974            first_gate_sender.send(()).unwrap();
975            wait_for_blocked(&context, &oracle, &peers[0], &peers[1]).await;
976
977            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
978            oracle.manager().track(
979                1,
980                Set::try_from([peers[0].clone(), peers[2].clone()]).unwrap(),
981            );
982
983            let started_key = started.recv().await.expect("retry delivery did not start");
984            assert_eq!(started_key, key);
985
986            second_gate_sender.send(()).unwrap();
987            let (key_actual, value) = cons_out1.recv().await.expect("consumer channel closed");
988            assert_eq!(key_actual, key);
989            assert_eq!(value, data);
990        });
991    }
992
993    #[test_traced]
994    fn test_ambiguous_delivery_retries_without_blocking_peer() {
995        let executor = deterministic::Runner::timed(Duration::from_secs(10));
996        executor.start(|context| async move {
997            let (mut oracle, mut schemes, peers, mut connections) =
998                setup_network_and_peers(&context, &[1, 2, 3]).await;
999
1000            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1001
1002            let key = Key(1);
1003            let data = Bytes::from("data for key 1");
1004            let mut prod2 = Producer::default();
1005            prod2.insert(key.clone(), data.clone());
1006            let mut prod3 = Producer::default();
1007            prod3.insert(key.clone(), data.clone());
1008
1009            let (first_gate_sender, first_gate_receiver) = oneshot::channel();
1010            let (second_gate_sender, second_gate_receiver) = oneshot::channel();
1011            let (cons1, mut cons_out1, mut started) = BlockingConsumer::new(
1012                context.child("consumer"),
1013                vec![
1014                    (first_gate_receiver, Outcome::Ambiguous),
1015                    (second_gate_receiver, Outcome::Complete),
1016                ],
1017            );
1018
1019            let scheme = schemes.remove(0);
1020            let mut mailbox1 = setup_and_spawn_actor(
1021                &context,
1022                oracle.manager(),
1023                oracle.control(scheme.public_key()),
1024                scheme,
1025                connections.remove(0),
1026                cons1,
1027                Producer::default(),
1028            );
1029
1030            let scheme = schemes.remove(0);
1031            let _mailbox2 = setup_and_spawn_actor(
1032                &context,
1033                oracle.manager(),
1034                oracle.control(scheme.public_key()),
1035                scheme,
1036                connections.remove(0),
1037                dummy_consumer(),
1038                prod2,
1039            );
1040
1041            let scheme = schemes.remove(0);
1042            let _mailbox3 = setup_and_spawn_actor(
1043                &context,
1044                oracle.manager(),
1045                oracle.control(scheme.public_key()),
1046                scheme,
1047                connections.remove(0),
1048                dummy_consumer(),
1049                prod3,
1050            );
1051
1052            mailbox1.fetch_targeted(
1053                key.clone(),
1054                non_empty_vec![peers[1].clone(), peers[2].clone()],
1055            );
1056            assert_eq!(started.recv().await.unwrap(), key);
1057            first_gate_sender.send(()).unwrap();
1058
1059            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
1060            oracle
1061                .remove_link(peers[0].clone(), peers[1].clone())
1062                .await
1063                .unwrap();
1064            oracle
1065                .remove_link(peers[1].clone(), peers[0].clone())
1066                .await
1067                .unwrap();
1068
1069            assert_eq!(started.recv().await.unwrap(), key);
1070            second_gate_sender.send(()).unwrap();
1071            assert_eq!(cons_out1.recv().await.unwrap(), (key, data));
1072            assert!(oracle.blocked().await.unwrap().is_empty());
1073
1074            let metrics = context.encode();
1075            assert_eq!(
1076                status_metric_total(&metrics, "actor_fetch_total", "Ambiguous"),
1077                1
1078            );
1079        });
1080    }
1081
1082    #[test_traced]
1083    fn test_dropped_verdict_retires_fetch_without_blocking_peer() {
1084        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1085        executor.start(|context| async move {
1086            let (mut oracle, mut schemes, peers, mut connections) =
1087                setup_network_and_peers(&context, &[1, 2]).await;
1088
1089            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1090
1091            let key = Key(1);
1092            let mut prod2 = Producer::default();
1093            prod2.insert(key.clone(), Bytes::from("data for key 1"));
1094
1095            let (cons1, mut deliveries) = HoldingConsumer::new();
1096
1097            let scheme = schemes.remove(0);
1098            let mut mailbox1 = setup_and_spawn_actor(
1099                &context,
1100                oracle.manager(),
1101                oracle.control(scheme.public_key()),
1102                scheme,
1103                connections.remove(0),
1104                cons1,
1105                Producer::default(),
1106            );
1107
1108            let scheme = schemes.remove(0);
1109            let _mailbox2 = setup_and_spawn_actor(
1110                &context,
1111                oracle.manager(),
1112                oracle.control(scheme.public_key()),
1113                scheme,
1114                connections.remove(0),
1115                dummy_consumer(),
1116                prod2,
1117            );
1118
1119            mailbox1.fetch(key.clone());
1120            let (delivered, verdict) = deliveries.recv().await.unwrap();
1121            assert_eq!(delivered, key);
1122
1123            // Drop the verdict, as a consumer does when it stops with the
1124            // delivery still queued. No one is waiting on the key, so the fetch
1125            // is retired rather than retried, and the peer stays unblocked
1126            // because nothing is known about what it served.
1127            drop(verdict);
1128            select! {
1129                _ = deliveries.recv() => panic!("retired fetch must not be retried"),
1130                _ = context.sleep(Duration::from_secs(1)) => {},
1131            };
1132            assert!(oracle.blocked().await.unwrap().is_empty());
1133
1134            // The retired key is not deduplicated against, so a fresh fetch runs.
1135            mailbox1.fetch(key.clone());
1136            let (delivered, verdict) = deliveries.recv().await.unwrap();
1137            assert_eq!(delivered, key);
1138            verdict.send_lossy(Outcome::Complete);
1139
1140            let metrics = context.encode();
1141            assert_eq!(
1142                status_metric_total(&metrics, "actor_fetch_total", "Dropped"),
1143                1
1144            );
1145        });
1146    }
1147
1148    #[test_traced]
1149    fn test_dropped_verdict_hands_response_to_late_subscriber() {
1150        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1151        executor.start(|context| async move {
1152            let (mut oracle, mut schemes, peers, mut connections) =
1153                setup_network_and_peers(&context, &[1, 2]).await;
1154
1155            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1156
1157            let key = Key(1);
1158            let data = Bytes::from("data for key 1");
1159            let mut prod2 = Producer::default();
1160            prod2.insert(key.clone(), data.clone());
1161
1162            let (cons1, mut deliveries) = HoldingSubscriberConsumer::new();
1163
1164            let scheme = schemes.remove(0);
1165            let mut mailbox1 = setup_and_spawn_actor(
1166                &context,
1167                oracle.manager(),
1168                oracle.control(scheme.public_key()),
1169                scheme,
1170                connections.remove(0),
1171                cons1,
1172                Producer::default(),
1173            );
1174
1175            let scheme = schemes.remove(0);
1176            let _mailbox2 = setup_and_spawn_actor(
1177                &context,
1178                oracle.manager(),
1179                oracle.control(scheme.public_key()),
1180                scheme,
1181                connections.remove(0),
1182                dummy_consumer(),
1183                prod2,
1184            );
1185
1186            let first = SubscriberTag(1);
1187            let late = SubscriberTag(2);
1188            mailbox1.fetch(Fetch {
1189                key: key.clone(),
1190                subscriber: first.clone(),
1191                span: tracing::Span::none(),
1192            });
1193            let (delivery, value, verdict) = deliveries.recv().await.unwrap();
1194            assert_eq!(
1195                delivery,
1196                Delivery {
1197                    key: key.clone(),
1198                    subscribers: non_empty_vec![(first, tracing::Span::none())],
1199                }
1200            );
1201            assert_eq!(value, data);
1202
1203            // A late subscriber joins while the first delivery is unjudged, then
1204            // that delivery's verdict is dropped. The late subscriber is handed
1205            // the same response to judge instead of the key being retired.
1206            mailbox1.fetch(Fetch {
1207                key: key.clone(),
1208                subscriber: late.clone(),
1209                span: tracing::Span::none(),
1210            });
1211            context.sleep(Duration::from_millis(100)).await;
1212            drop(verdict);
1213            let (delivery, value, verdict) = deliveries.recv().await.unwrap();
1214            assert_eq!(
1215                delivery,
1216                Delivery {
1217                    key,
1218                    subscribers: non_empty_vec![(late, tracing::Span::none())],
1219                }
1220            );
1221            assert_eq!(value, data);
1222            verdict.send_lossy(Outcome::Complete);
1223
1224            context.sleep(Duration::from_millis(100)).await;
1225            assert!(oracle.blocked().await.unwrap().is_empty());
1226            let metrics = context.encode();
1227            assert_eq!(
1228                status_metric_total(&metrics, "actor_fetch_total", "Success"),
1229                1
1230            );
1231            assert_eq!(
1232                status_metric_total(&metrics, "actor_fetch_total", "Dropped"),
1233                0
1234            );
1235        });
1236    }
1237
1238    #[test_traced]
1239    fn test_ignored_delivery_retires_late_subscribers_without_rating_or_retry() {
1240        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1241        executor.start(|context| async move {
1242            let (mut oracle, mut schemes, peers, mut connections) =
1243                setup_network_and_peers(&context, &[1, 2]).await;
1244
1245            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1246
1247            let key = Key(1);
1248            let first = Bytes::from("obsolete response");
1249            let second = Bytes::from("fresh response");
1250            let mut producer = SequencedProducer::default();
1251            producer.insert(key.clone(), [first, second.clone()]);
1252            let producer_observer = producer.clone();
1253
1254            let (gate_sender, gate_receiver) = oneshot::channel();
1255            let (consumer, mut deliveries, mut started) = BlockingSubscriberRecordingConsumer::new(
1256                context.child("consumer"),
1257                vec![(gate_receiver, Outcome::Ignored)],
1258            );
1259
1260            let scheme = schemes.remove(0);
1261            let mut mailbox = setup_and_spawn_actor(
1262                &context,
1263                oracle.manager(),
1264                oracle.control(scheme.public_key()),
1265                scheme,
1266                connections.remove(0),
1267                consumer,
1268                Producer::default(),
1269            );
1270
1271            let scheme = schemes.remove(0);
1272            let _responder = setup_and_spawn_actor_with_producer(
1273                &context,
1274                oracle.manager(),
1275                oracle.control(scheme.public_key()),
1276                scheme,
1277                connections.remove(0),
1278                dummy_consumer(),
1279                producer,
1280            );
1281
1282            let first_subscriber = SubscriberTag(1);
1283            let late_subscriber = SubscriberTag(2);
1284            let fresh_subscriber = SubscriberTag(3);
1285            mailbox.fetch(Fetch {
1286                key: key.clone(),
1287                subscriber: first_subscriber.clone(),
1288                span: tracing::Span::none(),
1289            });
1290            assert_eq!(
1291                started.recv().await.expect("delivery did not start"),
1292                Delivery {
1293                    key: key.clone(),
1294                    subscribers: non_empty_vec![(first_subscriber, tracing::Span::none())],
1295                }
1296            );
1297
1298            // A subscriber attached after the delivery snapshot is still retired by
1299            // the key-global ignored outcome.
1300            mailbox.fetch(Fetch {
1301                key: key.clone(),
1302                subscriber: late_subscriber,
1303                span: tracing::Span::none(),
1304            });
1305            context.sleep(Duration::from_millis(100)).await;
1306            assert_eq!(producer_observer.remaining(&key), vec![second.clone()]);
1307            gate_sender.send(()).expect("consumer gate dropped");
1308
1309            context
1310                .sleep(FETCH_RETRY_TIMEOUT + Duration::from_millis(100))
1311                .await;
1312            assert_eq!(producer_observer.remaining(&key), vec![second.clone()]);
1313            assert!(oracle.blocked().await.unwrap().is_empty());
1314
1315            let metrics = context.encode();
1316            assert_eq!(
1317                status_metric_total(&metrics, "actor_fetch_total", "Dropped"),
1318                1
1319            );
1320            assert!(
1321                metrics
1322                    .lines()
1323                    .filter(|line| line.contains("resolves_count"))
1324                    .all(|line| line.ends_with(" 0")),
1325                "ignored response was recorded as a peer resolve:\n{metrics}"
1326            );
1327
1328            // A new request for the same key starts cleanly after the ignored fetch is retired.
1329            mailbox.fetch(Fetch {
1330                key: key.clone(),
1331                subscriber: fresh_subscriber.clone(),
1332                span: tracing::Span::none(),
1333            });
1334            assert_eq!(
1335                started.recv().await.expect("fresh delivery did not start"),
1336                Delivery {
1337                    key: key.clone(),
1338                    subscribers: non_empty_vec![(fresh_subscriber.clone(), tracing::Span::none())],
1339                }
1340            );
1341            assert_eq!(
1342                deliveries.recv().await.expect("consumer channel closed"),
1343                (
1344                    Delivery {
1345                        key: key.clone(),
1346                        subscribers: non_empty_vec![(fresh_subscriber, tracing::Span::none())],
1347                    },
1348                    second
1349                )
1350            );
1351            assert!(producer_observer.remaining(&key).is_empty());
1352        });
1353    }
1354
1355    async fn run_pending_invalid_delivery_race(
1356        context: &deterministic::Context,
1357        validation_first: bool,
1358    ) {
1359        let (mut oracle, mut schemes, peers, mut connections) =
1360            setup_network_and_peers(context, &[1, 2]).await;
1361
1362        add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1363
1364        let key = Key(1);
1365        let mut prod2 = Producer::default();
1366        prod2.insert(key.clone(), Bytes::from("data for key 1"));
1367
1368        let (mut gate_sender, gate_receiver) = oneshot::channel();
1369        let (cons1, mut cons_out1, mut started) = BlockingConsumer::new(
1370            context.child("consumer"),
1371            vec![(gate_receiver, Outcome::Invalid)],
1372        );
1373
1374        let scheme = schemes.remove(0);
1375        let mut mailbox1 = setup_and_spawn_actor(
1376            context,
1377            oracle.manager(),
1378            oracle.control(scheme.public_key()),
1379            scheme,
1380            connections.remove(0),
1381            cons1,
1382            Producer::default(),
1383        );
1384
1385        let scheme = schemes.remove(0);
1386        let _mailbox2 = setup_and_spawn_actor(
1387            context,
1388            oracle.manager(),
1389            oracle.control(scheme.public_key()),
1390            scheme,
1391            connections.remove(0),
1392            dummy_consumer(),
1393            prod2,
1394        );
1395
1396        mailbox1.fetch(key.clone());
1397        let started_key = started.recv().await.expect("delivery did not start");
1398        assert_eq!(started_key, key);
1399
1400        if validation_first {
1401            gate_sender.send(()).unwrap();
1402            wait_for_blocked(context, &oracle, &peers[0], &peers[1]).await;
1403            mailbox1.retain(|_, _| false);
1404            let blocked = oracle.blocked().await.unwrap();
1405            assert_eq!(blocked.len(), 1);
1406            assert_eq!(blocked[0].0, peers[0]);
1407            assert_eq!(blocked[0].1, peers[1]);
1408        } else {
1409            mailbox1.retain(|_, _| false);
1410            gate_sender.closed().await;
1411            assert!(oracle.blocked().await.unwrap().is_empty());
1412        }
1413
1414        select! {
1415            _ = cons_out1.recv() => panic!("unexpected event"),
1416            _ = context.sleep(Duration::from_millis(100)) => {},
1417        };
1418    }
1419
1420    #[test_traced]
1421    fn test_retain_pending_invalid_delivery_race() {
1422        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1423        executor.start(|context| async move {
1424            run_pending_invalid_delivery_race(&context, false).await;
1425            run_pending_invalid_delivery_race(&context, true).await;
1426        });
1427    }
1428
1429    /// Tests that pruning a fetch leaves the consumer untouched.
1430    /// This test initiates a fetch and immediately prunes it, verifying
1431    /// that the consumer does not receive any event.
1432    #[test_traced]
1433    fn test_retain_drops_fetch() {
1434        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1435        executor.start(|context| async move {
1436            let (oracle, mut schemes, _peers, mut connections) =
1437                setup_network_and_peers(&context, &[1]).await;
1438
1439            let (cons1, mut cons_out1) = consumer();
1440            let prod1 = Producer::default();
1441
1442            let scheme = schemes.remove(0);
1443            let mut mailbox1 = setup_and_spawn_actor(
1444                &context,
1445                oracle.manager(),
1446                oracle.control(scheme.public_key()),
1447                scheme,
1448                connections.remove(0),
1449                cons1,
1450                prod1,
1451            );
1452
1453            let key = Key(3);
1454            mailbox1.fetch(key.clone());
1455            let canceled = key.clone();
1456            mailbox1.retain(move |key, _| key != &canceled);
1457
1458            select! {
1459                _ = cons_out1.recv() => panic!("unexpected event"),
1460                _ = context.sleep(Duration::from_millis(100)) => {},
1461            };
1462        });
1463    }
1464
1465    /// Tests fetching data from a peer when some peers lack the data.
1466    /// This test sets up three peers, where Peer 1 requests data that only Peer 3 has.
1467    /// It verifies that the resolver retries with another peer and successfully
1468    /// delivers the data to Peer 1's consumer.
1469    #[test_traced]
1470    fn test_peer_no_data() {
1471        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1472        executor.start(|context| async move {
1473            let (mut oracle, mut schemes, peers, mut connections) =
1474                setup_network_and_peers(&context, &[1, 2, 3]).await;
1475
1476            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1477            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
1478
1479            let prod1 = Producer::default();
1480            let prod2 = Producer::default();
1481            let mut prod3 = Producer::default();
1482            let key = Key(3);
1483            prod3.insert(key.clone(), Bytes::from("data for key 3"));
1484
1485            let (cons1, mut cons_out1) = consumer();
1486
1487            let scheme = schemes.remove(0);
1488            let mut mailbox1 = setup_and_spawn_actor(
1489                &context,
1490                oracle.manager(),
1491                oracle.control(scheme.public_key()),
1492                scheme,
1493                connections.remove(0),
1494                cons1,
1495                prod1,
1496            );
1497
1498            let scheme = schemes.remove(0);
1499            let _mailbox2 = setup_and_spawn_actor(
1500                &context,
1501                oracle.manager(),
1502                oracle.control(scheme.public_key()),
1503                scheme,
1504                connections.remove(0),
1505                dummy_consumer(),
1506                prod2,
1507            );
1508
1509            let scheme = schemes.remove(0);
1510            let _mailbox3 = setup_and_spawn_actor(
1511                &context,
1512                oracle.manager(),
1513                oracle.control(scheme.public_key()),
1514                scheme,
1515                connections.remove(0),
1516                dummy_consumer(),
1517                prod3,
1518            );
1519
1520            mailbox1.fetch(key.clone());
1521
1522            let (key_actual, value) = cons_out1.recv().await.unwrap();
1523            assert_eq!(key_actual, key);
1524            assert_eq!(value, Bytes::from("data for key 3"));
1525        });
1526    }
1527
1528    /// Tests fetching when no peers are available.
1529    /// This test sets up a single peer with an empty peer provider (no peers).
1530    /// It initiates a fetch, waits beyond the retry timeout, prunes the fetch,
1531    /// and verifies that the consumer receives a failure notification.
1532    #[test_traced]
1533    fn test_no_peers_available() {
1534        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1535        executor.start(|context| async move {
1536            let (oracle, mut schemes, _peers, mut connections) =
1537                setup_network_and_peers(&context, &[1]).await;
1538
1539            let (cons1, mut cons_out1) = consumer();
1540            let prod1 = Producer::default();
1541
1542            let scheme = schemes.remove(0);
1543            let mut mailbox1 = setup_and_spawn_actor(
1544                &context,
1545                oracle.manager(),
1546                oracle.control(scheme.public_key()),
1547                scheme,
1548                connections.remove(0),
1549                cons1,
1550                prod1,
1551            );
1552
1553            mailbox1.fetch(Key(4));
1554            context.sleep(Duration::from_secs(5)).await;
1555
1556            // With no peers, no event should arrive
1557            select! {
1558                _ = cons_out1.recv() => panic!("Fetch should have failed due to no peers"),
1559                _ = context.sleep(Duration::from_millis(100)) => {},
1560            };
1561        });
1562    }
1563
1564    /// Tests that fetches issued before the first peer set arrives stay pending and complete once
1565    /// the initial update is tracked.
1566    #[test_traced]
1567    fn test_fetch_before_initial_peer_set_waits_for_update() {
1568        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1569        executor.start(|context| async move {
1570            let (network, mut oracle) = Network::new(
1571                context.child("network"),
1572                commonware_p2p::simulated::Config {
1573                    max_size: 1024 * 1024,
1574                    max_peers_per_set: NZUsize!(2),
1575                    disconnect_on_block: true,
1576                    tracked_peer_sets: NZUsize!(1),
1577                },
1578            );
1579            network.start();
1580
1581            let mut schemes = [1_u64, 2]
1582                .into_iter()
1583                .map(PrivateKey::from_seed)
1584                .collect::<Vec<_>>();
1585            schemes.sort_by_key(|s| s.public_key());
1586            let peers: Vec<PublicKey> = schemes.iter().map(|s| s.public_key()).collect();
1587
1588            let mut connections = Vec::new();
1589            for peer in &peers {
1590                let (sender, receiver) = oracle
1591                    .control(peer.clone())
1592                    .register(0, Quota::per_second(RATE_LIMIT))
1593                    .await
1594                    .unwrap();
1595                connections.push((sender, receiver));
1596            }
1597
1598            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1599
1600            let key = Key(2);
1601            let mut prod2 = Producer::default();
1602            prod2.insert(key.clone(), Bytes::from("data for key 2"));
1603
1604            let (cons1, mut cons_out1) = consumer();
1605
1606            let scheme = schemes.remove(0);
1607            let mut mailbox1 = setup_and_spawn_actor(
1608                &context,
1609                oracle.manager(),
1610                oracle.control(scheme.public_key()),
1611                scheme,
1612                connections.remove(0),
1613                cons1,
1614                Producer::default(),
1615            );
1616
1617            let scheme = schemes.remove(0);
1618            let _mailbox2 = setup_and_spawn_actor(
1619                &context,
1620                oracle.manager(),
1621                oracle.control(scheme.public_key()),
1622                scheme,
1623                connections.remove(0),
1624                dummy_consumer(),
1625                prod2,
1626            );
1627
1628            mailbox1.fetch(key.clone());
1629
1630            select! {
1631                event = cons_out1.recv() => {
1632                    panic!("fetch should wait for the initial peer set, got {event:?}");
1633                },
1634                _ = context.sleep(Duration::from_millis(200)) => {},
1635            };
1636
1637            oracle
1638                .manager()
1639                .track(0, Set::try_from(peers.clone()).unwrap());
1640
1641            let (key_actual, value) = cons_out1.recv().await.unwrap();
1642            assert_eq!(key_actual, key);
1643            assert_eq!(value, Bytes::from("data for key 2"));
1644        });
1645    }
1646
1647    /// Tests that concurrent fetches are handled correctly.
1648    /// Also tests that the peer can recover from having no peers available.
1649    /// Also tests that the peer can get data from multiple peers that have different sets of data.
1650    #[test_traced]
1651    fn test_concurrent_fetch_requests() {
1652        let executor = deterministic::Runner::default();
1653        executor.start(|context| async move {
1654            let (mut oracle, mut schemes, peers, mut connections) =
1655                setup_network_and_peers(&context, &[1, 2, 3]).await;
1656
1657            let key2 = Key(2);
1658            let key3 = Key(3);
1659            let mut prod2 = Producer::default();
1660            prod2.insert(key2.clone(), Bytes::from("data for key 2"));
1661            let mut prod3 = Producer::default();
1662            prod3.insert(key3.clone(), Bytes::from("data for key 3"));
1663
1664            let (cons1, mut cons_out1) = consumer();
1665
1666            let scheme = schemes.remove(0);
1667            let mut mailbox1 = setup_and_spawn_actor(
1668                &context,
1669                oracle.manager(),
1670                oracle.control(scheme.public_key()),
1671                scheme,
1672                connections.remove(0),
1673                cons1,
1674                Producer::default(),
1675            );
1676
1677            let scheme = schemes.remove(0);
1678            let _mailbox2 = setup_and_spawn_actor(
1679                &context,
1680                oracle.manager(),
1681                oracle.control(scheme.public_key()),
1682                scheme,
1683                connections.remove(0),
1684                dummy_consumer(),
1685                prod2,
1686            );
1687
1688            let scheme = schemes.remove(0);
1689            let _mailbox3 = setup_and_spawn_actor(
1690                &context,
1691                oracle.manager(),
1692                oracle.control(scheme.public_key()),
1693                scheme,
1694                connections.remove(0),
1695                dummy_consumer(),
1696                prod3,
1697            );
1698
1699            // Add choppy links between the requester and the two producers
1700            add_link(&mut oracle, LINK_UNRELIABLE.clone(), &peers, 0, 1).await;
1701            add_link(&mut oracle, LINK_UNRELIABLE.clone(), &peers, 0, 2).await;
1702
1703            // Run the fetches multiple times to ensure that the peer tries both of its peers
1704            for _ in 0..10 {
1705                // Initiate concurrent fetches.
1706                mailbox1.fetch(key2.clone());
1707                mailbox1.fetch(key3.clone());
1708
1709                // Collect both events without assuming order
1710                let mut events = Vec::new();
1711                events.push(cons_out1.recv().await.expect("Consumer channel closed"));
1712                events.push(cons_out1.recv().await.expect("Consumer channel closed"));
1713
1714                // Check that both keys were successfully fetched
1715                let mut found_key2 = false;
1716                let mut found_key3 = false;
1717                for (key_actual, value) in events {
1718                    if key_actual == key2 {
1719                        assert_eq!(value, Bytes::from("data for key 2"));
1720                        found_key2 = true;
1721                    } else if key_actual == key3 {
1722                        assert_eq!(value, Bytes::from("data for key 3"));
1723                        found_key3 = true;
1724                    } else {
1725                        panic!("Unexpected key received");
1726                    }
1727                }
1728                assert!(found_key2 && found_key3,);
1729            }
1730        });
1731    }
1732
1733    /// Tests that pruning an inactive fetch has no effect.
1734    /// Prunes a key before, after, and during the fetch process.
1735    #[test_traced]
1736    fn test_retain_drops_key() {
1737        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1738        executor.start(|context| async move {
1739            let (mut oracle, mut schemes, peers, mut connections) =
1740                setup_network_and_peers(&context, &[1, 2]).await;
1741
1742            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1743
1744            let key = Key(6);
1745            let mut prod2 = Producer::default();
1746            prod2.insert(key.clone(), Bytes::from("data for key 6"));
1747
1748            let (cons1, mut cons_out1) = consumer();
1749
1750            let scheme = schemes.remove(0);
1751            let mut mailbox1 = setup_and_spawn_actor(
1752                &context,
1753                oracle.manager(),
1754                oracle.control(scheme.public_key()),
1755                scheme,
1756                connections.remove(0),
1757                cons1,
1758                Producer::default(),
1759            );
1760
1761            let scheme = schemes.remove(0);
1762            let _mailbox2 = setup_and_spawn_actor(
1763                &context,
1764                oracle.manager(),
1765                oracle.control(scheme.public_key()),
1766                scheme,
1767                connections.remove(0),
1768                dummy_consumer(),
1769                prod2,
1770            );
1771
1772            // Prune before sending the fetch, expecting no effect.
1773            let canceled = key.clone();
1774            mailbox1.retain(move |key, _| key != &canceled);
1775            select! {
1776                _ = cons_out1.recv() => {
1777                    panic!("unexpected event");
1778                },
1779                _ = context.sleep(Duration::from_millis(100)) => {},
1780            };
1781
1782            // Initiate fetch and wait for data to be delivered
1783            mailbox1.fetch(key.clone());
1784            let (key_actual, value) = cons_out1.recv().await.unwrap();
1785            assert_eq!(key_actual, key);
1786            assert_eq!(value, Bytes::from("data for key 6"));
1787
1788            // Attempt to prune after data has been delivered, expecting no effect
1789            let canceled = key.clone();
1790            mailbox1.retain(move |key, _| key != &canceled);
1791            select! {
1792                _ = cons_out1.recv() => {
1793                    panic!("unexpected event");
1794                },
1795                _ = context.sleep(Duration::from_millis(100)) => {},
1796            };
1797
1798            // Initiate and prune another fetch.
1799            let key = Key(7);
1800            mailbox1.fetch(key.clone());
1801            let canceled = key.clone();
1802            mailbox1.retain(move |key, _| key != &canceled);
1803
1804            // No event should arrive after pruning.
1805            select! {
1806                _ = cons_out1.recv() => panic!("unexpected event"),
1807                _ = context.sleep(Duration::from_millis(100)) => {},
1808            };
1809        });
1810    }
1811
1812    /// Tests that a peer is blocked after delivering invalid data,
1813    /// preventing further fetches from that peer.
1814    #[test_traced]
1815    fn test_blocking_peer() {
1816        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1817        executor.start(|context| async move {
1818            let (mut oracle, mut schemes, peers, mut connections) =
1819                setup_network_and_peers(&context, &[1, 2, 3]).await;
1820
1821            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1822            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
1823            add_link(&mut oracle, LINK.clone(), &peers, 1, 2).await;
1824
1825            let key_a = Key(1);
1826            let key_b = Key(2);
1827            let invalid_data_a = Bytes::from("invalid for A");
1828            let valid_data_a = Bytes::from("valid for A");
1829            let valid_data_b = Bytes::from("valid for B");
1830
1831            // Set up producers
1832            let mut prod2 = Producer::default();
1833            prod2.insert(key_a.clone(), invalid_data_a.clone());
1834            prod2.insert(key_b.clone(), valid_data_b.clone());
1835
1836            let mut prod3 = Producer::default();
1837            prod3.insert(key_a.clone(), valid_data_a.clone());
1838
1839            // Set up consumer for Peer1 with expected values
1840            let (mut cons1, mut cons_out1) = consumer();
1841            cons1.add_expected(key_a.clone(), valid_data_a.clone());
1842            cons1.add_expected(key_b.clone(), valid_data_b.clone());
1843
1844            // Spawn actors
1845            let scheme = schemes.remove(0);
1846            let mut mailbox1 = setup_and_spawn_actor(
1847                &context,
1848                oracle.manager(),
1849                oracle.control(scheme.public_key()),
1850                scheme,
1851                connections.remove(0),
1852                cons1,
1853                Producer::default(),
1854            );
1855
1856            let scheme = schemes.remove(0);
1857            let _mailbox2 = setup_and_spawn_actor(
1858                &context,
1859                oracle.manager(),
1860                oracle.control(scheme.public_key()),
1861                scheme,
1862                connections.remove(0),
1863                dummy_consumer(),
1864                prod2,
1865            );
1866
1867            let scheme = schemes.remove(0);
1868            let _mailbox3 = setup_and_spawn_actor(
1869                &context,
1870                oracle.manager(),
1871                oracle.control(scheme.public_key()),
1872                scheme,
1873                connections.remove(0),
1874                dummy_consumer(),
1875                prod3,
1876            );
1877
1878            // Fetch keyA multiple times to ensure that Peer2 is blocked.
1879            for _ in 0..20 {
1880                // Fetch keyA
1881                mailbox1.fetch(key_a.clone());
1882
1883                // Wait for success event for keyA
1884                let (key_actual, value) = cons_out1.recv().await.unwrap();
1885                assert_eq!(key_actual, key_a);
1886                assert_eq!(value, valid_data_a);
1887            }
1888
1889            // Fetch keyB
1890            mailbox1.fetch(key_b.clone());
1891
1892            // Wait for some time (longer than retry timeout)
1893            context.sleep(Duration::from_secs(5)).await;
1894
1895            // No success event should be received for keyB since the only peer with valid data is blocked
1896            select! {
1897                _ = cons_out1.recv() => panic!("unexpected event"),
1898                _ = context.sleep(Duration::from_millis(100)) => {},
1899            };
1900
1901            // Prune the fetch for keyB.
1902            let canceled = key_b.clone();
1903            mailbox1.retain(move |key, _| key != &canceled);
1904
1905            // Check oracle
1906            let blocked = oracle.blocked().await.unwrap();
1907            assert_eq!(blocked.len(), 1);
1908            assert_eq!(blocked[0].0, peers[0]);
1909            assert_eq!(blocked[0].1, peers[1]);
1910        });
1911    }
1912
1913    /// Tests that duplicate fetches for the same key are handled properly.
1914    /// The test verifies that when the same key is fetched multiple times,
1915    /// the data is correctly delivered once without errors.
1916    #[test_traced]
1917    fn test_duplicate_fetch_key() {
1918        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1919        executor.start(|context| async move {
1920            let (mut oracle, mut schemes, peers, mut connections) =
1921                setup_network_and_peers(&context, &[1, 2]).await;
1922
1923            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1924
1925            let key = Key(5);
1926            let mut prod2 = Producer::default();
1927            prod2.insert(key.clone(), Bytes::from("data for key 5"));
1928
1929            let (cons1, mut cons_out1) = consumer();
1930
1931            let scheme = schemes.remove(0);
1932            let mut mailbox1 = setup_and_spawn_actor(
1933                &context,
1934                oracle.manager(),
1935                oracle.control(scheme.public_key()),
1936                scheme,
1937                connections.remove(0),
1938                cons1,
1939                Producer::default(),
1940            );
1941
1942            let scheme = schemes.remove(0);
1943            let _mailbox2 = setup_and_spawn_actor(
1944                &context,
1945                oracle.manager(),
1946                oracle.control(scheme.public_key()),
1947                scheme,
1948                connections.remove(0),
1949                dummy_consumer(),
1950                prod2,
1951            );
1952
1953            // Send duplicate fetches for the same key.
1954            mailbox1.fetch(key.clone());
1955            mailbox1.fetch(key.clone());
1956
1957            // Should receive the data only once
1958            let (key_actual, value) = cons_out1.recv().await.unwrap();
1959            assert_eq!(key_actual, key);
1960            assert_eq!(value, Bytes::from("data for key 5"));
1961
1962            // Make sure we don't receive a second event for the duplicate fetch
1963            select! {
1964                _ = cons_out1.recv() => {
1965                    panic!("Unexpected second event received for duplicate fetch");
1966                },
1967                _ = context.sleep(Duration::from_millis(500)) => {
1968                    // This is expected - no additional events should be produced
1969                },
1970            };
1971        });
1972    }
1973
1974    /// Tests that changing peer sets is handled correctly using the update channel.
1975    /// This test verifies that when the peer set changes from peer A to peer B,
1976    /// the resolver correctly adapts and fetches from the new peer.
1977    #[test_traced]
1978    fn test_changing_peer_sets() {
1979        let executor = deterministic::Runner::timed(Duration::from_secs(10));
1980        executor.start(|context| async move {
1981            let (mut oracle, mut schemes, peers, mut connections) =
1982                setup_network_and_peers(&context, &[1, 2, 3]).await;
1983
1984            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
1985            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
1986
1987            let key1 = Key(1);
1988            let key2 = Key(2);
1989
1990            let mut prod2 = Producer::default();
1991            prod2.insert(key1.clone(), Bytes::from("data from peer 2"));
1992
1993            let mut prod3 = Producer::default();
1994            prod3.insert(key2.clone(), Bytes::from("data from peer 3"));
1995
1996            let (cons1, mut cons_out1) = consumer();
1997
1998            let scheme = schemes.remove(0);
1999            let mut mailbox1 = setup_and_spawn_actor(
2000                &context,
2001                oracle.manager(),
2002                oracle.control(scheme.public_key()),
2003                scheme,
2004                connections.remove(0),
2005                cons1,
2006                Producer::default(),
2007            );
2008
2009            let scheme = schemes.remove(0);
2010            let _mailbox2 = setup_and_spawn_actor(
2011                &context,
2012                oracle.manager(),
2013                oracle.control(scheme.public_key()),
2014                scheme,
2015                connections.remove(0),
2016                dummy_consumer(),
2017                prod2,
2018            );
2019
2020            // Fetch key1 from peer 2
2021            mailbox1.fetch(key1.clone());
2022
2023            // Wait for successful fetch
2024            let (key_actual, value) = cons_out1.recv().await.unwrap();
2025            assert_eq!(key_actual, key1);
2026            assert_eq!(value, Bytes::from("data from peer 2"));
2027
2028            // Change peer set to include peer 3
2029            let scheme = schemes.remove(0);
2030            let _mailbox3 = setup_and_spawn_actor(
2031                &context,
2032                oracle.manager(),
2033                oracle.control(scheme.public_key()),
2034                scheme,
2035                connections.remove(0),
2036                dummy_consumer(),
2037                prod3,
2038            );
2039
2040            // Need to wait for the peer set change to propagate
2041            context.sleep(Duration::from_millis(200)).await;
2042
2043            // Fetch key2 from peer 3
2044            mailbox1.fetch(key2.clone());
2045
2046            // Wait for successful fetch
2047            let (key_actual, value) = cons_out1.recv().await.unwrap();
2048            assert_eq!(key_actual, key2);
2049            assert_eq!(value, Bytes::from("data from peer 3"));
2050        });
2051    }
2052
2053    #[test_traced]
2054    fn test_fetch_targeted() {
2055        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2056        executor.start(|context| async move {
2057            let (mut oracle, mut schemes, peers, mut connections) =
2058                setup_network_and_peers(&context, &[1, 2, 3]).await;
2059
2060            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2061            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
2062
2063            let key = Key(1);
2064            let invalid_data = Bytes::from("invalid data");
2065            let valid_data = Bytes::from("valid data");
2066
2067            // Peer 2 has invalid data, peer 3 has valid data
2068            let mut prod2 = Producer::default();
2069            prod2.insert(key.clone(), invalid_data.clone());
2070
2071            let mut prod3 = Producer::default();
2072            prod3.insert(key.clone(), valid_data.clone());
2073
2074            // Consumer expects only valid_data
2075            let (mut cons1, mut cons_out1) = consumer();
2076            cons1.add_expected(key.clone(), valid_data.clone());
2077
2078            let scheme = schemes.remove(0);
2079            let mut mailbox1 = setup_and_spawn_actor(
2080                &context,
2081                oracle.manager(),
2082                oracle.control(scheme.public_key()),
2083                scheme,
2084                connections.remove(0),
2085                cons1,
2086                Producer::default(),
2087            );
2088
2089            let scheme = schemes.remove(0);
2090            let _mailbox2 = setup_and_spawn_actor(
2091                &context,
2092                oracle.manager(),
2093                oracle.control(scheme.public_key()),
2094                scheme,
2095                connections.remove(0),
2096                dummy_consumer(),
2097                prod2,
2098            );
2099
2100            let scheme = schemes.remove(0);
2101            let _mailbox3 = setup_and_spawn_actor(
2102                &context,
2103                oracle.manager(),
2104                oracle.control(scheme.public_key()),
2105                scheme,
2106                connections.remove(0),
2107                dummy_consumer(),
2108                prod3,
2109            );
2110
2111            // Wait for peer set to be established
2112            context.sleep(Duration::from_millis(100)).await;
2113
2114            // Start fetch with targets for both peer 2 (invalid data) and peer 3 (valid data)
2115            // When peer 2 returns invalid data, only peer 2 should be skipped
2116            // Peer 3 should still be tried as a target and succeed
2117            mailbox1.fetch_targeted(
2118                key.clone(),
2119                non_empty_vec![peers[1].clone(), peers[2].clone()],
2120            );
2121
2122            // Should eventually succeed from peer 3
2123            let (key_actual, value) = cons_out1.recv().await.unwrap();
2124            assert_eq!(key_actual, key);
2125            assert_eq!(value, valid_data);
2126
2127            // Verify peer 2 was blocked (sent invalid data)
2128            let blocked = oracle.blocked().await.unwrap();
2129            assert_eq!(blocked.len(), 1);
2130            assert_eq!(blocked[0].0, peers[0]);
2131            assert_eq!(blocked[0].1, peers[1]);
2132
2133            // Verify metrics: 1 successful fetch (from peer 3 after peer 2 was blocked)
2134            let metrics = context.encode();
2135            assert_eq!(
2136                status_metric_total(&metrics, "actor_fetch_total", "Success"),
2137                1
2138            );
2139        });
2140    }
2141
2142    /// A blocker whose blocked-set subscription closes immediately, as test
2143    /// mocks elsewhere in the workspace do.
2144    #[derive(Clone)]
2145    struct ClosedBlocker;
2146
2147    impl Blocker for ClosedBlocker {
2148        type PublicKey = PublicKey;
2149
2150        fn block(&mut self, _peer: Self::PublicKey) -> commonware_actor::Feedback {
2151            commonware_actor::Feedback::Ok
2152        }
2153
2154        fn blocked(&mut self) -> commonware_p2p::BlockedSubscription<Self::PublicKey> {
2155            let (_, receiver) = commonware_utils::channel::ring::channel(NZUsize!(1));
2156            receiver
2157        }
2158    }
2159
2160    #[test_traced]
2161    fn test_closed_blocked_subscription_keeps_fetching() {
2162        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2163        executor.start(|context| async move {
2164            let (mut oracle, mut schemes, peers, mut connections) =
2165                setup_network_and_peers(&context, &[1, 2]).await;
2166
2167            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2168
2169            let key = Key(1);
2170            let data = Bytes::from("data for key 1");
2171            let mut prod2 = Producer::default();
2172            prod2.insert(key.clone(), data.clone());
2173
2174            let (cons1, mut cons_out1) = consumer();
2175
2176            // The engine keeps serving fetches when the blocked-set stream ends.
2177            let scheme = schemes.remove(0);
2178            let mut mailbox1 = setup_and_spawn_actor(
2179                &context,
2180                oracle.manager(),
2181                ClosedBlocker,
2182                scheme,
2183                connections.remove(0),
2184                cons1,
2185                Producer::default(),
2186            );
2187
2188            let scheme = schemes.remove(0);
2189            let _mailbox2 = setup_and_spawn_actor(
2190                &context,
2191                oracle.manager(),
2192                oracle.control(scheme.public_key()),
2193                scheme,
2194                connections.remove(0),
2195                dummy_consumer(),
2196                prod2,
2197            );
2198
2199            mailbox1.fetch(key.clone());
2200            assert_eq!(cons_out1.recv().await.unwrap(), (key, data));
2201        });
2202    }
2203
2204    #[test_traced]
2205    fn test_unblocked_peer_becomes_eligible_again() {
2206        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2207        executor.start(|context| async move {
2208            let (mut oracle, mut schemes, peers, mut connections) =
2209                setup_network_and_peers(&context, &[1, 2]).await;
2210
2211            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2212
2213            let key = Key(1);
2214            let data = Bytes::from("data for key 1");
2215            let mut prod2 = Producer::default();
2216            prod2.insert(key.clone(), data.clone());
2217
2218            // The first response is judged invalid, the second one is accepted.
2219            let (first_gate_sender, first_gate_receiver) = oneshot::channel();
2220            let (second_gate_sender, second_gate_receiver) = oneshot::channel();
2221            let (cons1, mut cons_out1, mut started) = BlockingConsumer::new(
2222                context.child("consumer"),
2223                vec![
2224                    (first_gate_receiver, Outcome::Invalid),
2225                    (second_gate_receiver, Outcome::Complete),
2226                ],
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(
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            mailbox1.fetch_targeted(key.clone(), non_empty_vec![peers[1].clone()]);
2252            assert_eq!(started.recv().await.unwrap(), key);
2253            first_gate_sender.send(()).unwrap();
2254            wait_for_blocked(&context, &oracle, &peers[0], &peers[1]).await;
2255
2256            // The only target is blocked, so the fetch waits.
2257            select! {
2258                _ = started.recv() => panic!("blocked target was retried"),
2259                _ = context.sleep(Duration::from_secs(1)) => {},
2260            };
2261
2262            // Once the network lifts the block, the same target is tried again.
2263            oracle
2264                .unblock(peers[0].clone(), peers[1].clone())
2265                .await
2266                .unwrap();
2267            assert_eq!(started.recv().await.unwrap(), key);
2268            second_gate_sender.send(()).unwrap();
2269            assert_eq!(cons_out1.recv().await.unwrap(), (key, data));
2270        });
2271    }
2272
2273    #[test_traced]
2274    fn test_fetch_targeted_no_fallback() {
2275        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2276        executor.start(|context| async move {
2277            let (mut oracle, mut schemes, peers, mut connections) =
2278                setup_network_and_peers(&context, &[1, 2, 3, 4]).await;
2279
2280            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2281            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
2282            add_link(&mut oracle, LINK.clone(), &peers, 0, 3).await;
2283
2284            let key = Key(1);
2285
2286            // Only peer 4 has the data, peers 2 and 3 don't
2287            let mut prod4 = Producer::default();
2288            prod4.insert(key.clone(), Bytes::from("data from peer 4"));
2289
2290            let (cons1, mut cons_out1) = consumer();
2291
2292            let scheme = schemes.remove(0);
2293            let mut mailbox1 = setup_and_spawn_actor(
2294                &context,
2295                oracle.manager(),
2296                oracle.control(scheme.public_key()),
2297                scheme,
2298                connections.remove(0),
2299                cons1,
2300                Producer::default(),
2301            );
2302
2303            let scheme = schemes.remove(0);
2304            let _mailbox2 = setup_and_spawn_actor(
2305                &context,
2306                oracle.manager(),
2307                oracle.control(scheme.public_key()),
2308                scheme,
2309                connections.remove(0),
2310                dummy_consumer(),
2311                Producer::default(), // no data
2312            );
2313
2314            let scheme = schemes.remove(0);
2315            let _mailbox3 = setup_and_spawn_actor(
2316                &context,
2317                oracle.manager(),
2318                oracle.control(scheme.public_key()),
2319                scheme,
2320                connections.remove(0),
2321                dummy_consumer(),
2322                Producer::default(), // no data
2323            );
2324
2325            let scheme = schemes.remove(0);
2326            let _mailbox4 = setup_and_spawn_actor(
2327                &context,
2328                oracle.manager(),
2329                oracle.control(scheme.public_key()),
2330                scheme,
2331                connections.remove(0),
2332                dummy_consumer(),
2333                prod4,
2334            );
2335
2336            // Wait for peer set to be established
2337            context.sleep(Duration::from_millis(100)).await;
2338
2339            // Start fetch with targets for peers 2 and 3 (both don't have data)
2340            // Peer 4 has data but is NOT a target - it should NEVER be tried
2341            mailbox1.fetch_targeted(
2342                key.clone(),
2343                non_empty_vec![peers[1].clone(), peers[2].clone()],
2344            );
2345
2346            // Wait enough time for targets to fail and retry multiple times
2347            // The fetch should not succeed because peer 4 (which has data) is not targeted
2348            select! {
2349                event = cons_out1.recv() => {
2350                    panic!("Fetch should not succeed, but got: {event:?}");
2351                },
2352                _ = context.sleep(Duration::from_secs(3)) => {
2353                    // Expected: no success event because peer 4 is not targeted
2354                },
2355            };
2356        });
2357    }
2358
2359    #[test_traced]
2360    fn test_fetch_all_targeted() {
2361        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2362        executor.start(|context| async move {
2363            let (mut oracle, mut schemes, peers, mut connections) =
2364                setup_network_and_peers(&context, &[1, 2, 3, 4]).await;
2365
2366            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2367            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
2368            add_link(&mut oracle, LINK.clone(), &peers, 0, 3).await;
2369
2370            let key1 = Key(1);
2371            let key2 = Key(2);
2372            let key3 = Key(3);
2373
2374            // Peer 2 has key1
2375            let mut prod2 = Producer::default();
2376            prod2.insert(key1.clone(), Bytes::from("data for key 1"));
2377
2378            // Peer 3 has key3
2379            let mut prod3 = Producer::default();
2380            prod3.insert(key3.clone(), Bytes::from("data for key 3"));
2381
2382            // Peer 4 has key2
2383            let mut prod4 = Producer::default();
2384            prod4.insert(key2.clone(), Bytes::from("data for key 2"));
2385
2386            // Consumer expects all three keys
2387            let (mut cons1, mut cons_out1) = consumer();
2388            cons1.add_expected(key1.clone(), Bytes::from("data for key 1"));
2389            cons1.add_expected(key2.clone(), Bytes::from("data for key 2"));
2390            cons1.add_expected(key3.clone(), Bytes::from("data for key 3"));
2391
2392            let scheme = schemes.remove(0);
2393            let mut mailbox1 = setup_and_spawn_actor(
2394                &context,
2395                oracle.manager(),
2396                oracle.control(scheme.public_key()),
2397                scheme,
2398                connections.remove(0),
2399                cons1,
2400                Producer::default(),
2401            );
2402
2403            let scheme = schemes.remove(0);
2404            let _mailbox2 = setup_and_spawn_actor(
2405                &context,
2406                oracle.manager(),
2407                oracle.control(scheme.public_key()),
2408                scheme,
2409                connections.remove(0),
2410                dummy_consumer(),
2411                prod2,
2412            );
2413
2414            let scheme = schemes.remove(0);
2415            let _mailbox3 = setup_and_spawn_actor(
2416                &context,
2417                oracle.manager(),
2418                oracle.control(scheme.public_key()),
2419                scheme,
2420                connections.remove(0),
2421                dummy_consumer(),
2422                prod3,
2423            );
2424
2425            let scheme = schemes.remove(0);
2426            let _mailbox4 = setup_and_spawn_actor(
2427                &context,
2428                oracle.manager(),
2429                oracle.control(scheme.public_key()),
2430                scheme,
2431                connections.remove(0),
2432                dummy_consumer(),
2433                prod4,
2434            );
2435
2436            // Wait for peer set to be established
2437            context.sleep(Duration::from_millis(100)).await;
2438
2439            // Fetch keys with mixed targeting:
2440            // - key1 targeted to peer 2 (has data) -> should succeed from target
2441            // - key2 targeted to peer 4 (has data) -> should succeed from target
2442            // - key3 no targeting -> fetched from any peer (peer 3 has it)
2443            mailbox1.fetch_all_targeted(vec![
2444                (key1.clone(), non_empty_vec![peers[1].clone()]), // peer 2 has key1
2445                (key2.clone(), non_empty_vec![peers[3].clone()]), // peer 4 has key2
2446            ]);
2447            mailbox1.fetch(key3.clone()); // no targeting for key3
2448
2449            // Collect all three events
2450            let mut results = HashMap::new();
2451            for _ in 0..3 {
2452                let (key, value) = cons_out1.recv().await.unwrap();
2453                results.insert(key, value);
2454            }
2455
2456            // Verify all keys received correct data
2457            assert_eq!(results.len(), 3);
2458            assert_eq!(results.get(&key1).unwrap(), &Bytes::from("data for key 1"));
2459            assert_eq!(results.get(&key2).unwrap(), &Bytes::from("data for key 2"));
2460            assert_eq!(results.get(&key3).unwrap(), &Bytes::from("data for key 3"));
2461
2462            // Verify metrics: 3 successful fetches
2463            let metrics = context.encode();
2464            assert_eq!(
2465                status_metric_total(&metrics, "actor_fetch_total", "Success"),
2466                3
2467            );
2468        });
2469    }
2470
2471    /// Tests that calling fetch() on an in-progress targeted fetch clears the targets,
2472    /// allowing the fetch to succeed from any available peer.
2473    #[test_traced]
2474    fn test_fetch_clears_targets() {
2475        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2476        executor.start(|context| async move {
2477            let (mut oracle, mut schemes, peers, mut connections) =
2478                setup_network_and_peers(&context, &[1, 2, 3]).await;
2479
2480            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2481            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
2482
2483            let key = Key(1);
2484            let valid_data = Bytes::from("valid data");
2485
2486            // Peer 2 has no data, peer 3 has the data
2487            let mut prod3 = Producer::default();
2488            prod3.insert(key.clone(), valid_data.clone());
2489
2490            let (cons1, mut cons_out1) = consumer();
2491
2492            let scheme = schemes.remove(0);
2493            let mut mailbox1 = setup_and_spawn_actor(
2494                &context,
2495                oracle.manager(),
2496                oracle.control(scheme.public_key()),
2497                scheme,
2498                connections.remove(0),
2499                cons1,
2500                Producer::default(),
2501            );
2502
2503            let scheme = schemes.remove(0);
2504            let _mailbox2 = setup_and_spawn_actor(
2505                &context,
2506                oracle.manager(),
2507                oracle.control(scheme.public_key()),
2508                scheme,
2509                connections.remove(0),
2510                dummy_consumer(),
2511                Producer::default(), // no data
2512            );
2513
2514            let scheme = schemes.remove(0);
2515            let _mailbox3 = setup_and_spawn_actor(
2516                &context,
2517                oracle.manager(),
2518                oracle.control(scheme.public_key()),
2519                scheme,
2520                connections.remove(0),
2521                dummy_consumer(),
2522                prod3,
2523            );
2524
2525            // Wait for peer set to be established
2526            context.sleep(Duration::from_millis(100)).await;
2527
2528            // Start fetch with target for peer 2 only (who doesn't have data)
2529            mailbox1.fetch_targeted(key.clone(), non_empty_vec![peers[1].clone()]);
2530
2531            // Wait for the targeted fetch to fail a few times
2532            context.sleep(Duration::from_millis(500)).await;
2533
2534            // Call fetch() which should clear the targets and allow fallback to any peer
2535            mailbox1.fetch(key.clone());
2536
2537            // Should now succeed from peer 3 (who has data but wasn't originally targeted)
2538            let (key_actual, value) = cons_out1.recv().await.unwrap();
2539            assert_eq!(key_actual, key);
2540            assert_eq!(value, valid_data);
2541        });
2542    }
2543
2544    #[test_traced]
2545    fn test_fetch_targeted_does_not_restrict_all() {
2546        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2547        executor.start(|context| async move {
2548            let (mut oracle, mut schemes, peers, mut connections) =
2549                setup_network_and_peers(&context, &[1, 2, 3]).await;
2550
2551            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2552            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
2553
2554            let key = Key(1);
2555            let valid_data = Bytes::from("valid data");
2556
2557            // Peer 2 has no data, peer 3 has the data
2558            let mut prod3 = Producer::default();
2559            prod3.insert(key.clone(), valid_data.clone());
2560
2561            let (cons1, mut cons_out1) = consumer();
2562
2563            let scheme = schemes.remove(0);
2564            let mut mailbox1 = setup_and_spawn_actor(
2565                &context,
2566                oracle.manager(),
2567                oracle.control(scheme.public_key()),
2568                scheme,
2569                connections.remove(0),
2570                cons1,
2571                Producer::default(),
2572            );
2573
2574            let scheme = schemes.remove(0);
2575            let _mailbox2 = setup_and_spawn_actor(
2576                &context,
2577                oracle.manager(),
2578                oracle.control(scheme.public_key()),
2579                scheme,
2580                connections.remove(0),
2581                dummy_consumer(),
2582                Producer::default(), // no data
2583            );
2584
2585            let scheme = schemes.remove(0);
2586            let _mailbox3 = setup_and_spawn_actor(
2587                &context,
2588                oracle.manager(),
2589                oracle.control(scheme.public_key()),
2590                scheme,
2591                connections.remove(0),
2592                dummy_consumer(),
2593                prod3,
2594            );
2595
2596            // Wait for peer set to be established
2597            context.sleep(Duration::from_millis(100)).await;
2598
2599            // Start fetch without targets (can try any peer)
2600            mailbox1.fetch(key.clone());
2601
2602            // Wait a bit for the fetch to start
2603            context.sleep(Duration::from_millis(50)).await;
2604
2605            // Call fetch_targeted with peer 2 only (who doesn't have data)
2606            // This should NOT restrict the existing "all" fetch
2607            mailbox1.fetch_targeted(key.clone(), non_empty_vec![peers[1].clone()]);
2608
2609            // Should still succeed from peer 3 (who has data but wasn't in the targeted call)
2610            // because the original fetch was "all" and shouldn't be restricted
2611            let (key_actual, value) = cons_out1.recv().await.unwrap();
2612            assert_eq!(key_actual, key);
2613            assert_eq!(value, valid_data);
2614        });
2615    }
2616
2617    #[test_traced]
2618    fn test_retain() {
2619        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2620        executor.start(|context| async move {
2621            let (mut oracle, mut schemes, peers, mut connections) =
2622                setup_network_and_peers(&context, &[1, 2]).await;
2623
2624            let key = Key(5);
2625            let mut prod2 = Producer::default();
2626            prod2.insert(key.clone(), Bytes::from("data for key 5"));
2627
2628            let (cons1, mut cons_out1) = consumer();
2629
2630            let scheme = schemes.remove(0);
2631            let mut mailbox1 = setup_and_spawn_actor(
2632                &context,
2633                oracle.manager(),
2634                oracle.control(scheme.public_key()),
2635                scheme,
2636                connections.remove(0),
2637                cons1,
2638                Producer::default(),
2639            );
2640
2641            let scheme = schemes.remove(0);
2642            let _mailbox2 = setup_and_spawn_actor(
2643                &context,
2644                oracle.manager(),
2645                oracle.control(scheme.public_key()),
2646                scheme,
2647                connections.remove(0),
2648                dummy_consumer(),
2649                prod2,
2650            );
2651
2652            // Retain before fetching should have no effect
2653            mailbox1.retain(|_, _| true);
2654            select! {
2655                _ = cons_out1.recv() => {
2656                    panic!("unexpected event");
2657                },
2658                _ = context.sleep(Duration::from_millis(100)) => {},
2659            };
2660
2661            // Start a fetch (no link, so fetch stays in-flight)
2662            mailbox1.fetch(key.clone());
2663
2664            // Retain with predicate that excludes the key. This must clean up
2665            // the in-flight entry for the key.
2666            let key_clone = key.clone();
2667            mailbox1.retain(move |key, _| key != &key_clone);
2668
2669            // Now add link so fetches can complete
2670            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2671
2672            // Fetch same key again, if the in-flight entry wasn't cleaned up, this would
2673            // be treated as a duplicate and silently ignored
2674            mailbox1.fetch(key.clone());
2675
2676            // Should succeed
2677            let (key_actual, value) = cons_out1.recv().await.unwrap();
2678            assert_eq!(key_actual, key);
2679            assert_eq!(value, Bytes::from("data for key 5"));
2680        });
2681    }
2682
2683    #[test_traced]
2684    fn test_retain_uses_subscribers() {
2685        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2686        executor.start(|context| async move {
2687            let (mut oracle, mut schemes, peers, mut connections) =
2688                setup_network_and_peers(&context, &[1, 2]).await;
2689
2690            let key = Key(5);
2691            let mut prod2 = Producer::default();
2692            prod2.insert(key.clone(), Bytes::from("data for key 5"));
2693
2694            let (cons1, mut cons_out1): (Consumer<Key, Bytes, SubscriberTag>, _) = Consumer::new();
2695
2696            let scheme = schemes.remove(0);
2697            let mut mailbox1 = setup_and_spawn_actor(
2698                &context,
2699                oracle.manager(),
2700                oracle.control(scheme.public_key()),
2701                scheme,
2702                connections.remove(0),
2703                cons1,
2704                Producer::default(),
2705            );
2706
2707            let scheme = schemes.remove(0);
2708            let _mailbox2 = setup_and_spawn_actor(
2709                &context,
2710                oracle.manager(),
2711                oracle.control(scheme.public_key()),
2712                scheme,
2713                connections.remove(0),
2714                dummy_consumer(),
2715                prod2,
2716            );
2717
2718            let dropped_subscriber = SubscriberTag(50);
2719            let kept_subscriber = SubscriberTag(51);
2720            mailbox1.fetch(Fetch {
2721                key: key.clone(),
2722                subscriber: dropped_subscriber,
2723                span: tracing::Span::none(),
2724            });
2725            mailbox1.fetch(Fetch {
2726                key: key.clone(),
2727                subscriber: kept_subscriber.clone(),
2728                span: tracing::Span::none(),
2729            });
2730
2731            context.sleep(Duration::from_millis(100)).await;
2732            mailbox1.retain(move |_, subscriber| subscriber == &kept_subscriber);
2733            context.sleep(Duration::from_millis(100)).await;
2734
2735            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2736
2737            let (key_actual, value) = cons_out1.recv().await.unwrap();
2738            assert_eq!(key_actual, key);
2739            assert_eq!(value, Bytes::from("data for key 5"));
2740        });
2741    }
2742
2743    #[test_traced]
2744    fn test_deliver_receives_subscribers() {
2745        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2746        executor.start(|context| async move {
2747            let (mut oracle, mut schemes, peers, mut connections) =
2748                setup_network_and_peers(&context, &[1, 2]).await;
2749
2750            let key = Key(5);
2751            let mut prod2 = Producer::default();
2752            prod2.insert(key.clone(), Bytes::from("data for key 5"));
2753
2754            let (cons1, mut cons_out1) = SubscriberRecordingConsumer::new();
2755
2756            let scheme = schemes.remove(0);
2757            let mut mailbox1 = setup_and_spawn_actor(
2758                &context,
2759                oracle.manager(),
2760                oracle.control(scheme.public_key()),
2761                scheme,
2762                connections.remove(0),
2763                cons1,
2764                Producer::default(),
2765            );
2766
2767            let scheme = schemes.remove(0);
2768            let _mailbox2 = setup_and_spawn_actor(
2769                &context,
2770                oracle.manager(),
2771                oracle.control(scheme.public_key()),
2772                scheme,
2773                connections.remove(0),
2774                dummy_consumer(),
2775                prod2,
2776            );
2777
2778            let first_subscriber = SubscriberTag(50);
2779            let second_subscriber = SubscriberTag(51);
2780            mailbox1.fetch(Fetch {
2781                key: key.clone(),
2782                subscriber: second_subscriber.clone(),
2783                span: tracing::Span::none(),
2784            });
2785            mailbox1.fetch(Fetch {
2786                key: key.clone(),
2787                subscriber: first_subscriber.clone(),
2788                span: tracing::Span::none(),
2789            });
2790
2791            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2792
2793            let (delivery, value) = cons_out1.recv().await.unwrap();
2794            assert_eq!(
2795                delivery,
2796                Delivery {
2797                    key,
2798                    subscribers: non_empty_vec![
2799                        (first_subscriber, tracing::Span::none()),
2800                        (second_subscriber, tracing::Span::none())
2801                    ],
2802                }
2803            );
2804            assert_eq!(value, Bytes::from("data for key 5"));
2805        });
2806    }
2807
2808    #[test_traced]
2809    fn test_deliver_receives_multiple_subscribers() {
2810        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2811        executor.start(|context| async move {
2812            let (mut oracle, mut schemes, peers, mut connections) =
2813                setup_network_and_peers(&context, &[1, 2]).await;
2814
2815            let key = Key(5);
2816            let mut prod2 = Producer::default();
2817            prod2.insert(key.clone(), Bytes::from("data for key 5"));
2818
2819            let (cons1, mut cons_out1) = SubscriberRecordingConsumer::new();
2820
2821            let scheme = schemes.remove(0);
2822            let mut mailbox1 = setup_and_spawn_actor(
2823                &context,
2824                oracle.manager(),
2825                oracle.control(scheme.public_key()),
2826                scheme,
2827                connections.remove(0),
2828                cons1,
2829                Producer::default(),
2830            );
2831
2832            let scheme = schemes.remove(0);
2833            let _mailbox2 = setup_and_spawn_actor(
2834                &context,
2835                oracle.manager(),
2836                oracle.control(scheme.public_key()),
2837                scheme,
2838                connections.remove(0),
2839                dummy_consumer(),
2840                prod2,
2841            );
2842
2843            let first_subscriber = SubscriberTag(49);
2844            let second_subscriber = SubscriberTag(50);
2845            mailbox1.fetch(Fetch {
2846                key: key.clone(),
2847                subscriber: first_subscriber.clone(),
2848                span: tracing::Span::none(),
2849            });
2850            mailbox1.fetch(Fetch {
2851                key: key.clone(),
2852                subscriber: second_subscriber.clone(),
2853                span: tracing::Span::none(),
2854            });
2855
2856            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2857
2858            let (delivery, value) = cons_out1.recv().await.unwrap();
2859            assert_eq!(
2860                delivery,
2861                Delivery {
2862                    key: key.clone(),
2863                    subscribers: non_empty_vec![
2864                        (first_subscriber, tracing::Span::none()),
2865                        (second_subscriber, tracing::Span::none())
2866                    ],
2867                }
2868            );
2869            assert_eq!(value, Bytes::from("data for key 5"));
2870        });
2871    }
2872
2873    #[test_traced]
2874    fn test_fetch_during_validation_reuses_response_after_success() {
2875        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2876        executor.start(|context| async move {
2877            let (mut oracle, mut schemes, peers, mut connections) =
2878                setup_network_and_peers(&context, &[1, 2]).await;
2879
2880            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2881
2882            let key = Key(5);
2883            let first_response = Bytes::from("data for key 5");
2884            let second_response = Bytes::from("refetched data for key 5");
2885            let mut prod2 = SequencedProducer::default();
2886            prod2.insert(
2887                key.clone(),
2888                [first_response.clone(), second_response.clone()],
2889            );
2890            let prod2_observer = prod2.clone();
2891
2892            let (first_gate_sender, first_gate_receiver) = oneshot::channel();
2893            let (second_gate_sender, second_gate_receiver) = oneshot::channel();
2894            let (cons1, mut deliveries, mut started) = BlockingSubscriberRecordingConsumer::new(
2895                context.child("consumer"),
2896                vec![
2897                    (first_gate_receiver, Outcome::Complete),
2898                    (second_gate_receiver, Outcome::Complete),
2899                ],
2900            );
2901
2902            let scheme = schemes.remove(0);
2903            let mut mailbox1 = setup_and_spawn_actor(
2904                &context,
2905                oracle.manager(),
2906                oracle.control(scheme.public_key()),
2907                scheme,
2908                connections.remove(0),
2909                cons1,
2910                Producer::default(),
2911            );
2912
2913            let scheme = schemes.remove(0);
2914            let _mailbox2 = setup_and_spawn_actor_with_producer(
2915                &context,
2916                oracle.manager(),
2917                oracle.control(scheme.public_key()),
2918                scheme,
2919                connections.remove(0),
2920                dummy_consumer(),
2921                prod2,
2922            );
2923
2924            let first_subscriber = SubscriberTag(49);
2925            let second_subscriber = SubscriberTag(50);
2926            mailbox1.fetch(Fetch {
2927                key: key.clone(),
2928                subscriber: first_subscriber.clone(),
2929                span: tracing::Span::none(),
2930            });
2931
2932            let delivery = started.recv().await.expect("delivery did not start");
2933            assert_eq!(
2934                delivery,
2935                Delivery {
2936                    key: key.clone(),
2937                    subscribers: non_empty_vec![(first_subscriber.clone(), tracing::Span::none())],
2938                }
2939            );
2940
2941            mailbox1.fetch(Fetch {
2942                key: key.clone(),
2943                subscriber: second_subscriber.clone(),
2944                span: tracing::Span::none(),
2945            });
2946            context.sleep(Duration::from_millis(100)).await;
2947            assert_eq!(
2948                prod2_observer.remaining(&key),
2949                vec![second_response.clone()]
2950            );
2951
2952            first_gate_sender.send(()).unwrap();
2953            let (delivery, value) = deliveries.recv().await.expect("consumer channel closed");
2954            assert_eq!(
2955                delivery,
2956                Delivery {
2957                    key: key.clone(),
2958                    subscribers: non_empty_vec![(first_subscriber, tracing::Span::none())],
2959                }
2960            );
2961            assert_eq!(value, first_response);
2962
2963            let delivery = select! {
2964                delivery = started.recv() => delivery.expect("second delivery did not start"),
2965                _ = context.sleep(Duration::from_secs(2)) => {
2966                    panic!("late subscriber was not delivered");
2967                },
2968            };
2969            assert_eq!(
2970                delivery,
2971                Delivery {
2972                    key: key.clone(),
2973                    subscribers: non_empty_vec![(second_subscriber.clone(), tracing::Span::none())],
2974                }
2975            );
2976
2977            second_gate_sender.send(()).unwrap();
2978            let (delivery, value) = deliveries.recv().await.expect("consumer channel closed");
2979            assert_eq!(
2980                delivery,
2981                Delivery {
2982                    key: key.clone(),
2983                    subscribers: non_empty_vec![(second_subscriber, tracing::Span::none())],
2984                }
2985            );
2986            assert_eq!(value, first_response);
2987            assert_eq!(prod2_observer.remaining(&key), vec![second_response]);
2988        });
2989    }
2990
2991    #[test_traced]
2992    fn test_late_targeted_subscriber_joins_retry_after_ambiguous_delivery() {
2993        let executor = deterministic::Runner::timed(Duration::from_secs(10));
2994        executor.start(|context| async move {
2995            let (mut oracle, mut schemes, peers, mut connections) =
2996                setup_network_and_peers(&context, &[1, 2, 3]).await;
2997
2998            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2999
3000            let key = Key(5);
3001            let ambiguous_response = Bytes::from("ambiguous data for key 5");
3002            let unexpected_refetch = Bytes::from("unexpected refetch for key 5");
3003            let mut prod2 = SequencedProducer::default();
3004            prod2.insert(
3005                key.clone(),
3006                [ambiguous_response, unexpected_refetch.clone()],
3007            );
3008            let prod2_observer = prod2.clone();
3009
3010            let valid_response = Bytes::from("valid data for key 5");
3011            let mut prod3 = SequencedProducer::default();
3012            prod3.insert(key.clone(), [valid_response.clone()]);
3013            let prod3_observer = prod3.clone();
3014
3015            let (first_gate_sender, first_gate_receiver) = oneshot::channel();
3016            let (second_gate_sender, second_gate_receiver) = oneshot::channel();
3017            let (cons1, mut deliveries, mut started) = BlockingSubscriberRecordingConsumer::new(
3018                context.child("consumer"),
3019                vec![
3020                    (first_gate_receiver, Outcome::Ambiguous),
3021                    (second_gate_receiver, Outcome::Complete),
3022                ],
3023            );
3024
3025            let scheme = schemes.remove(0);
3026            let mut mailbox1 = setup_and_spawn_actor(
3027                &context,
3028                oracle.manager(),
3029                oracle.control(scheme.public_key()),
3030                scheme,
3031                connections.remove(0),
3032                cons1,
3033                Producer::default(),
3034            );
3035
3036            let scheme = schemes.remove(0);
3037            let _mailbox2 = setup_and_spawn_actor_with_producer(
3038                &context,
3039                oracle.manager(),
3040                oracle.control(scheme.public_key()),
3041                scheme,
3042                connections.remove(0),
3043                dummy_consumer(),
3044                prod2,
3045            );
3046
3047            let scheme = schemes.remove(0);
3048            let _mailbox3 = setup_and_spawn_actor_with_producer(
3049                &context,
3050                oracle.manager(),
3051                oracle.control(scheme.public_key()),
3052                scheme,
3053                connections.remove(0),
3054                dummy_consumer(),
3055                prod3,
3056            );
3057
3058            let first_subscriber = SubscriberTag(49);
3059            let second_subscriber = SubscriberTag(50);
3060
3061            // Start unrestricted repair and park its first response in validation.
3062            mailbox1.fetch(Fetch {
3063                key: key.clone(),
3064                subscriber: first_subscriber.clone(),
3065                span: tracing::Span::none(),
3066            });
3067
3068            let delivery = started.recv().await.expect("delivery did not start");
3069            assert_eq!(
3070                delivery,
3071                Delivery {
3072                    key: key.clone(),
3073                    subscribers: non_empty_vec![(first_subscriber.clone(), tracing::Span::none())],
3074                }
3075            );
3076
3077            // A targeted objection for the same key attaches to the parked fetch
3078            // without issuing another network request.
3079            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
3080            mailbox1.fetch_targeted(
3081                Fetch {
3082                    key: key.clone(),
3083                    subscriber: second_subscriber.clone(),
3084                    span: tracing::Span::none(),
3085                },
3086                non_empty_vec![peers[2].clone()],
3087            );
3088
3089            context.sleep(Duration::from_millis(100)).await;
3090            assert_eq!(
3091                prod2_observer.remaining(&key),
3092                vec![unexpected_refetch.clone()]
3093            );
3094            assert_eq!(prod3_observer.remaining(&key), vec![valid_response.clone()]);
3095
3096            oracle
3097                .remove_link(peers[0].clone(), peers[1].clone())
3098                .await
3099                .unwrap();
3100            oracle
3101                .remove_link(peers[1].clone(), peers[0].clone())
3102                .await
3103                .unwrap();
3104
3105            // An ambiguous response retries unrestricted repair with both
3106            // subscribers and does not penalize the serving peer.
3107            first_gate_sender.send(()).unwrap();
3108            oracle.manager().track(
3109                1,
3110                Set::try_from([peers[0].clone(), peers[2].clone()]).unwrap(),
3111            );
3112
3113            let delivery = select! {
3114                delivery = started.recv() => delivery.expect("retry delivery did not start"),
3115                _ = context.sleep(Duration::from_secs(2)) => {
3116                    panic!("ambiguous response was not retried with the late subscriber");
3117                },
3118            };
3119            assert_eq!(
3120                delivery,
3121                Delivery {
3122                    key: key.clone(),
3123                    subscribers: non_empty_vec![
3124                        (first_subscriber.clone(), tracing::Span::none()),
3125                        (second_subscriber.clone(), tracing::Span::none())
3126                    ],
3127                }
3128            );
3129
3130            second_gate_sender.send(()).unwrap();
3131            let (delivery, value) = deliveries.recv().await.expect("consumer channel closed");
3132            assert_eq!(
3133                delivery,
3134                Delivery {
3135                    key: key.clone(),
3136                    subscribers: non_empty_vec![
3137                        (first_subscriber, tracing::Span::none()),
3138                        (second_subscriber, tracing::Span::none())
3139                    ],
3140                }
3141            );
3142            assert_eq!(value, valid_response);
3143            assert_eq!(prod2_observer.remaining(&key), vec![unexpected_refetch]);
3144            assert!(prod3_observer.remaining(&key).is_empty());
3145            assert!(oracle.blocked().await.unwrap().is_empty());
3146        });
3147    }
3148
3149    #[test_traced]
3150    fn test_due_retry_precedes_queued_fresh_request() {
3151        let executor = deterministic::Runner::timed(Duration::from_secs(10));
3152        executor.start(|context| async move {
3153            let (mut oracle, mut schemes, peers, mut connections) =
3154                setup_network_and_peers_with_rate_limit(
3155                    &context,
3156                    &[1, 2],
3157                    Quota::per_second(NZU32!(2)),
3158                )
3159                .await;
3160            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3161
3162            let retry_key = Key(60);
3163            let fresh_key = Key(61);
3164
3165            let (requester_consumer, mut deliveries) = HoldingConsumer::new();
3166            let requester = schemes.remove(0);
3167            let requester_key = requester.public_key();
3168            let (requester_engine, mut requester_mailbox) = Engine::new(
3169                context.child("requester"),
3170                Config {
3171                    peer_provider: oracle.manager(),
3172                    blocker: oracle.control(requester_key.clone()),
3173                    consumer: requester_consumer,
3174                    producer: Producer::<Key, Bytes>::default(),
3175                    mailbox_size: MAILBOX_SIZE,
3176                    me: Some(requester_key),
3177                    timeout: TIMEOUT,
3178                    fetch_retry_timeout: Duration::ZERO,
3179                    priority_requests: false,
3180                    priority_responses: false,
3181                },
3182            );
3183
3184            let mut producer = Producer::default();
3185            producer.insert(retry_key.clone(), Bytes::from("retry"));
3186            producer.insert(fresh_key.clone(), Bytes::from("fresh"));
3187            let responder = schemes.remove(0);
3188            let _responder_mailbox = setup_and_spawn_actor(
3189                &context,
3190                oracle.manager(),
3191                oracle.control(responder.public_key()),
3192                responder,
3193                connections.remove(1),
3194                dummy_consumer(),
3195                producer,
3196            );
3197
3198            requester_engine.start(connections.remove(0));
3199            requester_mailbox.fetch(retry_key.clone());
3200
3201            let (delivered_key, verdict) =
3202                deliveries.recv().await.expect("requester consumer closed");
3203            assert_eq!(delivered_key, retry_key);
3204
3205            // Resolving the delivery makes its retry due before the queued
3206            // mailbox request is admitted. The two-token quota admits the
3207            // retry and holds the fresh request.
3208            requester_mailbox.fetch(fresh_key.clone());
3209            verdict.send_lossy(Outcome::Ambiguous);
3210
3211            let (delivered_key, verdict) =
3212                deliveries.recv().await.expect("requester consumer closed");
3213            assert_eq!(delivered_key, retry_key);
3214            verdict.send_lossy(Outcome::Complete);
3215
3216            let (delivered_key, verdict) =
3217                deliveries.recv().await.expect("requester consumer closed");
3218            assert_eq!(delivered_key, fresh_key);
3219            verdict.send_lossy(Outcome::Complete);
3220        });
3221    }
3222
3223    #[test_traced]
3224    fn test_late_subscriber_delivery_ignores_unrelated_waiter() {
3225        let executor = deterministic::Runner::timed(Duration::from_secs(10));
3226        executor.start(|context| async move {
3227            let (mut oracle, mut schemes, peers, mut connections) =
3228                setup_network_and_peers(&context, &[1, 2, 3]).await;
3229
3230            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3231            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
3232
3233            let blocked_key = Key(4);
3234            let waiting_key = Key(5);
3235            let main_key = Key(6);
3236            let data = Bytes::from("data for key 6");
3237
3238            let mut prod2 = Producer::default();
3239            prod2.insert(blocked_key.clone(), Bytes::from("bad data"));
3240
3241            let mut prod3 = Producer::default();
3242            prod3.insert(main_key.clone(), data.clone());
3243
3244            let (first_gate_sender, first_gate_receiver) = oneshot::channel();
3245            let (second_gate_sender, second_gate_receiver) = oneshot::channel();
3246            let (cons1, mut deliveries, mut started) = BlockingSubscriberRecordingConsumer::new(
3247                context.child("consumer"),
3248                vec![
3249                    (first_gate_receiver, Outcome::Invalid),
3250                    (second_gate_receiver, Outcome::Complete),
3251                ],
3252            );
3253
3254            let scheme = schemes.remove(0);
3255            let mut mailbox1 = setup_and_spawn_actor(
3256                &context,
3257                oracle.manager(),
3258                oracle.control(scheme.public_key()),
3259                scheme,
3260                connections.remove(0),
3261                cons1,
3262                Producer::default(),
3263            );
3264
3265            let scheme = schemes.remove(0);
3266            let _mailbox2 = setup_and_spawn_actor(
3267                &context,
3268                oracle.manager(),
3269                oracle.control(scheme.public_key()),
3270                scheme,
3271                connections.remove(0),
3272                dummy_consumer(),
3273                prod2,
3274            );
3275
3276            let scheme = schemes.remove(0);
3277            let _mailbox3 = setup_and_spawn_actor(
3278                &context,
3279                oracle.manager(),
3280                oracle.control(scheme.public_key()),
3281                scheme,
3282                connections.remove(0),
3283                dummy_consumer(),
3284                prod3,
3285            );
3286
3287            mailbox1.fetch(Fetch {
3288                key: blocked_key.clone(),
3289                subscriber: SubscriberTag(1),
3290                span: tracing::Span::none(),
3291            });
3292            started
3293                .recv()
3294                .await
3295                .expect("blocking delivery did not start");
3296            first_gate_sender.send(()).unwrap();
3297            wait_for_blocked(&context, &oracle, &peers[0], &peers[1]).await;
3298
3299            mailbox1.fetch_targeted(
3300                Fetch {
3301                    key: waiting_key,
3302                    subscriber: SubscriberTag(2),
3303                    span: tracing::Span::none(),
3304                },
3305                non_empty_vec![peers[1].clone()],
3306            );
3307            context.sleep(Duration::from_millis(100)).await;
3308
3309            let first_subscriber = SubscriberTag(3);
3310            let second_subscriber = SubscriberTag(4);
3311            mailbox1.fetch(Fetch {
3312                key: main_key.clone(),
3313                subscriber: first_subscriber.clone(),
3314                span: tracing::Span::none(),
3315            });
3316
3317            let delivery = started.recv().await.expect("delivery did not start");
3318            assert_eq!(
3319                delivery,
3320                Delivery {
3321                    key: main_key.clone(),
3322                    subscribers: non_empty_vec![(first_subscriber.clone(), tracing::Span::none())],
3323                }
3324            );
3325
3326            mailbox1.fetch(Fetch {
3327                key: main_key.clone(),
3328                subscriber: second_subscriber.clone(),
3329                span: tracing::Span::none(),
3330            });
3331            context.sleep(Duration::from_millis(100)).await;
3332
3333            second_gate_sender.send(()).unwrap();
3334            let (delivery, value) = deliveries.recv().await.expect("consumer channel closed");
3335            assert_eq!(
3336                delivery,
3337                Delivery {
3338                    key: main_key.clone(),
3339                    subscribers: non_empty_vec![(first_subscriber, tracing::Span::none())],
3340                }
3341            );
3342            assert_eq!(value, data);
3343
3344            let delivery = select! {
3345                delivery = started.recv() => delivery.expect("second delivery did not start"),
3346                _ = context.sleep(Duration::from_secs(2)) => {
3347                    panic!("late subscriber was not delivered while an unrelated waiter was armed");
3348                },
3349            };
3350            assert_eq!(
3351                delivery,
3352                Delivery {
3353                    key: main_key,
3354                    subscribers: non_empty_vec![(second_subscriber, tracing::Span::none())],
3355                }
3356            );
3357        });
3358    }
3359
3360    #[test_traced]
3361    fn test_deliver_receives_distinct_subscriber_type() {
3362        let executor = deterministic::Runner::timed(Duration::from_secs(10));
3363        executor.start(|context| async move {
3364            let (mut oracle, mut schemes, peers, mut connections) =
3365                setup_network_and_peers(&context, &[1, 2]).await;
3366
3367            let key = Key(5);
3368            let mut prod2 = Producer::default();
3369            prod2.insert(key.clone(), Bytes::from("data for key 5"));
3370
3371            let (cons1, mut cons_out1) = SubscriberRecordingConsumer::new();
3372
3373            let scheme = schemes.remove(0);
3374            let mut mailbox1 = setup_and_spawn_actor(
3375                &context,
3376                oracle.manager(),
3377                oracle.control(scheme.public_key()),
3378                scheme,
3379                connections.remove(0),
3380                cons1,
3381                Producer::default(),
3382            );
3383
3384            let scheme = schemes.remove(0);
3385            let _mailbox2 = setup_and_spawn_actor(
3386                &context,
3387                oracle.manager(),
3388                oracle.control(scheme.public_key()),
3389                scheme,
3390                connections.remove(0),
3391                dummy_consumer(),
3392                prod2,
3393            );
3394
3395            let subscriber = SubscriberTag(50);
3396            let retained = subscriber.clone();
3397            mailbox1.fetch(Fetch {
3398                key: key.clone(),
3399                subscriber: subscriber.clone(),
3400                span: tracing::Span::none(),
3401            });
3402
3403            context.sleep(Duration::from_millis(100)).await;
3404            mailbox1.retain(move |_, subscriber| subscriber == &retained);
3405            context.sleep(Duration::from_millis(100)).await;
3406
3407            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3408
3409            let (delivery, value) = cons_out1.recv().await.unwrap();
3410            assert_eq!(
3411                delivery,
3412                Delivery {
3413                    key,
3414                    subscribers: non_empty_vec![(subscriber, tracing::Span::none())],
3415                }
3416            );
3417            assert_eq!(value, Bytes::from("data for key 5"));
3418        });
3419    }
3420
3421    #[test_traced]
3422    fn test_deliver_receives_single_subscriber() {
3423        let executor = deterministic::Runner::timed(Duration::from_secs(10));
3424        executor.start(|context| async move {
3425            let (mut oracle, mut schemes, peers, mut connections) =
3426                setup_network_and_peers(&context, &[1, 2]).await;
3427
3428            let key = Key(5);
3429            let mut prod2 = Producer::default();
3430            prod2.insert(key.clone(), Bytes::from("data for key 5"));
3431
3432            let (cons1, mut cons_out1) = SubscriberRecordingConsumer::new();
3433
3434            let scheme = schemes.remove(0);
3435            let mut mailbox1 = setup_and_spawn_actor(
3436                &context,
3437                oracle.manager(),
3438                oracle.control(scheme.public_key()),
3439                scheme,
3440                connections.remove(0),
3441                cons1,
3442                Producer::default(),
3443            );
3444
3445            let scheme = schemes.remove(0);
3446            let _mailbox2 = setup_and_spawn_actor(
3447                &context,
3448                oracle.manager(),
3449                oracle.control(scheme.public_key()),
3450                scheme,
3451                connections.remove(0),
3452                dummy_consumer(),
3453                prod2,
3454            );
3455
3456            let subscriber = SubscriberTag(50);
3457            mailbox1.fetch(Fetch {
3458                key: key.clone(),
3459                subscriber: subscriber.clone(),
3460                span: tracing::Span::none(),
3461            });
3462            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3463
3464            let (delivery, value) = cons_out1.recv().await.unwrap();
3465            assert_eq!(
3466                delivery,
3467                Delivery {
3468                    key,
3469                    subscribers: non_empty_vec![(subscriber, tracing::Span::none())],
3470                }
3471            );
3472            assert_eq!(value, Bytes::from("data for key 5"));
3473        });
3474    }
3475
3476    #[test_traced]
3477    fn test_retain_drops_all() {
3478        let executor = deterministic::Runner::timed(Duration::from_secs(10));
3479        executor.start(|context| async move {
3480            let (mut oracle, mut schemes, peers, mut connections) =
3481                setup_network_and_peers(&context, &[1, 2]).await;
3482
3483            // No link yet - fetch will stay in-flight
3484            let key = Key(6);
3485            let mut prod2 = Producer::default();
3486            prod2.insert(key.clone(), Bytes::from("data for key 6"));
3487
3488            let (cons1, mut cons_out1) = consumer();
3489
3490            let scheme = schemes.remove(0);
3491            let mut mailbox1 = setup_and_spawn_actor(
3492                &context,
3493                oracle.manager(),
3494                oracle.control(scheme.public_key()),
3495                scheme,
3496                connections.remove(0),
3497                cons1,
3498                Producer::default(),
3499            );
3500
3501            let scheme = schemes.remove(0);
3502            let _mailbox2 = setup_and_spawn_actor(
3503                &context,
3504                oracle.manager(),
3505                oracle.control(scheme.public_key()),
3506                scheme,
3507                connections.remove(0),
3508                dummy_consumer(),
3509                prod2,
3510            );
3511
3512            // Pruning before fetching should have no effect.
3513            mailbox1.retain(|_, _| false);
3514            select! {
3515                _ = cons_out1.recv() => {
3516                    panic!("unexpected event");
3517                },
3518                _ = context.sleep(Duration::from_millis(100)) => {},
3519            };
3520
3521            // Start a fetch (no link, so fetch stays in-flight)
3522            mailbox1.fetch(key.clone());
3523
3524            // Prune all fetches.
3525            mailbox1.retain(|_, _| false);
3526
3527            // Now add link so fetches can complete
3528            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3529
3530            // Fetch same key again, if the in-flight entry wasn't cleaned up, this would
3531            // be treated as a duplicate and silently ignored
3532            mailbox1.fetch(key.clone());
3533
3534            // Should succeed
3535            let (key_actual, value) = cons_out1.recv().await.unwrap();
3536            assert_eq!(key_actual, key);
3537            assert_eq!(value, Bytes::from("data for key 6"));
3538        });
3539    }
3540
3541    /// Tests that when a peer is rate-limited, the fetcher spills over to another peer.
3542    /// With 2 peers and rate limit of 1/sec each, 2 requests issued simultaneously should
3543    /// both complete immediately (one to each peer) without waiting for rate limit reset.
3544    #[test_traced]
3545    fn test_rate_limit_spillover() {
3546        let executor = deterministic::Runner::timed(Duration::from_secs(30));
3547        executor.start(|context| async move {
3548            // Use a very restrictive rate limit: 1 request per second per peer
3549            let (mut oracle, mut schemes, peers, mut connections) =
3550                setup_network_and_peers_with_rate_limit(
3551                    &context,
3552                    &[1, 2, 3],
3553                    Quota::per_second(NZU32!(1)),
3554                )
3555                .await;
3556
3557            // Add links between peer 1 and both peer 2 and peer 3
3558            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3559            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
3560
3561            // Both peer 2 and peer 3 have the same data
3562            let mut prod2 = Producer::default();
3563            let mut prod3 = Producer::default();
3564            prod2.insert(Key(0), Bytes::from("data for key 0"));
3565            prod2.insert(Key(1), Bytes::from("data for key 1"));
3566            prod3.insert(Key(0), Bytes::from("data for key 0"));
3567            prod3.insert(Key(1), Bytes::from("data for key 1"));
3568
3569            let (cons1, mut cons_out1) = consumer();
3570
3571            // Set up peer 1 (the requester)
3572            let scheme = schemes.remove(0);
3573            let mut mailbox1 = setup_and_spawn_actor(
3574                &context,
3575                oracle.manager(),
3576                oracle.control(scheme.public_key()),
3577                scheme,
3578                connections.remove(0),
3579                cons1,
3580                Producer::default(),
3581            );
3582
3583            // Set up peer 2 (has data)
3584            let scheme = schemes.remove(0);
3585            let _mailbox2 = setup_and_spawn_actor(
3586                &context,
3587                oracle.manager(),
3588                oracle.control(scheme.public_key()),
3589                scheme,
3590                connections.remove(0),
3591                dummy_consumer(),
3592                prod2,
3593            );
3594
3595            // Set up peer 3 (also has data)
3596            let scheme = schemes.remove(0);
3597            let _mailbox3 = setup_and_spawn_actor(
3598                &context,
3599                oracle.manager(),
3600                oracle.control(scheme.public_key()),
3601                scheme,
3602                connections.remove(0),
3603                dummy_consumer(),
3604                prod3,
3605            );
3606
3607            // Wait for peer set to be established
3608            context.sleep(Duration::from_millis(100)).await;
3609            let start = context.current();
3610
3611            // Issue 2 fetches rapidly.
3612            // With rate limit of 1/sec per peer and 2 peers, both should complete
3613            // immediately via spill-over (one request to each peer)
3614            mailbox1.fetch(Key(0));
3615            mailbox1.fetch(Key(1));
3616
3617            // Collect results
3618            let mut results = HashMap::new();
3619            for _ in 0..2 {
3620                let (key, value) = cons_out1.recv().await.unwrap();
3621                results.insert(key.clone(), value);
3622            }
3623
3624            // Verify both keys were fetched successfully
3625            assert_eq!(results.len(), 2);
3626            assert_eq!(
3627                results.get(&Key(0)).unwrap(),
3628                &Bytes::from("data for key 0")
3629            );
3630            assert_eq!(
3631                results.get(&Key(1)).unwrap(),
3632                &Bytes::from("data for key 1")
3633            );
3634
3635            // Verify it completed quickly (well under 1 second) - proves spill-over worked
3636            // Without spill-over, the second request would wait ~1 second for rate limit reset
3637            let elapsed = context.current().duration_since(start).unwrap();
3638            assert!(
3639                elapsed < Duration::from_millis(500),
3640                "Expected quick completion via spill-over, but took {elapsed:?}"
3641            );
3642        });
3643    }
3644
3645    /// Tests that rate limiting causes retries to eventually succeed after the rate limit resets.
3646    /// This test uses a single peer with a restrictive rate limit and verifies that
3647    /// fetches eventually complete after waiting for the rate limit to reset.
3648    #[test_traced]
3649    fn test_rate_limit_retry_after_reset() {
3650        let executor = deterministic::Runner::timed(Duration::from_secs(30));
3651        executor.start(|context| async move {
3652            // Use a restrictive rate limit: 1 request per second
3653            let (mut oracle, mut schemes, peers, mut connections) =
3654                setup_network_and_peers_with_rate_limit(
3655                    &context,
3656                    &[1, 2],
3657                    Quota::per_second(NZU32!(1)),
3658                )
3659                .await;
3660
3661            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3662
3663            // Peer 2 has data for multiple keys
3664            let mut prod2 = Producer::default();
3665            prod2.insert(Key(1), Bytes::from("data for key 1"));
3666            prod2.insert(Key(2), Bytes::from("data for key 2"));
3667            prod2.insert(Key(3), Bytes::from("data for key 3"));
3668
3669            let (cons1, mut cons_out1) = consumer();
3670
3671            let scheme = schemes.remove(0);
3672            let mut mailbox1 = setup_and_spawn_actor(
3673                &context,
3674                oracle.manager(),
3675                oracle.control(scheme.public_key()),
3676                scheme,
3677                connections.remove(0),
3678                cons1,
3679                Producer::default(),
3680            );
3681
3682            let scheme = schemes.remove(0);
3683            let _mailbox2 = setup_and_spawn_actor(
3684                &context,
3685                oracle.manager(),
3686                oracle.control(scheme.public_key()),
3687                scheme,
3688                connections.remove(0),
3689                dummy_consumer(),
3690                prod2,
3691            );
3692
3693            // Wait for peer set to be established
3694            context.sleep(Duration::from_millis(100)).await;
3695            let start = context.current();
3696
3697            // Issue 3 fetches to a single peer with rate limit of 1/sec.
3698            // Only 1 can be sent immediately, the others must wait for rate limit reset
3699            mailbox1.fetch(Key(1));
3700            mailbox1.fetch(Key(2));
3701            mailbox1.fetch(Key(3));
3702
3703            // All 3 should eventually succeed (after rate limit resets)
3704            let mut results = HashMap::new();
3705            for _ in 0..3 {
3706                let (key, value) = cons_out1.recv().await.unwrap();
3707                results.insert(key.clone(), value);
3708            }
3709
3710            assert_eq!(results.len(), 3);
3711            for i in 1..=3 {
3712                assert_eq!(
3713                    results.get(&Key(i)).unwrap(),
3714                    &Bytes::from(format!("data for key {}", i))
3715                );
3716            }
3717
3718            // Verify it took significant time due to rate limiting
3719            // With 3 requests at 1/sec to a single peer, requests 2 and 3 must wait
3720            // for rate limit resets (~1 second each), so total should be > 2 seconds
3721            let elapsed = context.current().duration_since(start).unwrap();
3722            assert!(
3723                elapsed > Duration::from_secs(2),
3724                "Expected rate limiting to cause delay > 2s, but took {elapsed:?}"
3725            );
3726        });
3727    }
3728
3729    /// Tests that the resolver never sends fetches to itself (me exclusion).
3730    /// Even when the local peer has the data in its producer, it should fetch from
3731    /// another peer instead.
3732    #[test_traced]
3733    fn test_self_exclusion() {
3734        let executor = deterministic::Runner::timed(Duration::from_secs(10));
3735        executor.start(|context| async move {
3736            let (mut oracle, mut schemes, peers, mut connections) =
3737                setup_network_and_peers(&context, &[1, 2]).await;
3738
3739            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3740
3741            let key = Key(1);
3742            let data = Bytes::from("shared data");
3743
3744            // Both peers have the data - peer 1 (requester) and peer 2
3745            let mut prod1 = Producer::default();
3746            prod1.insert(key.clone(), data.clone());
3747            let mut prod2 = Producer::default();
3748            prod2.insert(key.clone(), data.clone());
3749
3750            let (cons1, mut cons_out1) = consumer();
3751
3752            // Set up peer 1 with `me` set - it has the data but should NOT fetch from itself
3753            let scheme = schemes.remove(0);
3754            let mut mailbox1 = setup_and_spawn_actor(
3755                &context,
3756                oracle.manager(),
3757                oracle.control(scheme.public_key()),
3758                scheme,
3759                connections.remove(0),
3760                cons1,
3761                prod1, // peer 1 has the data
3762            );
3763
3764            // Set up peer 2 - also has the data
3765            let scheme = schemes.remove(0);
3766            let _mailbox2 = setup_and_spawn_actor(
3767                &context,
3768                oracle.manager(),
3769                oracle.control(scheme.public_key()),
3770                scheme,
3771                connections.remove(0),
3772                dummy_consumer(),
3773                prod2,
3774            );
3775
3776            // Wait for peer set to be established
3777            context.sleep(Duration::from_millis(100)).await;
3778
3779            // Fetch the key - should get it from peer 2, not from self
3780            mailbox1.fetch(key.clone());
3781
3782            // Should succeed (from peer 2)
3783            let (key_actual, value) = cons_out1.recv().await.unwrap();
3784            assert_eq!(key_actual, key);
3785            assert_eq!(value, data);
3786        });
3787    }
3788
3789    #[test_traced]
3790    fn test_fetch_uses_primary_peers_only() {
3791        let executor = deterministic::Runner::timed(Duration::from_secs(10));
3792        executor.start(|context| async move {
3793            let (network, oracle) = Network::new(
3794                context.child("network"),
3795                commonware_p2p::simulated::Config {
3796                    max_size: 1024 * 1024,
3797                    max_peers_per_set: NZUsize!(2),
3798                    disconnect_on_block: true,
3799                    tracked_peer_sets: NZUsize!(1),
3800                },
3801            );
3802            network.start();
3803
3804            let schemes: Vec<PrivateKey> = [1u64, 2, 3]
3805                .into_iter()
3806                .map(PrivateKey::from_seed)
3807                .collect();
3808            let peers: Vec<PublicKey> = schemes.iter().map(|s| s.public_key()).collect();
3809            let mut schemes = schemes;
3810
3811            let mut connections = Vec::new();
3812            for peer in &peers {
3813                let (sender, receiver) = oracle
3814                    .control(peer.clone())
3815                    .register(0, Quota::per_second(RATE_LIMIT))
3816                    .await
3817                    .unwrap();
3818                connections.push((sender, receiver));
3819            }
3820
3821            // Topology: peer 1 (requester) linked to peers 2 and 3.
3822            // Peer 2 is primary (no data), peer 3 is secondary (has data).
3823            // Fetch should only query primary peers, so the request must time out.
3824            let mut oracle = oracle;
3825            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3826            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
3827
3828            oracle.manager().track(
3829                1,
3830                TrackedPeers::new(
3831                    Set::try_from([peers[1].clone()]).unwrap(),
3832                    Set::try_from([peers[2].clone()]).unwrap(),
3833                ),
3834            );
3835            context.sleep(Duration::from_millis(100)).await;
3836
3837            let key = Key(1);
3838            let data = Bytes::from("secondary only data");
3839
3840            let (cons1, mut cons_out1) = consumer();
3841
3842            // Peer 1: the requester, has no data.
3843            let scheme = schemes.remove(0);
3844            let mut mailbox1 = setup_and_spawn_actor(
3845                &context,
3846                oracle.manager(),
3847                oracle.control(scheme.public_key()),
3848                scheme,
3849                connections.remove(0),
3850                cons1,
3851                Producer::default(),
3852            );
3853
3854            // Peer 2: primary, has no data.
3855            let scheme = schemes.remove(0);
3856            let _mailbox2 = setup_and_spawn_actor(
3857                &context,
3858                oracle.manager(),
3859                oracle.control(scheme.public_key()),
3860                scheme,
3861                connections.remove(0),
3862                dummy_consumer(),
3863                Producer::default(),
3864            );
3865
3866            // Peer 3: secondary, has the data. Should not be queried.
3867            let mut prod3 = Producer::default();
3868            prod3.insert(key.clone(), data);
3869            let scheme = schemes.remove(0);
3870            let _mailbox3 = setup_and_spawn_actor(
3871                &context,
3872                oracle.manager(),
3873                oracle.control(scheme.public_key()),
3874                scheme,
3875                connections.remove(0),
3876                dummy_consumer(),
3877                prod3,
3878            );
3879
3880            // Fetch should time out because the only peer with data (peer 3)
3881            // is secondary and won't be queried.
3882            mailbox1.fetch(key.clone());
3883
3884            select! {
3885                event = cons_out1.recv() => {
3886                    panic!("fetch should not succeed from a secondary peer, got: {event:?}");
3887                },
3888                _ = context.sleep(Duration::from_secs(2)) => {},
3889            }
3890        });
3891    }
3892
3893    #[test_traced]
3894    fn test_fetch_uses_latest_primary_set_only() {
3895        let executor = deterministic::Runner::timed(Duration::from_secs(10));
3896        executor.start(|context| async move {
3897            let (network, oracle) = Network::new(
3898                context.child("network"),
3899                commonware_p2p::simulated::Config {
3900                    max_size: 1024 * 1024,
3901                    max_peers_per_set: NZUsize!(2),
3902                    disconnect_on_block: true,
3903                    tracked_peer_sets: NZUsize!(2),
3904                },
3905            );
3906            network.start();
3907
3908            let schemes: Vec<PrivateKey> = [1u64, 2, 3]
3909                .into_iter()
3910                .map(PrivateKey::from_seed)
3911                .collect();
3912            let peers: Vec<PublicKey> = schemes.iter().map(|s| s.public_key()).collect();
3913            let mut schemes = schemes;
3914
3915            let mut connections = Vec::new();
3916            for peer in &peers {
3917                let (sender, receiver) = oracle
3918                    .control(peer.clone())
3919                    .register(0, Quota::per_second(RATE_LIMIT))
3920                    .await
3921                    .unwrap();
3922                connections.push((sender, receiver));
3923            }
3924
3925            let mut oracle = oracle;
3926            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3927            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
3928
3929            // Keep the requester tracked across the cutover so the fetch path itself remains
3930            // active, while peer 2 is retained only through the overlap window after peer 3
3931            // becomes the newest primary set.
3932            oracle
3933                .manager()
3934                .track(
3935                    0,
3936                    Set::try_from([peers[0].clone(), peers[1].clone()]).unwrap(),
3937                );
3938            context.sleep(Duration::from_millis(100)).await;
3939
3940            let key = Key(7);
3941            let targeted_key = Key(8);
3942            let data = Bytes::from("old primary data");
3943
3944            let (cons1, mut cons_out1) = consumer();
3945
3946            // Peer 1: requester.
3947            let scheme = schemes.remove(0);
3948            let mut mailbox1 = setup_and_spawn_actor(
3949                &context,
3950                oracle.manager(),
3951                oracle.control(scheme.public_key()),
3952                scheme,
3953                connections.remove(0),
3954                cons1,
3955                Producer::default(),
3956            );
3957
3958            // Peer 2: old primary, still retained in `all.primary`, has the data.
3959            let mut prod2 = Producer::default();
3960            prod2.insert(key.clone(), data.clone());
3961            prod2.insert(targeted_key.clone(), data);
3962            let scheme = schemes.remove(0);
3963            let _mailbox2 = setup_and_spawn_actor(
3964                &context,
3965                oracle.manager(),
3966                oracle.control(scheme.public_key()),
3967                scheme,
3968                connections.remove(0),
3969                dummy_consumer(),
3970                prod2,
3971            );
3972
3973            // Peer 3: latest primary, has no data.
3974            let scheme = schemes.remove(0);
3975            let _mailbox3 = setup_and_spawn_actor(
3976                &context,
3977                oracle.manager(),
3978                oracle.control(scheme.public_key()),
3979                scheme,
3980                connections.remove(0),
3981                dummy_consumer(),
3982                Producer::default(),
3983            );
3984
3985            context.sleep(Duration::from_millis(100)).await;
3986
3987            // Track peer 3 as the latest primary while keeping the requester in the peer set.
3988            // Peer 2 remains in the provider's overlap window (`all.primary`), but new resolver traffic
3989            // should use only `latest.primary`.
3990            oracle
3991                .manager()
3992                .track(
3993                    1,
3994                    Set::try_from([peers[0].clone(), peers[2].clone()]).unwrap(),
3995                );
3996            context.sleep(Duration::from_millis(100)).await;
3997
3998            mailbox1.fetch(key);
3999
4000            select! {
4001                event = cons_out1.recv() => {
4002                    panic!(
4003                        "fetch should not succeed from an old primary retained only in the overlap window, got: {event:?}"
4004                    );
4005                },
4006                _ = context.sleep(Duration::from_secs(1)) => {},
4007            }
4008
4009            // Explicit targets still respect the latest-primary filter.
4010            mailbox1
4011                .fetch_targeted(targeted_key, non_empty_vec![peers[1].clone()]);
4012
4013            select! {
4014                event = cons_out1.recv() => {
4015                    panic!(
4016                        "targeted fetch should not bypass the latest-primary filter, got: {event:?}"
4017                    );
4018                },
4019                _ = context.sleep(Duration::from_secs(1)) => {},
4020            }
4021        });
4022    }
4023
4024    #[test_traced]
4025    fn test_fetch_after_cutover_relies_on_latest_primary_history() {
4026        let executor = deterministic::Runner::timed(Duration::from_secs(10));
4027        executor.start(|context| async move {
4028            let (network, oracle) = Network::new(
4029                context.child("network"),
4030                commonware_p2p::simulated::Config {
4031                    max_size: 1024 * 1024,
4032                    max_peers_per_set: NZUsize!(2),
4033                    disconnect_on_block: true,
4034                    tracked_peer_sets: NZUsize!(2),
4035                },
4036            );
4037            network.start();
4038
4039            let schemes: Vec<PrivateKey> = [1u64, 2, 3]
4040                .into_iter()
4041                .map(PrivateKey::from_seed)
4042                .collect();
4043            let peers: Vec<PublicKey> = schemes.iter().map(|s| s.public_key()).collect();
4044            let mut schemes = schemes;
4045
4046            let mut connections = Vec::new();
4047            for peer in &peers {
4048                let (sender, receiver) = oracle
4049                    .control(peer.clone())
4050                    .register(0, Quota::per_second(RATE_LIMIT))
4051                    .await
4052                    .unwrap();
4053                connections.push((sender, receiver));
4054            }
4055
4056            let mut oracle = oracle;
4057            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
4058            add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
4059
4060            // Keep the requester in the peer set across the cutover while peer 2 remains connected
4061            // only through the overlap window after the latest primary advances to peer 3.
4062            oracle.manager().track(
4063                0,
4064                Set::try_from([peers[0].clone(), peers[1].clone()]).unwrap(),
4065            );
4066            context.sleep(Duration::from_millis(100)).await;
4067
4068            let key = Key(9);
4069            let invalid_history = Bytes::from("stale overlap history");
4070            let valid_history = Bytes::from("latest primary history");
4071
4072            let (mut cons1, mut cons_out1) = consumer();
4073            cons1.add_expected(key.clone(), valid_history.clone());
4074
4075            // Peer 1: requester.
4076            let scheme = schemes.remove(0);
4077            let mut mailbox1 = setup_and_spawn_actor(
4078                &context,
4079                oracle.manager(),
4080                oracle.control(scheme.public_key()),
4081                scheme,
4082                connections.remove(0),
4083                cons1,
4084                Producer::default(),
4085            );
4086
4087            // Peer 2: old primary retained only via overlap. If queried, it would be blocked for
4088            // serving invalid history.
4089            let mut prod2 = Producer::default();
4090            prod2.insert(key.clone(), invalid_history);
4091            let scheme = schemes.remove(0);
4092            let _mailbox2 = setup_and_spawn_actor(
4093                &context,
4094                oracle.manager(),
4095                oracle.control(scheme.public_key()),
4096                scheme,
4097                connections.remove(0),
4098                dummy_consumer(),
4099                prod2,
4100            );
4101
4102            // Peer 3: latest primary and the only peer that should satisfy the fetch.
4103            let mut prod3 = Producer::default();
4104            prod3.insert(key.clone(), valid_history.clone());
4105            let scheme = schemes.remove(0);
4106            let _mailbox3 = setup_and_spawn_actor(
4107                &context,
4108                oracle.manager(),
4109                oracle.control(scheme.public_key()),
4110                scheme,
4111                connections.remove(0),
4112                dummy_consumer(),
4113                prod3,
4114            );
4115
4116            context.sleep(Duration::from_millis(100)).await;
4117
4118            oracle.manager().track(
4119                1,
4120                Set::try_from([peers[0].clone(), peers[2].clone()]).unwrap(),
4121            );
4122            context.sleep(Duration::from_millis(100)).await;
4123
4124            mailbox1.fetch(key.clone());
4125
4126            let (key_actual, value) = cons_out1.recv().await.unwrap();
4127            assert_eq!(key_actual, key);
4128            assert_eq!(value, valid_history);
4129
4130            assert!(
4131                oracle.blocked().await.unwrap().is_empty(),
4132                "overlap-only peers should not be queried for post-cutover history"
4133            );
4134        });
4135    }
4136
4137    #[test_traced]
4138    fn test_secondary_peer_requests_are_served() {
4139        let executor = deterministic::Runner::timed(Duration::from_secs(10));
4140        executor.start(|context| async move {
4141            let (mut oracle, mut schemes, peers, mut connections) =
4142                setup_network_and_peers(&context, &[1, 2]).await;
4143
4144            // Topology: peer 1 is primary (has data), peer 2 is secondary (requester).
4145            // Verifies that a primary peer serves requests from secondary peers
4146            // (i.e. secondary peers can't fetch via broadcast, but their direct
4147            // requests are still answered).
4148            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
4149
4150            oracle.manager().track(
4151                1,
4152                TrackedPeers::new(
4153                    Set::try_from([peers[0].clone()]).unwrap(),
4154                    Set::try_from([peers[1].clone()]).unwrap(),
4155                ),
4156            );
4157            context.sleep(Duration::from_millis(100)).await;
4158
4159            let key = Key(9);
4160            let data = Bytes::from("served to secondary");
4161
4162            // Peer 1: primary, has the data.
4163            let mut prod1 = Producer::default();
4164            prod1.insert(key.clone(), data.clone());
4165
4166            let scheme = schemes.remove(0);
4167            let _mailbox1 = setup_and_spawn_actor(
4168                &context,
4169                oracle.manager(),
4170                oracle.control(scheme.public_key()),
4171                scheme,
4172                connections.remove(0),
4173                dummy_consumer(),
4174                prod1,
4175            );
4176
4177            // Peer 2: secondary, uses fetch_targeted to explicitly request from peer 1.
4178            let (mut cons2, mut cons_out2) = consumer();
4179            cons2.add_expected(key.clone(), data.clone());
4180            let scheme = schemes.remove(0);
4181            let mut mailbox2 = setup_and_spawn_actor(
4182                &context,
4183                oracle.manager(),
4184                oracle.control(scheme.public_key()),
4185                scheme,
4186                connections.remove(0),
4187                cons2,
4188                Producer::default(),
4189            );
4190
4191            mailbox2.fetch_targeted(key.clone(), non_empty_vec![peers[0].clone()]);
4192
4193            let (key_actual, value) = cons_out2.recv().await.unwrap();
4194            assert_eq!(key_actual, key);
4195            assert_eq!(value, data);
4196        });
4197    }
4198
4199    #[test_traced]
4200    fn test_shutdown_aborts_pending_delivery_without_leaked_tasks() {
4201        let executor = deterministic::Runner::timed(Duration::from_secs(10));
4202        executor.start(|context| async move {
4203            let (mut oracle, mut schemes, peers, mut connections) =
4204                setup_network_and_peers(&context, &[1, 2]).await;
4205
4206            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
4207
4208            let key = Key(1);
4209            let data = Bytes::from("data for key 1");
4210            let mut prod2 = Producer::default();
4211            prod2.insert(key.clone(), data);
4212
4213            let (mut gate_sender, gate_receiver) = oneshot::channel();
4214            let (cons1, mut cons_out1, mut started) = BlockingConsumer::new(
4215                context.child("consumer"),
4216                vec![(gate_receiver, Outcome::Complete)],
4217            );
4218
4219            let actor_context = context.child("actor");
4220
4221            let scheme = schemes.remove(0);
4222            let public_key = scheme.public_key();
4223            let (engine, mut mailbox1): (_, Mailbox<Key, PublicKey>) = Engine::new(
4224                actor_context.child("peer").with_attribute("index", 0),
4225                Config {
4226                    peer_provider: oracle.manager(),
4227                    blocker: oracle.control(public_key.clone()),
4228                    consumer: cons1,
4229                    producer: Producer::<Key, Bytes>::default(),
4230                    mailbox_size: MAILBOX_SIZE,
4231                    me: Some(public_key),
4232                    timeout: TIMEOUT,
4233                    fetch_retry_timeout: FETCH_RETRY_TIMEOUT,
4234                    priority_requests: false,
4235                    priority_responses: false,
4236                },
4237            );
4238            let handle1 = engine.start(connections.remove(0));
4239
4240            let scheme = schemes.remove(0);
4241            let public_key = scheme.public_key();
4242            let (engine, _mailbox2): (_, Mailbox<Key, PublicKey>) = Engine::new(
4243                actor_context.child("peer").with_attribute("index", 1),
4244                Config {
4245                    peer_provider: oracle.manager(),
4246                    blocker: oracle.control(public_key.clone()),
4247                    consumer: dummy_consumer(),
4248                    producer: prod2,
4249                    mailbox_size: MAILBOX_SIZE,
4250                    me: Some(public_key),
4251                    timeout: TIMEOUT,
4252                    fetch_retry_timeout: FETCH_RETRY_TIMEOUT,
4253                    priority_requests: false,
4254                    priority_responses: false,
4255                },
4256            );
4257            let handle2 = engine.start(connections.remove(0));
4258
4259            mailbox1.fetch(key.clone());
4260            let started_key = started.recv().await.expect("delivery did not start");
4261            assert_eq!(started_key, key);
4262
4263            assert!(count_running_tasks(&context, "actor") > 0);
4264
4265            handle1.abort();
4266            handle2.abort();
4267
4268            context.sleep(Duration::from_millis(100)).await;
4269
4270            select! {
4271                _ = gate_sender.closed() => {},
4272                _ = context.sleep(Duration::from_secs(2)) => {
4273                    panic!("pending delivery was not aborted");
4274                },
4275            };
4276
4277            select! {
4278                event = cons_out1.recv() => assert!(event.is_none(), "unexpected event"),
4279                _ = context.sleep(Duration::from_millis(100)) => {},
4280            };
4281
4282            let running_after = count_running_tasks(&context, "actor");
4283            assert_eq!(
4284                running_after, 0,
4285                "all actor tasks should be stopped, but {running_after} still running"
4286            );
4287        });
4288    }
4289
4290    #[allow(clippy::type_complexity)]
4291    fn spawn_actors_with_handles(
4292        context: &deterministic::Context,
4293        oracle: &Oracle<PublicKey, deterministic::Context>,
4294        schemes: Vec<PrivateKey>,
4295        connections: Vec<(
4296            Sender<PublicKey, deterministic::Context>,
4297            Receiver<PublicKey>,
4298        )>,
4299        consumers: Vec<Consumer<Key, Bytes>>,
4300        producers: Vec<Producer<Key, Bytes>>,
4301    ) -> (
4302        Vec<Mailbox<Key, PublicKey>>,
4303        Vec<commonware_runtime::Handle<()>>,
4304    ) {
4305        let actor_context = context.child("actor");
4306        let mut mailboxes = Vec::new();
4307        let mut handles = Vec::new();
4308
4309        for (idx, ((scheme, conn), (consumer, producer))) in schemes
4310            .into_iter()
4311            .zip(connections)
4312            .zip(consumers.into_iter().zip(producers))
4313            .enumerate()
4314        {
4315            let ctx = actor_context.child("peer").with_attribute("index", idx);
4316            let public_key = scheme.public_key();
4317            let (engine, mailbox) = Engine::new(
4318                ctx,
4319                Config {
4320                    peer_provider: oracle.manager(),
4321                    blocker: oracle.control(public_key.clone()),
4322                    consumer,
4323                    producer,
4324                    mailbox_size: MAILBOX_SIZE,
4325                    me: Some(public_key),
4326                    timeout: TIMEOUT,
4327                    fetch_retry_timeout: FETCH_RETRY_TIMEOUT,
4328                    priority_requests: false,
4329                    priority_responses: false,
4330                },
4331            );
4332            handles.push(engine.start(conn));
4333            mailboxes.push(mailbox);
4334        }
4335
4336        (mailboxes, handles)
4337    }
4338
4339    #[test_traced]
4340    fn test_operations_after_shutdown_do_not_panic() {
4341        let executor = deterministic::Runner::timed(Duration::from_secs(10));
4342        executor.start(|context| async move {
4343            let (mut oracle, schemes, peers, connections) =
4344                setup_network_and_peers(&context, &[1, 2]).await;
4345
4346            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
4347
4348            let key = Key(1);
4349            let mut prod2 = Producer::default();
4350            prod2.insert(key.clone(), Bytes::from("data for key 1"));
4351
4352            let (cons1, mut cons_out1) = consumer();
4353
4354            let (mut mailboxes, handles) = spawn_actors_with_handles(
4355                &context,
4356                &oracle,
4357                schemes,
4358                connections,
4359                vec![cons1, dummy_consumer()],
4360                vec![Producer::default(), prod2],
4361            );
4362
4363            // Fetch to verify network is functional
4364            mailboxes[0].fetch(key.clone());
4365            let (_, value) = cons_out1.recv().await.unwrap();
4366            assert_eq!(value, Bytes::from("data for key 1"));
4367
4368            // Abort all actors
4369            for handle in handles {
4370                handle.abort();
4371            }
4372            context.sleep(Duration::from_millis(100)).await;
4373
4374            // All operations should not panic after shutdown
4375
4376            // Fetch should not panic
4377            let key2 = Key(2);
4378            mailboxes[0].fetch(key2.clone());
4379
4380            // Retain can prune a single key after shutdown without panicking.
4381            let canceled = key2;
4382            mailboxes[0].retain(move |key, _| key != &canceled);
4383
4384            // Retain should not panic
4385            mailboxes[0].retain(|_, _| true);
4386
4387            // Fetch targeted should not panic
4388            mailboxes[0].fetch_targeted(Key(3), non_empty_vec![peers[1].clone()]);
4389        });
4390    }
4391
4392    fn clean_shutdown(seed: u64) {
4393        let cfg = deterministic::Config::default()
4394            .with_seed(seed)
4395            .with_timeout(Some(Duration::from_secs(30)));
4396        let executor = deterministic::Runner::new(cfg);
4397        executor.start(|context| async move {
4398            let (mut oracle, schemes, peers, connections) =
4399                setup_network_and_peers(&context, &[1, 2]).await;
4400
4401            add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
4402
4403            let key = Key(1);
4404            let mut prod2 = Producer::default();
4405            prod2.insert(key.clone(), Bytes::from("data for key 1"));
4406
4407            let (cons1, mut cons_out1) = consumer();
4408
4409            let (mut mailboxes, handles) = spawn_actors_with_handles(
4410                &context,
4411                &oracle,
4412                schemes,
4413                connections,
4414                vec![cons1, dummy_consumer()],
4415                vec![Producer::default(), prod2],
4416            );
4417
4418            // Allow tasks to start
4419            context.sleep(Duration::from_millis(100)).await;
4420
4421            // Count running tasks under the actor prefix
4422            let running_before = count_running_tasks(&context, "actor");
4423            assert!(
4424                running_before > 0,
4425                "at least one actor task should be running"
4426            );
4427
4428            // Verify network is functional
4429            mailboxes[0].fetch(key.clone());
4430            let (_, value) = cons_out1.recv().await.unwrap();
4431            assert_eq!(value, Bytes::from("data for key 1"));
4432
4433            // Abort all actors
4434            for handle in handles {
4435                handle.abort();
4436            }
4437            context.sleep(Duration::from_millis(100)).await;
4438
4439            // Verify all actor tasks are stopped
4440            let running_after = count_running_tasks(&context, "actor");
4441            assert_eq!(
4442                running_after, 0,
4443                "all actor tasks should be stopped, but {running_after} still running"
4444            );
4445        });
4446    }
4447
4448    #[test]
4449    fn test_clean_shutdown() {
4450        for seed in 0..25 {
4451            clean_shutdown(seed);
4452        }
4453    }
4454}