diff --git a/Cargo.lock b/Cargo.lock index 17a5c04b28..c87e36e213 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -12533,6 +12533,7 @@ dependencies = [ "starcoin-vm1-vm-types", "starcoin-vm2-account-api", "starcoin-vm2-account-service", + "starcoin-vm2-state-api", "starcoin-vm2-statedb", "stest", "test-helper", diff --git a/miner/Cargo.toml b/miner/Cargo.toml index 2cc8395c6f..2de31922cf 100644 --- a/miner/Cargo.toml +++ b/miner/Cargo.toml @@ -42,6 +42,8 @@ starcoin-vm2-account-service = { workspace = true } starcoin-vm2-account-api = { workspace = true } starcoin-vm2-types = { workspace = true } starcoin-vm2-vm-types = { workspace = true } +starcoin-vm2-state-api = { workspace = true } +starcoin-vm2-statedb = { workspace = true } [dev-dependencies] starcoin-network-rpc = { workspace = true } diff --git a/miner/src/create_block_template/block_builder_service.rs b/miner/src/create_block_template/block_builder_service.rs index c0b86b9718..c2b55f7e59 100644 --- a/miner/src/create_block_template/block_builder_service.rs +++ b/miner/src/create_block_template/block_builder_service.rs @@ -42,6 +42,7 @@ use std::sync::RwLock; use crate::{MinerService, NewHeaderChannel}; use super::metrics::BlockBuilderMetrics; +use super::state_check::StateCheck; use sp_utils::thread_pool::RAYON_EXEC_POOL; use starcoin_dag::types::ghostdata::GhostdagData; use starcoin_types::U256; @@ -522,7 +523,7 @@ where now_millis, ); - let (txns, txns2) = self.fetch_transactions(&previous_header, &blue_blocks, max_txns)?; + let (txns, txns2) = self.fetch_transactions(&main, &blue_blocks, max_txns)?; info!( "[BlockProcess] VM1 txns len: {}, VM2 txns len: {}", txns.len(), @@ -668,46 +669,54 @@ where fn fetch_transactions( &self, - header: &BlockHeader, + chain: &BlockChain, blue_blocks: &[Block], max_txns: u64, ) -> Result<(Vec, Vec)> { - let pending_multi_transactions = self.tx_provider.get_txns_with_header(max_txns, header); + let header = chain.current_header(); - // Separate VM1 and VM2 transactions - let mut pending_transactions = vec![]; - let mut pending_transactions2 = vec![]; - pending_multi_transactions - .into_iter() - .for_each(|txn| match txn { - MultiSignedUserTransaction::VM1(txn) => pending_transactions.push(txn), - MultiSignedUserTransaction::VM2(txn) => pending_transactions2.push(txn), - }); + let state_db1 = chain.chain_state_reader(); + let state_db2 = chain.chain_state_reader2(); + let mut state_check = StateCheck::new(state_db1, state_db2); - if pending_transactions.len() + pending_transactions2.len() >= max_txns as usize { - return Ok((pending_transactions, pending_transactions2)); - } + let pending_multi_transactions = self.tx_provider.get_txns_with_header(max_txns, &header); + let mut blue_transactions: Vec = vec![]; blue_blocks.iter().for_each(|block| { block.transactions().iter().for_each(|transaction| { - pending_transactions.push(transaction.clone()); + blue_transactions.push(MultiSignedUserTransaction::VM1(transaction.clone())); }); block.transactions2().iter().for_each(|transaction| { - pending_transactions2.push(transaction.clone()); + blue_transactions.push(MultiSignedUserTransaction::VM2(transaction.clone())); }) }); - pending_transactions.sort_by(|a, b| match a.sender().cmp(&b.sender()) { + blue_transactions.sort_by(|a, b| match a.sender().to_hex().cmp(&b.sender().to_hex()) { std::cmp::Ordering::Equal => a.sequence_number().cmp(&b.sequence_number()), other => other, }); + let _ = state_check.filter_continuous_transactions(blue_transactions)?; - pending_transactions2.sort_by(|a, b| match a.sender().cmp(&b.sender()) { - std::cmp::Ordering::Equal => a.sequence_number().cmp(&b.sequence_number()), - other => other, + let mut pending_multi_transactions = pending_multi_transactions; + pending_multi_transactions.sort_by(|a, b| { + match a.sender().to_hex().cmp(&b.sender().to_hex()) { + std::cmp::Ordering::Equal => a.sequence_number().cmp(&b.sequence_number()), + other => other, + } }); + let filtered_transactions = + state_check.filter_continuous_transactions(pending_multi_transactions)?; + + let mut pending_transactions = vec![]; + let mut pending_transactions2 = vec![]; + for txn in filtered_transactions { + match txn { + MultiSignedUserTransaction::VM1(txn) => pending_transactions.push(txn), + MultiSignedUserTransaction::VM2(txn) => pending_transactions2.push(txn), + } + } Ok((pending_transactions, pending_transactions2)) } diff --git a/miner/src/create_block_template/mod.rs b/miner/src/create_block_template/mod.rs index febea0537d..1c196d5f36 100644 --- a/miner/src/create_block_template/mod.rs +++ b/miner/src/create_block_template/mod.rs @@ -6,3 +6,7 @@ pub mod new_header_service; //#[cfg(test)] //mod test_create_block_template; pub mod block_builder_service; +pub mod state_check; + +#[cfg(test)] +mod test_state_check; diff --git a/miner/src/create_block_template/state_check.rs b/miner/src/create_block_template/state_check.rs new file mode 100644 index 0000000000..1a879fe018 --- /dev/null +++ b/miner/src/create_block_template/state_check.rs @@ -0,0 +1,97 @@ +use anyhow::Result; +use starcoin_state_api::ChainStateReader as ChainStateReader1; +use starcoin_state_api::StateReaderExt as StateReaderExt1; +use starcoin_types::account_address::AccountAddress; +use starcoin_types::multi_transaction::{MultiAccountAddress, MultiSignedUserTransaction}; +use starcoin_vm2_state_api::ChainStateReader as ChainStateReader2; +use starcoin_vm2_state_api::StateReaderExt as StateReaderExt2; +use starcoin_vm2_types::account_address::AccountAddress as AccountAddress2; +use std::collections::HashMap; + +pub struct StateCheck<'a, R1: ?Sized + ChainStateReader1, R2: ?Sized + ChainStateReader2> { + state_reader1: &'a R1, + state_reader2: &'a R2, + sequence_cache1: HashMap, + sequence_cache2: HashMap, +} + +impl<'a, R1: ?Sized + ChainStateReader1, R2: ?Sized + ChainStateReader2> StateCheck<'a, R1, R2> { + pub fn new(state_reader1: &'a R1, state_reader2: &'a R2) -> Self { + Self { + state_reader1, + state_reader2, + sequence_cache1: HashMap::new(), + sequence_cache2: HashMap::new(), + } + } + + pub fn get_next_expected_sequence(&mut self, sender: MultiAccountAddress) -> Result { + match sender { + MultiAccountAddress::VM1(addr) => { + if let Some(&seq) = self.sequence_cache1.get(&addr) { + return Ok(seq); + } + let seq = self.state_reader1.get_sequence_number(addr)?; + Ok(seq) + } + MultiAccountAddress::VM2(addr) => { + if let Some(&seq) = self.sequence_cache2.get(&addr) { + return Ok(seq); + } + let seq = self.state_reader2.get_sequence_number(addr)?; + Ok(seq) + } + } + } + + pub fn check_and_accept(&mut self, txn: &MultiSignedUserTransaction) -> Result { + let sender = txn.sender(); + let txn_seq = txn.sequence_number(); + let expected_seq = self.get_next_expected_sequence(sender)?; + + if txn_seq == expected_seq { + match sender { + MultiAccountAddress::VM1(addr) => { + self.sequence_cache1.insert(addr, txn_seq + 1); + } + MultiAccountAddress::VM2(addr) => { + self.sequence_cache2.insert(addr, txn_seq + 1); + } + } + Ok(true) + } else { + Ok(false) + } + } + + pub fn filter_continuous_transactions( + &mut self, + transactions: Vec, + ) -> Result> { + let mut result = Vec::with_capacity(transactions.len()); + + for txn in transactions { + match self.check_and_accept(&txn) { + Ok(true) => result.push(txn), + Ok(false) => { + continue; + } + Err(e) => { + starcoin_logger::prelude::warn!( + "Failed to check sequence for sender {}: {:?}", + txn.sender().to_hex(), + e + ); + continue; + } + } + } + + Ok(result) + } + + pub fn clear_cache(&mut self) { + self.sequence_cache1.clear(); + self.sequence_cache2.clear(); + } +} diff --git a/miner/src/create_block_template/test_state_check.rs b/miner/src/create_block_template/test_state_check.rs new file mode 100644 index 0000000000..2f0e136828 --- /dev/null +++ b/miner/src/create_block_template/test_state_check.rs @@ -0,0 +1,302 @@ +use super::state_check::StateCheck; +use anyhow::{format_err, Result}; +use starcoin_state_api::{ + AccountStateSetIterator, ChainStateReader as ChainStateReader1, StateView as StateView1, + StateWithProof, StateWithTableItemProof, +}; +use starcoin_types::{ + access_path::AccessPath as AccessPath1, + account::peer_to_peer_txn as peer_to_peer_txn1, + account::Account as Account1, + account::DEFAULT_EXPIRATION_TIME as DEFAULT_EXPIRATION_TIME1, + account_address::AccountAddress, + account_state::AccountState, + event::EventHandle as EventHandle1, + multi_transaction::MultiSignedUserTransaction, + state_set::{AccountStateSet, ChainStateSet}, + transaction::SignedUserTransaction as SignedUserTransaction1, +}; +use starcoin_vm2_state_api::ChainStateReader as ChainStateReader2; +use starcoin_vm2_state_api::{ + StateWithProof as StateWithProof2, StateWithTableItemProof as StateWithTableItemProof2, +}; +use starcoin_vm2_types::{ + account::peer_to_peer_txn as peer_to_peer_txn2, account::Account as Account2, + account::DEFAULT_EXPIRATION_TIME as DEFAULT_EXPIRATION_TIME2, + account_address::AccountAddress as AccountAddress2, + account_state::AccountState as AccountState2, +}; +use starcoin_vm2_vm_types::{ + access_path::DataPath as DataPath2, account_config::AccountResource as AccountResource2, + event::EventKey as EventKey2, move_resource::MoveStructType, + on_chain_resource::ChainId as ChainId2, state_store::state_key::inner::StateKeyInner, + state_store::state_key::StateKey as StateKey2, + state_store::state_storage_usage::StateStorageUsage, state_store::TStateView, +}; +use starcoin_vm_types::{ + account_config::AccountResource as AccountResource1, genesis_config::ChainId as ChainId1, + move_resource::MoveResource, state_store::state_key::StateKey as StateKey1, + state_store::table::TableHandle as TableHandle1, +}; +use std::collections::HashMap; + +struct TestStateReader1 { + seqs: HashMap, +} + +impl TestStateReader1 { + fn new(seqs: HashMap) -> Self { + Self { seqs } + } +} + +impl StateView1 for TestStateReader1 { + fn get_state_value(&self, state_key: &StateKey1) -> Result>> { + let StateKey1::AccessPath(access_path) = state_key else { + return Ok(None); + }; + if access_path.path != AccountResource1::resource_path() { + return Ok(None); + } + let Some(seq) = self.seqs.get(&access_path.address) else { + return Ok(None); + }; + let handle = EventHandle1::random_handle(0); + let resource = AccountResource1::new( + *seq, + AccountResource1::DUMMY_AUTH_KEY.to_vec(), + None, + None, + handle.clone(), + handle.clone(), + handle, + ); + Ok(Some(bcs_ext::to_bytes(&resource)?)) + } + + fn is_genesis(&self) -> bool { + false + } +} + +impl ChainStateReader1 for TestStateReader1 { + fn get_with_proof(&self, _access_path: &AccessPath1) -> Result { + Err(format_err!("not used in test")) + } + + fn get_account_state(&self, _address: &AccountAddress) -> Result> { + Ok(None) + } + + fn get_account_state_set(&self, _address: &AccountAddress) -> Result> { + Ok(None) + } + + fn state_root(&self) -> starcoin_crypto::HashValue { + starcoin_crypto::HashValue::zero() + } + + fn dump(&self) -> Result { + Err(format_err!("not used in test")) + } + + fn dump_iter(&self) -> Result { + Err(format_err!("not used in test")) + } + + fn get_with_table_item_proof( + &self, + _handle: &TableHandle1, + _key: &[u8], + ) -> Result { + Err(format_err!("not used in test")) + } +} + +struct TestStateReader2 { + seqs: HashMap, +} + +impl TestStateReader2 { + fn new(seqs: HashMap) -> Self { + Self { seqs } + } +} + +impl TStateView for TestStateReader2 { + type Key = StateKey2; + + fn get_state_value( + &self, + state_key: &Self::Key, + ) -> starcoin_vm2_vm_types::state_store::Result< + Option, + > { + let StateKeyInner::AccessPath(access_path) = state_key.inner() else { + return Ok(None); + }; + let DataPath2::Resource(struct_tag) = &access_path.path else { + return Ok(None); + }; + if struct_tag != &AccountResource2::struct_tag() { + return Ok(None); + } + let Some(seq) = self.seqs.get(&access_path.address) else { + return Ok(None); + }; + let handle = starcoin_vm2_vm_types::event::EventHandle::new( + EventKey2::new(0, AccountAddress2::ZERO), + 0, + ); + let resource = AccountResource2::new( + *seq, + AccountResource2::DUMMY_AUTH_KEY.to_vec(), + handle.clone(), + handle, + ); + Ok(Some(bcs_ext::to_bytes(&resource)?.into())) + } + + fn get_usage(&self) -> starcoin_vm2_vm_types::state_store::Result { + Ok(StateStorageUsage::zero()) + } + + fn is_genesis(&self) -> bool { + false + } +} + +impl ChainStateReader2 for TestStateReader2 { + fn get_with_proof(&self, _state_key: &StateKey2) -> Result { + Err(format_err!("not used in test")) + } + + fn get_account_state(&self, _address: &AccountAddress2) -> Result { + Err(format_err!("not used in test")) + } + + fn get_account_state_set( + &self, + _address: &AccountAddress2, + ) -> Result> { + Ok(None) + } + + fn state_root(&self) -> starcoin_crypto::HashValue { + starcoin_crypto::HashValue::zero() + } + + fn dump(&self) -> Result { + Err(format_err!("not used in test")) + } + + fn dump_iter(&self) -> Result { + Err(format_err!("not used in test")) + } + + fn get_with_table_item_proof( + &self, + _handle: &starcoin_vm2_vm_types::state_store::table::TableHandle, + _key: &[u8], + ) -> Result { + Err(format_err!("not used in test")) + } +} + +fn make_vm1_txn(sender: &Account1, seq: u64) -> SignedUserTransaction1 { + let receiver = Account1::new(); + peer_to_peer_txn1( + sender, + &receiver, + seq, + 1, + DEFAULT_EXPIRATION_TIME1, + ChainId1::test(), + ) +} + +fn make_vm2_txn( + sender: &Account2, + seq: u64, +) -> starcoin_vm2_vm_types::transaction::SignedUserTransaction { + let receiver = Account2::new(); + peer_to_peer_txn2( + sender, + &receiver, + seq, + 1, + DEFAULT_EXPIRATION_TIME2, + ChainId2::test(), + ) +} + +fn sort_transactions(transactions: &mut [MultiSignedUserTransaction]) { + transactions.sort_by(|a, b| match a.sender().to_hex().cmp(&b.sender().to_hex()) { + std::cmp::Ordering::Equal => a.sequence_number().cmp(&b.sequence_number()), + other => other, + }); +} + +#[test] +fn filter_continuous_transactions_rejects_gaps() -> Result<()> { + let sender1 = Account1::new(); + let sender2 = Account2::new(); + + let mut seqs1 = HashMap::new(); + seqs1.insert(*sender1.address(), 1); + let mut seqs2 = HashMap::new(); + seqs2.insert(*sender2.address(), 10); + + let state_reader1 = TestStateReader1::new(seqs1); + let state_reader2 = TestStateReader2::new(seqs2); + let mut state_check = StateCheck::new(&state_reader1, &state_reader2); + + let mut transactions = vec![ + MultiSignedUserTransaction::VM1(make_vm1_txn(&sender1, 1)), + MultiSignedUserTransaction::VM1(make_vm1_txn(&sender1, 3)), + MultiSignedUserTransaction::VM2(make_vm2_txn(&sender2, 11)), + MultiSignedUserTransaction::VM2(make_vm2_txn(&sender2, 10)), + ]; + sort_transactions(&mut transactions); + + let filtered = state_check.filter_continuous_transactions(transactions)?; + let vm1_seqs: Vec = filtered + .iter() + .filter_map(|txn| match txn { + MultiSignedUserTransaction::VM1(txn) => Some(txn.sequence_number()), + _ => None, + }) + .collect(); + let vm2_seqs: Vec = filtered + .iter() + .filter_map(|txn| match txn { + MultiSignedUserTransaction::VM2(txn) => Some(txn.sequence_number()), + _ => None, + }) + .collect(); + + assert_eq!(vm1_seqs, vec![1]); + assert_eq!(vm2_seqs, vec![10, 11]); + Ok(()) +} + +#[test] +fn filter_continuous_transactions_uses_cache() -> Result<()> { + let sender1 = Account1::new(); + + let mut seqs1 = HashMap::new(); + seqs1.insert(*sender1.address(), 5); + let state_reader1 = TestStateReader1::new(seqs1); + let state_reader2 = TestStateReader2::new(HashMap::new()); + let mut state_check = StateCheck::new(&state_reader1, &state_reader2); + + let first = vec![MultiSignedUserTransaction::VM1(make_vm1_txn(&sender1, 5))]; + let second = vec![MultiSignedUserTransaction::VM1(make_vm1_txn(&sender1, 6))]; + + let filtered_first = state_check.filter_continuous_transactions(first)?; + let filtered_second = state_check.filter_continuous_transactions(second)?; + + assert_eq!(filtered_first.len(), 1); + assert_eq!(filtered_second.len(), 1); + Ok(()) +}