Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions miner/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand Down
51 changes: 30 additions & 21 deletions miner/src/create_block_template/block_builder_service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -668,46 +669,54 @@ where

fn fetch_transactions(
&self,
header: &BlockHeader,
chain: &BlockChain,
blue_blocks: &[Block],
max_txns: u64,
) -> Result<(Vec<SignedUserTransaction>, Vec<SignedUserTransaction2>)> {
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<MultiSignedUserTransaction> = 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))
}

Expand Down
4 changes: 4 additions & 0 deletions miner/src/create_block_template/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
97 changes: 97 additions & 0 deletions miner/src/create_block_template/state_check.rs
Original file line number Diff line number Diff line change
@@ -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<AccountAddress, u64>,
sequence_cache2: HashMap<AccountAddress2, u64>,
}

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<u64> {
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<bool> {
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<MultiSignedUserTransaction>,
) -> Result<Vec<MultiSignedUserTransaction>> {
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();
}
}
Loading
Loading