Skip to main content

commonware_resolver/p2p/
engine.rs

1use super::{
2    config::Config,
3    fetcher::{Config as FetcherConfig, Fetcher},
4    inflight::Inflight,
5    ingress::{FetchKey, Mailbox, Message},
6    metrics, wire, Producer,
7};
8use crate::{subscribers, Consumer, Delivery};
9use bytes::Bytes;
10use commonware_actor::mailbox;
11use commonware_cryptography::PublicKey;
12use commonware_macros::select_loop;
13use commonware_p2p::{
14    utils::codec::{wrap, WrappedSender},
15    Blocker, Provider, Receiver, Recipients, Sender,
16};
17use commonware_runtime::{
18    spawn_cell,
19    telemetry::metrics::{histogram, status::Status, GaugeExt},
20    BufferPooler, Clock, ContextCell, Handle, Metrics, Spawner,
21};
22use commonware_utils::{channel::oneshot, futures::Pool as FuturesPool, Span};
23use futures::future::{self, Either};
24use rand_core::Rng;
25use std::marker::PhantomData;
26use tracing::{debug, error, trace, warn};
27
28/// Represents a pending serve operation.
29struct Serve<P: PublicKey> {
30    timer: histogram::Timer,
31    peer: P,
32    id: u64,
33    result: Result<Bytes, oneshot::error::RecvError>,
34}
35
36/// Manages incoming and outgoing P2P requests, coordinating fetch and serve operations.
37pub struct Engine<E, P, D, B, Key, Con, Pro, NetS, NetR>
38where
39    E: BufferPooler + Clock + Spawner + Rng + Metrics,
40    P: PublicKey,
41    D: Provider<PublicKey = P>,
42    B: Blocker<PublicKey = P>,
43    Key: Span,
44    Con: Consumer<Key = Key, Value = Bytes>,
45    Pro: Producer<Key = Key>,
46    NetS: Sender<PublicKey = P>,
47    NetR: Receiver<PublicKey = P>,
48    Con::Subscriber: Eq,
49{
50    /// Context used to spawn tasks, manage time, etc.
51    context: ContextCell<E>,
52
53    /// Produces data for incoming requests
54    producer: Pro,
55
56    /// Manages the list of peers that can be used to fetch data
57    peer_provider: D,
58
59    /// The blocker that will be used to block peers that send invalid responses
60    blocker: B,
61
62    /// Used to detect changes in the peer set
63    last_peer_set_id: Option<u64>,
64
65    /// Mailbox that makes and prunes fetches
66    mailbox: mailbox::Receiver<Message<Key, P, Con::Subscriber>>,
67
68    /// Manages outgoing fetch requests
69    fetcher: Fetcher<E, P, Key, NetS>,
70
71    /// Tracks all in-flight fetch state
72    inflight: Inflight<Con, P>,
73
74    /// Subscribers that keep each fetch alive.
75    subscribers: subscribers::Tracker<Key, Con::Subscriber>,
76
77    /// Holds futures that resolve once the `Producer` has produced the data.
78    /// Once the future is resolved, the data (or an error) is sent to the peer.
79    /// Has unbounded size; the number of concurrent requests should be limited
80    /// by the `Producer` which may drop requests.
81    serves: FuturesPool<Serve<P>>,
82
83    /// Whether responses are sent with priority over other network messages
84    priority_responses: bool,
85
86    /// Metrics for the peer actor
87    metrics: metrics::Metrics,
88
89    /// Phantom data for networking types
90    _r: PhantomData<NetR>,
91}
92
93impl<E, P, D, B, Key, Con, Pro, NetS, NetR> Engine<E, P, D, B, Key, Con, Pro, NetS, NetR>
94where
95    E: BufferPooler + Clock + Spawner + Rng + Metrics,
96    P: PublicKey,
97    D: Provider<PublicKey = P>,
98    B: Blocker<PublicKey = P>,
99    Key: Span,
100    Con: Consumer<Key = Key, Value = Bytes>,
101    Pro: Producer<Key = Key>,
102    NetS: Sender<PublicKey = P>,
103    NetR: Receiver<PublicKey = P>,
104    Con::Subscriber: Clone + Ord + Send + 'static,
105{
106    /// Creates a new `Actor` with the given configuration.
107    ///
108    /// Returns the actor and a mailbox to send messages to it.
109    pub fn new(
110        context: E,
111        cfg: Config<P, D, B, Key, Con, Pro>,
112    ) -> (Self, Mailbox<Key, P, Con::Subscriber>) {
113        let (sender, receiver) = mailbox::new(context.child("mailbox"), cfg.mailbox_size);
114
115        let metrics = metrics::Metrics::init(&context);
116        let fetcher = Fetcher::new(
117            context.child("fetcher"),
118            FetcherConfig {
119                me: cfg.me,
120                initial: cfg.initial,
121                timeout: cfg.timeout,
122                retry_timeout: cfg.fetch_retry_timeout,
123                priority_requests: cfg.priority_requests,
124            },
125        );
126        (
127            Self {
128                context: ContextCell::new(context),
129                producer: cfg.producer,
130                peer_provider: cfg.peer_provider,
131                blocker: cfg.blocker,
132                last_peer_set_id: None,
133                mailbox: receiver,
134                fetcher,
135                inflight: Inflight::new(cfg.consumer),
136                subscribers: subscribers::Tracker::new(),
137                serves: FuturesPool::default(),
138                priority_responses: cfg.priority_responses,
139                metrics,
140                _r: PhantomData,
141            },
142            Mailbox::new(sender),
143        )
144    }
145
146    /// Runs the actor until the context is stopped.
147    ///
148    /// The actor will handle:
149    /// - Fetching data from other peers and notifying the `Consumer`
150    /// - Serving data to other peers by requesting it from the `Producer`
151    pub fn start(mut self, network: (NetS, NetR)) -> Handle<()> {
152        spawn_cell!(self.context, self.run(network))
153    }
154
155    /// Inner run loop called by `start`.
156    async fn run(mut self, network: (NetS, NetR)) {
157        // Wrap channel
158        let (mut sender, mut receiver) = wrap(
159            (),
160            self.context.network_buffer_pool().clone(),
161            network.0,
162            network.1,
163        );
164        let mut peer_set_subscription = self.peer_provider.subscribe().await;
165
166        select_loop! {
167            self.context,
168            on_start => {
169                // Update metrics
170                let _ = self
171                    .metrics
172                    .fetch_pending
173                    .try_set(self.fetcher.len_pending());
174                let _ = self.metrics.fetch_active.try_set(self.fetcher.len_active());
175                let _ = self
176                    .metrics
177                    .peers_blocked
178                    .try_set(self.fetcher.len_blocked());
179                let _ = self.metrics.serve_processing.try_set(self.serves.len());
180
181                // Get retry timeout (if any)
182                let deadline_pending = match self.fetcher.get_pending_deadline() {
183                    Some(deadline) => Either::Left(self.context.sleep_until(deadline)),
184                    None => Either::Right(future::pending()),
185                };
186
187                // Get requester timeout (if any)
188                let deadline_active = match self.fetcher.get_active_deadline() {
189                    Some(deadline) => Either::Left(self.context.sleep_until(deadline)),
190                    None => Either::Right(future::pending()),
191                };
192            },
193            on_stopped => {
194                debug!("shutdown");
195                self.inflight.drain();
196                self.subscribers.clear();
197                self.serves.cancel_all();
198            },
199            // Handle peer set updates
200            Some(update) = peer_set_subscription.recv() else {
201                debug!("peer set subscription closed");
202                return;
203            } => {
204                if self.last_peer_set_id < Some(update.index) {
205                    self.last_peer_set_id = Some(update.index);
206                    self.fetcher.reconcile(update.latest.primary.as_ref());
207                }
208            },
209            // Handle active deadline
210            _ = deadline_active => {
211                if let Some(key) = self.fetcher.pop_active() {
212                    debug!(?key, "requester timeout");
213                    self.metrics.fetch.inc(Status::Failure);
214                    self.fetcher.add_retry(key);
215                }
216            },
217            // Handle pending deadline
218            _ = deadline_pending => {
219                self.fetcher.fetch(&mut sender);
220            },
221            // Handle mailbox messages
222            Some(msg) = self.mailbox.recv() else {
223                error!("mailbox closed");
224                return;
225            } => {
226                match msg {
227                    Message::Fetch(keys) => {
228                        for FetchKey {
229                            key,
230                            subscribers,
231                            metadata: targets,
232                        } in keys
233                        {
234                            trace!(?key, "mailbox: fetch");
235
236                            // Check if the fetch is already in progress
237                            let is_new = !self.inflight.contains(&key);
238                            self.subscribers.insert(key.clone(), subscribers);
239
240                            // Update targets
241                            match targets {
242                                Some(targets) => {
243                                    // Only add targets if this is a new fetch OR the existing
244                                    // fetch already has targets. Don't restrict an "all" fetch
245                                    // (no targets) to specific targets.
246                                    if is_new || self.fetcher.has_targets(&key) {
247                                        self.fetcher.add_targets(key.clone(), targets);
248                                    }
249                                }
250                                None => self.fetcher.clear_targets(&key),
251                            }
252
253                            // Only start new fetch if not already in progress
254                            if is_new {
255                                self.inflight.insert(
256                                    key.clone(),
257                                    self.metrics.fetch_duration.timer(self.context.as_ref()),
258                                );
259                                self.fetcher.add_ready(key);
260                            } else {
261                                trace!(?key, "updated targets for existing fetch");
262                            }
263                        }
264                    }
265                    Message::Retain { predicate } => {
266                        trace!("mailbox: retain");
267
268                        self.subscribers
269                            .retain(|key, subscriber| predicate(key, subscriber));
270                        let subscribers = &self.subscribers;
271                        self.fetcher.retain(|key| subscribers.contains(key));
272                        let count = self.inflight.retain(|key| subscribers.contains(key)) as u64;
273                        self.record_cancellations(count);
274                    }
275                }
276            },
277            // Handle completed consumer deliveries
278            delivery = self.inflight.next_delivery() => {
279                // If the delivery was aborted, its inflight entry was dropped (via
280                // Retain or shutdown) before the consumer finished validating.
281                let (peer, delivery, result) = match delivery {
282                    Ok(delivery) => delivery,
283                    Err(_) => continue,
284                };
285                self.handle_delivery(peer, delivery, result);
286            },
287            // Handle completed server requests
288            serve = self.serves.next_completed() => {
289                let Serve {
290                    timer,
291                    peer,
292                    id,
293                    result,
294                } = serve;
295
296                // Metrics and logs
297                match result {
298                    Ok(_) => {
299                        timer.observe(self.context.as_ref());
300                        self.metrics.serve.inc(Status::Success);
301                    }
302                    Err(ref err) => {
303                        debug!(?err, ?peer, ?id, "serve failed");
304                        self.metrics.serve.inc(Status::Failure);
305                    }
306                }
307
308                // Send response to peer
309                self.handle_serve(&mut sender, peer, id, result, self.priority_responses);
310            },
311            // Handle network messages
312            msg = receiver.recv() => {
313                // Break if the receiver is closed
314                let (peer, msg) = match msg {
315                    Ok(msg) => msg,
316                    Err(err) => {
317                        error!(?err, "receiver closed");
318                        return;
319                    }
320                };
321
322                // Skip if there is a decoding error
323                let msg = match msg {
324                    Ok(msg) => msg,
325                    Err(err) => {
326                        trace!(?err, ?peer, "decode failed");
327                        continue;
328                    }
329                };
330                match msg.payload {
331                    wire::Payload::Request(key) => self.handle_network_request(peer, msg.id, key),
332                    wire::Payload::Response(response) => {
333                        self.handle_network_response(peer, msg.id, response)
334                    }
335                    wire::Payload::Error => self.handle_network_error_response(peer, msg.id),
336                };
337            },
338        }
339    }
340
341    /// Record cancellation metrics for a retain-style operation.
342    fn record_cancellations(&mut self, count: u64) {
343        if count == 0 {
344            self.metrics.cancel.inc(Status::Dropped);
345        } else {
346            self.metrics.cancel.inc_by(Status::Success, count);
347        }
348    }
349
350    /// Handles the case where the application responds to a request from an external peer.
351    fn handle_serve(
352        &mut self,
353        sender: &mut WrappedSender<NetS, wire::Message<Key>>,
354        peer: P,
355        id: u64,
356        response: Result<Bytes, oneshot::error::RecvError>,
357        priority: bool,
358    ) {
359        // Encode message
360        let payload: wire::Payload<Key> = response.map_or_else(
361            |_| wire::Payload::Error,
362            |data| wire::Payload::Response(data),
363        );
364        let msg = wire::Message { id, payload };
365
366        // Send message to peer
367        let result = sender.send(Recipients::One(peer.clone()), msg, priority);
368
369        // Log result, but do not handle errors.
370        if result.is_empty() {
371            warn!(?peer, ?id, "serve send failed");
372        } else {
373            trace!(?peer, ?id, "serve sent");
374        };
375    }
376
377    /// Handle a network request from a peer.
378    fn handle_network_request(&mut self, peer: P, id: u64, key: Key) {
379        // Serve the request
380        trace!(?peer, ?id, "peer request");
381        let mut producer = self.producer.clone();
382        let timer = self.metrics.serve_duration.timer(self.context.as_ref());
383        let receiver = producer.produce(key);
384        self.serves.push(async move {
385            let result = receiver.await;
386            Serve {
387                timer,
388                peer,
389                id,
390                result,
391            }
392        });
393    }
394
395    /// Handle a network response from a peer.
396    fn handle_network_response(&mut self, peer: P, id: u64, response: Bytes) {
397        trace!(?peer, ?id, "peer response: data");
398
399        // Get the key associated with the response, if any
400        let Some(key) = self.fetcher.pop_by_id(id, &peer, true) else {
401            // It's possible that the key does not exist if the request was pruned.
402            return;
403        };
404
405        let Some(subscribers) = self.subscribers.pending(&key) else {
406            warn!(?key, "response for fetch with no subscribers");
407            self.inflight.cancel(&key);
408            return;
409        };
410        let delivery = Delivery { key, subscribers };
411
412        // The peer had the data, so deliver it to the consumer without blocking the engine.
413        self.inflight.deliver(delivery, peer, response);
414    }
415
416    /// Handle completed delivery to the consumer.
417    fn handle_delivery(&mut self, peer: P, delivery: Delivery<Key, Con::Subscriber>, valid: bool) {
418        let Delivery {
419            key,
420            subscribers: delivered,
421            ..
422        } = delivery;
423
424        if valid {
425            let already_accepted = self.inflight.response_accepted(&key);
426
427            // Remove only the subscribers that accepted this response. If other
428            // subscribers still need the key, deliver the same accepted response
429            // locally with the remaining annotations.
430            let remaining = self
431                .subscribers
432                .remove_delivered(&key, delivered.map_into(|(subscriber, _)| subscriber));
433
434            if let Some(subscribers) = remaining {
435                if !already_accepted {
436                    self.metrics.fetch.inc(Status::Success);
437                    self.inflight.accept_response(&key, self.context.as_ref());
438                }
439                self.inflight.redeliver(Delivery { key, subscribers });
440            } else {
441                // All subscribers observed a valid response; clear any targeting
442                // state retained for this key.
443                if !already_accepted {
444                    self.metrics.fetch.inc(Status::Success);
445                }
446                self.inflight.complete(self.context.as_ref(), &key);
447                self.fetcher.clear_targets(&key);
448            }
449            return;
450        }
451
452        if self.inflight.response_accepted(&key) {
453            warn!(
454                ?key,
455                "previously accepted response was rejected during local redelivery"
456            );
457            self.metrics.fetch.inc(Status::Failure);
458            self.inflight.complete(self.context.as_ref(), &key);
459            self.subscribers.remove(&key);
460            self.fetcher.clear_targets(&key);
461            return;
462        }
463
464        // If the data is invalid, block the peer and try again. Blocking the
465        // peer also removes any targets associated with it.
466        commonware_p2p::block!(self.blocker, peer.clone(), "invalid data received");
467        self.fetcher.block(peer);
468        self.metrics.fetch.inc(Status::Failure);
469        self.inflight.discard_response(&key);
470        self.fetcher.add_retry(key);
471    }
472
473    /// Handle a network response from a peer that did not have the data.
474    fn handle_network_error_response(&mut self, peer: P, id: u64) {
475        trace!(?peer, ?id, "peer response: error");
476
477        // Get the key associated with the response, if any
478        let Some(key) = self.fetcher.pop_by_id(id, &peer, false) else {
479            // It's possible that the key does not exist if the request was pruned.
480            return;
481        };
482
483        // The peer did not have the data, so we need to try again
484        self.metrics.fetch.inc(Status::Failure);
485        self.fetcher.add_retry(key);
486    }
487}