use super::mailbox::{Mailbox, Message};
use commonware_actor::mailbox::Receiver as ActorReceiver;
use commonware_consensus::{marshal::core::Variant, simplex::scheme::Scheme, types::Epoch};
use commonware_cryptography::{certificate::Provider, PublicKey};
use commonware_p2p::{Blocker, Receiver, Sender};
use commonware_parallel::Strategy;
use commonware_runtime::{spawn_cell, Clock, ContextCell, Handle, Metrics, Spawner};
use commonware_utils::NonZeroDuration;
use discovery::Discovery;
use rand_core::CryptoRng;
use std::num::NonZeroUsize;
mod discovery;
mod service;
pub struct Config<E, D, T, P, B>
where
E: Spawner + CryptoRng + Clock + Metrics,
D: Provider<Scope = Epoch>,
T: Strategy,
P: PublicKey,
B: Blocker<PublicKey = P>,
{
pub context: E,
pub provider: D,
pub strategy: T,
pub capacity: NonZeroUsize,
pub blocker: B,
pub minimum_epoch: Epoch,
pub retry_timeout: NonZeroDuration,
}
pub struct Probe<E, S, D, V, T, P, B>
where
E: Spawner + CryptoRng + Clock + Metrics,
S: Scheme<V::Commitment>,
D: Provider<Scope = Epoch, Scheme = S>,
V: Variant,
T: Strategy,
P: PublicKey,
B: Blocker<PublicKey = P>,
{
context: ContextCell<E>,
mailbox: ActorReceiver<Message<S, V>>,
provider: D,
strategy: T,
blocker: B,
minimum_epoch: Epoch,
retry_timeout: NonZeroDuration,
}
impl<E, S, D, V, T, P, B> Probe<E, S, D, V, T, P, B>
where
E: Spawner + CryptoRng + Clock + Metrics,
S: Scheme<V::Commitment>,
D: Provider<Scope = Epoch, Scheme = S>,
V: Variant,
T: Strategy,
P: PublicKey,
B: Blocker<PublicKey = P>,
{
pub fn new(config: Config<E, D, T, P, B>) -> (Self, Mailbox<S, V>) {
let (sender, receiver) =
commonware_actor::mailbox::new(config.context.child("mailbox"), config.capacity);
let mailbox = Mailbox::new(sender);
(
Self {
context: ContextCell::new(config.context),
mailbox: receiver,
provider: config.provider,
strategy: config.strategy,
blocker: config.blocker,
minimum_epoch: config.minimum_epoch,
retry_timeout: config.retry_timeout,
},
mailbox,
)
}
pub fn start(
mut self,
net: (impl Sender<PublicKey = P>, impl Receiver<PublicKey = P>),
) -> Handle<()> {
spawn_cell!(self.context, self.run(net))
}
async fn run(
self,
(mut sender, mut receiver): (impl Sender<PublicKey = P>, impl Receiver<PublicKey = P>),
) {
Discovery {
context: self.context,
mailbox: self.mailbox,
provider: self.provider,
strategy: self.strategy,
blocker: self.blocker,
minimum_epoch: self.minimum_epoch,
retry_timeout: self.retry_timeout,
floor: None,
floor_subscribers: Vec::new(),
}
.run(&mut sender, &mut receiver)
.await;
}
}