1use 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
105pub trait Producer: Clone + Send + 'static {
107 type Key: Span;
109
110 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 #[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 #[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 #[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 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(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 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 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 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 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 #[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 #[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 #[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 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 #[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 #[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_link(&mut oracle, LINK_UNRELIABLE.clone(), &peers, 0, 1).await;
1701 add_link(&mut oracle, LINK_UNRELIABLE.clone(), &peers, 0, 2).await;
1702
1703 for _ in 0..10 {
1705 mailbox1.fetch(key2.clone());
1707 mailbox1.fetch(key3.clone());
1708
1709 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 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 #[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 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 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 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 let key = Key(7);
1800 mailbox1.fetch(key.clone());
1801 let canceled = key.clone();
1802 mailbox1.retain(move |key, _| key != &canceled);
1803
1804 select! {
1806 _ = cons_out1.recv() => panic!("unexpected event"),
1807 _ = context.sleep(Duration::from_millis(100)) => {},
1808 };
1809 });
1810 }
1811
1812 #[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 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 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 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 for _ in 0..20 {
1880 mailbox1.fetch(key_a.clone());
1882
1883 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 mailbox1.fetch(key_b.clone());
1891
1892 context.sleep(Duration::from_secs(5)).await;
1894
1895 select! {
1897 _ = cons_out1.recv() => panic!("unexpected event"),
1898 _ = context.sleep(Duration::from_millis(100)) => {},
1899 };
1900
1901 let canceled = key_b.clone();
1903 mailbox1.retain(move |key, _| key != &canceled);
1904
1905 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 #[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 mailbox1.fetch(key.clone());
1955 mailbox1.fetch(key.clone());
1956
1957 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 select! {
1964 _ = cons_out1.recv() => {
1965 panic!("Unexpected second event received for duplicate fetch");
1966 },
1967 _ = context.sleep(Duration::from_millis(500)) => {
1968 },
1970 };
1971 });
1972 }
1973
1974 #[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 mailbox1.fetch(key1.clone());
2022
2023 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 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 context.sleep(Duration::from_millis(200)).await;
2042
2043 mailbox1.fetch(key2.clone());
2045
2046 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 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 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 context.sleep(Duration::from_millis(100)).await;
2113
2114 mailbox1.fetch_targeted(
2118 key.clone(),
2119 non_empty_vec![peers[1].clone(), peers[2].clone()],
2120 );
2121
2122 let (key_actual, value) = cons_out1.recv().await.unwrap();
2124 assert_eq!(key_actual, key);
2125 assert_eq!(value, valid_data);
2126
2127 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 let metrics = context.encode();
2135 assert_eq!(
2136 status_metric_total(&metrics, "actor_fetch_total", "Success"),
2137 1
2138 );
2139 });
2140 }
2141
2142 #[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 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 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 select! {
2258 _ = started.recv() => panic!("blocked target was retried"),
2259 _ = context.sleep(Duration::from_secs(1)) => {},
2260 };
2261
2262 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 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(), );
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(), );
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 context.sleep(Duration::from_millis(100)).await;
2338
2339 mailbox1.fetch_targeted(
2342 key.clone(),
2343 non_empty_vec![peers[1].clone(), peers[2].clone()],
2344 );
2345
2346 select! {
2349 event = cons_out1.recv() => {
2350 panic!("Fetch should not succeed, but got: {event:?}");
2351 },
2352 _ = context.sleep(Duration::from_secs(3)) => {
2353 },
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 let mut prod2 = Producer::default();
2376 prod2.insert(key1.clone(), Bytes::from("data for key 1"));
2377
2378 let mut prod3 = Producer::default();
2380 prod3.insert(key3.clone(), Bytes::from("data for key 3"));
2381
2382 let mut prod4 = Producer::default();
2384 prod4.insert(key2.clone(), Bytes::from("data for key 2"));
2385
2386 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 context.sleep(Duration::from_millis(100)).await;
2438
2439 mailbox1.fetch_all_targeted(vec![
2444 (key1.clone(), non_empty_vec![peers[1].clone()]), (key2.clone(), non_empty_vec![peers[3].clone()]), ]);
2447 mailbox1.fetch(key3.clone()); 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 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 let metrics = context.encode();
2464 assert_eq!(
2465 status_metric_total(&metrics, "actor_fetch_total", "Success"),
2466 3
2467 );
2468 });
2469 }
2470
2471 #[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 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(), );
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 context.sleep(Duration::from_millis(100)).await;
2527
2528 mailbox1.fetch_targeted(key.clone(), non_empty_vec![peers[1].clone()]);
2530
2531 context.sleep(Duration::from_millis(500)).await;
2533
2534 mailbox1.fetch(key.clone());
2536
2537 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 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(), );
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 context.sleep(Duration::from_millis(100)).await;
2598
2599 mailbox1.fetch(key.clone());
2601
2602 context.sleep(Duration::from_millis(50)).await;
2604
2605 mailbox1.fetch_targeted(key.clone(), non_empty_vec![peers[1].clone()]);
2608
2609 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 mailbox1.retain(|_, _| true);
2654 select! {
2655 _ = cons_out1.recv() => {
2656 panic!("unexpected event");
2657 },
2658 _ = context.sleep(Duration::from_millis(100)) => {},
2659 };
2660
2661 mailbox1.fetch(key.clone());
2663
2664 let key_clone = key.clone();
2667 mailbox1.retain(move |key, _| key != &key_clone);
2668
2669 add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
2671
2672 mailbox1.fetch(key.clone());
2675
2676 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 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 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 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 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 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 mailbox1.retain(|_, _| false);
3514 select! {
3515 _ = cons_out1.recv() => {
3516 panic!("unexpected event");
3517 },
3518 _ = context.sleep(Duration::from_millis(100)) => {},
3519 };
3520
3521 mailbox1.fetch(key.clone());
3523
3524 mailbox1.retain(|_, _| false);
3526
3527 add_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3529
3530 mailbox1.fetch(key.clone());
3533
3534 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 #[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 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_link(&mut oracle, LINK.clone(), &peers, 0, 1).await;
3559 add_link(&mut oracle, LINK.clone(), &peers, 0, 2).await;
3560
3561 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 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 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 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 context.sleep(Duration::from_millis(100)).await;
3609 let start = context.current();
3610
3611 mailbox1.fetch(Key(0));
3615 mailbox1.fetch(Key(1));
3616
3617 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 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 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 #[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 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 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 context.sleep(Duration::from_millis(100)).await;
3695 let start = context.current();
3696
3697 mailbox1.fetch(Key(1));
3700 mailbox1.fetch(Key(2));
3701 mailbox1.fetch(Key(3));
3702
3703 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 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 #[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 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 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, );
3763
3764 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 context.sleep(Duration::from_millis(100)).await;
3778
3779 mailbox1.fetch(key.clone());
3781
3782 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 for handle in handles {
4370 handle.abort();
4371 }
4372 context.sleep(Duration::from_millis(100)).await;
4373
4374 let key2 = Key(2);
4378 mailboxes[0].fetch(key2.clone());
4379
4380 let canceled = key2;
4382 mailboxes[0].retain(move |key, _| key != &canceled);
4383
4384 mailboxes[0].retain(|_, _| true);
4386
4387 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 context.sleep(Duration::from_millis(100)).await;
4420
4421 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 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 for handle in handles {
4435 handle.abort();
4436 }
4437 context.sleep(Duration::from_millis(100)).await;
4438
4439 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}