diff --git a/crates/dojo/test-utils/src/sequencer.rs b/crates/dojo/test-utils/src/sequencer.rs index 4393bf3cd4..69966ea34e 100644 --- a/crates/dojo/test-utils/src/sequencer.rs +++ b/crates/dojo/test-utils/src/sequencer.rs @@ -6,6 +6,7 @@ use katana_core::constants::DEFAULT_SEQUENCER_ADDRESS; use katana_executor::implementation::blockifier::BlockifierFactory; use katana_node::config::dev::DevConfig; use katana_node::config::rpc::{RpcConfig, DEFAULT_RPC_ADDR}; +use katana_node::config::sequencing::SequencingConfig; pub use katana_node::config::*; use katana_node::LaunchedNode; use katana_primitives::chain::ChainId; diff --git a/crates/katana/chain-spec/src/rollup/utils.rs b/crates/katana/chain-spec/src/rollup/utils.rs index b152621c0e..a064bbecc0 100644 --- a/crates/katana/chain-spec/src/rollup/utils.rs +++ b/crates/katana/chain-spec/src/rollup/utils.rs @@ -336,7 +336,7 @@ mod tests { use alloy_primitives::U256; use katana_executor::implementation::blockifier::BlockifierFactory; - use katana_executor::ExecutorFactory; + use katana_executor::{BlockLimits, ExecutorFactory}; use katana_primitives::chain::ChainId; use katana_primitives::contract::Nonce; use katana_primitives::env::CfgEnv; @@ -386,6 +386,7 @@ mod tests { ..Default::default() }, Default::default(), + BlockLimits::max(), ) } diff --git a/crates/katana/cli/src/args.rs b/crates/katana/cli/src/args.rs index 28a0349702..abc0025e85 100644 --- a/crates/katana/cli/src/args.rs +++ b/crates/katana/cli/src/args.rs @@ -19,7 +19,8 @@ use katana_node::config::metrics::MetricsConfig; use katana_node::config::rpc::RpcConfig; #[cfg(feature = "server")] use katana_node::config::rpc::{RpcModuleKind, RpcModulesList}; -use katana_node::config::{Config, SequencingConfig}; +use katana_node::config::sequencing::SequencingConfig; +use katana_node::config::Config; use katana_primitives::chain::ChainId; use katana_primitives::genesis::allocation::DevAllocationsGenerator; use katana_primitives::genesis::constant::DEFAULT_PREFUNDED_ACCOUNT_BALANCE; @@ -58,6 +59,10 @@ pub struct NodeArgs { #[arg(value_name = "MILLISECONDS")] pub block_time: Option, + #[arg(long = "sequencing.block-max-cairo-steps")] + #[arg(value_name = "TOTAL")] + pub block_cairo_steps_limit: Option, + /// Directory path of the database to initialize from. /// /// The path must either be an empty directory or a directory which already contains a @@ -187,7 +192,11 @@ impl NodeArgs { } fn sequencer_config(&self) -> SequencingConfig { - SequencingConfig { block_time: self.block_time, no_mining: self.no_mining } + SequencingConfig { + block_time: self.block_time, + no_mining: self.no_mining, + block_cairo_steps_limit: self.block_cairo_steps_limit, + } } fn rpc_config(&self) -> Result { diff --git a/crates/katana/core/src/service/block_producer.rs b/crates/katana/core/src/service/block_producer.rs index 81af6d1887..dbf7d92934 100644 --- a/crates/katana/core/src/service/block_producer.rs +++ b/crates/katana/core/src/service/block_producer.rs @@ -1,3 +1,14 @@ +//! ********************************************************************************************** +//! +//! "We are all in the gutter, but some of us are looking at the stars." +//! — Oscar Wilde, Lady Windermere's Fan +//! +//! Within this imperfect realm lies a spark of aspiration. What you find may be tangled +//! and weathered, but in its heart beats the rhythm of possibility. Tread gently, dear +//! wanderer, and perhaps together we can guide it toward those distant stars. +//! +//! ********************************************************************************************** + use std::collections::VecDeque; use std::future::Future; use std::pin::Pin; @@ -46,6 +57,16 @@ pub enum BlockProductionError { TransactionExecutionError(#[from] katana_executor::ExecutorError), } +impl BlockProductionError { + /// Returns `true` if the error is caused by block limit being exhausted. + pub fn is_block_limit_exhausted(&self) -> bool { + matches!( + self, + Self::TransactionExecutionError(katana_executor::ExecutorError::LimitsExhausted) + ) + } +} + #[derive(Debug, Clone)] pub struct MinedBlockOutcome { pub block_number: u64, @@ -65,7 +86,8 @@ type ServiceFuture = Pin> + Sen type BlockProductionResult = Result; type BlockProductionFuture = ServiceFuture>; -type TxExecutionResult = Result, BlockProductionError>; +type TxExecutionResult = + Result<(Vec, Option>), BlockProductionError>; type TxExecutionFuture = ServiceFuture; type BlockProductionWithTxnsFuture = @@ -213,6 +235,8 @@ pub struct IntervalBlockProducer { /// - `block_time` is `Some`, /// - and, at least one transaction has been executed and thus a new block is opened. timer: Option, + + is_block_full: bool, } impl IntervalBlockProducer { @@ -236,6 +260,7 @@ impl IntervalBlockProducer { TxValidator::new(state, flags.clone(), cfg.clone(), block_env, permit.clone()); Self { + is_block_full: false, validator, permit, backend, @@ -309,12 +334,11 @@ impl IntervalBlockProducer { fn execute_transactions( executor: PendingExecutor, - transactions: Vec, + mut transactions: Vec, ) -> TxExecutionResult { let executor = &mut executor.write(); - let new_txs_count = transactions.len(); - executor.execute_transactions(transactions)?; + let (total_executed, is_full) = executor.execute_transactions(transactions.clone())?; let txs = executor.transactions(); let total_txs = txs.len(); @@ -322,7 +346,7 @@ impl IntervalBlockProducer { // Take only the results of the newly executed transactions let results = txs .iter() - .skip(total_txs - new_txs_count) + .skip(total_txs.saturating_sub(total_executed)) .filter_map(|(tx, res)| match res { ExecutionResult::Failed { .. } => None, ExecutionResult::Success { receipt, trace, .. } => Some(TxWithOutcome { @@ -333,7 +357,10 @@ impl IntervalBlockProducer { }) .collect::>(); - Ok(results) + let non_executed_txs = + if is_full.is_some() { Some(transactions.split_off(total_executed)) } else { None }; + + Ok((results, non_executed_txs)) } fn create_new_executor_for_next_block(&self) -> Result { @@ -393,7 +420,17 @@ impl Stream for IntervalBlockProducer { if let Some(mut timer) = pin.timer.take() { // Mine block if the interval is over - if timer.poll_tick(cx).is_ready() && pin.ongoing_mining.is_none() { + // + // if block is already full but the timer hasn't ready yet, we will still mine but we + // don't have to do anything to the timer as it will be dropped and reset once new + // transaction is executed. + if (timer.poll_tick(cx).is_ready() || pin.is_block_full) && pin.ongoing_mining.is_none() + { + if pin.is_block_full { + info!("Block has reached capacity! Closing block..."); + pin.is_block_full = false; + } + pin.ongoing_mining = Some(Box::pin({ let executor = pin.executor.clone(); let backend = pin.backend.clone(); @@ -402,9 +439,21 @@ impl Stream for IntervalBlockProducer { pin.blocking_task_spawner.spawn(|| Self::do_mine(permit, executor, backend)) })); } else { - // Unable to close the block due to ongoing mining. pin.timer = Some(timer); } + } else if pin.is_block_full && pin.ongoing_mining.is_none() { + info!("Block has reached capacity! Closing block..."); + + pin.ongoing_mining = Some(Box::pin({ + let executor = pin.executor.clone(); + let backend = pin.backend.clone(); + let permit = pin.permit.clone(); + + pin.blocking_task_spawner.spawn(|| Self::do_mine(permit, executor, backend)) + })); + + pin.is_block_full = false; + pin.timer = None; } loop { @@ -438,7 +487,18 @@ impl Stream for IntervalBlockProducer { if let Some(mut execution) = pin.ongoing_execution.take() { if let Poll::Ready(executor) = execution.poll_unpin(cx) { match executor { - Ok(Ok(txs)) => { + Ok(Ok((txs, leftovers))) => { + if let Some(leftovers) = leftovers { + pin.is_block_full = true; + + // Push leftover transactions back to front of queue + pin.queued.push_front(leftovers); + + // Schedule future poll if block is full + cx.waker().wake_by_ref(); + break; + } + pin.notify_listener(txs); continue; } @@ -486,6 +546,8 @@ impl Stream for IntervalBlockProducer { Err(e) => return Poll::Ready(Some(Err(e))), } + pin.is_block_full = false; + return Poll::Ready(Some(outcome)); } diff --git a/crates/katana/core/src/service/block_producer_tests.rs b/crates/katana/core/src/service/block_producer_tests.rs index 58971aea80..25c91eb10d 100644 --- a/crates/katana/core/src/service/block_producer_tests.rs +++ b/crates/katana/core/src/service/block_producer_tests.rs @@ -5,7 +5,6 @@ use katana_executor::implementation::noop::NoopExecutorFactory; use katana_primitives::transaction::{ExecutableTx, InvokeTx}; use katana_primitives::Felt; use katana_provider::providers::db::DbProvider; -use tokio::time; use super::*; use crate::backend::gas_oracle::GasOracle; @@ -47,7 +46,7 @@ async fn interval_force_mine_without_transactions() { async fn interval_mine_after_timer() { let backend = test_backend(); let mut producer = IntervalBlockProducer::new(backend.clone(), Some(1000)); - // Initial state + // no timer should be set when no block is opened. assert!(producer.timer.is_none()); producer.queued.push_back(vec![dummy_transaction()]); @@ -55,19 +54,23 @@ async fn interval_mine_after_timer() { let stream = producer; pin_mut!(stream); - // Process the transaction, the timer should be automatically started - let _ = stream.next().await; - assert!(stream.timer.is_some()); + let waker = futures::task::noop_waker(); + let mut context = Context::from_waker(&waker); - // Advance time to trigger mining - time::sleep(Duration::from_secs(1)).await; - let result = stream.next().await.expect("should mine block").unwrap(); + // mine the block + let poll_result = stream.as_mut().poll_next(&mut context); - assert_eq!(result.block_number, 1); - assert_eq!(backend.blockchain.provider().latest_number().unwrap(), 1); + // based on how the `Stream` trait is implemented, there is a possibility that a single + // call to `poll_next` can complete the whole production flow so we added this just in case. + if poll_result.is_pending() { + assert!(stream.timer.is_some(), "timer should start once we received a tx"); + } else { + assert!(stream.timer.is_none(), "no timer if block has been mined"); + } - // Final state - assert!(stream.timer.is_none()); + let outcome = stream.next().await.expect("should mine block").unwrap(); + assert_eq!(outcome.block_number, 1); + assert_eq!(backend.blockchain.provider().latest_number().unwrap(), 1); } // Helper functions to create test transactions diff --git a/crates/katana/core/tests/backend.rs b/crates/katana/core/tests/backend.rs index 1200949bb7..ee70c3478d 100644 --- a/crates/katana/core/tests/backend.rs +++ b/crates/katana/core/tests/backend.rs @@ -5,6 +5,7 @@ use katana_core::backend::gas_oracle::GasOracle; use katana_core::backend::storage::{Blockchain, Database}; use katana_core::backend::Backend; use katana_executor::implementation::blockifier::BlockifierFactory; +use katana_executor::BlockLimits; use katana_primitives::chain::ChainId; use katana_primitives::env::CfgEnv; use katana_primitives::felt; @@ -25,6 +26,7 @@ fn executor(chain_spec: &ChainSpec) -> BlockifierFactory { ..Default::default() }, Default::default(), + BlockLimits::max(), ) } diff --git a/crates/katana/executor/benches/concurrent.rs b/crates/katana/executor/benches/concurrent.rs index 50af452d31..75b9579a36 100644 --- a/crates/katana/executor/benches/concurrent.rs +++ b/crates/katana/executor/benches/concurrent.rs @@ -9,7 +9,7 @@ use std::time::Duration; use criterion::measurement::WallTime; use criterion::{criterion_group, criterion_main, BatchSize, BenchmarkGroup, Criterion}; use katana_executor::implementation::blockifier::BlockifierFactory; -use katana_executor::{ExecutionFlags, ExecutorFactory}; +use katana_executor::{BlockLimits, ExecutionFlags, ExecutorFactory}; use katana_primitives::env::{BlockEnv, CfgEnv}; use katana_primitives::transaction::ExecutableTxWithHash; use katana_provider::test_utils; @@ -44,7 +44,7 @@ fn blockifier( (block_env, cfg_env): (BlockEnv, CfgEnv), tx: ExecutableTxWithHash, ) { - let factory = Arc::new(BlockifierFactory::new(cfg_env, flags.clone())); + let factory = Arc::new(BlockifierFactory::new(cfg_env, flags.clone(), BlockLimits::max())); group.bench_function("Blockifier.1", |b| { b.iter_batched( diff --git a/crates/katana/executor/benches/execution.rs b/crates/katana/executor/benches/execution.rs index cdf19e59cd..1ddff343e2 100644 --- a/crates/katana/executor/benches/execution.rs +++ b/crates/katana/executor/benches/execution.rs @@ -49,7 +49,9 @@ fn blockifier( (state, &block_context, execution_flags, tx.clone()) }, - |(mut state, block_context, flags, tx)| transact(&mut state, block_context, flags, tx), + |(mut state, block_context, flags, tx)| { + transact(&mut state, block_context, flags, tx, None) + }, BatchSize::SmallInput, ) }); diff --git a/crates/katana/executor/src/abstraction/error.rs b/crates/katana/executor/src/abstraction/error.rs index b4d02d2271..b9e5ca6d53 100644 --- a/crates/katana/executor/src/abstraction/error.rs +++ b/crates/katana/executor/src/abstraction/error.rs @@ -4,7 +4,13 @@ use katana_primitives::Felt; /// Errors that can be returned by the executor. #[derive(Debug, thiserror::Error)] -pub enum ExecutorError {} +pub enum ExecutorError { + #[error("Limits exhausted")] + LimitsExhausted, + + #[error(transparent)] + Other(Box), +} /// Errors that can occur during the transaction execution. #[derive(Debug, Clone, thiserror::Error)] diff --git a/crates/katana/executor/src/abstraction/executor.rs b/crates/katana/executor/src/abstraction/executor.rs index 525880babe..9bc8213b8f 100644 --- a/crates/katana/executor/src/abstraction/executor.rs +++ b/crates/katana/executor/src/abstraction/executor.rs @@ -5,6 +5,7 @@ use katana_primitives::transaction::{ExecutableTxWithHash, TxWithHash}; use katana_primitives::Felt; use katana_provider::traits::state::StateProvider; +use super::ExecutorError; use crate::{ EntryPointCall, ExecutionError, ExecutionFlags, ExecutionOutput, ExecutionResult, ExecutorResult, ResultAndStates, @@ -38,10 +39,11 @@ pub trait BlockExecutor<'a>: ExecutorExt + Send + Sync + core::fmt::Debug { /// Executes the given block. fn execute_block(&mut self, block: ExecutableBlock) -> ExecutorResult<()>; + /// Execute transactions and returns the total number of transactions that was executed. fn execute_transactions( &mut self, transactions: Vec, - ) -> ExecutorResult<()>; + ) -> ExecutorResult<(usize, Option)>; /// Takes the output state of the executor. fn take_execution_output(&mut self) -> ExecutorResult; diff --git a/crates/katana/executor/src/abstraction/mod.rs b/crates/katana/executor/src/abstraction/mod.rs index d9a0e99af4..022ab83868 100644 --- a/crates/katana/executor/src/abstraction/mod.rs +++ b/crates/katana/executor/src/abstraction/mod.rs @@ -17,6 +17,19 @@ use katana_trie::MultiProof; pub type ExecutorResult = Result; +/// See . +#[derive(Debug, Clone, Default)] +pub struct BlockLimits { + /// The maximum number of Cairo steps that can be completed within each block. + pub cairo_steps: u64, +} + +impl BlockLimits { + pub fn max() -> Self { + Self { cairo_steps: u64::MAX } + } +} + /// Transaction execution simulation flags. /// /// These flags can be used to control the behavior of the transaction execution, such as skipping diff --git a/crates/katana/executor/src/implementation/blockifier/error.rs b/crates/katana/executor/src/implementation/blockifier/error.rs index 89b4d8ab4c..4e05c8afb7 100644 --- a/crates/katana/executor/src/implementation/blockifier/error.rs +++ b/crates/katana/executor/src/implementation/blockifier/error.rs @@ -1,3 +1,4 @@ +use blockifier::blockifier::transaction_executor::TransactionExecutorError; use blockifier::execution::errors::{EntryPointExecutionError, PreExecutionError}; use blockifier::execution::execution_utils::format_panic_data; use blockifier::state::errors::StateError; @@ -6,7 +7,7 @@ use blockifier::transaction::errors::{ }; use crate::implementation::blockifier::utils::to_address; -use crate::ExecutionError; +use crate::{ExecutionError, ExecutorError}; impl From for ExecutionError { fn from(error: TransactionExecutionError) -> Self { @@ -100,3 +101,13 @@ impl From for ExecutionError { } } } + +impl From for ExecutorError { + fn from(value: TransactionExecutorError) -> Self { + match value { + TransactionExecutorError::BlockFull => Self::LimitsExhausted, + TransactionExecutorError::StateError(e) => Self::Other(e.into()), + TransactionExecutorError::TransactionExecutionError(e) => Self::Other(e.into()), + } + } +} diff --git a/crates/katana/executor/src/implementation/blockifier/mod.rs b/crates/katana/executor/src/implementation/blockifier/mod.rs index d515beb5a9..453cb66c1c 100644 --- a/crates/katana/executor/src/implementation/blockifier/mod.rs +++ b/crates/katana/executor/src/implementation/blockifier/mod.rs @@ -1,5 +1,6 @@ // Re-export the blockifier crate. pub use blockifier; +use blockifier::bouncer::{Bouncer, BouncerConfig, BouncerWeights}; mod error; mod state; @@ -22,9 +23,9 @@ use tracing::info; use self::state::CachedState; use crate::{ - BlockExecutor, EntryPointCall, ExecutionError, ExecutionFlags, ExecutionOutput, - ExecutionResult, ExecutionStats, ExecutorExt, ExecutorFactory, ExecutorResult, ResultAndStates, - StateProviderDb, + BlockExecutor, BlockLimits, EntryPointCall, ExecutionError, ExecutionFlags, ExecutionOutput, + ExecutionResult, ExecutionStats, ExecutorError, ExecutorExt, ExecutorFactory, ExecutorResult, + ResultAndStates, StateProviderDb, }; pub(crate) const LOG_TARGET: &str = "katana::executor::blockifier"; @@ -33,12 +34,13 @@ pub(crate) const LOG_TARGET: &str = "katana::executor::blockifier"; pub struct BlockifierFactory { cfg: CfgEnv, flags: ExecutionFlags, + limits: BlockLimits, } impl BlockifierFactory { /// Create a new factory with the given configuration and simulation flags. - pub fn new(cfg: CfgEnv, flags: ExecutionFlags) -> Self { - Self { cfg, flags } + pub fn new(cfg: CfgEnv, flags: ExecutionFlags, limits: BlockLimits) -> Self { + Self { cfg, flags, limits } } } @@ -60,7 +62,8 @@ impl ExecutorFactory for BlockifierFactory { { let cfg_env = self.cfg.clone(); let flags = self.flags.clone(); - Box::new(StarknetVMProcessor::new(Box::new(state), block_env, cfg_env, flags)) + let limits = self.limits.clone(); + Box::new(StarknetVMProcessor::new(Box::new(state), block_env, cfg_env, flags, limits)) } fn cfg(&self) -> &CfgEnv { @@ -80,6 +83,7 @@ pub struct StarknetVMProcessor<'a> { transactions: Vec<(TxWithHash, ExecutionResult)>, simulation_flags: ExecutionFlags, stats: ExecutionStats, + bouncer: Bouncer, } impl<'a> StarknetVMProcessor<'a> { @@ -88,11 +92,24 @@ impl<'a> StarknetVMProcessor<'a> { block_env: BlockEnv, cfg_env: CfgEnv, simulation_flags: ExecutionFlags, + limits: BlockLimits, ) -> Self { let transactions = Vec::new(); let block_context = utils::block_context_from_envs(&block_env, &cfg_env); let state = state::CachedState::new(StateProviderDb::new(state)); - Self { block_context, state, transactions, simulation_flags, stats: Default::default() } + + let mut block_max_capacity = BouncerWeights::max(); + block_max_capacity.n_steps = limits.cairo_steps as usize; + let bouncer = Bouncer::new(BouncerConfig { block_max_capacity }); + + Self { + state, + transactions, + block_context, + simulation_flags, + stats: Default::default(), + bouncer, + } } fn fill_block_env_from_header(&mut self, header: &PartialHeader) { @@ -148,7 +165,9 @@ impl<'a> StarknetVMProcessor<'a> { let mut results = Vec::with_capacity(transactions.len()); for exec_tx in transactions { let tx = TxWithHash::from(&exec_tx); - let res = utils::transact(&mut state, block_context, flags, exec_tx); + // Safe to unwrap here because the only way the call to `transact` can return an error + // is when bouncer is `Some`. + let res = utils::transact(&mut state, block_context, flags, exec_tx, None).unwrap(); results.push(op(&mut state, (tx, res))); } @@ -166,11 +185,12 @@ impl<'a> BlockExecutor<'a> for StarknetVMProcessor<'a> { fn execute_transactions( &mut self, transactions: Vec, - ) -> ExecutorResult<()> { + ) -> ExecutorResult<(usize, Option)> { let block_context = &self.block_context; let flags = &self.simulation_flags; let mut state = self.state.0.lock(); + let mut total_executed = 0; for exec_tx in transactions { // Collect class artifacts if its a declare tx let class_decl_artifacts = if let ExecutableTx::Declare(tx) = exec_tx.as_ref() { @@ -182,34 +202,48 @@ impl<'a> BlockExecutor<'a> for StarknetVMProcessor<'a> { let tx = TxWithHash::from(&exec_tx); let hash = tx.hash; - let res = utils::transact(&mut state.inner, block_context, flags, exec_tx); - - match &res { - ExecutionResult::Success { receipt, trace } => { - self.stats.l1_gas_used += receipt.fee().gas_consumed; - self.stats.cairo_steps_used += - receipt.resources_used().vm_resources.n_steps as u128; - - if let Some(reason) = receipt.revert_reason() { - info!(target: LOG_TARGET, hash = format!("{hash:#x}"), %reason, "Transaction reverted."); + let result = utils::transact( + &mut state.inner, + block_context, + flags, + exec_tx, + Some(&mut self.bouncer), + ); + + match result { + Ok(exec_result) => { + match &exec_result { + ExecutionResult::Success { receipt, trace } => { + self.stats.l1_gas_used += receipt.fee().gas_consumed; + self.stats.cairo_steps_used += + receipt.resources_used().vm_resources.n_steps as u128; + + if let Some(reason) = receipt.revert_reason() { + info!(target: LOG_TARGET, hash = format!("{hash:#x}"), %reason, "Transaction reverted."); + } + + if let Some((class_hash, class)) = class_decl_artifacts { + state.declared_classes.insert(class_hash, class.as_ref().clone()); + } + + crate::utils::log_resources(&trace.actual_resources); + } + + ExecutionResult::Failed { error } => { + info!(target: LOG_TARGET, hash = format!("{hash:#x}"), %error, "Executing transaction."); + } } - if let Some((class_hash, class)) = class_decl_artifacts { - state.declared_classes.insert(class_hash, class.as_ref().clone()); - } - - crate::utils::log_resources(&trace.actual_resources); + total_executed += 1; + self.transactions.push((tx, exec_result)); } - ExecutionResult::Failed { error } => { - info!(target: LOG_TARGET, hash = format!("{hash:#x}"), %error, "Executing transaction."); - } + Err(e @ ExecutorError::LimitsExhausted) => return Ok((total_executed, Some(e))), + Err(e) => return Err(e), }; - - self.transactions.push((tx, res)); } - Ok(()) + Ok((total_executed, None)) } fn take_execution_output(&mut self) -> ExecutorResult { diff --git a/crates/katana/executor/src/implementation/blockifier/utils.rs b/crates/katana/executor/src/implementation/blockifier/utils.rs index 028590abaf..8b0cc8f9ab 100644 --- a/crates/katana/executor/src/implementation/blockifier/utils.rs +++ b/crates/katana/executor/src/implementation/blockifier/utils.rs @@ -3,7 +3,7 @@ use std::num::NonZeroU128; use std::sync::Arc; use blockifier::blockifier::block::{BlockInfo, GasPrices}; -use blockifier::bouncer::BouncerConfig; +use blockifier::bouncer::{Bouncer, BouncerConfig}; use blockifier::context::{BlockContext, ChainInfo, FeeTokenAddresses, TransactionContext}; use blockifier::execution::call_info::{ CallExecution, CallInfo, OrderedEvent, OrderedL2ToL1Message, @@ -14,8 +14,8 @@ use blockifier::execution::contract_class::{ }; use blockifier::execution::entry_point::{CallEntryPoint, CallType, EntryPointExecutionContext}; use blockifier::fee::fee_utils::get_fee_by_gas_vector; -use blockifier::state::cached_state; -use blockifier::state::state_api::StateReader; +use blockifier::state::cached_state::{self, TransactionalState}; +use blockifier::state::state_api::{StateReader, UpdatableState}; use blockifier::transaction::account_transaction::AccountTransaction; use blockifier::transaction::objects::{ DeprecatedTransactionInfo, FeeType, HasRelatedFeeType, TransactionExecutionInfo, @@ -58,16 +58,17 @@ use starknet::core::utils::parse_cairo_short_string; use super::state::{CachedState, StateDb}; use crate::abstraction::{EntryPointCall, ExecutionFlags}; use crate::utils::build_receipt; -use crate::{ExecutionError, ExecutionResult}; +use crate::{ExecutionError, ExecutionResult, ExecutorResult}; pub fn transact( state: &mut cached_state::CachedState, block_context: &BlockContext, simulation_flags: &ExecutionFlags, tx: ExecutableTxWithHash, -) -> ExecutionResult { - fn transact_inner( - state: &mut cached_state::CachedState, + bouncer: Option<&mut Bouncer>, +) -> ExecutorResult { + fn transact_inner( + state: &mut U, block_context: &BlockContext, simulation_flags: &ExecutionFlags, tx: Transaction, @@ -122,15 +123,36 @@ pub fn transact( Ok((info, fee_info)) } - match transact_inner(state, block_context, simulation_flags, to_executor_tx(tx.clone())) { + let transaction = to_executor_tx(tx.clone()); + let mut tx_state = TransactionalState::create_transactional(state); + let result = transact_inner(&mut tx_state, block_context, simulation_flags, transaction); + + match result { Ok((info, fee)) => { + if let Some(bouncer) = bouncer { + let tx_state_changes_keys = + tx_state.get_actual_state_changes().unwrap().into_keys(); + + bouncer.try_update( + &tx_state, + &tx_state_changes_keys, + &info.summarize(), + &info.transaction_receipt.resources, + )?; + } + + tx_state.commit(); + // get the trace and receipt from the execution info let trace = to_exec_info(info, tx.r#type()); let receipt = build_receipt(tx.tx_ref(), fee, &trace); - ExecutionResult::new_success(receipt, trace) + Ok(ExecutionResult::new_success(receipt, trace)) } - Err(e) => ExecutionResult::new_failed(e), + Err(e) => { + tx_state.commit(); + Ok(ExecutionResult::new_failed(e)) + } } } diff --git a/crates/katana/executor/src/implementation/noop.rs b/crates/katana/executor/src/implementation/noop.rs index 0392e633eb..c1fbfce33e 100644 --- a/crates/katana/executor/src/implementation/noop.rs +++ b/crates/katana/executor/src/implementation/noop.rs @@ -13,7 +13,7 @@ use crate::abstraction::{ BlockExecutor, EntryPointCall, ExecutionFlags, ExecutionOutput, ExecutionResult, ExecutorExt, ExecutorFactory, ExecutorResult, ResultAndStates, }; -use crate::ExecutionError; +use crate::{ExecutionError, ExecutorError}; /// A no-op executor factory. Creates an executor that does nothing. #[derive(Debug, Default)] @@ -101,9 +101,8 @@ impl<'a> BlockExecutor<'a> for NoopExecutor { fn execute_transactions( &mut self, transactions: Vec, - ) -> ExecutorResult<()> { - let _ = transactions; - Ok(()) + ) -> ExecutorResult<(usize, Option)> { + Ok((transactions.len(), None)) } fn take_execution_output(&mut self) -> ExecutorResult { diff --git a/crates/katana/executor/tests/fixtures/mod.rs b/crates/katana/executor/tests/fixtures/mod.rs index 0b21b0541c..903f858d94 100644 --- a/crates/katana/executor/tests/fixtures/mod.rs +++ b/crates/katana/executor/tests/fixtures/mod.rs @@ -266,12 +266,12 @@ pub fn executor_factory( #[cfg(feature = "blockifier")] pub mod blockifier { use katana_executor::implementation::blockifier::BlockifierFactory; - use katana_executor::ExecutionFlags; + use katana_executor::{BlockLimits, ExecutionFlags}; use super::{cfg, flags, CfgEnv}; #[rstest::fixture] pub fn factory(cfg: CfgEnv, #[with(true)] flags: ExecutionFlags) -> BlockifierFactory { - BlockifierFactory::new(cfg, flags) + BlockifierFactory::new(cfg, flags, BlockLimits::max()) } } diff --git a/crates/katana/node/src/config/mod.rs b/crates/katana/node/src/config/mod.rs index e2d32e51c7..0db09ddaff 100644 --- a/crates/katana/node/src/config/mod.rs +++ b/crates/katana/node/src/config/mod.rs @@ -6,6 +6,7 @@ pub mod execution; pub mod fork; pub mod metrics; pub mod rpc; +pub mod sequencing; use db::DbConfig; use dev::DevConfig; @@ -15,6 +16,7 @@ use katana_chain_spec::ChainSpec; use katana_core::service::messaging::MessagingConfig; use metrics::MetricsConfig; use rpc::RpcConfig; +use sequencing::SequencingConfig; /// Node configurations. /// @@ -48,15 +50,3 @@ pub struct Config { /// Development options. pub dev: DevConfig, } - -/// Configurations related to block production. -#[derive(Debug, Clone, Default)] -pub struct SequencingConfig { - /// The time in milliseconds for a block to be produced. - pub block_time: Option, - - /// Disable automatic block production. - /// - /// Allowing block to only be produced manually. - pub no_mining: bool, -} diff --git a/crates/katana/node/src/config/sequencing.rs b/crates/katana/node/src/config/sequencing.rs new file mode 100644 index 0000000000..f8b3aa2426 --- /dev/null +++ b/crates/katana/node/src/config/sequencing.rs @@ -0,0 +1,29 @@ +use katana_executor::BlockLimits; + +/// Configurations related to block production. +#[derive(Debug, Clone, Default)] +pub struct SequencingConfig { + /// The time in milliseconds for a block to be produced. + pub block_time: Option, + + /// Disable automatic block production. + /// + /// Allowing block to only be produced manually. + pub no_mining: bool, + + /// The maximum number of Cairo steps in a block. + // + /// The block will automatically be closed when the accumulated Cairo steps across all the + /// transactions has reached this limit. + /// + /// NOTE: This only affect interval block production. + /// + /// See . + pub block_cairo_steps_limit: Option, +} + +impl SequencingConfig { + pub fn block_limits(&self) -> BlockLimits { + BlockLimits { cairo_steps: self.block_cairo_steps_limit.unwrap_or(u64::MAX) } + } +} diff --git a/crates/katana/node/src/lib.rs b/crates/katana/node/src/lib.rs index f56c086480..83772a5ff3 100644 --- a/crates/katana/node/src/lib.rs +++ b/crates/katana/node/src/lib.rs @@ -182,7 +182,11 @@ pub async fn build(mut config: Config) -> Result { .with_account_validation(config.dev.account_validation) .with_fee(config.dev.fee); - let executor_factory = Arc::new(BlockifierFactory::new(cfg_env, execution_flags)); + let executor_factory = Arc::new(BlockifierFactory::new( + cfg_env, + execution_flags, + config.sequencing.block_limits(), + )); // --- build backend diff --git a/crates/katana/rpc/rpc/tests/dev.rs b/crates/katana/rpc/rpc/tests/dev.rs index 4708e19326..da80ea47eb 100644 --- a/crates/katana/rpc/rpc/tests/dev.rs +++ b/crates/katana/rpc/rpc/tests/dev.rs @@ -1,5 +1,5 @@ use dojo_test_utils::sequencer::{get_default_test_config, TestSequencer}; -use katana_node::config::SequencingConfig; +use katana_node::config::sequencing::SequencingConfig; use katana_provider::traits::block::{BlockNumberProvider, BlockProvider}; use katana_provider::traits::env::BlockEnvProvider; use katana_rpc_api::dev::DevApiClient; diff --git a/crates/katana/rpc/rpc/tests/forking.rs b/crates/katana/rpc/rpc/tests/forking.rs index fd2abf34be..eb98cb231f 100644 --- a/crates/katana/rpc/rpc/tests/forking.rs +++ b/crates/katana/rpc/rpc/tests/forking.rs @@ -3,7 +3,7 @@ use assert_matches::assert_matches; use cainome::rs::abigen_legacy; use dojo_test_utils::sequencer::{get_default_test_config, TestSequencer}; use katana_node::config::fork::ForkingConfig; -use katana_node::config::SequencingConfig; +use katana_node::config::sequencing::SequencingConfig; use katana_primitives::block::{BlockHash, BlockHashOrNumber, BlockIdOrTag, BlockNumber, BlockTag}; use katana_primitives::chain::NamedChainId; use katana_primitives::event::MaybeForkedContinuationToken; diff --git a/crates/katana/rpc/rpc/tests/messaging.rs b/crates/katana/rpc/rpc/tests/messaging.rs index ac11f7e8a3..6360cfd1c9 100644 --- a/crates/katana/rpc/rpc/tests/messaging.rs +++ b/crates/katana/rpc/rpc/tests/messaging.rs @@ -11,7 +11,7 @@ use cainome::rs::abigen; use dojo_test_utils::sequencer::{get_default_test_config, TestSequencer}; use dojo_utils::TransactionWaiter; use katana_core::service::messaging::MessagingConfig; -use katana_node::config::SequencingConfig; +use katana_node::config::sequencing::SequencingConfig; use katana_primitives::felt; use katana_primitives::utils::transaction::{ compute_l1_handler_tx_hash, compute_l1_to_l2_message_hash, compute_l2_to_l1_message_hash, diff --git a/crates/katana/rpc/rpc/tests/proofs.rs b/crates/katana/rpc/rpc/tests/proofs.rs index 3f2f75d261..a29501edd5 100644 --- a/crates/katana/rpc/rpc/tests/proofs.rs +++ b/crates/katana/rpc/rpc/tests/proofs.rs @@ -5,7 +5,7 @@ use dojo_test_utils::sequencer::{get_default_test_config, TestSequencer}; use jsonrpsee::http_client::HttpClientBuilder; use katana_chain_spec::ChainSpec; use katana_node::config::rpc::DEFAULT_RPC_MAX_PROOF_KEYS; -use katana_node::config::SequencingConfig; +use katana_node::config::sequencing::SequencingConfig; use katana_primitives::block::BlockIdOrTag; use katana_primitives::class::{ClassHash, CompiledClassHash}; use katana_primitives::contract::{StorageKey, StorageValue}; diff --git a/crates/katana/rpc/rpc/tests/saya.rs b/crates/katana/rpc/rpc/tests/saya.rs index 323e26439a..7dafa95c70 100644 --- a/crates/katana/rpc/rpc/tests/saya.rs +++ b/crates/katana/rpc/rpc/tests/saya.rs @@ -6,7 +6,7 @@ use std::sync::Arc; use dojo_test_utils::sequencer::{get_default_test_config, TestSequencer}; use dojo_utils::TransactionWaiter; use jsonrpsee::http_client::HttpClientBuilder; -use katana_node::config::SequencingConfig; +use katana_node::config::sequencing::SequencingConfig; use katana_primitives::block::{BlockIdOrTag, BlockTag}; use katana_rpc_api::dev::DevApiClient; use katana_rpc_api::saya::SayaApiClient; diff --git a/crates/katana/rpc/rpc/tests/starknet.rs b/crates/katana/rpc/rpc/tests/starknet.rs index 8e8064f27b..d0b5629f32 100644 --- a/crates/katana/rpc/rpc/tests/starknet.rs +++ b/crates/katana/rpc/rpc/tests/starknet.rs @@ -9,7 +9,7 @@ use common::split_felt; use dojo_test_utils::sequencer::{get_default_test_config, TestSequencer}; use indexmap::IndexSet; use jsonrpsee::http_client::HttpClientBuilder; -use katana_node::config::SequencingConfig; +use katana_node::config::sequencing::SequencingConfig; use katana_primitives::event::ContinuationToken; use katana_primitives::genesis::constant::{ DEFAULT_ACCOUNT_CLASS_HASH, DEFAULT_ETH_FEE_TOKEN_ADDRESS, DEFAULT_PREFUNDED_ACCOUNT_BALANCE, diff --git a/crates/katana/rpc/rpc/tests/torii.rs b/crates/katana/rpc/rpc/tests/torii.rs index 18ffa0e855..0644b5ca3d 100644 --- a/crates/katana/rpc/rpc/tests/torii.rs +++ b/crates/katana/rpc/rpc/tests/torii.rs @@ -7,7 +7,7 @@ use std::time::Duration; use dojo_test_utils::sequencer::{get_default_test_config, TestSequencer}; use dojo_utils::TransactionWaiter; use jsonrpsee::http_client::HttpClientBuilder; -use katana_node::config::SequencingConfig; +use katana_node::config::sequencing::SequencingConfig; use katana_rpc_api::dev::DevApiClient; use katana_rpc_api::starknet::StarknetApiClient; use katana_rpc_api::torii::ToriiApiClient;