use crate::stateful::{
Application,
actor::{
core::{mailbox::Message, processing::Processing, syncing::Syncing},
metrics::Metrics as StatefulMetrics,
processor::{PendingSyncTargets, Processor, Pruning},
syncer::{self, SyncPlan, SyncResult},
},
db::{AttachableResolverSet, DatabaseSet, StateSyncSet, SyncEngineConfig},
};
use commonware_actor::mailbox::{self as actor_mailbox};
use commonware_consensus::{
marshal::{
ancestry::BlockProvider,
core::{Floor, Mailbox as MarshalMailbox, Variant},
},
simplex::types::Finalization,
};
use commonware_cryptography::{Digestible, certificate::Scheme};
use commonware_runtime::{ContextCell, Handle, Spawner, spawn_cell, telemetry::metrics::GaugeExt};
use commonware_storage::Context;
use commonware_utils::channel::oneshot;
use futures::join;
use rand_core::Rng;
use std::num::NonZeroUsize;
mod mailbox;
pub use mailbox::Mailbox;
pub(super) use mailbox::Verification;
mod processing;
mod syncing;
mod verifications;
type BlockDigest<A, E> = <<A as Application<E>>::Block as Digestible>::Digest;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct PruneConfig {
pub maintenance_interval: NonZeroUsize,
pub retained_marshal_blocks: usize,
pub retained_qmdb_blocks: usize,
}
impl PruneConfig {
pub const fn assert_valid(self) {
assert!(
self.retained_marshal_blocks >= self.retained_qmdb_blocks,
"marshal must retain at least as many blocks as QMDB",
);
}
}
pub struct Config<E, A, S, V, R>
where
E: Rng + Spawner + Context,
A: Application<E>,
S: Scheme,
V: Variant<ApplicationBlock = A::Block>,
{
pub application: A,
pub db_config: <A::Databases as DatabaseSet<E>>::Config,
pub provider: A::Provider,
pub marshal: (MarshalMailbox<S, V>, Floor),
pub mailbox_size: NonZeroUsize,
pub plan: SyncPlan<E, S, V>,
pub resolvers: R,
pub sync_config: SyncEngineConfig,
pub prune_config: Option<PruneConfig>,
}
pub struct Stateful<E, A, S, V, R>
where
E: Rng + Spawner + Context,
A: Application<E>,
S: Scheme,
V: Variant<ApplicationBlock = A::Block>,
{
context: ContextCell<E>,
mailbox: actor_mailbox::Receiver<Message<E, A>>,
application: A,
provider: A::Provider,
marshal: (MarshalMailbox<S, V>, Floor),
db_config: <A::Databases as DatabaseSet<E>>::Config,
plan: SyncPlan<E, S, V>,
resolvers: R,
sync_config: SyncEngineConfig,
pruning: Option<Pruning<PendingSyncTargets<A, E>>>,
}
impl<E, A, S, V, R> Stateful<E, A, S, V, R>
where
E: Rng + Spawner + Context,
A: Application<E>,
A::Databases: StateSyncSet<E, R, BlockDigest<A, E>>,
S: Scheme,
V: Variant<ApplicationBlock = A::Block>,
R: AttachableResolverSet<A::Databases>,
MarshalMailbox<S, V>: BlockProvider<Block = A::Block>,
{
pub fn init(mut context: E, config: Config<E, A, S, V, R>) -> (Self, Mailbox<E, A>) {
let pruning = config.prune_config.map(|prune_config| {
Pruning::random(
prune_config,
config.marshal.0.max_pending_acks(),
&mut context,
)
});
let (sender, mailbox) = actor_mailbox::new(context.child("mailbox"), config.mailbox_size);
(
Self {
context: ContextCell::new(context),
mailbox,
application: config.application,
provider: config.provider,
marshal: config.marshal,
db_config: config.db_config,
plan: config.plan,
resolvers: config.resolvers,
sync_config: config.sync_config,
pruning,
},
Mailbox::new(sender),
)
}
pub fn start(mut self) -> Handle<()> {
spawn_cell!(self.context, self.run())
}
async fn run(self) {
if let Some(floor) = self.plan.floor().cloned() {
self.start_state_sync(floor).await;
} else if self.plan.requires_state_sync_floor() {
panic!("interrupted state sync is missing its persisted floor");
} else {
self.start_from_marshal().await;
}
}
async fn start_state_sync(self, finalization: Finalization<S, V::Commitment>) {
let (marshal, floor) = self.marshal;
let metrics = StatefulMetrics::new(self.context.as_present());
let sync_metadata = self
.plan
.into_sync_metadata()
.begin_sync(finalization.clone())
.await;
let (sync_complete, sync_completed) = oneshot::channel();
let (syncer, syncer_mailbox) = syncer::Syncer::new(syncer::Config {
context: self.context.child("syncer"),
db_config: self.db_config,
sync_config: self.sync_config,
resolvers: self.resolvers.clone(),
finalization,
marshal: (marshal.clone(), floor),
sync_complete,
});
let syncing = Syncing {
context: self.context,
mailbox: self.mailbox,
application: self.application,
provider: self.provider,
marshal,
sync_metadata,
syncer: syncer_mailbox,
deferred_verifications: Vec::new(),
database_subscribers: Vec::new(),
artifact: None,
resolvers: self.resolvers,
sync_completed,
pending_finalizations: Default::default(),
pruning: self.pruning,
metrics,
};
let _ = join!(syncer.start(), syncing.start());
}
async fn start_from_marshal(self) {
let (marshal, _) = self.marshal;
let syncer::StartupResult {
sync: SyncResult { databases, anchor },
skip_finalized_until,
} = syncer::init_databases_from_marshal::<E, A, S, V>(
self.context.as_present(),
&marshal,
self.db_config,
self.plan.into_sync_metadata(),
)
.await;
self.resolvers.attach_databases(databases.clone()).await;
let metrics = StatefulMetrics::new(self.context.as_present());
let _ = metrics.sync_done.try_set(1);
let processor = Processor::new(self.application, databases, anchor, metrics, self.pruning);
Processing {
context: self.context,
mailbox: self.mailbox,
provider: self.provider,
marshal,
processor,
deferred_verifications: Vec::new(),
skip_finalized_until,
}
.start()
.await
}
}
#[cfg(test)]
mod tests {
use super::{Config, Stateful};
use crate::stateful::{
actor::syncer::SyncPlan,
db::{AttachableResolver, Shared, StateSyncDb, SyncEngineConfig},
tests::{
fixtures,
mocks::{TestApp, TestBlock, TestDb},
},
};
use commonware_consensus::{
Application as _, CertifiableBlock as _, Reporter as _,
marshal::{Update, ancestry},
simplex::mocks::scheme as scheme_mocks,
};
use commonware_cryptography::sha256::Digest as Sha256Digest;
use commonware_macros::select;
use commonware_runtime::{Clock as _, Runner as _, Supervisor as _, deterministic};
use commonware_utils::{
Acknowledgement as _, NZU64, NZUsize,
acknowledgement::Exact,
channel::{mpsc, oneshot},
sync::Mutex,
};
use futures::poll;
use std::{convert::Infallible, sync::Arc, time::Duration};
struct StartupGate {
started: oneshot::Sender<()>,
release: oneshot::Receiver<()>,
}
#[derive(Clone, Default)]
struct NoopResolver {
startup_gate: Arc<Mutex<Option<StartupGate>>>,
}
impl NoopResolver {
fn gated() -> (Self, oneshot::Receiver<()>, oneshot::Sender<()>) {
let (started, started_rx) = oneshot::channel();
let (release, release_rx) = oneshot::channel();
(
Self {
startup_gate: Arc::new(Mutex::new(Some(StartupGate {
started,
release: release_rx,
}))),
},
started_rx,
release,
)
}
}
impl AttachableResolver<TestDb> for NoopResolver {
async fn attach_database(&self, _db: Shared<TestDb>) {
let Some(StartupGate {
started,
mut release,
}) = self.startup_gate.lock().take()
else {
return;
};
started
.send(())
.expect("test should await the startup gate");
let _ = (&mut release).await;
}
}
impl StateSyncDb<deterministic::Context, NoopResolver> for TestDb {
type SyncError = Infallible;
async fn sync_db(
_context: deterministic::Context,
_config: Self::Config,
_resolver: NoopResolver,
_target: Self::SyncTarget,
_tip_updates: mpsc::Receiver<Self::SyncTarget>,
_finish: Option<mpsc::Receiver<()>>,
_reached_target: Option<mpsc::Sender<Self::SyncTarget>>,
_sync_config: SyncEngineConfig,
) -> Result<Self, Self::SyncError> {
Ok(Self::default())
}
}
#[test]
fn mailbox_rejects_propose_while_floor_resolution_waits() {
deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move {
let mut signing_context = context.child("signing");
let fixture = scheme_mocks::fixture(&mut signing_context, b"pending-floor", 1);
let finalization = fixtures::finalization(&fixture, 1, Sha256Digest::from([7; 32]));
let marshal = fixtures::marshal_fixture(
context.child("marshal_fixture"),
"pending-floor",
fixture.schemes[0].clone(),
None,
NZUsize!(1),
false,
)
.await;
let plan = SyncPlan::init(&context, "pending-floor-stateful".to_string()).await;
let (stateful, mut mailbox) = Stateful::init(
context.child("stateful"),
Config {
application: TestApp::default(),
db_config: (),
provider: (),
marshal: (marshal.mailbox, marshal.floor),
mailbox_size: NZUsize!(8),
plan: plan.with_floor(finalization),
resolvers: NoopResolver::default(),
sync_config: SyncEngineConfig {
fetch_batch_size: NZU64!(1),
apply_batch_size: NZU64!(1),
max_outstanding_requests: 1,
update_channel_size: NZUsize!(1),
max_retained_roots: 1,
},
prune_config: None,
},
);
let handle = stateful.start();
select! {
result = mailbox.propose(
(context.child("proposal"), TestBlock::new(1, 1).context()),
ancestry::from_iter([]),
(),
) => {
assert!(result.is_none());
},
_ = context.sleep(Duration::from_millis(100)) => {
panic!("stateful mailbox stalled while resolving state sync floor");
},
}
handle.abort();
});
}
#[test]
fn startup_recovery_releases_cancelled_verify_ancestries() {
deterministic::Runner::timed(Duration::from_secs(5)).start(|mut context| async move {
let prefix = "startup-recovery-cancelled-verifications";
let scheme = scheme_mocks::fixture(&mut context, prefix.as_bytes(), 1);
let genesis = TestBlock::new(0, 0);
let finalized = TestBlock::child(&genesis, 1);
let marshal = fixtures::marshal_fixture(
context.child("marshal"),
prefix,
scheme.schemes[0].clone(),
None,
NZUsize!(8),
true,
)
.await;
let (resolver, startup_started, startup_release) = NoopResolver::gated();
let plan = SyncPlan::init(&context, format!("{prefix}-stateful")).await;
let (stateful, mut mailbox) = Stateful::init(
context.child("stateful"),
Config {
application: TestApp::default(),
db_config: (),
provider: (),
marshal: (marshal.mailbox.clone(), marshal.floor),
mailbox_size: NZUsize!(1),
plan,
resolvers: resolver,
sync_config: SyncEngineConfig {
fetch_batch_size: NZU64!(1),
apply_batch_size: NZU64!(1),
max_outstanding_requests: 1,
update_channel_size: NZUsize!(1),
max_retained_roots: 1,
},
prune_config: None,
},
);
let actor = stateful.start();
startup_started
.await
.expect("startup should reach resolver attachment before processing");
let owners = [
Arc::new(TestBlock::new(2, 2)),
Arc::new(TestBlock::new(3, 3)),
Arc::new(TestBlock::new(4, 4)),
];
let weak_owners = owners.iter().map(Arc::downgrade).collect::<Vec<_>>();
let mut first_mailbox = mailbox.clone();
let mut first = Box::pin(first_mailbox.verify(
(context.child("verify_first"), owners[0].context()),
ancestry::from_iter([Arc::clone(&owners[0])]),
));
assert!(poll!(&mut first).is_pending());
let mut second_mailbox = mailbox.clone();
let mut second = Box::pin(second_mailbox.verify(
(context.child("verify_second"), owners[1].context()),
ancestry::from_iter([Arc::clone(&owners[1])]),
));
assert!(poll!(&mut second).is_pending());
let mut third_mailbox = mailbox.clone();
let mut third = Box::pin(third_mailbox.verify(
(context.child("verify_third"), owners[2].context()),
ancestry::from_iter([Arc::clone(&owners[2])]),
));
assert!(poll!(&mut third).is_pending());
let (acknowledgement, mut acknowledgement_waiter) = Exact::handle();
let _ = mailbox.report(Update::Block(Arc::new(finalized), acknowledgement));
drop(first);
drop(second);
drop(third);
drop(owners);
context.sleep(Duration::from_millis(10)).await;
assert!(poll!(&mut acknowledgement_waiter).is_pending());
for (index, owner) in weak_owners.iter().enumerate() {
assert!(
owner.upgrade().is_none(),
"cancelled startup verification {index} retained its ancestry owner",
);
}
startup_release
.send(())
.expect("startup should remain gated");
select! {
result = acknowledgement_waiter => {
result.expect("finalized block should be acknowledged after startup");
},
_ = context.sleep(Duration::from_millis(100)) => {
panic!("finalized acknowledgement stalled after startup");
},
}
actor.abort();
drop(marshal.guards);
});
}
}