use crate::chain_sync::BadBlockCache;
use crate::db::DbImpl;
use crate::networks::Height;
use crate::prelude::*;
use crate::shim::clock::ALLOWABLE_CLOCK_DRIFT;
use crate::shim::crypto::SignatureType;
use crate::shim::version::NetworkVersion;
use crate::shim::{
address::Address,
crypto::verify_bls_aggregate,
econ::BLOCK_GAS_LIMIT,
gas::{PriceList, price_list_by_network_version},
state_tree::StateTree,
};
use crate::state_manager::ExecutedTipset;
use crate::state_manager::{Error as StateManagerError, StateManager, utils::is_valid_for_sending};
use crate::{
blocks::{Block, CachingBlockHeader, Error as ForestBlockError, FullTipset, Tipset},
fil_cns::{self, FilecoinConsensus, FilecoinConsensusError},
};
use crate::{
chain::{ChainStore, Error as ChainStoreError},
metrics::HistogramTimerExt,
};
use crate::{
eth::is_valid_eth_tx_for_sending,
message::{MessageRead, valid_for_block_inclusion},
};
use ahash::HashMap;
use futures::TryFutureExt;
use nunny::Vec as NonEmpty;
use thiserror::Error;
use tokio::task::JoinSet;
use tracing::{trace, warn};
use crate::chain_sync::{consensus::collect_errs, metrics, validation::TipsetValidator};
#[derive(Debug, Error)]
pub enum TipsetSyncerError {
#[error("Block must have a signature")]
BlockWithoutSignature,
#[error("Block without BLS aggregate signature")]
BlockWithoutBlsAggregate,
#[error("Block received from the future: now = {0}, block = {1}")]
TimeTravellingBlock(u64, u64),
#[error("Validation error: {0}")]
Validation(String),
#[error("Parent chain state mismatch: {0}")]
ParentChainStateMismatch(String),
#[error("Processing error: {0}")]
Calculation(String),
#[error("Chain store error: {0}")]
ChainStore(#[from] ChainStoreError),
#[error("StateManager error: {0}")]
StateManager(#[from] StateManagerError),
#[error("Block error: {0}")]
BlockError(#[from] ForestBlockError),
#[error("Querying tipsets from the network failed: {0}")]
NetworkTipsetQueryFailed(String),
#[error("BLS aggregate signature {0} was invalid for msgs {1}")]
BlsAggregateSignatureInvalid(String, String),
#[error("Message signature invalid: {0}")]
MessageSignatureInvalid(String),
#[error("Block message root does not match: expected {0}, computed {1}")]
BlockMessageRootInvalid(String, String),
#[error("Computing message root failed: {0}")]
ComputingMessageRoot(String),
#[error("Resolving address from message failed: {0}")]
ResolvingAddressFromMessage(String),
#[error("Loading tipset parent from the store failed: {0}")]
TipsetParentNotFound(ChainStoreError),
#[error("Consensus error: {0}")]
ConsensusError(FilecoinConsensusError),
#[error(
"Block had a signed message at index {0} whose signature type {1} is not allowed in the SECP message list"
)]
SecpSignatureTypeInvalid(usize, SignatureType),
}
impl From<tokio::task::JoinError> for TipsetSyncerError {
fn from(err: tokio::task::JoinError) -> Self {
TipsetSyncerError::NetworkTipsetQueryFailed(format!("{err}"))
}
}
impl TipsetSyncerError {
fn concat(errs: NonEmpty<TipsetSyncerError>) -> Self {
let msg = errs.iter().map(|e| e.to_string()).collect_vec().join(", ");
if errs
.iter()
.any(|e| matches!(e, TipsetSyncerError::ParentChainStateMismatch(_)))
{
TipsetSyncerError::ParentChainStateMismatch(msg)
} else {
TipsetSyncerError::Validation(msg)
}
}
}
pub async fn validate_tipset(
state_manager: &StateManager,
full_tipset: FullTipset,
bad_block_cache: Option<BadBlockCache>,
) -> Result<(), TipsetSyncerError> {
if full_tipset
.key()
.eq(state_manager.chain_store().genesis_tipset().key())
{
trace!("Skipping genesis tipset validation");
return Ok(());
}
let timer = metrics::TIPSET_PROCESSING_TIME.start_timer();
let epoch = full_tipset.epoch();
let parent_state = *full_tipset.parent_state();
let tipset_key = full_tipset.key();
trace!("Tipset keys: {tipset_key}");
let blocks = full_tipset.into_blocks();
let mut validations = JoinSet::new();
for b in blocks {
validations.spawn(validate_block(state_manager.shallow_clone(), Arc::new(b)));
}
while let Some(result) = validations.join_next().await {
match result? {
Ok(block) => {
state_manager
.chain_store()
.add_to_tipset_tracker(block.header());
}
Err(boxed) => {
let (cid, why) = *boxed;
warn!(
"Validating block [CID = {cid}, PARENT_STATE = {parent_state}] in EPOCH = {epoch} failed: {why}",
);
match &why {
TipsetSyncerError::TimeTravellingBlock(_, _) => {
}
_ => {
if StateTree::new_from_root(state_manager.db(), &parent_state).is_ok()
&& let Some(bad_block_cache) = bad_block_cache
{
bad_block_cache.push(cid);
}
}
};
return Err(why);
}
}
}
drop(timer);
Ok(())
}
async fn validate_block(
state_manager: StateManager,
block: Arc<Block>,
) -> Result<Arc<Block>, Box<(Cid, TipsetSyncerError)>> {
let consensus = FilecoinConsensus::new(state_manager.beacon_schedule().clone());
trace!(
"Validating block: epoch = {}, weight = {}, key = {}",
block.header().epoch,
block.header().weight,
block.header().cid(),
);
let chain_store = state_manager.chain_store().shallow_clone();
let block_cid = block.cid();
let is_validated = chain_store.is_block_validated(block_cid);
if is_validated {
return Ok(block);
}
let _timer = metrics::BLOCK_VALIDATION_TIME.start_timer();
let header = block.header();
block_sanity_checks(header).map_err(|e| Box::new((*block_cid, e)))?;
block_timestamp_checks(header).map_err(|e| Box::new((*block_cid, e)))?;
let base_tipset = chain_store
.chain_index()
.load_required_tipset(&header.parents)
.map_err(|why| Box::new((*block_cid, TipsetSyncerError::TipsetParentNotFound(why))))?;
let lookback_state = ChainStore::get_lookback_tipset_for_round(
state_manager.chain_store().chain_index().shallow_clone(),
state_manager.chain_config().shallow_clone(),
base_tipset.shallow_clone(),
block.header().epoch,
)
.await
.map_err(|e| Box::new((*block_cid, e.into())))
.map(|(_, s)| Arc::new(s))?;
let work_addr = state_manager
.get_miner_work_addr(*lookback_state, &header.miner_address)
.map_err(|e| Box::new((*block_cid, e.into())))?;
let mut validations = JoinSet::new();
validations.spawn(check_block_messages(
state_manager.shallow_clone(),
block.shallow_clone(),
base_tipset.shallow_clone(),
));
validations.spawn_blocking({
let smoke_height = state_manager.chain_config().epoch(Height::Smoke);
let firehorse_height = state_manager.chain_config().epoch(Height::FireHorse);
let base_tipset = base_tipset.shallow_clone();
let block_store = state_manager.db_owned();
let block = block.shallow_clone();
move || {
let base_fee = crate::chain::compute_base_fee(
&block_store,
&base_tipset,
smoke_height,
firehorse_height,
)
.map_err(|e| {
TipsetSyncerError::Validation(format!("Could not compute base fee: {e}"))
})?;
let parent_base_fee = &block.header.parent_base_fee;
if &base_fee != parent_base_fee {
return Err(TipsetSyncerError::Validation(format!(
"base fee doesn't match: {parent_base_fee} (header), {base_fee} (computed)"
)));
}
Ok(())
}
});
validations.spawn_blocking({
let block_store = state_manager.db_owned();
let base_tipset = base_tipset.shallow_clone();
let weight = header.weight.clone();
move || {
let calc_weight = fil_cns::weight(&block_store, &base_tipset).map_err(|e| {
TipsetSyncerError::Calculation(format!("Error calculating weight: {e:#}"))
})?;
if weight != calc_weight {
return Err(TipsetSyncerError::Validation(format!(
"Parent weight doesn't match: {weight} (header), {calc_weight} (computed)"
)));
}
Ok(())
}
});
validations.spawn({
let state_manager = state_manager.shallow_clone();
let block = block.shallow_clone();
async move {
let header = block.header();
let ExecutedTipset {
state_root,
receipt_root,
..
} = state_manager
.load_executed_tipset(&base_tipset)
.await
.map_err(|e| {
TipsetSyncerError::Calculation(format!("Failed to calculate state: {e:#}"))
})?;
if state_root != header.state_root {
return Err(TipsetSyncerError::ParentChainStateMismatch(format!(
"Parent state root did not match computed state: {} (header), {} (computed)",
header.state_root, state_root,
)));
}
if receipt_root != header.message_receipts {
return Err(TipsetSyncerError::ParentChainStateMismatch(format!(
"Parent receipt root did not match computed root: {} (header), {} (computed)",
header.message_receipts, receipt_root
)));
}
Ok(())
}
});
validations.spawn_blocking({
let block = block.shallow_clone();
move || {
block.header().verify_signature_against(&work_addr)?;
Ok(())
}
});
validations.spawn({
let block = block.shallow_clone();
async move {
consensus
.validate_block(state_manager, block)
.map_err(|errs| {
TipsetSyncerError::concat(
errs.into_iter_ne()
.map(TipsetSyncerError::ConsensusError)
.collect_vec(),
)
})
.await
}
});
if let Err(errs) = collect_errs(validations).await {
return Err(Box::new((*block_cid, TipsetSyncerError::concat(errs))));
}
chain_store.mark_block_as_validated(block_cid);
Ok(block)
}
struct MessageChecker {
price_list: PriceList,
network_version: NetworkVersion,
tree: StateTree<DbImpl>,
sum_gas_limit: u64,
account_sequences: HashMap<Address, u64>,
}
impl MessageChecker {
fn check(&mut self, msg: &impl MessageRead) -> anyhow::Result<()> {
let min_gas = self.price_list.on_chain_message(msg.chain_length()?);
let msg = msg.vm_message();
valid_for_block_inclusion(msg, min_gas.total(), self.network_version)
.map_err(|e| anyhow::anyhow!("{e}"))?;
self.sum_gas_limit += msg.gas_limit();
anyhow::ensure!(
self.sum_gas_limit <= BLOCK_GAS_LIMIT,
"block gas limit exceeded"
);
let sequence: u64 = match self.account_sequences.get(&msg.from()) {
Some(sequence) => *sequence,
None => {
let actor = self.tree.get_actor(&msg.from())?.ok_or_else(|| {
anyhow::anyhow!(
"Failed to retrieve nonce for addr: Actor does not exist in state"
)
})?;
anyhow::ensure!(
is_valid_for_sending(self.network_version, &actor),
"not valid for sending!"
);
actor.sequence
}
};
anyhow::ensure!(
sequence == msg.sequence(),
"Message has incorrect sequence (exp: {} got: {})",
sequence,
msg.sequence()
);
self.account_sequences.insert(msg.from(), sequence + 1);
Ok(())
}
}
async fn check_block_messages(
state_manager: StateManager,
block: Arc<Block>,
base_tipset: Tipset,
) -> Result<(), TipsetSyncerError> {
let network_version = state_manager
.chain_config()
.network_version(block.header.epoch);
let eth_chain_id = state_manager.chain_config().eth_chain_id;
if let Some(sig) = &block.header().bls_aggregate {
let mut pub_keys = Vec::with_capacity(block.bls_msgs().len());
let mut cids = Vec::with_capacity(block.bls_msgs().len());
let db = state_manager.db();
for m in block.bls_msgs() {
let pk = StateManager::get_bls_public_key(db, m.from(), *base_tipset.parent_state())?;
pub_keys.push(pk);
cids.push(m.cid().to_bytes());
}
if !verify_bls_aggregate(
&cids.iter().map(|x| x.as_slice()).collect_vec(),
&pub_keys,
sig,
) {
return Err(TipsetSyncerError::BlsAggregateSignatureInvalid(
format!("{sig:?}"),
format!("{cids:?}"),
));
}
} else {
return Err(TipsetSyncerError::BlockWithoutBlsAggregate);
}
let ExecutedTipset { state_root, .. } = state_manager
.load_executed_tipset(&base_tipset)
.await
.map_err(|e| TipsetSyncerError::Calculation(format!("Could not update state: {e:#}")))?;
let tree = StateTree::new_from_root(state_manager.db(), &state_root).map_err(|e| {
TipsetSyncerError::Calculation(format!(
"Could not load from new state root in state manager: {e:#}"
))
})?;
let mut checker = MessageChecker {
price_list: price_list_by_network_version(network_version),
network_version,
tree,
sum_gas_limit: 0,
account_sequences: HashMap::default(),
};
for (i, msg) in block.bls_msgs().iter().enumerate() {
checker.check(msg).map_err(|e| {
TipsetSyncerError::Validation(format!(
"Block had invalid BLS message at index {i}: {e:#}"
))
})?;
}
for (i, msg) in block.secp_msgs().iter().enumerate() {
if network_version >= NetworkVersion::V14
&& !msg.signature().is_valid_secpk_sig_type(network_version)
{
return Err(TipsetSyncerError::SecpSignatureTypeInvalid(
i,
msg.signature().signature_type(),
));
}
if msg.signature().signature_type() == SignatureType::Delegated
&& !is_valid_eth_tx_for_sending(eth_chain_id, network_version, msg)
{
return Err(TipsetSyncerError::Validation(
"Network version must be at least NV23 for legacy Ethereum transactions".to_owned(),
));
}
checker.check(msg).map_err(|e| {
TipsetSyncerError::Validation(format!(
"block had an invalid secp message at index {i}: {e:#}"
))
})?;
let key_addr = state_manager
.resolve_to_deterministic_address(msg.from(), &base_tipset)
.await
.map_err(|e| TipsetSyncerError::ResolvingAddressFromMessage(e.to_string()))?;
msg.signature()
.authenticate_msg(eth_chain_id, msg, &key_addr)
.map_err(|e| TipsetSyncerError::MessageSignatureInvalid(e.to_string()))?;
}
let msg_root =
TipsetValidator::compute_msg_root(state_manager.db(), block.bls_msgs(), block.secp_msgs())
.map_err(|err| TipsetSyncerError::ComputingMessageRoot(err.to_string()))?;
if block.header().messages != msg_root {
return Err(TipsetSyncerError::BlockMessageRootInvalid(
format!("{:?}", block.header().messages),
format!("{msg_root:?}"),
));
}
Ok(())
}
fn block_sanity_checks(header: &CachingBlockHeader) -> Result<(), TipsetSyncerError> {
if header.signature.is_none() {
return Err(TipsetSyncerError::BlockWithoutSignature);
}
if header.bls_aggregate.is_none() {
return Err(TipsetSyncerError::BlockWithoutBlsAggregate);
}
Ok(())
}
fn block_timestamp_checks(header: &CachingBlockHeader) -> Result<(), TipsetSyncerError> {
let time_now = chrono::Utc::now().timestamp() as u64;
if header.timestamp > time_now.saturating_add(ALLOWABLE_CLOCK_DRIFT) {
return Err(TipsetSyncerError::TimeTravellingBlock(
time_now,
header.timestamp,
));
} else if header.timestamp > time_now {
warn!(
"Got block from the future, but within clock drift threshold, {} > {}",
header.timestamp, time_now
);
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn concat_preserves_parent_chain_state_mismatch() {
let concatenated = TipsetSyncerError::concat(nunny::vec![
TipsetSyncerError::Validation("a".into()),
TipsetSyncerError::ParentChainStateMismatch("b".into()),
]);
assert!(matches!(
concatenated,
TipsetSyncerError::ParentChainStateMismatch(_)
));
let concatenated =
TipsetSyncerError::concat(nunny::vec![TipsetSyncerError::Validation("a".into())]);
assert!(matches!(concatenated, TipsetSyncerError::Validation(_)));
}
mod block_messages {
use super::*;
use crate::blocks::RawBlockHeader;
use crate::db::MemoryDB;
use crate::message::SignedMessage;
use crate::networks::ChainConfig;
use crate::shim::crypto::{BLS_SIG_LEN, SECP_SIG_LEN, Signature};
use crate::shim::message::Message;
use crate::shim::state_tree::StateTreeVersion;
const GAS_LIMIT_PLACEHOLDER: u64 = 1_000_000;
fn chain_config() -> Arc<ChainConfig> {
Arc::new(ChainConfig::default())
}
fn last_epoch_before_nv14(chain_config: &ChainConfig) -> ChainEpoch {
chain_config.epoch(Height::Chocolate)
}
fn first_epoch_at_nv14(chain_config: &ChainConfig) -> ChainEpoch {
last_epoch_before_nv14(chain_config) + 1
}
fn first_epoch_at_nv28(chain_config: &ChainConfig) -> ChainEpoch {
chain_config.epoch(Height::FireHorse) + 1
}
fn signed_message(gas_limit: u64, signature: Signature) -> SignedMessage {
SignedMessage::new_unchecked(
Message::builder()
.to(Address::new_id(1))
.from(Address::new_id(2))
.gas_limit(gas_limit)
.build(),
signature,
)
}
fn secp_message(gas_limit: u64) -> SignedMessage {
signed_message(gas_limit, Signature::new_secp256k1(vec![0; SECP_SIG_LEN]))
}
fn empty_state_manager(chain_config: Arc<ChainConfig>) -> (StateManager, Tipset) {
let db = Arc::new(MemoryDB::default());
let genesis = CachingBlockHeader::new(RawBlockHeader {
timestamp: 7777,
..Default::default()
});
let chain_store = ChainStore::new(db, chain_config, genesis).unwrap();
let state_manager = StateManager::new(chain_store).unwrap();
let state_root = StateTree::new(state_manager.db(), StateTreeVersion::V5)
.unwrap()
.flush()
.unwrap();
let base_tipset = state_manager.chain_store().heaviest_tipset();
state_manager.insert_executed_tipset(
base_tipset.key().clone(),
ExecutedTipset {
state_root,
receipt_root: Cid::default(),
executed_messages: Arc::new(vec![]),
},
);
(state_manager, base_tipset)
}
fn block_with(epoch: ChainEpoch, secp_message: SignedMessage) -> Arc<Block> {
Arc::new(Block {
header: CachingBlockHeader::new(RawBlockHeader {
epoch,
bls_aggregate: Some(Signature::new_bls(vec![0; BLS_SIG_LEN])),
..Default::default()
}),
bls_messages: vec![],
secp_messages: vec![secp_message],
})
}
async fn validate(
chain_config: Arc<ChainConfig>,
epoch: ChainEpoch,
signature: Signature,
) -> TipsetSyncerError {
let (state_manager, base_tipset) = empty_state_manager(chain_config);
let message = signed_message(u64::from(u32::MAX), signature);
check_block_messages(state_manager, block_with(epoch, message), base_tipset)
.await
.unwrap_err()
}
#[tokio::test]
async fn bls_signature_type_is_rejected_in_the_secp_message_list() {
let chain_config = chain_config();
let epoch = first_epoch_at_nv14(&chain_config);
let err = validate(
chain_config,
epoch,
Signature::new_bls(vec![0; BLS_SIG_LEN]),
)
.await;
assert!(
matches!(
err,
TipsetSyncerError::SecpSignatureTypeInvalid(0, SignatureType::Bls)
),
"got: {err}"
);
}
#[tokio::test]
async fn secp256k1_signature_type_passes_the_guard() {
let chain_config = chain_config();
let epoch = first_epoch_at_nv14(&chain_config);
let err = validate(
chain_config,
epoch,
Signature::new_secp256k1(vec![0; SECP_SIG_LEN]),
)
.await;
assert!(
!matches!(err, TipsetSyncerError::SecpSignatureTypeInvalid(..)),
"the guard must not fire, got: {err}"
);
}
#[tokio::test]
async fn signature_type_is_not_checked_before_nv14() {
let chain_config = chain_config();
let epoch = last_epoch_before_nv14(&chain_config);
let err = validate(
chain_config,
epoch,
Signature::new_bls(vec![0; BLS_SIG_LEN]),
)
.await;
assert!(
!matches!(err, TipsetSyncerError::SecpSignatureTypeInvalid(..)),
"the guard must not fire, got: {err}"
);
}
#[tokio::test]
async fn secp_gas_floor_is_charged_over_the_signed_encoding() {
let (state_manager, base_tipset) = empty_state_manager(chain_config());
let epoch = first_epoch_at_nv28(state_manager.chain_config());
let price_list =
price_list_by_network_version(state_manager.chain_config().network_version(epoch));
let floor = |length| price_list.on_chain_message(length).total().round_up();
let placeholder = secp_message(GAS_LIMIT_PLACEHOLDER);
let signed_floor = floor(placeholder.chain_length().unwrap());
let unsigned_floor = floor(placeholder.vm_message().chain_length().unwrap());
assert!(signed_floor > unsigned_floor);
let underpaying = secp_message(unsigned_floor);
let paying = secp_message(signed_floor);
for message in [&underpaying, &paying] {
assert_eq!(
message.chain_length().unwrap(),
placeholder.chain_length().unwrap(),
"the gas floors must be exact for the messages they are applied to"
);
}
let err = check_block_messages(
state_manager.shallow_clone(),
block_with(epoch, underpaying),
base_tipset.shallow_clone(),
)
.await
.unwrap_err();
assert!(
err.to_string().contains("less than cost"),
"a gas limit covering only the unsigned encoding must be rejected, got: {err}"
);
let err = check_block_messages(state_manager, block_with(epoch, paying), base_tipset)
.await
.unwrap_err();
assert!(
err.to_string().contains("Actor does not exist in state"),
"a gas limit covering the signed encoding must clear the floor, got: {err}"
);
}
}
}