From 6654d981eeacb627070dd3f6a8f41a9588848c53 Mon Sep 17 00:00:00 2001 From: gventino-cw Date: Fri, 25 Sep 2026 12:54:59 -0300 Subject: [PATCH 1/4] feat: unified EVM worker pool Replace the per-kind EVM worker pools (call-present, call-past, inspector) with a single pool shared by every execution kind, closing #2708. - One set of `executor.evm_workers` OS threads (default 150) pulling from a single bounded channel; each worker's EVM executes any kind, since the point-in-time is selected per task by its input. - Per-kind concurrency limits (`executor.call_present_limit`, `executor.call_past_limit`, `executor.inspector_limit`) enforced by counting semaphores acquired on the sender side: the permit rides inside the task and is released after execution, so no kind can be head-of-line blocked by unrelated kinds and FIFO fairness is preserved. - `call_present_limit` defaults to the pool capacity remaining after the other kinds, so the busiest kind can use every idle worker. - `executor.evm_flex_quota` (default 0) adds extra permits any saturated kind can borrow, sharing idle capacity across kinds. - Deprecated `executor.{call_present,call_past,inspector}_evms` fields are aliases that preserve the total capacity of pre-migration deployments. - Generalized `Semaphore` with `try_acquire` and shutdown-aware acquire (condvar polling), with pluggable metrics; new pool gauges `executor_workers_total`, `executor_pool_permits_waiting`, `executor_pool_permits_held` and `executor_pool_queue_len`. - De-genericized `Evm` (the `Input` parameter was PhantomData-only) and moved the worker loop into `EvmWorkerPool::worker`. --- config/stratus.example.toml | 17 +- crates/stratus_metrics/src/definitions.rs | 14 +- src/config/loader.rs | 5 + src/eth/executor/config.rs | 58 +++- src/eth/executor/evm/mod.rs | 11 +- src/eth/executor/evm_worker_pool.rs | 368 +++++++++++++++++----- src/eth/executor/mod.rs | 13 +- src/eth/executor/transaction_worker.rs | 4 +- src/eth/executor/types/mod.rs | 3 +- src/eth/executor/types/task.rs | 75 +++-- src/main.rs | 2 +- src/utils.rs | 107 ++++++- tests/config_loader.rs | 6 +- 13 files changed, 533 insertions(+), 150 deletions(-) diff --git a/config/stratus.example.toml b/config/stratus.example.toml index d3f57525e..5f02ad716 100644 --- a/config/stratus.example.toml +++ b/config/stratus.example.toml @@ -65,11 +65,22 @@ [executor] # Chain ID of the network. Required. chain_id = 2008 -# Number of EVMs to execute calls against the present state. +# Total number of EVM workers in the unified pool, shared by every execution kind. +# Defaults to the sum of the deprecated per-kind sizes, or 150 when none is set. +# 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 +# Extra permits that any kind can borrow when its own limit is exhausted (for example 20% of the pool). +# evm_flex_quota = 0 +# Deprecated: number of EVMs to execute calls against the present state. # call_present_evms = 50 -# Number of EVMs to execute calls against past states. +# Deprecated: number of EVMs to execute calls against past states. # call_past_evms = 50 -# Number of EVMs to execute inspector calls. +# Deprecated: number of EVMs to execute inspector calls. # inspector_evms = 50 # Should reject contract transactions and calls to accounts that are not contracts? # reject_not_contract = true diff --git a/crates/stratus_metrics/src/definitions.rs b/crates/stratus_metrics/src/definitions.rs index 9c04a010b..dd32f0c60 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 waiting to acquire an executor pool permit, by semaphore kind." + gauge executor_pool_permits_waiting{kind}, + + "Number of executor pool permits currently held, by semaphore kind." + gauge executor_pool_permits_held{kind}, + + "Length of the executor pool task queue." + gauge executor_pool_queue_len{} }, group: rocks { diff --git a/src/config/loader.rs b/src/config/loader.rs index 0fbd8de3a..c8f37948f 100644 --- a/src/config/loader.rs +++ b/src/config/loader.rs @@ -431,6 +431,11 @@ mod tests { [executor] chain_id = 100 + evm_workers = 10 + call_present_limit = 7 + call_past_limit = 4 + inspector_limit = 5 + evm_flex_quota = 1 call_present_evms = 1 call_past_evms = 2 inspector_evms = 3 diff --git a/src/eth/executor/config.rs b/src/eth/executor/config.rs index 40e9542a5..b0d1b7d44 100644 --- a/src/eth/executor/config.rs +++ b/src/eth/executor/config.rs @@ -16,14 +16,25 @@ 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, + /// Total number of EVM workers in the unified pool, shared by every execution kind. + #[arg(id = "executor.evm_workers", long = "executor-evm-workers")] + pub evm_workers: Option, - #[arg(id = "executor.call_past_evms", long = "executor-call-past-evms", default_value_t = 50)] - pub call_past_evms: usize, + /// Maximum number of concurrent call-present executions. + #[arg(id = "executor.call_present_limit", long = "executor-call-present-limit")] + pub call_present_limit: Option, - #[arg(id = "executor.inspector_evms", long = "executor-inspector-evms", default_value_t = 50)] - pub inspector_evms: usize, + /// Maximum number of concurrent call-past executions. + #[arg(id = "executor.call_past_limit", long = "executor-call-past-limit")] + pub call_past_limit: Option, + + /// Maximum number of concurrent inspector executions. + #[arg(id = "executor.inspector_limit", long = "executor-inspector-limit")] + pub inspector_limit: Option, + + /// Extra permits that any execution kind can borrow when its own limit is exhausted. + #[arg(id = "executor.evm_flex_quota", long = "executor-evm-flex-quota", default_value_t = 0)] + pub evm_flex_quota: usize, /// Should reject contract transactions and calls to accounts that are not contracts? #[arg( @@ -41,6 +52,18 @@ pub struct ExecutorConfig { #[arg(id = "executor.evm_spec", long = "executor-evm-spec", default_value = "Prague", value_parser = parse_evm_spec)] #[serde(rename = "evm_spec", with = "spec_id_serde")] pub executor_evm_spec: SpecId, + + /// Deprecated alias of `executor.call_present_limit`. + #[arg(id = "executor.call_present_evms", long = "executor-call-present-evms")] + pub call_present_evms: Option, + + /// Deprecated alias of `executor.call_past_limit`. + #[arg(id = "executor.call_past_evms", long = "executor-call-past-evms")] + pub call_past_evms: Option, + + /// Deprecated alias of `executor.inspector_limit`. + #[arg(id = "executor.inspector_evms", long = "executor-inspector-evms")] + pub inspector_evms: Option, } #[cfg(test)] @@ -48,11 +71,16 @@ impl Default for ExecutorConfig { fn default() -> Self { Self { executor_chain_id: 0, - call_present_evms: 50, - call_past_evms: 50, - inspector_evms: 50, + evm_workers: None, + call_present_limit: None, + call_past_limit: None, + inspector_limit: None, + evm_flex_quota: 0, executor_reject_not_contract: true, executor_evm_spec: SpecId::PRAGUE, + call_present_evms: None, + call_past_evms: None, + inspector_evms: None, } } } @@ -81,14 +109,20 @@ fn parse_evm_spec(input: &str) -> anyhow::Result { } impl ExecutorConfig { + /// Returns whether any deprecated per-kind pool size field is set + /// (`call_present_evms`, `call_past_evms` or `inspector_evms`). + pub fn has_deprecated_pool_sizes(&self) -> bool { + self.call_present_evms.is_some() || self.call_past_evms.is_some() || self.inspector_evms.is_some() + } + /// Initializes Executor. /// /// Note: Should be called only after async runtime is initialized. - pub fn init(&self, storage: Arc, miner: Arc) -> Arc { + pub fn init(&self, storage: Arc, miner: Arc) -> anyhow::Result> { let config = *self; tracing::info!(?config, "creating executor"); - let executor = Executor::new(storage, miner, config); - Arc::new(executor) + let executor = Executor::new(storage, miner, config)?; + Ok(Arc::new(executor)) } } 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..af8f76cbd 100644 --- a/src/eth/executor/evm_worker_pool.rs +++ b/src/eth/executor/evm_worker_pool.rs @@ -1,6 +1,8 @@ use std::sync::Arc; use alloy_rpc_types_trace::geth::GethTrace; +use anyhow::anyhow; +use anyhow::bail; use stratus_metrics as metrics; use crate::GlobalState; @@ -10,111 +12,207 @@ use crate::eth::executor::ExecutorError; use crate::eth::executor::evm::Evm; 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::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; +use crate::utils::Permit; +use crate::utils::Semaphore; +use crate::utils::SemaphoreMetrics; -/// 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>>, +/// Total capacity of the unified EVM pool task queue. +const TASK_QUEUE_CAPACITY: usize = 4096; + +/// Default number of EVM workers in the unified pool (sum of the old per-kind pool defaults). +const DEFAULT_WORKERS: usize = 150; + +/// Default maximum number of concurrent call-past and inspector executions. +const DEFAULT_KIND_LIMIT: usize = 50; + +/// Effective configuration of the unified EVM pool, resolved from [`ExecutorConfig`]. +#[derive(Clone, Copy, Debug)] +pub struct PoolConfig { + /// Total number of EVM workers, shared by every execution kind. + pub workers: usize, + + /// Maximum number of concurrent call-present executions. + pub call_present_limit: usize, - /// 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>>, + /// Maximum number of concurrent call-past executions. + pub call_past_limit: usize, - /// Pool for parallel execution of tx inspections (debug_traceTransaction). Usually contains multiple EVMs. - pub inspector: crossbeam_channel::Sender>, + /// Maximum number of concurrent inspector executions. + pub inspector_limit: usize, + + /// Extra permits that any execution kind can borrow when its own limit is exhausted. + pub flex_quota: usize, } -impl EvmWorkerPool { - /// Spawns EVM tasks in background. - 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); +impl PoolConfig { + /// Resolves the effective pool configuration, mapping deprecated fields to their new meaning. + pub fn resolve(config: &ExecutorConfig) -> anyhow::Result { + if config.evm_workers == Some(0) { + bail!("executor.evm_workers must be greater than zero"); } - // 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); - }); + for (field, value) in [ + ("executor.call_present_evms", config.call_present_evms), + ("executor.call_past_evms", config.call_past_evms), + ("executor.inspector_evms", config.inspector_evms), + ] { + if value.is_some() { + tracing::warn!( + field, + "deprecated executor pool field; use executor.evm_workers and the per-kind limit fields instead" + ); } - metrics::set_executor_workers_busy(0, kind); - evm_tx } - 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); + let call_past_limit = config.call_past_limit.or(config.call_past_evms).unwrap_or(DEFAULT_KIND_LIMIT); + let inspector_limit = config.inspector_limit.or(config.inspector_evms).unwrap_or(DEFAULT_KIND_LIMIT); + + let workers = match config.evm_workers { + Some(workers) => workers, + // no old field set: default pool size + None if !config.has_deprecated_pool_sizes() => DEFAULT_WORKERS, + // at least one old field set: preserve the total capacity of the old per-kind pools + None => + config.call_present_evms.unwrap_or(DEFAULT_KIND_LIMIT) + + config.call_past_evms.unwrap_or(DEFAULT_KIND_LIMIT) + + config.inspector_evms.unwrap_or(DEFAULT_KIND_LIMIT), + }; + + let call_present_limit = config + .call_present_limit + .or(config.call_present_evms) + .unwrap_or_else(|| workers.saturating_sub(call_past_limit + inspector_limit)); + + let resolved = Self { + workers, + call_present_limit, + call_past_limit, + inspector_limit, + flex_quota: config.evm_flex_quota, + }; + + let limits_sum = call_present_limit + call_past_limit + inspector_limit; + if limits_sum > workers { + bail!( + "executor pool kind limits ({call_present_limit} call-present + {call_past_limit} call-past + {inspector_limit} inspector = {limits_sum}) \ + exceed the total number of workers ({workers}); increase executor.evm_workers or lower the limits" + ); + } + + tracing::info!(?resolved, "unified EVM pool configuration resolved"); + Ok(resolved) + } +} + +/// Per-kind concurrency limits of the unified EVM pool. +struct PoolLimits { + call_present: Semaphore, + call_past: Semaphore, + inspector: Semaphore, + flex: Semaphore, +} + +impl PoolLimits { + fn new(config: PoolConfig) -> Self { + Self { + call_present: Semaphore::with_metrics(config.call_present_limit, SemaphoreMetrics::Pool("call_present")), + call_past: Semaphore::with_metrics(config.call_past_limit, SemaphoreMetrics::Pool("call_past")), + inspector: Semaphore::with_metrics(config.inspector_limit, SemaphoreMetrics::Pool("inspector")), + flex: Semaphore::with_metrics(config.flex_quota, SemaphoreMetrics::Pool("flex")), + } + } + + /// Own-limit semaphore of an execution kind. + fn own(&self, kind: EvmKind) -> &Semaphore { + 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"), + } + } + + /// Acquires a permit for the kind, borrowing from the flex quota when the own limit is exhausted. + fn acquire(&self, kind: EvmKind) -> Option { + if let Some(permit) = self.own(kind).try_acquire() { + return Some(permit); + } + + if let Some(permit) = self.flex.try_acquire() { + return Some(permit); + } + + self.own(kind).acquire_shutdown_aware() + } +} + +/// Manages the unified EVM pool: one shared set of workers serving every execution kind. +pub struct EvmWorkerPool { + tx: crossbeam_channel::Sender, + limits: Arc, +} + +impl EvmWorkerPool { + /// Spawns the unified EVM pool workers. + pub fn spawn(storage: Arc, config: &ExecutorConfig) -> anyhow::Result { + let pool = PoolConfig::resolve(config)?; + let (tx, rx) = crossbeam_channel::bounded::(TASK_QUEUE_CAPACITY); + + for worker_index in 1..=pool.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); + }); + } - EvmWorkerPool { - call_present, - call_past, - inspector, + // 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); } + metrics::set_executor_workers_total(pool.workers as u64); + + Ok(Self { + tx, + limits: Arc::new(PoolLimits::new(pool)), + }) } - /// Executes a transaction in the specified route. + /// Executes a call in the specified route. pub fn execute(&self, route: EvmRoute) -> Result<(Output, ExecutionMetrics), StratusError> where Output: TryFrom, { + let kind = match &route { + EvmRoute::CallPresent(_) => EvmKind::CallPresent, + EvmRoute::CallPast(_) => EvmKind::CallPast, + }; + + let Some(permit) = self.limits.acquire(kind) else { + return Err(UnexpectedError::Unexpected(anyhow!("executor pool is shutting down")).into()); + }; + 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 = match route { + EvmRoute::CallPresent(input) => PoolTask::call(ExecutionTask::new(input, execution_tx), kind, permit), + EvmRoute::CallPast(input) => PoolTask::call(ExecutionTask::new(input, execution_tx), kind, permit), }; + self.tx.send(task)?; + metrics::set_executor_pool_queue_len(self.tx.len() as u64); match execution_rx.recv() { Ok(result) => { @@ -125,13 +223,129 @@ impl EvmWorkerPool { } } + /// Executes a transaction inspection (debug_traceTransaction). pub fn inspect(&self, input: InspectorInput) -> Result { + let Some(permit) = self.limits.acquire(EvmKind::Inspect) else { + return Err(UnexpectedError::Unexpected(anyhow!("executor pool is shutting down")).into()); + }; + 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), permit); + 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 super::*; + + fn test_config() -> ExecutorConfig { + ExecutorConfig { + executor_chain_id: 1, + ..Default::default() + } + } + + #[test] + fn test_pool_config_resolves_defaults() { + let config = test_config(); + let pool = PoolConfig::resolve(&config).unwrap(); + assert_eq!(pool.workers, DEFAULT_WORKERS); + assert_eq!(pool.call_present_limit, 50); + assert_eq!(pool.call_past_limit, 50); + assert_eq!(pool.inspector_limit, 50); + assert_eq!(pool.flex_quota, 0); + } + + #[test] + fn test_pool_config_maps_deprecated_fields() { + let mut config = test_config(); + config.call_present_evms = Some(100); + config.call_past_evms = Some(20); + config.inspector_evms = Some(30); + let pool = PoolConfig::resolve(&config).unwrap(); + // total capacity preserved: 100 + 20 + 30 + assert_eq!(pool.workers, 150); + assert_eq!(pool.call_present_limit, 100); + assert_eq!(pool.call_past_limit, 20); + assert_eq!(pool.inspector_limit, 30); + } + + #[test] + fn test_pool_config_deprecated_fields_with_unset_kinds() { + let mut config = test_config(); + config.call_present_evms = Some(100); + let pool = PoolConfig::resolve(&config).unwrap(); + // unset deprecated fields keep their old defaults (50) when computing the total + assert_eq!(pool.workers, 200); + assert_eq!(pool.call_present_limit, 100); + assert_eq!(pool.call_past_limit, 50); + assert_eq!(pool.inspector_limit, 50); + } + + #[test] + fn test_pool_config_new_fields_take_precedence() { + let mut config = test_config(); + config.evm_workers = Some(200); + config.call_present_limit = Some(120); + config.call_past_limit = Some(20); + config.inspector_limit = Some(30); + config.evm_flex_quota = 40; + let pool = PoolConfig::resolve(&config).unwrap(); + assert_eq!(pool.workers, 200); + assert_eq!(pool.call_present_limit, 120); + assert_eq!(pool.call_past_limit, 20); + assert_eq!(pool.inspector_limit, 30); + assert_eq!(pool.flex_quota, 40); + } + + #[test] + fn test_pool_config_call_present_uses_remaining_capacity() { + let mut config = test_config(); + config.evm_workers = Some(200); + let pool = PoolConfig::resolve(&config).unwrap(); + assert_eq!(pool.call_present_limit, 200 - 50 - 50); + } + + #[test] + fn test_pool_config_rejects_limits_exceeding_workers() { + let mut config = test_config(); + config.evm_workers = Some(100); + config.call_present_limit = Some(60); + config.call_past_limit = Some(50); + config.inspector_limit = Some(50); + assert!(PoolConfig::resolve(&config).is_err()); + } + + #[test] + fn test_pool_config_rejects_zero_workers() { + let mut config = test_config(); + config.evm_workers = Some(0); + assert!(PoolConfig::resolve(&config).is_err()); + } } diff --git a/src/eth/executor/mod.rs b/src/eth/executor/mod.rs index 1aeae0268..6b6fa1b32 100644 --- a/src/eth/executor/mod.rs +++ b/src/eth/executor/mod.rs @@ -60,6 +60,7 @@ use crate::ext::OptionExt; use crate::ext::to_json_string; use crate::infra::tracing::SpanExt; use crate::utils::Semaphore; +use crate::utils::SemaphoreMetrics; // ----------------------------------------------------------------------------- // Executor @@ -83,18 +84,18 @@ pub struct Executor { } impl Executor { - pub fn new(storage: Arc, miner: Arc, config: ExecutorConfig) -> Self { + pub fn new(storage: Arc, miner: Arc, config: ExecutorConfig) -> anyhow::Result { tracing::info!(?config, "creating executor"); let reject_not_contract = config.executor_reject_not_contract; let transaction_worker = TransactionWorker::spawn(Arc::clone(&storage), Arc::clone(&miner), &config); - let evms = EvmWorkerPool::spawn(Arc::clone(&storage), &config); - Self { - transaction_warmup: Semaphore::new(100), + let evms = EvmWorkerPool::spawn(Arc::clone(&storage), &config)?; + Ok(Self { + transaction_warmup: Semaphore::with_metrics(100, SemaphoreMetrics::LocalTransaction), transaction_worker, evms, storage, reject_not_contract, - } + }) } // ------------------------------------------------------------------------- @@ -145,7 +146,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, 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..562cf2ddd 100644 --- a/src/eth/executor/types/mod.rs +++ b/src/eth/executor/types/mod.rs @@ -9,8 +9,7 @@ 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..d6b21390f 100644 --- a/src/eth/executor/types/task.rs +++ b/src/eth/executor/types/task.rs @@ -5,20 +5,16 @@ 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::types::error::ExecutorError; use crate::eth::types::StratusError; - -pub struct EvmTask { - pub span: Span, - task: T, -} +use crate::utils::Permit; #[derive(derive_new::new)] pub struct ExecutionTask { @@ -41,38 +37,67 @@ pub enum EvmRoute { CallPast(CallExecutionInput), } -impl From for EvmTask { - fn from(task: T) -> Self { - Self { span: Span::current(), task } - } +/// A task for the unified EVM pool. +pub struct PoolTask { + pub span: Span, + kind: EvmKind, + permit: Permit, + task: PoolTaskKind, } -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()) - } +enum PoolTaskKind { + Call(ExecutionTask), + Inspect(InspectionTask), } -pub trait Task { - type Input: EvmInput; +impl PoolTask { + pub fn call(task: ExecutionTask, kind: EvmKind, permit: Permit) -> Self { + debug_assert!(matches!(kind, EvmKind::CallPresent | EvmKind::CallPast)); + Self { + span: Span::current(), + kind, + permit, + task: PoolTaskKind::Call(task), + } + } + + pub fn inspect(task: InspectionTask, permit: Permit) -> Self { + Self { + span: Span::current(), + kind: EvmKind::Inspect, + permit, + task: PoolTaskKind::Inspect(task), + } + } - fn execute(self, evm: &mut Evm); -} + pub fn execute(self, evm: &mut Evm) -> anyhow::Result<(), StratusError> { + let Self { + span, + kind, + permit: _permit, + task, + } = self; + let _enter = span.enter(); + let _busy = kind.mark_executor_pool_busy(); -impl Task for ExecutionTask { - type Input = Input; + 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()) + } +} - 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/main.rs b/src/main.rs index 84aa5bbfe..c5f138c35 100644 --- a/src/main.rs +++ b/src/main.rs @@ -26,7 +26,7 @@ async fn run(config: StratusConfig) -> anyhow::Result<()> { let miner = config.miner.init(Arc::clone(&storage)).await?; // Init executor - let executor = config.executor.init(Arc::clone(&storage), Arc::clone(&miner)); + let executor = config.executor.init(Arc::clone(&storage), Arc::clone(&miner))?; let (consensus, importer_runtime) = if let Some(importer_config) = &config.importer { tracing::info!(?importer_config, "creating importer"); diff --git a/src/utils.rs b/src/utils.rs index d38e363b4..b05c02d3b 100644 --- a/src/utils.rs +++ b/src/utils.rs @@ -1,12 +1,13 @@ use std::sync::Arc; +use std::time::Duration; use derive_more::Deref; use parking_lot::Condvar; use parking_lot::Mutex; -#[cfg(feature = "metrics")] -use stratus_metrics as metrics; use tokio::time::Instant; +use crate::GlobalState; + /// Amount of bytes in one GB (technically, GiB). pub const GIGABYTE: usize = 1024 * 1024 * 1024; @@ -30,9 +31,11 @@ impl Drop for DropTimer { } } +/// Interval between shutdown checks while blocked acquiring a shutdown-aware semaphore permit. +const SHUTDOWN_POLL_INTERVAL: Duration = Duration::from_millis(250); + #[derive(Deref, Default)] pub struct Semaphore { - // refac to another file #[deref] sem: Arc, } @@ -41,6 +44,55 @@ pub struct Semaphore { pub struct SemaphoreInner { permits: Mutex, cvar: Condvar, + metrics: SemaphoreMetrics, +} + +/// Metrics recorded by a [`Semaphore`] when its permits change. +#[derive(Clone, Copy, Default)] +pub enum SemaphoreMetrics { + /// Do not record any metric. + #[default] + Disabled, + + /// Legacy metrics of the local transaction warmup semaphore (unlabeled). + LocalTransaction, + + /// Unified executor pool permit metrics, labeled by semaphore kind. + Pool(&'static str), +} + +impl SemaphoreMetrics { + fn waiting_added(&self) { + match *self { + Self::Disabled => {} + Self::LocalTransaction => stratus_metrics::inc_executor_local_transaction_semaphore_waiting(1), + Self::Pool(kind) => stratus_metrics::inc_executor_pool_permits_waiting(1, kind), + } + } + + fn waiting_removed(&self) { + match *self { + Self::Disabled => {} + Self::LocalTransaction => stratus_metrics::dec_executor_local_transaction_semaphore_waiting(1), + Self::Pool(kind) => stratus_metrics::dec_executor_pool_permits_waiting(1, kind), + } + } + + fn permit_acquired(&self) { + match *self { + Self::Disabled => {} + Self::LocalTransaction => stratus_metrics::inc_executor_local_transaction_permit_holders(1), + Self::Pool(kind) => stratus_metrics::inc_executor_pool_permits_held(1, kind), + } + } + + fn permit_released(&self) { + match *self { + Self::Disabled => {} + Self::LocalTransaction => stratus_metrics::dec_executor_local_transaction_permit_holders(1), + Self::Pool(kind) => stratus_metrics::dec_executor_pool_permits_held(1, kind), + } + } } pub struct Permit { @@ -49,29 +101,63 @@ pub struct Permit { impl Semaphore { pub fn new(permits: usize) -> Self { + Self::with_metrics(permits, SemaphoreMetrics::Disabled) + } + + pub fn with_metrics(permits: usize, metrics: SemaphoreMetrics) -> Self { Self { sem: Arc::new(SemaphoreInner { permits: Mutex::new(permits), cvar: Condvar::new(), + metrics, }), } } + /// Blocks until a permit is available. pub fn acquire(&self) -> Permit { - #[cfg(feature = "metrics")] - metrics::inc_executor_local_transaction_semaphore_waiting(1); + self.metrics.waiting_added(); let mut permits = self.permits.lock(); while *permits == 0 { self.cvar.wait(&mut permits); } *permits -= 1; drop(permits); - #[cfg(feature = "metrics")] - metrics::dec_executor_local_transaction_semaphore_waiting(1); - #[cfg(feature = "metrics")] - metrics::inc_executor_local_transaction_permit_holders(1); + self.metrics.waiting_removed(); + self.metrics.permit_acquired(); Permit { sem: Arc::clone(&self.sem) } } + + /// Tries to acquire a permit without blocking. + pub fn try_acquire(&self) -> Option { + let mut permits = self.permits.lock(); + if *permits == 0 { + return None; + } + *permits -= 1; + drop(permits); + self.metrics.permit_acquired(); + Some(Permit { sem: Arc::clone(&self.sem) }) + } + + /// Blocks until a permit is available or the application starts shutting down. + pub fn acquire_shutdown_aware(&self) -> Option { + self.metrics.waiting_added(); + let mut permits = self.permits.lock(); + while *permits == 0 { + if GlobalState::is_shutdown() { + drop(permits); + self.metrics.waiting_removed(); + return None; + } + self.cvar.wait_for(&mut permits, SHUTDOWN_POLL_INTERVAL); + } + *permits -= 1; + drop(permits); + self.metrics.waiting_removed(); + self.metrics.permit_acquired(); + Some(Permit { sem: Arc::clone(&self.sem) }) + } } impl Drop for Permit { @@ -79,8 +165,7 @@ impl Drop for Permit { let mut permits = self.sem.permits.lock(); *permits += 1; self.sem.cvar.notify_one(); - #[cfg(feature = "metrics")] - metrics::dec_executor_local_transaction_permit_holders(1); + self.sem.metrics.permit_released(); } } diff --git a/tests/config_loader.rs b/tests/config_loader.rs index 254dc1f4f..97b52e1cf 100644 --- a/tests/config_loader.rs +++ b/tests/config_loader.rs @@ -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.call_present_evms, Some(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.call_present_evms, Some(11)); assert_eq!(config.rpc_server.rpc_address.to_string(), "0.0.0.0:3002"); assert_eq!(config.miner.block_mode, MinerMode::Automine); } @@ -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.call_present_evms, Some(11)); } #[test] From 5bf82f6b91848c1c24551aaf13862cdc32cf77e0 Mon Sep 17 00:00:00 2001 From: gventino-cw Date: Fri, 25 Sep 2026 14:55:53 -0300 Subject: [PATCH 2/4] refactor: rename PoolTask fields to evm_kind and task_kind --- src/eth/executor/types/task.rs | 20 ++++++++++---------- 1 file changed, 10 insertions(+), 10 deletions(-) diff --git a/src/eth/executor/types/task.rs b/src/eth/executor/types/task.rs index d6b21390f..6e819f8bd 100644 --- a/src/eth/executor/types/task.rs +++ b/src/eth/executor/types/task.rs @@ -40,9 +40,9 @@ pub enum EvmRoute { /// A task for the unified EVM pool. pub struct PoolTask { pub span: Span, - kind: EvmKind, + evm_kind: EvmKind, permit: Permit, - task: PoolTaskKind, + task_kind: PoolTaskKind, } enum PoolTaskKind { @@ -55,32 +55,32 @@ impl PoolTask { debug_assert!(matches!(kind, EvmKind::CallPresent | EvmKind::CallPast)); Self { span: Span::current(), - kind, + evm_kind: kind, permit, - task: PoolTaskKind::Call(task), + task_kind: PoolTaskKind::Call(task), } } pub fn inspect(task: InspectionTask, permit: Permit) -> Self { Self { span: Span::current(), - kind: EvmKind::Inspect, + evm_kind: EvmKind::Inspect, permit, - task: PoolTaskKind::Inspect(task), + task_kind: PoolTaskKind::Inspect(task), } } pub fn execute(self, evm: &mut Evm) -> anyhow::Result<(), StratusError> { let Self { span, - kind, + evm_kind, permit: _permit, - task, + task_kind, } = self; let _enter = span.enter(); - let _busy = kind.mark_executor_pool_busy(); + let _busy = evm_kind.mark_executor_pool_busy(); - catch_unwind(AssertUnwindSafe(move || match task { + catch_unwind(AssertUnwindSafe(move || match task_kind { PoolTaskKind::Call(task) => task.execute(evm), PoolTaskKind::Inspect(task) => task.execute(evm), })) From 58edf25b6761150a5fe309c4cbe8f7df4333fa0c Mon Sep 17 00:00:00 2001 From: gventino-cw Date: Mon, 28 Sep 2026 10:42:21 -0300 Subject: [PATCH 3/4] refac: busy/relaxed admission pool to solve comments --- config/stratus.example.toml | 11 +- crates/stratus_metrics/src/definitions.rs | 14 +- src/config/loader.rs | 15 +- src/eth/executor/config.rs | 52 ++--- src/eth/executor/evm_worker_pool.rs | 212 ++++++++----------- src/eth/executor/mod.rs | 1 + src/eth/executor/pool_admission.rs | 245 ++++++++++++++++++++++ src/eth/executor/types/task.rs | 8 +- src/utils.rs | 44 ---- tests/config_loader.rs | 10 +- 10 files changed, 374 insertions(+), 238 deletions(-) create mode 100644 src/eth/executor/pool_admission.rs diff --git a/config/stratus.example.toml b/config/stratus.example.toml index 5f02ad716..27b185a89 100644 --- a/config/stratus.example.toml +++ b/config/stratus.example.toml @@ -66,7 +66,6 @@ # Chain ID of the network. Required. chain_id = 2008 # Total number of EVM workers in the unified pool, shared by every execution kind. -# Defaults to the sum of the deprecated per-kind sizes, or 150 when none is set. # evm_workers = 150 # Maximum number of concurrent calls against the present state. Defaults to the remaining pool capacity. # call_present_limit = 50 @@ -74,14 +73,8 @@ chain_id = 2008 # call_past_limit = 50 # Maximum number of concurrent inspector calls. # inspector_limit = 50 -# Extra permits that any kind can borrow when its own limit is exhausted (for example 20% of the pool). -# evm_flex_quota = 0 -# Deprecated: number of EVMs to execute calls against the present state. -# call_present_evms = 50 -# Deprecated: number of EVMs to execute calls against past states. -# call_past_evms = 50 -# Deprecated: number of EVMs to execute inspector calls. -# inspector_evms = 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 dd32f0c60..e040f7f93 100644 --- a/crates/stratus_metrics/src/definitions.rs +++ b/crates/stratus_metrics/src/definitions.rs @@ -112,11 +112,17 @@ metrics! { "Total number of EVM workers in the unified executor pool." gauge executor_workers_total{}, - "Number of tasks waiting to acquire an executor pool permit, by semaphore kind." - gauge executor_pool_permits_waiting{kind}, + "Number of tasks currently executing in the unified executor pool, by kind." + gauge executor_pool_inflight{kind}, - "Number of executor pool permits currently held, by semaphore kind." - gauge executor_pool_permits_held{kind}, + "Number of tasks waiting to be admitted to the unified executor pool, by kind." + gauge executor_pool_waiting{kind}, + + "Number of tasks currently executing in the unified executor pool." + gauge executor_pool_inflight_total{}, + + "Number of tasks admitted through the relaxed path while the pool was below the busy threshold, by kind." + counter executor_pool_relaxed_admissions{kind}, "Length of the executor pool task queue." gauge executor_pool_queue_len{} diff --git a/src/config/loader.rs b/src/config/loader.rs index 9ea7dbc0b..d72f746d5 100644 --- a/src/config/loader.rs +++ b/src/config/loader.rs @@ -375,9 +375,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" @@ -573,13 +575,10 @@ mod tests { [executor] chain_id = 100 evm_workers = 10 - call_present_limit = 7 + call_present_limit = 1 call_past_limit = 4 inspector_limit = 5 - evm_flex_quota = 1 - call_present_evms = 1 - call_past_evms = 2 - inspector_evms = 3 + evm_busy_threshold = 75 reject_not_contract = false evm_spec = "Cancun" diff --git a/src/eth/executor/config.rs b/src/eth/executor/config.rs index b0d1b7d44..31e2e2156 100644 --- a/src/eth/executor/config.rs +++ b/src/eth/executor/config.rs @@ -6,6 +6,9 @@ use display_json::DebugAsJson; use revm::primitives::hardfork::SpecId; use crate::eth::executor::Executor; +use crate::eth::executor::evm_worker_pool::DEFAULT_BUSY_THRESHOLD; +use crate::eth::executor::evm_worker_pool::DEFAULT_KIND_LIMIT; +use crate::eth::executor::evm_worker_pool::DEFAULT_WORKERS; use crate::eth::miner::Miner; use crate::eth::storage::StratusStorage; @@ -17,24 +20,26 @@ pub struct ExecutorConfig { pub executor_chain_id: u64, /// Total number of EVM workers in the unified pool, shared by every execution kind. - #[arg(id = "executor.evm_workers", long = "executor-evm-workers")] - pub evm_workers: Option, + #[arg(id = "executor.evm_workers", long = "executor-evm-workers", default_value_t = DEFAULT_WORKERS)] + 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")] pub call_present_limit: Option, /// Maximum number of concurrent call-past executions. - #[arg(id = "executor.call_past_limit", long = "executor-call-past-limit")] - pub call_past_limit: Option, + #[arg(id = "executor.call_past_limit", long = "executor-call-past-limit", default_value_t = DEFAULT_KIND_LIMIT)] + pub call_past_limit: usize, /// Maximum number of concurrent inspector executions. - #[arg(id = "executor.inspector_limit", long = "executor-inspector-limit")] - pub inspector_limit: Option, + #[arg(id = "executor.inspector_limit", long = "executor-inspector-limit", default_value_t = DEFAULT_KIND_LIMIT)] + pub inspector_limit: usize, - /// Extra permits that any execution kind can borrow when its own limit is exhausted. - #[arg(id = "executor.evm_flex_quota", long = "executor-evm-flex-quota", default_value_t = 0)] - pub evm_flex_quota: 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 = DEFAULT_BUSY_THRESHOLD)] + pub evm_busy_threshold: usize, /// Should reject contract transactions and calls to accounts that are not contracts? #[arg( @@ -52,18 +57,6 @@ pub struct ExecutorConfig { #[arg(id = "executor.evm_spec", long = "executor-evm-spec", default_value = "Prague", value_parser = parse_evm_spec)] #[serde(rename = "evm_spec", with = "spec_id_serde")] pub executor_evm_spec: SpecId, - - /// Deprecated alias of `executor.call_present_limit`. - #[arg(id = "executor.call_present_evms", long = "executor-call-present-evms")] - pub call_present_evms: Option, - - /// Deprecated alias of `executor.call_past_limit`. - #[arg(id = "executor.call_past_evms", long = "executor-call-past-evms")] - pub call_past_evms: Option, - - /// Deprecated alias of `executor.inspector_limit`. - #[arg(id = "executor.inspector_evms", long = "executor-inspector-evms")] - pub inspector_evms: Option, } #[cfg(test)] @@ -71,16 +64,13 @@ impl Default for ExecutorConfig { fn default() -> Self { Self { executor_chain_id: 0, - evm_workers: None, + evm_workers: DEFAULT_WORKERS, call_present_limit: None, - call_past_limit: None, - inspector_limit: None, - evm_flex_quota: 0, + call_past_limit: DEFAULT_KIND_LIMIT, + inspector_limit: DEFAULT_KIND_LIMIT, + evm_busy_threshold: DEFAULT_BUSY_THRESHOLD, executor_reject_not_contract: true, executor_evm_spec: SpecId::PRAGUE, - call_present_evms: None, - call_past_evms: None, - inspector_evms: None, } } } @@ -109,12 +99,6 @@ fn parse_evm_spec(input: &str) -> anyhow::Result { } impl ExecutorConfig { - /// Returns whether any deprecated per-kind pool size field is set - /// (`call_present_evms`, `call_past_evms` or `inspector_evms`). - pub fn has_deprecated_pool_sizes(&self) -> bool { - self.call_present_evms.is_some() || self.call_past_evms.is_some() || self.inspector_evms.is_some() - } - /// Initializes Executor. /// /// Note: Should be called only after async runtime is initialized. diff --git a/src/eth/executor/evm_worker_pool.rs b/src/eth/executor/evm_worker_pool.rs index af8f76cbd..6e86b8c22 100644 --- a/src/eth/executor/evm_worker_pool.rs +++ b/src/eth/executor/evm_worker_pool.rs @@ -13,6 +13,7 @@ use crate::eth::executor::evm::Evm; use crate::eth::executor::evm::EvmKind; use crate::eth::executor::evm::RevmResultAndState; use crate::eth::executor::evm::types::InspectorInput; +use crate::eth::executor::pool_admission::PoolAdmission; use crate::eth::executor::types::EvmRoute; use crate::eth::executor::types::ExecutionTask; use crate::eth::executor::types::InspectionTask; @@ -22,18 +23,18 @@ use crate::eth::types::StratusError; use crate::eth::types::UnexpectedError; use crate::ext::spawn_thread; use crate::infra::tracing::warn_task_tx_closed; -use crate::utils::Permit; -use crate::utils::Semaphore; -use crate::utils::SemaphoreMetrics; /// Total capacity of the unified EVM pool task queue. const TASK_QUEUE_CAPACITY: usize = 4096; /// Default number of EVM workers in the unified pool (sum of the old per-kind pool defaults). -const DEFAULT_WORKERS: usize = 150; +pub const DEFAULT_WORKERS: usize = 150; /// Default maximum number of concurrent call-past and inspector executions. -const DEFAULT_KIND_LIMIT: usize = 50; +pub const DEFAULT_KIND_LIMIT: usize = 50; + +/// Default pool busy percentage above which per-kind limits are enforced. +pub const DEFAULT_BUSY_THRESHOLD: usize = 80; /// Effective configuration of the unified EVM pool, resolved from [`ExecutorConfig`]. #[derive(Clone, Copy, Debug)] @@ -50,116 +51,74 @@ pub struct PoolConfig { /// Maximum number of concurrent inspector executions. pub inspector_limit: usize, - /// Extra permits that any execution kind can borrow when its own limit is exhausted. - pub flex_quota: usize, + /// Pool busy percentage above which per-kind limits are enforced. + pub busy_threshold: usize, } impl PoolConfig { - /// Resolves the effective pool configuration, mapping deprecated fields to their new meaning. + /// Resolves the effective pool configuration. pub fn resolve(config: &ExecutorConfig) -> anyhow::Result { - if config.evm_workers == Some(0) { - bail!("executor.evm_workers must be greater than zero"); - } - for (field, value) in [ - ("executor.call_present_evms", config.call_present_evms), - ("executor.call_past_evms", config.call_past_evms), - ("executor.inspector_evms", config.inspector_evms), + ("executor.evm_workers", config.evm_workers), + ("executor.call_past_limit", config.call_past_limit), + ("executor.inspector_limit", config.inspector_limit), ] { - if value.is_some() { - tracing::warn!( - field, - "deprecated executor pool field; use executor.evm_workers and the per-kind limit fields instead" - ); + if value == 0 { + bail!("{field} must be greater than zero"); } } - let call_past_limit = config.call_past_limit.or(config.call_past_evms).unwrap_or(DEFAULT_KIND_LIMIT); - let inspector_limit = config.inspector_limit.or(config.inspector_evms).unwrap_or(DEFAULT_KIND_LIMIT); - - let workers = match config.evm_workers { - Some(workers) => workers, - // no old field set: default pool size - None if !config.has_deprecated_pool_sizes() => DEFAULT_WORKERS, - // at least one old field set: preserve the total capacity of the old per-kind pools - None => - config.call_present_evms.unwrap_or(DEFAULT_KIND_LIMIT) - + config.call_past_evms.unwrap_or(DEFAULT_KIND_LIMIT) - + config.inspector_evms.unwrap_or(DEFAULT_KIND_LIMIT), - }; + if let Some(0) = config.call_present_limit { + bail!("executor.call_present_limit must be greater than zero"); + } - let call_present_limit = config - .call_present_limit - .or(config.call_present_evms) - .unwrap_or_else(|| workers.saturating_sub(call_past_limit + inspector_limit)); + if config.evm_busy_threshold > 100 { + bail!("executor.evm_busy_threshold must be a percentage between 0 and 100"); + } + + // defaults to the remaining pool capacity + let call_present_limit = config.call_present_limit.unwrap_or_else(|| { + let remaining = config.evm_workers.saturating_sub(config.call_past_limit + config.inspector_limit); + if remaining == 0 { + tracing::warn!("call-present limit defaults to zero; call-present tasks will only be admitted while the pool is below the busy threshold"); + } + remaining + }); let resolved = Self { - workers, + workers: config.evm_workers, call_present_limit, - call_past_limit, - inspector_limit, - flex_quota: config.evm_flex_quota, + call_past_limit: config.call_past_limit, + inspector_limit: config.inspector_limit, + busy_threshold: config.evm_busy_threshold, }; - let limits_sum = call_present_limit + call_past_limit + inspector_limit; - if limits_sum > workers { + let limits_sum = resolved.call_present_limit + resolved.call_past_limit + resolved.inspector_limit; + if limits_sum > resolved.workers { bail!( - "executor pool kind limits ({call_present_limit} call-present + {call_past_limit} call-past + {inspector_limit} inspector = {limits_sum}) \ - exceed the total number of workers ({workers}); increase executor.evm_workers or lower the limits" + "executor pool kind limits ({} call-present + {} call-past + {} inspector = {limits_sum}) \ + exceed the total number of workers ({}); increase executor.evm_workers or lower the limits", + resolved.call_present_limit, + resolved.call_past_limit, + resolved.inspector_limit, + resolved.workers ); } tracing::info!(?resolved, "unified EVM pool configuration resolved"); Ok(resolved) } -} - -/// Per-kind concurrency limits of the unified EVM pool. -struct PoolLimits { - call_present: Semaphore, - call_past: Semaphore, - inspector: Semaphore, - flex: Semaphore, -} - -impl PoolLimits { - fn new(config: PoolConfig) -> Self { - Self { - call_present: Semaphore::with_metrics(config.call_present_limit, SemaphoreMetrics::Pool("call_present")), - call_past: Semaphore::with_metrics(config.call_past_limit, SemaphoreMetrics::Pool("call_past")), - inspector: Semaphore::with_metrics(config.inspector_limit, SemaphoreMetrics::Pool("inspector")), - flex: Semaphore::with_metrics(config.flex_quota, SemaphoreMetrics::Pool("flex")), - } - } - /// Own-limit semaphore of an execution kind. - fn own(&self, kind: EvmKind) -> &Semaphore { - 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"), - } - } - - /// Acquires a permit for the kind, borrowing from the flex quota when the own limit is exhausted. - fn acquire(&self, kind: EvmKind) -> Option { - if let Some(permit) = self.own(kind).try_acquire() { - return Some(permit); - } - - if let Some(permit) = self.flex.try_acquire() { - return Some(permit); - } - - self.own(kind).acquire_shutdown_aware() + /// In-flight task count at which relaxed admission ends and per-kind limits are enforced. + pub fn relaxed_limit(&self) -> usize { + self.workers * self.busy_threshold / 100 } } /// Manages the unified EVM pool: one shared set of workers serving every execution kind. pub struct EvmWorkerPool { tx: crossbeam_channel::Sender, - limits: Arc, + admission: Arc, } impl EvmWorkerPool { @@ -167,6 +126,7 @@ impl EvmWorkerPool { pub fn spawn(storage: Arc, config: &ExecutorConfig) -> anyhow::Result { let pool = PoolConfig::resolve(config)?; let (tx, rx) = crossbeam_channel::bounded::(TASK_QUEUE_CAPACITY); + let admission = Arc::new(PoolAdmission::new(pool)); for worker_index in 1..=pool.workers { let task_name = format!("evm-pool-{worker_index}"); @@ -185,10 +145,7 @@ impl EvmWorkerPool { } metrics::set_executor_workers_total(pool.workers as u64); - Ok(Self { - tx, - limits: Arc::new(PoolLimits::new(pool)), - }) + Ok(Self { tx, admission }) } /// Executes a call in the specified route. @@ -201,7 +158,7 @@ impl EvmWorkerPool { EvmRoute::CallPast(_) => EvmKind::CallPast, }; - let Some(permit) = self.limits.acquire(kind) else { + let Some(permit) = self.admission.acquire(kind) else { return Err(UnexpectedError::Unexpected(anyhow!("executor pool is shutting down")).into()); }; @@ -225,7 +182,7 @@ impl EvmWorkerPool { /// Executes a transaction inspection (debug_traceTransaction). pub fn inspect(&self, input: InspectorInput) -> Result { - let Some(permit) = self.limits.acquire(EvmKind::Inspect) else { + let Some(permit) = self.admission.acquire(EvmKind::Inspect) else { return Err(UnexpectedError::Unexpected(anyhow!("executor pool is shutting down")).into()); }; @@ -279,55 +236,28 @@ mod tests { assert_eq!(pool.call_present_limit, 50); assert_eq!(pool.call_past_limit, 50); assert_eq!(pool.inspector_limit, 50); - assert_eq!(pool.flex_quota, 0); + assert_eq!(pool.busy_threshold, DEFAULT_BUSY_THRESHOLD); + assert_eq!(pool.relaxed_limit(), 120); } #[test] - fn test_pool_config_maps_deprecated_fields() { + fn test_pool_config_explicit_call_present_limit() { let mut config = test_config(); - config.call_present_evms = Some(100); - config.call_past_evms = Some(20); - config.inspector_evms = Some(30); - let pool = PoolConfig::resolve(&config).unwrap(); - // total capacity preserved: 100 + 20 + 30 - assert_eq!(pool.workers, 150); - assert_eq!(pool.call_present_limit, 100); - assert_eq!(pool.call_past_limit, 20); - assert_eq!(pool.inspector_limit, 30); - } - - #[test] - fn test_pool_config_deprecated_fields_with_unset_kinds() { - let mut config = test_config(); - config.call_present_evms = Some(100); - let pool = PoolConfig::resolve(&config).unwrap(); - // unset deprecated fields keep their old defaults (50) when computing the total - assert_eq!(pool.workers, 200); - assert_eq!(pool.call_present_limit, 100); - assert_eq!(pool.call_past_limit, 50); - assert_eq!(pool.inspector_limit, 50); - } - - #[test] - fn test_pool_config_new_fields_take_precedence() { - let mut config = test_config(); - config.evm_workers = Some(200); + config.evm_workers = 200; config.call_present_limit = Some(120); - config.call_past_limit = Some(20); - config.inspector_limit = Some(30); - config.evm_flex_quota = 40; + config.call_past_limit = 20; + config.inspector_limit = 30; let pool = PoolConfig::resolve(&config).unwrap(); assert_eq!(pool.workers, 200); assert_eq!(pool.call_present_limit, 120); assert_eq!(pool.call_past_limit, 20); assert_eq!(pool.inspector_limit, 30); - assert_eq!(pool.flex_quota, 40); } #[test] fn test_pool_config_call_present_uses_remaining_capacity() { let mut config = test_config(); - config.evm_workers = Some(200); + config.evm_workers = 200; let pool = PoolConfig::resolve(&config).unwrap(); assert_eq!(pool.call_present_limit, 200 - 50 - 50); } @@ -335,17 +265,39 @@ mod tests { #[test] fn test_pool_config_rejects_limits_exceeding_workers() { let mut config = test_config(); - config.evm_workers = Some(100); + config.evm_workers = 100; config.call_present_limit = Some(60); - config.call_past_limit = Some(50); - config.inspector_limit = Some(50); + config.call_past_limit = 50; + config.inspector_limit = 50; assert!(PoolConfig::resolve(&config).is_err()); } #[test] fn test_pool_config_rejects_zero_workers() { let mut config = test_config(); - config.evm_workers = Some(0); + config.evm_workers = 0; assert!(PoolConfig::resolve(&config).is_err()); } + + #[test] + fn test_pool_config_rejects_zero_limit() { + let mut config = test_config(); + config.call_past_limit = 0; + assert!(PoolConfig::resolve(&config).is_err()); + } + + #[test] + fn test_pool_config_rejects_threshold_above_100() { + let mut config = test_config(); + config.evm_busy_threshold = 101; + assert!(PoolConfig::resolve(&config).is_err()); + } + + #[test] + fn test_pool_config_zero_threshold_is_strict() { + let mut config = test_config(); + config.evm_busy_threshold = 0; + let pool = PoolConfig::resolve(&config).unwrap(); + assert_eq!(pool.relaxed_limit(), 0); + } } diff --git a/src/eth/executor/mod.rs b/src/eth/executor/mod.rs index 6b6fa1b32..a36d48040 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; diff --git a/src/eth/executor/pool_admission.rs b/src/eth/executor/pool_admission.rs new file mode 100644 index 000000000..fef7e0a5d --- /dev/null +++ b/src/eth/executor/pool_admission.rs @@ -0,0 +1,245 @@ +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::evm::EvmKind; +use crate::eth::executor::evm_worker_pool::PoolConfig; + +/// Interval between shutdown checks while blocked waiting for a kind slot. +const SHUTDOWN_POLL_INTERVAL: Duration = Duration::from_millis(250); + +struct KindState { + inflight: usize, +} + +/// Per-kind admission state, counting in-flight tasks admitted by both the relaxed and throttled paths. +struct KindGate { + kind: EvmKind, + limit: usize, + state: 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, + state: Mutex::new(KindState { inflight: 0 }), + cvar: Condvar::new(), + } + } + + /// Relaxed admission: increments the in-flight count without checking the limit. + fn admit_relaxed(&self) { + let mut state = self.state.lock(); + state.inflight += 1; + drop(state); + metrics::inc_executor_pool_inflight(1, self.kind); + } + + /// Throttled admission: blocks until the kind's in-flight count drops below its limit. + /// Returns `None` when the application starts shutting down. + fn admit_throttled(&self) -> Option<()> { + metrics::inc_executor_pool_waiting(1, self.kind); + + let mut state = self.state.lock(); + while state.inflight >= self.limit { + if GlobalState::is_shutdown() { + drop(state); + metrics::dec_executor_pool_waiting(1, self.kind); + return None; + } + self.cvar.wait_for(&mut state, SHUTDOWN_POLL_INTERVAL); + } + + state.inflight += 1; + drop(state); + metrics::dec_executor_pool_waiting(1, self.kind); + metrics::inc_executor_pool_inflight(1, self.kind); + Some(()) + } + + /// Releases an admission slot of the kind. + fn release(&self) { + let mut state = self.state.lock(); + state.inflight -= 1; + drop(state); + 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 { + metrics::set_executor_pool_inflight_total(0); + 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 `None` on shutdown. + pub fn acquire(self: &Arc, kind: EvmKind) -> Option { + 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); + metrics::inc_executor_pool_inflight_total(1); + Some(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); + metrics::dec_executor_pool_inflight_total(1); + } +} + +/// Admission slot of a task in the unified EVM pool, released on drop. +pub struct PoolPermit { + admission: Arc, + kind: EvmKind, +} + +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 { + workers, + call_present_limit: 0, + call_past_limit, + inspector_limit: workers, + 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_returns_none_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/types/task.rs b/src/eth/executor/types/task.rs index 6e819f8bd..f5ce9553c 100644 --- a/src/eth/executor/types/task.rs +++ b/src/eth/executor/types/task.rs @@ -12,9 +12,9 @@ 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::PoolPermit; use crate::eth::executor::types::error::ExecutorError; use crate::eth::types::StratusError; -use crate::utils::Permit; #[derive(derive_new::new)] pub struct ExecutionTask { @@ -41,7 +41,7 @@ pub enum EvmRoute { pub struct PoolTask { pub span: Span, evm_kind: EvmKind, - permit: Permit, + permit: PoolPermit, task_kind: PoolTaskKind, } @@ -51,7 +51,7 @@ enum PoolTaskKind { } impl PoolTask { - pub fn call(task: ExecutionTask, kind: EvmKind, permit: Permit) -> Self { + pub fn call(task: ExecutionTask, kind: EvmKind, permit: PoolPermit) -> Self { debug_assert!(matches!(kind, EvmKind::CallPresent | EvmKind::CallPast)); Self { span: Span::current(), @@ -61,7 +61,7 @@ impl PoolTask { } } - pub fn inspect(task: InspectionTask, permit: Permit) -> Self { + pub fn inspect(task: InspectionTask, permit: PoolPermit) -> Self { Self { span: Span::current(), evm_kind: EvmKind::Inspect, diff --git a/src/utils.rs b/src/utils.rs index b05c02d3b..dc9153f86 100644 --- a/src/utils.rs +++ b/src/utils.rs @@ -1,13 +1,10 @@ use std::sync::Arc; -use std::time::Duration; use derive_more::Deref; use parking_lot::Condvar; use parking_lot::Mutex; use tokio::time::Instant; -use crate::GlobalState; - /// Amount of bytes in one GB (technically, GiB). pub const GIGABYTE: usize = 1024 * 1024 * 1024; @@ -31,9 +28,6 @@ impl Drop for DropTimer { } } -/// Interval between shutdown checks while blocked acquiring a shutdown-aware semaphore permit. -const SHUTDOWN_POLL_INTERVAL: Duration = Duration::from_millis(250); - #[derive(Deref, Default)] pub struct Semaphore { #[deref] @@ -56,9 +50,6 @@ pub enum SemaphoreMetrics { /// Legacy metrics of the local transaction warmup semaphore (unlabeled). LocalTransaction, - - /// Unified executor pool permit metrics, labeled by semaphore kind. - Pool(&'static str), } impl SemaphoreMetrics { @@ -66,7 +57,6 @@ impl SemaphoreMetrics { match *self { Self::Disabled => {} Self::LocalTransaction => stratus_metrics::inc_executor_local_transaction_semaphore_waiting(1), - Self::Pool(kind) => stratus_metrics::inc_executor_pool_permits_waiting(1, kind), } } @@ -74,7 +64,6 @@ impl SemaphoreMetrics { match *self { Self::Disabled => {} Self::LocalTransaction => stratus_metrics::dec_executor_local_transaction_semaphore_waiting(1), - Self::Pool(kind) => stratus_metrics::dec_executor_pool_permits_waiting(1, kind), } } @@ -82,7 +71,6 @@ impl SemaphoreMetrics { match *self { Self::Disabled => {} Self::LocalTransaction => stratus_metrics::inc_executor_local_transaction_permit_holders(1), - Self::Pool(kind) => stratus_metrics::inc_executor_pool_permits_held(1, kind), } } @@ -90,7 +78,6 @@ impl SemaphoreMetrics { match *self { Self::Disabled => {} Self::LocalTransaction => stratus_metrics::dec_executor_local_transaction_permit_holders(1), - Self::Pool(kind) => stratus_metrics::dec_executor_pool_permits_held(1, kind), } } } @@ -127,37 +114,6 @@ impl Semaphore { self.metrics.permit_acquired(); Permit { sem: Arc::clone(&self.sem) } } - - /// Tries to acquire a permit without blocking. - pub fn try_acquire(&self) -> Option { - let mut permits = self.permits.lock(); - if *permits == 0 { - return None; - } - *permits -= 1; - drop(permits); - self.metrics.permit_acquired(); - Some(Permit { sem: Arc::clone(&self.sem) }) - } - - /// Blocks until a permit is available or the application starts shutting down. - pub fn acquire_shutdown_aware(&self) -> Option { - self.metrics.waiting_added(); - let mut permits = self.permits.lock(); - while *permits == 0 { - if GlobalState::is_shutdown() { - drop(permits); - self.metrics.waiting_removed(); - return None; - } - self.cvar.wait_for(&mut permits, SHUTDOWN_POLL_INTERVAL); - } - *permits -= 1; - drop(permits); - self.metrics.waiting_removed(); - self.metrics.permit_acquired(); - Some(Permit { sem: Arc::clone(&self.sem) }) - } } impl Drop for Permit { diff --git a/tests/config_loader.rs b/tests/config_loader.rs index 7ea340a4b..61f30ca25 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, Some(11)); + assert_eq!(config.executor.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, Some(11)); + assert_eq!(config.executor.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, Some(11)); + assert_eq!(config.executor.call_past_limit, 11); } #[test] From 8f9f6bb522c78727f20669f23b2f67127d5a1772 Mon Sep 17 00:00:00 2001 From: gventino-cw Date: Tue, 29 Sep 2026 13:31:20 -0300 Subject: [PATCH 4/4] refac: solving comments --- crates/stratus_metrics/src/definitions.rs | 8 +- src/config/loader.rs | 10 +- src/config/mod.rs | 1 - src/config/validate/mod.rs | 11 ++ src/eth/executor/config.rs | 100 ++++++----- src/eth/executor/evm_worker_pool.rs | 202 +++++----------------- src/eth/executor/mod.rs | 20 +-- src/eth/executor/pool_admission.rs | 65 ++++--- src/eth/executor/types/mod.rs | 1 - src/eth/executor/types/task.rs | 62 +++---- src/main.rs | 2 +- src/utils.rs | 62 ++----- tests/config_loader.rs | 6 +- 13 files changed, 202 insertions(+), 348 deletions(-) diff --git a/crates/stratus_metrics/src/definitions.rs b/crates/stratus_metrics/src/definitions.rs index e040f7f93..2ea7ae036 100644 --- a/crates/stratus_metrics/src/definitions.rs +++ b/crates/stratus_metrics/src/definitions.rs @@ -118,14 +118,8 @@ metrics! { "Number of tasks waiting to be admitted to the unified executor pool, by kind." gauge executor_pool_waiting{kind}, - "Number of tasks currently executing in the unified executor pool." - gauge executor_pool_inflight_total{}, - "Number of tasks admitted through the relaxed path while the pool was below the busy threshold, by kind." - counter executor_pool_relaxed_admissions{kind}, - - "Length of the executor pool task queue." - gauge executor_pool_queue_len{} + counter executor_pool_relaxed_admissions{kind} }, group: rocks { diff --git a/src/config/loader.rs b/src/config/loader.rs index d72f746d5..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; @@ -703,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 31e2e2156..a86b1ad60 100644 --- a/src/eth/executor/config.rs +++ b/src/eth/executor/config.rs @@ -6,9 +6,6 @@ use display_json::DebugAsJson; use revm::primitives::hardfork::SpecId; use crate::eth::executor::Executor; -use crate::eth::executor::evm_worker_pool::DEFAULT_BUSY_THRESHOLD; -use crate::eth::executor::evm_worker_pool::DEFAULT_KIND_LIMIT; -use crate::eth::executor::evm_worker_pool::DEFAULT_WORKERS; use crate::eth::miner::Miner; use crate::eth::storage::StratusStorage; @@ -19,27 +16,9 @@ pub struct ExecutorConfig { #[serde(rename = "chain_id")] pub executor_chain_id: u64, - /// 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 = DEFAULT_WORKERS)] - 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")] - 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 = DEFAULT_KIND_LIMIT)] - pub call_past_limit: usize, - - /// Maximum number of concurrent inspector executions. - #[arg(id = "executor.inspector_limit", long = "executor-inspector-limit", default_value_t = DEFAULT_KIND_LIMIT)] - 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 = DEFAULT_BUSY_THRESHOLD)] - pub evm_busy_threshold: usize, + #[command(flatten)] + #[serde(flatten)] + pub pool: PoolConfig, /// Should reject contract transactions and calls to accounts that are not contracts? #[arg( @@ -59,19 +38,62 @@ pub struct ExecutorConfig { pub executor_evm_spec: SpecId, } -#[cfg(test)] -impl Default for ExecutorConfig { - fn default() -> Self { - Self { - executor_chain_id: 0, - evm_workers: DEFAULT_WORKERS, - call_present_limit: None, - call_past_limit: DEFAULT_KIND_LIMIT, - inspector_limit: DEFAULT_KIND_LIMIT, - evm_busy_threshold: DEFAULT_BUSY_THRESHOLD, - 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 } } @@ -102,11 +124,11 @@ impl ExecutorConfig { /// Initializes Executor. /// /// Note: Should be called only after async runtime is initialized. - pub fn init(&self, storage: Arc, miner: Arc) -> anyhow::Result> { + pub fn init(&self, storage: Arc, miner: Arc) -> Arc { let config = *self; tracing::info!(?config, "creating executor"); - let executor = Executor::new(storage, miner, config)?; - Ok(Arc::new(executor)) + let executor = Executor::new(storage, miner, config); + Arc::new(executor) } } diff --git a/src/eth/executor/evm_worker_pool.rs b/src/eth/executor/evm_worker_pool.rs index 6e86b8c22..453fd74ea 100644 --- a/src/eth/executor/evm_worker_pool.rs +++ b/src/eth/executor/evm_worker_pool.rs @@ -1,8 +1,6 @@ use std::sync::Arc; use alloy_rpc_types_trace::geth::GethTrace; -use anyhow::anyhow; -use anyhow::bail; use stratus_metrics as metrics; use crate::GlobalState; @@ -12,9 +10,9 @@ use crate::eth::executor::ExecutorError; use crate::eth::executor::evm::Evm; 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::pool_admission::PoolAdmission; -use crate::eth::executor::types::EvmRoute; use crate::eth::executor::types::ExecutionTask; use crate::eth::executor::types::InspectionTask; use crate::eth::executor::types::PoolTask; @@ -27,94 +25,6 @@ use crate::infra::tracing::warn_task_tx_closed; /// Total capacity of the unified EVM pool task queue. const TASK_QUEUE_CAPACITY: usize = 4096; -/// Default number of EVM workers in the unified pool (sum of the old per-kind pool defaults). -pub const DEFAULT_WORKERS: usize = 150; - -/// Default maximum number of concurrent call-past and inspector executions. -pub const DEFAULT_KIND_LIMIT: usize = 50; - -/// Default pool busy percentage above which per-kind limits are enforced. -pub const DEFAULT_BUSY_THRESHOLD: usize = 80; - -/// Effective configuration of the unified EVM pool, resolved from [`ExecutorConfig`]. -#[derive(Clone, Copy, Debug)] -pub struct PoolConfig { - /// Total number of EVM workers, shared by every execution kind. - pub workers: usize, - - /// Maximum number of concurrent call-present executions. - pub call_present_limit: usize, - - /// Maximum number of concurrent call-past executions. - pub call_past_limit: usize, - - /// Maximum number of concurrent inspector executions. - pub inspector_limit: usize, - - /// Pool busy percentage above which per-kind limits are enforced. - pub busy_threshold: usize, -} - -impl PoolConfig { - /// Resolves the effective pool configuration. - pub fn resolve(config: &ExecutorConfig) -> anyhow::Result { - for (field, value) in [ - ("executor.evm_workers", config.evm_workers), - ("executor.call_past_limit", config.call_past_limit), - ("executor.inspector_limit", config.inspector_limit), - ] { - if value == 0 { - bail!("{field} must be greater than zero"); - } - } - - if let Some(0) = config.call_present_limit { - bail!("executor.call_present_limit must be greater than zero"); - } - - if config.evm_busy_threshold > 100 { - bail!("executor.evm_busy_threshold must be a percentage between 0 and 100"); - } - - // defaults to the remaining pool capacity - let call_present_limit = config.call_present_limit.unwrap_or_else(|| { - let remaining = config.evm_workers.saturating_sub(config.call_past_limit + config.inspector_limit); - if remaining == 0 { - tracing::warn!("call-present limit defaults to zero; call-present tasks will only be admitted while the pool is below the busy threshold"); - } - remaining - }); - - let resolved = Self { - workers: config.evm_workers, - call_present_limit, - call_past_limit: config.call_past_limit, - inspector_limit: config.inspector_limit, - busy_threshold: config.evm_busy_threshold, - }; - - let limits_sum = resolved.call_present_limit + resolved.call_past_limit + resolved.inspector_limit; - if limits_sum > resolved.workers { - bail!( - "executor pool kind limits ({} call-present + {} call-past + {} inspector = {limits_sum}) \ - exceed the total number of workers ({}); increase executor.evm_workers or lower the limits", - resolved.call_present_limit, - resolved.call_past_limit, - resolved.inspector_limit, - resolved.workers - ); - } - - tracing::info!(?resolved, "unified EVM pool configuration resolved"); - Ok(resolved) - } - - /// In-flight task count at which relaxed admission ends and per-kind limits are enforced. - pub fn relaxed_limit(&self) -> usize { - self.workers * self.busy_threshold / 100 - } -} - /// Manages the unified EVM pool: one shared set of workers serving every execution kind. pub struct EvmWorkerPool { tx: crossbeam_channel::Sender, @@ -123,12 +33,12 @@ pub struct EvmWorkerPool { impl EvmWorkerPool { /// Spawns the unified EVM pool workers. - pub fn spawn(storage: Arc, config: &ExecutorConfig) -> anyhow::Result { - let pool = PoolConfig::resolve(config)?; + pub fn spawn(storage: Arc, config: &ExecutorConfig) -> Self { + 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.workers { + 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; @@ -143,33 +53,20 @@ impl EvmWorkerPool { for kind in [EvmKind::CallPresent, EvmKind::CallPast, EvmKind::Inspect] { metrics::set_executor_workers_busy(0, kind); } - metrics::set_executor_workers_total(pool.workers as u64); + metrics::set_executor_workers_total(pool.evm_workers as u64); - Ok(Self { tx, admission }) + Self { tx, admission } } /// Executes a call in the specified route. - pub fn execute(&self, route: EvmRoute) -> Result<(Output, ExecutionMetrics), StratusError> + pub fn execute(&self, input: CallExecutionInput) -> Result<(Output, ExecutionMetrics), StratusError> where Output: TryFrom, { - let kind = match &route { - EvmRoute::CallPresent(_) => EvmKind::CallPresent, - EvmRoute::CallPast(_) => EvmKind::CallPast, - }; - - let Some(permit) = self.admission.acquire(kind) else { - return Err(UnexpectedError::Unexpected(anyhow!("executor pool is shutting down")).into()); - }; - let (execution_tx, execution_rx) = oneshot::channel::>(); - let task = match route { - EvmRoute::CallPresent(input) => PoolTask::call(ExecutionTask::new(input, execution_tx), kind, permit), - EvmRoute::CallPast(input) => PoolTask::call(ExecutionTask::new(input, execution_tx), kind, permit), - }; + let task = PoolTask::call(ExecutionTask::new(input, execution_tx), &self.admission)?; self.tx.send(task)?; - metrics::set_executor_pool_queue_len(self.tx.len() as u64); match execution_rx.recv() { Ok(result) => { @@ -182,12 +79,8 @@ impl EvmWorkerPool { /// Executes a transaction inspection (debug_traceTransaction). pub fn inspect(&self, input: InspectorInput) -> Result { - let Some(permit) = self.admission.acquire(EvmKind::Inspect) else { - return Err(UnexpectedError::Unexpected(anyhow!("executor pool is shutting down")).into()); - }; - let (inspector_tx, inspector_rx) = oneshot::channel::>(); - let task = PoolTask::inspect(InspectionTask::new(input, inspector_tx), permit); + let task = PoolTask::inspect(InspectionTask::new(input, inspector_tx), &self.admission)?; let _ = self.tx.send(task); match inspector_rx.recv() { Ok(result) => result, @@ -219,85 +112,80 @@ impl EvmWorkerPool { #[cfg(test)] mod tests { + use clap::Parser; + use super::*; - fn test_config() -> ExecutorConfig { - ExecutorConfig { - executor_chain_id: 1, - ..Default::default() - } + 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 config = test_config(); - let pool = PoolConfig::resolve(&config).unwrap(); - assert_eq!(pool.workers, DEFAULT_WORKERS); - assert_eq!(pool.call_present_limit, 50); + 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.busy_threshold, DEFAULT_BUSY_THRESHOLD); + 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 mut config = test_config(); - config.evm_workers = 200; - config.call_present_limit = Some(120); - config.call_past_limit = 20; - config.inspector_limit = 30; - let pool = PoolConfig::resolve(&config).unwrap(); - assert_eq!(pool.workers, 200); - assert_eq!(pool.call_present_limit, 120); + 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 mut config = test_config(); - config.evm_workers = 200; - let pool = PoolConfig::resolve(&config).unwrap(); - assert_eq!(pool.call_present_limit, 200 - 50 - 50); + let pool = test_config(&["--executor-evm-workers", "200"]).pool; + assert_eq!(pool.call_present_limit(), 200 - 50 - 50); } #[test] - fn test_pool_config_rejects_limits_exceeding_workers() { - let mut config = test_config(); - config.evm_workers = 100; - config.call_present_limit = Some(60); - config.call_past_limit = 50; - config.inspector_limit = 50; - assert!(PoolConfig::resolve(&config).is_err()); + 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 mut config = test_config(); - config.evm_workers = 0; - assert!(PoolConfig::resolve(&config).is_err()); + 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 mut config = test_config(); - config.call_past_limit = 0; - assert!(PoolConfig::resolve(&config).is_err()); + 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 mut config = test_config(); - config.evm_busy_threshold = 101; - assert!(PoolConfig::resolve(&config).is_err()); + 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 mut config = test_config(); - config.evm_busy_threshold = 0; - let pool = PoolConfig::resolve(&config).unwrap(); + 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 a36d48040..12da4d7e0 100644 --- a/src/eth/executor/mod.rs +++ b/src/eth/executor/mod.rs @@ -40,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; @@ -53,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")] @@ -61,7 +59,6 @@ use crate::ext::OptionExt; use crate::ext::to_json_string; use crate::infra::tracing::SpanExt; use crate::utils::Semaphore; -use crate::utils::SemaphoreMetrics; // ----------------------------------------------------------------------------- // Executor @@ -85,18 +82,18 @@ pub struct Executor { } impl Executor { - pub fn new(storage: Arc, miner: Arc, config: ExecutorConfig) -> anyhow::Result { + pub fn new(storage: Arc, miner: Arc, config: ExecutorConfig) -> Self { tracing::info!(?config, "creating executor"); let reject_not_contract = config.executor_reject_not_contract; let transaction_worker = TransactionWorker::spawn(Arc::clone(&storage), Arc::clone(&miner), &config); - let evms = EvmWorkerPool::spawn(Arc::clone(&storage), &config)?; - Ok(Self { - transaction_warmup: Semaphore::with_metrics(100, SemaphoreMetrics::LocalTransaction), + let evms = EvmWorkerPool::spawn(Arc::clone(&storage), &config); + Self { + transaction_warmup: Semaphore::new(100), transaction_worker, evms, storage, reject_not_contract, - }) + } } // ------------------------------------------------------------------------- @@ -295,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 index fef7e0a5d..72639e670 100644 --- a/src/eth/executor/pool_admission.rs +++ b/src/eth/executor/pool_admission.rs @@ -8,21 +8,18 @@ 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::executor::evm_worker_pool::PoolConfig; +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); -struct KindState { - inflight: usize, -} - /// Per-kind admission state, counting in-flight tasks admitted by both the relaxed and throttled paths. struct KindGate { kind: EvmKind, limit: usize, - state: Mutex, + inflight: Mutex, cvar: Condvar, } @@ -33,46 +30,42 @@ impl KindGate { Self { kind, limit, - state: Mutex::new(KindState { inflight: 0 }), + inflight: Mutex::new(0), cvar: Condvar::new(), } } /// Relaxed admission: increments the in-flight count without checking the limit. fn admit_relaxed(&self) { - let mut state = self.state.lock(); - state.inflight += 1; - drop(state); + *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 `None` when the application starts shutting down. - fn admit_throttled(&self) -> Option<()> { + /// 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 state = self.state.lock(); - while state.inflight >= self.limit { + let mut inflight = self.inflight.lock(); + while *inflight >= self.limit { if GlobalState::is_shutdown() { - drop(state); + drop(inflight); metrics::dec_executor_pool_waiting(1, self.kind); - return None; + return Err(StateError::StratusShutdown); } - self.cvar.wait_for(&mut state, SHUTDOWN_POLL_INTERVAL); + self.cvar.wait_for(&mut inflight, SHUTDOWN_POLL_INTERVAL); } - state.inflight += 1; - drop(state); + *inflight += 1; + drop(inflight); metrics::dec_executor_pool_waiting(1, self.kind); metrics::inc_executor_pool_inflight(1, self.kind); - Some(()) + Ok(()) } /// Releases an admission slot of the kind. fn release(&self) { - let mut state = self.state.lock(); - state.inflight -= 1; - drop(state); + *self.inflight.lock() -= 1; self.cvar.notify_one(); metrics::dec_executor_pool_inflight(1, self.kind); } @@ -94,9 +87,8 @@ pub struct PoolAdmission { impl PoolAdmission { pub fn new(config: PoolConfig) -> Self { - metrics::set_executor_pool_inflight_total(0); Self { - call_present: Arc::new(KindGate::new(EvmKind::CallPresent, config.call_present_limit)), + 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), @@ -114,8 +106,8 @@ impl PoolAdmission { } /// 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 `None` on shutdown. - pub fn acquire(self: &Arc, kind: EvmKind) -> Option { + /// 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(); @@ -123,8 +115,7 @@ impl PoolAdmission { self.gate(kind).admit_throttled()?; } self.inflight_total.fetch_add(1, Ordering::Relaxed); - metrics::inc_executor_pool_inflight_total(1); - Some(PoolPermit { + Ok(PoolPermit { admission: Arc::clone(self), kind, }) @@ -134,7 +125,6 @@ impl PoolAdmission { fn release(&self, kind: EvmKind) { self.gate(kind).release(); self.inflight_total.fetch_sub(1, Ordering::Relaxed); - metrics::dec_executor_pool_inflight_total(1); } } @@ -144,6 +134,13 @@ pub struct PoolPermit { 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); @@ -158,11 +155,11 @@ mod tests { fn admission(workers: usize, busy_threshold: usize, call_past_limit: usize) -> Arc { let config = PoolConfig { - workers, - call_present_limit: 0, + evm_workers: workers, + call_present_limit: None, call_past_limit, inspector_limit: workers, - busy_threshold, + evm_busy_threshold: busy_threshold, }; Arc::new(PoolAdmission::new(config)) } @@ -229,7 +226,7 @@ mod tests { /// test would poison the other admission tests when run in parallel. Run with `--ignored`. #[test] #[ignore = "triggers process-global shutdown"] - fn test_throttled_admission_returns_none_on_shutdown() { + fn test_throttled_admission_fails_on_shutdown() { let admission = admission(10, 80, 1); let _first = admission.acquire(EvmKind::CallPast).unwrap(); diff --git a/src/eth/executor/types/mod.rs b/src/eth/executor/types/mod.rs index 562cf2ddd..b0f2b7804 100644 --- a/src/eth/executor/types/mod.rs +++ b/src/eth/executor/types/mod.rs @@ -8,7 +8,6 @@ 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::ExecutionTask; pub use task::InspectionTask; pub use task::PoolTask; diff --git a/src/eth/executor/types/task.rs b/src/eth/executor/types/task.rs index f5ce9553c..d711c20b0 100644 --- a/src/eth/executor/types/task.rs +++ b/src/eth/executor/types/task.rs @@ -1,5 +1,6 @@ use std::panic::AssertUnwindSafe; use std::panic::catch_unwind; +use std::sync::Arc; use alloy_rpc_types_trace::geth::GethTrace; use anyhow::anyhow; @@ -12,8 +13,11 @@ 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; #[derive(derive_new::new)] @@ -28,21 +32,11 @@ pub struct InspectionTask { pub response_tx: oneshot::Sender>, } -#[derive(Debug, Clone, strum::Display)] -pub enum EvmRoute { - #[strum(to_string = "call_present")] - CallPresent(CallExecutionInput), - - #[strum(to_string = "call_past")] - CallPast(CallExecutionInput), -} - /// A task for the unified EVM pool. pub struct PoolTask { - pub span: Span, - evm_kind: EvmKind, + span: Span, permit: PoolPermit, - task_kind: PoolTaskKind, + task: PoolTaskKind, } enum PoolTaskKind { @@ -51,36 +45,33 @@ enum PoolTaskKind { } impl PoolTask { - pub fn call(task: ExecutionTask, kind: EvmKind, permit: PoolPermit) -> Self { - debug_assert!(matches!(kind, EvmKind::CallPresent | EvmKind::CallPast)); - Self { + /// 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(), - evm_kind: kind, permit, - task_kind: PoolTaskKind::Call(task), - } + task: PoolTaskKind::Call(task), + }) } - pub fn inspect(task: InspectionTask, permit: PoolPermit) -> Self { - Self { + /// 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(), - evm_kind: EvmKind::Inspect, permit, - task_kind: PoolTaskKind::Inspect(task), - } + task: PoolTaskKind::Inspect(task), + }) } + /// 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, - evm_kind, - permit: _permit, - task_kind, - } = self; + let Self { span, permit, task } = self; let _enter = span.enter(); - let _busy = evm_kind.mark_executor_pool_busy(); + let _busy = permit.evm_kind().mark_executor_pool_busy(); - catch_unwind(AssertUnwindSafe(move || match task_kind { + catch_unwind(AssertUnwindSafe(move || match task { PoolTaskKind::Call(task) => task.execute(evm), PoolTaskKind::Inspect(task) => task.execute(evm), })) @@ -88,6 +79,15 @@ impl PoolTask { } } +/// 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, + } +} + impl ExecutionTask { fn execute(self, evm: &mut Evm) { if let Err(e) = self.response_tx.send(evm.execute(self.input)) { diff --git a/src/main.rs b/src/main.rs index c5f138c35..84aa5bbfe 100644 --- a/src/main.rs +++ b/src/main.rs @@ -26,7 +26,7 @@ async fn run(config: StratusConfig) -> anyhow::Result<()> { let miner = config.miner.init(Arc::clone(&storage)).await?; // Init executor - let executor = config.executor.init(Arc::clone(&storage), Arc::clone(&miner))?; + let executor = config.executor.init(Arc::clone(&storage), Arc::clone(&miner)); let (consensus, importer_runtime) = if let Some(importer_config) = &config.importer { tracing::info!(?importer_config, "creating importer"); diff --git a/src/utils.rs b/src/utils.rs index dc9153f86..6809c0414 100644 --- a/src/utils.rs +++ b/src/utils.rs @@ -3,6 +3,8 @@ use std::sync::Arc; use derive_more::Deref; use parking_lot::Condvar; use parking_lot::Mutex; +#[cfg(feature = "metrics")] +use stratus_metrics as metrics; use tokio::time::Instant; /// Amount of bytes in one GB (technically, GiB). @@ -38,48 +40,6 @@ pub struct Semaphore { pub struct SemaphoreInner { permits: Mutex, cvar: Condvar, - metrics: SemaphoreMetrics, -} - -/// Metrics recorded by a [`Semaphore`] when its permits change. -#[derive(Clone, Copy, Default)] -pub enum SemaphoreMetrics { - /// Do not record any metric. - #[default] - Disabled, - - /// Legacy metrics of the local transaction warmup semaphore (unlabeled). - LocalTransaction, -} - -impl SemaphoreMetrics { - fn waiting_added(&self) { - match *self { - Self::Disabled => {} - Self::LocalTransaction => stratus_metrics::inc_executor_local_transaction_semaphore_waiting(1), - } - } - - fn waiting_removed(&self) { - match *self { - Self::Disabled => {} - Self::LocalTransaction => stratus_metrics::dec_executor_local_transaction_semaphore_waiting(1), - } - } - - fn permit_acquired(&self) { - match *self { - Self::Disabled => {} - Self::LocalTransaction => stratus_metrics::inc_executor_local_transaction_permit_holders(1), - } - } - - fn permit_released(&self) { - match *self { - Self::Disabled => {} - Self::LocalTransaction => stratus_metrics::dec_executor_local_transaction_permit_holders(1), - } - } } pub struct Permit { @@ -88,30 +48,27 @@ pub struct Permit { impl Semaphore { pub fn new(permits: usize) -> Self { - Self::with_metrics(permits, SemaphoreMetrics::Disabled) - } - - pub fn with_metrics(permits: usize, metrics: SemaphoreMetrics) -> Self { Self { sem: Arc::new(SemaphoreInner { permits: Mutex::new(permits), cvar: Condvar::new(), - metrics, }), } } - /// Blocks until a permit is available. pub fn acquire(&self) -> Permit { - self.metrics.waiting_added(); + #[cfg(feature = "metrics")] + metrics::inc_executor_local_transaction_semaphore_waiting(1); let mut permits = self.permits.lock(); while *permits == 0 { self.cvar.wait(&mut permits); } *permits -= 1; drop(permits); - self.metrics.waiting_removed(); - self.metrics.permit_acquired(); + #[cfg(feature = "metrics")] + metrics::dec_executor_local_transaction_semaphore_waiting(1); + #[cfg(feature = "metrics")] + metrics::inc_executor_local_transaction_permit_holders(1); Permit { sem: Arc::clone(&self.sem) } } } @@ -121,7 +78,8 @@ impl Drop for Permit { let mut permits = self.sem.permits.lock(); *permits += 1; self.sem.cvar.notify_one(); - self.sem.metrics.permit_released(); + #[cfg(feature = "metrics")] + metrics::dec_executor_local_transaction_permit_holders(1); } } diff --git a/tests/config_loader.rs b/tests/config_loader.rs index 61f30ca25..66ab03bf1 100644 --- a/tests/config_loader.rs +++ b/tests/config_loader.rs @@ -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_past_limit, 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_past_limit, 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); } @@ -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_past_limit, 11); + assert_eq!(config.executor.pool.call_past_limit, 11); } #[test]