diff --git a/config/stratus.example.toml b/config/stratus.example.toml index d3f57525e..27b185a89 100644 --- a/config/stratus.example.toml +++ b/config/stratus.example.toml @@ -65,12 +65,16 @@ [executor] # Chain ID of the network. Required. chain_id = 2008 -# Number of EVMs to execute calls against the present state. -# call_present_evms = 50 -# Number of EVMs to execute calls against past states. -# call_past_evms = 50 -# Number of EVMs to execute inspector calls. -# inspector_evms = 50 +# Total number of EVM workers in the unified pool, shared by every execution kind. +# evm_workers = 150 +# Maximum number of concurrent calls against the present state. Defaults to the remaining pool capacity. +# call_present_limit = 50 +# Maximum number of concurrent calls against past states. +# call_past_limit = 50 +# Maximum number of concurrent inspector calls. +# inspector_limit = 50 +# Pool busy percentage above which per-kind limits are enforced; below it, tasks of any kind are admitted freely. +# evm_busy_threshold = 80 # Should reject contract transactions and calls to accounts that are not contracts? # reject_not_contract = true # EVM hardfork spec, e.g. Prague, Cancun, Shanghai. diff --git a/crates/stratus_metrics/src/definitions.rs b/crates/stratus_metrics/src/definitions.rs index 9c04a010b..2ea7ae036 100644 --- a/crates/stratus_metrics/src/definitions.rs +++ b/crates/stratus_metrics/src/definitions.rs @@ -107,7 +107,19 @@ metrics! { histogram_duration executor_inspect{trace_type}, "Number of EVM pool workers busy executing right now." - gauge executor_workers_busy{pool} + gauge executor_workers_busy{pool}, + + "Total number of EVM workers in the unified executor pool." + gauge executor_workers_total{}, + + "Number of tasks currently executing in the unified executor pool, by kind." + gauge executor_pool_inflight{kind}, + + "Number of tasks waiting to be admitted to the unified executor pool, by kind." + gauge executor_pool_waiting{kind}, + + "Number of tasks admitted through the relaxed path while the pool was below the busy threshold, by kind." + counter executor_pool_relaxed_admissions{kind} }, group: rocks { diff --git a/src/config/loader.rs b/src/config/loader.rs index dacf4e79c..f325c1435 100644 --- a/src/config/loader.rs +++ b/src/config/loader.rs @@ -272,6 +272,7 @@ mod tests { use std::ffi::OsString; use clap::CommandFactory; + use clap::Parser; use crate::config::Environment; use crate::config::StratusConfig; @@ -375,9 +376,11 @@ mod tests { [executor] chain_id = 100 - call_present_evms = 1 - call_past_evms = 2 - inspector_evms = 3 + evm_workers = 10 + call_present_limit = 1 + call_past_limit = 4 + inspector_limit = 5 + evm_busy_threshold = 75 reject_not_contract = false evm_spec = "Cancun" @@ -572,9 +575,11 @@ mod tests { [executor] chain_id = 100 - call_present_evms = 1 - call_past_evms = 2 - inspector_evms = 3 + evm_workers = 10 + call_present_limit = 1 + call_past_limit = 4 + inspector_limit = 5 + evm_busy_threshold = 75 reject_not_contract = false evm_spec = "Cancun" @@ -699,14 +704,7 @@ mod tests { fn test_empty_config_uses_defaults() { // an empty file plus the minimal valid arguments must yield exactly the defaults let config = load_with(&["--leader", "--executor-chain-id", "1"], "").unwrap(); - let default = StratusConfig { - leader: true, - executor: crate::eth::executor::ExecutorConfig { - executor_chain_id: 1, - ..Default::default() - }, - ..Default::default() - }; + let default = StratusConfig::try_parse_from(["stratus", "--leader", "--executor-chain-id", "1"]).unwrap(); assert_eq!(serde_json::to_value(&config).unwrap(), serde_json::to_value(&default).unwrap()); } } diff --git a/src/config/mod.rs b/src/config/mod.rs index d4e3a26bf..27c206b63 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -173,7 +173,6 @@ impl CommonConfig { /// Configuration for main Stratus service. #[derive(DebugAsJson, Clone, Parser, derive_more::Deref, serde::Serialize)] -#[cfg_attr(test, derive(Default))] #[clap(group = ArgGroup::new("mode").args(&["leader", "follower", "fake_leader"]).required(true))] pub struct StratusConfig { #[arg(id = "leader", long = "leader", conflicts_with_all = ["follower", "fake_leader"])] diff --git a/src/config/validate/mod.rs b/src/config/validate/mod.rs index b3ae1e987..bf6ee0f09 100644 --- a/src/config/validate/mod.rs +++ b/src/config/validate/mod.rs @@ -27,6 +27,9 @@ enum Step { /// Shares the block-mode conflict check of `MinerConfig::init` MinerExternalBlockMode, + /// Shares the suboptimal-configuration warnings of `PoolConfig::validate` + ExecutorPoolLimits, + /// A genesis file that does not exist is silently replaced by the default genesis block at runtime. #[cfg(feature = "dev")] GenesisFileExists, @@ -36,6 +39,7 @@ impl Step { const WORKFLOW: &'static [Self] = &[ Self::KafkaCompleteness, Self::MinerExternalBlockMode, + Self::ExecutorPoolLimits, #[cfg(feature = "dev")] Self::GenesisFileExists, ]; @@ -45,6 +49,7 @@ impl Step { match self { Self::KafkaCompleteness => Severity::Error, Self::MinerExternalBlockMode => Severity::Warning, + Self::ExecutorPoolLimits => Severity::Warning, #[cfg(feature = "dev")] Self::GenesisFileExists => Severity::Warning, } @@ -55,6 +60,7 @@ impl Step { match self { Self::KafkaCompleteness => Self::check_kafka_completeness(config), Self::MinerExternalBlockMode => Self::check_miner_external_block_mode(config), + Self::ExecutorPoolLimits => Self::check_executor_pool_limits(config), #[cfg(feature = "dev")] Self::GenesisFileExists => Self::check_genesis_file_exists(config), } @@ -79,6 +85,11 @@ impl Step { Vec::new() } + /// Returns warnings about suboptimal executor pool configurations. + fn check_executor_pool_limits(config: &StratusConfig) -> Vec { + config.executor.pool.validate() + } + /// Returns a warning when the configured genesis file does not exist. #[cfg(feature = "dev")] fn check_genesis_file_exists(config: &StratusConfig) -> Vec { diff --git a/src/eth/executor/config.rs b/src/eth/executor/config.rs index 40e9542a5..a86b1ad60 100644 --- a/src/eth/executor/config.rs +++ b/src/eth/executor/config.rs @@ -16,14 +16,9 @@ pub struct ExecutorConfig { #[serde(rename = "chain_id")] pub executor_chain_id: u64, - #[arg(id = "executor.call_present_evms", long = "executor-call-present-evms", default_value_t = 50)] - pub call_present_evms: usize, - - #[arg(id = "executor.call_past_evms", long = "executor-call-past-evms", default_value_t = 50)] - pub call_past_evms: usize, - - #[arg(id = "executor.inspector_evms", long = "executor-inspector-evms", default_value_t = 50)] - pub inspector_evms: usize, + #[command(flatten)] + #[serde(flatten)] + pub pool: PoolConfig, /// Should reject contract transactions and calls to accounts that are not contracts? #[arg( @@ -43,17 +38,62 @@ pub struct ExecutorConfig { pub executor_evm_spec: SpecId, } -#[cfg(test)] -impl Default for ExecutorConfig { - fn default() -> Self { - Self { - executor_chain_id: 0, - call_present_evms: 50, - call_past_evms: 50, - inspector_evms: 50, - executor_reject_not_contract: true, - executor_evm_spec: SpecId::PRAGUE, +/// Configuration of the unified EVM worker pool: one shared set of workers serving every execution kind. +#[derive(Parser, DebugAsJson, Clone, Copy, serde::Serialize)] +pub struct PoolConfig { + /// Total number of EVM workers in the unified pool, shared by every execution kind. + #[arg(id = "executor.evm_workers", long = "executor-evm-workers", default_value_t = 150, value_parser = clap::builder::RangedU64ValueParser::::new().range(1..))] + pub evm_workers: usize, + + /// Maximum number of concurrent call-present executions. + /// Defaults to the remaining pool capacity (`evm_workers` minus the other limits). + #[arg(id = "executor.call_present_limit", long = "executor-call-present-limit", value_parser = clap::builder::RangedU64ValueParser::::new().range(1..))] + #[serde(skip_serializing_if = "Option::is_none")] + pub call_present_limit: Option, + + /// Maximum number of concurrent call-past executions. + #[arg(id = "executor.call_past_limit", long = "executor-call-past-limit", default_value_t = 50, value_parser = clap::builder::RangedU64ValueParser::::new().range(1..))] + pub call_past_limit: usize, + + /// Maximum number of concurrent inspector executions. + #[arg(id = "executor.inspector_limit", long = "executor-inspector-limit", default_value_t = 50, value_parser = clap::builder::RangedU64ValueParser::::new().range(1..))] + pub inspector_limit: usize, + + /// Pool busy percentage above which per-kind limits are enforced: while the pool is below this + /// threshold, tasks are admitted even above their kind's limit. + #[arg(id = "executor.evm_busy_threshold", long = "executor-evm-busy-threshold", default_value_t = 80, value_parser = clap::builder::RangedU64ValueParser::::new().range(0..=100))] + pub evm_busy_threshold: usize, +} + +impl PoolConfig { + /// Effective call-present limit: the configured value, or the remaining pool capacity by default. + pub fn call_present_limit(&self) -> usize { + self.call_present_limit + .unwrap_or_else(|| self.evm_workers.saturating_sub(self.call_past_limit + self.inspector_limit)) + } + + /// In-flight task count at which relaxed admission ends and per-kind limits are enforced. + pub fn relaxed_limit(&self) -> usize { + self.evm_workers * self.evm_busy_threshold / 100 + } + + /// Returns warnings about suboptimal configurations the pool can still run with. + pub fn validate(&self) -> Vec { + let mut warnings = Vec::new(); + + let limits_sum = self.call_present_limit() + self.call_past_limit + self.inspector_limit; + if limits_sum > self.evm_workers { + warnings.push(format!( + "executor pool kind limits ({}) exceed the total number of workers ({}); saturating every kind makes tasks queue instead of execute", + limits_sum, self.evm_workers + )); + } + + if self.call_present_limit.is_none() && self.call_present_limit() == 0 { + warnings.push("call-present limit defaults to zero; call-present tasks are only admitted while the pool is below the busy threshold".to_string()); } + + warnings } } diff --git a/src/eth/executor/evm/mod.rs b/src/eth/executor/evm/mod.rs index d9958097c..d0c238d2c 100644 --- a/src/eth/executor/evm/mod.rs +++ b/src/eth/executor/evm/mod.rs @@ -2,7 +2,6 @@ mod session; pub mod types; mod util; -use std::marker::PhantomData; use std::sync::Arc; use alloy_consensus::transaction::TransactionInfo; @@ -48,13 +47,12 @@ use crate::eth::types::StratusError; pub type RevmResultAndState = ExecResultAndState; /// Implementation of EVM using [`revm`](https://crates.io/crates/revm). -pub struct Evm { +pub struct Evm { evm: GeneralRevm, kind: EvmKind, - _input_type: PhantomData, } -impl Evm { +impl Evm { /// Creates a new instance of the Evm. pub fn new(storage: Arc, config: &ExecutorConfig, kind: EvmKind) -> Self { tracing::info!(?config, "creating revm"); @@ -65,12 +63,11 @@ impl Evm { Self { evm: create_evm(chain_id, config.executor_evm_spec, RevmSession::new(storage), kind), kind, - _input_type: PhantomData, } } /// Execute a transaction that deploys a contract or call a contract function. - pub fn execute(&mut self, input: Input) -> Result<(RevmResultAndState, ExecutionMetrics), StratusError> { + pub fn execute(&mut self, input: Input) -> Result<(RevmResultAndState, ExecutionMetrics), StratusError> { let metrics_context = input.metrics_context(); // configure session @@ -94,7 +91,7 @@ impl Evm { } } -impl Evm { +impl Evm { /// Execute a transaction using a tracer. pub fn inspect(&mut self, input: InspectorInput) -> Result { let InspectorInput { diff --git a/src/eth/executor/evm_worker_pool.rs b/src/eth/executor/evm_worker_pool.rs index 85ab2611d..453fd74ea 100644 --- a/src/eth/executor/evm_worker_pool.rs +++ b/src/eth/executor/evm_worker_pool.rs @@ -12,109 +12,61 @@ use crate::eth::executor::evm::EvmKind; use crate::eth::executor::evm::RevmResultAndState; use crate::eth::executor::evm::types::CallExecutionInput; use crate::eth::executor::evm::types::InspectorInput; -use crate::eth::executor::types::EvmRoute; -use crate::eth::executor::types::EvmTask; +use crate::eth::executor::pool_admission::PoolAdmission; use crate::eth::executor::types::ExecutionTask; use crate::eth::executor::types::InspectionTask; -use crate::eth::executor::types::Task; +use crate::eth::executor::types::PoolTask; use crate::eth::storage::StratusStorage; use crate::eth::types::StratusError; use crate::eth::types::UnexpectedError; use crate::ext::spawn_thread; use crate::infra::tracing::warn_task_tx_closed; -/// Manages EVM pool and communication channels. -pub struct EvmWorkerPool { - /// Pool for parallel execution of calls (eth_call and eth_estimateGas) reading from current state. Usually contains multiple EVMs. - pub call_present: crossbeam_channel::Sender>>, - - /// Pool for parallel execution of calls (eth_call and eth_estimateGas) reading from past state. Usually contains multiple EVMs. - pub call_past: crossbeam_channel::Sender>>, +/// Total capacity of the unified EVM pool task queue. +const TASK_QUEUE_CAPACITY: usize = 4096; - /// Pool for parallel execution of tx inspections (debug_traceTransaction). Usually contains multiple EVMs. - pub inspector: crossbeam_channel::Sender>, +/// Manages the unified EVM pool: one shared set of workers serving every execution kind. +pub struct EvmWorkerPool { + tx: crossbeam_channel::Sender, + admission: Arc, } impl EvmWorkerPool { - /// Spawns EVM tasks in background. + /// Spawns the unified EVM pool workers. pub fn spawn(storage: Arc, config: &ExecutorConfig) -> Self { - // function executed by evm threads - fn worker( - task_name: &str, - storage: Arc, - config: ExecutorConfig, - task_rx: crossbeam_channel::Receiver>, - kind: EvmKind, - ) { - let mut evm = Evm::new(Arc::clone(&storage), &config, kind); - - // keep executing transactions until the channel is closed - while let Ok(task) = task_rx.recv() { - if GlobalState::is_shutdown_warn(task_name) { - return; - } - - let _guard = kind.mark_executor_pool_busy(); - if let Err(StratusError::Executor(ExecutorError::Panic { err: panic_err })) = task.execute(&mut evm) { - tracing::error!(?panic_err, "executor panicked; recreating EVM"); - evm = Evm::new(Arc::clone(&storage), &config, kind); - } - } - warn_task_tx_closed(task_name); + let pool = config.pool; + let (tx, rx) = crossbeam_channel::bounded::(TASK_QUEUE_CAPACITY); + let admission = Arc::new(PoolAdmission::new(pool)); + + for worker_index in 1..=pool.evm_workers { + let task_name = format!("evm-pool-{worker_index}"); + let worker_storage = Arc::clone(&storage); + let worker_config = *config; + let worker_rx = rx.clone(); + let thread_name = task_name.clone(); + spawn_thread(&thread_name, move || { + Self::worker(&task_name, worker_storage, worker_config, worker_rx); + }); } - // function that spawn evm threads - fn spawn_evms( - task_name: &str, - num_evms: usize, - kind: EvmKind, - storage: &Arc, - config: &ExecutorConfig, - ) -> crossbeam_channel::Sender> { - let (evm_tx, evm_rx) = crossbeam_channel::bounded::>(4096); - - for evm_index in 1..=num_evms { - let evm_task_name = format!("{task_name}-{evm_index}"); - let evm_storage = Arc::clone(storage); - let evm_config = *config; - let evm_rx = evm_rx.clone(); - let thread_name = evm_task_name.clone(); - spawn_thread(&thread_name, move || { - worker(&evm_task_name, evm_storage, evm_config, evm_rx, kind); - }); - } + // initialize the gauges so every kind series exists from startup + for kind in [EvmKind::CallPresent, EvmKind::CallPast, EvmKind::Inspect] { metrics::set_executor_workers_busy(0, kind); - evm_tx } + metrics::set_executor_workers_total(pool.evm_workers as u64); - let call_present = spawn_evms("evm-call-present", config.call_present_evms, EvmKind::CallPresent, &storage, config); - let call_past = spawn_evms("evm-call-past", config.call_past_evms, EvmKind::CallPast, &storage, config); - let inspector = spawn_evms("inspector", config.inspector_evms, EvmKind::Inspect, &storage, config); - - EvmWorkerPool { - call_present, - call_past, - inspector, - } + Self { tx, admission } } - /// Executes a transaction in the specified route. - pub fn execute(&self, route: EvmRoute) -> Result<(Output, ExecutionMetrics), StratusError> + /// Executes a call in the specified route. + pub fn execute(&self, input: CallExecutionInput) -> Result<(Output, ExecutionMetrics), StratusError> where Output: TryFrom, { let (execution_tx, execution_rx) = oneshot::channel::>(); - match route { - EvmRoute::CallPresent(input) => { - let task = ExecutionTask::new(input, execution_tx).into(); - self.call_present.send(task)?; - } - EvmRoute::CallPast(input) => { - let task = ExecutionTask::new(input, execution_tx).into(); - self.call_past.send(task)?; - } - }; + let task = PoolTask::call(ExecutionTask::new(input, execution_tx), &self.admission)?; + self.tx.send(task)?; match execution_rx.recv() { Ok(result) => { @@ -125,13 +77,115 @@ impl EvmWorkerPool { } } + /// Executes a transaction inspection (debug_traceTransaction). pub fn inspect(&self, input: InspectorInput) -> Result { let (inspector_tx, inspector_rx) = oneshot::channel::>(); - let task = InspectionTask::new(input, inspector_tx).into(); - let _ = self.inspector.send(task); + let task = PoolTask::inspect(InspectionTask::new(input, inspector_tx), &self.admission)?; + let _ = self.tx.send(task); match inspector_rx.recv() { Ok(result) => result, Err(_) => Err(UnexpectedError::ChannelClosed { channel: "evm" }.into()), } } + + /// Function executed by the unified EVM pool worker threads. + fn worker(task_name: &str, storage: Arc, config: ExecutorConfig, task_rx: crossbeam_channel::Receiver) { + let mut evm = Evm::new(Arc::clone(&storage), &config, EvmKind::CallPresent); + + while let Ok(task) = task_rx.recv() { + if GlobalState::is_shutdown_warn(task_name) { + return; + } + + if let Err(StratusError::Executor(ExecutorError::Panic { err: panic_err })) = task.execute(&mut evm) { + tracing::error!(?panic_err, "executor panicked; recreating EVM"); + evm = Evm::new(Arc::clone(&storage), &config, EvmKind::CallPresent); + } + } + warn_task_tx_closed(task_name); + } +} + +// ----------------------------------------------------------------------------- +// Tests +// ----------------------------------------------------------------------------- + +#[cfg(test)] +mod tests { + use clap::Parser; + + use super::*; + + fn test_config(args: &[&str]) -> ExecutorConfig { + let mut arguments = vec!["stratus", "--executor-chain-id", "1"]; + arguments.extend_from_slice(args); + ExecutorConfig::parse_from(arguments) + } + + #[test] + fn test_pool_config_resolves_defaults() { + let pool = test_config(&[]).pool; + assert_eq!(pool.evm_workers, 150); + assert_eq!(pool.call_present_limit(), 50); + assert_eq!(pool.call_past_limit, 50); + assert_eq!(pool.inspector_limit, 50); + assert_eq!(pool.evm_busy_threshold, 80); + assert_eq!(pool.relaxed_limit(), 120); + assert!(pool.validate().is_empty()); + } + + #[test] + fn test_pool_config_explicit_call_present_limit() { + let pool = test_config(&[ + "--executor-evm-workers", + "200", + "--executor-call-present-limit", + "120", + "--executor-call-past-limit", + "20", + "--executor-inspector-limit", + "30", + ]) + .pool; + assert_eq!(pool.evm_workers, 200); + assert_eq!(pool.call_present_limit(), 120); + assert_eq!(pool.call_past_limit, 20); + assert_eq!(pool.inspector_limit, 30); + } + + #[test] + fn test_pool_config_call_present_uses_remaining_capacity() { + let pool = test_config(&["--executor-evm-workers", "200"]).pool; + assert_eq!(pool.call_present_limit(), 200 - 50 - 50); + } + + #[test] + fn test_pool_config_warns_on_limits_exceeding_workers() { + let pool = test_config(&["--executor-evm-workers", "100", "--executor-call-present-limit", "60"]).pool; + assert!(!pool.validate().is_empty()); + } + + #[test] + fn test_pool_config_rejects_zero_workers() { + let result = ExecutorConfig::try_parse_from(["stratus", "--executor-chain-id", "1", "--executor-evm-workers", "0"]); + assert!(result.is_err(), "must reject zero workers"); + } + + #[test] + fn test_pool_config_rejects_zero_limit() { + let result = ExecutorConfig::try_parse_from(["stratus", "--executor-chain-id", "1", "--executor-call-past-limit", "0"]); + assert!(result.is_err(), "must reject zero limit"); + } + + #[test] + fn test_pool_config_rejects_threshold_above_100() { + let result = ExecutorConfig::try_parse_from(["stratus", "--executor-chain-id", "1", "--executor-evm-busy-threshold", "101"]); + assert!(result.is_err(), "must reject threshold above 100"); + } + + #[test] + fn test_pool_config_zero_threshold_is_strict() { + let pool = test_config(&["--executor-evm-busy-threshold", "0"]).pool; + assert_eq!(pool.relaxed_limit(), 0); + } } diff --git a/src/eth/executor/mod.rs b/src/eth/executor/mod.rs index 1aeae0268..12da4d7e0 100644 --- a/src/eth/executor/mod.rs +++ b/src/eth/executor/mod.rs @@ -1,6 +1,7 @@ mod config; mod evm; mod evm_worker_pool; +mod pool_admission; mod transaction_worker; pub mod types; @@ -39,7 +40,6 @@ use crate::eth::executor::evm::types::CallExecutionInput; use crate::eth::executor::evm::types::InspectorInput; use crate::eth::executor::evm_worker_pool::EvmWorkerPool; use crate::eth::executor::transaction_worker::TransactionWorker; -use crate::eth::executor::types::EvmRoute; use crate::eth::miner::Miner; use crate::eth::storage::ExecutionKind; use crate::eth::storage::StorageError; @@ -52,7 +52,6 @@ use crate::eth::types::ExternalReceipt; use crate::eth::types::ExternalReceipts; use crate::eth::types::ExternalTransaction; use crate::eth::types::Hash; -use crate::eth::types::PointInTime; use crate::eth::types::StratusError; use crate::eth::types::TransactionInput; #[cfg(feature = "metrics")] @@ -145,7 +144,7 @@ impl Executor { fn execute_external_transaction_inner( storage: &StratusStorage, miner: &Miner, - evm: &mut Evm, + evm: &mut Evm, tx: ExternalTransaction, receipt: ExternalReceipt, block_number: BlockNumber, @@ -293,12 +292,7 @@ impl Executor { let evm_input = CallExecutionInput::create(call_input, block_info, kind); - let evm_route = match kind.point_in_time() { - PointInTime::Pending | PointInTime::Latest => EvmRoute::CallPresent(evm_input), - PointInTime::Past(_) => EvmRoute::CallPast(evm_input), - }; - - self.evms.execute::(evm_route).map(|(output, _metrics)| output) + self.evms.execute::(evm_input).map(|(output, _metrics)| output) } #[timed(executor_inspect, labels( diff --git a/src/eth/executor/pool_admission.rs b/src/eth/executor/pool_admission.rs new file mode 100644 index 000000000..72639e670 --- /dev/null +++ b/src/eth/executor/pool_admission.rs @@ -0,0 +1,242 @@ +use std::sync::Arc; +use std::sync::atomic::AtomicUsize; +use std::sync::atomic::Ordering; +use std::time::Duration; + +use parking_lot::Condvar; +use parking_lot::Mutex; +use stratus_metrics as metrics; + +use crate::GlobalState; +use crate::eth::executor::config::PoolConfig; +use crate::eth::executor::evm::EvmKind; +use crate::eth::types::StateError; + +/// Interval between shutdown checks while blocked waiting for a kind slot. +const SHUTDOWN_POLL_INTERVAL: Duration = Duration::from_millis(250); + +/// Per-kind admission state, counting in-flight tasks admitted by both the relaxed and throttled paths. +struct KindGate { + kind: EvmKind, + limit: usize, + inflight: Mutex, + cvar: Condvar, +} + +impl KindGate { + fn new(kind: EvmKind, limit: usize) -> Self { + // initialize the gauge so the kind series exists from startup + metrics::set_executor_pool_inflight(0, kind); + Self { + kind, + limit, + inflight: Mutex::new(0), + cvar: Condvar::new(), + } + } + + /// Relaxed admission: increments the in-flight count without checking the limit. + fn admit_relaxed(&self) { + *self.inflight.lock() += 1; + metrics::inc_executor_pool_inflight(1, self.kind); + } + + /// Throttled admission: blocks until the kind's in-flight count drops below its limit. + /// Returns an error when the application starts shutting down. + fn admit_throttled(&self) -> Result<(), StateError> { + metrics::inc_executor_pool_waiting(1, self.kind); + + let mut inflight = self.inflight.lock(); + while *inflight >= self.limit { + if GlobalState::is_shutdown() { + drop(inflight); + metrics::dec_executor_pool_waiting(1, self.kind); + return Err(StateError::StratusShutdown); + } + self.cvar.wait_for(&mut inflight, SHUTDOWN_POLL_INTERVAL); + } + + *inflight += 1; + drop(inflight); + metrics::dec_executor_pool_waiting(1, self.kind); + metrics::inc_executor_pool_inflight(1, self.kind); + Ok(()) + } + + /// Releases an admission slot of the kind. + fn release(&self) { + *self.inflight.lock() -= 1; + self.cvar.notify_one(); + metrics::dec_executor_pool_inflight(1, self.kind); + } +} + +/// Admission control of the unified EVM pool: tasks of any kind are admitted immediately while the +/// pool is below the busy threshold, and limited per kind above it. +/// +/// Both paths are counted by the same per-kind in-flight counters, so when the busy threshold is +/// crossed and throttling starts, executions that were admitted without a limit check are already +/// accounted for. +pub struct PoolAdmission { + call_present: Arc, + call_past: Arc, + inspector: Arc, + inflight_total: AtomicUsize, + relaxed_limit: usize, +} + +impl PoolAdmission { + pub fn new(config: PoolConfig) -> Self { + Self { + call_present: Arc::new(KindGate::new(EvmKind::CallPresent, config.call_present_limit())), + call_past: Arc::new(KindGate::new(EvmKind::CallPast, config.call_past_limit)), + inspector: Arc::new(KindGate::new(EvmKind::Inspect, config.inspector_limit)), + inflight_total: AtomicUsize::new(0), + relaxed_limit: config.relaxed_limit(), + } + } + + fn gate(&self, kind: EvmKind) -> &Arc { + match kind { + EvmKind::CallPresent => &self.call_present, + EvmKind::CallPast => &self.call_past, + EvmKind::Inspect => &self.inspector, + EvmKind::Transaction => unreachable!("transaction execution is not managed by the unified EVM pool"), + } + } + + /// Admits a task of the kind: immediately while the pool is below the busy threshold (even above + /// the kind's limit), blocking on the kind's limit otherwise. Returns an error on shutdown. + pub fn acquire(self: &Arc, kind: EvmKind) -> Result { + if self.inflight_total.load(Ordering::Relaxed) < self.relaxed_limit { + metrics::inc_executor_pool_relaxed_admissions(kind); + self.gate(kind).admit_relaxed(); + } else { + self.gate(kind).admit_throttled()?; + } + self.inflight_total.fetch_add(1, Ordering::Relaxed); + Ok(PoolPermit { + admission: Arc::clone(self), + kind, + }) + } + + /// Releases a slot taken by [`PoolAdmission::acquire`]. + fn release(&self, kind: EvmKind) { + self.gate(kind).release(); + self.inflight_total.fetch_sub(1, Ordering::Relaxed); + } +} + +/// Admission slot of a task in the unified EVM pool, released on drop. +pub struct PoolPermit { + admission: Arc, + kind: EvmKind, +} + +impl PoolPermit { + /// Kind of the task holding the permit. + pub(crate) fn evm_kind(&self) -> EvmKind { + self.kind + } +} + +impl Drop for PoolPermit { + fn drop(&mut self) { + self.admission.release(self.kind); + } +} + +#[cfg(test)] +mod tests { + use std::sync::mpsc; + + use super::*; + + fn admission(workers: usize, busy_threshold: usize, call_past_limit: usize) -> Arc { + let config = PoolConfig { + evm_workers: workers, + call_present_limit: None, + call_past_limit, + inspector_limit: workers, + evm_busy_threshold: busy_threshold, + }; + Arc::new(PoolAdmission::new(config)) + } + + /// Spawns a thread acquiring a permit of the kind; returns a channel the thread signals on completion. + fn spawn_acquire(admission: &Arc, kind: EvmKind) -> mpsc::Receiver<()> { + let (tx, rx) = mpsc::channel(); + let admission = Arc::clone(admission); + std::thread::spawn(move || { + let _permit = admission.acquire(kind); + let _ = tx.send(()); + }); + rx + } + + #[test] + fn test_relaxed_admission_above_kind_limit() { + // busy threshold of 80% of 10 workers: any kind is admitted while total in-flight < 8 + let admission = admission(10, 80, 1); + let permits: Vec<_> = (0..5).map(|_| admission.acquire(EvmKind::CallPast).unwrap()).collect(); + assert_eq!(permits.len(), 5); + } + + #[test] + fn test_throttled_admission_blocks_at_kind_limit() { + // threshold 80% of 1 worker: relaxed limit 0, so every admission checks the kind limit + let admission = admission(1, 80, 1); + + let first = admission.acquire(EvmKind::CallPast).unwrap(); + + let rx = spawn_acquire(&admission, EvmKind::CallPast); + assert!(rx.recv_timeout(Duration::from_millis(100)).is_err(), "must block while the limit is held"); + + drop(first); + assert!(rx.recv_timeout(Duration::from_secs(5)).is_ok(), "must proceed after the permit is released"); + } + + #[test] + fn test_relaxed_admissions_count_toward_kind_limit() { + // bypass executions must already be accounted for when throttling starts + let admission = admission(10, 80, 2); + + // 3 call-past tasks bypass the limit of 2 while total in-flight < 8 + let mut permits: Vec<_> = (0..3).map(|_| admission.acquire(EvmKind::CallPast).unwrap()).collect(); + + // saturate the pool to the relaxed limit with inspector tasks + let saturating: Vec<_> = (0..5).map(|_| admission.acquire(EvmKind::Inspect).unwrap()).collect(); + assert_eq!(permits.len() + saturating.len(), 8); + + // new call-past tasks are throttled and the kind is already over its limit + let rx = spawn_acquire(&admission, EvmKind::CallPast); + assert!(rx.recv_timeout(Duration::from_millis(100)).is_err(), "must block while over the limit"); + + // still over the limit after releasing one bypass permit + drop(permits.swap_remove(0)); + assert!(rx.recv_timeout(Duration::from_millis(100)).is_err(), "must still block while over the limit"); + + // below the limit after releasing a second one + drop(permits.swap_remove(0)); + assert!(rx.recv_timeout(Duration::from_secs(5)).is_ok(), "must proceed below the limit"); + } + + /// Ignored by default: `GlobalState` shutdown is process-global and cannot be reset, so this + /// test would poison the other admission tests when run in parallel. Run with `--ignored`. + #[test] + #[ignore = "triggers process-global shutdown"] + fn test_throttled_admission_fails_on_shutdown() { + let admission = admission(10, 80, 1); + + let _first = admission.acquire(EvmKind::CallPast).unwrap(); + let saturating: Vec<_> = (0..7).map(|_| admission.acquire(EvmKind::Inspect).unwrap()).collect(); + + let rx = spawn_acquire(&admission, EvmKind::CallPast); + assert!(rx.recv_timeout(Duration::from_millis(100)).is_err(), "must block while over the limit"); + + GlobalState::shutdown_from("test", "test_throttled_admission_returns_none_on_shutdown"); + assert!(rx.recv_timeout(Duration::from_secs(5)).is_ok(), "must return on shutdown"); + drop(saturating); + } +} diff --git a/src/eth/executor/transaction_worker.rs b/src/eth/executor/transaction_worker.rs index 81cae68df..fb6ee28bd 100644 --- a/src/eth/executor/transaction_worker.rs +++ b/src/eth/executor/transaction_worker.rs @@ -101,7 +101,7 @@ impl TransactionWorker { fn execute_local_transaction_attempts( storage: &StratusStorage, miner: &Miner, - evm: &mut Evm, + evm: &mut Evm, tx_input: TransactionInput, max_attempts: usize, ) -> LocalTransactionResult { @@ -172,7 +172,7 @@ impl TransactionTask { } } - fn execute(self, storage: &StratusStorage, miner: &Miner, evm: &mut Evm) -> anyhow::Result<(), StratusError> { + fn execute(self, storage: &StratusStorage, miner: &Miner, evm: &mut Evm) -> anyhow::Result<(), StratusError> { let Self { span, kind } = self; let _enter = span.enter(); diff --git a/src/eth/executor/types/mod.rs b/src/eth/executor/types/mod.rs index e8ff7e59d..b0f2b7804 100644 --- a/src/eth/executor/types/mod.rs +++ b/src/eth/executor/types/mod.rs @@ -8,9 +8,7 @@ pub use error::ExecutorError; pub use execution_result::ExecutionResult; pub use execution_result::RevertReason; pub use state::State; -pub use task::EvmRoute; -pub use task::EvmTask; pub use task::ExecutionTask; pub use task::InspectionTask; -pub use task::Task; +pub use task::PoolTask; pub use transaction_execution::TransactionExecution; diff --git a/src/eth/executor/types/task.rs b/src/eth/executor/types/task.rs index cfee25a3c..d711c20b0 100644 --- a/src/eth/executor/types/task.rs +++ b/src/eth/executor/types/task.rs @@ -1,25 +1,25 @@ use std::panic::AssertUnwindSafe; use std::panic::catch_unwind; +use std::sync::Arc; use alloy_rpc_types_trace::geth::GethTrace; use anyhow::anyhow; use tracing::Span; -use crate::eth::executor::TransactionExecutionInput; use crate::eth::executor::evm::Evm; use crate::eth::executor::evm::RevmResultAndState; use crate::eth::executor::evm::types::CallExecutionInput; use crate::eth::executor::evm::types::EvmInput; +use crate::eth::executor::evm::types::EvmKind; use crate::eth::executor::evm::types::ExecutionMetrics; use crate::eth::executor::evm::types::InspectorInput; +use crate::eth::executor::pool_admission::PoolAdmission; +use crate::eth::executor::pool_admission::PoolPermit; use crate::eth::executor::types::error::ExecutorError; +use crate::eth::types::PointInTime; +use crate::eth::types::StateError; use crate::eth::types::StratusError; -pub struct EvmTask { - pub span: Span, - task: T, -} - #[derive(derive_new::new)] pub struct ExecutionTask { pub input: Input, @@ -32,47 +32,72 @@ pub struct InspectionTask { pub response_tx: oneshot::Sender>, } -#[derive(Debug, Clone, strum::Display)] -pub enum EvmRoute { - #[strum(to_string = "call_present")] - CallPresent(CallExecutionInput), +/// A task for the unified EVM pool. +pub struct PoolTask { + span: Span, + permit: PoolPermit, + task: PoolTaskKind, +} - #[strum(to_string = "call_past")] - CallPast(CallExecutionInput), +enum PoolTaskKind { + Call(ExecutionTask), + Inspect(InspectionTask), } -impl From for EvmTask { - fn from(task: T) -> Self { - Self { span: Span::current(), task } +impl PoolTask { + /// Creates a call task, acquiring an admission slot for it. Fails when the pool is shutting down. + pub fn call(task: ExecutionTask, admission: &Arc) -> Result { + let permit = admission.acquire(call_evm_kind(&task.input))?; + Ok(Self { + span: Span::current(), + permit, + task: PoolTaskKind::Call(task), + }) } -} -impl EvmTask { - pub fn execute(self, evm: &mut Evm) -> anyhow::Result<(), StratusError> { - let _enter = self.span.enter(); - catch_unwind(AssertUnwindSafe(|| self.task.execute(evm))).map_err(|err| ExecutorError::Panic { err: anyhow!("{err:?}") }.into()) + /// Creates an inspection task, acquiring an admission slot for it. Fails when the pool is shutting down. + pub fn inspect(task: InspectionTask, admission: &Arc) -> Result { + let permit = admission.acquire(EvmKind::Inspect)?; + Ok(Self { + span: Span::current(), + permit, + task: PoolTaskKind::Inspect(task), + }) } -} -pub trait Task { - type Input: EvmInput; + /// Executes the task on the EVM. The admission slot is released when the task finishes. + pub fn execute(self, evm: &mut Evm) -> anyhow::Result<(), StratusError> { + let Self { span, permit, task } = self; + let _enter = span.enter(); + let _busy = permit.evm_kind().mark_executor_pool_busy(); - fn execute(self, evm: &mut Evm); + catch_unwind(AssertUnwindSafe(move || match task { + PoolTaskKind::Call(task) => task.execute(evm), + PoolTaskKind::Inspect(task) => task.execute(evm), + })) + .map_err(|err| ExecutorError::Panic { err: anyhow!("{err:?}") }.into()) + } } -impl Task for ExecutionTask { - type Input = Input; +/// Returns the pool kind of a call: calls against the latest state and calls against a past state +/// are admitted by different pool gates. +fn call_evm_kind(input: &CallExecutionInput) -> EvmKind { + match input.kind.point_in_time() { + PointInTime::Pending | PointInTime::Latest => EvmKind::CallPresent, + PointInTime::Past(_) => EvmKind::CallPast, + } +} - fn execute(self, evm: &mut Evm) { +impl ExecutionTask { + fn execute(self, evm: &mut Evm) { if let Err(e) = self.response_tx.send(evm.execute(self.input)) { tracing::error!(reason = ?e, "failed to send evm task execution result"); } } } -impl Task for InspectionTask { - type Input = TransactionExecutionInput; - fn execute(self, evm: &mut Evm) { +impl InspectionTask { + fn execute(self, evm: &mut Evm) { if let Err(e) = self.response_tx.send(evm.inspect(self.input)) { tracing::error!(reason = ?e, "failed to send evm task execution result"); } diff --git a/src/utils.rs b/src/utils.rs index d38e363b4..6809c0414 100644 --- a/src/utils.rs +++ b/src/utils.rs @@ -32,7 +32,6 @@ impl Drop for DropTimer { #[derive(Deref, Default)] pub struct Semaphore { - // refac to another file #[deref] sem: Arc, } diff --git a/tests/config_loader.rs b/tests/config_loader.rs index 86275aca7..66ab03bf1 100644 --- a/tests/config_loader.rs +++ b/tests/config_loader.rs @@ -19,7 +19,7 @@ fn test_cli_overrides_file() { [executor] chain_id = 2008 - call_present_evms = 11 + call_past_limit = 11 [rpc] address = "0.0.0.0:3001" @@ -32,7 +32,7 @@ fn test_cli_overrides_file() { let config = load_with(&[], file).unwrap(); assert!(config.leader); assert_eq!(config.executor.executor_chain_id, 2008); - assert_eq!(config.executor.call_present_evms, 11); + assert_eq!(config.executor.pool.call_past_limit, 11); assert_eq!(config.rpc_server.rpc_address.to_string(), "0.0.0.0:3001"); assert_eq!(config.miner.block_mode, MinerMode::Interval(std::time::Duration::from_secs(1))); @@ -40,7 +40,7 @@ fn test_cli_overrides_file() { let config = load_with(&["--executor-chain-id", "9999", "-a", "0.0.0.0:3002", "--block-mode", "automine"], file).unwrap(); assert_eq!(config.executor.executor_chain_id, 9999); // file value preserved when not overridden in the CLI - assert_eq!(config.executor.call_present_evms, 11); + assert_eq!(config.executor.pool.call_past_limit, 11); assert_eq!(config.rpc_server.rpc_address.to_string(), "0.0.0.0:3002"); assert_eq!(config.miner.block_mode, MinerMode::Automine); } @@ -53,7 +53,7 @@ fn test_clap_defaults_do_not_override_file() { [executor] chain_id = 2008 - call_present_evms = 11 + call_past_limit = 11 [common] blocking_threads = 64 @@ -62,7 +62,7 @@ fn test_clap_defaults_do_not_override_file() { let config = load_with(&["--async-threads", "8"], file).unwrap(); assert_eq!(config.common.num_async_threads, 8); assert_eq!(config.common.num_blocking_threads, 64); - assert_eq!(config.executor.call_present_evms, 11); + assert_eq!(config.executor.pool.call_past_limit, 11); } #[test]