diff --git a/blockprod/src/detail/tests/collect_transactions.rs b/blockprod/src/detail/tests/collect_transactions.rs index 7c4ef653d..b0913371c 100644 --- a/blockprod/src/detail/tests/collect_transactions.rs +++ b/blockprod/src/detail/tests/collect_transactions.rs @@ -16,7 +16,6 @@ use common::{ chain::block::timestamp::BlockTimestamp, primitives::{H256, Id}, - time_getter::TimeGetter, }; use mempool::{ error::{BlockConstructionError, TxValidationError}, @@ -24,10 +23,11 @@ use mempool::{ }; use mocks::MockMempoolInterface; use subsystem::error::ResponseError; +use test_utils::assert_matches; use utils::once_destructor::OnceDestructor; use crate::{ - BlockProductionError, detail::collect_transactions, tests::helpers::setup_blockprod_test, + BlockProductionError, detail::collect_transactions, tests::helpers::BlockprodTestSetupBuilder, }; // A dummy timestamp for tests where the block timestamp is irrelevant @@ -37,8 +37,7 @@ const DUMMY_TIMESTAMP: BlockTimestamp = BlockTimestamp::from_int_seconds(0u64); #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn collect_txs_failed() { - let (mut manager, chain_config, _chainstate, _mempool, _p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, mut manager) = BlockprodTestSetupBuilder::new().build(); let mut mock_mempool = MockMempoolInterface::default(); mock_mempool.expect_collect_txs().return_once(|_, _, _| { @@ -55,7 +54,7 @@ async fn collect_txs_failed() { let tester = tokio::spawn(async move { let transactions = collect_transactions( &mock_mempool_subsystem, - &chain_config, + &blockprod_setup.chain_config, current_tip, DUMMY_TIMESTAMP, vec![], @@ -64,12 +63,12 @@ async fn collect_txs_failed() { ) .await; - match transactions { + assert_matches!( + transactions, Err(BlockProductionError::MempoolBlockConstruction( BlockConstructionError::Validity(TxValidationError::SubsystemCallError(_)), - )) => {} - _ => panic!("Expected collect_tx() to fail"), - }; + )) + ); shutdown.initiate(); }); @@ -79,8 +78,7 @@ async fn collect_txs_failed() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn subsystem_error() { - let (mut manager, chain_config, _chainstate, _mempool, _p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, mut manager) = BlockprodTestSetupBuilder::new().build(); let mock_mempool = MockMempoolInterface::default(); let mock_mempool_subsystem = manager.add_subsystem("mock-mempool", mock_mempool); @@ -102,7 +100,7 @@ async fn subsystem_error() { tokio::spawn(async move { let transactions = collect_transactions( &mock_mempool_subsystem, - &chain_config, + &blockprod_setup.chain_config, current_tip, DUMMY_TIMESTAMP, vec![], @@ -118,13 +116,12 @@ async fn subsystem_error() { }; }) .await - .expect("Subsystem error thread failed"); + .unwrap(); } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn succeeded() { - let (mut manager, chain_config, _chainstate, _mempool, _p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, mut manager) = BlockprodTestSetupBuilder::new().build(); let mut mock_mempool = MockMempoolInterface::default(); @@ -153,7 +150,7 @@ async fn succeeded() { let transactions = collect_transactions( &mock_mempool_subsystem, - &chain_config, + &blockprod_setup.chain_config, current_tip, DUMMY_TIMESTAMP, vec![], diff --git a/blockprod/src/detail/tests/process_block_with_custom_id.rs b/blockprod/src/detail/tests/process_block_with_custom_id.rs index dc0a9a211..1e1ac4373 100644 --- a/blockprod/src/detail/tests/process_block_with_custom_id.rs +++ b/blockprod/src/detail/tests/process_block_with_custom_id.rs @@ -13,21 +13,17 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::sync::Arc; - use rstest::rstest; -use common::time_getter::TimeGetter; use mempool::tx_accumulator::PackingStrategy; use randomness::RngExt as _; use test_utils::random::{Seed, make_seedable_rng}; use utils::once_destructor::OnceDestructor; use crate::{ - BlockProduction, BlockProductionError, + BlockProductionError, detail::{GenerateBlockInputData, job_manager::JobManagerError}, - prepare_thread_pool, test_blockprod_config, - tests::helpers::setup_blockprod_test, + tests::helpers::BlockprodTestSetupBuilder, }; #[rstest] @@ -35,23 +31,13 @@ use crate::{ #[case(Seed::from_entropy())] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn multiple_jobs_with_wait(#[case] seed: Seed) { - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new().build(); let mut rng = make_seedable_rng(seed); let jobs_to_create = rng.random_range(1..=20); - let block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate, - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let block_production = blockprod_setup.make_blockprod_builder().build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -95,23 +81,13 @@ async fn multiple_jobs_with_wait(#[case] seed: Seed) { #[case(Seed::from_entropy())] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn multiple_jobs_without_wait_same_jobkey(#[case] seed: Seed) { - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new().build(); let mut rng = make_seedable_rng(seed); let jobs_to_create = 10 + rng.random_range(1..=20); - let block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate, - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let block_production = blockprod_setup.make_blockprod_builder().build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); diff --git a/blockprod/src/detail/tests/produce_block/mod.rs b/blockprod/src/detail/tests/produce_block/mod.rs index fa61912fc..ec855b278 100644 --- a/blockprod/src/detail/tests/produce_block/mod.rs +++ b/blockprod/src/detail/tests/produce_block/mod.rs @@ -24,8 +24,7 @@ use tokio::{ }; use chainstate::{ - ChainstateError, ChainstateHandle, GenBlockIndex, PropertyQueryError, - chainstate_interface::ChainstateInterface, + ChainstateError, GenBlockIndex, PropertyQueryError, chainstate_interface::ChainstateInterface, }; use common::{ Uint256, @@ -68,19 +67,18 @@ use crate::{ CustomId, GenerateBlockInputData, job_manager::{JobManagerError, JobManagerImpl, tests::MockJobManager}, }, - prepare_thread_pool, test_blockprod_config, + test_blockprod_config, tests::helpers::{ - assert_process_block, build_chain_config_for_pos, make_genesis_timestamp, - setup_blockprod_test, setup_pos, setup_pos_with_genesis_timestamp, + BlockprodTestSetupBuilder, PoSTestSetupBuilder, make_chain_config_builder, + make_genesis_timestamp, }, }; #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn initial_block_download() { - let (mut manager, chain_config, _, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, mut manager) = BlockprodTestSetupBuilder::new().build(); - let chainstate_subsystem: ChainstateHandle = { + let mock_chainstate = { let mut mock_chainstate = MockChainstateInterface::new(); mock_chainstate.expect_is_initial_block_download().returning(|| true); @@ -100,30 +98,22 @@ async fn initial_block_download() { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate_subsystem, - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let block_production = blockprod_setup + .make_blockprod_builder() + .with_chainstate(mock_chainstate) + .build(); - let result = block_production + let err = block_production .produce_block( GenerateBlockInputData::None, vec![], vec![], PackingStrategy::FillSpaceFromMempool, ) - .await; + .await + .unwrap_err(); - match result { - Err(BlockProductionError::ChainstateWaitForSync) => {} - _ => panic!("Unexpected return value"), - } + assert_eq!(err, BlockProductionError::ChainstateWaitForSync); } }); @@ -133,8 +123,7 @@ async fn initial_block_download() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn below_peer_count() { - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new().build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -147,30 +136,25 @@ async fn below_peer_count() { let mut blockprod_config = test_blockprod_config(); blockprod_config.min_peers_to_produce_blocks = 100; - let block_production = BlockProduction::new( - chain_config, - Arc::new(blockprod_config), - chainstate, - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let block_production = blockprod_setup + .make_blockprod_builder() + .with_blockprod_config(blockprod_config) + .build(); - let result = block_production + let err = block_production .produce_block( GenerateBlockInputData::None, vec![], vec![], PackingStrategy::LeaveEmptySpace, ) - .await; + .await + .unwrap_err(); - match result { - Err(BlockProductionError::PeerCountBelowRequiredThreshold(0, 100)) => {} - _ => panic!("Unexpected return value"), - } + assert_eq!( + err, + BlockProductionError::PeerCountBelowRequiredThreshold(0, 100) + ); } }); @@ -180,10 +164,9 @@ async fn below_peer_count() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn pull_best_block_index_error() { - let (mut manager, chain_config, _, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, mut manager) = BlockprodTestSetupBuilder::new().build(); - let chainstate_subsystem: ChainstateHandle = { + let mock_chainstate = { let mut mock_chainstate = Box::new(MockChainstateInterface::new()); mock_chainstate .expect_subscribe_to_subsystem_events() @@ -209,16 +192,10 @@ async fn pull_best_block_index_error() { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate_subsystem, - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let block_production = blockprod_setup + .make_blockprod_builder() + .with_chainstate(mock_chainstate) + .build(); let result = block_production .produce_block( @@ -229,12 +206,12 @@ async fn pull_best_block_index_error() { ) .await; - match result { + assert_matches!( + result, Err(BlockProductionError::ChainstateError( consensus::ChainstateError::FailedToObtainBestBlockIndex(_), - )) => {} - _ => panic!("Unexpected return value"), - } + )) + ); } }); @@ -244,8 +221,7 @@ async fn pull_best_block_index_error() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn add_job_error() { - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new().build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -255,16 +231,7 @@ async fn add_job_error() { shutdown_trigger.initiate(); }); - let mut block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let mut block_production = blockprod_setup.make_blockprod_builder().build(); let mut mock_job_manager = Box::new(MockJobManager::new()); @@ -275,21 +242,20 @@ async fn add_job_error() { block_production.set_job_manager(mock_job_manager); - let result = block_production + let err = block_production .produce_block( GenerateBlockInputData::None, vec![], vec![], PackingStrategy::LeaveEmptySpace, ) - .await; + .await + .unwrap_err(); - match result { - Err(BlockProductionError::JobManagerError( - JobManagerError::FailedToSendNewJobEvent, - )) => {} - _ => panic!("Unexpected return value"), - } + assert_eq!( + err, + BlockProductionError::JobManagerError(JobManagerError::FailedToSendNewJobEvent,) + ); } }); @@ -304,24 +270,12 @@ async fn add_job_error() { async fn overflow_tip_plus_one(#[case] seed: Seed) { let mut rng = make_seedable_rng(seed); - let ( - chain_config_builder, - genesis_stake_private_key, - genesis_vrf_private_key, - create_genesis_pool_txoutput, - ) = setup_pos_with_genesis_timestamp( - BlockTimestamp::from_int_seconds(u64::MAX), - BlockHeight::new(1), - &[], - &mut rng, - ); - - let (manager, chain_config, chainstate, mempool, p2p) = { - setup_blockprod_test( - Some(build_chain_config_for_pos(chain_config_builder)), - TimeGetter::default(), - ) - }; + let pos_setup = + PoSTestSetupBuilder::new().build(BlockTimestamp::from_int_seconds(u64::MAX), &mut rng); + + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new() + .with_chain_config(Arc::clone(&pos_setup.chain_config)) + .build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -332,27 +286,8 @@ async fn overflow_tip_plus_one(#[case] seed: Seed) { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - Arc::clone(&chain_config), - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); - - let input_data = GenerateBlockInputData::PoS(Box::new(PoSGenerateBlockInputData::new( - genesis_stake_private_key, - genesis_vrf_private_key, - PoolId::new(H256::zero()), - vec![TxInput::from_utxo( - OutPointSourceId::BlockReward(chain_config.genesis_block_id()), - 0, - )], - vec![create_genesis_pool_txoutput], - ))); + let block_production = blockprod_setup.make_blockprod_builder().build(); + let input_data = pos_setup.make_first_pos_block_input_data(); let result = block_production .produce_block(input_data, vec![], vec![], PackingStrategy::LeaveEmptySpace) @@ -375,18 +310,16 @@ async fn overflow_tip_plus_one(#[case] seed: Seed) { async fn overflow_max_blocktimestamp(#[case] seed: Seed) { let mut rng = make_seedable_rng(seed); let time_getter = TimeGetter::default(); - let ( - chain_config_builder, - genesis_stake_private_key, - genesis_vrf_private_key, - create_genesis_pool_txoutput, - ) = setup_pos(&time_getter, BlockHeight::new(1), &[], &mut rng); - let chain_config = build_chain_config_for_pos( - chain_config_builder.max_future_block_time_offset(Some(Duration::MAX)), - ); - - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(Some(chain_config), TimeGetter::default()); + let pos_setup = PoSTestSetupBuilder::new() + .with_chain_config_builder( + make_chain_config_builder().max_future_block_time_offset(Some(Duration::MAX)), + ) + .build(make_genesis_timestamp(&time_getter, &mut rng), &mut rng); + + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new() + .with_chain_config(Arc::clone(&pos_setup.chain_config)) + .with_time_getter(time_getter) + .build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -396,27 +329,8 @@ async fn overflow_max_blocktimestamp(#[case] seed: Seed) { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - Arc::clone(&chain_config), - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); - - let input_data = GenerateBlockInputData::PoS(Box::new(PoSGenerateBlockInputData::new( - genesis_stake_private_key, - genesis_vrf_private_key, - PoolId::new(H256::zero()), - vec![TxInput::from_utxo( - OutPointSourceId::BlockReward(chain_config.genesis_block_id()), - 0, - )], - vec![create_genesis_pool_txoutput], - ))); + let block_production = blockprod_setup.make_blockprod_builder().build(); + let input_data = pos_setup.make_first_pos_block_input_data(); let result = block_production .produce_block(input_data, vec![], vec![], PackingStrategy::LeaveEmptySpace) @@ -439,17 +353,13 @@ async fn overflow_max_blocktimestamp(#[case] seed: Seed) { async fn update_last_used_block_timestamp(#[case] seed: Seed) { let mut rng = make_seedable_rng(seed); let time_getter = TimeGetter::default(); - let ( - chain_config_builder, - genesis_stake_private_key, - genesis_vrf_private_key, - create_genesis_pool_txoutput, - ) = setup_pos(&time_getter, BlockHeight::new(1), &[], &mut rng); - - let (manager, chain_config, chainstate, mempool, p2p) = setup_blockprod_test( - Some(build_chain_config_for_pos(chain_config_builder)), - time_getter, - ); + let pos_setup = + PoSTestSetupBuilder::new().build(make_genesis_timestamp(&time_getter, &mut rng), &mut rng); + + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new() + .with_chain_config(Arc::clone(&pos_setup.chain_config)) + .with_time_getter(time_getter) + .build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -459,27 +369,8 @@ async fn update_last_used_block_timestamp(#[case] seed: Seed) { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - chain_config.clone(), - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); - - let input_data = GenerateBlockInputData::PoS(Box::new(PoSGenerateBlockInputData::new( - genesis_stake_private_key, - genesis_vrf_private_key, - PoolId::new(H256::zero()), - vec![TxInput::from_utxo( - OutPointSourceId::BlockReward(chain_config.genesis_block_id()), - 0, - )], - vec![create_genesis_pool_txoutput], - ))); + let block_production = blockprod_setup.make_blockprod_builder().build(); + let input_data = pos_setup.make_first_pos_block_input_data(); let _ = block_production .job_manager_handle @@ -513,29 +404,23 @@ async fn try_again_later(#[case] seed: Seed) { let default_time_getter = TimeGetter::default(); let genesis_time = default_time_getter.get_time(); - let ( - chain_config_builder, - genesis_stake_private_key, - genesis_vrf_private_key, - create_genesis_pool_txoutput, - ) = setup_pos_with_genesis_timestamp( - BlockTimestamp::from_time(genesis_time), - BlockHeight::new(1), - &[], - &mut rng, - ); - - let chain_config = build_chain_config_for_pos(chain_config_builder); + let pos_setup = + PoSTestSetupBuilder::new().build(BlockTimestamp::from_time(genesis_time), &mut rng); + let time_getter = { let cur_time_secs = genesis_time - .saturating_duration_sub(chain_config.max_future_block_time_offset(BlockHeight::zero())) + .saturating_duration_sub( + pos_setup.chain_config.max_future_block_time_offset(BlockHeight::zero()), + ) .as_secs_since_epoch(); let time_value = Arc::new(SeqCstAtomicU64::new(cur_time_secs)); mocked_time_getter_seconds(time_value) }; - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(Some(chain_config), time_getter.clone()); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new() + .with_chain_config(Arc::clone(&pos_setup.chain_config)) + .with_time_getter(time_getter) + .build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -545,33 +430,15 @@ async fn try_again_later(#[case] seed: Seed) { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - Arc::clone(&chain_config), - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool, - p2p, - time_getter, - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); - - let input_data = GenerateBlockInputData::PoS(Box::new(PoSGenerateBlockInputData::new( - genesis_stake_private_key, - genesis_vrf_private_key, - PoolId::new(H256::zero()), - vec![TxInput::from_utxo( - OutPointSourceId::BlockReward(chain_config.genesis_block_id()), - 0, - )], - vec![create_genesis_pool_txoutput], - ))); + let block_production = blockprod_setup.make_blockprod_builder().build(); + let input_data = pos_setup.make_first_pos_block_input_data(); - let result = block_production + let err = block_production .produce_block(input_data, vec![], vec![], PackingStrategy::LeaveEmptySpace) - .await; + .await + .unwrap_err(); - assert_matches!(result, Err(BlockProductionError::TryAgainLater)); + assert_eq!(err, BlockProductionError::TryAgainLater); assert_job_count(&block_production, 0).await; } @@ -583,10 +450,9 @@ async fn try_again_later(#[case] seed: Seed) { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn pull_consensus_data_error() { - let (mut manager, chain_config, _, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, mut manager) = BlockprodTestSetupBuilder::new().build(); - let chainstate_subsystem: ChainstateHandle = { + let mock_chainstate = { let mut mock_chainstate = MockChainstateInterface::new(); mock_chainstate .expect_subscribe_to_subsystem_events() @@ -595,7 +461,7 @@ async fn pull_consensus_data_error() { mock_chainstate.expect_is_initial_block_download().returning(|| false); let mut expected_return_values = vec![ - Ok(GenBlockIndex::genesis(&chain_config)), + Ok(GenBlockIndex::genesis(&blockprod_setup.chain_config)), Err(ChainstateError::FailedToReadProperty( PropertyQueryError::BestBlockIndexNotFound, )), @@ -617,16 +483,10 @@ async fn pull_consensus_data_error() { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate_subsystem, - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let block_production = blockprod_setup + .make_blockprod_builder() + .with_chainstate(mock_chainstate) + .build(); let result = block_production .produce_block( @@ -652,18 +512,19 @@ async fn pull_consensus_data_error() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn transaction_source_mempool_error() { - let (mut manager, chain_config, chainstate, _mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, mut manager) = BlockprodTestSetupBuilder::new().build(); - let mut mock_mempool = MockMempoolInterface::default(); + let mock_mempool = { + let mut mock_mempool = MockMempoolInterface::default(); - mock_mempool.expect_collect_txs().return_once(|_, _, _| { - Err(BlockConstructionError::Validity( - TxValidationError::SubsystemCallError(ResponseError::NoResponse.into()), - )) - }); + mock_mempool.expect_collect_txs().return_once(|_, _, _| { + Err(BlockConstructionError::Validity( + TxValidationError::SubsystemCallError(ResponseError::NoResponse.into()), + )) + }); - let mempool_subsystem = manager.add_subsystem("mock-mempool", mock_mempool); + manager.add_subsystem("mock-mempool", mock_mempool) + }; let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -673,16 +534,8 @@ async fn transaction_source_mempool_error() { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool_subsystem, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let block_production = + blockprod_setup.make_blockprod_builder().with_mempool(mock_mempool).build(); let result = block_production .produce_block( @@ -693,12 +546,12 @@ async fn transaction_source_mempool_error() { ) .await; - match result { + assert_matches!( + result, Err(BlockProductionError::MempoolBlockConstruction( BlockConstructionError::Validity(TxValidationError::SubsystemCallError(_)), - )) => {} - _ => panic!("Unexpected return value: {result:?}"), - } + )) + ); } }); @@ -708,8 +561,7 @@ async fn transaction_source_mempool_error() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn transaction_source_mempool() { - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new().build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -719,16 +571,7 @@ async fn transaction_source_mempool() { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool.clone(), - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let block_production = blockprod_setup.make_blockprod_builder().build(); let (new_block, job_finished_receiver) = block_production // TODO: Add transactions to the mempool @@ -739,12 +582,12 @@ async fn transaction_source_mempool() { PackingStrategy::FillSpaceFromMempool, ) .await - .expect("Failed to produce a block: {:?}"); + .unwrap(); - job_finished_receiver.await.expect("Job finished receiver closed"); + job_finished_receiver.await.unwrap(); assert_job_count(&block_production, 0).await; - assert_process_block(&chainstate, &mempool, new_block).await; + blockprod_setup.assert_process_block(new_block).await; } }); @@ -754,8 +597,7 @@ async fn transaction_source_mempool() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn transaction_source_provided() { - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new().build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -765,16 +607,7 @@ async fn transaction_source_provided() { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool.clone(), - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let block_production = blockprod_setup.make_blockprod_builder().build(); let (new_block, job_finished_receiver) = block_production // TODO: Add transactions to the parameters @@ -785,12 +618,12 @@ async fn transaction_source_provided() { PackingStrategy::LeaveEmptySpace, ) .await - .expect("Failed to produce a block: {:?}"); + .unwrap(); - job_finished_receiver.await.expect("Job finished receiver closed"); + job_finished_receiver.await.unwrap(); assert_job_count(&block_production, 0).await; - assert_process_block(&chainstate, &mempool, new_block).await; + blockprod_setup.assert_process_block(new_block).await; } }); @@ -803,7 +636,7 @@ async fn transaction_source_provided() { #[case(Seed::from_entropy())] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn cancel_received(#[case] seed: Seed) { - let override_chain_config = { + let chain_config = { let net_upgrades = NetUpgrades::initialize(vec![( BlockHeight::new(0), ConsensusUpgrade::PoW { @@ -813,13 +646,14 @@ async fn cancel_received(#[case] seed: Seed) { initial_difficulty: Uint256::ZERO.into(), }, )]) - .expect("Net upgrade is valid"); + .unwrap(); - Builder::new(ChainType::Regtest).consensus_upgrades(net_upgrades).build() + Arc::new(Builder::new(ChainType::Regtest).consensus_upgrades(net_upgrades).build()) }; - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(Some(override_chain_config), TimeGetter::default()); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new() + .with_chain_config(Arc::clone(&chain_config)) + .build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -829,16 +663,7 @@ async fn cancel_received(#[case] seed: Seed) { shutdown_trigger.initiate(); }); - let mut block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate, - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let mut block_production = blockprod_setup.make_blockprod_builder().build(); let mut mock_job_manager = Box::::default(); @@ -861,7 +686,7 @@ async fn cancel_received(#[case] seed: Seed) { block_production.set_job_manager(mock_job_manager); - let result = block_production + let err = block_production .produce_block( GenerateBlockInputData::PoW(Box::new(PoWGenerateBlockInputData::new( Destination::AnyoneCanSpend, @@ -870,12 +695,10 @@ async fn cancel_received(#[case] seed: Seed) { vec![], PackingStrategy::LeaveEmptySpace, ) - .await; + .await + .unwrap_err(); - match result { - Err(BlockProductionError::Cancelled) => {} - _ => panic!("Unexpected return value"), - } + assert_eq!(err, BlockProductionError::Cancelled); } }); @@ -885,8 +708,7 @@ async fn cancel_received(#[case] seed: Seed) { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn solved_ignore_consensus() { - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new().build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -896,16 +718,7 @@ async fn solved_ignore_consensus() { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool.clone(), - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let block_production = blockprod_setup.make_blockprod_builder().build(); let (new_block, job_finished_receiver) = block_production .produce_block( @@ -915,12 +728,12 @@ async fn solved_ignore_consensus() { PackingStrategy::LeaveEmptySpace, ) .await - .expect("Failed to produce a block: {:?}"); + .unwrap(); - job_finished_receiver.await.expect("Job finished receiver closed"); + job_finished_receiver.await.unwrap(); assert_job_count(&block_production, 0).await; - assert_process_block(&chainstate, &mempool, new_block).await; + blockprod_setup.assert_process_block(new_block).await; } }); @@ -930,20 +743,21 @@ async fn solved_ignore_consensus() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn solved_pow_consensus() { - let override_chain_config = { + let chain_config = { let net_upgrades = NetUpgrades::initialize(vec![( BlockHeight::new(0), ConsensusUpgrade::PoW { initial_difficulty: Uint256::MAX.into(), }, )]) - .expect("Net upgrade is valid"); + .unwrap(); - Builder::new(ChainType::Regtest).consensus_upgrades(net_upgrades).build() + Arc::new(Builder::new(ChainType::Regtest).consensus_upgrades(net_upgrades).build()) }; - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(Some(override_chain_config), TimeGetter::default()); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new() + .with_chain_config(Arc::clone(&chain_config)) + .build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -953,16 +767,7 @@ async fn solved_pow_consensus() { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool.clone(), - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let block_production = blockprod_setup.make_blockprod_builder().build(); let (new_block, job_finished_receiver) = block_production .produce_block( @@ -974,12 +779,12 @@ async fn solved_pow_consensus() { PackingStrategy::LeaveEmptySpace, ) .await - .expect("Failed to produce a block: {:?}"); + .unwrap(); - job_finished_receiver.await.expect("Job finished receiver closed"); + job_finished_receiver.await.unwrap(); assert_job_count(&block_production, 0).await; - assert_process_block(&chainstate, &mempool, new_block).await; + blockprod_setup.assert_process_block(new_block).await; } }); @@ -994,17 +799,13 @@ async fn solved_pow_consensus() { async fn solved_pos_consensus(#[case] seed: Seed) { let mut rng = make_seedable_rng(seed); let time_getter = TimeGetter::default(); - let ( - chain_config_builder, - genesis_stake_private_key, - genesis_vrf_private_key, - create_genesis_pool_txoutput, - ) = setup_pos(&time_getter, BlockHeight::new(1), &[], &mut rng); - - let (manager, chain_config, chainstate, mempool, p2p) = setup_blockprod_test( - Some(build_chain_config_for_pos(chain_config_builder)), - time_getter, - ); + let pos_setup = + PoSTestSetupBuilder::new().build(make_genesis_timestamp(&time_getter, &mut rng), &mut rng); + + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new() + .with_chain_config(Arc::clone(&pos_setup.chain_config)) + .with_time_getter(time_getter) + .build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -1014,42 +815,18 @@ async fn solved_pos_consensus(#[case] seed: Seed) { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - chain_config.clone(), - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool.clone(), - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); - - let input_data = Box::new(PoSGenerateBlockInputData::new( - genesis_stake_private_key, - genesis_vrf_private_key, - PoolId::new(H256::zero()), - vec![TxInput::from_utxo( - OutPointSourceId::BlockReward(chain_config.genesis_block_id()), - 0, - )], - vec![create_genesis_pool_txoutput], - )); + let block_production = blockprod_setup.make_blockprod_builder().build(); + let input_data = pos_setup.make_first_pos_block_input_data(); let (new_block, job_finished_receiver) = block_production - .produce_block( - GenerateBlockInputData::PoS(input_data), - vec![], - vec![], - PackingStrategy::LeaveEmptySpace, - ) + .produce_block(input_data, vec![], vec![], PackingStrategy::LeaveEmptySpace) .await - .expect("Failed to produce a block: {:?}"); + .unwrap(); - job_finished_receiver.await.expect("Job finished receiver closed"); + job_finished_receiver.await.unwrap(); assert_job_count(&block_production, 0).await; - assert_process_block(&chainstate, &mempool, new_block).await; + blockprod_setup.assert_process_block(new_block).await; } }); @@ -1086,7 +863,7 @@ async fn solve_lots_of_blocks_with_differing_consensus(#[case] seed: Seed) { Destination::PublicKey(genesis_stake_public_key.clone()), genesis_vrf_public_key, Destination::PublicKey(genesis_stake_public_key.clone()), - PerThousand::new(1000).expect("Valid per thousand"), + PerThousand::new(1000).unwrap(), Amount::ZERO, )), ) @@ -1094,7 +871,7 @@ async fn solve_lots_of_blocks_with_differing_consensus(#[case] seed: Seed) { let blocks_to_generate = rng.random_range(100..=1000); - let override_chain_config = { + let chain_config = { let genesis_block = Genesis::new( "blockprod-testing".into(), make_genesis_timestamp(&time_getter, &mut rng), @@ -1130,17 +907,20 @@ async fn solve_lots_of_blocks_with_differing_consensus(#[case] seed: Seed) { next_height_consensus_change += rng.random_range(1..50); } - let net_upgrades = - NetUpgrades::initialize(randomized_net_upgrades).expect("Net upgrades are valid"); + let net_upgrades = NetUpgrades::initialize(randomized_net_upgrades).unwrap(); - Builder::new(ChainType::Regtest) - .genesis_custom(genesis_block) - .consensus_upgrades(net_upgrades) - .build() + Arc::new( + Builder::new(ChainType::Regtest) + .genesis_custom(genesis_block) + .consensus_upgrades(net_upgrades) + .build(), + ) }; - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(Some(override_chain_config), time_getter); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new() + .with_chain_config(Arc::clone(&chain_config)) + .with_time_getter(time_getter) + .build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -1150,16 +930,7 @@ async fn solve_lots_of_blocks_with_differing_consensus(#[case] seed: Seed) { shutdown_trigger.initiate(); }); - let mut block_production = BlockProduction::new( - chain_config.clone(), - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool.clone(), - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let mut block_production = blockprod_setup.make_blockprod_builder().build(); let no_chainstate_job_manager = Box::new(JobManagerImpl::new(None)); block_production.set_job_manager(no_chainstate_job_manager); @@ -1196,52 +967,54 @@ async fn solve_lots_of_blocks_with_differing_consensus(#[case] seed: Seed) { PackingStrategy::LeaveEmptySpace, ) .await - .expect("Failed to produce a block: {:?}"); + .unwrap(); - job_finished_receiver.await.expect("Job finished receiver closed"); + job_finished_receiver.await.unwrap(); - assert_process_block(&chainstate, &mempool, new_block.clone()).await; + blockprod_setup.assert_process_block(new_block.clone()).await; } RequiredConsensus::PoS(_) => { // Try no input data for PoS consensus - let input_data_none_result = block_production + let input_data_none_err = block_production .produce_block( GenerateBlockInputData::None, vec![], vec![], PackingStrategy::LeaveEmptySpace, ) - .await; + .await + .unwrap_err(); - match input_data_none_result { - Err(BlockProductionError::FailedConsensusInitialization( + assert_eq!( + input_data_none_err, + BlockProductionError::FailedConsensusInitialization( ConsensusCreationError::StakingError( ConsensusPoSError::NoInputDataProvided, ), - )) => {} - _ => panic!("Unexpected return value"), - } + ) + ); // Try PoW input data for PoS consensus - let input_data_pow_result = block_production + let input_data_pow_err = block_production .produce_block( input_data_pow, vec![], vec![], PackingStrategy::LeaveEmptySpace, ) - .await; + .await + .unwrap_err(); - match input_data_pow_result { - Err(BlockProductionError::FailedConsensusInitialization( + assert_eq!( + input_data_pow_err, + BlockProductionError::FailedConsensusInitialization( ConsensusCreationError::StakingError( ConsensusPoSError::PoWInputDataProvided, ), - )) => {} - _ => panic!("Unexpected return value"), - } + ) + ); // Try PoS input data for PoS consensus @@ -1253,11 +1026,11 @@ async fn solve_lots_of_blocks_with_differing_consensus(#[case] seed: Seed) { PackingStrategy::LeaveEmptySpace, ) .await - .expect("Failed to produce a job: {:?}"); + .unwrap(); - job_finished_receiver.await.expect("Job finished receiver closed"); + job_finished_receiver.await.unwrap(); - let result = assert_process_block(&chainstate, &mempool, new_block).await; + let result = blockprod_setup.assert_process_block(new_block).await; // Update kernel input parameters for future PoS blocks @@ -1274,43 +1047,45 @@ async fn solve_lots_of_blocks_with_differing_consensus(#[case] seed: Seed) { RequiredConsensus::PoW(_) => { // Try no input data for PoW consensus - let input_data_none_result = block_production + let input_data_none_err = block_production .produce_block( GenerateBlockInputData::None, vec![], vec![], PackingStrategy::LeaveEmptySpace, ) - .await; + .await + .unwrap_err(); - match input_data_none_result { - Err(BlockProductionError::FailedConsensusInitialization( + assert_eq!( + input_data_none_err, + BlockProductionError::FailedConsensusInitialization( ConsensusCreationError::MiningError( ConsensusPoWError::NoInputDataProvided, ), - )) => {} - _ => panic!("Unexpected return value"), - } + ) + ); // Try PoS input data for PoW consensus - let input_data_pos_result = block_production + let input_data_pos_err = block_production .produce_block( input_data_pos, vec![], vec![], PackingStrategy::LeaveEmptySpace, ) - .await; + .await + .unwrap_err(); - match input_data_pos_result { - Err(BlockProductionError::FailedConsensusInitialization( + assert_eq!( + input_data_pos_err, + BlockProductionError::FailedConsensusInitialization( ConsensusCreationError::MiningError( ConsensusPoWError::PoSInputDataProvided, ), - )) => {} - _ => panic!("Unexpected return value"), - } + ) + ); // Try PoW input data for PoW consensus @@ -1322,11 +1097,11 @@ async fn solve_lots_of_blocks_with_differing_consensus(#[case] seed: Seed) { PackingStrategy::LeaveEmptySpace, ) .await - .expect("Failed to produce a block: {:?}"); + .unwrap(); - job_finished_receiver.await.expect("Job finished receiver closed"); + job_finished_receiver.await.unwrap(); - assert_process_block(&chainstate, &mempool, new_block.clone()).await; + blockprod_setup.assert_process_block(new_block.clone()).await; } } } @@ -1342,8 +1117,7 @@ async fn solve_lots_of_blocks_with_differing_consensus(#[case] seed: Seed) { #[case(Seed::from_entropy())] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn multiple_jobs_with_wait(#[case] seed: Seed) { - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new().build(); let join_handle = tokio::spawn({ let shutdown_trigger = manager.make_shutdown_trigger(); @@ -1353,16 +1127,7 @@ async fn multiple_jobs_with_wait(#[case] seed: Seed) { shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate, - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let block_production = blockprod_setup.make_blockprod_builder().build(); let mut rng = make_seedable_rng(seed); let jobs_to_create = rng.random_range(1..=20); diff --git a/blockprod/src/detail/tests/produce_block/tx_selection_mtp.rs b/blockprod/src/detail/tests/produce_block/tx_selection_mtp.rs index e57ab581f..342615ad4 100644 --- a/blockprod/src/detail/tests/produce_block/tx_selection_mtp.rs +++ b/blockprod/src/detail/tests/produce_block/tx_selection_mtp.rs @@ -23,7 +23,7 @@ use common::{ Uint256, chain::{ ChainConfig, CoinUnit, ConsensusUpgrade, Destination, Genesis, NetUpgrades, - OutPointSourceId, PoolId, TxOutput, + OutPointSourceId, TxOutput, block::timestamp::BlockTimestamp, config::{Builder, ChainType}, output_value::OutputValue, @@ -31,25 +31,20 @@ use common::{ timelock::OutputTimeLock, transaction::TxInput, }, - primitives::{Amount, BlockHeight, H256, Idable}, + primitives::{Amount, BlockHeight, Idable}, time_getter::TimeGetter, }; -use consensus::{PoSGenerateBlockInputData, PoWGenerateBlockInputData}; +use consensus::PoWGenerateBlockInputData; use mempool::{TxOptions, tx_accumulator::PackingStrategy, tx_origin::LocalTxOrigin}; use test_utils::{ - mock_time_getter::mocked_time_getter_seconds, + BasicTestTimeGetter, random::{Seed, make_seedable_rng}, }; -use utils::{atomics::SeqCstAtomicU64, once_destructor::OnceDestructor}; +use utils::once_destructor::OnceDestructor; use crate::{ - BlockProduction, detail::{GenerateBlockInputData, tests::produce_block::assert_job_count}, - prepare_thread_pool, test_blockprod_config, - tests::helpers::{ - assert_process_block, build_chain_config_for_pos, make_genesis_timestamp, - setup_blockprod_test, setup_pos, - }, + tests::helpers::{BlockprodTestSetupBuilder, PoSTestSetupBuilder, make_genesis_timestamp}, }; // The height at which the transaction_selection_mtp_xxx tests will create their test block. @@ -76,13 +71,15 @@ const_assert!(TRANSACTION_SELECTION_MTP_TESTS_BLOCK_HEIGHT > chainstate::MEDIAN_ // b) The block contains the main tx and all dependent txs up to and including the one at // the "median time past" time. async fn transaction_selection_mtp_test_impl( - chain_config: ChainConfig, + chain_config: Arc, input_data: GenerateBlockInputData, time_getter: TimeGetter, genesis_premint_output_index: u32, ) { - let (manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(Some(chain_config), time_getter.clone()); + let (blockprod_setup, manager) = BlockprodTestSetupBuilder::new() + .with_chain_config(Arc::clone(&chain_config)) + .with_time_getter(time_getter.clone()) + .build(); let genesis_timestamp = chain_config.genesis_block().timestamp(); let expected_median_time_past = genesis_timestamp.add_int_seconds(9).unwrap(); @@ -95,16 +92,7 @@ async fn transaction_selection_mtp_test_impl( shutdown_trigger.initiate(); }); - let block_production = BlockProduction::new( - chain_config.clone(), - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool.clone(), - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .unwrap(); + let block_production = blockprod_setup.make_blockprod_builder().build(); for i in 1..TRANSACTION_SELECTION_MTP_TESTS_BLOCK_HEIGHT { let (new_block, job_finished_receiver) = block_production @@ -123,10 +111,11 @@ async fn transaction_selection_mtp_test_impl( assert_eq!(new_block.timestamp(), expected_timestamp); assert_job_count(&block_production, 0).await; - assert_process_block(&chainstate, &mempool, new_block).await; + blockprod_setup.assert_process_block(new_block).await; } - let median_time_past = chainstate + let median_time_past = blockprod_setup + .chainstate .call(|cs| cs.calculate_median_time_past(&cs.get_best_block_id().unwrap())) .await .unwrap() @@ -182,7 +171,8 @@ async fn transaction_selection_mtp_test_impl( txs }; - mempool + blockprod_setup + .mempool .call_mut({ let dependent_txs = dependent_txs.clone(); |mp| { @@ -221,7 +211,7 @@ async fn transaction_selection_mtp_test_impl( assert_job_count(&block_production, 0).await; // First ensure that the produced block is actually correct. - assert_process_block(&chainstate, &mempool, new_block).await; + blockprod_setup.assert_process_block(new_block).await; // Now check the transaction ids. let expected_tx_ids = dependent_txs[..=timestamp_offsets_count as usize] @@ -243,40 +233,23 @@ async fn transaction_selection_mtp_test_impl( #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn transaction_selection_mtp_test_pos(#[case] seed: Seed) { let mut rng = make_seedable_rng(seed); - let initial_time_value_secs = TimeGetter::default().get_time().as_secs_since_epoch(); - let initial_time_value = Arc::new(SeqCstAtomicU64::new(initial_time_value_secs)); - let time_getter = mocked_time_getter_seconds(Arc::clone(&initial_time_value)); + let time_getter = BasicTestTimeGetter::new().get_time_getter(); - let extra_genesis_txs = [TxOutput::Transfer( + let extra_genesis_txos = vec![TxOutput::Transfer( OutputValue::Coin(Amount::from_atoms(1000 * CoinUnit::ATOMS_PER_COIN)), Destination::AnyoneCanSpend, )]; - let ( - chain_config_builder, - genesis_stake_private_key, - genesis_vrf_private_key, - create_genesis_pool_txoutput, - ) = setup_pos( - &time_getter, - BlockHeight::new(TRANSACTION_SELECTION_MTP_TESTS_BLOCK_HEIGHT as u64), - &extra_genesis_txs, - &mut rng, - ); - let chain_config = build_chain_config_for_pos(chain_config_builder); - - let input_data = GenerateBlockInputData::PoS(Box::new(PoSGenerateBlockInputData::new( - genesis_stake_private_key, - genesis_vrf_private_key, - PoolId::new(H256::zero()), - vec![TxInput::from_utxo( - OutPointSourceId::BlockReward(chain_config.genesis_block_id()), - 0, - )], - vec![create_genesis_pool_txoutput], - ))); + let pos_setup = PoSTestSetupBuilder::new() + .with_extra_genesis_txos(extra_genesis_txos) + .with_pos_switch_height(BlockHeight::new( + TRANSACTION_SELECTION_MTP_TESTS_BLOCK_HEIGHT as u64, + )) + .build(make_genesis_timestamp(&time_getter, &mut rng), &mut rng); + + let input_data = pos_setup.make_first_pos_block_input_data(); - transaction_selection_mtp_test_impl(chain_config, input_data, time_getter, 1).await; + transaction_selection_mtp_test_impl(pos_setup.chain_config, input_data, time_getter, 1).await; } #[rstest] @@ -285,9 +258,7 @@ async fn transaction_selection_mtp_test_pos(#[case] seed: Seed) { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn transaction_selection_mtp_test_pow(#[case] seed: Seed) { let mut rng = make_seedable_rng(seed); - let initial_time_value_secs = TimeGetter::default().get_time().as_secs_since_epoch(); - let initial_time_value = Arc::new(SeqCstAtomicU64::new(initial_time_value_secs)); - let time_getter = mocked_time_getter_seconds(Arc::clone(&initial_time_value)); + let time_getter = BasicTestTimeGetter::new().get_time_getter(); let extra_genesis_txs = vec![TxOutput::Transfer( OutputValue::Coin(Amount::from_atoms(1000 * CoinUnit::ATOMS_PER_COIN)), @@ -313,10 +284,12 @@ async fn transaction_selection_mtp_test_pow(#[case] seed: Seed) { ]) .unwrap(); - Builder::new(ChainType::Regtest) - .genesis_custom(genesis) - .consensus_upgrades(net_upgrades) - .build() + Arc::new( + Builder::new(ChainType::Regtest) + .genesis_custom(genesis) + .consensus_upgrades(net_upgrades) + .build(), + ) }; let input_data = GenerateBlockInputData::PoW(Box::new(PoWGenerateBlockInputData::new( @@ -332,9 +305,7 @@ async fn transaction_selection_mtp_test_pow(#[case] seed: Seed) { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn transaction_selection_mtp_test_ignore_consensus(#[case] seed: Seed) { let mut rng = make_seedable_rng(seed); - let initial_time_value_secs = TimeGetter::default().get_time().as_secs_since_epoch(); - let initial_time_value = Arc::new(SeqCstAtomicU64::new(initial_time_value_secs)); - let time_getter = mocked_time_getter_seconds(Arc::clone(&initial_time_value)); + let time_getter = BasicTestTimeGetter::new().get_time_getter(); let extra_genesis_txs = vec![TxOutput::Transfer( OutputValue::Coin(Amount::from_atoms(1000 * CoinUnit::ATOMS_PER_COIN)), @@ -355,10 +326,12 @@ async fn transaction_selection_mtp_test_ignore_consensus(#[case] seed: Seed) { )]) .unwrap(); - Builder::new(ChainType::Regtest) - .genesis_custom(genesis) - .consensus_upgrades(net_upgrades) - .build() + Arc::new( + Builder::new(ChainType::Regtest) + .genesis_custom(genesis) + .consensus_upgrades(net_upgrades) + .build(), + ) }; transaction_selection_mtp_test_impl(chain_config, GenerateBlockInputData::None, time_getter, 0) diff --git a/blockprod/src/detail/tests/stop_jobs.rs b/blockprod/src/detail/tests/stop_jobs.rs index 2e32cdb52..aaee3085e 100644 --- a/blockprod/src/detail/tests/stop_jobs.rs +++ b/blockprod/src/detail/tests/stop_jobs.rs @@ -13,22 +13,21 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::sync::Arc; - use rstest::rstest; -use common::time_getter::TimeGetter; use randomness::RngExt as _; -use test_utils::random::{Seed, make_seedable_rng}; +use test_utils::{ + assert_matches, + random::{Seed, make_seedable_rng}, +}; use crate::{ - BlockProduction, BlockProductionError, JobKey, + BlockProductionError, JobKey, detail::{ CustomId, job_manager::{JobManagerError, tests::MockJobManager}, }, - prepare_thread_pool, test_blockprod_config, - tests::helpers::setup_blockprod_test, + tests::helpers::BlockprodTestSetupBuilder, }; mod stop_all_jobs { @@ -36,19 +35,9 @@ mod stop_all_jobs { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn error() { - let (_manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); - - let mut block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let (blockprod_setup, _manager) = BlockprodTestSetupBuilder::new().build(); + + let mut block_production = blockprod_setup.make_blockprod_builder().build(); let mut mock_job_manager = Box::::default(); @@ -61,10 +50,7 @@ mod stop_all_jobs { let result = block_production.stop_all_jobs().await; - match result { - Err(BlockProductionError::JobManagerError(_)) => {} - _ => panic!("Unexpected return value"), - } + assert_matches!(result, Err(BlockProductionError::JobManagerError(_))); } #[rstest] @@ -72,21 +58,11 @@ mod stop_all_jobs { #[case(Seed::from_entropy())] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn ok(#[case] seed: Seed) { - let (_manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, _manager) = BlockprodTestSetupBuilder::new().build(); let mut rng = make_seedable_rng(seed); - let mut block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate, - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let mut block_production = blockprod_setup.make_blockprod_builder().build(); let (_other_job_key, _other_last_used_block_timestamp, _other_job_cancel_receiver) = block_production @@ -114,19 +90,9 @@ mod stop_all_jobs { #[case(Seed::from_entropy())] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn mocked_ok(#[case] seed: Seed) { - let (_manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); - - let mut block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let (blockprod_setup, _manager) = BlockprodTestSetupBuilder::new().build(); + + let mut block_production = blockprod_setup.make_blockprod_builder().build(); let mut mock_job_manager = Box::::default(); let return_value = make_seedable_rng(seed).random_range(0..=usize::MAX); @@ -153,19 +119,9 @@ mod stop_job { #[case(Seed::from_entropy())] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn error(#[case] seed: Seed) { - let (_manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); - - let mut block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let (blockprod_setup, _manager) = BlockprodTestSetupBuilder::new().build(); + + let mut block_production = blockprod_setup.make_blockprod_builder().build(); let mut mock_job_manager = Box::::default(); @@ -181,10 +137,7 @@ mod stop_job { let result = block_production.stop_job(job_key).await; - match result { - Err(BlockProductionError::JobManagerError(_)) => {} - _ => panic!("Unexpected return value"), - } + assert_matches!(result, Err(BlockProductionError::JobManagerError(_))); } #[rstest] @@ -192,21 +145,11 @@ mod stop_job { #[case(Seed::from_entropy())] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn existing_job_ok(#[case] seed: Seed) { - let (_manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, _manager) = BlockprodTestSetupBuilder::new().build(); let mut rng = make_seedable_rng(seed); - let mut block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate, - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let mut block_production = blockprod_setup.make_blockprod_builder().build(); let (_other_job_key, _other_last_used_block_timestamp, _other_job_cancel_receiver) = block_production @@ -234,21 +177,11 @@ mod stop_job { #[case(Seed::from_entropy())] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn multiple_jobs_ok(#[case] seed: Seed) { - let (_manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, _manager) = BlockprodTestSetupBuilder::new().build(); let mut rng = make_seedable_rng(seed); - let mut block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate, - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let mut block_production = blockprod_setup.make_blockprod_builder().build(); let mut job_keys = Vec::new(); let jobs_to_create = rng.random_range(1..=20); @@ -291,21 +224,11 @@ mod stop_job { #[case(Seed::from_entropy())] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn non_existent_job_ok(#[case] seed: Seed) { - let (_manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); + let (blockprod_setup, _manager) = BlockprodTestSetupBuilder::new().build(); let mut rng = make_seedable_rng(seed); - let mut block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate, - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let mut block_production = blockprod_setup.make_blockprod_builder().build(); let (_other_job_key, _other_last_used_block_timestamp, _other_job_cancel_receiver) = block_production @@ -328,19 +251,9 @@ mod stop_job { #[case(Seed::from_entropy())] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn mocked_ok(#[case] seed: Seed) { - let (_manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, TimeGetter::default()); - - let mut block_production = BlockProduction::new( - chain_config, - Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool, - p2p, - Default::default(), - prepare_thread_pool(1), - ) - .expect("Error initializing blockprod"); + let (blockprod_setup, _manager) = BlockprodTestSetupBuilder::new().build(); + + let mut block_production = blockprod_setup.make_blockprod_builder().build(); let mut mock_job_manager = Box::::default(); diff --git a/blockprod/src/tests/helpers.rs b/blockprod/src/tests/helpers.rs index 99db18aaa..c0d377760 100644 --- a/blockprod/src/tests/helpers.rs +++ b/blockprod/src/tests/helpers.rs @@ -27,8 +27,8 @@ use chainstate_storage::inmemory::Store; use common::{ Uint256, Uint512, chain::{ - self, Block, ConsensusUpgrade, Destination, Genesis, NetUpgrades, PoSChainConfigBuilder, - TxOutput, + self, Block, ConsensusUpgrade, Destination, Genesis, NetUpgrades, OutPointSourceId, + PoSChainConfigBuilder, PoolId, TxInput, TxOutput, block::timestamp::BlockTimestamp, config::{ChainConfig, ChainType, create_unit_test_config}, pos_initial_difficulty, @@ -37,7 +37,10 @@ use common::{ primitives::{Amount, BlockHeight, H256, Idable, per_thousand::PerThousand}, time_getter::{MonotonicTimeGetter, TimeGetter}, }; -use consensus::{calculate_effective_pool_balance, compact_target_to_target}; +use consensus::{ + GenerateBlockInputData, PoSGenerateBlockInputData, calculate_effective_pool_balance, + compact_target_to_target, +}; use crypto::{ key::{KeyKind, PrivateKey}, vrf::{VRFKeyKind, VRFPrivateKey}, @@ -50,6 +53,273 @@ use randomness::{CryptoRng, Rng, RngExt as _}; use storage_inmemory::InMemory; use subsystem::Manager; +use crate::{ + config::BlockProdConfig, detail::BlockProduction, prepare_thread_pool, test_blockprod_config, +}; + +/// A collection of objects needed to run a blockprod test. +/// +/// Note that the subsystem manager is not a part of it; this is because in most tests +/// the test setup object will be moved into a separate tokio task where BlockProduction +/// will be created and tested, while the manager object has to remain outside the task. +pub struct BlockprodTestSetup { + pub chain_config: Arc, + pub time_getter: TimeGetter, + pub chainstate: ChainstateHandle, + pub mempool: MempoolHandle, + pub p2p: P2pHandle, +} + +impl BlockprodTestSetup { + pub async fn assert_process_block(&self, new_block: Block) -> BlockIndex { + assert_process_block(&self.chainstate, &self.mempool, new_block).await + } + + pub fn make_blockprod_builder(&self) -> TestBlockProdBuilder<'_> { + TestBlockProdBuilder { + test_setup: self, + blockprod_config: None, + chainstate: None, + mempool: None, + } + } +} + +/// A builder that produces `BlockprodTestSetup` and the subsystem manager. +pub struct BlockprodTestSetupBuilder { + chain_config: Option>, + time_getter: Option, +} + +impl BlockprodTestSetupBuilder { + pub fn new() -> Self { + Self { + chain_config: None, + time_getter: None, + } + } + + pub fn with_chain_config(mut self, chain_config: Arc) -> Self { + self.chain_config = Some(chain_config); + self + } + + pub fn with_time_getter(mut self, time_getter: TimeGetter) -> Self { + self.time_getter = Some(time_getter); + self + } + + pub fn build(self) -> (BlockprodTestSetup, Manager) { + let chain_config = self.chain_config.unwrap_or_else(|| Arc::new(create_unit_test_config())); + let time_getter = self.time_getter.unwrap_or_default(); + + let manager_config = + subsystem::ManagerConfig::new("blockprod-unit-test").enable_signal_handlers(); + let mut manager = Manager::new_with_config(manager_config); + + let chainstate_config = ChainstateConfig { + max_tip_age: Duration::from_secs(60 * 60 * 24 * 365 * 100).into(), + // There is at least one long test in blockprod that gets significantly slowed down + // by the heavy checks in chainstate. But since the checks are not very useful in blockprod + // tests in general, we disable them globally. + enable_heavy_checks: Some(false), + + max_db_commit_attempts: Default::default(), + enable_db_reckless_mode_in_ibd: Default::default(), + max_orphan_blocks: Default::default(), + allow_checkpoints_mismatch: Default::default(), + }; + + let mempool_config = MempoolConfig::new(); + + let chainstate = chainstate::make_chainstate( + Arc::clone(&chain_config), + chainstate_config, + Store::new_empty().unwrap(), + DefaultTransactionVerificationStrategy::new(), + None, + time_getter.clone(), + None, + ) + .unwrap(); + + let chainstate = manager.add_subsystem("chainstate", chainstate); + + let mempool_init = MempoolInit::new( + Arc::clone(&chain_config), + mempool_config, + subsystem::Handle::clone(&chainstate), + time_getter.clone(), + ) + .unwrap(); + let mempool = manager.add_custom_subsystem("mempool", |hdl, _| mempool_init.init(hdl)); + + let mut p2p_config = test_p2p_config(); + p2p_config.bind_addresses = vec![SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0).into()]; + + let p2p = p2p::make_p2p( + true, + Arc::clone(&chain_config), + Arc::new(p2p_config), + subsystem::Handle::clone(&chainstate), + mempool.clone(), + time_getter.clone(), + MonotonicTimeGetter::default(), + PeerDbStorageImpl::new(InMemory::new()).unwrap(), + ) + .unwrap() + .add_to_manager("p2p", &mut manager); + + ( + BlockprodTestSetup { + chain_config, + time_getter, + chainstate, + mempool, + p2p, + }, + manager, + ) + } +} + +/// A builder that produces `BlockProduction` from `BlockprodTestSetup` with some optional overrides. +pub struct TestBlockProdBuilder<'a> { + test_setup: &'a BlockprodTestSetup, + blockprod_config: Option, + chainstate: Option, + mempool: Option, +} + +impl<'a> TestBlockProdBuilder<'a> { + pub fn with_blockprod_config(mut self, blockprod_config: BlockProdConfig) -> Self { + self.blockprod_config = Some(blockprod_config); + self + } + + pub fn with_chainstate(mut self, chainstate: ChainstateHandle) -> Self { + self.chainstate = Some(chainstate); + self + } + + pub fn with_mempool(mut self, mempool: MempoolHandle) -> Self { + self.mempool = Some(mempool); + self + } + + pub fn build(self) -> BlockProduction { + let blockprod_config = self.blockprod_config.unwrap_or_else(test_blockprod_config); + let chainstate = self.chainstate.unwrap_or_else(|| self.test_setup.chainstate.clone()); + let mempool = self.mempool.unwrap_or_else(|| self.test_setup.mempool.clone()); + + BlockProduction::new( + Arc::clone(&self.test_setup.chain_config), + Arc::new(blockprod_config), + chainstate, + mempool, + self.test_setup.p2p.clone(), + self.test_setup.time_getter.clone(), + prepare_thread_pool(1), + ) + .unwrap() + } +} + +/// A bunch of data specific to PoS tests. +pub struct PoSTestSetup { + pub chain_config: Arc, + pub genesis_stake_private_key: PrivateKey, + pub genesis_vrf_private_key: VRFPrivateKey, + pub create_genesis_pool_utxo: TxOutput, +} + +impl PoSTestSetup { + pub fn make_first_pos_block_input_data(&self) -> GenerateBlockInputData { + GenerateBlockInputData::PoS(Box::new(PoSGenerateBlockInputData::new( + self.genesis_stake_private_key.clone(), + self.genesis_vrf_private_key.clone(), + PoolId::new(H256::zero()), + vec![TxInput::from_utxo( + OutPointSourceId::BlockReward(self.chain_config.genesis_block_id()), + 0, + )], + vec![self.create_genesis_pool_utxo.clone()], + ))) + } +} + +/// A builder that produces `PoSTestSetup`. +pub struct PoSTestSetupBuilder { + chain_config_builder: Option, + extra_genesis_txos: Vec, + pos_switch_height: BlockHeight, +} + +impl PoSTestSetupBuilder { + pub fn new() -> Self { + Self { + chain_config_builder: None, + extra_genesis_txos: Vec::new(), + pos_switch_height: BlockHeight::new(1), + } + } + + pub fn with_chain_config_builder( + mut self, + chain_config_builder: chain::config::Builder, + ) -> Self { + self.chain_config_builder = Some(chain_config_builder); + self + } + + pub fn with_extra_genesis_txos(mut self, extra_genesis_txos: Vec) -> Self { + self.extra_genesis_txos = extra_genesis_txos; + self + } + + pub fn with_pos_switch_height(mut self, pos_switch_height: BlockHeight) -> Self { + self.pos_switch_height = pos_switch_height; + self + } + + pub fn build( + self, + genesis_timestamp: BlockTimestamp, + rng: &mut impl CryptoRng, + ) -> PoSTestSetup { + let initial_target = pos_initial_difficulty(ChainType::Regtest); + + let (genesis, genesis_stake_private_key, genesis_vrf_private_key, create_genesis_pool_utxo) = + create_genesis_for_pos_tests(genesis_timestamp, &self.extra_genesis_txos, rng); + + let net_upgrades = NetUpgrades::initialize(vec![ + (BlockHeight::new(0), ConsensusUpgrade::IgnoreConsensus), + ( + self.pos_switch_height, + ConsensusUpgrade::PoS { + initial_difficulty: Some(initial_target.into()), + config: PoSChainConfigBuilder::new_for_unit_test().build(), + }, + ), + ]) + .unwrap(); + + let chain_config_builder = self + .chain_config_builder + .unwrap_or_else(make_chain_config_builder) + .genesis_custom(genesis) + .consensus_upgrades(net_upgrades); + let chain_config = Arc::new(build_chain_config_for_pos(chain_config_builder)); + + PoSTestSetup { + chain_config, + genesis_stake_private_key, + genesis_vrf_private_key, + create_genesis_pool_utxo, + } + } +} + pub async fn assert_process_block( chainstate: &ChainstateHandle, mempool: &MempoolHandle, @@ -57,8 +327,10 @@ pub async fn assert_process_block( ) -> BlockIndex { let block_id = new_block.get_id(); - // Wait for mempool to be up-to-date with the new block. The subscriptions are not cleaned - // up but hopefully it's not too bad just for testing. + // Subscribe to mempool events, so that we can wait for it to become up-to-date with + // the new block. + // Note that currently we don't have a mechanism to remove a subscription, so "dead" event + // handlers will accumulate each time this function is called. But it's not a big deal in tests. let (tip_sx, tip_rx) = tokio::sync::oneshot::channel(); let tip_sx = utils::sync::Mutex::new(Some(tip_sx)); mempool @@ -68,7 +340,7 @@ pub async fn assert_process_block( mempool::event::MempoolEvent::NewTip(tip) => { if let Some(tip_sx) = tip_sx.lock().unwrap().take() { assert_eq!(tip.block_id(), &block_id); - tip_sx.send(()).unwrap(); + let _ = tip_sx.send(()); } } mempool::event::MempoolEvent::TransactionProcessed(_) => (), @@ -80,108 +352,33 @@ pub async fn assert_process_block( let block_index = chainstate .call_mut(move |this| { - let new_block_index = this - .process_block(new_block.clone(), BlockSource::Local) - .expect("Failed to process block") - .expect("Failed to activate best chain"); + let new_block_index = + this.process_block(new_block.clone(), BlockSource::Local).unwrap().unwrap(); assert_eq!( new_block.header().header().block_id(), *new_block_index.block_id(), - "The new block's Id is different to the new block index's block Id", + "The new block's id is different from the new block index's block id", ); - let best_block_index = - this.get_best_block_index().expect("Failed to get best block index"); + let best_block_index = this.get_best_block_index().unwrap(); assert_eq!( new_block_index.clone().into_gen_block_index().block_id(), best_block_index.block_id(), - "The new block index not the best block index" + "The new block index is not the best block index" ); new_block_index }) .await - .expect("New block is not the new tip"); + .unwrap(); tip_rx.await.unwrap(); block_index } -pub fn setup_blockprod_test( - chain_config: Option, - time_getter: TimeGetter, -) -> ( - Manager, - Arc, - ChainstateHandle, - MempoolHandle, - P2pHandle, -) { - let manager_config = - subsystem::ManagerConfig::new("blockprod-unit-test").enable_signal_handlers(); - let mut manager = Manager::new_with_config(manager_config); - - let chain_config = Arc::new(chain_config.unwrap_or_else(create_unit_test_config)); - - let chainstate_config = ChainstateConfig { - max_tip_age: Duration::from_secs(60 * 60 * 24 * 365 * 100).into(), - // There is at least one long test in blockprod that gets significantly slowed down - // by the heavy checks in chainstate. But since the checks are not very useful in blockprod - // tests in general, we disable them globally. - enable_heavy_checks: Some(false), - - max_db_commit_attempts: Default::default(), - enable_db_reckless_mode_in_ibd: Default::default(), - max_orphan_blocks: Default::default(), - allow_checkpoints_mismatch: Default::default(), - }; - - let mempool_config = MempoolConfig::new(); - - let chainstate = chainstate::make_chainstate( - Arc::clone(&chain_config), - chainstate_config, - Store::new_empty().expect("Error initializing empty store"), - DefaultTransactionVerificationStrategy::new(), - None, - time_getter.clone(), - None, - ) - .expect("Error initializing chainstate"); - - let chainstate = manager.add_subsystem("chainstate", chainstate); - - let mempool_init = MempoolInit::new( - Arc::clone(&chain_config), - mempool_config, - subsystem::Handle::clone(&chainstate), - time_getter.clone(), - ) - .unwrap(); - let mempool = manager.add_custom_subsystem("mempool", |hdl, _| mempool_init.init(hdl)); - - let mut p2p_config = test_p2p_config(); - p2p_config.bind_addresses = vec![SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0).into()]; - - let p2p = p2p::make_p2p( - true, - Arc::clone(&chain_config), - Arc::new(p2p_config), - subsystem::Handle::clone(&chainstate), - mempool.clone(), - time_getter, - MonotonicTimeGetter::default(), - PeerDbStorageImpl::new(InMemory::new()).unwrap(), - ) - .expect("P2p initialization was successful") - .add_to_manager("p2p", &mut manager); - - (manager, chain_config, chainstate, mempool, p2p) -} - pub fn make_genesis_timestamp(time_getter: &TimeGetter, rng: &mut impl Rng) -> BlockTimestamp { BlockTimestamp::from_int_seconds( (time_getter.get_time() @@ -190,7 +387,7 @@ pub fn make_genesis_timestamp(time_getter: &TimeGetter, rng: &mut impl Rng) -> B rng.random_range(60 * 60 * 24..60 * 60 * 24 * 14), 0, )) - .expect("No time underflow") + .unwrap() .as_secs_since_epoch(), ) } @@ -222,7 +419,7 @@ pub fn ensure_reasonable_initial_target_for_pos_tests( pub fn create_genesis_for_pos_tests( timestamp: BlockTimestamp, - extra_txs: &[TxOutput], + extra_txos: &[TxOutput], rng: &mut impl CryptoRng, ) -> ( Genesis, @@ -250,16 +447,16 @@ pub fn create_genesis_for_pos_tests( Destination::PublicKey(stake_public_key.clone()), vrf_public_key, Destination::PublicKey(stake_public_key), - PerThousand::new(1000).expect("Valid per thousand"), + PerThousand::new(1000).unwrap(), Amount::ZERO, )), ) }; - let mut txs = vec![create_pool_txoutput.clone()]; - txs.extend_from_slice(extra_txs); + let mut txos = vec![create_pool_txoutput.clone()]; + txos.extend_from_slice(extra_txos); - let genesis = Genesis::new("blockprod-testing".into(), timestamp, txs); + let genesis = Genesis::new("blockprod-testing".into(), timestamp, txos); ( genesis, @@ -269,49 +466,8 @@ pub fn create_genesis_for_pos_tests( ) } -pub fn setup_pos( - time_getter: &TimeGetter, - switch_to_pos_at: BlockHeight, - extra_genesis_txs: &[TxOutput], - rng: &mut impl CryptoRng, -) -> (chain::config::Builder, PrivateKey, VRFPrivateKey, TxOutput) { - let genesis_timestamp = make_genesis_timestamp(time_getter, rng); - setup_pos_with_genesis_timestamp(genesis_timestamp, switch_to_pos_at, extra_genesis_txs, rng) -} - -pub fn setup_pos_with_genesis_timestamp( - genesis_timestamp: BlockTimestamp, - switch_to_pos_at: BlockHeight, - extra_genesis_txs: &[TxOutput], - rng: &mut impl CryptoRng, -) -> (chain::config::Builder, PrivateKey, VRFPrivateKey, TxOutput) { - let initial_target = pos_initial_difficulty(ChainType::Regtest); - - let (genesis, genesis_stake_private_key, genesis_vrf_private_key, create_genesis_pool_txoutput) = - create_genesis_for_pos_tests(genesis_timestamp, extra_genesis_txs, rng); - - let net_upgrades = NetUpgrades::initialize(vec![ - (BlockHeight::new(0), ConsensusUpgrade::IgnoreConsensus), - ( - switch_to_pos_at, - ConsensusUpgrade::PoS { - initial_difficulty: Some(initial_target.into()), - config: PoSChainConfigBuilder::new_for_unit_test().build(), - }, - ), - ]) - .expect("Net upgrades are valid"); - - let chain_config_builder = chain::config::Builder::new(ChainType::Regtest) - .genesis_custom(genesis) - .consensus_upgrades(net_upgrades); - - ( - chain_config_builder, - genesis_stake_private_key, - genesis_vrf_private_key, - create_genesis_pool_txoutput, - ) +pub fn make_chain_config_builder() -> chain::config::Builder { + chain::config::Builder::new(ChainType::Regtest) } pub fn build_chain_config_for_pos(builder: chain::config::Builder) -> ChainConfig { diff --git a/blockprod/src/tests/mod.rs b/blockprod/src/tests/mod.rs index 40979c326..f100a53e8 100644 --- a/blockprod/src/tests/mod.rs +++ b/blockprod/src/tests/mod.rs @@ -17,25 +17,23 @@ pub mod helpers; use std::sync::Arc; -use common::time_getter::TimeGetter; - -use crate::{make_blockproduction, test_blockprod_config, tests::helpers::setup_blockprod_test}; +use crate::{ + make_blockproduction, test_blockprod_config, tests::helpers::BlockprodTestSetupBuilder, +}; #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn test_make_blockproduction() { - let time_getter = TimeGetter::default(); - let (mut manager, chain_config, chainstate, mempool, p2p) = - setup_blockprod_test(None, time_getter.clone()); + let (blockprod_setup, mut manager) = BlockprodTestSetupBuilder::new().build(); let blockprod = make_blockproduction( - Arc::clone(&chain_config), + Arc::clone(&blockprod_setup.chain_config), Arc::new(test_blockprod_config()), - chainstate.clone(), - mempool.clone(), - p2p.clone(), - time_getter, + blockprod_setup.chainstate.clone(), + blockprod_setup.mempool.clone(), + blockprod_setup.p2p.clone(), + blockprod_setup.time_getter, ) - .expect("Error initializing blockprod"); + .unwrap(); let blockprod = manager.add_direct_subsystem("blockprod", blockprod); let shutdown = manager.make_shutdown_trigger(); @@ -55,7 +53,7 @@ async fn test_make_blockproduction() { }) }) .await - .expect("Error initializing block production"); + .unwrap(); }); manager.main().await; diff --git a/mempool/src/pool/mod.rs b/mempool/src/pool/mod.rs index 35e26a673..5a0d23413 100644 --- a/mempool/src/pool/mod.rs +++ b/mempool/src/pool/mod.rs @@ -372,7 +372,7 @@ impl Mempool { log::trace!("Performing orphan processing work"); let orphan = state.work_queue.pick(|peer, orphan_id| { - log::debug!("Processing orphan tx {orphan_id:?} coming from peer {peer}"); + log::debug!("Processing orphan tx {orphan_id:x} coming from peer {peer}"); match state.orphans.entry(&orphan_id) { Some(orphan) if orphan.is_ready() => { @@ -386,7 +386,7 @@ impl Mempool { None => { // The orphan may have been kicked out of the pool in the meantime. // Return with `None` in that case to indicate we're not really doing any work. - log::debug!("Orphan tx {orphan_id:?} no longer in the pool"); + log::debug!("Orphan tx {orphan_id:x} no longer in the pool"); None } } @@ -396,12 +396,12 @@ impl Mempool { Some(Ok(orphan)) => { let orphan = orphan.map_origin(TxOrigin::from); let orphan_id = *orphan.tx_id(); - log::trace!("Re-processing orphan transaction {orphan_id:?}"); + log::trace!("Re-processing orphan transaction {orphan_id:x}"); if let Err(err) = state.add_transaction(orphan) { - log::debug!("Orphan transaction {orphan_id:?} evicted: {err}"); + log::debug!("Orphan transaction {orphan_id:x} evicted: {err}"); } } - Some(Err(orphan_id)) => log::trace!("Orphan tx {orphan_id:?} not ready"), + Some(Err(orphan_id)) => log::trace!("Orphan tx {orphan_id:x} not ready"), None => log::trace!("No orphan processing work left to do"), } } @@ -681,7 +681,7 @@ impl<'a> TxFinalizer<'a> { let orphan_id = *orphan.tx_id(); let peer_id = orphan.origin().peer_id(); if self.work_queue.insert(peer_id, orphan_id) { - log::trace!("Added orphan {orphan_id:?} to peer{peer_id}'s work queue"); + log::trace!("Added orphan {orphan_id:x} to peer{peer_id}'s work queue"); } } } diff --git a/mempool/src/pool/tx_pool/mod.rs b/mempool/src/pool/tx_pool/mod.rs index f877902f2..2497ad088 100644 --- a/mempool/src/pool/tx_pool/mod.rs +++ b/mempool/src/pool/tx_pool/mod.rs @@ -634,7 +634,7 @@ impl TxPool { let expired = now.saturating_sub(creation_time) > self.max_tx_age; if expired { log::trace!( - "Evicting tx {} which was created at {:?}. It is now {:?}", + "Evicting tx {:x} which was created at {:?}. It is now {:?}", tx_id, creation_time, now @@ -665,7 +665,7 @@ impl TxPool { let removed = self.store.txs_by_id.get(&removed_id).expect("tx with id should exist"); log::debug!( - "Mempool trim: Evicting tx {} which has a descendant score of {:?} and has size {}", + "Mempool trim: Evicting tx {:x} which has a descendant score of {:?} and has size {}", removed_id, removed.descendant_score(), removed.size() @@ -697,7 +697,9 @@ impl TxPool { // dependencies. However, it does not appear to be easy to extract that information // from the transaction verifier at the moment. To be addressed in the future. - log::error!("Disconnecting {disc_id} failed with '{err}' during eviction of {tx_id}"); + log::error!( + "Disconnecting {disc_id:x} failed with '{err}' during eviction of {tx_id:x}" + ); if let Err(refresh_err) = reorg::refresh_mempool(self, |_, _| ()) { log::error!("Refreshing mempool failed: {refresh_err}"); @@ -792,7 +794,7 @@ impl TxPool { self.check_preliminary_mempool_policy(&transaction)?; for attempt_no in 1..=config::MAX_TX_ADDITION_ATTEMPTS { - log::trace!("Adding {tx_id:?} attempt #{attempt_no}"); + log::trace!("Adding {tx_id:x} attempt #{attempt_no}"); transaction = match self.try_add_transaction(transaction)? { TxAdditionAttemptOutcome::Added => { let transaction = self.store.get_entry(&tx_id).expect("just added"); @@ -809,7 +811,7 @@ impl TxPool { current_tip, } => { log::debug!( - "Tip moved from {start_tip:?} to {current_tip:?} while verifying {tx_id:?}" + "Tip moved from {start_tip:x} to {current_tip:x} while verifying {tx_id:x}" ); transaction } @@ -872,7 +874,7 @@ impl TxPool { let mut tx_verifier = self.tx_verifier.derive_child(); log::trace!( - "Verifying {tx_id:?}, tip = {start_tip:?}, tx_verifier's best block for utxos = {:?}", + "Verifying {tx_id:x}, tip = {start_tip:x}, tx_verifier's best block for utxos = {:x}", tx_verifier.get_best_block_for_utxos()? ); diff --git a/mempool/src/pool/tx_pool/reorg.rs b/mempool/src/pool/tx_pool/reorg.rs index 8ace67c7e..c54aa1c8d 100644 --- a/mempool/src/pool/tx_pool/reorg.rs +++ b/mempool/src/pool/tx_pool/reorg.rs @@ -117,7 +117,7 @@ fn fetch_disconnected_txs( .get_best_block_for_utxos() .map_err(|_| ReorgError::BestBlockForUtxos)?; - log::debug!("Fetching disconnected txs, old_tip = {old_tip:?}"); + log::debug!("Fetching disconnected txs, old_tip = {old_tip:x}"); let now = tx_pool.clock.get_time(); @@ -148,7 +148,7 @@ pub fn handle_new_tip( if new_tip != actual_tip { log::debug!( - "Not updating mempool because actual tip differs: new_tip = {new_tip:?}, actual_tip = {actual_tip:?}" + "Not updating mempool because actual tip differs: new_tip = {new_tip:x}, actual_tip = {actual_tip:x}" ); return Ok(()); @@ -174,7 +174,7 @@ fn reorg_mempool_transactions( let old_transactions = tx_pool.reset(); log::debug!( - "Reorging mempool txs, tx_verifier's best block for utxos after mempool reset: {:?}", + "Reorging mempool txs, tx_verifier's best block for utxos after mempool reset: {:x}", tx_pool .tx_verifier .get_best_block_for_utxos() @@ -183,18 +183,18 @@ fn reorg_mempool_transactions( for tx in txs_to_insert { let tx_id = *tx.tx_id(); - log::trace!("Adding {tx_id} after reorg"); + log::trace!("Adding {tx_id:x} after reorg"); if let Err(e) = tx_pool.add_transaction(tx, &mut finalizer) { - log::debug!("Disconnected transaction {tx_id:?} no longer validates: {e:?}") + log::debug!("Disconnected transaction {tx_id:x} no longer validates: {e:?}") } } // Re-populate the verifier with transactions from mempool for tx in old_transactions { let tx_id = *tx.tx_id(); - log::trace!("Adding {tx_id} after reorg"); + log::trace!("Adding {tx_id:x} after reorg"); if let Err(e) = tx_pool.add_transaction(tx, &mut finalizer) { - log::debug!("Evicting {tx_id:?} from mempool: {e:?}") + log::debug!("Evicting {tx_id:x} from mempool: {e:?}") } } diff --git a/mempool/src/pool/tx_pool/store/mod.rs b/mempool/src/pool/tx_pool/store/mod.rs index 3a605ffe2..862e33393 100644 --- a/mempool/src/pool/tx_pool/store/mod.rs +++ b/mempool/src/pool/tx_pool/store/mod.rs @@ -543,7 +543,7 @@ impl MempoolStore { tx_id: &Id, reason: MempoolRemovalReason, ) -> Option { - log::info!("remove_tx: {}", tx_id.to_hash()); + log::debug!("remove_tx: {:x}", tx_id.to_hash()); let entry = self.mem_tracker.modify(&mut self.txs_by_id, |by_id, _| by_id.remove(tx_id)); if let Some(entry) = entry { @@ -773,7 +773,7 @@ impl TxMempoolEntry { } pub fn ancestor_score(&self) -> AncestorScore { - log::debug!("ancestor score for {:?}", self.tx_id()); + log::debug!("ancestor score for {:x}", self.tx_id()); log::debug!( "fees with ancestors: {:?}, size_with_ancestors: {}, fee: {:?}, size: {}", self.fees_with_ancestors,