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