#![allow(clippy::unwrap_in_result)]
use std::{
collections::{HashMap, HashSet},
iter,
sync::{
atomic::{AtomicUsize, Ordering},
Arc,
},
task::{Context, Poll},
time::Duration,
};
use color_eyre::Report;
use futures::{Future, FutureExt, StreamExt};
use tower::{timeout::Timeout, Service};
use zakura_chain::{
block::{self, Block, Height},
chain_tip::mock::{MockChainTip, MockChainTipSender},
serialization::ZcashDeserializeInto,
};
use zakura_consensus::{
Config as ConsensusConfig, RouterError, VerifyBlockError, VerifyCheckpointError,
};
use zakura_network::{InventoryResponse, PeerSocketAddr};
use zakura_state::Config as StateConfig;
use zakura_test::mock_service::{MockService, PanicAssertion};
use zakura_network as zn;
use zakura_state as zs;
use crate::{
components::{
auth_download_height::poison_coinbase_height,
sync::{
self,
downloads::{BlockDownloadVerifyError, Downloads},
legacy_trace::LegacySyncTrace,
SyncStatus,
},
ChainSync,
},
config::ZakuradConfig,
};
use InventoryResponse::*;
type TestChainSync = ChainSync<
MockService<zn::Request, zn::Response, PanicAssertion>,
MockService<zs::Request, zs::Response, PanicAssertion>,
MockService<zakura_consensus::Request, block::Hash, PanicAssertion>,
MockChainTip,
>;
#[derive(Clone, Debug)]
struct NeverReadyNetwork;
impl Service<zn::Request> for NeverReadyNetwork {
type Response = zn::Response;
type Error = crate::BoxError;
type Future = futures::future::Pending<Result<Self::Response, Self::Error>>;
fn poll_ready(&mut self, _context: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
Poll::Pending
}
fn call(&mut self, _request: zn::Request) -> Self::Future {
futures::future::pending()
}
}
const MAX_SERVICE_REQUEST_DELAY: Duration = Duration::from_millis(1000);
const STALLED_SERVICE_REQUEST_DELAY: Duration = Duration::from_secs(30 * 60);
#[test]
fn oversized_find_blocks_response_is_rejected() {
let hash = block::Hash([0; 32]);
assert!(sync::has_valid_tips_response_hash_count(&vec![
hash;
sync::MAX_TIPS_RESPONSE_HASH_COUNT
]));
assert!(!sync::has_valid_tips_response_hash_count(&vec![
hash;
sync::MAX_TIPS_RESPONSE_HASH_COUNT
+ 1
]));
let raw = vec![hash; sync::MAX_TIPS_RESPONSE_HASH_COUNT + 1];
let stripped = &raw[..raw.len() - 1];
assert_eq!(stripped.len(), sync::MAX_TIPS_RESPONSE_HASH_COUNT);
assert!(sync::has_valid_tips_response_hash_count(stripped));
}
#[tokio::test(start_paused = true)]
async fn sync_blocks_ok() -> Result<(), crate::BoxError> {
let (
chain_sync_future,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup();
let block0: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES.zcash_deserialize_into()?;
let block0_hash = block0.hash();
let block1: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?;
let block1_hash = block1.hash();
let block2: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_2_BYTES.zcash_deserialize_into()?;
let block2_hash = block2.hash();
let block3: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_3_BYTES.zcash_deserialize_into()?;
let block3_hash = block3.hash();
let block4: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_4_BYTES.zcash_deserialize_into()?;
let block4_hash = block4.hash();
let block5: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_5_BYTES.zcash_deserialize_into()?;
let block5_hash = block5.hash();
let chain_sync_task_handle = tokio::spawn(chain_sync_future);
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block0_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block0.clone(),
None,
))]));
block_verifier_router
.expect_request(zakura_consensus::Request::Commit(block0))
.await
.respond(block0_hash);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(Some(zs::KnownBlock::BestChain)));
state_service
.expect_request(zs::Request::BlockLocator)
.await
.respond(zs::Response::BlockLocator(vec![block0_hash]));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block0_hash],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
block1_hash, block2_hash, block3_hash, ]));
state_service
.expect_request(zs::Request::KnownBlock(block1_hash))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(block2_hash))
.await
.respond(zs::Response::KnownBlock(None));
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block0_hash],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test obtain tips error")));
}
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service
.expect_request(zs::Request::KnownBlock(block1_hash))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(block2_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block1_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block1.clone(),
None,
))]));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block2_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block2.clone(),
None,
))]));
let mut remaining_blocks: HashMap<block::Hash, Arc<Block>> =
[(block1_hash, block1), (block2_hash, block2)]
.iter()
.cloned()
.collect();
for _ in 1..=2 {
block_verifier_router
.expect_request_that(|req| remaining_blocks.remove(&req.block().hash()).is_some())
.await
.respond_with(|req| req.block().hash());
}
assert_eq!(
remaining_blocks,
HashMap::new(),
"expected all non-tip blocks to be verified by obtain tips"
);
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block1_hash],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
block2_hash, block3_hash, block4_hash,
block5_hash, ]));
for hash in [block3_hash, block4_hash] {
state_service
.expect_request(zs::Request::KnownBlock(hash))
.await
.respond(zs::Response::KnownBlock(None));
}
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block1_hash],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test extend tips error")));
}
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block3_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block3.clone(),
None,
))]));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block4_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block4.clone(),
None,
))]));
let mut remaining_blocks: HashMap<block::Hash, Arc<Block>> =
[(block3_hash, block3), (block4_hash, block4)]
.iter()
.cloned()
.collect();
for _ in 3..=4 {
block_verifier_router
.expect_request_that(|req| remaining_blocks.remove(&req.block().hash()).is_some())
.await
.respond_with(|req| req.block().hash());
}
assert_eq!(
remaining_blocks,
HashMap::new(),
"expected all non-tip blocks to be verified by extend tips"
);
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
let chain_sync_result = chain_sync_task_handle.now_or_never();
assert!(
chain_sync_result.is_none(),
"unexpected error or panic in chain sync task: {chain_sync_result:?}",
);
Ok(())
}
#[tokio::test(start_paused = true)]
async fn sync_singleton_obtain_tips_ok() -> Result<(), crate::BoxError> {
let (
chain_sync_future,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup();
let block0: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES.zcash_deserialize_into()?;
let block0_hash = block0.hash();
let block1: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?;
let block1_hash = block1.hash();
let chain_sync_task_handle = tokio::spawn(chain_sync_future);
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block0_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block0.clone(),
None,
))]));
block_verifier_router
.expect_request(zakura_consensus::Request::Commit(block0))
.await
.respond(block0_hash);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(Some(zs::KnownBlock::BestChain)));
state_service
.expect_request(zs::Request::BlockLocator)
.await
.respond(zs::Response::BlockLocator(vec![block0_hash]));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block0_hash],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![block1_hash]));
state_service
.expect_request(zs::Request::KnownBlock(block1_hash))
.await
.respond(zs::Response::KnownBlock(None));
for _ in 1..sync::FANOUT {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block0_hash],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test obtain tips error")));
}
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service
.expect_request(zs::Request::KnownBlock(block1_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block1_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block1.clone(),
None,
))]));
block_verifier_router
.expect_request(zakura_consensus::Request::Commit(block1))
.await
.respond(block1_hash);
let chain_sync_result = chain_sync_task_handle.now_or_never();
assert!(
chain_sync_result.is_none(),
"unexpected error or panic in chain sync task: {chain_sync_result:?}",
);
Ok(())
}
#[tokio::test(start_paused = true)]
async fn sync_singleton_extend_tips_ok() -> Result<(), crate::BoxError> {
let (
chain_sync_future,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup();
let block0: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES.zcash_deserialize_into()?;
let block0_hash = block0.hash();
let block1: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?;
let block1_hash = block1.hash();
let block2: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_2_BYTES.zcash_deserialize_into()?;
let block2_hash = block2.hash();
let block3: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_3_BYTES.zcash_deserialize_into()?;
let block3_hash = block3.hash();
let chain_sync_task_handle = tokio::spawn(chain_sync_future);
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block0_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block0.clone(),
None,
))]));
block_verifier_router
.expect_request(zakura_consensus::Request::Commit(block0))
.await
.respond(block0_hash);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(Some(zs::KnownBlock::BestChain)));
state_service
.expect_request(zs::Request::BlockLocator)
.await
.respond(zs::Response::BlockLocator(vec![block0_hash]));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block0_hash],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
block1_hash,
block2_hash,
block3_hash, ]));
state_service
.expect_request(zs::Request::KnownBlock(block1_hash))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(block2_hash))
.await
.respond(zs::Response::KnownBlock(None));
for _ in 1..sync::FANOUT {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block0_hash],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test obtain tips error")));
}
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service
.expect_request(zs::Request::KnownBlock(block1_hash))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(block2_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block1_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block1.clone(),
None,
))]));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block2_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block2.clone(),
None,
))]));
let mut remaining_blocks: HashMap<block::Hash, Arc<Block>> =
[(block1_hash, block1), (block2_hash, block2)]
.iter()
.cloned()
.collect();
for _ in 1..=2 {
block_verifier_router
.expect_request_that(|req| {
matches!(req, zakura_consensus::Request::Commit(_))
&& remaining_blocks.remove(&req.block().hash()).is_some()
})
.await
.respond_with(|req| req.block().hash());
}
assert!(
remaining_blocks.is_empty(),
"expected obtain tips to verify blocks 1 and 2; remaining blocks: {:?}",
remaining_blocks.keys().collect::<Vec<_>>(),
);
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block1_hash],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
block2_hash, block3_hash, ]));
state_service
.expect_request(zs::Request::KnownBlock(block3_hash))
.await
.respond(zs::Response::KnownBlock(None));
for _ in 1..sync::FANOUT {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block1_hash],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test extend tips error")));
}
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block3_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block3.clone(),
None,
))]));
block_verifier_router
.expect_request(zakura_consensus::Request::Commit(block3))
.await
.respond(block3_hash);
let chain_sync_result = chain_sync_task_handle.now_or_never();
assert!(
chain_sync_result.is_none(),
"unexpected error or panic in chain sync task: {chain_sync_result:?}",
);
Ok(())
}
#[tokio::test(start_paused = true)]
async fn incomplete_checkpoint_range_refreshes_tips_without_verifier_timeout(
) -> Result<(), crate::BoxError> {
let (
chain_sync,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
mock_chain_tip_sender,
) = setup_chain_sync_with_options(Height(4), STALLED_SERVICE_REQUEST_DELAY);
mock_chain_tip_sender.send_best_tip_height(Height(0));
let blocks: Vec<Arc<Block>> = vec![
zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES.zcash_deserialize_into()?,
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?,
zakura_test::vectors::BLOCK_MAINNET_2_BYTES.zcash_deserialize_into()?,
zakura_test::vectors::BLOCK_MAINNET_3_BYTES.zcash_deserialize_into()?,
zakura_test::vectors::BLOCK_MAINNET_4_BYTES.zcash_deserialize_into()?,
zakura_test::vectors::BLOCK_MAINNET_5_BYTES.zcash_deserialize_into()?,
];
let hashes: Vec<_> = blocks.iter().map(|block| block.hash()).collect();
let sync_task = tokio::spawn(chain_sync.sync());
state_service
.expect_request(zs::Request::KnownBlock(hashes[0]))
.await
.respond(zs::Response::KnownBlock(Some(zs::KnownBlock::BestChain)));
state_service
.expect_request(zs::Request::BlockLocator)
.await
.respond(zs::Response::BlockLocator(vec![hashes[0]]));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![hashes[0]],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
hashes[1], hashes[2], hashes[3],
]));
state_service
.expect_request(zs::Request::KnownBlock(hashes[1]))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(hashes[2]))
.await
.respond(zs::Response::KnownBlock(None));
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![hashes[0]],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic short-tip peer error")));
}
for hash in &hashes[1..=2] {
state_service
.expect_request(zs::Request::KnownBlock(*hash))
.await
.respond(zs::Response::KnownBlock(None));
}
for (hash, block) in hashes[1..=2].iter().zip(&blocks[1..=2]) {
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(*hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((block.clone(), None))]));
}
let mut pending_verifications = Vec::new();
let mut expected_verifications: HashSet<_> = hashes[1..=2].iter().copied().collect();
for _ in 0..2 {
pending_verifications.push(
block_verifier_router
.expect_request_that(|request| {
expected_verifications.remove(&request.block().hash())
})
.await,
);
}
assert!(expected_verifications.is_empty());
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![hashes[1]],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![hashes[2], hashes[3]]));
state_service
.expect_request(zs::Request::KnownBlock(hashes[3]))
.await
.respond(zs::Response::KnownBlock(None));
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![hashes[1]],
stop: None,
})
.await
.respond(Err(zn::BoxError::from(
"synthetic empty-extension peer error",
)));
}
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(hashes[3]).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
blocks[3].clone(),
None,
))]));
pending_verifications.push(
block_verifier_router
.expect_request_that(|request| request.block().hash() == hashes[3])
.await,
);
let refresh_requested_at = tokio::time::Instant::now();
let refreshed_locator = tokio::time::timeout(
sync::BLOCK_VERIFY_TIMEOUT,
state_service.expect_request(zs::Request::BlockLocator),
)
.await
.expect(
"syncer must discover the rest of the checkpoint range while its blocks are parked in the \
verifier: waiting out BLOCK_VERIFY_TIMEOUT and restarting discards the partial range",
);
refreshed_locator.respond(zs::Response::BlockLocator(vec![hashes[0]]));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![hashes[0]],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
hashes[1], hashes[2], hashes[3], hashes[4], hashes[5],
]));
state_service
.expect_request(zs::Request::KnownBlock(hashes[1]))
.await
.respond(zs::Response::KnownBlock(None));
for hash in &hashes[2..=4] {
state_service
.expect_request(zs::Request::KnownBlock(*hash))
.await
.respond(zs::Response::KnownBlock(None));
}
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![hashes[0]],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic refreshed-tip error")));
}
for hash in &hashes[1..=4] {
state_service
.expect_request(zs::Request::KnownBlock(*hash))
.await
.respond(zs::Response::KnownBlock(None));
}
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(hashes[4]).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
blocks[4].clone(),
None,
))]));
pending_verifications.push(
block_verifier_router
.expect_request_that(|request| request.block().hash() == hashes[4])
.await,
);
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![hashes[3]],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![hashes[4]]));
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![hashes[3]],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic final-tip error")));
}
for verification in pending_verifications {
verification.respond_with(|request| request.block().hash());
}
assert!(
refresh_requested_at.elapsed() < sync::BLOCK_VERIFY_TIMEOUT,
"the checkpoint range must complete without waiting out BLOCK_VERIFY_TIMEOUT, but it took {:?}",
refresh_requested_at.elapsed(),
);
tokio::task::yield_now().await;
assert!(
!sync_task.is_finished(),
"legacy sync should continue after the checkpoint range completes"
);
sync_task.abort();
Ok(())
}
#[tokio::test(start_paused = true)]
async fn sync_blocks_trailing_hashes_ok() -> Result<(), crate::BoxError> {
let (
chain_sync_future,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup();
let block0: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES.zcash_deserialize_into()?;
let block0_hash = block0.hash();
let block1: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?;
let block1_hash = block1.hash();
let block2: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_2_BYTES.zcash_deserialize_into()?;
let block2_hash = block2.hash();
let block3: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_3_BYTES.zcash_deserialize_into()?;
let block3_hash = block3.hash();
let block4: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_4_BYTES.zcash_deserialize_into()?;
let block4_hash = block4.hash();
let block5: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_5_BYTES.zcash_deserialize_into()?;
let block5_hash = block5.hash();
let chain_sync_task_handle = tokio::spawn(chain_sync_future);
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block0_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block0.clone(),
None,
))]));
block_verifier_router
.expect_request(zakura_consensus::Request::Commit(block0))
.await
.respond(block0_hash);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(Some(zs::KnownBlock::BestChain)));
state_service
.expect_request(zs::Request::BlockLocator)
.await
.respond(zs::Response::BlockLocator(vec![block0_hash]));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block0_hash],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
block1_hash, block2_hash, block3_hash, ]));
state_service
.expect_request(zs::Request::KnownBlock(block1_hash))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(block2_hash))
.await
.respond(zs::Response::KnownBlock(None));
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block0_hash],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test obtain tips error")));
}
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service
.expect_request(zs::Request::KnownBlock(block1_hash))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(block2_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block1_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block1.clone(),
None,
))]));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block2_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block2.clone(),
None,
))]));
let mut remaining_blocks: HashMap<block::Hash, Arc<Block>> =
[(block1_hash, block1), (block2_hash, block2)]
.iter()
.cloned()
.collect();
for _ in 1..=2 {
block_verifier_router
.expect_request_that(|req| remaining_blocks.remove(&req.block().hash()).is_some())
.await
.respond_with(|req| req.block().hash());
}
assert_eq!(
remaining_blocks,
HashMap::new(),
"expected all non-tip blocks to be verified by obtain tips"
);
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block1_hash],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
block2_hash, block3_hash, block4_hash,
block5_hash, ]));
for hash in [block3_hash, block4_hash] {
state_service
.expect_request(zs::Request::KnownBlock(hash))
.await
.respond(zs::Response::KnownBlock(None));
}
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block1_hash],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test extend tips error")));
}
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block3_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block3.clone(),
None,
))]));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block4_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block4.clone(),
None,
))]));
let mut remaining_blocks: HashMap<block::Hash, Arc<Block>> =
[(block3_hash, block3), (block4_hash, block4)]
.iter()
.cloned()
.collect();
for _ in 3..=4 {
block_verifier_router
.expect_request_that(|req| remaining_blocks.remove(&req.block().hash()).is_some())
.await
.respond_with(|req| req.block().hash());
}
assert_eq!(
remaining_blocks,
HashMap::new(),
"expected all non-tip blocks to be verified by extend tips"
);
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
let chain_sync_result = chain_sync_task_handle.now_or_never();
assert!(
chain_sync_result.is_none(),
"unexpected error or panic in chain sync task: {chain_sync_result:?}",
);
Ok(())
}
#[tokio::test]
async fn sync_block_lookahead_drop() -> Result<(), crate::BoxError> {
let (
chain_sync_future,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup();
let block0: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES.zcash_deserialize_into()?;
let block0_hash = block0.hash();
let block982k: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_982681_BYTES.zcash_deserialize_into()?;
let chain_sync_task_handle = tokio::spawn(chain_sync_future);
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block0_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block982k.clone(),
None,
))]));
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
let chain_sync_result = chain_sync_task_handle.now_or_never();
assert!(
chain_sync_result.is_none(),
"unexpected error or panic in chain sync task: {chain_sync_result:?}",
);
Ok(())
}
#[tokio::test]
async fn sync_block_too_high_obtain_tips() -> Result<(), crate::BoxError> {
let (
chain_sync_future,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup();
let block0: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES.zcash_deserialize_into()?;
let block0_hash = block0.hash();
let block1: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?;
let block1_hash = block1.hash();
let block2: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_2_BYTES.zcash_deserialize_into()?;
let block2_hash = block2.hash();
let block3: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_3_BYTES.zcash_deserialize_into()?;
let block3_hash = block3.hash();
let block982k: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_982681_BYTES.zcash_deserialize_into()?;
let block982k_hash = block982k.hash();
let chain_sync_task_handle = tokio::spawn(chain_sync_future);
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block0_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block0.clone(),
None,
))]));
block_verifier_router
.expect_request(zakura_consensus::Request::Commit(block0))
.await
.respond(block0_hash);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(Some(zs::KnownBlock::BestChain)));
state_service
.expect_request(zs::Request::BlockLocator)
.await
.respond(zs::Response::BlockLocator(vec![block0_hash]));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block0_hash],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
block982k_hash,
block1_hash, block2_hash, block3_hash, ]));
state_service
.expect_request(zs::Request::KnownBlock(block982k_hash))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(block1_hash))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(block2_hash))
.await
.respond(zs::Response::KnownBlock(None));
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block0_hash],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test obtain tips error")));
}
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service
.expect_request(zs::Request::KnownBlock(block982k_hash))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(block1_hash))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(block2_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(
iter::once(block982k_hash).collect(),
))
.await
.respond(zn::Response::Blocks(vec![Available((
block982k.clone(),
None,
))]));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block1_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block1.clone(),
None,
))]));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block2_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block2.clone(),
None,
))]));
let chain_sync_result = chain_sync_task_handle.now_or_never();
assert!(
chain_sync_result.is_none(),
"unexpected error or panic in chain sync task: {chain_sync_result:?}",
);
Ok(())
}
#[tokio::test(start_paused = true)]
async fn sync_block_too_high_extend_tips() -> Result<(), crate::BoxError> {
let (
chain_sync_future,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup();
let block0: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES.zcash_deserialize_into()?;
let block0_hash = block0.hash();
let block1: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?;
let block1_hash = block1.hash();
let block2: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_2_BYTES.zcash_deserialize_into()?;
let block2_hash = block2.hash();
let block3: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_3_BYTES.zcash_deserialize_into()?;
let block3_hash = block3.hash();
let block4: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_4_BYTES.zcash_deserialize_into()?;
let block4_hash = block4.hash();
let block5: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_5_BYTES.zcash_deserialize_into()?;
let block5_hash = block5.hash();
let block982k: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_982681_BYTES.zcash_deserialize_into()?;
let block982k_hash = block982k.hash();
let chain_sync_task_handle = tokio::spawn(chain_sync_future);
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block0_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block0.clone(),
None,
))]));
block_verifier_router
.expect_request(zakura_consensus::Request::Commit(block0))
.await
.respond(block0_hash);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service
.expect_request(zs::Request::KnownBlock(block0_hash))
.await
.respond(zs::Response::KnownBlock(Some(zs::KnownBlock::BestChain)));
state_service
.expect_request(zs::Request::BlockLocator)
.await
.respond(zs::Response::BlockLocator(vec![block0_hash]));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block0_hash],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
block1_hash, block2_hash, block3_hash, ]));
state_service
.expect_request(zs::Request::KnownBlock(block1_hash))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(block2_hash))
.await
.respond(zs::Response::KnownBlock(None));
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block0_hash],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test obtain tips error")));
}
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service
.expect_request(zs::Request::KnownBlock(block1_hash))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(block2_hash))
.await
.respond(zs::Response::KnownBlock(None));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block1_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block1.clone(),
None,
))]));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block2_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block2.clone(),
None,
))]));
let mut remaining_blocks: HashMap<block::Hash, Arc<Block>> =
[(block1_hash, block1), (block2_hash, block2)]
.iter()
.cloned()
.collect();
for _ in 1..=2 {
block_verifier_router
.expect_request_that(|req| remaining_blocks.remove(&req.block().hash()).is_some())
.await
.respond_with(|req| req.block().hash());
}
assert_eq!(
remaining_blocks,
HashMap::new(),
"expected all non-tip blocks to be verified by obtain tips"
);
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block1_hash],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
block2_hash, block3_hash, block4_hash,
block982k_hash,
block5_hash, ]));
for hash in [block3_hash, block4_hash, block982k_hash] {
state_service
.expect_request(zs::Request::KnownBlock(hash))
.await
.respond(zs::Response::KnownBlock(None));
}
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block1_hash],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test extend tips error")));
}
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block3_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block3.clone(),
None,
))]));
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block4_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
block4.clone(),
None,
))]));
peer_set
.expect_request(zn::Request::BlocksByHash(
iter::once(block982k_hash).collect(),
))
.await
.respond(zn::Response::Blocks(vec![Available((
block982k.clone(),
None,
))]));
let chain_sync_result = chain_sync_task_handle.now_or_never();
assert!(
chain_sync_result.is_none(),
"unexpected error or panic in chain sync task: {chain_sync_result:?}",
);
Ok(())
}
#[tokio::test]
async fn should_restart_sync_returns_false() {
let commit_error = zs::CommitBlockError::Duplicate {
hash_or_height: None,
location: zakura_state::KnownBlock::BestChain,
};
let verify_block_error = VerifyBlockError::Commit(commit_error);
let router_error = RouterError::Block {
source: Box::new(verify_block_error),
};
let err = BlockDownloadVerifyError::Invalid {
error: router_error,
height: block::Height(42),
hash: block::Hash::from([0xAA; 32]),
advertiser_addr: None,
};
let restart = ChainSync::<
MockService<zn::Request, zn::Response, PanicAssertion>,
MockService<zs::Request, zs::Response, PanicAssertion>,
MockService<zakura_consensus::Request, block::Hash, PanicAssertion>,
MockChainTip,
>::should_restart_sync(&err, false);
assert!(
!restart,
"duplicate commit block errors should NOT trigger sync restart"
);
}
#[test]
fn invalid_ancestor_does_not_score_descendant_advertiser() {
let (
mut chain_sync,
_sync_status,
_block_verifier_router,
_peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let (misbehavior_tx, mut misbehavior_rx) = tokio::sync::mpsc::channel(2);
chain_sync.misbehavior_sender = misbehavior_tx;
let parent_advertiser: PeerSocketAddr = "127.0.0.1:8233".parse().unwrap();
let child_advertiser: PeerSocketAddr = "127.0.0.2:8233".parse().unwrap();
let parent_hash = block::Hash([0xA1; 32]);
let child_hash = block::Hash([0xB2; 32]);
let parent_error = zs::CommitBlockError::ValidateContextError(Box::new(
zs::ValidateContextError::InvalidBlockCommitment(
zakura_chain::block::CommitmentError::InvalidChainHistoryActivationReserved {
actual: [1; 32],
},
),
));
let parent_error = RouterError::Block {
source: Box::new(VerifyBlockError::Commit(parent_error)),
};
let parent_response = BlockDownloadVerifyError::Invalid {
error: parent_error,
height: Height(42),
hash: parent_hash,
advertiser_addr: Some(parent_advertiser),
};
assert!(chain_sync
.handle_block_response(Err(parent_response))
.is_err());
assert_eq!(
misbehavior_rx.try_recv(),
Ok((parent_advertiser, 100)),
"the peer that served the invalid block must be scored"
);
let child_error = zs::CommitBlockError::ValidateContextError(Box::new(
zs::ValidateContextError::InvalidAncestorBlock(parent_hash),
));
let child_error = RouterError::Block {
source: Box::new(VerifyBlockError::Commit(child_error)),
};
let child_response = BlockDownloadVerifyError::Invalid {
error: child_error,
height: Height(43),
hash: child_hash,
advertiser_addr: Some(child_advertiser),
};
assert!(chain_sync
.handle_block_response(Err(child_response))
.is_err());
assert!(
matches!(
misbehavior_rx.try_recv(),
Err(tokio::sync::mpsc::error::TryRecvError::Empty)
),
"the different peer that served the descendant must not be scored"
);
}
#[test]
fn far_ahead_block_does_not_score_serving_peer() {
let (
mut chain_sync,
_sync_status,
_block_verifier_router,
_peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let (misbehavior_tx, mut misbehavior_rx) = tokio::sync::mpsc::channel(2);
chain_sync.misbehavior_sender = misbehavior_tx;
let serving_peer: PeerSocketAddr = "127.0.0.1:8233".parse().unwrap();
let far_ahead = BlockDownloadVerifyError::AboveLookaheadHeightLimit {
height: Height(60_000),
hash: block::Hash([0xC1; 32]),
advertiser_addr: Some(serving_peer),
};
assert!(
chain_sync.handle_block_response(Err(far_ahead)).is_ok(),
"a far-ahead block is dropped without restarting sync"
);
assert!(
matches!(
misbehavior_rx.try_recv(),
Err(tokio::sync::mpsc::error::TryRecvError::Empty)
),
"the peer that served a far-ahead block we requested must not be scored"
);
let invalid_height = BlockDownloadVerifyError::InvalidHeight {
hash: block::Hash([0xC2; 32]),
advertiser_addr: Some(serving_peer),
};
assert!(chain_sync
.handle_block_response(Err(invalid_height))
.is_ok());
assert_eq!(
misbehavior_rx.try_recv(),
Ok((serving_peer, 100)),
"a block with no valid height is still the serving peer's fault"
);
}
#[tokio::test]
async fn request_genesis_accepts_duplicate_finalized_genesis() -> Result<(), crate::BoxError> {
let block0: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES.zcash_deserialize_into()?;
let block0_hash = block0.hash();
let state_requests = Arc::new(AtomicUsize::new(0));
let state_requests_in_service = Arc::clone(&state_requests);
let state_service = tower::service_fn(move |request| {
state_requests_in_service.fetch_add(1, Ordering::SeqCst);
async move {
assert_eq!(request, zs::Request::KnownBlock(block0_hash));
Ok::<_, crate::BoxError>(zs::Response::KnownBlock(None))
}
});
let peer_requests = Arc::new(AtomicUsize::new(0));
let peer_requests_in_service = Arc::clone(&peer_requests);
let peer_block = block0.clone();
let peer_set = tower::service_fn(move |request| {
peer_requests_in_service.fetch_add(1, Ordering::SeqCst);
let peer_block = peer_block.clone();
async move {
assert_eq!(
request,
zn::Request::BlocksByHash(iter::once(block0_hash).collect())
);
Ok::<_, crate::BoxError>(zn::Response::Blocks(vec![Available((peer_block, None))]))
}
});
let verifier_requests = Arc::new(AtomicUsize::new(0));
let verifier_requests_in_service = Arc::clone(&verifier_requests);
let verifier_service = tower::service_fn(move |request| {
verifier_requests_in_service.fetch_add(1, Ordering::SeqCst);
async move {
let zakura_consensus::Request::Commit(block) = request else {
unreachable!("no other verifier request is allowed")
};
assert_eq!(block.hash(), block0_hash);
let duplicate = zs::CommitBlockError::Duplicate {
hash_or_height: None,
location: zs::KnownBlock::Finalized,
};
let duplicate = zs::CommitCheckpointVerifiedError::from(duplicate);
let router_error = RouterError::Checkpoint {
source: Box::new(VerifyCheckpointError::CommitCheckpointVerified(Box::new(
duplicate,
))),
};
Err::<block::Hash, crate::BoxError>(Box::new(router_error) as crate::BoxError)
}
});
let consensus_config = ConsensusConfig::default();
let state_config = StateConfig::ephemeral();
let config = ZakuradConfig {
consensus: consensus_config,
state: state_config,
..Default::default()
};
let (mock_chain_tip, _mock_chain_tip_sender) = MockChainTip::new();
let (misbehavior_tx, _misbehavior_rx) = tokio::sync::mpsc::channel(1);
let (mut chain_sync, _sync_status) = ChainSync::new(
&config,
Height(0),
peer_set,
verifier_service,
state_service,
mock_chain_tip,
misbehavior_tx,
);
tokio::time::timeout(Duration::from_secs(2), chain_sync.request_genesis())
.await
.expect("duplicate finalized genesis should not sleep and retry")
.expect("duplicate finalized genesis is accepted");
assert_eq!(state_requests.load(Ordering::SeqCst), 1);
assert_eq!(peer_requests.load(Ordering::SeqCst), 1);
assert_eq!(verifier_requests.load(Ordering::SeqCst), 1);
Ok(())
}
#[test]
fn duplicate_finalized_checkpoint_block_does_not_restart_sync() -> Result<(), crate::BoxError> {
let block1: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?;
let block1_hash = block1.hash();
let duplicate = zs::CommitBlockError::Duplicate {
hash_or_height: None,
location: zs::KnownBlock::Finalized,
};
let duplicate = zs::CommitCheckpointVerifiedError::from(duplicate);
let router_error = RouterError::Checkpoint {
source: Box::new(VerifyCheckpointError::CommitCheckpointVerified(Box::new(
duplicate,
))),
};
let err = BlockDownloadVerifyError::Invalid {
error: router_error,
height: Height(1),
hash: block1_hash,
advertiser_addr: None,
};
let restart = TestChainSync::should_restart_sync(&err, false);
assert!(
!restart,
"duplicate finalized checkpoint blocks are stale in-flight work, not sync restarts"
);
Ok(())
}
#[tokio::test]
async fn above_lookahead_does_not_restart_sync() {
let err = BlockDownloadVerifyError::AboveLookaheadHeightLimit {
height: block::Height(60_000),
hash: block::Hash::from([0xBB; 32]),
advertiser_addr: None,
};
let restart = ChainSync::<
MockService<zn::Request, zn::Response, PanicAssertion>,
MockService<zs::Request, zs::Response, PanicAssertion>,
MockService<zakura_consensus::Request, block::Hash, PanicAssertion>,
MockChainTip,
>::should_restart_sync(&err, false);
assert!(
!restart,
"AboveLookaheadHeightLimit should NOT trigger sync restart (GHSA-gvjc-3w7c-92jx fix)"
);
}
#[tokio::test]
async fn above_lookahead_has_peer_attribution() {
let addr: PeerSocketAddr = "127.0.0.1:8233".parse().unwrap();
let err = BlockDownloadVerifyError::AboveLookaheadHeightLimit {
height: block::Height(60_000),
hash: block::Hash::from([0xCC; 32]),
advertiser_addr: Some(addr),
};
assert_eq!(
err.advertiser_addr(),
Some(addr),
"AboveLookaheadHeightLimit should carry advertiser_addr for drop logs \
(GHSA-gvjc-3w7c-92jx fix)"
);
}
#[tokio::test]
async fn both_height_limits_do_not_restart_sync() {
let below = BlockDownloadVerifyError::BehindTipHeightLimit {
height: block::Height(1),
hash: block::Hash::from([0xDD; 32]),
};
let above = BlockDownloadVerifyError::AboveLookaheadHeightLimit {
height: block::Height(60_000),
hash: block::Hash::from([0xEE; 32]),
advertiser_addr: None,
};
let restart_below = ChainSync::<
MockService<zn::Request, zn::Response, PanicAssertion>,
MockService<zs::Request, zs::Response, PanicAssertion>,
MockService<zakura_consensus::Request, block::Hash, PanicAssertion>,
MockChainTip,
>::should_restart_sync(&below, false);
let restart_above = ChainSync::<
MockService<zn::Request, zn::Response, PanicAssertion>,
MockService<zs::Request, zs::Response, PanicAssertion>,
MockService<zakura_consensus::Request, block::Hash, PanicAssertion>,
MockChainTip,
>::should_restart_sync(&above, false);
assert!(
!restart_below,
"BehindTipHeightLimit should NOT restart sync"
);
assert!(
!restart_above,
"AboveLookaheadHeightLimit should NOT restart sync (GHSA-gvjc-3w7c-92jx fix)"
);
}
#[tokio::test]
async fn invalid_height_does_not_restart_sync() {
let addr: PeerSocketAddr = "127.0.0.1:8233".parse().unwrap();
let err = BlockDownloadVerifyError::InvalidHeight {
hash: block::Hash::from([0xFF; 32]),
advertiser_addr: Some(addr),
};
let restart = ChainSync::<
MockService<zn::Request, zn::Response, PanicAssertion>,
MockService<zs::Request, zs::Response, PanicAssertion>,
MockService<zakura_consensus::Request, block::Hash, PanicAssertion>,
MockChainTip,
>::should_restart_sync(&err, false);
assert!(
!restart,
"InvalidHeight should NOT trigger sync restart (GHSA-rj6c-83wx-jxf2 fix)"
);
assert_eq!(
err.advertiser_addr(),
Some(addr),
"InvalidHeight should carry advertiser_addr for peer scoring"
);
}
#[tokio::test]
async fn not_found_download_requeues_missing_block() -> Result<(), crate::BoxError> {
let (
mut chain_sync,
_sync_status,
mut block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let block1: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?;
let block1_hash = block1.hash();
let error = BlockDownloadVerifyError::DownloadFailed {
error: not_found_block_error(block1_hash),
hash: block1_hash,
};
let requeue = tokio::spawn(async move {
chain_sync
.handle_block_response_with_missing_retry(Err(error))
.await
});
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block1_hash).collect()))
.await
.respond(Err(not_found_block_error(block1_hash)));
requeue
.await
.expect("missing block retry task should not panic")?;
block_verifier_router.expect_no_requests().await;
Ok(())
}
#[tokio::test]
async fn not_found_download_restarts_after_queue_retry_limit() {
let (
mut chain_sync,
_sync_status,
_block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let block_hash = block::Hash::from([0xAB; 32]);
chain_sync
.missing_block_retry_counts
.insert(block_hash, sync::MISSING_BLOCK_DOWNLOAD_RETRY_LIMIT);
let error = BlockDownloadVerifyError::DownloadFailed {
error: not_found_block_error(block_hash),
hash: block_hash,
};
let result = chain_sync
.handle_block_response_with_missing_retry(Err(error))
.await;
assert!(
result.is_err(),
"notfound downloads should restart sync after queue retry limit"
);
peer_set.expect_no_requests().await;
}
#[tokio::test]
async fn not_found_download_restarts_sync() {
let block_hash = block::Hash::from([0xCD; 32]);
let err = BlockDownloadVerifyError::DownloadFailed {
error: not_found_block_error(block_hash),
hash: block_hash,
};
let restart = TestChainSync::should_restart_sync(&err, false);
assert!(
restart,
"notfound block downloads should restart sync after queue retries"
);
}
#[tokio::test]
async fn transient_download_failure_requeues_and_clears_on_success() -> Result<(), crate::BoxError>
{
let (
mut chain_sync,
_sync_status,
mut block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let block: Arc<Block> = zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?;
let block_hash = block.hash();
let error = BlockDownloadVerifyError::DownloadFailed {
error: client_dropped_error(),
hash: block_hash,
};
chain_sync
.handle_block_response_with_missing_retry(Err(error))
.await
.expect("a transient block failure within budget should preserve the round");
assert_eq!(
chain_sync.transient_block_retry_counts.get(&block_hash),
Some(&1),
"the transient block retry should consume one retry"
);
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((block.clone(), None))]));
block_verifier_router
.expect_request(zakura_consensus::Request::Commit(block))
.await
.respond(block_hash);
let response = chain_sync
.downloads
.next()
.await
.expect("the replacement download should complete");
chain_sync
.handle_block_response_with_missing_retry(response)
.await?;
assert!(
!chain_sync
.transient_block_retry_counts
.contains_key(&block_hash),
"a successful replacement should clear its transient retry count"
);
assert_eq!(chain_sync.downloads.in_flight(), 0);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
Ok(())
}
#[tokio::test]
async fn transient_download_failure_restarts_after_retry_limit() {
let (
mut chain_sync,
_sync_status,
_block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let block_hash = block::Hash::from([0xCF; 32]);
chain_sync
.transient_block_retry_counts
.insert(block_hash, sync::TRANSIENT_BLOCK_DOWNLOAD_RETRY_LIMIT);
let error = BlockDownloadVerifyError::DownloadFailed {
error: client_dropped_error(),
hash: block_hash,
};
let result = chain_sync
.handle_block_response_with_missing_retry(Err(error))
.await;
assert!(
result.is_err(),
"a transient block failure should restart sync after its retry budget"
);
assert!(
!chain_sync
.transient_block_retry_counts
.contains_key(&block_hash),
"an exhausted transient retry budget should be cleared"
);
peer_set.expect_no_requests().await;
}
#[tokio::test(start_paused = true)]
async fn transient_download_failure_has_eight_peer_request_ceiling() -> Result<(), crate::BoxError>
{
let (
mut chain_sync,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let prime_block: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?;
let prime_hash = prime_block.hash();
for _ in 0..sync::BLOCK_DOWNLOAD_HEDGE_MIN_DATA_POINTS {
chain_sync.downloads.download_and_verify(prime_hash).await?;
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(prime_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((
prime_block.clone(),
None,
))]));
block_verifier_router
.expect_request_that(|request| request.block().hash() == prime_hash)
.await
.respond(prime_hash);
assert!(
chain_sync
.downloads
.next()
.await
.expect("priming download should complete")
.is_ok(),
"priming download should succeed"
);
}
tokio::time::advance(2 * sync::SYNC_RESTART_DELAY).await;
let failed_hash = block::Hash::from([0xCF; 32]);
let sync_round = chain_sync.sync_round(iter::once(failed_hash).collect(), None);
let fail_peer_requests = async {
let request = zn::Request::BlocksByHash(iter::once(failed_hash).collect());
let queue_attempts = sync::TRANSIENT_BLOCK_DOWNLOAD_RETRY_LIMIT + 1;
let mut peer_request_count = 0;
for _ in 0..queue_attempts {
let original = peer_set.expect_request(request.clone()).await;
let hedge = peer_set.expect_request(request.clone()).await;
peer_request_count += 2;
original.respond(Err(client_dropped_error()));
hedge.respond(Err(client_dropped_error()));
}
peer_request_count
};
let (result, peer_request_count) = tokio::join!(sync_round, fail_peer_requests);
assert!(
result.is_err(),
"the exhausted transient retry budget should restart the sync round"
);
assert_eq!(
peer_request_count,
sync::MAX_TRANSIENT_BLOCK_PEER_REQUESTS_PER_SYNC_ROUND
);
assert!(chain_sync.transient_block_retry_counts.is_empty());
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
Ok(())
}
#[tokio::test]
async fn transient_download_failure_preserves_sync_round() -> Result<(), crate::BoxError> {
let (
mut chain_sync,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let retried_block: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?;
let retried_hash = retried_block.hash();
let unrelated_block: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_2_BYTES.zcash_deserialize_into()?;
let unrelated_hash = unrelated_block.hash();
let reserve = [retried_hash, unrelated_hash].into_iter().collect();
let sync_round = chain_sync.sync_round(reserve, None);
let drive_services = async {
let mut failed_attempt_sent = false;
let mut unrelated_response_sent = false;
while !failed_attempt_sent || !unrelated_response_sent {
let response = peer_set
.expect_request_that(|request| match request {
zn::Request::BlocksByHash(hashes) => {
hashes == &HashSet::from([retried_hash])
|| hashes == &HashSet::from([unrelated_hash])
}
_ => false,
})
.await;
let requested_hash = match response.request() {
zn::Request::BlocksByHash(hashes) => *hashes
.iter()
.next()
.expect("single-hash block request is nonempty"),
_ => unreachable!("request matcher accepts only block requests"),
};
if requested_hash == retried_hash {
assert!(
!failed_attempt_sent,
"the failed block should use one request"
);
failed_attempt_sent = true;
response.respond(Err(client_dropped_error()));
} else {
assert!(
!unrelated_response_sent,
"the unrelated block should be requested once"
);
unrelated_response_sent = true;
response.respond(zn::Response::Blocks(vec![Available((
unrelated_block.clone(),
None,
))]));
}
}
block_verifier_router
.expect_request_that(|request| request.block().hash() == unrelated_hash)
.await
.respond(unrelated_hash);
peer_set
.expect_request(zn::Request::BlocksByHash(
iter::once(retried_hash).collect(),
))
.await
.respond(zn::Response::Blocks(vec![Available((
retried_block.clone(),
None,
))]));
block_verifier_router
.expect_request_that(|request| request.block().hash() == retried_hash)
.await
.respond(retried_hash);
};
let (result, ()) = tokio::join!(sync_round, drive_services);
result.expect("replacement download should complete the original sync round");
assert!(
chain_sync.transient_block_retry_counts.is_empty(),
"a successful replacement should clear its transient retry state"
);
assert_eq!(chain_sync.downloads.in_flight(), 0);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
Ok(())
}
#[tokio::test]
async fn build_extend_discovers_hashes_without_dispatching() -> Result<(), crate::BoxError> {
let (
_chain_sync,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let block1: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into()?;
let block1_hash = block1.hash();
let block2: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_2_BYTES.zcash_deserialize_into()?;
let block2_hash = block2.hash();
let block3: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_3_BYTES.zcash_deserialize_into()?;
let block3_hash = block3.hash();
let block4: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_4_BYTES.zcash_deserialize_into()?;
let block4_hash = block4.hash();
let block5: Arc<Block> =
zakura_test::vectors::BLOCK_MAINNET_5_BYTES.zcash_deserialize_into()?;
let block5_hash = block5.hash();
let tip = sync::CheckedTip {
tip: block1_hash,
expected_next: block2_hash,
};
let tips: HashSet<_> = iter::once(tip).collect();
let tip_network = Timeout::new(peer_set.clone(), sync::TIPS_RESPONSE_TIMEOUT);
let extend_handle = tokio::spawn(TestChainSync::build_extend(
tip_network,
state_service.clone(),
tips,
));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block1_hash],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
block2_hash, block3_hash,
block4_hash,
block5_hash, ]));
for hash in [block3_hash, block4_hash] {
state_service
.expect_request(zs::Request::KnownBlock(hash))
.await
.respond(zs::Response::KnownBlock(None));
}
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![block1_hash],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test extend tips error")));
}
let (download_set, prospective_tips, discovered) = extend_handle
.await
.expect("build_extend task should not panic")?;
assert_eq!(
download_set.into_iter().collect::<Vec<_>>(),
vec![block3_hash, block4_hash],
"build_extend should discover the inner hashes in response order",
);
assert_eq!(
discovered, 2,
"discovered count should match the download set length",
);
let expected_tip = sync::CheckedTip {
tip: block3_hash,
expected_next: block4_hash,
};
assert_eq!(
prospective_tips,
iter::once(expected_tip).collect(),
"build_extend should return the next prospective tip",
);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
Ok(())
}
#[tokio::test]
async fn obtain_tips_ignores_known_hash_after_first_unknown() -> Result<(), crate::BoxError> {
let (
mut chain_sync,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let locator = block::Hash::from([0x01; 32]);
let unknown = block::Hash::from([0x02; 32]);
let known_not_in_locator = block::Hash::from([0x03; 32]);
let trailing = block::Hash::from([0x04; 32]);
let respond_to_requests = async {
state_service
.expect_request(zs::Request::BlockLocator)
.await
.respond(zs::Response::BlockLocator(vec![locator]));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![locator],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
unknown,
known_not_in_locator,
trailing,
]));
state_service
.expect_request(zs::Request::KnownBlock(unknown))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(known_not_in_locator))
.await
.respond(zs::Response::KnownBlock(Some(zs::KnownBlock::BestChain)));
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![locator],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test obtain tips error")));
}
Ok::<_, crate::BoxError>(())
};
let (extra_hashes, responded) =
futures::join!(chain_sync.obtain_tips(false), respond_to_requests);
responded?;
let extra_hashes = extra_hashes?;
assert!(extra_hashes.is_empty());
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
Ok(())
}
#[tokio::test]
async fn build_extend_ignores_malformed_find_blocks_responses() -> Result<(), crate::BoxError> {
let (
_chain_sync,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let tip = block::Hash::from([0x10; 32]);
let expected_next = block::Hash::from([0x11; 32]);
let unknown = block::Hash::from([0x12; 32]);
let known_suffix = block::Hash::from([0x13; 32]);
let trailing = block::Hash::from([0x14; 32]);
let random = block::Hash::from([0x15; 32]);
let tips = HashSet::from([sync::CheckedTip { tip, expected_next }]);
let extend_handle = tokio::spawn(TestChainSync::build_extend(
Timeout::new(peer_set.clone(), sync::TIPS_RESPONSE_TIMEOUT),
state_service.clone(),
tips,
));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![tip],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
expected_next,
unknown,
known_suffix,
trailing,
]));
state_service
.expect_request(zs::Request::KnownBlock(unknown))
.await
.respond(zs::Response::KnownBlock(None));
state_service
.expect_request(zs::Request::KnownBlock(known_suffix))
.await
.respond(zs::Response::KnownBlock(Some(zs::KnownBlock::BestChain)));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![tip],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
expected_next,
unknown,
unknown,
trailing,
]));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![tip],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![random]));
let (download_set, prospective_tips, discovered) = extend_handle
.await
.expect("build_extend task should not panic")?;
assert!(download_set.is_empty());
assert!(prospective_tips.is_empty());
assert_eq!(discovered, 0);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
Ok(())
}
#[tokio::test]
async fn build_extend_rejects_oversized_response_before_state_queries(
) -> Result<(), crate::BoxError> {
let (
_chain_sync,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let tip = block::Hash::from([0x20; 32]);
let expected_next = block::Hash::from([0x21; 32]);
let trailing = block::Hash::from([0x22; 32]);
let unknown_hashes = (0..=sync::MAX_TIPS_RESPONSE_HASH_COUNT).map(|index| {
let index = u64::try_from(index).expect("test hash count fits in u64");
let mut bytes = [0; 32];
bytes[..8].copy_from_slice(&index.to_le_bytes());
block::Hash::from(bytes)
});
let response_hashes = std::iter::once(expected_next)
.chain(unknown_hashes)
.chain(std::iter::once(trailing))
.collect();
let tips = HashSet::from([sync::CheckedTip { tip, expected_next }]);
let extend_handle = tokio::spawn(TestChainSync::build_extend(
Timeout::new(peer_set.clone(), sync::TIPS_RESPONSE_TIMEOUT),
state_service.clone(),
tips,
));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![tip],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(response_hashes));
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![tip],
stop: None,
})
.await
.respond(Err(zn::BoxError::from(
"synthetic test oversized response error",
)));
}
let (download_set, prospective_tips, discovered) = extend_handle
.await
.expect("build_extend task should not panic")?;
assert!(download_set.is_empty());
assert!(prospective_tips.is_empty());
assert_eq!(discovered, 0);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
Ok(())
}
#[tokio::test]
async fn build_extend_rejects_locator_echo_before_state_queries() -> Result<(), crate::BoxError> {
let (
_chain_sync,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let tip = block::Hash::from([0x20; 32]);
let expected_next = block::Hash::from([0x21; 32]);
let unknown = block::Hash::from([0x22; 32]);
let trailing = block::Hash::from([0x23; 32]);
let tips = HashSet::from([sync::CheckedTip { tip, expected_next }]);
let extend_handle = tokio::spawn(TestChainSync::build_extend(
Timeout::new(peer_set.clone(), sync::TIPS_RESPONSE_TIMEOUT),
state_service.clone(),
tips,
));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![tip],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
expected_next,
unknown,
tip,
trailing,
]));
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![tip],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test locator echo error")));
}
let (download_set, prospective_tips, discovered) = extend_handle
.await
.expect("build_extend task should not panic")?;
assert!(download_set.is_empty());
assert!(prospective_tips.is_empty());
assert_eq!(discovered, 0);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
Ok(())
}
#[tokio::test]
async fn build_extend_ignores_known_trailing_find_blocks_hash() -> Result<(), crate::BoxError> {
let (
_chain_sync,
_sync_status,
mut block_verifier_router,
mut peer_set,
mut state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let tip = block::Hash::from([0x20; 32]);
let expected_next = block::Hash::from([0x21; 32]);
let unknown_a = block::Hash::from([0x22; 32]);
let unknown_b = block::Hash::from([0x23; 32]);
let tips = HashSet::from([sync::CheckedTip { tip, expected_next }]);
let extend_handle = tokio::spawn(TestChainSync::build_extend(
Timeout::new(peer_set.clone(), sync::TIPS_RESPONSE_TIMEOUT),
state_service.clone(),
tips,
));
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![tip],
stop: None,
})
.await
.respond(zn::Response::BlockHashes(vec![
expected_next,
unknown_a,
unknown_b,
tip, ]));
for hash in [unknown_a, unknown_b] {
state_service
.expect_request(zs::Request::KnownBlock(hash))
.await
.respond(zs::Response::KnownBlock(None));
}
for _ in 0..(sync::FANOUT - 1) {
peer_set
.expect_request(zn::Request::FindBlocks {
known_blocks: vec![tip],
stop: None,
})
.await
.respond(Err(zn::BoxError::from("synthetic test extend tips error")));
}
let (download_set, prospective_tips, discovered) = extend_handle
.await
.expect("build_extend task should not panic")?;
assert_eq!(
download_set.into_iter().collect::<Vec<_>>(),
vec![unknown_a, unknown_b],
"build_extend should keep valid inner hashes and discard the trailing known hash",
);
assert_eq!(discovered, 2);
assert_eq!(
prospective_tips,
HashSet::from([sync::CheckedTip {
tip: unknown_a,
expected_next: unknown_b,
}]),
);
peer_set.expect_no_requests().await;
block_verifier_router.expect_no_requests().await;
state_service.expect_no_requests().await;
Ok(())
}
#[tokio::test]
async fn registry_miss_schedules_backoff_retry() {
let (
mut chain_sync,
_sync_status,
_block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let block_hash = block::Hash::from([0xAB; 32]);
let error = BlockDownloadVerifyError::DownloadFailed {
error: not_found_registry_error(block_hash),
hash: block_hash,
};
let result = chain_sync
.handle_block_response_with_missing_retry(Err(error))
.await;
assert!(
result.is_ok(),
"a registry miss within budget should keep the round alive, not restart"
);
assert!(
chain_sync.registry_miss_retry.contains_key(&block_hash),
"the missing block should be scheduled for a backoff retry"
);
assert_eq!(
chain_sync.registry_miss_retry_counts.get(&block_hash),
Some(&1),
"the registry-miss retry budget should be consumed once",
);
peer_set.expect_no_requests().await;
}
#[tokio::test]
async fn registry_miss_restarts_after_retry_limit() {
let (
mut chain_sync,
_sync_status,
_block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let block_hash = block::Hash::from([0xCD; 32]);
chain_sync
.registry_miss_retry_counts
.insert(block_hash, sync::MISSING_BLOCK_REGISTRY_RETRY_LIMIT);
let error = BlockDownloadVerifyError::DownloadFailed {
error: not_found_registry_error(block_hash),
hash: block_hash,
};
let result = chain_sync
.handle_block_response_with_missing_retry(Err(error))
.await;
assert!(
result.is_err(),
"a registry miss should restart sync once the retry budget is exhausted"
);
assert!(
!chain_sync.registry_miss_retry.contains_key(&block_hash),
"exhausted retry schedule should be cleared"
);
assert!(
!chain_sync
.registry_miss_retry_counts
.contains_key(&block_hash),
"exhausted retry budget should be cleared"
);
peer_set.expect_no_requests().await;
}
#[tokio::test]
async fn registry_miss_schedules_multiple_blocks() {
let (
mut chain_sync,
_sync_status,
_block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let first_hash = block::Hash::from([0x11; 32]);
let second_hash = block::Hash::from([0x22; 32]);
for hash in [first_hash, second_hash] {
let error = BlockDownloadVerifyError::DownloadFailed {
error: not_found_registry_error(hash),
hash,
};
chain_sync
.handle_block_response_with_missing_retry(Err(error))
.await
.expect("a registry miss within budget should not restart");
}
assert!(
chain_sync.registry_miss_retry.contains_key(&first_hash)
&& chain_sync.registry_miss_retry.contains_key(&second_hash),
"both registry-missed blocks should stay scheduled for retry",
);
peer_set.expect_no_requests().await;
}
#[tokio::test]
async fn registry_miss_retry_clears_on_successful_block() {
let (
mut chain_sync,
_sync_status,
_block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let block_hash = block::Hash::from([0xAB; 32]);
chain_sync
.handle_block_response_with_missing_retry(Err(BlockDownloadVerifyError::DownloadFailed {
error: not_found_registry_error(block_hash),
hash: block_hash,
}))
.await
.expect("a registry miss within budget should not restart");
assert!(chain_sync.registry_miss_retry.contains_key(&block_hash));
chain_sync
.handle_block_response_with_missing_retry(Ok((Height(42), block_hash)))
.await
.expect("a successful response should not restart");
assert!(
!chain_sync.registry_miss_retry.contains_key(&block_hash),
"a successful block should clear its scheduled registry-miss retry"
);
assert!(
!chain_sync
.registry_miss_retry_counts
.contains_key(&block_hash),
"a successful block should clear its consumed registry-miss budget"
);
peer_set.expect_no_requests().await;
}
#[tokio::test]
async fn registry_miss_retry_clears_only_the_responded_block() {
let (
mut chain_sync,
_sync_status,
_block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let kept_hash = block::Hash::from([0x11; 32]);
let arrived_hash = block::Hash::from([0x22; 32]);
for hash in [kept_hash, arrived_hash] {
chain_sync
.handle_block_response_with_missing_retry(Err(
BlockDownloadVerifyError::DownloadFailed {
error: not_found_registry_error(hash),
hash,
},
))
.await
.expect("a registry miss within budget should not restart");
}
chain_sync
.handle_block_response_with_missing_retry(Ok((Height(7), arrived_hash)))
.await
.expect("a successful response should not restart");
assert!(
!chain_sync.registry_miss_retry.contains_key(&arrived_hash),
"the arrived block's retry should be cleared"
);
assert!(
chain_sync.registry_miss_retry.contains_key(&kept_hash),
"a different block's retry must not be cleared by an unrelated success"
);
peer_set.expect_no_requests().await;
}
#[tokio::test]
async fn registry_miss_retry_is_deferred_by_the_backoff() {
let (
mut chain_sync,
_sync_status,
_block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let block_hash = block::Hash::from([0xAB; 32]);
let before = tokio::time::Instant::now();
chain_sync
.handle_block_response_with_missing_retry(Err(BlockDownloadVerifyError::DownloadFailed {
error: not_found_registry_error(block_hash),
hash: block_hash,
}))
.await
.expect("a registry miss within budget should not restart");
let after = tokio::time::Instant::now();
let deadline = chain_sync
.registry_miss_retry
.get(&block_hash)
.copied()
.expect("the missing block should be scheduled");
assert!(
deadline >= before + sync::REGISTRY_MISS_RETRY_BACKOFF
&& deadline <= after + sync::REGISTRY_MISS_RETRY_BACKOFF,
"the retry deadline should be one backoff interval in the future, not immediate"
);
peer_set.expect_no_requests().await;
}
#[tokio::test]
async fn registry_miss_retry_accumulates_budget_for_the_same_block() {
let (
mut chain_sync,
_sync_status,
_block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let block_hash = block::Hash::from([0xAB; 32]);
let miss = || BlockDownloadVerifyError::DownloadFailed {
error: not_found_registry_error(block_hash),
hash: block_hash,
};
chain_sync
.handle_block_response_with_missing_retry(Err(miss()))
.await
.expect("a registry miss within budget should not restart");
let first_deadline = chain_sync
.registry_miss_retry
.get(&block_hash)
.copied()
.expect("the missing block should be scheduled");
chain_sync
.handle_block_response_with_missing_retry(Err(miss()))
.await
.expect("a second registry miss within budget should not restart");
let second_deadline = chain_sync
.registry_miss_retry
.get(&block_hash)
.copied()
.expect("the missing block should still be scheduled");
assert_eq!(
chain_sync.registry_miss_retry_counts.get(&block_hash),
Some(&2),
"each registry miss for the same block should consume one more retry from its budget"
);
assert!(
second_deadline >= first_deadline,
"each miss should re-arm the backoff deadline"
);
peer_set.expect_no_requests().await;
}
fn setup() -> (
// ChainSync
impl Future<Output = Result<(), Report>> + Send,
SyncStatus,
// BlockVerifierRouter
MockService<zakura_consensus::Request, block::Hash, PanicAssertion>,
// PeerSet
MockService<zakura_network::Request, zakura_network::Response, PanicAssertion>,
// StateService
MockService<zakura_state::Request, zakura_state::Response, PanicAssertion>,
MockChainTipSender,
) {
let (
chain_sync,
sync_status,
block_verifier_router,
peer_set,
state_service,
mock_chain_tip_sender,
) = setup_chain_sync();
let chain_sync_future = chain_sync.sync();
(
chain_sync_future,
sync_status,
block_verifier_router,
peer_set,
state_service,
mock_chain_tip_sender,
)
}
fn setup_chain_sync() -> (
TestChainSync,
SyncStatus,
MockService<zakura_consensus::Request, block::Hash, PanicAssertion>,
MockService<zakura_network::Request, zakura_network::Response, PanicAssertion>,
MockService<zakura_state::Request, zakura_state::Response, PanicAssertion>,
MockChainTipSender,
) {
setup_chain_sync_with_options(Height(0), MAX_SERVICE_REQUEST_DELAY)
}
fn setup_chain_sync_with_options(
max_checkpoint_height: Height,
max_service_request_delay: Duration,
) -> (
TestChainSync,
SyncStatus,
MockService<zakura_consensus::Request, block::Hash, PanicAssertion>,
MockService<zakura_network::Request, zakura_network::Response, PanicAssertion>,
MockService<zakura_state::Request, zakura_state::Response, PanicAssertion>,
MockChainTipSender,
) {
let _init_guard = zakura_test::init();
let consensus_config = ConsensusConfig::default();
let state_config = StateConfig::ephemeral();
let config = ZakuradConfig {
consensus: consensus_config,
state: state_config,
..Default::default()
};
let peer_set = MockService::build()
.with_max_request_delay(max_service_request_delay)
.for_unit_tests();
let block_verifier_router = MockService::build()
.with_max_request_delay(max_service_request_delay)
.for_unit_tests();
let state_service = MockService::build()
.with_max_request_delay(max_service_request_delay)
.for_unit_tests();
let (mock_chain_tip, mock_chain_tip_sender) = MockChainTip::new();
let (misbehavior_tx, _misbehavior_rx) = tokio::sync::mpsc::channel(1);
let (chain_sync, sync_status) = ChainSync::new(
&config,
max_checkpoint_height,
peer_set.clone(),
block_verifier_router.clone(),
state_service.clone(),
mock_chain_tip,
misbehavior_tx,
);
(
chain_sync,
sync_status,
block_verifier_router,
peer_set,
state_service,
mock_chain_tip_sender,
)
}
fn not_found_block_error(_hash: block::Hash) -> crate::BoxError {
zn::SharedPeerError::from(zn::PeerError::NotFoundResponse(Vec::new())).into()
}
fn not_found_registry_error(_hash: block::Hash) -> crate::BoxError {
zn::SharedPeerError::from(zn::PeerError::NotFoundRegistry(Vec::new())).into()
}
fn client_dropped_error() -> crate::BoxError {
zn::SharedPeerError::from(zn::PeerError::ClientDropped).into()
}
#[test]
fn debug_skip_regtest_genesis_self_seed_defaults_off_and_is_opt_in() {
use crate::components::sync::Config;
assert!(!Config::default().debug_skip_regtest_genesis_self_seed);
assert_eq!(
Config::default().debug_blocksync_throughput_target_height,
None
);
let config: Config = toml::from_str("debug_skip_regtest_genesis_self_seed = true")
.expect("sync config with the genesis-bootstrap flag parses");
assert!(config.debug_skip_regtest_genesis_self_seed);
let config: Config = toml::from_str("debug_blocksync_throughput_target_height = 100")
.expect("sync config with the block-sync throughput flag parses");
assert_eq!(config.debug_blocksync_throughput_target_height, Some(100));
let serialized = toml::to_string(&Config::default()).expect("sync config serializes");
assert!(
!serialized.contains("debug_skip_regtest_genesis_self_seed"),
"debug bootstrap flag must not appear in generated config output"
);
assert!(
!serialized.contains("debug_blocksync_throughput_target_height"),
"debug block-sync throughput flag must not appear in generated config output"
);
}
#[tokio::test]
async fn empty_block_response_is_retryable_download_failure() {
let _init_guard = zakura_test::init();
let mut peer_set = MockService::build().for_unit_tests::<zn::Request, zn::Response, _>();
let verifier =
MockService::build().for_unit_tests::<zakura_consensus::Request, block::Hash, _>();
let (chain_tip, _chain_tip_sender) = MockChainTip::new();
let (past_lookahead_limit_sender, _past_lookahead_limit_receiver) =
tokio::sync::watch::channel(false);
let mut downloads = Downloads::new(
peer_set.clone(),
verifier,
chain_tip,
past_lookahead_limit_sender,
sync::MIN_CONCURRENCY_LIMIT,
Height(0),
LegacySyncTrace::new(None, false),
);
let block0: Arc<Block> = zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES
.zcash_deserialize_into()
.expect("test vector deserializes");
let hash = block0.hash();
downloads
.download_and_verify(hash)
.await
.expect("queuing a fresh hash succeeds");
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(hash).collect()))
.await
.respond(zn::Response::Blocks(vec![]));
let result = downloads
.next()
.await
.expect("the download task produces a result instead of panicking");
assert!(
matches!(result, Err(BlockDownloadVerifyError::DownloadFailed { .. })),
"an empty block response must be a retryable DownloadFailed, got {result:?}",
);
}
#[tokio::test(start_paused = true)]
async fn block_download_network_readiness_times_out() {
let _init_guard = zakura_test::init();
let verifier =
MockService::build().for_unit_tests::<zakura_consensus::Request, block::Hash, _>();
let (chain_tip, _chain_tip_sender) = MockChainTip::new();
let (past_lookahead_limit_sender, _past_lookahead_limit_receiver) =
tokio::sync::watch::channel(false);
let mut downloads = Downloads::new(
NeverReadyNetwork,
verifier,
chain_tip,
past_lookahead_limit_sender,
sync::MIN_CONCURRENCY_LIMIT,
Height(0),
LegacySyncTrace::new(None, false),
);
let hash = block::Hash::from([0xCE; 32]);
let result = downloads.download_and_verify(hash).await;
assert!(
matches!(result, Err(BlockDownloadVerifyError::Timeout)),
"network readiness should time out, got {result:?}"
);
assert_eq!(downloads.in_flight(), 0);
}
fn setup_downloads(
peer_set: MockService<zn::Request, zn::Response, PanicAssertion>,
verifier: MockService<zakura_consensus::Request, block::Hash, PanicAssertion>,
chain_tip: MockChainTip,
) -> Downloads<
MockService<zn::Request, zn::Response, PanicAssertion>,
MockService<zakura_consensus::Request, block::Hash, PanicAssertion>,
MockChainTip,
> {
let (past_lookahead_limit_sender, _past_lookahead_limit_receiver) =
tokio::sync::watch::channel(false);
Downloads::new(
peer_set,
verifier,
chain_tip,
past_lookahead_limit_sender,
sync::MIN_CONCURRENCY_LIMIT,
Height(0),
LegacySyncTrace::new(None, false),
)
}
#[tokio::test]
async fn tip_child_rejects_poisoned_coinbase_height() {
let _init_guard = zakura_test::init();
let canonical: Arc<Block> = zakura_test::vectors::BLOCK_MAINNET_1687107_BYTES
.zcash_deserialize_into()
.expect("test vector deserializes");
let canonical_hash = canonical.hash();
let parent_hash = canonical.header.previous_block_hash;
let canonical_coinbase_hash = canonical.transactions[0].hash();
let canonical_auth_digest = canonical.transactions[0]
.auth_digest()
.expect("an NU5 coinbase has an authorizing data digest");
let poisoned = poison_coinbase_height(&canonical, Height(1));
assert_eq!(
poisoned.hash(),
canonical_hash,
"the poisoned body still answers a request for the canonical hash"
);
assert_eq!(
poisoned.transactions[0].hash(),
canonical_coinbase_hash,
"V5 authorizing data is excluded from the mined transaction ID, \
so the transaction merkle root is unchanged"
);
assert_ne!(
poisoned.transactions[0].auth_digest(),
Some(canonical_auth_digest),
"the mutation does change the authorizing data digest, \
which is what full consensus validation would have caught"
);
let mut peer_set = MockService::build().for_unit_tests::<zn::Request, zn::Response, _>();
let mut verifier =
MockService::build().for_unit_tests::<zakura_consensus::Request, block::Hash, _>();
let (chain_tip, chain_tip_sender) = MockChainTip::new();
chain_tip_sender.send_best_tip_height(Height(1_687_106));
chain_tip_sender.send_best_tip_hash(parent_hash);
let addr: PeerSocketAddr = "127.0.0.1:8233".parse().expect("valid peer address");
let mut downloads = setup_downloads(peer_set.clone(), verifier.clone(), chain_tip);
downloads
.download_and_verify(canonical_hash)
.await
.expect("queuing a fresh hash succeeds");
peer_set
.expect_request(zn::Request::BlocksByHash(
iter::once(canonical_hash).collect(),
))
.await
.respond(zn::Response::Blocks(vec![Available((
poisoned,
Some(addr),
))]));
let result = downloads
.next()
.await
.expect("the download task produces a result instead of panicking");
assert!(
matches!(
result,
Err(BlockDownloadVerifyError::TipChildHeightMismatch {
height: Height(1),
expected_height: Height(1_687_107),
hash,
advertiser_addr: Some(error_addr),
}) if hash == canonical_hash && error_addr == addr
),
"a poisoned tip child must be attributed to its supplier, got {result:?}"
);
verifier.expect_no_requests().await;
}
#[tokio::test]
async fn tip_child_rejects_forged_high_coinbase_height() {
let _init_guard = zakura_test::init();
let canonical: Arc<Block> = zakura_test::vectors::BLOCK_MAINNET_1687107_BYTES
.zcash_deserialize_into()
.expect("test vector deserializes");
let canonical_hash = canonical.hash();
let parent_hash = canonical.header.previous_block_hash;
let poisoned = poison_coinbase_height(&canonical, Height(2_000_000));
let mut peer_set = MockService::build().for_unit_tests::<zn::Request, zn::Response, _>();
let mut verifier =
MockService::build().for_unit_tests::<zakura_consensus::Request, block::Hash, _>();
let (chain_tip, chain_tip_sender) = MockChainTip::new();
chain_tip_sender.send_best_tip_height(Height(1_687_106));
chain_tip_sender.send_best_tip_hash(parent_hash);
let addr: PeerSocketAddr = "127.0.0.1:8233".parse().expect("valid peer address");
let mut downloads = setup_downloads(peer_set.clone(), verifier.clone(), chain_tip);
downloads
.download_and_verify(canonical_hash)
.await
.expect("queuing a fresh hash succeeds");
peer_set
.expect_request(zn::Request::BlocksByHash(
iter::once(canonical_hash).collect(),
))
.await
.respond(zn::Response::Blocks(vec![Available((
poisoned,
Some(addr),
))]));
let result = downloads
.next()
.await
.expect("the download task produces a result instead of panicking");
assert!(
matches!(
result,
Err(BlockDownloadVerifyError::TipChildHeightMismatch {
height: Height(2_000_000),
expected_height: Height(1_687_107),
..
})
),
"a forged high height on a tip child must be caught by the tip check, got {result:?}"
);
verifier.expect_no_requests().await;
}
#[tokio::test]
async fn canonical_tip_child_reaches_the_verifier() {
let _init_guard = zakura_test::init();
let canonical: Arc<Block> = zakura_test::vectors::BLOCK_MAINNET_1687107_BYTES
.zcash_deserialize_into()
.expect("test vector deserializes");
let canonical_hash = canonical.hash();
let parent_hash = canonical.header.previous_block_hash;
let mut peer_set = MockService::build().for_unit_tests::<zn::Request, zn::Response, _>();
let mut verifier =
MockService::build().for_unit_tests::<zakura_consensus::Request, block::Hash, _>();
let (chain_tip, chain_tip_sender) = MockChainTip::new();
chain_tip_sender.send_best_tip_height(Height(1_687_106));
chain_tip_sender.send_best_tip_hash(parent_hash);
let mut downloads = setup_downloads(peer_set.clone(), verifier.clone(), chain_tip);
downloads
.download_and_verify(canonical_hash)
.await
.expect("queuing a fresh hash succeeds");
peer_set
.expect_request(zn::Request::BlocksByHash(
iter::once(canonical_hash).collect(),
))
.await
.respond(zn::Response::Blocks(vec![Available((
canonical.clone(),
None,
))]));
verifier
.expect_request_that(|req| matches!(req, zakura_consensus::Request::Commit(_)))
.await
.respond(canonical_hash);
let result = downloads
.next()
.await
.expect("the download task produces a result instead of panicking");
assert_eq!(
result.expect("a canonical tip child verifies"),
(Height(1_687_107), canonical_hash),
"the canonical tip child must be unaffected by the height check"
);
}
#[tokio::test]
async fn non_tip_child_keeps_the_behind_tip_policy() {
let _init_guard = zakura_test::init();
let block1: Arc<Block> = zakura_test::vectors::BLOCK_MAINNET_1_BYTES
.zcash_deserialize_into()
.expect("test vector deserializes");
let block1_hash = block1.hash();
let mut peer_set = MockService::build().for_unit_tests::<zn::Request, zn::Response, _>();
let mut verifier =
MockService::build().for_unit_tests::<zakura_consensus::Request, block::Hash, _>();
let (chain_tip, chain_tip_sender) = MockChainTip::new();
chain_tip_sender.send_best_tip_height(Height(1_000_000));
chain_tip_sender.send_best_tip_hash(block::Hash([0x11; 32]));
let mut downloads = setup_downloads(peer_set.clone(), verifier.clone(), chain_tip);
downloads
.download_and_verify(block1_hash)
.await
.expect("queuing a fresh hash succeeds");
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block1_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((block1, None))]));
let result = downloads
.next()
.await
.expect("the download task produces a result instead of panicking");
assert!(
matches!(
result,
Err(BlockDownloadVerifyError::BehindTipHeightLimit { .. })
),
"an unauthenticated old block keeps the existing behind-tip policy, got {result:?}"
);
verifier.expect_no_requests().await;
}
#[tokio::test]
async fn poisoned_tip_child_requeues_and_scores_its_supplier() -> Result<(), crate::BoxError> {
let (
mut chain_sync,
_sync_status,
mut block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let (misbehavior_tx, mut misbehavior_rx) = tokio::sync::mpsc::channel(1);
chain_sync.misbehavior_sender = misbehavior_tx;
let block_hash = block::Hash::from([0xAB; 32]);
let addr: PeerSocketAddr = "127.0.0.1:8233".parse().expect("valid peer address");
let error = BlockDownloadVerifyError::TipChildHeightMismatch {
height: Height(1),
expected_height: Height(1_687_107),
hash: block_hash,
advertiser_addr: Some(addr),
};
let requeue = tokio::spawn(async move {
let result = chain_sync
.handle_block_response_with_missing_retry(Err(error))
.await;
(result, chain_sync)
});
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block_hash).collect()))
.await
.respond(Err(not_found_block_error(block_hash)));
let (result, chain_sync) = requeue.await.expect("the retry task should not panic");
result?;
assert_eq!(
misbehavior_rx.recv().await,
Some((addr, 100)),
"the supplier of a poisoned body must be scored for a ban"
);
assert_eq!(
chain_sync.poisoned_block_retry_counts.get(&block_hash),
Some(&1),
"the requeue must be counted against the retry budget"
);
block_verifier_router.expect_no_requests().await;
Ok(())
}
#[tokio::test]
async fn poisoned_tip_child_restarts_after_retry_limit() {
let (
mut chain_sync,
_sync_status,
_block_verifier_router,
mut peer_set,
_state_service,
_mock_chain_tip_sender,
) = setup_chain_sync();
let block_hash = block::Hash::from([0xAB; 32]);
chain_sync
.poisoned_block_retry_counts
.insert(block_hash, sync::POISONED_BLOCK_RETRY_LIMIT);
let error = BlockDownloadVerifyError::TipChildHeightMismatch {
height: Height(1),
expected_height: Height(1_687_107),
hash: block_hash,
advertiser_addr: Some("127.0.0.1:8233".parse().expect("valid peer address")),
};
let result = chain_sync
.handle_block_response_with_missing_retry(Err(error))
.await;
assert!(
result.is_err(),
"an exhausted poisoned-body budget must restart sync"
);
peer_set.expect_no_requests().await;
}
#[tokio::test]
async fn tip_height_without_a_tip_hash_keeps_the_behind_tip_policy() {
let _init_guard = zakura_test::init();
let block1: Arc<Block> = zakura_test::vectors::BLOCK_MAINNET_1_BYTES
.zcash_deserialize_into()
.expect("test vector deserializes");
let block1_hash = block1.hash();
let mut peer_set = MockService::build().for_unit_tests::<zn::Request, zn::Response, _>();
let mut verifier =
MockService::build().for_unit_tests::<zakura_consensus::Request, block::Hash, _>();
let (chain_tip, chain_tip_sender) = MockChainTip::new();
chain_tip_sender.send_best_tip_height(Height(1_000_000));
let mut downloads = setup_downloads(peer_set.clone(), verifier.clone(), chain_tip);
downloads
.download_and_verify(block1_hash)
.await
.expect("queuing a fresh hash succeeds");
peer_set
.expect_request(zn::Request::BlocksByHash(iter::once(block1_hash).collect()))
.await
.respond(zn::Response::Blocks(vec![Available((block1, None))]));
let result = downloads
.next()
.await
.expect("the download task produces a result instead of panicking");
assert!(
matches!(
result,
Err(BlockDownloadVerifyError::BehindTipHeightLimit { .. })
),
"a height-only tip must still apply the behind-tip policy, got {result:?}"
);
verifier.expect_no_requests().await;
}