use crate::stateful::{ Application, Input, actor::{ core::{ mailbox::Message, verifications::{Handler as Verifications, Request as VerificationRequest}, }, processor::{Applied, Processor}, }, db::{Barrier, DatabaseSet}, }; use commonware_actor::mailbox as actor_mailbox; use commonware_consensus::{ Heightable, marshal::{ ancestry::BlockProvider, core::{Mailbox as MarshalMailbox, Variant}, }, types::Height, }; use commonware_cryptography::certificate::Scheme; use commonware_macros::{select, select_loop}; use commonware_runtime::{Clock, ContextCell, Handle, Metrics, Spawner}; use commonware_utils::{ Acknowledgement as _, acknowledgement::Exact, channel::fallible::OneshotExt, }; use futures::{ FutureExt as _, future::{Either, pending, ready}, }; use rand_core::Rng; use std::{collections::VecDeque, sync::mpsc::TryRecvError}; use tracing::{Instrument as _, debug, info_span}; /// Work selected for one iteration of the processing actor. enum Step { /// A message received from the actor mailbox. Message(M), /// Deferred pruning work ready for its database mutation boundary. Prune(P), /// Completion of the active database durability barrier. Sync((Height, bool)), } /// Tracks the durable database prefix and marshal acknowledgements awaiting it. /// /// At most one sync covers a captured prefix. Applied heights beyond that prefix remain queued /// for a successor sync. struct Durability { /// Highest applied height known to be durable. durable: Height, /// Applied heights whose marshal acknowledgements await durability. acknowledgements: VecDeque<(Height, Exact)>, /// Active barrier, whose output includes the height of its captured prefix. sync: Option>, } impl Durability { /// Initialize tracking at a height already known to be durable. const fn new(height: Height) -> Self { Self { durable: height, acknowledgements: VecDeque::new(), sync: None, } } /// Return the highest applied height, or the durable floor when none are pending. fn latest_applied(&self) -> Height { self.acknowledgements .back() .map_or(self.durable, |(height, _)| *height) } /// Record a newly applied height and retain its acknowledgement until durability. /// /// Heights must be recorded in strictly increasing order. fn applied(&mut self, height: Height, acknowledgement: Exact) { assert!( height > self.latest_applied(), "finalized heights must increase" ); self.acknowledgements.push_back((height, acknowledgement)); } /// Return whether applied state remains uncovered and no sync is active. fn needs_sync(&self) -> bool { self.sync.is_none() && self.durable < self.latest_applied() } /// Record a barrier covering applied state through `height`. /// /// Only one barrier may be active, and `height` must extend the durable prefix without /// exceeding the latest applied height. fn started(&mut self, height: Height, barrier: Barrier) { assert!(self.sync.is_none(), "sync already active"); assert!(height > self.durable && height <= self.latest_applied()); self.sync = Some(Handle::from_future(async move { Ok((height, barrier.durable().await)) })); } /// Complete the active sync and acknowledge every height it made durable. /// /// Returns false without advancing the durable prefix when durability was not established. fn complete(&mut self, (height, durable): (Height, bool)) -> bool { assert!(self.sync.take().is_some(), "sync not active"); if !durable { return false; } assert!(height > self.durable && height <= self.latest_applied()); self.durable = height; let covered = self .acknowledgements .iter() .take_while(|(height, _)| *height <= self.durable) .count(); for (_, acknowledgement) in self.acknowledgements.drain(..covered) { acknowledgement.acknowledge(); } true } /// Return whether `height` lies within the known durable prefix. fn covers(&self, height: Height) -> bool { self.durable >= height } } /// Await the active sync, remaining pending so callers can select unconditionally when none exists. async fn sync_completion(sync: &mut Option>) -> (Height, bool) { let Some(sync) = sync else { return pending().await; }; sync.await.expect("internal sync handle cannot fail") } /// Start a durability barrier for pending applied state. /// /// Verification work remains driven while the database writer is acquired. Returns false if the /// actor stops before the barrier starts. async fn start_sync( context: &E, durability: &mut Durability, verifications: &mut Verifications, databases: &A::Databases, ) -> bool where E: Rng + Spawner + Metrics + Clock, A: Application, S: Scheme, V: Variant, MarshalMailbox: BlockProvider, { // A requested successor is a no-op when no applied suffix remains uncovered. if !durability.needs_sync() { return true; } // Capture the dirty prefix before waiting for the database writer. Drive verification readers // until the barrier starts, then bind its completion to exactly the prefix it captured. let height = durability.latest_applied(); let barrier = select! { _ = context.stopped() => return false, barrier = verifications.drive(databases.finalize()) => barrier, }; durability.started(height, barrier); true } fn requeue_verifications( mailbox: &(dyn Fn(Message) + Send + Sync), requests: Vec>, ) where E: Rng + Spawner + Metrics + Clock, A: Application, { // Re-enter each live request through FIFO. Work accepted during the mutation // precedes its next attempt, while later arrivals remain behind it. for VerificationRequest { span, context, ancestry, verification, } in requests { if verification.is_cancelled() { continue; } mailbox(Message::Verify { span, context, ancestry, verification, }); } } pub(super) struct Processing where E: Rng + Spawner + Metrics + Clock, A: Application, S: Scheme, V: Variant, { /// Runtime context. pub(super) context: ContextCell, /// Actor ingress. pub(super) mailbox: actor_mailbox::Receiver>, /// Provider cloned into each proposal. pub(super) provider: A::Provider, /// Marshal mailbox used for lazy block lookup. pub(super) marshal: MarshalMailbox, /// The processing state of the actor. pub(super) processor: Processor, /// Verification requests deferred until processing starts. pub(super) deferred_verifications: Vec>, /// Finalized marshal blocks at or below this height were already reflected /// in the selected database anchor and should be acknowledged only. pub(super) skip_finalized_until: Option, } impl Processing where E: Rng + Spawner + Metrics + Clock, A: Application, S: Scheme, V: Variant, MarshalMailbox: BlockProvider, { pub async fn start(mut self) { let mut pending_prune = None; let mut deferred_message = None; let mut verifications = Verifications::new(self.marshal.clone()); for request in std::mem::take(&mut self.deferred_verifications) { verifications.schedule(self.processor.verifier(), request); } // One database sync stays active while later finalized state accumulates behind it. // Completion starts a successor for that suffix unless a pending prune must establish // the next storage-mutation boundary first. let mut durability = Durability::new(self.processor.last_processed().height); select_loop! { self.context, on_start => { // Observe completed durability before taking more work. A queued prune suppresses // an automatic dirty-suffix successor until it has released database readers. if let Some(completion) = sync_completion(&mut durability.sync).now_or_never() && !durability.complete(completion) { return; } if pending_prune.is_none() && !start_sync::( &self.context, &mut durability, &mut verifications, self.processor.databases(), ).await { return; } // Publish completed verdicts before admitting another message. // A later finalization cannot retroactively invalidate them. verifications.complete_ready(); // A message deferred by an active proposal is the FIFO barrier // for subsequent mailbox work, so handle it before later arrivals. let prune_needs_sync = pending_prune.is_some() && durability.needs_sync(); let message = if prune_needs_sync { // The prune must release verification readers before this sync can acquire // its writer. Run that boundary now so durability does not wait for idle. Err(TryRecvError::Empty) } else { match deferred_message.take() { Some(message) => Ok(message), None => self.mailbox.try_recv(), } }; // A prune remains idle work unless it owns the next dirty-suffix mutation boundary. let next = match message { Ok(message) => Either::Left(ready(Some(Step::Message(message)))), Err(TryRecvError::Empty) => match pending_prune.take() { Some(prune) => Either::Left(ready(Some(Step::Prune(prune)))), None => { let mailbox = &mut self.mailbox; let sync = &mut durability.sync; let verifications = &mut verifications; Either::Right(async move { loop { select! { message = mailbox.recv() => { break message.map(Step::Message); }, completion = sync_completion(sync) => { break Some(Step::Sync(completion)); }, _ = verifications.next_completed() => { continue; }, } } }) } }, Err(TryRecvError::Disconnected) => { debug!("mailbox closed, stopping processing"); return; } }; }, on_stopped => { debug!("shutdown signal received, stopping processing"); }, Some(step) = next else { debug!("mailbox closed, stopping processing"); break; } => match step { Step::Message(Message::Propose { span, context, ancestry, upstream, response, }) => { let process = info_span!(parent: &span, "stateful.actor.propose"); let input = Input { upstream, provider: self.provider.clone(), }; let verifier = self.processor.verifier(); let actor_context = self.context.as_present(); let marshal = self.marshal.clone(); let proposal = self .processor .propose( actor_context, marshal.clone(), context, ancestry, input, response, ) .instrument(process); futures::pin_mut!(proposal); let mut receive_messages = true; loop { if receive_messages { select! { _ = &mut proposal => break, message = self.mailbox.recv() => match message { Some(Message::Verify { span, context, ancestry, verification, }) => verifications.schedule( verifier.clone(), VerificationRequest { span, context, ancestry, verification, }, ), Some(message) => { // Only verification may overtake an active proposal. The // first other message becomes a FIFO barrier for later // mailbox work. deferred_message = Some(message); receive_messages = false; } None => receive_messages = false, }, _ = verifications.next_completed() => {}, } } else { select! { _ = &mut proposal => break, _ = verifications.next_completed() => {}, } } } } Step::Message(Message::Verify { span, context, ancestry, verification, }) => { verifications.schedule( self.processor.verifier(), VerificationRequest { span, context, ancestry, verification, }, ); } Step::Message(Message::Finalized { span, block, acknowledgement, retry_mailbox, }) => { if skip_finalized_block(&mut self.skip_finalized_until, block.height()) { // The block is already reflected in the database set by a // completed state sync, so there is nothing to capture or apply. acknowledgement.acknowledge(); } else { let process = info_span!(parent: &span, "stateful.actor.finalized"); let boundary = self.processor.finalization_boundary(block.as_ref()); let (retry, reject) = verifications .quiesce_where(|progress| boundary.disposition(progress)) .await; drop(boundary); async { let should_start_sync = durability.sync.is_none(); let applied = verifications .drive(self.processor.finalize( &self.context, block.as_ref(), should_start_sync, )) .await; let Some(Applied { barrier, prune }) = applied else { // A duplicate report is the startup anchor redelivered by // marshal: genesis on a fresh boot or a newly installed // floor. Its state is durable before the actor starts, so // no barrier is needed. acknowledgement.acknowledge(); return; }; debug!( height = block.height().get(), "applied finalized database batch" ); // Retain marshal acknowledgements until a barrier makes their database // prefix durable. This keeps marshal's processed floor within // recoverable database state while later work proceeds. The // acknowledgement window bounds the queue; a barrier that returns false // leaves the suffix unacknowledged for restart replay. let height = block.height(); durability.applied(height, acknowledgement); if let Some(barrier) = barrier { durability.started(height, barrier); } // Defer pruning to the loop so it can settle durability and quiesce // verification readers at one database mutation boundary. if let Some(prune) = prune { pending_prune = Some((prune, retry_mailbox.clone())); } } .instrument(process) .await; for verification in reject { verification.respond(false); } requeue_verifications(retry_mailbox.as_ref(), retry); } } Step::Message(Message::SubscribeDatabases { response }) => { response.send_lossy(self.processor.databases().clone()); } Step::Prune((prune, retry_mailbox)) => { // Pruning owns a strict database mutation boundary. Observe an existing sync // before quiescing readers, then run storage maintenance with no sync active. while durability.sync.is_some() { select! { completion = sync_completion(&mut durability.sync) => { if !durability.complete(completion) { return; } }, _ = verifications.next_completed() => {}, } } let retry = verifications.quiesce().await; assert!( self.processor.replays_idle(), "verification replay remained active after quiescence" ); // A prune target applied behind an earlier sync may still need durability. if !durability.covers(prune.barrier_height) { assert!( durability.needs_sync(), "uncovered prune target must have unapplied durability", ); if !start_sync::( &self.context, &mut durability, &mut verifications, self.processor.databases(), ).await { return; } let completion = sync_completion(&mut durability.sync).await; if !durability.complete(completion) { return; } assert!(durability.covers(prune.barrier_height)); } prune .run(self.processor.databases(), &self.marshal) .await; requeue_verifications(retry_mailbox.as_ref(), retry); } Step::Sync(completion) => { if !durability.complete(completion) { return; } } }, } } } fn skip_finalized_block(skip_until: &mut Option, height: Height) -> bool { let Some(target) = *skip_until else { return false; }; if height > target { *skip_until = None; return false; } if height == target { *skip_until = None; } true } #[cfg(test)] mod tests { use super::{Message, Processing, VerificationRequest, skip_finalized_block}; use crate::stateful::{ Application, Input, Proposed, PruneConfig, actor::{ core::mailbox::Mailbox, metrics::Metrics as StatefulMetrics, processor::{Processor, Pruning}, }, db::{DatabaseSet, Shared}, tests::{ fixtures, mocks::{ FlushControl, TestApp, TestBlock, TestDatabases, TestDb, TestMerkleized, TestScheme, TestUnmerkleized, anchor, test_databases, }, }, }; use commonware_actor::mailbox as actor_mailbox; use commonware_consensus::{ Application as _, CertifiableBlock as _, Heightable as _, Reporter as _, marshal::{ Update, ancestry::{self, Ancestry}, }, simplex::mocks::scheme as scheme_mocks, types::Height, }; use commonware_macros::select; use commonware_runtime::{ Clock as _, ContextCell, Error as RuntimeError, Handle, Name, Runner as _, Spawner as _, Supervisor as _, deterministic, }; use commonware_utils::{ NZUsize, acknowledgement::{Acknowledgement as _, Exact}, channel::oneshot, sync::Mutex, }; use futures::{StreamExt as _, poll}; use std::{ collections::VecDeque, sync::{ Arc, atomic::{AtomicUsize, Ordering}, }, time::Duration, }; struct ApplicationGate { started: oneshot::Sender<()>, release: oneshot::Receiver<()>, } #[derive(Clone)] struct GatedApp { verify_gates: Arc>>, proposal_gate: Arc>>, verify_valid: bool, observed_contexts: Arc>>, } impl Application for GatedApp { type SigningScheme = TestScheme; type Context = >::Context; type Block = TestBlock; type Databases = TestDatabases; type Captured = (); type Provider = (); type Input = (); fn sync_targets(block: &Self::Block) -> u64 { block.height().get() } async fn genesis(&mut self) -> Self::Block { panic!("gated application genesis is not used") } async fn propose( &mut self, _context: (deterministic::Context, Self::Context), _ancestry: impl Ancestry, _batches: TestUnmerkleized, _input: Input, ) -> Option> { let gate = self.proposal_gate.lock().take(); if let Some(mut gate) = gate { let _ = gate.started.send(()); let _ = (&mut gate.release).await; } None } async fn verify( &mut self, context: (deterministic::Context, Self::Context), ancestry: impl Ancestry, _batches: TestUnmerkleized, ) -> Option { self.observed_contexts.lock().push(context.0.name()); let mut ancestry = Box::pin(ancestry); let _block = ancestry.next().await?; let mut gate = self .verify_gates .lock() .pop_front() .expect("unexpected verification"); let _ = gate.started.send(()); let _ = (&mut gate.release).await; self.verify_valid.then_some(TestMerkleized) } async fn apply( &mut self, _context: (deterministic::Context, Self::Context), _block: &Self::Block, _batches: TestUnmerkleized, ) -> Option { Some(TestMerkleized) } async fn capture( &mut self, _context: (deterministic::Context, Self::Context), _block: &Self::Block, _batches: &TestMerkleized, _readers: >::Readers, ) { } async fn finalized( &mut self, _context: (deterministic::Context, Self::Context), _block: &Self::Block, _captured: Self::Captured, _readers: >::Readers, ) { } } #[derive(Clone)] struct ReadGatedApp { database: Shared, verify_gate_height: Height, verify_gate: Arc>>, } impl Application for ReadGatedApp { type SigningScheme = TestScheme; type Context = >::Context; type Block = TestBlock; type Databases = TestDatabases; type Captured = (); type Provider = (); type Input = (); fn sync_targets(block: &Self::Block) -> u64 { block.height().get() } async fn genesis(&mut self) -> Self::Block { panic!("read-gated application genesis is not used") } async fn propose( &mut self, _context: (deterministic::Context, Self::Context), _ancestry: impl Ancestry, _batches: TestUnmerkleized, _input: Input, ) -> Option> { panic!("read-gated application proposal is not used") } async fn verify( &mut self, _context: (deterministic::Context, Self::Context), ancestry: impl Ancestry, _batches: TestUnmerkleized, ) -> Option { let mut ancestry = Box::pin(ancestry); let block = ancestry.next().await?; if block.height() != self.verify_gate_height { return Some(TestMerkleized); } let database = self.database.read().await; let Some(mut gate) = self.verify_gate.lock().take() else { return std::future::pending().await; }; let _ = gate.started.send(()); let _ = (&mut gate.release).await; drop(database); Some(TestMerkleized) } async fn apply( &mut self, _context: (deterministic::Context, Self::Context), _block: &Self::Block, _batches: TestUnmerkleized, ) -> Option { Some(TestMerkleized) } async fn capture( &mut self, _context: (deterministic::Context, Self::Context), _block: &Self::Block, _batches: &TestMerkleized, _readers: >::Readers, ) { } async fn finalized( &mut self, _context: (deterministic::Context, Self::Context), _block: &Self::Block, _captured: Self::Captured, _readers: >::Readers, ) { } } #[derive(Clone)] struct ReplayGatedApp { gates: Arc>>, verify_gate: Arc>>, finalized_gate: Arc>>, gate_height: Height, unexecutable: Option, apply_calls: Arc, capture_calls: Arc, verify_calls: Arc, applied_finalizations: Arc>>, } impl Application for ReplayGatedApp { type SigningScheme = TestScheme; type Context = >::Context; type Block = TestBlock; type Databases = TestDatabases; type Captured = Height; type Provider = (); type Input = (); fn sync_targets(block: &Self::Block) -> u64 { block.height().get() } async fn genesis(&mut self) -> Self::Block { panic!("replay-gated application genesis is not used") } async fn propose( &mut self, _context: (deterministic::Context, Self::Context), _ancestry: impl Ancestry, _batches: TestUnmerkleized, _input: Input, ) -> Option> { panic!("replay-gated application proposal is not used") } async fn verify( &mut self, _context: (deterministic::Context, Self::Context), _ancestry: impl Ancestry, _batches: TestUnmerkleized, ) -> Option { self.verify_calls.fetch_add(1, Ordering::SeqCst); let gate = self.verify_gate.lock().take(); if let Some(mut gate) = gate { let _ = gate.started.send(()); let _ = (&mut gate.release).await; } Some(TestMerkleized) } async fn apply( &mut self, _context: (deterministic::Context, Self::Context), block: &Self::Block, _batches: TestUnmerkleized, ) -> Option { self.apply_calls.fetch_add(1, Ordering::SeqCst); if self.unexecutable == Some(block.height()) { return None; } let gate = (block.height() == self.gate_height) .then(|| self.gates.lock().pop_front()) .flatten(); if let Some(mut gate) = gate { let _ = gate.started.send(()); let _ = (&mut gate.release).await; } Some(TestMerkleized) } async fn capture( &mut self, _context: (deterministic::Context, Self::Context), block: &Self::Block, _batches: &TestMerkleized, _readers: >::Readers, ) -> Self::Captured { self.capture_calls.fetch_add(1, Ordering::SeqCst); block.height() } async fn finalized( &mut self, _context: (deterministic::Context, Self::Context), _block: &Self::Block, height: Self::Captured, _readers: >::Readers, ) { self.applied_finalizations.lock().push(height); let gate = self.finalized_gate.lock().take(); if let Some(mut gate) = gate { let _ = gate.started.send(()); let _ = (&mut gate.release).await; } } } fn application_gate() -> (ApplicationGate, oneshot::Receiver<()>, oneshot::Sender<()>) { let (started, started_rx) = oneshot::channel(); let (release, release_rx) = oneshot::channel(); ( ApplicationGate { started, release: release_rx, }, started_rx, release, ) } async fn spawn_gated_application( context: &deterministic::Context, prefix: &str, app: GatedApp, ) -> ( Mailbox, Box, Handle<()>, ) { let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"gated-application", 1).schemes[0].clone(); let marshal = fixtures::marshal_fixture( context.child("marshal"), prefix, scheme, None, NZUsize!(1), false, ) .await; let processor = Processor::new( app, test_databases(), anchor(0, 0), StatefulMetrics::new(context), None, ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); (Mailbox::new(sender), marshal.guards, actor) } /// Spawn a [`Processing`] loop over a gated [`TestDb`], returning its /// mailbox, flush controls, a guard keeping the (never-started) marshal /// actor's mailbox open, and the processing actor handle. async fn spawn_processing( context: &deterministic::Context, prefix: &str, prune_config: Option, ) -> ( Mailbox, FlushControl, Box, Handle<()>, ) { spawn_processing_with_gates(context, prefix, prune_config, VecDeque::new()).await } async fn spawn_processing_with_gates( context: &deterministic::Context, prefix: &str, prune_config: Option, verify_gates: VecDeque, ) -> ( Mailbox, FlushControl, Box, Handle<()>, ) { let mut signing = context.child("signing"); let scheme_fixture = scheme_mocks::fixture(&mut signing, b"gated", 1); let marshal = fixtures::marshal_fixture( context.child("marshal_fixture"), prefix, scheme_fixture.schemes[0].clone(), None, NZUsize!(1), false, ) .await; let control = FlushControl::default(); let databases = Shared::new("test", TestDb::gated(control.clone())); let pruning = prune_config .map(|config| Pruning::build(config, marshal.mailbox.max_pending_acks(), 0)); let app = GatedApp { verify_gates: Arc::new(Mutex::new(verify_gates)), proposal_gate: Arc::new(Mutex::new(None)), verify_valid: true, observed_contexts: Arc::default(), }; let processor = Processor::new( app, databases, anchor(0, 0), StatefulMetrics::new(context), pruning, ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); (Mailbox::new(sender), control, marshal.guards, actor) } async fn spawn_read_gated_processing( context: &deterministic::Context, prefix: &str, verify_gate: ApplicationGate, prune_config: Option, ) -> ( Mailbox, FlushControl, Box, Handle<()>, ) { let mut signing = context.child("signing"); let scheme_fixture = scheme_mocks::fixture(&mut signing, b"read-gated", 1); let marshal = fixtures::marshal_fixture( context.child("marshal_fixture"), prefix, scheme_fixture.schemes[0].clone(), None, NZUsize!(1), false, ) .await; let control = FlushControl::default(); let databases = Shared::new("test", TestDb::gated(control.clone())); let app = ReadGatedApp { database: databases.clone(), verify_gate_height: Height::new(3), verify_gate: Arc::new(Mutex::new(Some(verify_gate))), }; let pruning = prune_config .map(|config| Pruning::build(config, marshal.mailbox.max_pending_acks(), 0)); let processor = Processor::new( app, databases, anchor(0, 0), StatefulMetrics::new(context), pruning, ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); (Mailbox::new(sender), control, marshal.guards, actor) } #[test] fn independent_verifications_do_not_block_each_other() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (first_gate, first_started, first_release) = application_gate(); let (second_gate, second_started, second_release) = application_gate(); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([first_gate, second_gate]))), proposal_gate: Arc::new(Mutex::new(None)), verify_valid: true, observed_contexts: Arc::default(), }; let (mut mailbox, _marshal, actor) = spawn_gated_application(&context, "concurrent-verify", app).await; let genesis = TestBlock::new(0, 0); let first_block = TestBlock::child(&genesis, 1); let first_genesis = genesis.clone(); let mut first_mailbox = mailbox.clone(); let first = context.child("first").spawn(move |task_context| { let consensus_context = first_block.context(); async move { first_mailbox .verify( (task_context, consensus_context), ancestry::from_iter([Arc::new(first_block), Arc::new(first_genesis)]), ) .await } }); first_started .await .expect("first verification should start"); let second_block = TestBlock::child(&genesis, 2); let second = context.child("second").spawn(move |task_context| { let consensus_context = second_block.context(); async move { mailbox .verify( (task_context, consensus_context), ancestry::from_iter([Arc::new(second_block), Arc::new(genesis)]), ) .await } }); select! { result = second_started => { result.expect("second verification should start while first remains pending"); }, _ = context.sleep(Duration::from_millis(100)) => { panic!("pending verification blocked unrelated verification"); }, } first_release .send(()) .expect("first verification should remain active"); second_release .send(()) .expect("second verification should remain active"); assert!(first.await.expect("first verification failed")); assert!(second.await.expect("second verification failed")); actor.abort(); }); } #[test] fn verification_preserves_request_attributes() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (gate, started, release) = application_gate(); let observed_contexts = Arc::new(Mutex::new(Vec::new())); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([gate]))), proposal_gate: Arc::new(Mutex::new(None)), verify_valid: true, observed_contexts: observed_contexts.clone(), }; let (mut mailbox, _marshal, actor) = spawn_gated_application(&context, "verify-attributes", app).await; let genesis = TestBlock::new(0, 0); let block = TestBlock::child(&genesis, 1); let block_context = block.context(); let request_context = context .child("request") .with_attribute("round", "request-round") .with_attribute("owner", "request") .with_attribute("shard", 4); let mut verify = Box::pin(mailbox.verify( (request_context, block_context), ancestry::from_iter([Arc::new(block), Arc::new(genesis)]), )); assert!(poll!(&mut verify).is_pending()); started.await.expect("verification should start"); { let observed = observed_contexts.lock(); assert_eq!(observed.len(), 1); assert_eq!( observed[0].attributes, vec![ ("owner".to_string(), "request".to_string()), ("round".to_string(), "request-round".to_string()), ("shard".to_string(), "4".to_string()), ] ); } release.send(()).expect("verification should remain active"); assert!(verify.await); actor.abort(); }); } #[test] fn abandoned_verification_cancels_with_caller() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (gate, started, release) = application_gate(); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([gate]))), proposal_gate: Arc::new(Mutex::new(None)), verify_valid: true, observed_contexts: Arc::default(), }; let (mut mailbox, _marshal, actor) = spawn_gated_application(&context, "caller-cancellation", app).await; let genesis = TestBlock::new(0, 0); let block = TestBlock::child(&genesis, 1); let block_context = block.context(); let mut verify = Box::pin(mailbox.verify( (context.child("caller"), block_context), ancestry::from_iter([Arc::new(block), Arc::new(genesis)]), )); assert!(poll!(&mut verify).is_pending()); started.await.expect("application task should start"); drop(verify); context.sleep(Duration::from_millis(10)).await; assert!( release.send(()).is_err(), "application verification should stop with its caller" ); actor.abort(); }); } #[test] fn abandoned_incomplete_verifications_do_not_block_later_work() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (first_gate, first_started, first_release) = application_gate(); let (second_gate, second_started, second_release) = application_gate(); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([first_gate, second_gate]))), proposal_gate: Arc::new(Mutex::new(None)), verify_valid: true, observed_contexts: Arc::default(), }; let (mut mailbox, _marshal, actor) = spawn_gated_application(&context, "incomplete-verify", app).await; let genesis = TestBlock::new(0, 0); let block1 = TestBlock::child(&genesis, 1); let mut incomplete = Box::pin(mailbox.verify( (context.child("empty"), block1.context()), ancestry::from_iter([]), )); assert!(poll!(&mut incomplete).is_pending()); context.sleep(Duration::from_millis(10)).await; drop(incomplete); let mut first = Box::pin(mailbox.verify( (context.child("first"), block1.context()), ancestry::from_iter([Arc::new(block1.clone()), Arc::new(genesis)]), )); assert!(poll!(&mut first).is_pending()); first_started .await .expect("later verification should start"); first_release .send(()) .expect("later verification should remain active"); assert!(first.await); let block2 = TestBlock::child(&block1, 2); let mut incomplete = Box::pin(mailbox.verify( (context.child("missing_parent"), block2.context()), ancestry::from_iter([Arc::new(block2.clone())]), )); assert!(poll!(&mut incomplete).is_pending()); context.sleep(Duration::from_millis(10)).await; drop(incomplete); let mut second = Box::pin(mailbox.verify( (context.child("second"), block2.context()), ancestry::from_iter([Arc::new(block2), Arc::new(block1)]), )); assert!(poll!(&mut second).is_pending()); second_started .await .expect("later verification should start"); second_release .send(()) .expect("later verification should remain active"); assert!(second.await); actor.abort(); }); } #[test] fn application_rejection_returns_false() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (gate, started, release) = application_gate(); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([gate]))), proposal_gate: Arc::new(Mutex::new(None)), verify_valid: false, observed_contexts: Arc::default(), }; let (mut mailbox, _marshal, actor) = spawn_gated_application(&context, "rejected-verify", app).await; let genesis = TestBlock::new(0, 0); let block = TestBlock::child(&genesis, 1); let mut verify = Box::pin(mailbox.verify( (context.child("verify"), block.context()), ancestry::from_iter([Arc::new(block), Arc::new(genesis)]), )); assert!(poll!(&mut verify).is_pending()); started.await.expect("verification should start"); release.send(()).expect("verification should remain active"); assert!(!verify.await); actor.abort(); }); } #[test] fn replayed_parent_is_not_a_verdict() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let genesis = TestBlock::new(0, 0); let parent = TestBlock::child(&genesis, 1); let child = TestBlock::child(&parent, 2); let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"replayed-parent", 1).schemes[0].clone(); let marshal = fixtures::marshal_fixture_with_finalized_block( context.child("marshal"), "replayed-parent", scheme, &genesis, NZUsize!(1), true, ) .await; let apply_calls = Arc::new(AtomicUsize::new(0)); let verify_calls = Arc::new(AtomicUsize::new(0)); let app = ReplayGatedApp { gates: Arc::default(), verify_gate: Arc::default(), finalized_gate: Arc::default(), gate_height: parent.height(), unexecutable: None, apply_calls: apply_calls.clone(), capture_calls: Arc::new(AtomicUsize::new(0)), verify_calls: verify_calls.clone(), applied_finalizations: Arc::default(), }; let processor = Processor::new( app, test_databases(), anchor(0, 0), StatefulMetrics::new(&context), None, ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let mut mailbox = Mailbox::new(sender); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); // Verifying the child replays its missing parent through apply. assert!( mailbox .verify( (context.child("verify_child"), child.context()), ancestry::from_iter([Arc::new(child), Arc::new(parent.clone())]), ) .await ); assert_eq!(apply_calls.load(Ordering::SeqCst), 1); assert_eq!(verify_calls.load(Ordering::SeqCst), 1); // Replayed state is reusable parent state, not a verdict: verifying the // parent asks the application once and then settles from the cache. for label in ["verify_parent", "verify_parent_again"] { assert!( mailbox .verify( (context.child(label), parent.context()), ancestry::from_iter([ Arc::new(parent.clone()), Arc::new(genesis.clone()) ]), ) .await ); } assert_eq!(apply_calls.load(Ordering::SeqCst), 1); assert_eq!(verify_calls.load(Ordering::SeqCst), 2); actor.abort(); drop(marshal.guards); }); } #[test] fn unexecutable_parent_rejects_child() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let genesis = TestBlock::new(0, 0); let parent = TestBlock::child(&genesis, 1); let child = TestBlock::child(&parent, 2); let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"unexecutable-parent", 1).schemes[0].clone(); let marshal = fixtures::marshal_fixture_with_finalized_block( context.child("marshal"), "unexecutable-parent", scheme, &genesis, NZUsize!(1), true, ) .await; let apply_calls = Arc::new(AtomicUsize::new(0)); let verify_calls = Arc::new(AtomicUsize::new(0)); let app = ReplayGatedApp { gates: Arc::default(), verify_gate: Arc::default(), finalized_gate: Arc::default(), gate_height: parent.height(), unexecutable: Some(parent.height()), apply_calls: apply_calls.clone(), capture_calls: Arc::new(AtomicUsize::new(0)), verify_calls: verify_calls.clone(), applied_finalizations: Arc::default(), }; let processor = Processor::new( app, test_databases(), anchor(0, 0), StatefulMetrics::new(&context), None, ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let mut mailbox = Mailbox::new(sender); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); // A parent that cannot be executed invalidates the child's ancestry // before the application is asked to verify the child. assert!( !mailbox .verify( (context.child("verify_child"), child.context()), ancestry::from_iter([Arc::new(child), Arc::new(parent)]), ) .await ); assert_eq!(apply_calls.load(Ordering::SeqCst), 1); assert_eq!(verify_calls.load(Ordering::SeqCst), 0); actor.abort(); drop(marshal.guards); }); } #[test] fn conflicting_processed_block_is_rejected() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (mut mailbox, control, _marshal, actor) = spawn_processing(&context, "conflicting-processed", None).await; let genesis = TestBlock::new(0, 0); let canonical = TestBlock::child(&genesis, 1); let conflicting = TestBlock::child(&genesis, 2); let (acknowledgement, waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(canonical), acknowledgement)); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } control .flushes .lock() .remove(0) .send(Ok(())) .expect("finalized block sync should remain pending"); waiter .await .expect("finalized block should be acknowledged"); assert!( !mailbox .verify( (context.child("verify"), conflicting.context()), ancestry::from_iter([Arc::new(conflicting), Arc::new(genesis)]), ) .await, "conflicting block at the processed height must be rejected", ); actor.abort(); }); } #[test] fn pending_proposal_does_not_block_active_verification() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (verify_gate, verify_started, verify_release) = application_gate(); let (proposal_gate, proposal_started, proposal_release) = application_gate(); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([verify_gate]))), proposal_gate: Arc::new(Mutex::new(Some(proposal_gate))), verify_valid: true, observed_contexts: Arc::default(), }; let (mut mailbox, _marshal, actor) = spawn_gated_application(&context, "propose-verify", app).await; let genesis = TestBlock::new(0, 0); let block = TestBlock::child(&genesis, 1); let mut verifier = mailbox.clone(); let verify_genesis = genesis.clone(); let consensus_context = block.context(); let mut verify = Box::pin(verifier.verify( (context.child("verify"), consensus_context), ancestry::from_iter([Arc::new(block), Arc::new(verify_genesis)]), )); let subscriber = mailbox.clone(); assert!(poll!(&mut verify).is_pending()); verify_started.await.expect("verification should start"); let proposal_context = TestBlock::child(&genesis, 2).context(); let mut proposal = Box::pin(mailbox.propose( (context.child("propose"), proposal_context), ancestry::from_iter([Arc::new(genesis)]), (), )); assert!(poll!(&mut proposal).is_pending()); proposal_started.await.expect("proposal should start"); let mut databases = Box::pin(subscriber.subscribe_databases()); assert!(poll!(&mut databases).is_pending()); context.sleep(Duration::from_millis(10)).await; verify_release .send(()) .expect("verification should remain active"); select! { result = &mut verify => { assert!(result); }, _ = context.sleep(Duration::from_millis(100)) => { panic!("pending proposal blocked active verification"); }, } assert!(poll!(&mut databases).is_pending()); proposal_release .send(()) .expect("proposal should remain active"); assert!(proposal.await.is_none()); drop(databases.await); actor.abort(); }); } #[test] fn pending_proposal_does_not_block_new_verification() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (verify_gate, verify_started, verify_release) = application_gate(); let (proposal_gate, proposal_started, proposal_release) = application_gate(); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([verify_gate]))), proposal_gate: Arc::new(Mutex::new(Some(proposal_gate))), verify_valid: true, observed_contexts: Arc::default(), }; let (mut mailbox, _marshal, actor) = spawn_gated_application(&context, "propose-new-verify", app).await; let genesis = TestBlock::new(0, 0); let proposal_context = TestBlock::child(&genesis, 1).context(); let mut proposer = mailbox.clone(); let mut proposal = Box::pin(proposer.propose( (context.child("propose"), proposal_context), ancestry::from_iter([Arc::new(genesis.clone())]), (), )); assert!(poll!(&mut proposal).is_pending()); proposal_started.await.expect("proposal should start"); let block = TestBlock::child(&genesis, 2); let consensus_context = block.context(); let mut verify = Box::pin(mailbox.verify( (context.child("verify"), consensus_context), ancestry::from_iter([Arc::new(block), Arc::new(genesis)]), )); assert!(poll!(&mut verify).is_pending()); select! { result = verify_started => { result.expect("verification should start while proposal remains pending"); }, _ = context.sleep(Duration::from_millis(100)) => { panic!("pending proposal blocked new verification"); }, } verify_release .send(()) .expect("verification should remain active"); assert!(verify.await); proposal_release .send(()) .expect("proposal should remain active"); assert!(proposal.await.is_none()); actor.abort(); }); } #[test] fn deferred_finalization_does_not_block_completed_verification() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (parent_gate, parent_started, parent_release) = application_gate(); let (child_gate, child_started, child_release) = application_gate(); let (proposal_gate, proposal_started, proposal_release) = application_gate(); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([parent_gate, child_gate]))), proposal_gate: Arc::new(Mutex::new(Some(proposal_gate))), verify_valid: true, observed_contexts: Arc::default(), }; let (mut mailbox, _marshal, actor) = spawn_gated_application(&context, "proposal-finalization", app).await; let genesis = TestBlock::new(0, 0); let winner = TestBlock::child(&genesis, 1); let losing_parent = TestBlock::child(&genesis, 2); let losing_child = TestBlock::child(&losing_parent, 3); let mut parent_verifier = mailbox.clone(); let mut verify_parent = Box::pin(parent_verifier.verify( (context.child("verify_parent"), losing_parent.context()), ancestry::from_iter([Arc::new(losing_parent.clone()), Arc::new(genesis.clone())]), )); assert!(poll!(&mut verify_parent).is_pending()); parent_started .await .expect("losing parent verification should start"); parent_release .send(()) .expect("losing parent verification should remain active"); assert!(verify_parent.await); let mut proposer = mailbox.clone(); let mut proposal = Box::pin(proposer.propose( ( context.child("propose"), TestBlock::child(&genesis, 4).context(), ), ancestry::from_iter([Arc::new(genesis)]), (), )); assert!(poll!(&mut proposal).is_pending()); proposal_started.await.expect("proposal should start"); let mut child_verifier = mailbox.clone(); let mut verify_child = Box::pin(child_verifier.verify( (context.child("verify_child"), losing_child.context()), ancestry::from_iter([Arc::new(losing_child), Arc::new(losing_parent)]), )); assert!(poll!(&mut verify_child).is_pending()); child_started .await .expect("losing child verification should start"); let (acknowledgement, mut waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(winner), acknowledgement)); context.sleep(Duration::from_millis(10)).await; assert!(poll!(&mut waiter).is_pending()); child_release .send(()) .expect("losing child verification should remain active"); select! { valid = &mut verify_child => { assert!(valid, "completed branch-relative verification must remain valid"); }, _ = context.sleep(Duration::from_millis(100)) => { panic!("deferred finalization blocked completed verification"); }, } proposal_release .send(()) .expect("proposal should remain active"); assert!(proposal.await.is_none()); waiter .await .expect("conflicting finalized block should be acknowledged"); actor.abort(); }); } #[test] fn finalization_keeps_compatible_verification_active() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (parent_gate, parent_started, parent_release) = application_gate(); let (child_gate, child_started, child_release) = application_gate(); let (retry_gate, _retry_started, _retry_release) = application_gate(); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([ parent_gate, child_gate, retry_gate, ]))), proposal_gate: Arc::new(Mutex::new(None)), verify_valid: true, observed_contexts: Arc::default(), }; let (mut mailbox, _marshal, actor) = spawn_gated_application(&context, "finalize-compatible", app).await; let genesis = TestBlock::new(0, 0); let parent = TestBlock::child(&genesis, 1); let child = TestBlock::child(&parent, 2); let mut parent_verifier = mailbox.clone(); let parent_genesis = genesis.clone(); let parent_context = parent.context(); let mut verify_parent = Box::pin(parent_verifier.verify( (context.child("verify_parent"), parent_context), ancestry::from_iter([Arc::new(parent.clone()), Arc::new(parent_genesis)]), )); assert!(poll!(&mut verify_parent).is_pending()); parent_started .await .expect("parent verification should start"); parent_release .send(()) .expect("parent verification should remain active"); assert!(verify_parent.await); let child_context = child.context(); let mut child_verifier = mailbox.clone(); let mut verify_child = Box::pin(child_verifier.verify( (context.child("verify_child"), child_context), ancestry::from_iter([Arc::new(child), Arc::new(parent.clone())]), )); assert!(poll!(&mut verify_child).is_pending()); child_started .await .expect("child verification should start"); let (acknowledgement, waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(parent), acknowledgement)); waiter .await .expect("finalized parent should be acknowledged"); child_release .send(()) .expect("compatible verification should remain active"); assert!(verify_child.await); actor.abort(); }); } #[test] fn finalization_invalidates_incompatible_verification() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (fork_gate, fork_started, fork_release) = application_gate(); let (child_gate, child_started, mut child_release) = application_gate(); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([fork_gate, child_gate]))), proposal_gate: Arc::new(Mutex::new(None)), verify_valid: true, observed_contexts: Arc::default(), }; let (mut mailbox, _marshal, actor) = spawn_gated_application(&context, "finalize-incompatible", app).await; let genesis = TestBlock::new(0, 0); let winner = TestBlock::child(&genesis, 1); let losing_parent = TestBlock::child(&genesis, 2); let losing_child = TestBlock::child(&losing_parent, 3); let mut fork_verifier = mailbox.clone(); let fork_genesis = genesis.clone(); let fork_context = losing_parent.context(); let mut verify_fork = Box::pin(fork_verifier.verify( (context.child("verify_fork"), fork_context), ancestry::from_iter([Arc::new(losing_parent.clone()), Arc::new(fork_genesis)]), )); assert!(poll!(&mut verify_fork).is_pending()); fork_started.await.expect("fork verification should start"); fork_release .send(()) .expect("fork verification should remain active"); assert!(verify_fork.await); let child_context = losing_child.context(); let mut child_verifier = mailbox.clone(); let mut verify_child = Box::pin(child_verifier.verify( (context.child("verify_child"), child_context), ancestry::from_iter([Arc::new(losing_child), Arc::new(losing_parent)]), )); assert!(poll!(&mut verify_child).is_pending()); child_started .await .expect("losing child verification should start"); let (acknowledgement, waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(winner), acknowledgement)); select! { _ = child_release.closed() => {}, _ = context.sleep(Duration::from_millis(100)) => { panic!("incompatible verification was not cancelled"); }, } let mut waiter = Box::pin(waiter); select! { result = &mut waiter => { result.expect("winning block should be acknowledged"); }, _ = context.sleep(Duration::from_millis(100)) => { panic!("winning block was not acknowledged"); }, } select! { valid = &mut verify_child => { assert!(!valid, "verification on a finalized-away fork must fail"); }, _ = context.sleep(Duration::from_millis(100)) => { panic!("incompatible verification retry did not resolve"); }, } actor.abort(); }); } #[test] fn finalization_rejects_deep_incompatible_verification() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (parent_gate, parent_started, parent_release) = application_gate(); let (child_gate, child_started, child_release) = application_gate(); let (grandchild_gate, grandchild_started, mut grandchild_release) = application_gate(); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([ parent_gate, child_gate, grandchild_gate, ]))), proposal_gate: Arc::new(Mutex::new(None)), verify_valid: true, observed_contexts: Arc::default(), }; let (mut mailbox, _marshal, actor) = spawn_gated_application(&context, "finalize-deep-incompatible", app).await; let genesis = TestBlock::new(0, 0); let winner = TestBlock::child(&genesis, 1); let losing_parent = TestBlock::child(&genesis, 2); let losing_child = TestBlock::child(&losing_parent, 3); let losing_grandchild = TestBlock::child(&losing_child, 4); let mut parent_verifier = mailbox.clone(); let mut verify_parent = Box::pin(parent_verifier.verify( (context.child("verify_parent"), losing_parent.context()), ancestry::from_iter([Arc::new(losing_parent.clone()), Arc::new(genesis.clone())]), )); assert!(poll!(&mut verify_parent).is_pending()); parent_started .await .expect("losing parent verification should start"); parent_release .send(()) .expect("losing parent verification should remain active"); assert!(verify_parent.await); let mut child_verifier = mailbox.clone(); let mut verify_child = Box::pin(child_verifier.verify( (context.child("verify_child"), losing_child.context()), ancestry::from_iter([Arc::new(losing_child.clone()), Arc::new(losing_parent)]), )); assert!(poll!(&mut verify_child).is_pending()); child_started .await .expect("losing child verification should start"); child_release .send(()) .expect("losing child verification should remain active"); assert!(verify_child.await); let mut grandchild_verifier = mailbox.clone(); let mut verify_grandchild = Box::pin(grandchild_verifier.verify( ( context.child("verify_grandchild"), losing_grandchild.context(), ), ancestry::from_iter([Arc::new(losing_grandchild), Arc::new(losing_child)]), )); assert!(poll!(&mut verify_grandchild).is_pending()); grandchild_started .await .expect("losing grandchild verification should start"); let (acknowledgement, waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(winner), acknowledgement)); grandchild_release.closed().await; waiter.await.expect("winning block should be acknowledged"); let result = select! { valid = &mut verify_grandchild => Some(valid), _ = context.sleep(Duration::from_millis(100)) => None, }; actor.abort(); assert_eq!( result, Some(false), "verification on a pruned deep fork must resolve false", ); }); } #[test] fn skipped_finalization_keeps_retained_verification_progressing() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let genesis = TestBlock::new(0, 0); let finalized = TestBlock::child(&genesis, 1); let child = TestBlock::child(&finalized, 2); let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"skip-finalized", 1).schemes[0].clone(); let marshal = fixtures::marshal_fixture( context.child("marshal"), "skip-finalized", scheme, None, NZUsize!(1), false, ) .await; let (verify_gate, verify_started, verify_release) = application_gate(); let apply_calls = Arc::new(AtomicUsize::new(0)); let capture_calls = Arc::new(AtomicUsize::new(0)); let applied_finalizations: Arc>> = Arc::default(); let app = ReplayGatedApp { gates: Arc::new(Mutex::new(VecDeque::new())), verify_gate: Arc::new(Mutex::new(Some(verify_gate))), finalized_gate: Arc::default(), gate_height: finalized.height(), unexecutable: None, apply_calls: apply_calls.clone(), capture_calls: capture_calls.clone(), verify_calls: Arc::new(AtomicUsize::new(0)), applied_finalizations: applied_finalizations.clone(), }; let processor = Processor::new( app, test_databases(), anchor(1, 1), StatefulMetrics::new(&context), None, ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let mut mailbox = Mailbox::new(sender); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: Some(finalized.height()), }; let actor = context.child("loop").spawn(move |_| processing.start()); let mut verifier = mailbox.clone(); let mut verify_child = Box::pin(verifier.verify( (context.child("verify_child"), child.context()), ancestry::from_iter([Arc::new(child), Arc::new(finalized.clone())]), )); assert!(poll!(&mut verify_child).is_pending()); verify_started .await .expect("child verification should start"); let (acknowledgement, waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(finalized), acknowledgement)); waiter .await .expect("skipped finalized block should be acknowledged"); assert!( poll!(&mut verify_child).is_pending(), "skipped finalization must not resolve a retained verification", ); verify_release .send(()) .expect("child verification should remain active"); let valid = verify_child.await; actor.abort(); drop(marshal.guards); assert!(valid); assert_eq!(apply_calls.load(Ordering::SeqCst), 0); assert_eq!(capture_calls.load(Ordering::SeqCst), 0); assert!(applied_finalizations.lock().is_empty()); }); } #[test] fn fresh_boot_genesis_skips_hooks() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let genesis = TestBlock::new(0, 0); let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"fresh-boot-genesis", 1).schemes[0].clone(); let marshal = fixtures::marshal_fixture( context.child("marshal"), "fresh-boot-genesis", scheme, None, NZUsize!(1), false, ) .await; let apply_calls = Arc::new(AtomicUsize::new(0)); let capture_calls = Arc::new(AtomicUsize::new(0)); let applied_finalizations: Arc>> = Arc::default(); let app = ReplayGatedApp { gates: Arc::default(), verify_gate: Arc::default(), finalized_gate: Arc::default(), gate_height: genesis.height(), unexecutable: None, apply_calls: apply_calls.clone(), capture_calls: capture_calls.clone(), verify_calls: Arc::new(AtomicUsize::new(0)), applied_finalizations: applied_finalizations.clone(), }; let processor = Processor::new( app, test_databases(), anchor(0, 0), StatefulMetrics::new(&context), None, ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(1)); let mut mailbox = Mailbox::new(sender); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); let (acknowledgement, waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(genesis), acknowledgement)); waiter.await.expect("genesis should be acknowledged"); assert_eq!(apply_calls.load(Ordering::SeqCst), 0); assert_eq!(capture_calls.load(Ordering::SeqCst), 0); assert!(applied_finalizations.lock().is_empty()); actor.abort(); drop(marshal.guards); }); } #[test] fn deferred_verification_resumes_after_sync() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let (gate, started, release) = application_gate(); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([gate]))), proposal_gate: Arc::new(Mutex::new(None)), verify_valid: true, observed_contexts: Arc::default(), }; let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"deferred-verify", 1).schemes[0].clone(); let marshal = fixtures::marshal_fixture( context.child("marshal"), "deferred-verify", scheme, None, NZUsize!(1), false, ) .await; let processor = Processor::new( app, test_databases(), anchor(0, 0), StatefulMetrics::new(&context), None, ); // Defer a verification as the syncing actor does before its // database set is ready. let (sender, mut receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let mut mailbox = Mailbox::new(sender); let genesis = TestBlock::new(0, 0); let block = TestBlock::child(&genesis, 1); let deferred = context.child("deferred").spawn(move |task_context| { let consensus_context = block.context(); async move { mailbox .verify( (task_context, consensus_context), ancestry::from_iter([Arc::new(block), Arc::new(genesis)]), ) .await } }); let request = match receiver.recv().await { Some(Message::Verify { span, context: request_context, ancestry, verification, }) => VerificationRequest { span, context: request_context, ancestry, verification, }, _ => panic!("deferred verification request must arrive"), }; // Resume the deferred verification after state sync. let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: vec![request], skip_finalized_until: Some(Height::new(0)), }; let actor = context.child("loop").spawn(move |_| processing.start()); started.await.expect("deferred verification should resume"); release .send(()) .expect("deferred verification should remain active"); assert!( deferred .await .expect("deferred verification should resolve") ); actor.abort(); drop(marshal.guards); }); } #[test] fn finalization_reuses_active_replay() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let genesis = TestBlock::new(0, 0); let parent = TestBlock::child(&genesis, 1); let first_child = TestBlock::child(&parent, 2); let second_child = TestBlock::child(&parent, 3); let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"finalize-replay", 1).schemes[0].clone(); let marshal = fixtures::marshal_fixture_with_finalized_block( context.child("marshal"), "finalize-replay", scheme, &genesis, NZUsize!(1), true, ) .await; let (gate, apply_started, apply_release) = application_gate(); let (verify_gate, verify_started, verify_release) = application_gate(); let (finalized_gate, finalized_started, finalized_release) = application_gate(); let apply_calls = Arc::new(AtomicUsize::new(0)); let verify_calls = Arc::new(AtomicUsize::new(0)); let app = ReplayGatedApp { gates: Arc::new(Mutex::new(VecDeque::from([gate]))), verify_gate: Arc::new(Mutex::new(Some(verify_gate))), finalized_gate: Arc::new(Mutex::new(Some(finalized_gate))), gate_height: parent.height(), unexecutable: None, apply_calls: apply_calls.clone(), capture_calls: Arc::new(AtomicUsize::new(0)), verify_calls: verify_calls.clone(), applied_finalizations: Arc::default(), }; let processor = Processor::new( app, test_databases(), anchor(0, 0), StatefulMetrics::new(&context), None, ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let mut mailbox = Mailbox::new(sender); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); let consensus_context = first_child.context(); let mut first_verifier = mailbox.clone(); let mut first = Box::pin(first_verifier.verify( (context.child("first_verify"), consensus_context.clone()), ancestry::from_iter([Arc::new(first_child), Arc::new(parent.clone())]), )); let mut second_verifier = mailbox.clone(); let mut second = Box::pin(second_verifier.verify( (context.child("second_verify"), consensus_context), ancestry::from_iter([Arc::new(second_child), Arc::new(parent.clone())]), )); assert!(poll!(&mut first).is_pending()); assert!(poll!(&mut second).is_pending()); apply_started.await.expect("replay should start"); context.sleep(Duration::from_millis(10)).await; let (acknowledgement, waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(parent), acknowledgement)); context.sleep(Duration::from_millis(10)).await; apply_release .send(()) .expect("finalization should keep the replay active"); verify_started .await .expect("compatible verification should start"); finalized_started .await .expect("finalization hook should start"); context.sleep(Duration::from_millis(10)).await; assert_eq!(apply_calls.load(Ordering::SeqCst), 1); verify_release .send(()) .expect("compatible verification should remain active"); assert!(first.await); finalized_release .send(()) .expect("finalization hook should remain active"); waiter .await .expect("finalized parent should be acknowledged"); assert!(second.await); assert_eq!(apply_calls.load(Ordering::SeqCst), 1); assert_eq!(verify_calls.load(Ordering::SeqCst), 2); actor.abort(); drop(marshal.guards); }); } #[test] fn finalization_does_not_bypass_active_winner_replay() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let genesis = TestBlock::new(0, 0); let finalized = TestBlock::child(&genesis, 1); let child = TestBlock::child(&finalized, 2); let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"finalize-pending-replay", 1).schemes [0] .clone(); let marshal = fixtures::marshal_fixture_with_finalized_block( context.child("marshal"), "finalize-pending-replay", scheme, &genesis, NZUsize!(1), true, ) .await; let (replay_gate, replay_started, replay_release) = application_gate(); let apply_calls = Arc::new(AtomicUsize::new(0)); let verify_calls = Arc::new(AtomicUsize::new(0)); let app = ReplayGatedApp { gates: Arc::new(Mutex::new(VecDeque::from([replay_gate]))), verify_gate: Arc::new(Mutex::new(None)), finalized_gate: Arc::new(Mutex::new(None)), gate_height: finalized.height(), unexecutable: None, apply_calls: apply_calls.clone(), capture_calls: Arc::new(AtomicUsize::new(0)), verify_calls: verify_calls.clone(), applied_finalizations: Arc::default(), }; let processor = Processor::new( app, test_databases(), anchor(0, 0), StatefulMetrics::new(&context), None, ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let mut mailbox = Mailbox::new(sender); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); let mut child_verifier = mailbox.clone(); let mut verify_child = Box::pin(child_verifier.verify( (context.child("verify_child"), child.context()), ancestry::from_iter([Arc::new(child), Arc::new(finalized.clone())]), )); assert!(poll!(&mut verify_child).is_pending()); replay_started.await.expect("winner replay should start"); let mut winner_verifier = mailbox.clone(); assert!( winner_verifier .verify( (context.child("verify_winner"), finalized.context()), ancestry::from_iter([Arc::new(finalized.clone()), Arc::new(genesis),]), ) .await, "independent winner verification should cache its batch", ); let (acknowledgement, waiter) = Exact::handle(); let mut waiter = Box::pin(waiter); let _ = mailbox.report(Update::Block(Arc::new(finalized), acknowledgement)); assert!( poll!(&mut waiter).is_pending(), "finalization must wait for the existing winner computation", ); replay_release .send(()) .expect("winner replay should remain active"); waiter .await .expect("cached winner should finalize after replay completes"); let valid = verify_child.await; actor.abort(); drop(marshal.guards); assert!(valid, "late winner replay must not invalidate its child"); assert_eq!(apply_calls.load(Ordering::SeqCst), 1); assert_eq!(verify_calls.load(Ordering::SeqCst), 2); }); } #[test] fn consecutive_finalizations_preserve_descendant_replay() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let genesis = TestBlock::new(0, 0); let first = TestBlock::child(&genesis, 1); let second = TestBlock::child(&first, 2); let child = TestBlock::child(&second, 3); let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"finalize-previous-replay", 1) .schemes[0] .clone(); let marshal = fixtures::marshal_fixture_with_finalized_block( context.child("marshal"), "finalize-previous-replay", scheme, &first, NZUsize!(1), true, ) .await; let (replay_gate, replay_started, replay_release) = application_gate(); let verify_gate = Arc::new(Mutex::new(None)); let apply_calls = Arc::new(AtomicUsize::new(0)); let verify_calls = Arc::new(AtomicUsize::new(0)); let app = ReplayGatedApp { gates: Arc::new(Mutex::new(VecDeque::from([replay_gate]))), verify_gate: verify_gate.clone(), finalized_gate: Arc::new(Mutex::new(None)), gate_height: first.height(), unexecutable: None, apply_calls: apply_calls.clone(), capture_calls: Arc::new(AtomicUsize::new(0)), verify_calls: verify_calls.clone(), applied_finalizations: Arc::default(), }; let processor = Processor::new( app, test_databases(), anchor(0, 0), StatefulMetrics::new(&context), None, ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let mut mailbox = Mailbox::new(sender); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); let mut child_verifier = mailbox.clone(); let mut verify_child = Box::pin(child_verifier.verify( (context.child("verify_child"), child.context()), ancestry::from_iter([Arc::new(child), Arc::new(second.clone())]), )); assert!(poll!(&mut verify_child).is_pending()); replay_started .await .expect("first-block replay should start"); let mut first_verifier = mailbox.clone(); assert!( first_verifier .verify( (context.child("verify_first"), first.context()), ancestry::from_iter([Arc::new(first.clone()), Arc::new(genesis)]), ) .await, "independent verification should cache the first finalized block", ); let (gate, verify_started, verify_release) = application_gate(); assert!(verify_gate.lock().replace(gate).is_none()); let (acknowledgement, first_waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(first), acknowledgement)); replay_release .send(()) .expect("first-block replay should remain active"); first_waiter .await .expect("first finalized block should be acknowledged"); verify_started .await .expect("descendant verification should start"); assert!(poll!(&mut verify_child).is_pending()); let (acknowledgement, second_waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(second), acknowledgement)); second_waiter .await .expect("second finalized block should be acknowledged"); verify_release .send(()) .expect("descendant verification should remain active"); let valid = verify_child.await; actor.abort(); drop(marshal.guards); assert!( valid, "descendant replay must remain valid across consecutive finalizations", ); assert_eq!(apply_calls.load(Ordering::SeqCst), 2); assert_eq!(verify_calls.load(Ordering::SeqCst), 2); }); } #[test] fn retained_verification_can_finish_before_queued_finalization() { deterministic::Runner::timed(Duration::from_secs(5)).start(|context| async move { let genesis = TestBlock::new(0, 0); let first = TestBlock::child(&genesis, 1); let losing = TestBlock::child(&first, 2); let winner = TestBlock::child(&first, 3); let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"finalize-retry-order", 1).schemes[0].clone(); let marshal = fixtures::marshal_fixture_with_finalized_block( context.child("marshal"), "finalize-retry-order", scheme, &genesis, NZUsize!(2), true, ) .await; let (replay_gate, replay_started, replay_release) = application_gate(); let (verify_gate, verify_started, verify_release) = application_gate(); let (finalized_gate, finalized_started, finalized_release) = application_gate(); let app = ReplayGatedApp { gates: Arc::new(Mutex::new(VecDeque::from([replay_gate]))), verify_gate: Arc::new(Mutex::new(Some(verify_gate))), finalized_gate: Arc::new(Mutex::new(Some(finalized_gate))), gate_height: first.height(), unexecutable: None, apply_calls: Arc::new(AtomicUsize::new(0)), capture_calls: Arc::new(AtomicUsize::new(0)), verify_calls: Arc::new(AtomicUsize::new(0)), applied_finalizations: Arc::default(), }; let processor = Processor::new( app, test_databases(), anchor(0, 0), StatefulMetrics::new(&context), None, ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let mut mailbox = Mailbox::new(sender); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); let mut first_verifier = mailbox.clone(); let mut first_attempt = Box::pin(first_verifier.verify( (context.child("first_attempt"), losing.context()), ancestry::from_iter([Arc::new(losing.clone()), Arc::new(first.clone())]), )); assert!(poll!(&mut first_attempt).is_pending()); replay_started.await.expect("winner replay should start"); let mut retried_verifier = mailbox.clone(); let mut retried = Box::pin(retried_verifier.verify( (context.child("retried"), losing.context()), ancestry::from_iter([Arc::new(losing), Arc::new(first.clone())]), )); assert!(poll!(&mut retried).is_pending()); context.sleep(Duration::from_millis(10)).await; let (acknowledgement, first_waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(first), acknowledgement)); let (acknowledgement, winner_waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(winner), acknowledgement)); replay_release .send(()) .expect("finalization should retain the replay owner"); verify_release .send(()) .expect("retained verification should remain active"); verify_started .await .expect("retained verification should start"); finalized_started .await .expect("first finalization hook should start"); let valid = select! { valid = &mut first_attempt => { valid }, _ = context.sleep(Duration::from_millis(100)) => { panic!("queued finalization blocked retained verification"); }, }; assert!( valid, "retained branch-relative verification must remain valid" ); finalized_release .send(()) .expect("first finalization hook should remain active"); assert!( retried.await, "completed branch-relative verdict must remain valid" ); first_waiter .await .expect("first block should be acknowledged"); winner_waiter.await.expect("winner should be acknowledged"); actor.abort(); drop(marshal.guards); }); } #[test] fn pruning_quiesces_replay_before_database_prune() { deterministic::Runner::timed(Duration::from_secs(10)).start(|context| async move { let genesis = TestBlock::new(0, 0); let block1 = TestBlock::child(&genesis, 1); let block2 = TestBlock::child(&block1, 2); let parent = TestBlock::child(&block2, 3); let child = TestBlock::child(&parent, 4); let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"prune-replay", 1).schemes[0].clone(); let marshal = fixtures::marshal_fixture_with_finalized_block( context.child("marshal"), "prune-replay", scheme, &block2, NZUsize!(1), true, ) .await; let (first_gate, first_started, mut first_release) = application_gate(); let (second_gate, second_started, second_release) = application_gate(); let apply_calls = Arc::new(AtomicUsize::new(0)); let verify_calls = Arc::new(AtomicUsize::new(0)); let app = ReplayGatedApp { gates: Arc::new(Mutex::new(VecDeque::from([first_gate, second_gate]))), verify_gate: Arc::new(Mutex::new(None)), finalized_gate: Arc::new(Mutex::new(None)), gate_height: parent.height(), unexecutable: None, apply_calls: apply_calls.clone(), capture_calls: Arc::new(AtomicUsize::new(0)), verify_calls: verify_calls.clone(), applied_finalizations: Arc::default(), }; let control = FlushControl::default(); let (prune_started, prune_release) = control.gate_prune(); let databases = Shared::new("prune-replay", TestDb::gated(control.clone())); let pruning = Pruning::build( PruneConfig { maintenance_interval: NZUsize!(1), retained_marshal_blocks: 0, retained_qmdb_blocks: 0, }, 1, 0, ); let processor = Processor::new( app, databases, anchor(0, 0), StatefulMetrics::new(&context), Some(pruning), ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let mut mailbox = Mailbox::new(sender); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox.clone(), processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); let (acknowledgement, waiter1) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(block1), acknowledgement)); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } let (acknowledgement, waiter2) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(block2.clone()), acknowledgement)); let consensus_context = child.context(); let mut verify = Box::pin(mailbox.verify( (context.child("verify"), consensus_context), ancestry::from_iter([Arc::new(child), Arc::new(parent)]), )); assert!(poll!(&mut verify).is_pending()); select! { result = first_started => result.expect("verification should start before pruning"), _ = context.sleep(Duration::from_millis(100)) => { panic!( "verification did not start: flushes={} pruned={}", control.flushes.lock().len(), control.pruned.lock().len(), ); }, } assert!(control.pruned.lock().is_empty()); let release = control.flushes.lock().remove(0); release .send(Ok(())) .expect("target flush should be pending"); waiter1.await.expect("target block should be acknowledged"); first_release.closed().await; prune_started.await.expect("prune should start"); assert_eq!( control.flushes.lock().len(), 0, "the dirty successor must not overlap the conservative database prune", ); prune_release.send(()).expect("prune should remain active"); while control.pruned.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } assert_eq!( control.flushes.lock().len(), 1, "the dirty successor must start before verification retries resume", ); second_started .await .expect("live replay should restart after pruning"); second_release .send(()) .expect("restarted replay should remain active"); assert!(verify.await); assert_eq!(control.pruned.lock().clone(), vec![1]); assert_eq!(apply_calls.load(Ordering::SeqCst), 4); assert_eq!(verify_calls.load(Ordering::SeqCst), 1); let release = control.flushes.lock().remove(0); release.send(Ok(())).expect("newer flush should be pending"); waiter2.await.expect("newer block should be acknowledged"); actor.abort(); marshal.abort(); }); } #[test] fn prune_retries_wait_for_queued_finalization() { deterministic::Runner::timed(Duration::from_secs(10)).start(|context| async move { let genesis = TestBlock::new(0, 0); let block1 = TestBlock::child(&genesis, 1); let block2 = TestBlock::child(&block1, 2); let losing = TestBlock::child(&block2, 3); let winner = TestBlock::child(&block2, 4); let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"prune-retry", 1).schemes[0].clone(); let marshal = fixtures::marshal_fixture( context.child("marshal"), "prune-retry", scheme, None, NZUsize!(1), false, ) .await; let (verify_gate, verify_started, mut verify_release) = application_gate(); let (proposal_gate, proposal_started, proposal_release) = application_gate(); let observed_contexts: Arc>> = Arc::default(); let app = GatedApp { verify_gates: Arc::new(Mutex::new(VecDeque::from([verify_gate]))), proposal_gate: Arc::new(Mutex::new(Some(proposal_gate))), verify_valid: true, observed_contexts: observed_contexts.clone(), }; let control = FlushControl::default(); let (prune_started, prune_release) = control.gate_prune(); let databases = Shared::new("prune-retry", TestDb::gated(control.clone())); let pruning = Pruning::build( PruneConfig { maintenance_interval: NZUsize!(1), retained_marshal_blocks: 0, retained_qmdb_blocks: 0, }, 1, 0, ); let processor = Processor::new( app, databases, anchor(0, 0), StatefulMetrics::new(&context), Some(pruning), ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(8)); let mut mailbox = Mailbox::new(sender); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); let (acknowledgement, waiter1) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(block1), acknowledgement)); let (acknowledgement, waiter2) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(block2.clone()), acknowledgement)); let mut verifier = mailbox.clone(); let mut verify = Box::pin(verifier.verify( (context.child("verify"), losing.context()), ancestry::from_iter([Arc::new(losing), Arc::new(block2)]), )); assert!(poll!(&mut verify).is_pending()); verify_started .await .expect("verification should start before pruning"); while control.applied.load(Ordering::Relaxed) < 2 { context.sleep(Duration::from_millis(10)).await; } assert_eq!(control.flushes.lock().len(), 1); control .flushes .lock() .remove(0) .send(Ok(())) .expect("target flush should remain pending"); waiter1.await.expect("target block should be acknowledged"); prune_started.await.expect("prune should start"); verify_release.closed().await; let (acknowledgement, winner_waiter) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(winner.clone()), acknowledgement)); let proposal_context = TestBlock::child(&winner, 5).context(); let mut proposer = mailbox.clone(); let mut proposal = Box::pin(proposer.propose( (context.child("propose"), proposal_context), ancestry::from_iter([Arc::new(winner)]), (), )); assert!(poll!(&mut proposal).is_pending()); prune_release.send(()).expect("prune should remain active"); proposal_started .await .expect("proposal queued behind finalization should start"); let subscriber = mailbox.clone(); let mut databases = Box::pin(subscriber.subscribe_databases()); assert!(poll!(&mut databases).is_pending()); let result = select! { valid = &mut verify => Some(valid), _ = context.sleep(Duration::from_millis(100)) => None, }; assert_eq!( result, Some(false), "prune retry must observe the queued finalization", ); assert!(poll!(&mut databases).is_pending()); assert_eq!(observed_contexts.lock().len(), 1); assert_eq!(control.pruned.lock().as_slice(), [1]); proposal_release .send(()) .expect("proposal should remain active"); assert!(proposal.await.is_none()); drop(databases.await); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } control .flushes .lock() .remove(0) .send(Ok(())) .expect("block 2 sync should remain pending"); waiter2.await.expect("block 2 should be acknowledged"); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } control .flushes .lock() .remove(0) .send(Ok(())) .expect("winner sync should remain pending"); winner_waiter.await.expect("winner should be acknowledged"); actor.abort(); drop(marshal.guards); }); } /// Pruning waits for the flush that covers its target without waiting for newer state. #[test] fn prune_starts_after_target_sync() { deterministic::Runner::timed(Duration::from_secs(10)).start(|context| async move { // Marshal only receives prune requests here. Its actor never runs. let (verify_gate, verify_started, verify_release) = application_gate(); let (mut mailbox, control, _marshal, _actor) = spawn_processing_with_gates( &context, "gated-prune", Some(PruneConfig { maintenance_interval: NZUsize!(1), retained_marshal_blocks: 0, retained_qmdb_blocks: 0, }), VecDeque::from([verify_gate]), ) .await; let genesis = TestBlock::new(0, 0); let block1 = TestBlock::child(&genesis, 1); let block2 = TestBlock::child(&block1, 2); let block3 = TestBlock::child(&block2, 3); // Apply blocks 1 and 2 without releasing any flush: the loop must // stay live (both blocks applied) while no acknowledgement fires. let (acknowledgement, mut waiter1) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(block1), acknowledgement)); let (acknowledgement, mut waiter2) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(block2.clone()), acknowledgement)); // Queue a verification before pruning starts, then hold it in the // application until the prune is waiting on durability. let consensus_context = block3.context(); let mut verifier = mailbox.clone(); let mut verify = Box::pin(verifier.verify( (context.child("verify"), consensus_context), ancestry::from_iter([Arc::new(block3), Arc::new(block2)]), )); assert!(poll!(&mut verify).is_pending()); verify_started .await .expect("verification should start before pruning"); while control.applied.load(Ordering::Relaxed) < 2 { context.sleep(Duration::from_millis(10)).await; } assert_eq!(control.flushes.lock().len(), 1); assert!( poll!(&mut waiter1).is_pending() && poll!(&mut waiter2).is_pending(), "acknowledgements must wait for pending flushes", ); // Block 2 filled the retention window, but pruning must remain blocked behind the // target at block 1. context.sleep(Duration::from_millis(50)).await; assert!(control.pruned.lock().is_empty()); assert!( poll!(&mut waiter1).is_pending() && poll!(&mut waiter2).is_pending(), "acknowledgements must keep waiting for pending flushes", ); verify_release .send(()) .expect("verification should remain active"); select! { result = &mut verify => assert!(result), _ = context.sleep(Duration::from_millis(100)) => { panic!("pending prune blocked active verification"); }, } assert!(control.pruned.lock().is_empty()); // Releasing block 1 makes the prune target durable. Glue prunes before starting the // tracked successor for replayable block 2. let release = control.flushes.lock().remove(0); let _ = release.send(Ok(())); waiter1.await.expect("block 1 acknowledgement"); while control.pruned.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } assert_eq!(control.pruned.lock().clone(), vec![1]); assert!( poll!(&mut waiter2).is_pending(), "block 2 must stay unacknowledged while its flush is pending", ); // Releasing block 2's flush releases its acknowledgement. while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } let release = control.flushes.lock().remove(0); let _ = release.send(Ok(())); waiter2.await.expect("block 2 acknowledgement"); }); } #[test] fn finalized_handoff_survives_pending_sync() { deterministic::Runner::timed(Duration::from_secs(10)).start(|context| async move { let genesis = TestBlock::new(0, 0); let block1 = TestBlock::child(&genesis, 1); let block2 = TestBlock::child(&block1, 2); let mut signing = context.child("signing"); let scheme = scheme_mocks::fixture(&mut signing, b"handoff-sync-order", 1).schemes[0].clone(); let marshal = fixtures::marshal_fixture( context.child("marshal"), "handoff-sync-order", scheme, None, NZUsize!(2), false, ) .await; let (finalized_gate, finalized_started, finalized_release) = application_gate(); let capture_calls = Arc::new(AtomicUsize::new(0)); let applied_finalizations: Arc>> = Arc::default(); let app = ReplayGatedApp { gates: Arc::default(), verify_gate: Arc::default(), finalized_gate: Arc::new(Mutex::new(Some(finalized_gate))), gate_height: block1.height(), unexecutable: None, apply_calls: Arc::new(AtomicUsize::new(0)), capture_calls: capture_calls.clone(), verify_calls: Arc::new(AtomicUsize::new(0)), applied_finalizations: applied_finalizations.clone(), }; let control = FlushControl::default(); let processor = Processor::new( app, Shared::new("test", TestDb::gated(control.clone())), anchor(0, 0), StatefulMetrics::new(&context), None, ); let (sender, receiver) = actor_mailbox::new(context.child("mailbox"), NZUsize!(2)); let mut mailbox = Mailbox::new(sender); let processing = Processing { context: ContextCell::new(context.child("processing")), mailbox: receiver, provider: (), marshal: marshal.mailbox, processor, deferred_verifications: Vec::new(), skip_finalized_until: None, }; let actor = context.child("loop").spawn(move |_| processing.start()); let (acknowledgement, mut waiter1) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(block1), acknowledgement)); finalized_started .await .expect("finalized handoff should start"); assert_eq!(control.applied.load(Ordering::Relaxed), 1); assert_eq!(capture_calls.load(Ordering::SeqCst), 1); assert_eq!(applied_finalizations.lock().as_slice(), &[Height::new(1)]); assert_eq!(control.flushes.lock().len(), 1); assert!(poll!(&mut waiter1).is_pending()); finalized_release .send(()) .expect("finalized handoff should remain active"); let (acknowledgement, mut waiter2) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(block2), acknowledgement)); while control.applied.load(Ordering::Relaxed) < 2 { context.sleep(Duration::from_millis(10)).await; } assert_eq!(capture_calls.load(Ordering::SeqCst), 2); assert_eq!( applied_finalizations.lock().as_slice(), &[Height::new(1), Height::new(2)], ); assert_eq!(control.flushes.lock().len(), 1); assert!(poll!(&mut waiter1).is_pending()); assert!(poll!(&mut waiter2).is_pending()); control .flushes .lock() .remove(0) .send(Ok(())) .expect("first sync should remain pending"); waiter1.await.expect("block 1 should be acknowledged"); assert!(poll!(&mut waiter2).is_pending()); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } control .flushes .lock() .remove(0) .send(Ok(())) .expect("successor sync should remain pending"); waiter2.await.expect("block 2 should be acknowledged"); actor.abort(); drop(marshal.guards); }); } #[test] fn stable_leader_finalizations_coalesce_while_sync_pending() { deterministic::Runner::timed(Duration::from_secs(10)).start(|context| async move { let (mut mailbox, control, _marshal, _actor) = spawn_processing(&context, "gated-coalesced-sync", None).await; const BLOCKS: u64 = 3; let (acknowledgement, mut waiter1) = Exact::handle(); let _ = mailbox.report(Update::Block( Arc::new(TestBlock::new(1, 1)), acknowledgement, )); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } let mut waiters = Vec::with_capacity(BLOCKS as usize - 1); for height in 2..=BLOCKS { let (acknowledgement, waiter) = Exact::handle(); let _ = mailbox.report(Update::Block( Arc::new(TestBlock::new(height, height as u8)), acknowledgement, )); waiters.push(waiter); } while control.applied.load(Ordering::Relaxed) < BLOCKS as usize { context.sleep(Duration::from_millis(10)).await; } assert_eq!( control.flushes.lock().len(), 1, "a pending sync must coalesce later finalized state instead of starting a second sync", ); assert!(poll!(&mut waiter1).is_pending()); for waiter in &mut waiters { assert!(poll!(waiter).is_pending()); } control .flushes .lock() .remove(0) .send(Ok(())) .expect("first sync should remain pending"); waiter1.await.expect("first block acknowledgement"); for waiter in &mut waiters { assert!(poll!(waiter).is_pending()); } while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } assert_eq!(control.flushes.lock().len(), 1); control .flushes .lock() .remove(0) .send(Ok(())) .expect("successor sync should remain pending"); for acknowledgement in futures::future::join_all(waiters).await { acknowledgement.expect("stable-leader block acknowledgement"); } assert!(control.flushes.lock().is_empty()); }); } #[test] fn successor_sync_drives_verification_holding_database_read() { deterministic::Runner::timed(Duration::from_secs(10)).start(|context| async move { let (verify_gate, verify_started, verify_release) = application_gate(); let (mut mailbox, control, _marshal, _actor) = spawn_read_gated_processing( &context, "successor-sync-read-owner", verify_gate, None, ) .await; let genesis = TestBlock::new(0, 0); let block1 = TestBlock::child(&genesis, 1); let block2 = TestBlock::child(&block1, 2); let (acknowledgement, waiter1) = Exact::handle(); let _ = mailbox.report(Update::Block( Arc::new(block1.clone()), acknowledgement, )); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } let (acknowledgement, mut waiter2) = Exact::handle(); let _ = mailbox.report(Update::Block( Arc::new(block2.clone()), acknowledgement, )); while control.applied.load(Ordering::Relaxed) < 2 { context.sleep(Duration::from_millis(10)).await; } let block3 = TestBlock::child(&block2, 3); let consensus_context = block3.context(); let mut verifier = mailbox.clone(); let mut verify = Box::pin(verifier.verify( (context.child("verify"), consensus_context), ancestry::from_iter([Arc::new(block3), Arc::new(block2)]), )); assert!(poll!(&mut verify).is_pending()); verify_started .await .expect("verification should acquire the database read"); control .flushes .lock() .remove(0) .send(Ok(())) .expect("first sync should remain pending"); waiter1.await.expect("first block acknowledgement"); assert!(poll!(&mut waiter2).is_pending()); verify_release .send(()) .expect("verification should remain active"); select! { result = &mut verify => assert!(result), _ = context.sleep(Duration::from_millis(100)) => { panic!("successor sync stopped polling the verification that owned its read lock"); }, } while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } control .flushes .lock() .remove(0) .send(Ok(())) .expect("successor sync should remain pending"); waiter2.await.expect("second block acknowledgement"); }); } #[test] fn shutdown_preempts_successor_sync_waiting_for_verification_reader() { deterministic::Runner::timed(Duration::from_secs(10)).start(|context| async move { let (verify_gate, verify_started, _verify_release) = application_gate(); let (mut mailbox, control, _marshal, actor) = spawn_read_gated_processing(&context, "successor-sync-shutdown", verify_gate, None) .await; let genesis = TestBlock::new(0, 0); let block1 = TestBlock::child(&genesis, 1); let block2 = TestBlock::child(&block1, 2); let (acknowledgement, waiter1) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(block1.clone()), acknowledgement)); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } let (acknowledgement, waiter2) = Exact::handle(); let _ = mailbox.report(Update::Block(Arc::new(block2.clone()), acknowledgement)); while control.applied.load(Ordering::Relaxed) < 2 { context.sleep(Duration::from_millis(10)).await; } let block3 = TestBlock::child(&block2, 3); let mut verifier = mailbox.clone(); let verify = verifier.verify( (context.child("verify"), block3.context()), ancestry::from_iter([Arc::new(block3), Arc::new(block2)]), ); futures::pin_mut!(verify); assert!(poll!(&mut verify).is_pending()); verify_started .await .expect("verification should acquire the database read"); control .flushes .lock() .remove(0) .send(Ok(())) .expect("first sync should remain pending"); waiter1.await.expect("first block acknowledgement"); let stopper = context.child("stopper"); let stop = context .child("stop") .spawn(|_| async move { stopper.stop(0, Some(Duration::from_millis(100))).await }); assert!( stop.await.expect("stop task should finish").is_ok(), "shutdown must preempt successor sync acquisition", ); actor.await.expect("processing actor should stop cleanly"); assert!( waiter2.await.is_err(), "shutdown must cancel the dirty acknowledgement", ); }); } /// An aborted target flush must stop processing before pruning can discard its recovery state. #[test] fn aborted_target_flush_prevents_prune() { deterministic::Runner::timed(Duration::from_secs(10)).start(|context| async move { let (mut mailbox, control, _marshal, actor) = spawn_processing( &context, "gated-aborted-prune", Some(PruneConfig { maintenance_interval: NZUsize!(1), retained_marshal_blocks: 0, retained_qmdb_blocks: 0, }), ) .await; let (acknowledgement, waiter1) = Exact::handle(); let _ = mailbox.report(Update::Block( Arc::new(TestBlock::new(1, 1)), acknowledgement, )); let (acknowledgement, waiter2) = Exact::handle(); let _ = mailbox.report(Update::Block( Arc::new(TestBlock::new(2, 2)), acknowledgement, )); while control.applied.load(Ordering::Relaxed) < 2 { context.sleep(Duration::from_millis(10)).await; } assert_eq!(control.flushes.lock().len(), 1); drop(control.flushes.lock().remove(0)); actor.await.expect("processing actor should stop"); assert!( waiter1.await.is_err(), "aborted target flush must cancel the first acknowledgement", ); assert!( waiter2.await.is_err(), "aborted target flush must cancel the second acknowledgement", ); assert!( control.pruned.lock().is_empty(), "aborted flush must prevent pruning", ); }); } /// While the loop is idle, a completed flush must release its acknowledgement without /// displacing a simultaneously reported block, while an incomplete flush must cancel its /// acknowledgement when processing stops. #[test] fn idle_acks_follow_flush_outcome() { deterministic::Runner::timed(Duration::from_secs(10)).start(|context| async move { let (mut mailbox, control, _marshal, actor) = spawn_processing(&context, "gated-idle", None).await; // Park the loop idle with block 1's flush pending. let (acknowledgement, mut waiter1) = Exact::handle(); let _ = mailbox.report(Update::Block( Arc::new(TestBlock::new(1, 1)), acknowledgement, )); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } assert!(poll!(&mut waiter1).is_pending()); context.sleep(Duration::from_millis(50)).await; // Release the flush and report block 2 in the same scheduling // window: the completion must fire block 1's acknowledgement // without displacing the new message. let release = control.flushes.lock().remove(0); let _ = release.send(Ok(())); let (acknowledgement, waiter2) = Exact::handle(); let _ = mailbox.report(Update::Block( Arc::new(TestBlock::new(2, 2)), acknowledgement, )); waiter1.await.expect("block 1 acknowledgement"); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } context.sleep(Duration::from_millis(50)).await; // Dropping block 2's release resolves its flush as shutdown. The // acknowledgement is canceled so marshal stops without advancing // its floor past unflushed state. drop(control.flushes.lock().remove(0)); actor.await.expect("processing actor should stop"); assert!( waiter2.await.is_err(), "unflushed block acknowledgement must be canceled", ); }); } #[test] fn ready_aborted_flush_stops_processing() { deterministic::Runner::timed(Duration::from_secs(10)).start(|context| async move { let (mut mailbox, control, _marshal, actor) = spawn_processing(&context, "gated-ready-abort", None).await; let (acknowledgement, waiter1) = Exact::handle(); let _ = mailbox.report(Update::Block( Arc::new(TestBlock::new(1, 1)), acknowledgement, )); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } let (acknowledgement, waiter2) = Exact::handle(); let _ = mailbox.report(Update::Block( Arc::new(TestBlock::new(2, 2)), acknowledgement, )); while control.applied.load(Ordering::Relaxed) < 2 { context.sleep(Duration::from_millis(10)).await; } assert_eq!(control.flushes.lock().len(), 1); drop(control.flushes.lock().remove(0)); actor.await.expect("processing actor should stop"); assert!(control.flushes.lock().is_empty()); assert!( waiter1.await.is_err(), "the active unflushed acknowledgement must be canceled", ); assert!( waiter2.await.is_err(), "the queued unflushed acknowledgement must be canceled", ); }); } /// Stopping processing with a flush in flight must cancel marshal's acknowledgement. #[test] fn shutdown_cancels_pending_flush_ack() { deterministic::Runner::timed(Duration::from_secs(10)).start(|context| async move { let (mut mailbox, control, _marshal, actor) = spawn_processing(&context, "gated-shutdown", None).await; let (acknowledgement, waiter) = Exact::handle(); let _ = mailbox.report(Update::Block( Arc::new(TestBlock::new(1, 1)), acknowledgement, )); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } drop(mailbox); actor.await.expect("processing actor should stop"); assert!( waiter.await.is_err(), "shutdown must cancel in-flight acknowledgements", ); }); } /// A flush failure must panic the processing loop with the database identified and leave the /// block unacknowledged. #[test] #[should_panic(expected = "database sync failed (type")] fn flush_failure_panics_processing() { deterministic::Runner::timed(Duration::from_secs(10)).start(|context| async move { let (mut mailbox, control, _marshal, _actor) = spawn_processing(&context, "gated-failure", None).await; let (acknowledgement, _waiter) = Exact::handle(); let _ = mailbox.report(Update::Block( Arc::new(TestBlock::new(1, 1)), acknowledgement, )); while control.flushes.lock().is_empty() { context.sleep(Duration::from_millis(10)).await; } let release = control.flushes.lock().remove(0); let _ = release.send(Err(RuntimeError::WriteFailed)); // The active sync panics when the loop next polls it. loop { context.sleep(Duration::from_millis(100)).await; } }); } #[test] fn skip_finalized_block_skips_through_target_height() { let mut skip_until = Some(Height::new(3)); assert!(skip_finalized_block(&mut skip_until, Height::new(1))); assert_eq!(skip_until, Some(Height::new(3))); assert!(skip_finalized_block(&mut skip_until, Height::new(3))); assert_eq!(skip_until, None); assert!(!skip_finalized_block(&mut skip_until, Height::new(4))); } #[test] fn skip_finalized_block_clears_stale_target() { let mut skip_until = Some(Height::new(3)); assert!(!skip_finalized_block(&mut skip_until, Height::new(4))); assert_eq!(skip_until, None); } }