diff --git a/chain/mock/src/mock_chain.rs b/chain/mock/src/mock_chain.rs index 27be3bacef..e69c9ff113 100644 --- a/chain/mock/src/mock_chain.rs +++ b/chain/mock/src/mock_chain.rs @@ -89,7 +89,8 @@ impl MockChain { ) -> Result { let storage = Arc::new(Storage::new(StorageInstance::new_cache_instance())?); let storage2 = Arc::new(Storage2(storage.clone())); - let dag = BlockDAG::create_for_testing_with_parameters(k)?; + let genesis_hash = genesis.block().id(); + let dag = BlockDAG::create_for_testing_with_parameters(k, genesis_hash)?; let chain_info = genesis.execute_genesis_block(&net, storage.clone(), storage2.clone(), dag.clone())?; diff --git a/chain/tests/block_test_utils.rs b/chain/tests/block_test_utils.rs index 8fda3f11ce..02c4987217 100644 --- a/chain/tests/block_test_utils.rs +++ b/chain/tests/block_test_utils.rs @@ -39,8 +39,9 @@ pub fn genesis_strategy(storage: Arc) -> impl Strategy { BuiltinNetworkID::Test.genesis_config2().clone(), ); let genesis = Genesis::load_or_build(&net).unwrap(); + let genesis_hash = genesis.block().id(); let storage2 = Arc::new(Storage2(storage.clone())); - let dag = starcoin_dag::blockdag::BlockDAG::create_for_testing().unwrap(); + let dag = starcoin_dag::blockdag::BlockDAG::create_for_testing(genesis_hash).unwrap(); genesis .execute_genesis_block(&net, storage, storage2, dag) .unwrap(); diff --git a/cmd/generator/src/lib.rs b/cmd/generator/src/lib.rs index 791f15b112..743216fdc1 100644 --- a/cmd/generator/src/lib.rs +++ b/cmd/generator/src/lib.rs @@ -6,12 +6,13 @@ use starcoin_account::account_storage::AccountStorage; use starcoin_account::AccountManager; use starcoin_account_api::AccountInfo; use starcoin_config::{NodeConfig, StarcoinOpt}; +use starcoin_crypto::HashValue; use starcoin_dag::blockdag::BlockDAG; use starcoin_genesis::Genesis; use starcoin_storage::cache_storage::CacheStorage; use starcoin_storage::db_storage::DBStorage; use starcoin_storage::storage::StorageInstance; -use starcoin_storage::{Storage, Storage2}; +use starcoin_storage::{BlockStore, Storage, Storage2}; use starcoin_types::startup_info::ChainInfo; use std::sync::Arc; @@ -49,11 +50,13 @@ pub fn init_or_load_data_dir( .genesis_config2() .consensus_config .base_max_uncles_per_block; - let dag = starcoin_dag::blockdag::BlockDAG::new( + let genesis_hash = storage.get_genesis()?.unwrap_or(HashValue::zero()); + let mut dag = starcoin_dag::blockdag::BlockDAG::new( starcoin_types::blockhash::KType::try_from(k)?, config.miner.dag_merge_depth(), config.miner.maximum_parents_count(), dag_storage.clone(), + genesis_hash, ); let (chain_info, _genesis) = Genesis::init_and_check_storage( config.net(), @@ -61,6 +64,9 @@ pub fn init_or_load_data_dir( dag.clone(), config.data_dir(), )?; + if genesis_hash == HashValue::zero() { + dag.set_genesis(chain_info.genesis_hash()); + } let vault_config = &config.vault; let account_storage = AccountStorage::create_from_path(vault_config.dir(), config.storage.rocksdb_config())?; diff --git a/config/src/txpool_config.rs b/config/src/txpool_config.rs index cf9d8ff9db..be1fdfca7d 100644 --- a/config/src/txpool_config.rs +++ b/config/src/txpool_config.rs @@ -33,6 +33,11 @@ pub struct TxPoolConfig { /// interval(s) of tx propagation timer. default to 2. tx_propagate_interval: Option, + #[serde(skip_serializing_if = "Option::is_none")] + #[clap(name = "txpool-cull-interval", long)] + /// interval(s) of cull expired transactions timer. default to 1. + cull_interval: Option, + #[serde(skip_serializing_if = "Option::is_none")] #[clap(name = "txpool-min-gas-price", long)] /// reject transaction whose gas_price is less than the min_gas_price. default to 1. @@ -80,6 +85,9 @@ impl TxPoolConfig { pub fn tx_propagate_interval(&self) -> u64 { self.tx_propagate_interval.unwrap_or(2) } + pub fn cull_interval(&self) -> u64 { + self.cull_interval.unwrap_or(1).max(1) + } pub fn min_gas_price(&self) -> u64 { self.min_gas_price.unwrap_or(1) } @@ -109,6 +117,9 @@ impl ConfigModule for TxPoolConfig { if let Some(m) = txpool_opt.tx_propagate_interval.as_ref() { self.tx_propagate_interval = Some(*m); } + if let Some(m) = txpool_opt.cull_interval.as_ref() { + self.cull_interval = Some(*m); + } if let Some(m) = txpool_opt.min_gas_price.as_ref() { self.min_gas_price = Some(*m); } diff --git a/flexidag/src/blockdag.rs b/flexidag/src/blockdag.rs index 5a8ab1d4b8..262c78ffb1 100644 --- a/flexidag/src/blockdag.rs +++ b/flexidag/src/blockdag.rs @@ -65,15 +65,22 @@ pub struct BlockDAG { block_depth_manager: BlockDepthManager, max_parents_count: usize, commit_lock: Arc>, + genesis: Hash, } impl BlockDAG { - pub fn create_blockdag(dag_storage: FlexiDagStorage) -> Self { + pub fn create_blockdag(dag_storage: FlexiDagStorage, genesis: Hash) -> Self { // Test defaults: k=8, merge_depth=3600, max_parents=8 (k >= max_parents) - Self::new(8, 3600, 8, dag_storage) + Self::new(8, 3600, 8, dag_storage, genesis) } - pub fn new(k: KType, merge_depth: u64, max_parents_count: usize, db: FlexiDagStorage) -> Self { + pub fn new( + k: KType, + merge_depth: u64, + max_parents_count: usize, + db: FlexiDagStorage, + genesis: Hash, + ) -> Self { // Ensure k >= max_parents_count to prevent protocol violations assert!( k as usize >= max_parents_count, @@ -109,37 +116,36 @@ impl BlockDAG { block_depth_manager, max_parents_count, commit_lock: Arc::new(Mutex::new(db)), + genesis, } } /// For testing only - do not use in production code - pub fn create_for_testing() -> anyhow::Result { + pub fn create_for_testing(genesis: Hash) -> anyhow::Result { let config = FlexiDagStorageConfig { cache_size: 1024, ..Default::default() }; let dag_storage = FlexiDagStorage::create_from_path(temp_dir(), config)?; - // Test defaults: k=8, merge_depth=3600, max_parents=8 (k >= max_parents) - Ok(Self::new(8, 3600, 8, dag_storage)) + Ok(Self::new(8, 3600, 8, dag_storage, genesis)) } /// For testing only - do not use in production code - pub fn create_for_testing_with_parameters(k: KType) -> anyhow::Result { + pub fn create_for_testing_with_parameters(k: KType, genesis: Hash) -> anyhow::Result { let dag_storage = FlexiDagStorage::create_from_path(temp_dir(), FlexiDagStorageConfig::default())?; - // Test defaults: merge_depth=3600, max_parents=3 - Ok(Self::new(k, 3600, 3, dag_storage)) + Ok(Self::new(k, 3600, 3, dag_storage, genesis)) } /// For testing only - do not use in production code pub fn create_for_testing_with_k_and_merge_depth( k: KType, merge_depth: u64, + genesis: Hash, ) -> anyhow::Result { let dag_storage = FlexiDagStorage::create_from_path(temp_dir(), FlexiDagStorageConfig::default())?; - // Test default: max_parents=3 (safe for small k values) - Ok(Self::new(k, merge_depth, 3, dag_storage)) + Ok(Self::new(k, merge_depth, 3, dag_storage, genesis)) } pub fn has_block_connected(&self, block_header: &BlockHeader) -> anyhow::Result { @@ -540,7 +546,24 @@ impl BlockDAG { } pub fn get_dag_state(&self, hash: Hash) -> anyhow::Result { - Ok(self.storage.state_store.read().get_state_by_hash(hash)?) + let query_hash = if hash == Hash::zero() { + self.genesis + } else { + hash + }; + Ok(self + .storage + .state_store + .read() + .get_state_by_hash(query_hash)?) + } + + pub fn genesis(&self) -> Hash { + self.genesis + } + + pub fn set_genesis(&mut self, genesis: Hash) { + self.genesis = genesis; } pub fn save_dag_state_directly(&self, hash: Hash, state: DagState) -> anyhow::Result<()> { diff --git a/flexidag/tests/test_commit_atomicity.rs b/flexidag/tests/test_commit_atomicity.rs index 64d24b66c1..fb87192b59 100644 --- a/flexidag/tests/test_commit_atomicity.rs +++ b/flexidag/tests/test_commit_atomicity.rs @@ -10,13 +10,13 @@ fn test_commit_atomicity() -> Result<()> { let db_tempdir = tempfile::tempdir()?; let config = FlexiDagStorageConfig::new(); let dag_storage = FlexiDagStorage::create_from_path(db_tempdir.path(), config)?; - let mut dag = BlockDAG::new(8, 8, 3, dag_storage); // Create and initialize genesis let genesis = BlockHeaderBuilder::new() .with_number(0) .with_parent_hash(HashValue::zero()) .build(); + let mut dag = BlockDAG::new(8, 8, 3, dag_storage, genesis.id()); // Initialize DAG with genesis dag.init_with_genesis(genesis.clone())?; @@ -68,13 +68,13 @@ fn test_partial_write_detection() -> Result<()> { let db_tempdir = tempfile::tempdir()?; let config = FlexiDagStorageConfig::new(); let dag_storage = FlexiDagStorage::create_from_path(db_tempdir.path(), config)?; - let mut dag = BlockDAG::new(8, 8, 3, dag_storage); // Create and initialize genesis let genesis = BlockHeaderBuilder::new() .with_number(0) .with_parent_hash(HashValue::zero()) .build(); + let mut dag = BlockDAG::new(8, 8, 3, dag_storage, genesis.id()); // Initialize DAG with genesis dag.init_with_genesis(genesis.clone())?; diff --git a/flexidag/tests/tests.rs b/flexidag/tests/tests.rs index 6e1e09f481..693177ade4 100644 --- a/flexidag/tests/tests.rs +++ b/flexidag/tests/tests.rs @@ -34,11 +34,11 @@ use std::{ #[test] fn test_dag_commit() -> Result<()> { - let mut dag = BlockDAG::create_for_testing().unwrap(); let genesis = BlockHeader::random() .as_builder() .with_difficulty(0.into()) .build(); + let mut dag = BlockDAG::create_for_testing(genesis.id()).unwrap(); let mut parents_hash = vec![genesis.id()]; let _origin = dag.init_with_genesis(genesis.clone())?; @@ -94,7 +94,7 @@ fn test_dag_1() -> Result<()> { .build(); let mut latest_id = block6.id(); let genesis_id = genesis.id(); - let mut dag = BlockDAG::create_for_testing().unwrap(); + let mut dag = BlockDAG::create_for_testing(genesis.id()).unwrap(); let expect_selected_parented = [block5.id(), block3.id(), block3_1.id(), genesis_id]; let _origin = dag.init_with_genesis(genesis.clone())?; @@ -147,7 +147,7 @@ async fn test_with_spawn() { .with_difficulty(2.into()) .with_parents_hash(vec![genesis.id()]) .build(); - let mut dag = BlockDAG::create_for_testing().unwrap(); + let mut dag = BlockDAG::create_for_testing(genesis.id()).unwrap(); let _origin = dag.init_with_genesis(genesis.clone()).unwrap(); dag.commit_trusted_block( @@ -200,11 +200,11 @@ async fn test_with_spawn() { #[test] fn test_write_asynchronization() -> anyhow::Result<()> { - let mut dag = BlockDAG::create_for_testing()?; let genesis = BlockHeader::random() .as_builder() .with_difficulty(0.into()) .build(); + let mut dag = BlockDAG::create_for_testing(genesis.id())?; let _real_origin = dag.init_with_genesis(genesis.clone())?; let parent = BlockHeaderBuilder::random() @@ -273,12 +273,12 @@ fn test_write_asynchronization() -> anyhow::Result<()> { #[test] fn test_dag_genesis_fork() { // initialzie the dag firstly - let mut dag = BlockDAG::create_for_testing().unwrap(); - let genesis = BlockHeader::random() .as_builder() .with_difficulty(0.into()) .build(); + let mut dag = BlockDAG::create_for_testing(genesis.id()).unwrap(); + dag.init_with_genesis(genesis.clone()).unwrap(); // normally add the dag blocks @@ -348,7 +348,7 @@ fn test_dag_genesis_fork() { #[test] fn test_dag_tips_store() { - let dag = BlockDAG::create_for_testing().unwrap(); + let dag = BlockDAG::create_for_testing(Hash::random()).unwrap(); let state = DagState { tips: vec![Hash::random()], @@ -372,9 +372,8 @@ fn test_dag_tips_store() { #[test] fn test_dag_multiple_commits() -> anyhow::Result<()> { // initialzie the dag firstly - let mut dag = BlockDAG::create_for_testing().unwrap(); - let genesis = BlockHeader::random(); + let mut dag = BlockDAG::create_for_testing(genesis.id()).unwrap(); dag.init_with_genesis(genesis.clone()).unwrap(); @@ -405,7 +404,7 @@ fn test_dag_multiple_commits() -> anyhow::Result<()> { #[test] fn test_reachability_abort_add_block() -> anyhow::Result<()> { - let dag = BlockDAG::create_for_testing().unwrap(); + let dag = BlockDAG::create_for_testing(Hash::random()).unwrap(); let reachability_store = dag.storage.reachability_store.clone(); let mut parent = Hash::random(); @@ -462,7 +461,7 @@ fn test_reachability_abort_add_block() -> anyhow::Result<()> { #[test] fn test_reachability_check_ancestor() -> anyhow::Result<()> { - let dag = BlockDAG::create_for_testing().unwrap(); + let dag = BlockDAG::create_for_testing(Hash::random()).unwrap(); let reachability_store = dag.storage.reachability_store.clone(); let mut parent = Hash::random(); @@ -588,7 +587,7 @@ fn print_reachability_data(reachability: &DbReachabilityStore, key: &[Hash]) { #[test] fn test_reachability_not_ancestor() -> anyhow::Result<()> { - let dag = BlockDAG::create_for_testing().unwrap(); + let dag = BlockDAG::create_for_testing(Hash::random()).unwrap(); let reachability_store = dag.storage.reachability_store.clone(); let origin = Hash::random(); @@ -655,7 +654,7 @@ fn test_reachability_not_ancestor() -> anyhow::Result<()> { #[test] #[ignore = "maxmum data testing for dev"] fn test_hint_virtaul_selected_parent() -> anyhow::Result<()> { - let dag = BlockDAG::create_for_testing().unwrap(); + let dag = BlockDAG::create_for_testing(Hash::random()).unwrap(); let reachability_store = dag.storage.reachability_store.clone(); let origin = Hash::random(); @@ -711,7 +710,7 @@ fn test_hint_virtaul_selected_parent() -> anyhow::Result<()> { #[test] fn test_reachability_algorithm() -> anyhow::Result<()> { - let dag = BlockDAG::create_for_testing().unwrap(); + let dag = BlockDAG::create_for_testing(Hash::random()).unwrap(); let reachability_store = dag.storage.reachability_store.clone(); let origin = Hash::random(); @@ -887,9 +886,8 @@ fn add_and_print( #[test] fn test_dag_mergeset() -> anyhow::Result<()> { // initialzie the dag firstly - let mut dag = BlockDAG::create_for_testing().unwrap(); - let genesis = BlockHeader::random(); + let mut dag = BlockDAG::create_for_testing(genesis.id()).unwrap(); dag.init_with_genesis(genesis.clone()).unwrap(); @@ -925,9 +923,8 @@ fn test_dag_mergeset() -> anyhow::Result<()> { #[ignore = "this is the large amount of data testing for performance, dev only"] fn test_big_data_commit() -> anyhow::Result<()> { // initialzie the dag firstly - let mut dag = BlockDAG::create_for_testing().unwrap(); - let genesis = BlockHeader::random(); + let mut dag = BlockDAG::create_for_testing(genesis.id()).unwrap(); dag.init_with_genesis(genesis.clone()).unwrap(); @@ -979,10 +976,9 @@ fn test_prune() -> anyhow::Result<()> { let pruning_depth = 4; let pruning_finality = 3; - let mut dag = BlockDAG::create_for_testing_with_parameters(k).unwrap(); - let genesis = BlockHeader::random(); println!("genesis: {}", genesis.id()); + let mut dag = BlockDAG::create_for_testing_with_parameters(k, genesis.id()).unwrap(); dag.init_with_genesis(genesis.clone()).unwrap(); @@ -1144,9 +1140,8 @@ fn test_verification_blue_block() -> anyhow::Result<()> { // initialzie the dag firstly let k = 5; - let mut dag = BlockDAG::create_for_testing_with_parameters(k).unwrap(); - let genesis = BlockHeader::random(); + let mut dag = BlockDAG::create_for_testing_with_parameters(k, genesis.id()).unwrap(); dag.init_with_genesis(genesis.clone()).unwrap(); @@ -1569,9 +1564,8 @@ fn test_verification_blue_block() -> anyhow::Result<()> { #[test] fn test_check_ancestor_of() -> anyhow::Result<()> { // initialzie the dag firstly - let mut dag = BlockDAG::create_for_testing().unwrap(); - let genesis = BlockHeader::random(); + let mut dag = BlockDAG::create_for_testing(genesis.id()).unwrap(); dag.init_with_genesis(genesis.clone()).unwrap(); @@ -1631,9 +1625,8 @@ fn test_check_ancestor_of() -> anyhow::Result<()> { #[test] fn test_get_blocks_in_batch() -> anyhow::Result<()> { // initialzie the dag firstly - let mut dag = BlockDAG::create_for_testing().unwrap(); - let genesis = BlockHeader::random(); + let mut dag = BlockDAG::create_for_testing(genesis.id()).unwrap(); dag.init_with_genesis(genesis.clone()).unwrap(); diff --git a/genesis/src/lib.rs b/genesis/src/lib.rs index ffe5b7af1b..92f190a6af 100644 --- a/genesis/src/lib.rs +++ b/genesis/src/lib.rs @@ -452,7 +452,8 @@ impl Genesis { let storage = Arc::new(Storage::new(StorageInstance::new_cache_instance())?); let storage2 = Arc::new(Storage2(storage.clone())); let genesis = Genesis::load_or_build(net)?; - let dag = BlockDAG::create_for_testing()?; + let genesis_hash = genesis.block().id(); + let dag = BlockDAG::create_for_testing(genesis_hash)?; let chain_info = genesis.execute_genesis_block(net, storage.clone(), storage2.clone(), dag.clone())?; @@ -466,7 +467,8 @@ impl Genesis { let storage = Arc::new(Storage::new(StorageInstance::new_cache_instance())?); let storage2 = Arc::new(Storage2(storage.clone())); let genesis = Genesis::load_or_build(net)?; - let dag = BlockDAG::create_for_testing_with_parameters(k)?; + let genesis_hash = genesis.block().id(); + let dag = BlockDAG::create_for_testing_with_parameters(k, genesis_hash)?; let chain_info = genesis.execute_genesis_block(net, storage.clone(), storage2.clone(), dag.clone())?; Ok((storage, storage2, chain_info, genesis, dag)) @@ -483,7 +485,8 @@ impl Genesis { )?); let storage2 = Arc::new(Storage2(storage.clone())); let genesis = Genesis::load_or_build(net)?; - let dag = BlockDAG::create_for_testing()?; + let genesis_hash = genesis.block().id(); + let dag = BlockDAG::create_for_testing(genesis_hash)?; let chain_info = genesis.execute_genesis_block(net, storage.clone(), storage2.clone(), dag.clone())?; Ok((storage, storage2, chain_info, genesis, dag)) @@ -558,12 +561,14 @@ mod tests { pub fn do_test_genesis(net: &ChainNetwork, data_dir: &Path) -> Result<()> { let storage1 = Arc::new(Storage::new(StorageInstance::new_cache_instance())?); - let dag1 = BlockDAG::create_for_testing()?; + let genesis = Genesis::load_or_build(net)?; + let genesis_hash = genesis.block().id(); + let dag1 = BlockDAG::create_for_testing(genesis_hash)?; let (chain_info1, genesis1) = Genesis::init_and_check_storage(net, storage1.clone(), dag1, data_dir)?; let storage1_2 = Arc::new(Storage::new(StorageInstance::new_cache_instance())?); let storage2_2 = Arc::new(Storage2(storage1_2.clone())); - let dag2 = BlockDAG::create_for_testing()?; + let dag2 = BlockDAG::create_for_testing(genesis_hash)?; let (chain_info2, genesis2) = Genesis::init_and_check_storage(net, storage1_2.clone(), dag2, data_dir)?; diff --git a/node/src/node.rs b/node/src/node.rs index 1e485d7227..a2c47dd66c 100644 --- a/node/src/node.rs +++ b/node/src/node.rs @@ -16,6 +16,7 @@ use starcoin_block_relayer::BlockRelayer; use starcoin_chain_notify::ChainNotifyHandlerService; use starcoin_chain_service::ChainReaderService; use starcoin_config::NodeConfig; +use starcoin_crypto::HashValue; use starcoin_genesis::{Genesis, GenesisError}; use starcoin_logger::prelude::*; use starcoin_logger::structured_log::init_slog_logger; @@ -299,13 +300,14 @@ impl NodeService { .genesis_config2() .consensus_config .base_max_uncles_per_block; - let dag = starcoin_dag::blockdag::BlockDAG::new( + let genesis_hash = storage.get_genesis()?.unwrap_or(HashValue::zero()); + let mut dag = starcoin_dag::blockdag::BlockDAG::new( starcoin_types::blockhash::KType::try_from(k)?, config.miner.dag_merge_depth(), config.miner.maximum_parents_count(), dag_storage.clone(), + genesis_hash, ); - registry.put_shared(dag.clone()).await?; let (chain_info, genesis) = Genesis::init_and_check_storage( config.net(), @@ -314,6 +316,11 @@ impl NodeService { config.data_dir(), )?; + if genesis_hash == HashValue::zero() { + dag.set_genesis(chain_info.genesis_hash()); + } + registry.put_shared(dag.clone()).await?; + info!( "Start node with chain info: {}, number {}, dragon fork disabled, upgrade_time cost {} secs, ", chain_info, diff --git a/simnet/src/scene/mod.rs b/simnet/src/scene/mod.rs index 969ad8e8e4..4ea8328fd6 100644 --- a/simnet/src/scene/mod.rs +++ b/simnet/src/scene/mod.rs @@ -40,13 +40,13 @@ impl GhostAdpter { .with_difficulty(0.into()) .with_timestamp(time.now_millis()) .build(); + let genesis_id = genesis.id(); let db_tempdir = tempfile::tempdir()?; let config = FlexiDagStorageConfig::new(); let dag_storage = FlexiDagStorage::create_from_path(db_tempdir.path(), config)?; - let mut dag = BlockDAG::new(k, merge_depth, max_parents_count, dag_storage); + let mut dag = BlockDAG::new(k, merge_depth, max_parents_count, dag_storage, genesis_id); dag.init_with_genesis(genesis.clone())?; - let genesis_id = genesis.id(); Ok(Self { genesis, dag, diff --git a/sync/src/tasks/test_tools.rs b/sync/src/tasks/test_tools.rs index 24778c74f5..ce6262724e 100644 --- a/sync/src/tasks/test_tools.rs +++ b/sync/src/tasks/test_tools.rs @@ -55,13 +55,14 @@ impl SyncTestSystem { // Create storage2 using cache instance for simplicity in tests let storage2 = Arc::new(Storage2(storage.clone())); let genesis = Genesis::load_or_build(config.net())?; + let genesis_hash = genesis.block().id(); // init dag let dag_storage = starcoin_dag::consensusdb::prelude::FlexiDagStorage::create_from_path( dag_path.as_path(), FlexiDagStorageConfig::new(), ) .expect("init dag storage fail."); - let dag = starcoin_dag::blockdag::BlockDAG::create_blockdag(dag_storage); // local dag + let dag = starcoin_dag::blockdag::BlockDAG::create_blockdag(dag_storage, genesis_hash); // local dag let chain_info = genesis.execute_genesis_block( config.net(), diff --git a/test-helper/src/chain.rs b/test-helper/src/chain.rs index 42bfd0519a..f126031ec8 100644 --- a/test-helper/src/chain.rs +++ b/test-helper/src/chain.rs @@ -53,9 +53,10 @@ pub fn init_storage_for_test_with_temp_dir( // Load or build genesis let genesis = Genesis::load_or_build(net)?; + let genesis_hash = genesis.block().id(); // Create DAG for testing - let dag = BlockDAG::create_for_testing()?; + let dag = BlockDAG::create_for_testing(genesis_hash)?; // Execute genesis block let chain_info = diff --git a/txpool/src/lib.rs b/txpool/src/lib.rs index 34db1690c4..c039cbc8d8 100644 --- a/txpool/src/lib.rs +++ b/txpool/src/lib.rs @@ -8,23 +8,9 @@ extern crate log; extern crate trace_time; extern crate transaction_pool as tx_pool; -use anyhow::{format_err, Result}; -use network_api::messages::PeerTransactionsMessage; pub use pool::queue::Pool; pub use pool::TxStatus; -use starcoin_config::NodeConfig; -use starcoin_executor::VMMetrics; -use starcoin_service_registry::{ActorService, EventHandler, ServiceContext, ServiceFactory}; -use starcoin_storage::Storage2; -use starcoin_storage::{BlockStore, Storage}; -use starcoin_txpool_api::{PropagateTransactions, TxnStatusFullEvent}; -use starcoin_types::multi_transaction::MultiSignedUserTransaction; -use starcoin_types::{sync_status::SyncStatus, system_events::SyncStatusChangeEvent}; -use starcoin_vm2_state_api::AccountStateReader; -use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::Arc; -use std::time::Duration; -use tx_pool_service_impl::Inner; +pub use tx_pool_actor_service::TxPoolActorService; pub use tx_pool_service_impl::TxPoolService; mod metrics; @@ -33,218 +19,5 @@ mod pool; mod pool_client; #[cfg(test)] mod test; +mod tx_pool_actor_service; mod tx_pool_service_impl; - -//TODO refactor TxPoolService and rename. -#[derive(Clone)] -pub struct TxPoolActorService { - inner: Inner, - new_txs_received: Arc, - sync_status: Option, -} - -impl std::fmt::Debug for TxPoolActorService { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!(f, "pool: {:?}", &self.inner,) - } -} - -const MIN_TXN_TO_PROPAGATE: usize = 256; -const PROPAGATE_FOR_BLOCKS: u64 = 4; - -impl TxPoolActorService { - fn new(inner: Inner) -> Self { - Self { - inner, - sync_status: None, - new_txs_received: Arc::new(AtomicBool::new(false)), - } - } - - pub fn is_synced(&self) -> bool { - match self.sync_status.as_ref() { - Some(sync_status) => sync_status.is_synced(), - None => false, - } - } - - fn transactions_to_propagate(&self) -> Result> { - let statedb = self.inner.get_chain_reader()?; - let reader = AccountStateReader::new(&statedb); - - // TODO: fetch from a gas constants - //TODO optimize broadcast txn by hash, then calculate max length by block gas limit - // currently use a small size for reduce broadcast message size. - // let block_gas_limit = reader.get_epoch()?.block_gas_limit(); - //let min_tx_gas = 200; - // let max_len = std::cmp::max( - // MIN_TXN_TO_PROPAGATE, - // (block_gas_limit / min_tx_gas * PROPAGATE_FOR_BLOCKS) as usize, - // ); - let max_len = 100; - let current_timestamp = reader.get_timestamp()?.seconds(); - Ok(self - .inner - .get_pending(max_len, current_timestamp)? - .into_iter() - .map(|t| t.signed().clone()) - .collect()) - } -} - -impl ServiceFactory for TxPoolActorService { - fn create(ctx: &mut ServiceContext) -> Result { - let storage = ctx.get_shared::>()?; - let storage2 = ctx.get_shared::>()?; - let node_config = ctx.get_shared::>()?; - let vm_metrics = ctx.get_shared_opt::()?; - let txpool_service = ctx.get_shared_or_put(|| { - let startup_info = storage - .get_startup_info()? - .ok_or_else(|| format_err!("StartupInfo should exist when service init."))?; - let best_block = storage - .get_block_by_hash(startup_info.main)? - .ok_or_else(|| { - format_err!( - "best block id {} should exists in storage", - startup_info.main - ) - })?; - let best_block_header = best_block.into_inner().0; - Ok(TxPoolService::new( - node_config, - storage, - storage2, - best_block_header, - vm_metrics, - )) - })?; - Ok(Self::new(txpool_service.get_inner())) - } -} -impl TxPoolActorService { - fn try_propagate_txns(&self, ctx: &mut ServiceContext) { - // only propagate when new txns enter pool. - if self.new_txs_received.load(Ordering::Relaxed) { - match self.transactions_to_propagate() { - Err(e) => { - log::error!("txpool: fail to get txn to propagate, err: {}", &e) - } - Ok(txs) if !txs.is_empty() => { - if self - .new_txs_received - .compare_exchange(true, false, Ordering::Relaxed, Ordering::Relaxed) - .unwrap_or_else(|x| x) - { - let request = PropagateTransactions::new(txs); - ctx.broadcast(request); - } - } - Ok(_) => {} - } - } - } -} -impl ActorService for TxPoolActorService { - fn started(&mut self, ctx: &mut ServiceContext) -> Result<()> { - ctx.subscribe::(); - ctx.add_stream(self.inner.subscribe_txns()); - - // every x seconds, we tick a txn propagation. - let myself = self.clone(); - let interval = self.inner.node_config.tx_pool.tx_propagate_interval(); - ctx.run_interval(Duration::from_secs(interval), move |ctx| { - myself.try_propagate_txns(ctx) - }); - - Ok(()) - } - - fn stopped(&mut self, ctx: &mut ServiceContext) -> Result<()> { - ctx.unsubscribe::(); - Ok(()) - } -} - -impl EventHandler for TxPoolActorService { - fn handle_event(&mut self, msg: SyncStatusChangeEvent, _ctx: &mut ServiceContext) { - self.sync_status = Some(msg.0); - } -} - -/// Listen to txn status, and propagate to remote peers if necessary. -impl EventHandler for TxPoolActorService { - fn handle_event(&mut self, item: TxnStatusFullEvent, _ctx: &mut ServiceContext) { - // do metrics. - if let Some(metrics) = self.inner.metrics.as_ref() { - let status = self.inner.pool_status().status; - let mem_usage = status.mem_usage; - let senders = status.senders; - let txn_count = status.transaction_count; - - metrics - .txpool_status - .with_label_values(&["mem_usage"]) - .set(mem_usage as u64); - metrics - .txpool_status - .with_label_values(&["senders"]) - .set(senders as u64); - metrics - .txpool_status - .with_label_values(&["count"]) - .set(txn_count as u64); - } - let mut has_new_txns = false; - for (_, s) in item.iter() { - if let Some(metrics) = self.inner.metrics.as_ref() { - metrics - .txpool_txn_event_total - .with_label_values(&[format!("{}", s).as_str()]) - .inc(); - } - - if *s == TxStatus::Added { - has_new_txns = true; - } - } - if has_new_txns { - // notify txn-broadcaster. - self.new_txs_received.store(true, Ordering::Relaxed); - } - } -} - -impl EventHandler for TxPoolActorService { - fn handle_event(&mut self, msg: PeerTransactionsMessage, _ctx: &mut ServiceContext) { - if self.is_synced() { - // JUST need to keep at most once delivery. - let bypass_vm1_limit = msg - .message - .txns - .iter() - .all(|txn| matches!(txn, MultiSignedUserTransaction::VM2(_))); - let _ = self.inner.import_txns( - msg.message.txns, - bypass_vm1_limit, - Some(msg.peer_id.to_string()), - ); - } else { - //TODO should keep txn in a buffer, then execute after sync finished. - debug!("[txpool] Ignore PeerTransactions event because the node has not been synchronized yet."); - } - } -} - -#[cfg(test)] -mod test_sync_and_send { - fn assert_send() {} - fn assert_sync() {} - fn assert_static() {} - #[test] - fn test_sync_and_send() { - assert_send::(); - assert_sync::(); - assert_static::(); - } -} diff --git a/txpool/src/tx_pool_actor_service.rs b/txpool/src/tx_pool_actor_service.rs new file mode 100644 index 0000000000..0ff0121201 --- /dev/null +++ b/txpool/src/tx_pool_actor_service.rs @@ -0,0 +1,249 @@ +// Copyright (c) The Starcoin Core Contributors +// SPDX-License-Identifier: Apache-2.0 + +use anyhow::{format_err, Result}; +use network_api::messages::PeerTransactionsMessage; +use starcoin_config::NodeConfig; +use starcoin_executor::VMMetrics; +use starcoin_service_registry::{ActorService, EventHandler, ServiceContext, ServiceFactory}; +use starcoin_storage::Storage2; +use starcoin_storage::{BlockStore, Storage}; +use starcoin_txpool_api::{PropagateTransactions, TxnStatusFullEvent}; +use starcoin_types::multi_transaction::MultiSignedUserTransaction; +use starcoin_types::sync_status::SyncStatus; +use starcoin_types::system_events::SyncStatusChangeEvent; +use starcoin_vm2_state_api::AccountStateReader; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; +use std::time::Duration; + +use crate::pool::TxStatus; +use crate::tx_pool_service_impl::Inner; +use crate::TxPoolService; + +//TODO refactor TxPoolService and rename. +#[derive(Clone)] +pub struct TxPoolActorService { + inner: Inner, + new_txs_received: Arc, + sync_status: Option, +} + +impl std::fmt::Debug for TxPoolActorService { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "pool: {:?}", &self.inner,) + } +} + +impl TxPoolActorService { + pub(crate) fn new(inner: Inner) -> Self { + Self { + inner, + sync_status: None, + new_txs_received: Arc::new(AtomicBool::new(false)), + } + } + + pub fn is_synced(&self) -> bool { + match self.sync_status.as_ref() { + Some(sync_status) => sync_status.is_synced(), + None => false, + } + } + + fn transactions_to_propagate(&self) -> Result> { + let statedb = self.inner.get_chain_reader()?; + let reader = AccountStateReader::new(&statedb); + + // TODO: fetch from a gas constants + //TODO optimize broadcast txn by hash, then calculate max length by block gas limit + // currently use a small size for reduce broadcast message size. + // let block_gas_limit = reader.get_epoch()?.block_gas_limit(); + //let min_tx_gas = 200; + // let max_len = std::cmp::max( + // MIN_TXN_TO_PROPAGATE, + // (block_gas_limit / min_tx_gas * PROPAGATE_FOR_BLOCKS) as usize, + // ); + let max_len = 100; + let current_timestamp = reader.get_timestamp()?.seconds(); + Ok(self + .inner + .get_pending(max_len, current_timestamp)? + .into_iter() + .map(|t| t.signed().clone()) + .collect()) + } + + fn try_propagate_txns(&self, ctx: &mut ServiceContext) { + // only propagate when new txns enter pool. + if self.new_txs_received.load(Ordering::Relaxed) { + match self.transactions_to_propagate() { + Err(e) => { + log::error!("txpool: fail to get txn to propagate, err: {}", &e) + } + Ok(txs) if !txs.is_empty() => { + if self + .new_txs_received + .compare_exchange(true, false, Ordering::Relaxed, Ordering::Relaxed) + .unwrap_or_else(|x| x) + { + let request = PropagateTransactions::new(txs); + ctx.broadcast(request); + } + } + Ok(_) => {} + } + } + } + + fn try_cull(&self) { + if let Err(e) = self.inner.cull() { + log::error!("txpool: fail to cull expired transactions, err: {}", e); + } + } +} + +impl ServiceFactory for TxPoolActorService { + fn create(ctx: &mut ServiceContext) -> Result { + let storage = ctx.get_shared::>()?; + let storage2 = ctx.get_shared::>()?; + let node_config = ctx.get_shared::>()?; + let vm_metrics = ctx.get_shared_opt::()?; + let txpool_service = ctx.get_shared_or_put(|| { + let startup_info = storage + .get_startup_info()? + .ok_or_else(|| format_err!("StartupInfo should exist when service init."))?; + let best_block = storage + .get_block_by_hash(startup_info.main)? + .ok_or_else(|| { + format_err!( + "best block id {} should exists in storage", + startup_info.main + ) + })?; + let best_block_header = best_block.into_inner().0; + Ok(TxPoolService::new( + node_config, + storage, + storage2, + best_block_header, + vm_metrics, + )) + })?; + Ok(Self::new(txpool_service.get_inner())) + } +} + +impl ActorService for TxPoolActorService { + fn started(&mut self, ctx: &mut ServiceContext) -> Result<()> { + ctx.subscribe::(); + ctx.add_stream(self.inner.subscribe_txns()); + + // every x seconds, we tick a txn propagation. + let myself = self.clone(); + let interval = self.inner.node_config.tx_pool.tx_propagate_interval(); + ctx.run_interval(Duration::from_secs(interval), move |ctx| { + myself.try_propagate_txns(ctx) + }); + + // every x seconds, we cull expired transactions. + let myself_for_cull = self.clone(); + let cull_interval = self.inner.node_config.tx_pool.cull_interval(); + ctx.run_interval(Duration::from_secs(cull_interval), move |_ctx| { + myself_for_cull.try_cull() + }); + + Ok(()) + } + + fn stopped(&mut self, ctx: &mut ServiceContext) -> Result<()> { + ctx.unsubscribe::(); + Ok(()) + } +} + +impl EventHandler for TxPoolActorService { + fn handle_event(&mut self, msg: SyncStatusChangeEvent, _ctx: &mut ServiceContext) { + self.sync_status = Some(msg.0); + } +} + +/// Listen to txn status, and propagate to remote peers if necessary. +impl EventHandler for TxPoolActorService { + fn handle_event(&mut self, item: TxnStatusFullEvent, _ctx: &mut ServiceContext) { + // do metrics. + if let Some(metrics) = self.inner.metrics.as_ref() { + let status = self.inner.pool_status().status; + let mem_usage = status.mem_usage; + let senders = status.senders; + let txn_count = status.transaction_count; + + metrics + .txpool_status + .with_label_values(&["mem_usage"]) + .set(mem_usage as u64); + metrics + .txpool_status + .with_label_values(&["senders"]) + .set(senders as u64); + metrics + .txpool_status + .with_label_values(&["count"]) + .set(txn_count as u64); + } + let mut has_new_txns = false; + for (_, s) in item.iter() { + if let Some(metrics) = self.inner.metrics.as_ref() { + metrics + .txpool_txn_event_total + .with_label_values(&[format!("{}", s).as_str()]) + .inc(); + } + + if *s == TxStatus::Added { + has_new_txns = true; + } + } + if has_new_txns { + // notify txn-broadcaster. + self.new_txs_received.store(true, Ordering::Relaxed); + } + } +} + +impl EventHandler for TxPoolActorService { + fn handle_event(&mut self, msg: PeerTransactionsMessage, _ctx: &mut ServiceContext) { + if self.is_synced() { + // JUST need to keep at most once delivery. + let bypass_vm1_limit = msg + .message + .txns + .iter() + .all(|txn| matches!(txn, MultiSignedUserTransaction::VM2(_))); + let _ = self.inner.import_txns( + msg.message.txns, + bypass_vm1_limit, + Some(msg.peer_id.to_string()), + ); + } else { + //TODO should keep txn in a buffer, then execute after sync finished. + debug!("[txpool] Ignore PeerTransactions event because the node has not been synchronized yet."); + } + } +} + +#[cfg(test)] +mod test_sync_and_send { + use super::TxPoolActorService; + + fn assert_send() {} + fn assert_sync() {} + fn assert_static() {} + + #[test] + fn test_sync_and_send() { + assert_send::(); + assert_sync::(); + assert_static::(); + } +} diff --git a/txpool/src/tx_pool_service_impl.rs b/txpool/src/tx_pool_service_impl.rs index eff3f2381c..bd207f7fb7 100644 --- a/txpool/src/tx_pool_service_impl.rs +++ b/txpool/src/tx_pool_service_impl.rs @@ -400,12 +400,16 @@ impl Inner { bypass_vm1_limit: bool, peer_id: Option, ) -> Result>> { + let now_seconds = self.chain_header.read().timestamp() / 1000; + let pool_client = self.get_pool_client()?; + self.queue.cull(pool_client.clone(), now_seconds); + let txns = txns .into_iter() .map(|t| PoolTransaction::Unverified(UnverifiedUserTransaction::from(t))); Ok(self .queue - .import(self.get_pool_client()?, txns, bypass_vm1_limit, peer_id)) + .import(pool_client, txns, bypass_vm1_limit, peer_id)) } pub(crate) fn remove_txn( &self,