Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 10 additions & 6 deletions config/stratus.example.toml
Original file line number Diff line number Diff line change
Expand Up @@ -65,12 +65,16 @@
[executor]
# Chain ID of the network. Required.
chain_id = 2008
# Number of EVMs to execute calls against the present state.
# call_present_evms = 50
# Number of EVMs to execute calls against past states.
# call_past_evms = 50
# Number of EVMs to execute inspector calls.
# inspector_evms = 50
# Total number of EVM workers in the unified pool, shared by every execution kind.
# evm_workers = 150
# Maximum number of concurrent calls against the present state. Defaults to the remaining pool capacity.
# call_present_limit = 50
# Maximum number of concurrent calls against past states.
# call_past_limit = 50
# Maximum number of concurrent inspector calls.
# inspector_limit = 50
# Pool busy percentage above which per-kind limits are enforced; below it, tasks of any kind are admitted freely.
# evm_busy_threshold = 80
# Should reject contract transactions and calls to accounts that are not contracts?
# reject_not_contract = true
# EVM hardfork spec, e.g. Prague, Cancun, Shanghai.
Expand Down
14 changes: 13 additions & 1 deletion crates/stratus_metrics/src/definitions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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},
Comment thread
gventino-cw marked this conversation as resolved.

"Total number of EVM workers in the unified executor pool."
gauge executor_workers_total{},

"Number of tasks currently executing in the unified executor pool, by kind."
gauge executor_pool_inflight{kind},

"Number of tasks waiting to be admitted to the unified executor pool, by kind."
gauge executor_pool_waiting{kind},

"Number of tasks admitted through the relaxed path while the pool was below the busy threshold, by kind."
counter executor_pool_relaxed_admissions{kind}
},

