#![cfg(feature = "reactive")]
mod common;
use alloy_consensus::{BlockHeader as _, Header};
use alloy_eips::BlockId;
use alloy_primitives::{Address, B256, Bytes, Log as PrimitiveLog};
use alloy_rpc_types_eth::Log;
use anyhow::Result;
use common::setup_cache;
use evm_fork_cache::cache::BlockContextRequirements;
fn header(number: u64, basefee: Option<u64>) -> Header {
Header {
number,
timestamp: 1_700_000_000 + number,
base_fee_per_gas: basefee,
beneficiary: Address::repeat_byte(0xcb),
gas_limit: 30_000_000,
mix_hash: B256::repeat_byte(0xab),
..Default::default()
}
}
#[test]
fn strict_requirements_reject_header_missing_basefee() {
let reqs = BlockContextRequirements::strict();
assert!(
reqs.validate_header(&header(100, Some(7))).is_ok(),
"a complete header must satisfy strict requirements"
);
let err = reqs
.validate_header(&header(100, None))
.expect_err("strict must reject a header with no base fee");
assert!(
err.to_string().to_lowercase().contains("basefee")
|| err.to_string().to_lowercase().contains("base fee"),
"the error must name the missing base-fee field, got: {err}"
);
}
#[test]
fn lenient_requirements_accept_incomplete_header() {
let reqs = BlockContextRequirements::lenient();
assert!(reqs.validate_header(&header(100, None)).is_ok());
assert!(reqs.validate_header(&header(100, Some(7))).is_ok());
}
#[test]
fn per_field_requirements_allow_opting_out_of_basefee() {
let mut reqs = BlockContextRequirements::strict();
reqs.require_basefee = false;
assert!(
reqs.validate_header(&header(100, None)).is_ok(),
"opting out of the base-fee requirement must accept a header without one"
);
}
#[tokio::test]
async fn advance_block_refreshes_all_block_env_fields() -> Result<()> {
let mut cache = setup_cache().await?;
let h = header(12_345, Some(42));
cache
.advance_block(&h)
.expect("lenient advance_block over a complete header succeeds");
assert_eq!(cache.block_number(), Some(12_345));
assert_eq!(cache.basefee(), Some(42));
assert_eq!(cache.coinbase(), Some(Address::repeat_byte(0xcb)));
assert_eq!(cache.prevrandao(), Some(B256::repeat_byte(0xab)));
assert_eq!(cache.block_gas_limit(), Some(30_000_000));
assert_eq!(cache.timestamp(), Some(1_700_000_000 + 12_345));
assert_eq!(
cache.block(),
alloy_eips::BlockId::number(12_345),
"advance_block must re-pin RPC fetches to the advanced block"
);
Ok(())
}
#[tokio::test]
async fn set_block_clears_every_stale_header_field_on_repin() -> Result<()> {
let mut cache = setup_cache().await?;
cache
.advance_block(&header(100, Some(42)))
.expect("install complete old header");
cache.set_block(BlockId::number(101));
assert_eq!(cache.block(), BlockId::number(101));
assert_eq!(cache.block_number(), Some(101));
assert_eq!(cache.basefee(), None);
assert_eq!(cache.coinbase(), None);
assert_eq!(cache.prevrandao(), None);
assert_eq!(cache.block_gas_limit(), None);
assert_eq!(cache.timestamp(), None);
Ok(())
}
#[tokio::test]
async fn advance_block_strict_rejects_incomplete_header() -> Result<()> {
let mut cache = setup_cache().await?;
cache.set_block_context_requirements(BlockContextRequirements::strict());
let err = cache
.advance_block(&header(200, None))
.expect_err("strict advance_block must reject a header with no base fee");
assert!(
err.to_string().to_lowercase().contains("basefee")
|| err.to_string().to_lowercase().contains("base fee"),
"the error must name the missing base-fee field, got: {err}"
);
cache
.advance_block(&header(200, Some(9)))
.expect("strict advance_block over a complete header succeeds");
assert_eq!(cache.basefee(), Some(9));
Ok(())
}
use std::sync::Arc;
use alloy_network::{Ethereum, primitives::HeaderResponse as _};
use alloy_provider::RootProvider;
use alloy_provider::network::AnyNetwork;
use alloy_rpc_client::RpcClient;
use alloy_transport::mock::Asserter;
use evm_fork_cache::EvmCacheBuilder;
use evm_fork_cache::reactive::{
BlockRef, ChainControl, ChainStatus, InputSource, ReactiveConfig, ReactiveContext,
ReactiveInput, ReactiveInputBatch, ReactiveInputRecord, ReactiveReport, ReactiveRuntime,
};
fn mock_provider() -> Arc<RootProvider<AnyNetwork>> {
let client = RpcClient::mocked(Asserter::new());
Arc::new(RootProvider::<AnyNetwork>::new(client))
}
#[tokio::test]
async fn pinned_builder_captures_timestamp_from_fetched_header() -> Result<()> {
let asserter = Asserter::new();
let expected = header(12_345, Some(42));
let block: alloy_rpc_types_eth::Block =
alloy_rpc_types_eth::Block::empty(alloy_rpc_types_eth::Header::new(expected.clone()));
asserter.push_success(&Some(block));
let provider = Arc::new(RootProvider::<AnyNetwork>::new(RpcClient::mocked(asserter)));
let cache = EvmCacheBuilder::new(provider)
.block(BlockId::number(expected.number))
.chain_id(1)
.build()
.await;
assert_eq!(cache.block_number(), Some(expected.number));
assert_eq!(cache.timestamp(), Some(expected.timestamp));
Ok(())
}
#[tokio::test]
async fn try_build_strict_fails_when_header_unavailable() {
let result = EvmCacheBuilder::new(mock_provider())
.strict_block_context(true)
.try_build()
.await;
let err = match result {
Ok(_) => panic!("strict try_build over a header-less mock provider must error"),
Err(err) => err,
};
assert!(
err.to_string().to_lowercase().contains("fetch failed"),
"expected a fetch-failure error, got: {err}"
);
}
#[tokio::test]
async fn try_build_lenient_succeeds_without_header() -> Result<()> {
let cache = EvmCacheBuilder::new(mock_provider())
.strict_block_context(false)
.try_build()
.await?;
assert_eq!(cache.block_number(), None);
let _cache = EvmCacheBuilder::new(mock_provider()).try_build().await?;
Ok(())
}
fn rpc_header(number: u64, basefee: Option<u64>) -> alloy_rpc_types_eth::Header {
alloy_rpc_types_eth::Header::new(header(number, basefee))
}
fn included_header_context(header: &alloy_rpc_types_eth::Header) -> ReactiveContext {
let block = evm_fork_cache::reactive::BlockRef {
number: header.number(),
hash: header.hash(),
parent_hash: Some(header.parent_hash()),
timestamp: Some(header.timestamp()),
};
ReactiveContext {
chain_id: Some(1),
source: InputSource::Batch,
chain_status: ChainStatus::Included {
block,
confirmations: 0,
},
block: Some(block),
transaction_index: None,
log_index: None,
}
}
#[tokio::test]
async fn reactive_ingest_of_canonical_header_refreshes_block_env() -> Result<()> {
let mut cache = setup_cache().await?;
let mut runtime = ReactiveRuntime::<Ethereum>::new(ReactiveConfig::default());
let header = rpc_header(7_777, Some(123));
let context = included_header_context(&header);
let input = ReactiveInput::BlockHeader(header);
let batch = ReactiveInputBatch::new(vec![ReactiveInputRecord::new(input, context)]);
let report = runtime.ingest_batch(&mut cache, batch)?;
assert_eq!(cache.block_number(), Some(7_777));
assert_eq!(cache.basefee(), Some(123));
assert_eq!(cache.timestamp(), Some(1_700_000_000 + 7_777));
assert!(
!report
.reports
.iter()
.any(|r| matches!(r.as_ref(), ReactiveReport::Error(_))),
"a lenient canonical drive must not surface an error report"
);
Ok(())
}
#[tokio::test]
async fn post_record_compact_barrier_preserves_full_header_environment() -> Result<()> {
let mut cache = setup_cache().await?;
let mut runtime = ReactiveRuntime::<Ethereum>::new(ReactiveConfig::default());
let header = rpc_header(7_778, Some(124));
let context = included_header_context(&header);
let exact_hash = header.hash();
let compact = BlockRef {
number: header.number(),
hash: exact_hash,
parent_hash: None,
timestamp: None,
};
runtime.ingest_batch(
&mut cache,
ReactiveInputBatch::new(vec![ReactiveInputRecord::new(
ReactiveInput::BlockHeader(header),
context,
)])
.with_chain_controls([ChainControl::Barrier {
id: b"header-complete".to_vec(),
block: Some(compact),
}]),
)?;
assert_eq!(cache.block(), BlockId::from((exact_hash, Some(true))));
assert_eq!(cache.block_number(), Some(7_778));
assert_eq!(cache.basefee(), Some(124));
assert_eq!(cache.coinbase(), Some(Address::repeat_byte(0xcb)));
assert_eq!(cache.prevrandao(), Some(B256::repeat_byte(0xab)));
assert_eq!(cache.block_gas_limit(), Some(30_000_000));
assert_eq!(cache.timestamp(), Some(1_700_000_000 + 7_778));
Ok(())
}
#[tokio::test]
async fn zero_depth_runtime_preserves_full_header_env_for_same_block_compact_records() -> Result<()>
{
let mut cache = setup_cache().await?;
let mut runtime = ReactiveRuntime::<Ethereum>::new(ReactiveConfig {
journal_depth: 0,
..ReactiveConfig::default()
});
let header = rpc_header(7_779, Some(125));
let context = included_header_context(&header);
let block = context.block.expect("canonical block");
let log = Log {
inner: PrimitiveLog::new_unchecked(
Address::repeat_byte(0xdd),
vec![B256::repeat_byte(0xee)],
Bytes::new(),
),
block_hash: Some(block.hash),
block_number: Some(block.number),
block_timestamp: block.timestamp,
transaction_hash: Some(B256::repeat_byte(0xef)),
transaction_index: Some(0),
log_index: Some(0),
removed: false,
};
let log_context = ReactiveContext {
transaction_index: Some(0),
log_index: Some(0),
..context.clone()
};
runtime.ingest_batch(
&mut cache,
ReactiveInputBatch::new(vec![
ReactiveInputRecord::new(ReactiveInput::BlockHeader(header), context),
ReactiveInputRecord::new(ReactiveInput::Log(log), log_context),
]),
)?;
assert_eq!(runtime.last_canonical_block(), Some(block));
assert_eq!(cache.basefee(), Some(125));
assert_eq!(cache.coinbase(), Some(Address::repeat_byte(0xcb)));
assert_eq!(cache.prevrandao(), Some(B256::repeat_byte(0xab)));
assert_eq!(cache.block_gas_limit(), Some(30_000_000));
assert_eq!(cache.timestamp(), Some(1_700_000_000 + 7_779));
Ok(())
}
#[tokio::test]
async fn reactive_ingest_of_pending_header_does_not_refresh_block_env() -> Result<()> {
let mut cache = setup_cache().await?;
let mut runtime = ReactiveRuntime::<Ethereum>::new(ReactiveConfig::default());
let ctx = ReactiveContext {
chain_id: Some(1),
source: InputSource::Subscription,
chain_status: ChainStatus::Pending,
block: None,
transaction_index: None,
log_index: None,
};
let input = ReactiveInput::BlockHeader(rpc_header(9_999, Some(55)));
let batch = ReactiveInputBatch::new(vec![ReactiveInputRecord::new(input, ctx)]);
runtime.ingest_batch(&mut cache, batch)?;
assert_eq!(cache.block_number(), None);
assert_eq!(cache.basefee(), None);
Ok(())
}
#[tokio::test]
async fn reactive_strict_drive_surfaces_error_report_for_incomplete_header() -> Result<()> {
let mut cache = setup_cache().await?;
cache.set_block_context_requirements(BlockContextRequirements::strict());
let mut runtime = ReactiveRuntime::<Ethereum>::new(ReactiveConfig::default());
let header = rpc_header(4_242, None);
let context = included_header_context(&header);
let input = ReactiveInput::BlockHeader(header);
let batch = ReactiveInputBatch::new(vec![ReactiveInputRecord::new(input, context)]);
let report = runtime.ingest_batch(&mut cache, batch)?;
let error_message = report
.reports
.iter()
.find_map(|r| match r.as_ref() {
ReactiveReport::Error(e) => Some(e.message.clone()),
_ => None,
})
.expect("strict drive over an incomplete header must surface an error report");
assert!(
error_message.to_lowercase().contains("basefee"),
"the error report must name the missing base-fee field, got: {error_message}"
);
Ok(())
}
#[tokio::test]
async fn canonical_block_records_are_sorted_before_advancing_runtime_and_cache_heads() -> Result<()>
{
let mut cache = setup_cache().await?;
let mut runtime = ReactiveRuntime::<Ethereum>::new(ReactiveConfig::default());
let older = rpc_header(50, Some(5));
let newer = rpc_header(51, Some(6));
let older_context = included_header_context(&older);
let newer_context = included_header_context(&newer);
runtime.ingest_batch(
&mut cache,
ReactiveInputBatch::new(vec![
ReactiveInputRecord::new(ReactiveInput::BlockHeader(newer), newer_context),
ReactiveInputRecord::new(ReactiveInput::BlockHeader(older), older_context),
]),
)?;
assert_eq!(
runtime.last_canonical_block().map(|block| block.number),
Some(51)
);
assert_eq!(cache.block_number(), Some(51));
assert_eq!(cache.basefee(), Some(6));
Ok(())
}