use super::mailbox::{Mailbox, Message};
use crate::{
dkg::{ReshareBlock, network::Manager, probe::Bootstrap, types::EpochInfo},
stateful::probe::sample::Sample,
};
use commonware_actor::mailbox::{self as actor_mailbox, Receiver as ActorReceiver};
use commonware_codec::Read;
use commonware_consensus::{marshal::core::Variant, simplex::scheme::Scheme, types::FixedEpocher};
use commonware_cryptography::Signer;
use commonware_p2p::{Blocker, Receiver, Sender};
use commonware_parallel::Strategy;
use commonware_runtime::{Clock, ContextCell, Handle, Metrics, Spawner, spawn_cell};
use commonware_utils::NonZeroDuration;
use discovery::Discovery;
use rand_core::CryptoRng;
use std::num::{NonZeroU64, NonZeroUsize};
mod discovery;
mod service;
pub struct Config<E, M, S, V, T, B>
where
E: Spawner + CryptoRng + Clock + Metrics,
M: Manager<
PublicKey = S::PublicKey,
Directory = <V::ApplicationBlock as ReshareBlock>::Directory,
>,
S: Scheme<V::Commitment>,
V: Variant,
V::ApplicationBlock: ReshareBlock,
<V::ApplicationBlock as ReshareBlock>::Signer: Signer<PublicKey = S::PublicKey>,
T: Strategy,
B: Blocker<PublicKey = S::PublicKey>,
{
pub context: E,
pub manager: M,
pub bootstrap: Bootstrap<S::PublicKey, <V::ApplicationBlock as ReshareBlock>::Directory>,
pub verifier: S,
pub genesis: EpochInfo<
<V::ApplicationBlock as ReshareBlock>::Variant,
S::PublicKey,
<V::ApplicationBlock as ReshareBlock>::Directory,
>,
pub strategy: T,
pub blocker: B,
pub blocks_per_epoch: NonZeroU64,
pub retry_timeout: NonZeroDuration,
pub mailbox_size: NonZeroUsize,
pub block_codec_config: <V::ApplicationBlock as Read>::Cfg,
}
pub struct Actor<E, M, S, V, T, B>
where
E: Spawner + CryptoRng + Clock + Metrics,
M: Manager<
PublicKey = S::PublicKey,
Directory = <V::ApplicationBlock as ReshareBlock>::Directory,
>,
S: Scheme<V::Commitment>,
V: Variant,
V::ApplicationBlock: ReshareBlock,
<V::ApplicationBlock as ReshareBlock>::Signer: Signer<PublicKey = S::PublicKey>,
T: Strategy,
B: Blocker<PublicKey = S::PublicKey>,
{
context: ContextCell<E>,
mailbox: ActorReceiver<Message<S, V>>,
manager: M,
bootstrap: Bootstrap<S::PublicKey, <V::ApplicationBlock as ReshareBlock>::Directory>,
verifier: S,
genesis: EpochInfo<
<V::ApplicationBlock as ReshareBlock>::Variant,
S::PublicKey,
<V::ApplicationBlock as ReshareBlock>::Directory,
>,
strategy: T,
blocker: B,
blocks_per_epoch: NonZeroU64,
retry_timeout: NonZeroDuration,
block_codec_config: <V::ApplicationBlock as Read>::Cfg,
}
impl<E, M, S, V, T, B> Actor<E, M, S, V, T, B>
where
E: Spawner + CryptoRng + Clock + Metrics,
M: Manager<
PublicKey = S::PublicKey,
Directory = <V::ApplicationBlock as ReshareBlock>::Directory,
>,
S: Scheme<V::Commitment>,
V: Variant,
V::ApplicationBlock: ReshareBlock,
<V::ApplicationBlock as ReshareBlock>::Signer: Signer<PublicKey = S::PublicKey>,
T: Strategy,
B: Blocker<PublicKey = S::PublicKey>,
{
pub fn new(config: Config<E, M, S, V, T, B>) -> (Self, Mailbox<S, V>) {
let (sender, mailbox) =
actor_mailbox::new(config.context.child("mailbox"), config.mailbox_size);
let mailbox_handle = Mailbox::new(sender);
(
Self {
context: ContextCell::new(config.context),
mailbox,
manager: config.manager,
bootstrap: config.bootstrap,
verifier: config.verifier,
genesis: config.genesis,
strategy: config.strategy,
blocker: config.blocker,
blocks_per_epoch: config.blocks_per_epoch,
retry_timeout: config.retry_timeout,
block_codec_config: config.block_codec_config,
},
mailbox_handle,
)
}
pub fn start<BSE, BRE>(mut self, boundaries: (BSE, BRE)) -> Handle<()>
where
BSE: Sender<PublicKey = S::PublicKey>,
BRE: Receiver<PublicKey = S::PublicKey>,
{
spawn_cell!(self.context, self.run(boundaries,))
}
async fn run<BSE, BRE>(self, (boundary_sender, boundary_receiver): (BSE, BRE))
where
BSE: Sender<PublicKey = S::PublicKey>,
BRE: Receiver<PublicKey = S::PublicKey>,
{
Discovery {
context: self.context,
mailbox: self.mailbox,
manager: self.manager,
sample: Sample::new(self.bootstrap.epoch),
bootstrap_participants: self.bootstrap.participants,
bootstrap_directory: self.bootstrap.directory,
verifier: self.verifier,
genesis: self.genesis,
strategy: self.strategy,
blocker: self.blocker,
epocher: FixedEpocher::new(self.blocks_per_epoch),
block_codec_config: self.block_codec_config,
retry_timeout: self.retry_timeout,
artifact: None,
subscribers: Vec::new(),
pending: None,
}
.run(boundary_sender, boundary_receiver)
.await;
}
}