1use crate::{
17 Consumer, Delivery, Fetch, Outcome, TargetedResolver,
18 delivery::{Completion as DeliveryCompletion, Tracker as DeliveryTracker},
19 ingress::{self, FetchKey, Message},
20 subscribers,
21};
22use commonware_actor::{Feedback, mailbox};
23use commonware_cryptography::PublicKey;
24use commonware_macros::select_loop;
25use commonware_runtime::{
26 Clock, ContextCell, Handle, Metrics, Spawner, spawn_cell,
27 telemetry::metrics::{MetricsExt as _, status, status::Status},
28};
29use commonware_utils::{
30 Span,
31 futures::{AbortablePool, Aborter},
32 vec::NonEmptyVec,
33};
34use futures::future::{self, Either};
35use std::{
36 collections::{BTreeMap, BTreeSet},
37 future::Future,
38 marker::PhantomData,
39 num::NonZeroUsize,
40 time::{Duration, SystemTime},
41};
42use tracing::{debug, trace, warn};
43
44pub trait Fetcher {
46 type Key: Span;
48
49 type Value;
51
52 fn fetch(&self, key: Self::Key) -> impl Future<Output = Option<Self::Value>> + Send;
58}
59
60pub struct Resolver<K, S, P>
62where
63 K: Span,
64 S: Clone + Eq + Send + 'static,
65 P: PublicKey,
66{
67 mailbox: mailbox::Sender<Message<K, S>>,
68 _peer: PhantomData<P>,
69}
70
71impl<K, S, P> Clone for Resolver<K, S, P>
72where
73 K: Span,
74 S: Clone + Eq + Send + 'static,
75 P: PublicKey,
76{
77 fn clone(&self) -> Self {
78 Self {
79 mailbox: self.mailbox.clone(),
80 _peer: PhantomData,
81 }
82 }
83}
84
85impl<K, S, P> crate::Resolver for Resolver<K, S, P>
86where
87 K: Span,
88 S: Clone + Eq + Send + 'static,
89 P: PublicKey,
90{
91 type Key = K;
92 type Subscriber = S;
93
94 fn fetch<F>(&mut self, fetch: F) -> Feedback
95 where
96 F: Into<Fetch<Self::Key, Self::Subscriber>> + Send,
97 {
98 self.send(Message::Fetch(vec![FetchKey::from(fetch.into())]))
99 }
100
101 fn fetch_all<F>(&mut self, fetches: Vec<F>) -> Feedback
102 where
103 F: Into<Fetch<Self::Key, Self::Subscriber>> + Send,
104 {
105 self.send(Message::Fetch(
106 fetches
107 .into_iter()
108 .map(|fetch| FetchKey::from(fetch.into()))
109 .collect(),
110 ))
111 }
112
113 fn retain(
114 &mut self,
115 predicate: impl Fn(&Self::Key, &Self::Subscriber) -> bool + Send + 'static,
116 ) -> Feedback {
117 self.send(Message::Retain {
118 predicate: Box::new(predicate),
119 })
120 }
121}
122
123impl<K, S, P> TargetedResolver for Resolver<K, S, P>
124where
125 K: Span,
126 S: Clone + Eq + Send + 'static,
127 P: PublicKey,
128{
129 type PublicKey = P;
130
131 fn fetch_targeted(
132 &mut self,
133 fetch: impl Into<Fetch<Self::Key, Self::Subscriber>> + Send,
134 _targets: NonEmptyVec<Self::PublicKey>,
135 ) -> Feedback {
136 <Self as crate::Resolver>::fetch(self, fetch)
137 }
138
139 fn fetch_all_targeted<F>(&mut self, fetches: Vec<(F, NonEmptyVec<Self::PublicKey>)>) -> Feedback
140 where
141 F: Into<Fetch<Self::Key, Self::Subscriber>> + Send,
142 {
143 <Self as crate::Resolver>::fetch_all(
144 self,
145 fetches.into_iter().map(|(fetch, _)| fetch).collect(),
146 )
147 }
148}
149
150impl<K, S, P> Resolver<K, S, P>
151where
152 K: Span,
153 S: Clone + Eq + Send + 'static,
154 P: PublicKey,
155{
156 const fn new(mailbox: mailbox::Sender<Message<K, S>>) -> Self {
157 Self {
158 mailbox,
159 _peer: PhantomData,
160 }
161 }
162
163 fn send(&self, message: Message<K, S>) -> Feedback {
164 self.mailbox.enqueue(message)
165 }
166}
167
168pub fn init<E, F, Con, P>(
170 context: E,
171 fetcher: F,
172 consumer: Con,
173 mailbox_size: NonZeroUsize,
174 fetch_retry_timeout: Duration,
175) -> Resolver<F::Key, Con::Subscriber, P>
176where
177 E: Clock + Spawner + Metrics,
178 F: Fetcher + Clone + Send + 'static,
179 F::Value: Clone + Send + 'static,
180 Con: Consumer<Key = F::Key, Value = F::Value>,
181 Con::Subscriber: Ord,
182 P: PublicKey,
183{
184 let (mailbox_tx, mailbox_rx) = mailbox::new(context.child("mailbox"), mailbox_size);
185 Actor::new(
186 context.child("actor"),
187 fetcher,
188 mailbox_rx,
189 consumer,
190 fetch_retry_timeout,
191 )
192 .start();
193 Resolver::new(mailbox_tx)
194}
195
196struct Actor<E, F, Con>
198where
199 E: Clock + Spawner + Metrics,
200 F: Fetcher,
201 F::Value: Clone + Send + 'static,
202 Con: Consumer<Key = F::Key, Value = F::Value>,
203 Con::Subscriber: Ord,
204{
205 context: ContextCell<E>,
206 fetcher: F,
207 mailbox: mailbox::Receiver<Message<F::Key, Con::Subscriber>>,
208 fetches: AbortablePool<'static, FetchCompletion<F::Key, F::Value>>,
209 deliveries: DeliveryTracker<Con, u64>,
210 requests: BTreeMap<F::Key, Attempt>,
211 subscribers: subscribers::Tracker<F::Key, Con::Subscriber>,
212 fetch: status::Counter,
213 retry_schedule: BTreeSet<(SystemTime, F::Key)>,
214 fetch_retry_timeout: Duration,
215 next_id: u64,
216}
217
218enum Attempt {
219 Fetching { id: u64, _aborter: Aborter },
221
222 Delivering { id: u64 },
224
225 Scheduled(SystemTime),
227}
228
229struct FetchCompletion<K, V> {
230 key: K,
231 id: u64,
232 value: Option<V>,
233}
234
235impl<E, F, Con> Actor<E, F, Con>
236where
237 E: Clock + Spawner + Metrics,
238 F: Fetcher + Clone + Send + 'static,
239 F::Value: Clone + Send + 'static,
240 Con: Consumer<Key = F::Key, Value = F::Value>,
241 Con::Subscriber: Ord,
242{
243 fn new(
244 context: E,
245 fetcher: F,
246 mailbox: mailbox::Receiver<Message<F::Key, Con::Subscriber>>,
247 consumer: Con,
248 fetch_retry_timeout: Duration,
249 ) -> Self {
250 let fetch = context.family("fetch", "Number of fetches by status");
251 Self {
252 context: ContextCell::new(context),
253 fetcher,
254 mailbox,
255 fetches: AbortablePool::default(),
256 deliveries: DeliveryTracker::new(consumer),
257 requests: BTreeMap::new(),
258 subscribers: subscribers::Tracker::new(),
259 fetch,
260 retry_schedule: BTreeSet::new(),
261 fetch_retry_timeout,
262 next_id: 0,
263 }
264 }
265
266 fn start(mut self) -> Handle<()> {
267 spawn_cell!(self.context, self.run())
268 }
269
270 async fn run(mut self) {
271 select_loop! {
272 self.context,
273 on_stopped => {},
274 Ok(result) = self.fetches.next_completed() else continue => {
275 self.handle_fetch_completed(result);
276 },
277 delivery = self.deliveries.next_completion() => {
278 let delivery = match delivery {
279 Ok(delivery) => delivery,
280 Err(_) => continue,
281 };
282 self.handle_delivery_completed(delivery);
283 },
284 _ = match self.retry_schedule.first() {
285 Some((deadline, _)) => Either::Left(self.context.sleep_until(*deadline)),
286 None => Either::Right(future::pending()),
287 } => {
288 self.process_retries();
289 },
290 Some(message) = self.mailbox.recv() else break => {
291 self.handle_message(message);
292 },
293 }
294 }
295
296 fn handle_message(&mut self, message: Message<F::Key, Con::Subscriber>) {
298 match message {
299 Message::Fetch(fetches) => {
300 for fetch in fetches {
301 self.add_fetch(fetch);
302 }
303 }
304 Message::Retain { predicate } => self.retain(predicate),
305 }
306 }
307
308 fn add_fetch(&mut self, fetch: FetchKey<F::Key, Con::Subscriber>) {
310 let FetchKey {
311 key, subscribers, ..
312 } = fetch;
313 let is_new = self.subscribers.insert(key.clone(), subscribers);
314
315 if is_new {
316 assert!(self.deliveries.insert(key.clone()), "delivery entry");
317 self.requests
318 .insert(key.clone(), Attempt::Scheduled(self.context.current()));
319 self.start_fetch(key);
320 }
321 }
322
323 fn retain(&mut self, predicate: ingress::Predicate<F::Key, Con::Subscriber>) {
325 for key in self
326 .subscribers
327 .retain(|key, subscriber| predicate(key, subscriber))
328 {
329 self.deliveries.remove(&key);
330 if let Some(attempt) = self.requests.remove(&key) {
331 match attempt {
332 Attempt::Fetching { .. } | Attempt::Delivering { .. } => {}
333 Attempt::Scheduled(deadline) => {
334 self.retry_schedule.remove(&(deadline, key));
335 }
336 }
337 }
338 }
339 }
340
341 fn start_fetch(&mut self, key: F::Key) {
343 let id = self.next_id;
344 self.next_id = self.next_id.wrapping_add(1);
345 let future = Self::fetch(key.clone(), id, self.fetcher.clone());
346 let aborter = self.fetches.push(future);
347 self.requests.insert(
348 key,
349 Attempt::Fetching {
350 id,
351 _aborter: aborter,
352 },
353 );
354 }
355
356 fn start_delivery(
358 &mut self,
359 key: F::Key,
360 value: F::Value,
361 delivered: NonEmptyVec<(Con::Subscriber, tracing::Span)>,
362 ) {
363 let id = self.next_id;
364 self.next_id = self.next_id.wrapping_add(1);
365 self.deliveries.deliver(
366 Delivery {
367 key: key.clone(),
368 subscribers: delivered,
369 },
370 id,
371 value,
372 );
373 self.requests.insert(key, Attempt::Delivering { id });
374 }
375
376 fn redeliver(&mut self, key: F::Key, delivered: NonEmptyVec<(Con::Subscriber, tracing::Span)>) {
378 self.deliveries.redeliver(Delivery {
379 key,
380 subscribers: delivered,
381 });
382 }
383
384 fn handle_fetch_completed(&mut self, completion: FetchCompletion<F::Key, F::Value>) {
386 let FetchCompletion { key, id, value } = completion;
387 if !self.current_fetch(&key, id) {
388 return;
389 }
390 self.handle_fetched(key, value);
391 }
392
393 fn handle_delivery_completed(
395 &mut self,
396 completion: DeliveryCompletion<F::Key, Con::Subscriber, u64>,
397 ) {
398 let DeliveryCompletion {
399 context: id,
400 delivery,
401 outcome,
402 } = completion;
403 let Delivery {
404 key,
405 subscribers: delivered,
406 ..
407 } = delivery;
408 if !self.current_delivery(&key, id) {
409 return;
410 }
411 self.handle_delivered(key, delivered, outcome);
412 }
413
414 fn current_fetch(&self, key: &F::Key, id: u64) -> bool {
416 let Some(attempt) = self.requests.get(key) else {
417 trace!(?key, id, "ignoring stale fetch completion");
418 return false;
419 };
420 match attempt {
421 Attempt::Fetching { id: active_id, .. } if *active_id == id => true,
422 Attempt::Fetching { id: active_id, .. } => {
423 trace!(
424 ?key,
425 completed_id = id,
426 active_id,
427 "ignoring replaced fetch completion",
428 );
429 false
430 }
431 Attempt::Delivering { id: active_id } => {
432 trace!(
433 ?key,
434 completed_id = id,
435 active_id,
436 "ignoring fetch completion for delivery attempt",
437 );
438 false
439 }
440 Attempt::Scheduled(deadline) => {
441 trace!(?key, id, ?deadline, "ignoring scheduled fetch completion");
442 false
443 }
444 }
445 }
446
447 fn current_delivery(&self, key: &F::Key, id: u64) -> bool {
449 let Some(attempt) = self.requests.get(key) else {
450 trace!(?key, id, "ignoring stale delivery completion");
451 return false;
452 };
453 match attempt {
454 Attempt::Delivering { id: active_id } if *active_id == id => true,
455 Attempt::Delivering { id: active_id } => {
456 trace!(
457 ?key,
458 completed_id = id,
459 active_id,
460 "ignoring replaced delivery completion",
461 );
462 false
463 }
464 Attempt::Fetching { id: active_id, .. } => {
465 trace!(
466 ?key,
467 completed_id = id,
468 active_id,
469 "ignoring delivery completion for fetch attempt",
470 );
471 false
472 }
473 Attempt::Scheduled(deadline) => {
474 trace!(
475 ?key,
476 id,
477 ?deadline,
478 "ignoring scheduled delivery completion"
479 );
480 false
481 }
482 }
483 }
484
485 fn handle_fetched(&mut self, key: F::Key, value: Option<F::Value>) {
487 match value {
488 None => self.schedule_retry(key),
489 Some(value) => {
490 if let Some(subscribers) = self.subscribers.pending(&key) {
491 self.start_delivery(key, value, subscribers);
492 } else {
493 self.requests.remove(&key);
494 self.subscribers.remove(&key);
495 self.deliveries.remove(&key);
496 }
497 }
498 }
499 }
500
501 fn handle_delivered(
503 &mut self,
504 key: F::Key,
505 delivered: NonEmptyVec<(Con::Subscriber, tracing::Span)>,
506 outcome: Option<Outcome>,
507 ) {
508 let accepted = self.deliveries.response_accepted(&key);
509
510 let Some(outcome) = outcome else {
514 let remaining = self
515 .subscribers
516 .remove_delivered(&key, delivered.map_into(|(subscriber, _)| subscriber));
517 if let Some(subscribers) = remaining {
518 self.redeliver(key, subscribers);
519 return;
520 }
521 if !accepted {
522 self.fetch.inc(Status::Dropped);
523 }
524 self.requests.remove(&key);
525 self.deliveries.remove(&key);
526 return;
527 };
528
529 match outcome {
530 Outcome::Complete => {
531 let remaining = self
532 .subscribers
533 .remove_delivered(&key, delivered.map_into(|(subscriber, _)| subscriber));
534
535 if let Some(subscribers) = remaining {
539 if !accepted {
540 self.deliveries.accept_response(&key);
541 }
542 self.redeliver(key, subscribers);
543 } else {
544 self.requests.remove(&key);
545 self.subscribers.remove(&key);
546 self.deliveries.remove(&key);
547 }
548 }
549 Outcome::Ambiguous => {
550 self.fetch.inc(Status::Ambiguous);
554 self.deliveries.discard_response(&key);
555 self.schedule_retry(key);
556 }
557 Outcome::Invalid => {
558 if accepted {
562 warn!(
563 ?key,
564 "previously accepted resolver response rejected during opaque redelivery"
565 );
566 self.requests.remove(&key);
567 self.subscribers.remove(&key);
568 self.deliveries.remove(&key);
569 return;
570 }
571
572 warn!(?key, "consumer rejected opaque resolver delivery");
573 self.deliveries.discard_response(&key);
574 self.schedule_retry(key);
575 }
576 Outcome::Ignored => {
577 self.fetch.inc(Status::Dropped);
580 self.requests.remove(&key);
581 self.subscribers.remove(&key);
582 self.deliveries.remove(&key);
583 }
584 }
585 }
586
587 fn schedule_retry(&mut self, key: F::Key) {
589 let deadline = self.context.current() + self.fetch_retry_timeout;
590 let Some(attempt) = self.requests.get_mut(&key) else {
591 return;
592 };
593 *attempt = Attempt::Scheduled(deadline);
594 debug!(?key, ?deadline, "scheduled opaque resolver retry");
595 self.retry_schedule.insert((deadline, key));
596 }
597
598 fn process_retries(&mut self) {
600 let now = self.context.current();
601 while let Some((deadline, key)) = self.retry_schedule.pop_first() {
602 if deadline > now {
603 self.retry_schedule.insert((deadline, key));
604 break;
605 }
606
607 let Some(state) = self.requests.get(&key) else {
608 continue;
609 };
610 match state {
611 Attempt::Scheduled(state_deadline) if *state_deadline == deadline => {
612 debug!(?key, "retrying opaque resolver fetch");
613 self.start_fetch(key);
614 }
615 Attempt::Scheduled(_) | Attempt::Fetching { .. } | Attempt::Delivering { .. } => {}
616 }
617 }
618 }
619
620 async fn fetch(key: F::Key, id: u64, fetcher: F) -> FetchCompletion<F::Key, F::Value> {
622 let value = fetcher.fetch(key.clone()).await;
623 FetchCompletion { key, id, value }
624 }
625}
626
627#[cfg(test)]
628mod tests {
629 use super::*;
630 use crate::Resolver as _;
631 use bytes::Bytes;
632 use commonware_cryptography::{
633 Signer,
634 ed25519::{PrivateKey, PublicKey},
635 };
636 use commonware_runtime::{Runner as _, Supervisor as _, deterministic, deterministic::Runner};
637 use commonware_utils::{channel::oneshot, non_empty_vec, sync::Mutex};
638 use std::{
639 collections::{HashMap, VecDeque},
640 sync::{
641 Arc,
642 atomic::{AtomicU32, Ordering},
643 },
644 };
645
646 const RETRY_TIMEOUT: Duration = Duration::from_millis(100);
647
648 #[derive(Clone, Default)]
649 struct MockFetcher {
650 responses: Arc<Mutex<HashMap<u8, VecDeque<Option<Bytes>>>>>,
651 calls: Arc<AtomicU32>,
652 }
653
654 impl MockFetcher {
655 fn push(&self, key: u8, response: Option<Bytes>) {
656 self.responses
657 .lock()
658 .entry(key)
659 .or_default()
660 .push_back(response);
661 }
662
663 fn calls(&self) -> u32 {
664 self.calls.load(Ordering::Relaxed)
665 }
666 }
667
668 impl Fetcher for MockFetcher {
669 type Key = u8;
670 type Value = Bytes;
671
672 fn fetch(&self, key: Self::Key) -> impl Future<Output = Option<Self::Value>> + Send {
673 let responses = self.responses.clone();
674 let calls = self.calls.clone();
675 async move {
676 calls.fetch_add(1, Ordering::Relaxed);
677 responses
678 .lock()
679 .get_mut(&key)
680 .and_then(VecDeque::pop_front)
681 .flatten()
682 }
683 }
684 }
685
686 #[derive(Clone)]
687 struct BlockingFetcher {
688 started: Arc<Mutex<Option<oneshot::Sender<()>>>>,
689 response: Arc<Mutex<Option<oneshot::Receiver<Option<Bytes>>>>>,
690 }
691
692 impl BlockingFetcher {
693 fn new() -> (Self, oneshot::Receiver<()>, oneshot::Sender<Option<Bytes>>) {
694 let (started_tx, started_rx) = oneshot::channel();
695 let (response_tx, response_rx) = oneshot::channel();
696 (
697 Self {
698 started: Arc::new(Mutex::new(Some(started_tx))),
699 response: Arc::new(Mutex::new(Some(response_rx))),
700 },
701 started_rx,
702 response_tx,
703 )
704 }
705 }
706
707 impl Fetcher for BlockingFetcher {
708 type Key = u8;
709 type Value = Bytes;
710
711 fn fetch(&self, _key: Self::Key) -> impl Future<Output = Option<Self::Value>> + Send {
712 let started = self.started.clone();
713 let response = self.response.clone();
714 async move {
715 if let Some(started) = started.lock().take() {
716 let _ = started.send(());
717 }
718 let response = response.lock().take().expect("missing response");
719 response.await.unwrap_or(None)
720 }
721 }
722 }
723
724 struct CapturedDelivery {
725 delivery: Delivery<u8, u16>,
726 value: Bytes,
727 response: oneshot::Sender<Outcome>,
728 }
729
730 #[derive(Clone, Default)]
731 struct MockConsumer {
732 deliveries: Arc<Mutex<VecDeque<CapturedDelivery>>>,
733 }
734
735 impl MockConsumer {
736 fn pop(&self) -> Option<CapturedDelivery> {
737 self.deliveries.lock().pop_front()
738 }
739
740 fn len(&self) -> usize {
741 self.deliveries.lock().len()
742 }
743 }
744
745 impl Consumer for MockConsumer {
746 type Key = u8;
747 type Value = Bytes;
748 type Subscriber = u16;
749 type Outcome = Outcome;
750
751 fn deliver(
752 &mut self,
753 delivery: Delivery<Self::Key, Self::Subscriber>,
754 value: Self::Value,
755 ) -> oneshot::Receiver<Self::Outcome> {
756 let (response, receiver) = oneshot::channel();
757 self.deliveries.lock().push_back(CapturedDelivery {
758 delivery,
759 value,
760 response,
761 });
762 receiver
763 }
764 }
765
766 fn start_resolver<F>(
767 context: deterministic::Context,
768 fetcher: F,
769 consumer: MockConsumer,
770 ) -> Resolver<u8, u16, PublicKey>
771 where
772 F: Fetcher<Key = u8, Value = Bytes> + Clone + Send + 'static,
773 {
774 init(
775 context,
776 fetcher,
777 consumer,
778 NonZeroUsize::new(16).unwrap(),
779 RETRY_TIMEOUT,
780 )
781 }
782
783 async fn wait_for_delivery(
784 context: &deterministic::Context,
785 consumer: &MockConsumer,
786 ) -> CapturedDelivery {
787 for _ in 0..50 {
788 if let Some(delivery) = consumer.pop() {
789 return delivery;
790 }
791 context.sleep(Duration::from_millis(10)).await;
792 }
793 panic!("timed out waiting for delivery");
794 }
795
796 #[test]
797 fn fetch_during_validation_reuses_response_after_success() {
798 Runner::default().start(|context| async move {
799 let fetcher = MockFetcher::default();
800 fetcher.push(1, Some(Bytes::from_static(b"value")));
801 let consumer = MockConsumer::default();
802 let mut resolver =
803 start_resolver(context.child("resolver"), fetcher.clone(), consumer.clone());
804
805 assert!(
806 resolver
807 .fetch(Fetch {
808 key: 1,
809 subscriber: 10,
810 span: tracing::Span::none(),
811 })
812 .accepted()
813 );
814 let first = wait_for_delivery(&context, &consumer).await;
815 assert_eq!(first.value, Bytes::from_static(b"value"));
816
817 assert!(
818 resolver
819 .fetch(Fetch {
820 key: 1,
821 subscriber: 11,
822 span: tracing::Span::none(),
823 })
824 .accepted()
825 );
826 context.sleep(Duration::from_millis(10)).await;
827 first
828 .response
829 .send(Outcome::Complete)
830 .expect("response dropped");
831
832 let second = wait_for_delivery(&context, &consumer).await;
833 assert_eq!(second.value, Bytes::from_static(b"value"));
834 assert_eq!(
835 second
836 .delivery
837 .subscribers
838 .iter()
839 .map(|(subscriber, _)| *subscriber)
840 .collect::<Vec<_>>(),
841 vec![11]
842 );
843 second
844 .response
845 .send(Outcome::Complete)
846 .expect("response dropped");
847
848 context.sleep(Duration::from_millis(10)).await;
849 assert_eq!(fetcher.calls(), 1);
850 });
851 }
852
853 #[test]
854 fn missing_fetch_retries_until_value_is_available() {
855 Runner::default().start(|context| async move {
856 let fetcher = MockFetcher::default();
857 fetcher.push(1, None);
858 fetcher.push(1, Some(Bytes::from_static(b"value")));
859 let consumer = MockConsumer::default();
860 let mut resolver =
861 start_resolver(context.child("resolver"), fetcher.clone(), consumer.clone());
862
863 assert!(
864 resolver
865 .fetch(Fetch {
866 key: 1,
867 subscriber: 10,
868 span: tracing::Span::none(),
869 })
870 .accepted()
871 );
872 context
873 .sleep(RETRY_TIMEOUT + Duration::from_millis(10))
874 .await;
875
876 let delivery = wait_for_delivery(&context, &consumer).await;
877 assert_eq!(delivery.value, Bytes::from_static(b"value"));
878 delivery
879 .response
880 .send(Outcome::Complete)
881 .expect("response dropped");
882 assert_eq!(fetcher.calls(), 2);
883 });
884 }
885
886 #[test]
887 fn ambiguous_delivery_retries_without_retiring_subscriber() {
888 Runner::default().start(|context| async move {
889 let fetcher = MockFetcher::default();
890 fetcher.push(1, Some(Bytes::from_static(b"ambiguous")));
891 fetcher.push(1, Some(Bytes::from_static(b"complete")));
892 let consumer = MockConsumer::default();
893 let mut resolver =
894 start_resolver(context.child("resolver"), fetcher.clone(), consumer.clone());
895
896 assert!(
897 resolver
898 .fetch(Fetch {
899 key: 1,
900 subscriber: 10,
901 span: tracing::Span::none(),
902 })
903 .accepted()
904 );
905
906 let first = wait_for_delivery(&context, &consumer).await;
907 assert_eq!(first.value, Bytes::from_static(b"ambiguous"));
908 assert_eq!(
909 first
910 .delivery
911 .subscribers
912 .iter()
913 .map(|(subscriber, _)| *subscriber)
914 .collect::<Vec<_>>(),
915 vec![10]
916 );
917 first
918 .response
919 .send(Outcome::Ambiguous)
920 .expect("response dropped");
921
922 context
923 .sleep(RETRY_TIMEOUT + Duration::from_millis(10))
924 .await;
925 let second = wait_for_delivery(&context, &consumer).await;
926 assert_eq!(second.value, Bytes::from_static(b"complete"));
927 assert_eq!(
928 second
929 .delivery
930 .subscribers
931 .iter()
932 .map(|(subscriber, _)| *subscriber)
933 .collect::<Vec<_>>(),
934 vec![10]
935 );
936 assert_eq!(fetcher.calls(), 2);
937
938 second
939 .response
940 .send(Outcome::Complete)
941 .expect("response dropped");
942 context.sleep(Duration::from_millis(10)).await;
943 assert_eq!(consumer.len(), 0);
944 assert_eq!(fetcher.calls(), 2);
945 assert!(
946 context
947 .encode()
948 .contains("resolver_actor_fetch_total{status=\"Ambiguous\"} 1")
949 );
950 });
951 }
952
953 #[test]
954 fn ignored_delivery_retires_without_retry() {
955 Runner::default().start(|context| async move {
956 let fetcher = MockFetcher::default();
957 fetcher.push(1, Some(Bytes::from_static(b"obsolete")));
958 fetcher.push(1, Some(Bytes::from_static(b"fresh")));
959 let consumer = MockConsumer::default();
960 let mut resolver =
961 start_resolver(context.child("resolver"), fetcher.clone(), consumer.clone());
962
963 assert!(
964 resolver
965 .fetch(Fetch {
966 key: 1,
967 subscriber: 10,
968 span: tracing::Span::none(),
969 })
970 .accepted()
971 );
972 let delivery = wait_for_delivery(&context, &consumer).await;
973 assert_eq!(delivery.value, Bytes::from_static(b"obsolete"));
974 delivery
975 .response
976 .send(Outcome::Ignored)
977 .expect("response dropped");
978
979 context
980 .sleep(RETRY_TIMEOUT + Duration::from_millis(10))
981 .await;
982 assert_eq!(fetcher.calls(), 1);
983 assert_eq!(consumer.len(), 0);
984 assert!(
985 context
986 .encode()
987 .contains("resolver_actor_fetch_total{status=\"Dropped\"} 1")
988 );
989
990 assert!(
992 resolver
993 .fetch(Fetch {
994 key: 1,
995 subscriber: 10,
996 span: tracing::Span::none(),
997 })
998 .accepted()
999 );
1000 let delivery = wait_for_delivery(&context, &consumer).await;
1001 assert_eq!(delivery.value, Bytes::from_static(b"fresh"));
1002 delivery
1003 .response
1004 .send(Outcome::Complete)
1005 .expect("response dropped");
1006 assert_eq!(fetcher.calls(), 2);
1007 });
1008 }
1009
1010 #[test]
1011 fn dropped_verdict_retires_without_retry() {
1012 Runner::default().start(|context| async move {
1013 let fetcher = MockFetcher::default();
1014 fetcher.push(1, Some(Bytes::from_static(b"unjudged")));
1015 fetcher.push(1, Some(Bytes::from_static(b"fresh")));
1016 let consumer = MockConsumer::default();
1017 let mut resolver =
1018 start_resolver(context.child("resolver"), fetcher.clone(), consumer.clone());
1019
1020 assert!(
1021 resolver
1022 .fetch(Fetch {
1023 key: 1,
1024 subscriber: 10,
1025 span: tracing::Span::none(),
1026 })
1027 .accepted()
1028 );
1029 let delivery = wait_for_delivery(&context, &consumer).await;
1030 assert_eq!(delivery.value, Bytes::from_static(b"unjudged"));
1031
1032 drop(delivery);
1036
1037 context
1038 .sleep(RETRY_TIMEOUT + Duration::from_millis(10))
1039 .await;
1040 assert_eq!(fetcher.calls(), 1);
1041 assert_eq!(consumer.len(), 0);
1042 assert!(
1043 context
1044 .encode()
1045 .contains("resolver_actor_fetch_total{status=\"Dropped\"} 1")
1046 );
1047
1048 assert!(
1050 resolver
1051 .fetch(Fetch {
1052 key: 1,
1053 subscriber: 10,
1054 span: tracing::Span::none(),
1055 })
1056 .accepted()
1057 );
1058 let delivery = wait_for_delivery(&context, &consumer).await;
1059 assert_eq!(delivery.value, Bytes::from_static(b"fresh"));
1060 delivery
1061 .response
1062 .send(Outcome::Complete)
1063 .expect("response dropped");
1064 assert_eq!(fetcher.calls(), 2);
1065 });
1066 }
1067
1068 #[test]
1069 fn dropped_verdict_hands_response_to_late_subscriber() {
1070 Runner::default().start(|context| async move {
1071 let fetcher = MockFetcher::default();
1072 fetcher.push(1, Some(Bytes::from_static(b"unjudged")));
1073 let consumer = MockConsumer::default();
1074 let mut resolver =
1075 start_resolver(context.child("resolver"), fetcher.clone(), consumer.clone());
1076
1077 assert!(
1078 resolver
1079 .fetch(Fetch {
1080 key: 1,
1081 subscriber: 10,
1082 span: tracing::Span::none(),
1083 })
1084 .accepted()
1085 );
1086 let first = wait_for_delivery(&context, &consumer).await;
1087 assert_eq!(first.value, Bytes::from_static(b"unjudged"));
1088
1089 assert!(
1093 resolver
1094 .fetch(Fetch {
1095 key: 1,
1096 subscriber: 11,
1097 span: tracing::Span::none(),
1098 })
1099 .accepted()
1100 );
1101 context.sleep(Duration::from_millis(10)).await;
1102 drop(first);
1103 let second = wait_for_delivery(&context, &consumer).await;
1104 assert_eq!(
1105 second
1106 .delivery
1107 .subscribers
1108 .iter()
1109 .map(|(subscriber, _)| *subscriber)
1110 .collect::<Vec<_>>(),
1111 vec![11]
1112 );
1113 assert_eq!(second.value, Bytes::from_static(b"unjudged"));
1114 second
1115 .response
1116 .send(Outcome::Complete)
1117 .expect("response dropped");
1118
1119 context
1120 .sleep(RETRY_TIMEOUT + Duration::from_millis(10))
1121 .await;
1122 assert_eq!(fetcher.calls(), 1);
1123 assert_eq!(consumer.len(), 0);
1124 });
1125 }
1126
1127 #[test]
1128 fn accepted_redelivery_rejection_does_not_refetch() {
1129 Runner::default().start(|context| async move {
1130 let fetcher = MockFetcher::default();
1131 fetcher.push(1, Some(Bytes::from_static(b"value")));
1132 let consumer = MockConsumer::default();
1133 let mut resolver =
1134 start_resolver(context.child("resolver"), fetcher.clone(), consumer.clone());
1135
1136 assert!(
1137 resolver
1138 .fetch(Fetch {
1139 key: 1,
1140 subscriber: 10,
1141 span: tracing::Span::none(),
1142 })
1143 .accepted()
1144 );
1145 let first = wait_for_delivery(&context, &consumer).await;
1146
1147 assert!(
1148 resolver
1149 .fetch(Fetch {
1150 key: 1,
1151 subscriber: 11,
1152 span: tracing::Span::none(),
1153 })
1154 .accepted()
1155 );
1156 context.sleep(Duration::from_millis(10)).await;
1157 first
1158 .response
1159 .send(Outcome::Complete)
1160 .expect("response dropped");
1161
1162 let second = wait_for_delivery(&context, &consumer).await;
1163 second
1164 .response
1165 .send(Outcome::Invalid)
1166 .expect("response dropped");
1167
1168 context
1169 .sleep(RETRY_TIMEOUT + Duration::from_millis(10))
1170 .await;
1171 assert_eq!(fetcher.calls(), 1);
1172 assert_eq!(consumer.len(), 0);
1173 });
1174 }
1175
1176 #[test]
1177 fn retain_prunes_active_fetch_subscribers() {
1178 Runner::default().start(|context| async move {
1179 let (fetcher, started, response) = BlockingFetcher::new();
1180 let consumer = MockConsumer::default();
1181 let mut resolver = start_resolver(context.child("resolver"), fetcher, consumer.clone());
1182
1183 assert!(
1184 resolver
1185 .fetch(Fetch {
1186 key: 1,
1187 subscriber: 10,
1188 span: tracing::Span::none(),
1189 })
1190 .accepted()
1191 );
1192 assert!(
1193 resolver
1194 .fetch(Fetch {
1195 key: 1,
1196 subscriber: 11,
1197 span: tracing::Span::none(),
1198 })
1199 .accepted()
1200 );
1201 started.await.expect("fetch did not start");
1202 assert!(
1203 resolver
1204 .retain(|_, subscriber| *subscriber == 11)
1205 .accepted()
1206 );
1207 context.sleep(Duration::from_millis(10)).await;
1208 response
1209 .send(Some(Bytes::from_static(b"value")))
1210 .expect("fetcher dropped");
1211
1212 let delivery = wait_for_delivery(&context, &consumer).await;
1213 assert_eq!(
1214 delivery
1215 .delivery
1216 .subscribers
1217 .iter()
1218 .map(|(subscriber, _)| *subscriber)
1219 .collect::<Vec<_>>(),
1220 vec![11]
1221 );
1222 delivery
1223 .response
1224 .send(Outcome::Complete)
1225 .expect("response dropped");
1226 });
1227 }
1228
1229 #[test]
1230 fn retain_drops_last_subscriber_aborts_active_fetch() {
1231 Runner::default().start(|context| async move {
1232 let (fetcher, started, response) = BlockingFetcher::new();
1233 let consumer = MockConsumer::default();
1234 let mut resolver = start_resolver(context.child("resolver"), fetcher, consumer.clone());
1235
1236 assert!(
1237 resolver
1238 .fetch(Fetch {
1239 key: 1,
1240 subscriber: 10,
1241 span: tracing::Span::none(),
1242 })
1243 .accepted()
1244 );
1245 started.await.expect("fetch did not start");
1246 assert!(resolver.retain(|_, _| false).accepted());
1247 context.sleep(Duration::from_millis(10)).await;
1248
1249 assert!(
1250 response.send(Some(Bytes::from_static(b"value"))).is_err(),
1251 "fetch future should be aborted after its last subscriber is pruned"
1252 );
1253 context
1254 .sleep(RETRY_TIMEOUT + Duration::from_millis(10))
1255 .await;
1256 assert_eq!(consumer.len(), 0);
1257 });
1258 }
1259
1260 #[test]
1261 fn retain_drops_last_subscriber_aborts_active_delivery() {
1262 Runner::default().start(|context| async move {
1263 let fetcher = MockFetcher::default();
1264 fetcher.push(1, Some(Bytes::from_static(b"value")));
1265 let consumer = MockConsumer::default();
1266 let mut resolver =
1267 start_resolver(context.child("resolver"), fetcher.clone(), consumer.clone());
1268
1269 assert!(
1270 resolver
1271 .fetch(Fetch {
1272 key: 1,
1273 subscriber: 10,
1274 span: tracing::Span::none(),
1275 })
1276 .accepted()
1277 );
1278 let delivery = wait_for_delivery(&context, &consumer).await;
1279 assert!(resolver.retain(|_, _| false).accepted());
1280 context.sleep(Duration::from_millis(10)).await;
1281
1282 assert!(
1283 delivery.response.send(Outcome::Invalid).is_err(),
1284 "delivery future should be aborted after its last subscriber is pruned"
1285 );
1286 context
1287 .sleep(RETRY_TIMEOUT + Duration::from_millis(10))
1288 .await;
1289 assert_eq!(fetcher.calls(), 1);
1290 assert_eq!(consumer.len(), 0);
1291 });
1292 }
1293
1294 #[test]
1295 fn targeted_fetch_uses_same_opaque_fetch_path() {
1296 Runner::default().start(|context| async move {
1297 let fetcher = MockFetcher::default();
1298 fetcher.push(1, Some(Bytes::from_static(b"value")));
1299 let consumer = MockConsumer::default();
1300 let mut resolver =
1301 start_resolver(context.child("resolver"), fetcher.clone(), consumer.clone());
1302 let target = PrivateKey::from_seed(0).public_key();
1303
1304 assert!(
1305 resolver
1306 .fetch_targeted(
1307 Fetch {
1308 key: 1,
1309 subscriber: 10,
1310 span: tracing::Span::none(),
1311 },
1312 non_empty_vec![target]
1313 )
1314 .accepted()
1315 );
1316 let delivery = wait_for_delivery(&context, &consumer).await;
1317 assert_eq!(delivery.value, Bytes::from_static(b"value"));
1318 delivery
1319 .response
1320 .send(Outcome::Complete)
1321 .expect("response dropped");
1322 assert_eq!(fetcher.calls(), 1);
1323 });
1324 }
1325}