use std::{sync::Arc, time::Duration};
use tokio::{
sync::{mpsc, watch},
task::JoinHandle,
};
use tokio_util::sync::CancellationToken;
use zakura_chain::block;
use super::mock_blocksync::{
mainnet_genesis_hash, MockApplyFrontier, SyntheticBlockCorpus, SyntheticBlockShape,
};
use super::{SyntheticBlockSyncPeers, TraceCapture};
use crate::zakura::{
BlockApplyResult, BlockSyncAction, BlockSyncEvent, BlockSyncFrontiers, BlockSyncHandle,
ZakuraTrace,
};
use crate::BoxError;
mod invariants;
mod peer;
mod scenario;
#[cfg(test)]
mod tests;
pub(crate) use invariants::{
assert_core as assert_core_invariants, report as invariant_report, InvariantReport,
};
pub(crate) use scenario::*;
pub(crate) async fn run_scenario(
scenario: &Scenario,
trace: ZakuraTrace,
) -> Result<FuzzOutcome, BoxError> {
let corpus = SyntheticBlockCorpus::generate(
scenario.blocks,
scenario.seed,
SyntheticBlockShape {
target_block_bytes: scenario.target_block_bytes,
},
);
let target = corpus.target_height();
let genesis_hash = mainnet_genesis_hash();
let initial_header = scenario.initial_best_header.min(target);
let initial_header_hash = corpus_hash(&corpus, initial_header);
const EMPTY_STATE_HEADER_QUIET_MIN_LAG: u32 = 400;
let initial_verified = if initial_header.0 >= EMPTY_STATE_HEADER_QUIET_MIN_LAG {
target.min(block::Height(1)).min(initial_header)
} else {
block::Height(0)
};
let initial_verified_hash = corpus_hash(&corpus, initial_verified);
let apply = MockApplyFrontier::with_committed_height(corpus.clone(), initial_verified);
let initial = fuzz_snapshot(
1,
1,
1,
zakura_header_chain::Frontier::new(block::Height(0), genesis_hash),
zakura_header_chain::Frontier::new(initial_verified, initial_verified_hash),
zakura_header_chain::Frontier::new(initial_header, initial_header_hash),
);
let (snapshots, committed_views) =
watch::channel(Some(zakura_header_chain::CommittedHeaderChainView::new(
initial,
zakura_header_chain::BodyWorkEpoch::default(),
)));
let shutdown = CancellationToken::new();
let mut startup = crate::zakura::BlockSyncStartup::new_with_committed_views(
BlockSyncFrontiers {
finalized_height: block::Height(0),
verified_block_tip: initial_verified,
verified_block_hash: initial_verified_hash,
},
(initial_header, initial_header_hash),
committed_views,
scenario.config.clone(),
);
startup.trace = trace.clone();
startup.shutdown = shutdown.clone();
let (handle, actions, reactor_task) = crate::zakura::spawn_block_sync_reactor(startup);
let (committed_tx, mut committed_rx) = watch::channel(initial_verified);
let mut tasks = Vec::new();
tasks.push(spawn_action_driver(
handle.clone(),
actions,
corpus.clone(),
target,
apply.clone(),
scenario.commit,
committed_tx,
shutdown.clone(),
));
if !scenario.timeline.is_empty() {
tasks.push(spawn_timeline_driver(
snapshots.clone(),
corpus.clone(),
apply.clone(),
scenario.timeline.clone(),
shutdown.clone(),
));
}
let peers = Arc::new(SyntheticBlockSyncPeers::new(
scenario.config.clone(),
handle.clone(),
scenario.transport_queue_depth.unwrap_or(1024),
));
for spec in &scenario.peers {
tasks.push(peer::spawn_peer_lifecycle(
peers.clone(),
corpus.clone(),
*spec,
scenario.seed,
corpus_hash(&corpus, spec.servable_high),
shutdown.clone(),
));
}
let running = RunningHarness {
shutdown: shutdown.clone(),
reactor_task,
tasks,
_peers: peers,
};
let reached = tokio::time::timeout(
scenario.deadline,
committed_rx.wait_for(|height| *height >= target),
)
.await
.ok()
.and_then(|result| result.ok())
.map(|height| *height);
let committed_tip = reached.unwrap_or_else(|| *committed_rx.borrow());
running.stop().await;
Ok(FuzzOutcome {
committed_tip,
target,
})
}
struct RunningHarness {
shutdown: CancellationToken,
reactor_task: JoinHandle<()>,
tasks: Vec<JoinHandle<()>>,
_peers: Arc<SyntheticBlockSyncPeers>,
}
impl Drop for RunningHarness {
fn drop(&mut self) {
self.shutdown.cancel();
self.reactor_task.abort();
for task in &self.tasks {
task.abort();
}
}
}
impl RunningHarness {
async fn stop(mut self) {
self.shutdown.cancel();
stop_task(&mut self.reactor_task).await;
for task in &mut self.tasks {
stop_task(task).await;
}
}
}
async fn stop_task(task: &mut JoinHandle<()>) {
if tokio::time::timeout(Duration::from_secs(2), &mut *task)
.await
.is_err()
{
task.abort();
let _ = task.await;
}
}
#[allow(clippy::too_many_arguments)]
fn spawn_action_driver(
handle: BlockSyncHandle,
mut actions: mpsc::Receiver<BlockSyncAction>,
corpus: SyntheticBlockCorpus,
target: block::Height,
apply: MockApplyFrontier,
commit: CommitProfile,
committed_tx: watch::Sender<block::Height>,
shutdown: CancellationToken,
) -> JoinHandle<()> {
tokio::spawn(async move {
let mut applied = 0u64;
loop {
let action = tokio::select! {
_ = shutdown.cancelled() => break,
action = actions.recv() => match action {
Some(action) => action,
None => break,
},
};
match action {
BlockSyncAction::QueryNeededBlocks {
query_id,
from,
limit,
best_header_tip,
scope,
} => {
let start = from;
let metas = if limit == 0 {
Vec::new()
} else {
let end = (start + i64::from(limit.saturating_sub(1)))
.unwrap_or(block::Height::MAX)
.min(best_header_tip)
.min(target);
if start <= end {
corpus.metas_between(start, end)
} else {
Vec::new()
}
};
if handle
.send(BlockSyncEvent::ScopedNeededBlocks {
query_id,
scope,
body_anchor: {
let frontiers = apply.frontiers();
zakura_header_chain::Frontier::new(
frontiers.verified_block_tip,
frontiers.verified_block_hash,
)
},
blocks: metas,
})
.await
.is_err()
{
break;
}
}
BlockSyncAction::QueryBlocksByHeightRange { peer, start, count } => {
let blocks = corpus.blocks_in_range(start, count, target);
if handle
.send(BlockSyncEvent::BlockRangeResponseReady {
peer,
start_height: start,
requested_count: count,
blocks,
})
.await
.is_err()
{
break;
}
}
BlockSyncAction::SubmitBlock {
owner,
source,
token,
block,
} => {
if !commit.per_commit_delay.is_zero()
&& sleep_or_cancel(&shutdown, commit.per_commit_delay).await
{
break;
}
let height = block
.coinbase_height()
.expect("synthetic submitted block has height");
let outcome = apply.apply(block.as_ref());
if outcome.result == BlockApplyResult::Committed {
let _ = committed_tx.send(outcome.frontiers.verified_block_tip);
}
if handle
.send(BlockSyncEvent::BlockApplyFinished {
owner,
source,
token,
height,
hash: block.hash(),
outcome: crate::zakura::block_sync::test_block_apply_outcome(
outcome.result,
),
})
.await
.is_err()
{
break;
}
if outcome.result == BlockApplyResult::Committed {
let _ = handle
.send(BlockSyncEvent::ChainTipGrow(outcome.frontiers))
.await;
}
applied = applied.saturating_add(1);
if let Some(burst) = commit.burst {
if burst.every_commits > 0
&& applied.is_multiple_of(burst.every_commits)
&& !burst.duration.is_zero()
&& sleep_or_cancel(&shutdown, burst.duration).await
{
break;
}
}
}
BlockSyncAction::RecordBodyUnavailable { .. }
| BlockSyncAction::RecordBodyInvalid { .. }
| BlockSyncAction::RestartBodyAvailability { .. }
| BlockSyncAction::RetryBodyAvailability { .. } => {}
BlockSyncAction::Misbehavior { .. } => {}
}
}
})
}
fn spawn_timeline_driver(
snapshots: watch::Sender<Option<zakura_header_chain::CommittedHeaderChainView>>,
corpus: SyntheticBlockCorpus,
apply: MockApplyFrontier,
mut timeline: Vec<TipEvent>,
shutdown: CancellationToken,
) -> JoinHandle<()> {
tokio::spawn(async move {
timeline.sort_by_key(|event| event.at);
let mut elapsed = Duration::ZERO;
for event in timeline {
let wait = event.at.saturating_sub(elapsed);
if sleep_or_cancel(&shutdown, wait).await {
return;
}
elapsed = event.at;
let apply_frontiers = if let TipEventKind::VerifiedReset(height) = event.kind {
apply.reset_to(height)
} else {
apply.frontiers()
};
let current_view = snapshots
.borrow()
.clone()
.expect("the fuzz harness starts after semantic handoff");
let mut current = current_view.snapshot;
current.state_version =
zakura_header_chain::StateVersion::new(current.state_version.get() + 1);
current.frontiers.finalized = zakura_header_chain::Frontier::new(
apply_frontiers.finalized_height,
apply_frontiers.verified_block_hash,
);
current.frontiers.verified_best = zakura_header_chain::Frontier::new(
apply_frontiers.verified_block_tip,
apply_frontiers.verified_block_hash,
);
apply_tip_event(&corpus, &mut current, event.kind);
let body_work_epoch = if matches!(
event.kind,
TipEventKind::HeaderReanchor(_) | TipEventKind::VerifiedReset(_)
) {
current_view
.body_work_epoch
.checked_next()
.expect("the fuzz timeline cannot exhaust its body-work epoch")
} else {
current_view.body_work_epoch
};
snapshots
.send(Some(zakura_header_chain::CommittedHeaderChainView::new(
current,
body_work_epoch,
)))
.expect("the fuzz reactor keeps its committed-snapshot receiver");
}
})
}
fn apply_tip_event(
corpus: &SyntheticBlockCorpus,
snapshot: &mut zakura_header_chain::EngineSnapshot,
kind: TipEventKind,
) {
match kind {
TipEventKind::GrowTo(height) => {
snapshot.header_generation =
zakura_header_chain::HeaderGeneration::new(snapshot.header_generation.get() + 1);
snapshot.frontiers.header_best =
zakura_header_chain::Frontier::new(height, corpus_hash(corpus, height));
}
TipEventKind::HeaderReanchor(height) => {
snapshot.header_generation =
zakura_header_chain::HeaderGeneration::new(snapshot.header_generation.get() + 1);
snapshot.frontiers.header_best =
zakura_header_chain::Frontier::new(height, corpus_hash(corpus, height));
}
TipEventKind::VerifiedReset(height) => {
snapshot.verified_generation = zakura_header_chain::VerifiedGeneration::new(
snapshot.verified_generation.get() + 1,
);
snapshot.frontiers.verified_best =
zakura_header_chain::Frontier::new(height, corpus_hash(corpus, height));
}
}
snapshot.header_best_score = zakura_header_chain::ChainScore::new(
zakura_header_chain::SuffixWork::zero(),
snapshot.frontiers.header_best.hash,
);
}
fn fuzz_snapshot(
state_version: u64,
header_generation: u64,
verified_generation: u64,
finalized: zakura_header_chain::Frontier,
verified_best: zakura_header_chain::Frontier,
header_best: zakura_header_chain::Frontier,
) -> zakura_header_chain::EngineSnapshot {
zakura_header_chain::EngineSnapshot {
mode: zakura_header_chain::EngineMode::Integrated,
state_version: zakura_header_chain::StateVersion::new(state_version),
header_generation: zakura_header_chain::HeaderGeneration::new(header_generation),
verified_generation: zakura_header_chain::VerifiedGeneration::new(verified_generation),
frontiers: zakura_header_chain::FrontierSet {
finalized,
header_best,
verified_best,
},
header_best_score: zakura_header_chain::ChainScore::new(
zakura_header_chain::SuffixWork::zero(),
header_best.hash,
),
oldest_retained_height: finalized.height,
alarms: zakura_header_chain::AlarmSet::default(),
}
}
fn corpus_hash(corpus: &SyntheticBlockCorpus, height: block::Height) -> block::Hash {
if height == block::Height(0) {
mainnet_genesis_hash()
} else {
corpus
.block_at(height)
.map(|block| block.hash())
.unwrap_or_else(mainnet_genesis_hash)
}
}
pub(crate) async fn sleep_or_cancel(shutdown: &CancellationToken, duration: Duration) -> bool {
if duration.is_zero() {
return shutdown.is_cancelled();
}
tokio::select! {
_ = shutdown.cancelled() => true,
_ = tokio::time::sleep(duration) => false,
}
}
pub(crate) fn run_trace(name: &str) -> std::io::Result<(TraceCapture, ZakuraTrace)> {
let mut capture = TraceCapture::for_test(name)?;
let trace = ZakuraTrace::new(capture.tracer_for_node(0), "00");
Ok((capture, trace))
}