group: rocks {
Expand Down
26 changes: 12 additions & 14 deletions src/config/loader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -272,6 +272,7 @@ mod tests {
use std::ffi::OsString;

use clap::CommandFactory;
use clap::Parser;

use crate::config::Environment;
use crate::config::StratusConfig;
Expand Down Expand Up @@ -375,9 +376,11 @@ mod tests {

[executor]
chain_id = 100
call_present_evms = 1
call_past_evms = 2
inspector_evms = 3
evm_workers = 10
call_present_limit = 1
call_past_limit = 4
inspector_limit = 5
evm_busy_threshold = 75
reject_not_contract = false
evm_spec = "Cancun"

Expand Down Expand Up @@ -572,9 +575,11 @@ mod tests {

[executor]
chain_id = 100
call_present_evms = 1
call_past_evms = 2
inspector_evms = 3
evm_workers = 10
call_present_limit = 1
call_past_limit = 4
inspector_limit = 5
evm_busy_threshold = 75
reject_not_contract = false
evm_spec = "Cancun"

Expand Down Expand Up @@ -699,14 +704,7 @@ mod tests {
fn test_empty_config_uses_defaults() {
// an empty file plus the minimal valid arguments must yield exactly the defaults
let config = load_with(&["--leader", "--executor-chain-id", "1"], "").unwrap();
let default = StratusConfig {
leader: true,
executor: crate::eth::executor::ExecutorConfig {
executor_chain_id: 1,
..Default::default()
},
..Default::default()
};
let default = StratusConfig::try_parse_from(["stratus", "--leader", "--executor-chain-id", "1"]).unwrap();
assert_eq!(serde_json::to_value(&config).unwrap(), serde_json::to_value(&default).unwrap());
}
}
1 change: 0 additions & 1 deletion src/config/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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"])]
Expand Down
11 changes: 11 additions & 0 deletions src/config/validate/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -36,6 +39,7 @@ impl Step {
const WORKFLOW: &'static [Self] = &[
Self::KafkaCompleteness,
Self::MinerExternalBlockMode,
Self::ExecutorPoolLimits,
#[cfg(feature = "dev")]
Self::GenesisFileExists,
];
Expand All @@ -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,
}
Expand All @@ -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),
}
Expand All @@ -79,6 +85,11 @@ impl Step {
Vec::new()
}

/// Returns warnings about suboptimal executor pool configurations.
fn check_executor_pool_limits(config: &StratusConfig) -> Vec<String> {
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<String> {
Expand Down
76 changes: 58 additions & 18 deletions src/eth/executor/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,14 +16,9 @@ pub struct ExecutorConfig {
#[serde(rename = "chain_id")]
pub executor_chain_id: u64,

#[arg(id = "executor.call_present_evms", long = "executor-call-present-evms", default_value_t = 50)]
pub call_present_evms: usize,

#[arg(id = "executor.call_past_evms", long = "executor-call-past-evms", default_value_t = 50)]
pub call_past_evms: usize,

#[arg(id = "executor.inspector_evms", long = "executor-inspector-evms", default_value_t = 50)]
pub inspector_evms: usize,
#[command(flatten)]
#[serde(flatten)]
pub pool: PoolConfig,

/// Should reject contract transactions and calls to accounts that are not contracts?
#[arg(
Expand All @@ -43,17 +38,62 @@ pub struct ExecutorConfig {
pub executor_evm_spec: SpecId,
}

#[cfg(test)]
impl Default for ExecutorConfig {
fn default() -> Self {
Self {
executor_chain_id: 0,
call_present_evms: 50,
call_past_evms: 50,
inspector_evms: 50,
executor_reject_not_contract: true,
executor_evm_spec: SpecId::PRAGUE,
/// Configuration of the unified EVM worker pool: one shared set of workers serving every execution kind.
#[derive(Parser, DebugAsJson, Clone, Copy, serde::Serialize)]
pub struct PoolConfig {
/// Total number of EVM workers in the unified pool, shared by every execution kind.
#[arg(id = "executor.evm_workers", long = "executor-evm-workers", default_value_t = 150, value_parser = clap::builder::RangedU64ValueParser::<usize>::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::<usize>::new().range(1..))]
#[serde(skip_serializing_if = "Option::is_none")]
pub call_present_limit: Option<usize>,

/// 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::<usize>::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::<usize>::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::<usize>::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<String> {
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
}
}

Expand Down
11 changes: 4 additions & 7 deletions src/eth/executor/evm/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@ mod session;
pub mod types;
mod util;

use std::marker::PhantomData;
use std::sync::Arc;

use alloy_consensus::transaction::TransactionInfo;
Expand Down Expand Up @@ -48,13 +47,12 @@ use crate::eth::types::StratusError;
pub type RevmResultAndState = ExecResultAndState<RevmExecResult>;

/// Implementation of EVM using [`revm`](https://crates.io/crates/revm).
pub struct Evm<Input: EvmInput> {
pub struct Evm {
evm: GeneralRevm<RevmSession>,
kind: EvmKind,
_input_type: PhantomData<Input>,
}

impl<Input: EvmInput> Evm<Input> {
impl Evm {
/// Creates a new instance of the Evm.
pub fn new(storage: Arc<StratusStorage>, config: &ExecutorConfig, kind: EvmKind) -> Self {
tracing::info!(?config, "creating revm");
Expand All @@ -65,12 +63,11 @@ impl<Input: EvmInput> Evm<Input> {
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<Input: EvmInput>(&mut self, input: Input) -> Result<(RevmResultAndState, ExecutionMetrics), StratusError> {
let metrics_context = input.metrics_context();

// configure session
Expand All @@ -94,7 +91,7 @@ impl<Input: EvmInput> Evm<Input> {
}
}

impl Evm<TransactionExecutionInput> {
impl Evm {
/// Execute a transaction using a tracer.
pub fn inspect(&mut self, input: InspectorInput) -> Result<GethTrace, StratusError> {
let InspectorInput {
Expand Down
Loading
Loading