Skip to main content

commonware_resolver/
opaque.rs

1//! Resolve keys from an opaque asynchronous fetcher.
2//!
3//! This module owns the generic resolver actor used when fetching data only
4//! requires asking an application-provided source for raw bytes or objects.
5//! Implementations provide [`Fetcher::fetch`]; this module handles request
6//! coalescing, retain pruning, retry scheduling, consumer delivery, and
7//! cached-response redelivery. An ignored consumer outcome retires the key
8//! without accepting the value or retrying the source. A verdict the consumer
9//! drops without answering hands the response to the remaining subscribers, or
10//! retires the key when none remain.
11//!
12//! Target hints supplied through [`crate::TargetedResolver::fetch_targeted`] and
13//! [`crate::TargetedResolver::fetch_all_targeted`] are ignored because opaque
14//! fetchers do not have peer-specific routing.
15
16use 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
44/// Fetches raw values for resolver keys.
45pub trait Fetcher {
46    /// Key requested by the resolver.
47    type Key: Span;
48
49    /// Raw value delivered to the consumer for validation.
50    type Value;
51
52    /// Fetch the value for `key`.
53    ///
54    /// Return `None` for transient failures, missing data, or unexpected source
55    /// responses. The resolver will retry while the key still has retained
56    /// subscribers.
57    fn fetch(&self, key: Self::Key) -> impl Future<Output = Option<Self::Value>> + Send;
58}
59
60/// Handle to an opaque-fetcher resolver actor.
61pub 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
168/// Spawn an opaque-fetcher resolver actor.
169pub 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
196/// Actor that coalesces opaque fetches, retries failures, and delivers accepted values.
197struct 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    /// Fetch future is active for this key.
220    Fetching { id: u64, _aborter: Aborter },
221
222    /// Consumer validation is active for this key.
223    Delivering { id: u64 },
224
225    /// Fetch is sleeping until the retry deadline.
226    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    /// Apply a mailbox message to active resolver state.
297    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    /// Add subscribers for a key and start the first fetch if needed.
309    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    /// Prune subscribers, deliveries, active fetches, and scheduled retries.
324    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    /// Spawn one fetch attempt for `key`.
342    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    /// Deliver a fetched value to currently retained subscribers.
357    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    /// Deliver the cached response to subscribers that arrived later.
377    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    /// Handle a completed fetch future if it is still the active attempt.
385    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    /// Handle a completed consumer delivery if it is still the active attempt.
394    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    /// Return whether a fetch completion matches the current attempt id.
415    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    /// Return whether a delivery completion matches the current attempt id.
448    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    /// Deliver fetched values or schedule a retry when the source had no data.
486    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    /// Complete, redeliver, or retry a key after consumer validation.
502    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        // A dropped verdict says nothing about the response, only that the consumer
511        // did not judge it for these subscribers. Hand the response to the
512        // remaining subscribers, or retire the key when none remain.
513        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                // The first accepted response is reused for subscribers that joined
536                // while validation was pending, avoiding a duplicate source fetch
537                // for the same key.
538                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                // The fetcher returned one of multiple valid responses, but this response did not
551                // satisfy every subscriber. Discard it and retain the fetch so another response
552                // can be tried.
553                self.fetch.inc(Status::Ambiguous);
554                self.deliveries.discard_response(&key);
555                self.schedule_retry(key);
556            }
557            Outcome::Invalid => {
558                // A cached response already satisfied at least one subscriber. Treat a
559                // later rejection during redelivery as stale application feedback rather
560                // than re-fetching data that was accepted once.
561                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                // The consumer no longer needs this key. Retire it without accepting the
578                // response or scheduling another source fetch.
579                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    /// Schedule the next fetch attempt for `key`.
588    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    /// Start all retry attempts whose deadlines have passed.
599    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    /// Run the user-supplied fetcher and preserve the attempt id.
621    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            // Ignoring retires the old fetch rather than leaving a tombstone.
991            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 the verdict, as a consumer does when it stops with the delivery
1033            // still queued. No one is waiting on the key, so the fetch is retired
1034            // rather than retried.
1035            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            // The retired key is not deduplicated against, so a fresh fetch runs.
1049            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            // A late subscriber joins while the first delivery is unjudged, then
1090            // that delivery's verdict is dropped. The late subscriber is handed
1091            // the same response without another source fetch.
1092            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}