From f3e23b61afdd7a489acbe59f64f3b0b353ca6778 Mon Sep 17 00:00:00 2001 From: "liquan.eth" Date: Thu, 3 Sep 2026 19:05:32 +0800 Subject: [PATCH 1/9] refactor: behavior-preserving cleanups and unused-API removals Deletes the second sources of truth and the dead parameters/accessors the /simplify review found, with no observable change on canonical-block paths. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01CHVpMMX9N69sUNbuKgpBVY --- AGENTS.md | 4 +- Cargo.lock | 4 +- README.md | 2 +- bin/debug-trace-server/Cargo.toml | 1 - bin/debug-trace-server/src/chain_sync.rs | 14 +-- bin/debug-trace-server/src/metrics.rs | 13 +- bin/debug-trace-server/src/server_db.rs | 9 -- bin/stateless-validator/Cargo.toml | 1 - bin/stateless-validator/src/app.rs | 8 +- bin/stateless-validator/src/chain_sync.rs | 27 +---- bin/stateless-validator/src/lib.rs | 21 +++- bin/stateless-validator/src/metrics.rs | 17 +-- .../src/{workers.rs => runner.rs} | 10 +- bin/stateless-validator/src/validator_db.rs | 25 +--- bin/stateless-validator/tests/integration.rs | 113 ++++++++---------- crates/stateless-common/Cargo.toml | 2 +- crates/stateless-common/src/metrics.rs | 29 ++++- crates/stateless-common/src/rpc_client.rs | 27 +++-- crates/stateless-core/src/chain_spec.rs | 94 +++++++++++---- crates/stateless-core/src/db.rs | 19 ++- crates/stateless-core/src/evm_database.rs | 21 ++-- crates/stateless-core/src/executor.rs | 11 +- crates/stateless-core/src/light_witness.rs | 7 -- crates/stateless-core/src/pipeline/fetcher.rs | 27 ++--- crates/stateless-core/src/pipeline/mod.rs | 18 +-- crates/stateless-core/src/pipeline/tests.rs | 10 -- crates/stateless-r2/src/fetch.rs | 45 +++++-- 27 files changed, 294 insertions(+), 285 deletions(-) rename bin/stateless-validator/src/{workers.rs => runner.rs} (97%) diff --git a/AGENTS.md b/AGENTS.md index e9b176f9..b68e91e1 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -41,7 +41,7 @@ The project uses nightly `2026-02-03` toolchain (edition 2024, rust-version 1.95 | `stateless-common` | `crates/stateless-common` | RPC client, metrics/logging utilities, witness size estimation | | `stateless-test-utils` | `crates/stateless-test-utils` | Test fixtures (blocks, witnesses, contracts) and env-var lock for integration tests | | `stateless-r2` | `crates/stateless-r2` | Shared R2 witness primitives: SigV4 signer, object-key layout, endpoint parsing, signed PUT, and the retrying witness-object GET fetcher over either the signed S3 API or an unsigned Cloudflare custom domain; consumed by mega-reth's uploaders (write) and both binaries' R2 witness sources (read) | -| `stateless-validator` | `bin/stateless-validator` | Main binary: chain sync, parallel validation workers (`app.rs` / `workers.rs` / `main.rs`) | +| `stateless-validator` | `bin/stateless-validator` | Main binary: chain sync, parallel validation workers (`app.rs` / `runner.rs` / `main.rs`) | | `debug-trace-server` | `bin/debug-trace-server` | Standalone RPC server for debug/trace methods | Additional directories: `test_data/` (integration test fixtures including genesis config), `audits/` (security audit reports). @@ -172,7 +172,7 @@ The background chain-sync prefetch routes by freshness against the last observed | `crates/stateless-db/src/{lib,tables,helpers,serialize,cache}.rs` | Shared redb tables, helpers, serialization, and `ContractCache` | | `crates/stateless-common/src/rpc_client.rs` | RPC client for blocks, witnesses, and bytecode | | `crates/stateless-common/src/metrics.rs` | RpcMethod, RpcMetrics, RpcClientConfig | -| `bin/stateless-validator/src/{main,app,workers,chain_sync,validator_db,metrics}.rs` | Thin entry, CLI/startup wiring, pipeline+reporter, fetcher/processor, DB | +| `bin/stateless-validator/src/{main,app,runner,chain_sync,validator_db,metrics}.rs` | Thin entry, CLI/startup wiring, pipeline+reporter, fetcher/processor, DB | | `bin/debug-trace-server/src/chain_sync.rs` | TraceFetcher, TraceProcessor, TraceHooks | | `bin/debug-trace-server/src/rpc_service.rs` | RPC method definitions and handlers | | `bin/debug-trace-server/src/rpc_middleware.rs` | Concurrent execution of inbound JSON-RPC batch entries | diff --git a/Cargo.lock b/Cargo.lock index 96e8d4fb..6bd286ab 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1911,7 +1911,6 @@ dependencies = [ "mega-evm", "metrics", "metrics-derive", - "metrics-exporter-prometheus", "metrics-util", "op-alloy-consensus", "op-alloy-network", @@ -5682,10 +5681,10 @@ dependencies = [ "bincode 2.0.1", "clap", "eyre", - "fastrand", "futures", "jsonrpsee", "kanal", + "metrics-exporter-prometheus", "op-alloy-network", "op-alloy-rpc-types", "reqwest", @@ -5814,7 +5813,6 @@ dependencies = [ "jsonrpsee", "jsonrpsee-types", "metrics", - "metrics-exporter-prometheus", "op-alloy-rpc-types", "redb", "revm", diff --git a/README.md b/README.md index 65f36168..04792fcc 100644 --- a/README.md +++ b/README.md @@ -347,7 +347,7 @@ The pipeline is configured via `PipelineConfig` and customized through trait imp | `crates/stateless-common/src/metrics.rs` | `RpcMethod`, `RpcMetrics`, `RpcClientConfig` | | `crates/stateless-common/src/witness_size.rs` | `WitnessSizeBreakdown` + `estimate_witness_size` for RPC and trace-server metrics | | `crates/stateless-test-utils/src/fixtures.rs` | `TestFixtures` loader (blocks, SALT/MPT witnesses, contracts, genesis) | -| `bin/stateless-validator/src/{main,app,workers,chain_sync,validator_db,metrics}.rs` | Thin entry, CLI/startup wiring, pipeline+reporter, fetcher/processor, DB | +| `bin/stateless-validator/src/{main,app,runner,chain_sync,validator_db,metrics}.rs` | Thin entry, CLI/startup wiring, pipeline+reporter, fetcher/processor, DB | | `bin/debug-trace-server/src/chain_sync.rs` | `TraceFetcher`, `TraceProcessor`, `TraceHooks` | | `bin/debug-trace-server/src/rpc_service.rs` | RPC method definitions and handlers | | `bin/debug-trace-server/src/data_provider.rs` | Block data fetching with single-flight coalescing | diff --git a/bin/debug-trace-server/Cargo.toml b/bin/debug-trace-server/Cargo.toml index a1e01edf..91ee8386 100644 --- a/bin/debug-trace-server/Cargo.toml +++ b/bin/debug-trace-server/Cargo.toml @@ -46,7 +46,6 @@ jsonrpsee.workspace = true libc.workspace = true metrics.workspace = true metrics-derive.workspace = true -metrics-exporter-prometheus.workspace = true pin-project-lite.workspace = true quick_cache.workspace = true rayon.workspace = true diff --git a/bin/debug-trace-server/src/chain_sync.rs b/bin/debug-trace-server/src/chain_sync.rs index 1f929dbe..661776c9 100644 --- a/bin/debug-trace-server/src/chain_sync.rs +++ b/bin/debug-trace-server/src/chain_sync.rs @@ -157,12 +157,7 @@ impl BlockFetcher for TraceFetcher { async fn latest_block_meta(&self) -> Result { let header = self.rpc_client.get_header(BlockId::Number(BlockNumberOrTag::Latest), false).await; - Ok(BlockMeta { - block_number: header.number, - block_hash: header.hash, - post_state_root: header.state_root, - post_withdrawals_root: header.withdrawals_root.unwrap_or_default(), - }) + Ok(BlockMeta::from_header(&header)) } } @@ -207,12 +202,7 @@ impl BlockProcessor for TraceProcessor { &self, (block, witness): Self::Input, ) -> std::result::Result { - let meta = BlockMeta { - block_number: block.header.number, - block_hash: block.header.hash, - post_state_root: block.header.state_root, - post_withdrawals_root: block.header.withdrawals_root.unwrap_or_default(), - }; + let meta = BlockMeta::from_header(&block.header); Ok(TraceProcessedBlock { block, witness, meta }) } } diff --git a/bin/debug-trace-server/src/metrics.rs b/bin/debug-trace-server/src/metrics.rs index c12191b4..4fc647a8 100644 --- a/bin/debug-trace-server/src/metrics.rs +++ b/bin/debug-trace-server/src/metrics.rs @@ -10,10 +10,9 @@ use std::net::SocketAddr; use eyre::Result; use metrics::{Counter, Gauge, Histogram, counter, gauge, histogram}; use metrics_derive::Metrics; -use metrics_exporter_prometheus::{Matcher, PrometheusBuilder}; pub use stateless_common::{ DEFAULT_METRICS_PORT, - metrics::{BYTE_BUCKETS, REORG_DEPTH_BUCKETS}, + metrics::{BYTE_BUCKETS, REORG_DEPTH_BUCKETS, install_prometheus_exporter}, }; /// Prefix for timed RPC method aliases. @@ -1026,15 +1025,7 @@ const BUCKET_SPECS: &[(&str, &[f64])] = &[ /// Initializes the Prometheus metrics exporter. pub fn init_metrics(addr: SocketAddr) -> Result<()> { - let builder = BUCKET_SPECS.iter().fold(PrometheusBuilder::new(), |b, &(name, buckets)| { - b.set_buckets_for_metric(Matcher::Full(name.to_owned()), buckets) - .expect("valid bucket config") - }); - - builder - .with_http_listener(addr) - .install() - .map_err(|e| eyre::eyre!("Failed to install metrics exporter: {}", e))?; + install_prometheus_exporter(addr, BUCKET_SPECS)?; // Pre-register all metrics pre_register_all_metrics(); diff --git a/bin/debug-trace-server/src/server_db.rs b/bin/debug-trace-server/src/server_db.rs index deb7d39f..41ac890e 100644 --- a/bin/debug-trace-server/src/server_db.rs +++ b/bin/debug-trace-server/src/server_db.rs @@ -308,15 +308,6 @@ pub(crate) mod test_support { pub tip_reads: AtomicUsize, } - impl ContractStore for StubBlockStore { - fn get_contracts(&self, _: &[B256]) -> StoreResult<(HashMap, Vec)> { - Ok((HashMap::default(), vec![])) - } - fn add_contracts(&self, _: &[(B256, Bytecode)]) -> StoreResult<()> { - Ok(()) - } - } - impl ChainStore for StubBlockStore { fn get_canonical_tip(&self) -> StoreResult> { self.tip_reads.fetch_add(1, Ordering::Relaxed); diff --git a/bin/stateless-validator/Cargo.toml b/bin/stateless-validator/Cargo.toml index 51beca81..b2a53ccb 100644 --- a/bin/stateless-validator/Cargo.toml +++ b/bin/stateless-validator/Cargo.toml @@ -36,7 +36,6 @@ stateless-r2 = { path = "../../crates/stateless-r2" } clap = { workspace = true, features = ["env"] } eyre.workspace = true metrics.workspace = true -metrics-exporter-prometheus.workspace = true redb.workspace = true serde_json.workspace = true thiserror.workspace = true diff --git a/bin/stateless-validator/src/app.rs b/bin/stateless-validator/src/app.rs index 90a2ac3e..d5347f58 100644 --- a/bin/stateless-validator/src/app.rs +++ b/bin/stateless-validator/src/app.rs @@ -15,7 +15,7 @@ use stateless_core::{ChainStore, ContractStore, chain_spec::ChainSpec, db::Block use stateless_db::ContractCache; use tracing::{info, warn}; -use crate::{metrics, r2_witness::R2WitnessClient, validator_db::ValidatorDB, workers}; +use crate::{metrics, r2_witness::R2WitnessClient, runner, validator_db::ValidatorDB}; /// Where the validator sources witnesses from. #[derive(ValueEnum, Clone, Debug, PartialEq, Eq, Default)] @@ -291,7 +291,7 @@ pub struct CommandLineArgs { /// /// Parses CLI args, initializes tracing and metrics, constructs the RPC client and /// validator DB, loads or initializes the chain spec + anchor, then hands off to -/// [`workers::run_with_signals`]. +/// [`runner::run_with_signals`]. pub async fn run() -> Result<()> { let args = CommandLineArgs::parse(); let _log_guard = args.log.init_tracing()?; @@ -434,13 +434,13 @@ pub async fn run() -> Result<()> { info!(end_block = end, "Validating up to end block, then stopping"); } - let result = workers::run_with_signals( + let result = runner::run_with_signals( client, r2_witness, validator_db, contract_cache, chain_spec, - args.report_validation_endpoint, + args.report_validation_endpoint.is_some(), pipeline_config, ) .await; diff --git a/bin/stateless-validator/src/chain_sync.rs b/bin/stateless-validator/src/chain_sync.rs index 82cf2529..b91436fb 100644 --- a/bin/stateless-validator/src/chain_sync.rs +++ b/bin/stateless-validator/src/chain_sync.rs @@ -4,7 +4,7 @@ //! [`ValidatorHooks`] (metrics integration) for the shared pipeline in //! [`stateless_core::pipeline::run_pipeline`]. -use std::{collections::HashSet, sync::Arc}; +use std::sync::Arc; use alloy_primitives::{B256, BlockHash, BlockNumber}; use alloy_rpc_types_eth::{Block, BlockId}; @@ -15,7 +15,7 @@ use salt::SaltWitness; use stateless_common::{CodeFetchError, RpcClient}; use stateless_core::{ chain_spec::ChainSpec, - data_types::iter_code_hashes, + data_types::collect_code_hashes, db::BlockMeta, executor::validate_block, pipeline::{BlockFetcher, BlockProcessor, ErrorAction, PipelineHooks, ProcessedBlock}, @@ -36,7 +36,6 @@ pub struct ValidatorFetcher { pub rpc_client: Arc, /// `Some` ⇒ fetch witnesses directly from R2; `None` ⇒ RPC. pub r2_witness: Option>, - pub on_remote_height: fn(u64), } impl BlockFetcher for ValidatorFetcher { @@ -63,7 +62,7 @@ impl BlockFetcher for ValidatorFetcher { async fn latest_block_number(&self) -> Result { let n = self.rpc_client.get_latest_block_number().await; - (self.on_remote_height)(n); + metrics::set_remote_chain_height(n); Ok(n) } @@ -73,12 +72,7 @@ impl BlockFetcher for ValidatorFetcher { async fn latest_block_meta(&self) -> Result { let header = self.rpc_client.get_header(BlockId::latest(), false).await; - Ok(BlockMeta { - block_number: header.number, - block_hash: header.hash, - post_state_root: header.state_root, - post_withdrawals_root: header.withdrawals_root.unwrap_or_default(), - }) + Ok(BlockMeta::from_header(&header)) } } @@ -186,8 +180,7 @@ impl BlockProcessor for ValidatorProcessor { // Resolve contract codes via the shared three-tier chain. Memory/disk hits // are trusted; the RPC tier verifies each bytecode's hash inside `get_codes`. - let codehashes: Vec = - iter_code_hashes(&task.salt_witness.kvs).collect::>().into_iter().collect(); + let codehashes: Vec = collect_code_hashes(&task.salt_witness.kvs); let (mut contracts, missing_contracts) = self .contract_cache .get(&codehashes) @@ -313,15 +306,7 @@ mod tests { use stateless_core::pipeline::ProcessedBlock; use super::*; - - fn make_block_meta(num: u64) -> BlockMeta { - BlockMeta { - block_number: num, - block_hash: BlockHash::from([num as u8; 32]), - post_state_root: B256::from([(num.wrapping_add(100)) as u8; 32]), - post_withdrawals_root: B256::from([(num.wrapping_add(200)) as u8; 32]), - } - } + use crate::test_support::make_block_meta; #[test] fn test_verify_continuity_success() { diff --git a/bin/stateless-validator/src/lib.rs b/bin/stateless-validator/src/lib.rs index a5a8ea2b..22d2a11b 100644 --- a/bin/stateless-validator/src/lib.rs +++ b/bin/stateless-validator/src/lib.rs @@ -7,13 +7,30 @@ pub(crate) mod app; pub(crate) mod chain_sync; pub(crate) mod metrics; pub(crate) mod r2_witness; +pub(crate) mod runner; pub(crate) mod validator_db; -pub(crate) mod workers; pub use app::{ CommandLineArgs, VALIDATOR_DB_FILENAME, WitnessSource, load_or_create_chain_spec, run, }; pub use chain_sync::{ValidationTask, ValidatorFetcher, ValidatorHooks, ValidatorProcessor}; pub use r2_witness::{R2WitnessClient, R2WitnessError}; +pub use runner::run_with_signals; pub use validator_db::ValidatorDB; -pub use workers::run_with_signals; + +/// Fixtures shared by the unit tests of several modules. +#[cfg(test)] +pub(crate) mod test_support { + use alloy_primitives::{B256, BlockHash}; + use stateless_core::db::BlockMeta; + + /// A deterministic `BlockMeta` derived from `num` alone. + pub(crate) fn make_block_meta(num: u64) -> BlockMeta { + BlockMeta { + block_number: num, + block_hash: BlockHash::from([num as u8; 32]), + post_state_root: B256::from([(num.wrapping_add(100)) as u8; 32]), + post_withdrawals_root: B256::from([(num.wrapping_add(200)) as u8; 32]), + } + } +} diff --git a/bin/stateless-validator/src/metrics.rs b/bin/stateless-validator/src/metrics.rs index 6be652b4..2e68ac6d 100644 --- a/bin/stateless-validator/src/metrics.rs +++ b/bin/stateless-validator/src/metrics.rs @@ -10,10 +10,12 @@ use std::{ use eyre::Result; use metrics::{counter, describe_counter, describe_gauge, describe_histogram, gauge, histogram}; -use metrics_exporter_prometheus::{Matcher, PrometheusBuilder}; pub use stateless_common::{ DEFAULT_METRICS_PORT, WitnessSizeBreakdown, - metrics::{BYTE_BUCKETS, REORG_DEPTH_BUCKETS, RpcAttemptOutcome, RpcMethod, RpcMetrics}, + metrics::{ + BYTE_BUCKETS, REORG_DEPTH_BUCKETS, RpcAttemptOutcome, RpcMethod, RpcMetrics, + install_prometheus_exporter, + }, }; use tracing::info; @@ -127,15 +129,7 @@ const BUCKET_SPECS: &[(&str, &[f64])] = &[ /// Initialize the Prometheus metrics exporter at the given address. pub fn init_metrics(addr: SocketAddr) -> Result<()> { - let builder = BUCKET_SPECS.iter().fold(PrometheusBuilder::new(), |b, &(name, buckets)| { - b.set_buckets_for_metric(Matcher::Full(name.to_owned()), buckets) - .expect("valid bucket config") - }); - - builder - .with_http_listener(addr) - .install() - .map_err(|e| eyre::eyre!("Failed to install Prometheus exporter: {}", e))?; + install_prometheus_exporter(addr, BUCKET_SPECS)?; register_metric_descriptions(); init_rpc_method_counters(); @@ -228,7 +222,6 @@ fn init_rpc_method_counters() { RpcMethod::EthGetBlock, RpcMethod::EthBlockNumber, RpcMethod::EthGetHeader, - RpcMethod::EthGetTransactionByHash, RpcMethod::MegaGetBlockWitness, RpcMethod::MegaSetValidatedBlocks, ]; diff --git a/bin/stateless-validator/src/workers.rs b/bin/stateless-validator/src/runner.rs similarity index 97% rename from bin/stateless-validator/src/workers.rs rename to bin/stateless-validator/src/runner.rs index 8abd80ec..30b06a7e 100644 --- a/bin/stateless-validator/src/workers.rs +++ b/bin/stateless-validator/src/runner.rs @@ -21,7 +21,6 @@ use tracing::{debug, error, info, warn}; use crate::{ chain_sync::{ValidatorFetcher, ValidatorHooks, ValidatorProcessor}, - metrics, r2_witness::R2WitnessClient, validator_db::ValidatorDB, }; @@ -44,10 +43,9 @@ pub async fn run_with_signals( validator_db: Arc, contract_cache: Arc, chain_spec: Arc, - report_validation_endpoint: Option, + report_validation: bool, pipeline_config: PipelineConfig, ) -> Result<()> { - let report_validation = report_validation_endpoint.is_some(); let config = Arc::new(pipeline_config); let is_slice_run = config.sync_target.is_some(); info!( @@ -62,11 +60,7 @@ pub async fn run_with_signals( let mut sigterm = signal::unix::signal(signal::unix::SignalKind::terminate()) .map_err(|e| eyre::eyre!("Failed to register SIGTERM handler: {e}"))?; - let fetcher = Arc::new(ValidatorFetcher { - rpc_client: client.clone(), - r2_witness, - on_remote_height: metrics::set_remote_chain_height, - }); + let fetcher = Arc::new(ValidatorFetcher { rpc_client: client.clone(), r2_witness }); let processor = Arc::new(ValidatorProcessor { chain_spec, contract_cache, rpc_client: client.clone() }); let hooks = Arc::new(ValidatorHooks); diff --git a/bin/stateless-validator/src/validator_db.rs b/bin/stateless-validator/src/validator_db.rs index 0545b68b..66e2780d 100644 --- a/bin/stateless-validator/src/validator_db.rs +++ b/bin/stateless-validator/src/validator_db.rs @@ -58,18 +58,6 @@ impl ValidatorDB { Ok(Self { database, max_chain_length }) } - - #[cfg(test)] - fn set_anchor_block(&self, tip: &BlockMeta) -> StoreResult<()> { - use stateless_db::block_meta_to_tuple; - let write_txn = self.database.begin_write().store_err()?; - { - let mut table = write_txn.open_table(ANCHOR_BLOCK).store_err()?; - table.insert("anchor", block_meta_to_tuple(tip)).store_err()?; - } - write_txn.commit().store_err()?; - Ok(()) - } } impl ContractStore for ValidatorDB { @@ -158,6 +146,7 @@ mod tests { use stateless_db::ContractCache; use super::*; + use crate::test_support::make_block_meta; fn temp_store() -> (tempfile::TempDir, ValidatorDB) { let dir = tempfile::tempdir().unwrap(); @@ -165,15 +154,6 @@ mod tests { (dir, store) } - fn make_block_meta(number: u64) -> BlockMeta { - BlockMeta { - block_number: number, - block_hash: BlockHash::from([number as u8; 32]), - post_state_root: B256::from([(number + 100) as u8; 32]), - post_withdrawals_root: B256::from([(number + 200) as u8; 32]), - } - } - #[test] fn test_anchor_block_roundtrip() { let (_dir, store) = temp_store(); @@ -186,7 +166,8 @@ mod tests { post_state_root: B256::from([2u8; 32]), post_withdrawals_root: B256::from([3u8; 32]), }; - store.set_anchor_block(&tip).unwrap(); + // The production anchor write path — no test-only shortcut around the helper layer. + store.reset_to_anchor(&tip).unwrap(); let loaded = ChainStore::get_anchor(&store).unwrap().unwrap(); assert_eq!(loaded, tip); diff --git a/bin/stateless-validator/tests/integration.rs b/bin/stateless-validator/tests/integration.rs index c69347f1..69592abc 100644 --- a/bin/stateless-validator/tests/integration.rs +++ b/bin/stateless-validator/tests/integration.rs @@ -332,12 +332,47 @@ impl MockServerState { reject_reports: Arc::default(), } } + + /// Fixture block for a `0x…` hex block number, or the RPC error the handlers return + /// for unknown blocks. Shared by the by-number block and header handlers. + fn block_by_number_hex( + &self, + hex_number: &str, + ) -> Result<&Block, ErrorObject<'static>> { + let block_number = parse_hex_u64(hex_number); + self.fixtures + .block_numbers + .get(&block_number) + .and_then(|hash| self.fixtures.blocks.get(hash)) + .ok_or_else(|| { + make_rpc_error( + CALL_EXECUTION_FAILED_CODE, + format!("Block {block_number} not found"), + ) + }) + } + + /// Fixture block for a block hash, or the RPC error the handlers return for unknown + /// blocks. Shared by the by-hash block and header handlers. + fn block_by_hash( + &self, + hash: B256, + ) -> Result<&Block, ErrorObject<'static>> { + self.fixtures.blocks.get(&BlockHash::from(hash.0)).ok_or_else(|| { + make_rpc_error(CALL_EXECUTION_FAILED_CODE, format!("Block {hash} not found")) + }) + } } fn make_rpc_error(code: i32, msg: String) -> ErrorObject<'static> { ErrorObject::owned(code, msg, None::<()>) } +/// Invalid-params RPC error for a failed `params.parse()`. +fn invalid_params(e: impl std::fmt::Display) -> ErrorObject<'static> { + make_rpc_error(INVALID_PARAMS_CODE, format!("Invalid params: {e}")) +} + fn shape_block( block: &Block, full_block: bool, @@ -383,39 +418,17 @@ async fn setup_mock_rpc_server( serve_with_config(cfg, state, |module| { module .register_method("eth_getBlockByNumber", |params, ctx, _| { - let (hex_number, full_block): (String, bool) = params.parse().map_err(|e| { - make_rpc_error(INVALID_PARAMS_CODE, format!("Invalid params: {e}")) - })?; - let block_number = parse_hex_u64(&hex_number); - - let block = ctx - .fixtures - .block_numbers - .get(&block_number) - .and_then(|hash| ctx.fixtures.blocks.get(hash)) - .ok_or_else(|| { - make_rpc_error( - CALL_EXECUTION_FAILED_CODE, - format!("Block {block_number} not found"), - ) - })?; - + let (hex_number, full_block): (String, bool) = + params.parse().map_err(invalid_params)?; + let block = ctx.block_by_number_hex(&hex_number)?; Ok::<_, ErrorObject<'static>>(shape_block(block, full_block)) }) .unwrap(); module .register_method("eth_getBlockByHash", |params, ctx, _| { - let (hash, full_block): (B256, bool) = params.parse().map_err(|e| { - make_rpc_error(INVALID_PARAMS_CODE, format!("Invalid params: {e}")) - })?; - - let block_hash = BlockHash::from(hash.0); - let block = ctx.fixtures.blocks.get(&block_hash).ok_or_else(|| { - make_rpc_error(CALL_EXECUTION_FAILED_CODE, format!("Block {hash} not found")) - })?; - - Ok::<_, ErrorObject<'static>>(shape_block(block, full_block)) + let (hash, full_block): (B256, bool) = params.parse().map_err(invalid_params)?; + Ok::<_, ErrorObject<'static>>(shape_block(ctx.block_by_hash(hash)?, full_block)) }) .unwrap(); @@ -428,45 +441,21 @@ async fn setup_mock_rpc_server( module .register_method("eth_getHeaderByNumber", |params, ctx, _| { - let (hex_number,): (String,) = params.parse().unwrap(); - let block_number = parse_hex_u64(&hex_number); - - let block = ctx - .fixtures - .block_numbers - .get(&block_number) - .and_then(|hash| ctx.fixtures.blocks.get(hash)) - .ok_or_else(|| { - make_rpc_error( - CALL_EXECUTION_FAILED_CODE, - format!("Block {block_number} not found"), - ) - })?; - - Ok::<_, ErrorObject<'static>>(block.header.clone()) + let (hex_number,): (String,) = params.parse().map_err(invalid_params)?; + Ok::<_, ErrorObject<'static>>(ctx.block_by_number_hex(&hex_number)?.header.clone()) }) .unwrap(); module .register_method("eth_getHeaderByHash", |params, ctx, _| { - let (hash,): (B256,) = params.parse().map_err(|e| { - make_rpc_error(INVALID_PARAMS_CODE, format!("Invalid params: {e}")) - })?; - - let block_hash = BlockHash::from(hash.0); - let block = ctx.fixtures.blocks.get(&block_hash).ok_or_else(|| { - make_rpc_error(CALL_EXECUTION_FAILED_CODE, format!("Block {hash} not found")) - })?; - - Ok::<_, ErrorObject<'static>>(block.header.clone()) + let (hash,): (B256,) = params.parse().map_err(invalid_params)?; + Ok::<_, ErrorObject<'static>>(ctx.block_by_hash(hash)?.header.clone()) }) .unwrap(); module .register_method("eth_getCodeByHash", |params, ctx, _| { - let (hash,): (B256,) = params.parse().map_err(|e| { - make_rpc_error(INVALID_PARAMS_CODE, format!("Invalid params: {e}")) - })?; + let (hash,): (B256,) = params.parse().map_err(invalid_params)?; let code = ctx.fixtures.contracts.get(&hash).cloned().unwrap_or_default(); Ok::<_, ErrorObject<'static>>(code.original_bytes()) @@ -475,9 +464,7 @@ async fn setup_mock_rpc_server( module .register_method("mega_getBlockWitness", |params, ctx, _| { - let (keys,): (WitnessRequestKeys,) = params.parse().map_err(|e| { - make_rpc_error(INVALID_PARAMS_CODE, format!("Invalid params: {e}")) - })?; + let (keys,): (WitnessRequestKeys,) = params.parse().map_err(invalid_params)?; let block_hash = BlockHash::from(keys.block_hash.0); let salt_witness = @@ -562,11 +549,7 @@ async fn integration_test() { let config = Arc::new(cfg); let shutdown = CancellationToken::new(); - let fetcher = Arc::new(ValidatorFetcher { - rpc_client: client.clone(), - r2_witness: None, - on_remote_height: |_| {}, - }); + let fetcher = Arc::new(ValidatorFetcher { rpc_client: client.clone(), r2_witness: None }); let processor = Arc::new(ValidatorProcessor { chain_spec, contract_cache, rpc_client: client }); let hooks = Arc::new(ValidatorHooks); @@ -634,7 +617,7 @@ async fn run_end_block_slice( Arc::clone(&validator_db), contract_cache, chain_spec, - Some(url.clone()), + true, cfg, ) .await; diff --git a/crates/stateless-common/Cargo.toml b/crates/stateless-common/Cargo.toml index 16a6130e..60bf5eb7 100644 --- a/crates/stateless-common/Cargo.toml +++ b/crates/stateless-common/Cargo.toml @@ -35,8 +35,8 @@ base64.workspace = true bincode.workspace = true clap.workspace = true eyre.workspace = true -fastrand = { workspace = true, features = ["std"] } futures.workspace = true +metrics-exporter-prometheus.workspace = true # Also enables gzip/brotli on the shared reqwest 0.12 client (Cargo feature unification) for witness/data fetches. reqwest = { workspace = true, features = ["gzip", "brotli"] } rolling-file.workspace = true diff --git a/crates/stateless-common/src/metrics.rs b/crates/stateless-common/src/metrics.rs index 20d61a99..89f9aaae 100644 --- a/crates/stateless-common/src/metrics.rs +++ b/crates/stateless-common/src/metrics.rs @@ -1,10 +1,35 @@ //! RPC metrics types shared by both binaries. //! -//! Provides [`RpcMethod`] for identifying RPC calls and [`RpcMetrics`] as a -//! callback trait for tracking RPC performance. +//! Provides [`RpcMethod`] for identifying RPC calls, [`RpcMetrics`] as a +//! callback trait for tracking RPC performance, and the shared Prometheus +//! exporter installer ([`install_prometheus_exporter`]). + +use std::net::SocketAddr; + +use metrics_exporter_prometheus::{Matcher, PrometheusBuilder}; use crate::witness_size::WitnessSizeBreakdown; +/// Installs the Prometheus exporter with an HTTP listener on `addr`, applying the given +/// per-metric histogram buckets (`(metric_name, buckets)` pairs) before install. +/// +/// Shared by both binaries; each keeps its own metric names, descriptions, and +/// pre-registration after this returns. +pub fn install_prometheus_exporter( + addr: SocketAddr, + bucket_specs: &[(&str, &[f64])], +) -> eyre::Result<()> { + let builder = bucket_specs.iter().fold(PrometheusBuilder::new(), |b, &(name, buckets)| { + b.set_buckets_for_metric(Matcher::Full(name.to_owned()), buckets) + .expect("valid bucket config") + }); + + builder + .with_http_listener(addr) + .install() + .map_err(|e| eyre::eyre!("Failed to install Prometheus exporter: {e}")) +} + /// Byte-size histogram buckets: 1 KB, 10 KB, 50 KB, 200 KB, 1 MB, 5 MB, 20 MB. pub const BYTE_BUCKETS: &[f64] = &[1_024.0, 10_240.0, 51_200.0, 204_800.0, 1_048_576.0, 5_242_880.0, 20_971_520.0]; diff --git a/crates/stateless-common/src/rpc_client.rs b/crates/stateless-common/src/rpc_client.rs index dfe2f38e..c6e272c8 100644 --- a/crates/stateless-common/src/rpc_client.rs +++ b/crates/stateless-common/src/rpc_client.rs @@ -52,6 +52,7 @@ use revm::state::Bytecode; use salt::SaltWitness; use serde::{Deserialize, Serialize}; use stateless_core::{LightWitness, withdrawals::MptWitness}; +use stateless_r2::fetch::{BackoffSchedule, RetryPacing}; use tokio::sync::Semaphore; use tracing::{instrument, trace, warn}; @@ -63,9 +64,9 @@ use crate::{ /// Exponential-backoff policy used by [`RpcClient`]'s round-level retry loop. /// -/// `initial` is the first sleep duration; each round doubles it up to `max`. -/// The loop itself lives in [`round_robin_with_backoff`]; this type only describes -/// the sleep schedule. +/// `initial` is the first sleep duration; each round doubles it up to `max`. The loop +/// itself lives in [`round_robin_with_backoff`]; this type only describes the sleep +/// schedule, which it steps through [`Self::schedule`]. #[derive(Debug, Clone)] pub struct BackoffPolicy { /// First retry sleep. Each subsequent retry doubles up to `max`. @@ -79,6 +80,15 @@ impl BackoffPolicy { pub const fn new(initial: Duration, max: Duration) -> Self { Self { initial, max } } + + /// Starts executing the schedule from `initial`. + /// + /// The schedule itself lives in `stateless-r2` next to [`RetryPacing`], the pacing pair + /// its GET loop steps: that crate must stay free of upward dependencies, so it is the + /// only home both retry loops can reach. + pub fn schedule(&self) -> BackoffSchedule { + RetryPacing { initial: self.initial, max: self.max }.schedule() + } } /// Error returned by the `_with_deadline` RPC methods when a caller-supplied @@ -1151,9 +1161,7 @@ where const WARN_AT_ROUND: u32 = 3; let n = providers.len(); - let max_backoff_ms = policy.max.as_millis() as u64; - let initial_backoff_ms = policy.initial.as_millis() as u64; - let mut round_backoff_ms = initial_backoff_ms; + let mut backoff = policy.schedule(); let mut round = 0u32; let call_start = Instant::now(); // Records the logical-call deadline give-up (once) and builds the typed error. Called @@ -1393,11 +1401,7 @@ where // `last_err` is always `Some` here: `n >= 1` is enforced by the `RpcClient` // constructor and we only reach this point after `n` iterations that each set it. let last_err = last_err.expect("last_err set when every provider failed this round"); - let jitter_ms = fastrand::u64(0..=round_backoff_ms / 2); - // `.max(1)` prevents a hot-spin loop if a caller constructs a zero-backoff policy - // (`BackoffPolicy::new(Duration::ZERO, Duration::ZERO)`): the computed sleep would - // otherwise be `0` and the retry loop would busy-wait on every round. - let mut sleep_ms = (round_backoff_ms + jitter_ms).min(max_backoff_ms).max(1); + let mut sleep_ms = backoff.next_sleep_ms(); // Clamp the sleep so it doesn't overshoot the caller's deadline — if no time is // left we bail immediately rather than sleeping past the deadline and then bailing. if let Some(d) = deadline { @@ -1419,7 +1423,6 @@ where "All providers failed this round, backing off", ); tokio::time::sleep(std::time::Duration::from_millis(sleep_ms)).await; - round_backoff_ms = (round_backoff_ms * 2).min(max_backoff_ms); round += 1; } } diff --git a/crates/stateless-core/src/chain_spec.rs b/crates/stateless-core/src/chain_spec.rs index adba1fc3..54d5c9d3 100644 --- a/crates/stateless-core/src/chain_spec.rs +++ b/crates/stateless-core/src/chain_spec.rs @@ -15,9 +15,6 @@ use mega_evm::{ use reth_ethereum_forks::ChainHardforks; use reth_optimism_chainspec::OpChainSpec; -/// Default blob gas price update fraction for Cancun (from EIP-4844) -pub const BLOB_GASPRICE_UPDATE_FRACTION: u64 = 3338477; - /// Chain specification for the Optimism network. /// /// Defines when various Ethereum and Optimism hardforks are activated. @@ -74,10 +71,8 @@ impl ChainSpec { /// Ordering rules: /// - [`OpChainSpec`] already yields Optimism/Ethereum hardforks in the correct order, so they /// do not require reordering. - /// - MegaETH hardforks are extracted from the genesis `extra_fields` and explicitly ordered to - /// match the canonical sequence defined by [`mega_mainnet_hardforks()`]. Any remaining, - /// unknown MegaETH hardforks are preserved and appended after the known ones so nothing is - /// dropped. + /// - MegaETH hardforks are extracted from the genesis `extra_fields`; + /// [`MegaethGenesisHardforks::into_vec`] yields them in canonical activation order. /// - The MegaETH set is then merged with the Optimism/Ethereum set to build a single /// [`ChainHardforks`] that drives fork activation. /// @@ -120,7 +115,7 @@ impl ChainSpec { ); } - let mut megaeth_hardforks = megaeth_hardforks.into_vec(); + let megaeth_hardforks = megaeth_hardforks.into_vec(); // Rex5 SequencerRegistry bootstrap, required iff `rex5Time` is scheduled. Parsed from // the same flat schema mega-reth uses (`rex5InitialSequencer` / `rex5InitialAdmin` as @@ -167,20 +162,9 @@ impl ChainSpec { .map(|(f, b)| (dyn_clone::clone_box(f), b)) .collect(); - let hardfork_order = mega_mainnet_hardforks(); - let mut all_hardforks = Vec::with_capacity(op_hardforks.len() + megaeth_hardforks.len()); - for (order, _) in hardfork_order.forks_iter() { - if let Some(mega_hardfork_index) = - megaeth_hardforks.iter().position(|(hardfork, _)| **hardfork == *order) - { - all_hardforks.push(megaeth_hardforks.remove(mega_hardfork_index)); - } - } - - // append the remaining unknown hardforks to ensure we don't filter any out - all_hardforks.append(&mut megaeth_hardforks); - - // we merge megaeth_hardforks with op_hardforks + // `into_vec` yields the MegaETH hardforks already in canonical activation order, + // so the merge is a straight concatenation. + let mut all_hardforks = megaeth_hardforks; all_hardforks.append(&mut op_hardforks); Self { @@ -227,6 +211,13 @@ impl MegaethGenesisHardforks { } /// Convert the MegaETH genesis hardforks into a vector of hardforks and their conditions. + /// + /// The literal below is the single source of the canonical MegaETH activation order — + /// [`ChainSpec::from_genesis`] merges it as-is, so new hardforks must be inserted at + /// their activation position and appended to the expected list in + /// `test_mega_hardforks_iterate_in_activation_order`, which pins the order (fork + /// selection by timestamp walks `forks_iter()` in insertion order, so a wrong order + /// here means wrong hardfork params — consensus divergence — with lookups still green). pub fn into_vec(self) -> Vec<(Box, ForkCondition)> { vec![ (MegaHardfork::MiniRex.boxed(), self.mini_rex_time.map(ForkCondition::Timestamp)), @@ -303,7 +294,13 @@ impl MegaethGenesisSequencerRegistryRex6Config { } } -/// Build a fresh `ChainHardforks` describing MegaETH's canonical hardfork sequence. +/// MegaETH's canonical hardfork ladder: every fork the pinned mega-evm knows about, at a +/// placeholder activation. +/// +/// Not an activation schedule — the real activations come from genesis +/// ([`MegaethGenesisHardforks::into_vec`], which also fixes their order). This is the +/// membership list `mainnet_genesis_schedules_every_canonical_hardfork` checks the shipped +/// mainnet genesis against, so a fork the executor supports cannot go unscheduled unnoticed. pub fn mega_mainnet_hardforks() -> ChainHardforks { ChainHardforks::new(vec![ (MegaHardfork::MiniRex.boxed(), ForkCondition::Timestamp(0)), @@ -382,6 +379,57 @@ mod tests { assert_eq!(spec.hardforks.fork(MegaHardfork::MiniRex), ForkCondition::Timestamp(3)); } + /// Fork selection by timestamp walks `forks_iter()` in insertion order, so the + /// `into_vec` literal IS the consensus activation order. This pins the relative order + /// of every MegaETH fork end-to-end through `from_genesis`; a new hardfork must be + /// inserted at its activation position in `into_vec` and appended to the expected + /// list here. + #[test] + fn test_mega_hardforks_iterate_in_activation_order() { + let mut genesis = Genesis::default(); + for (field, ts) in [ + ("miniRexTime", 1u64), + ("miniRex1Time", 2), + ("miniRex2Time", 3), + ("rexTime", 4), + ("rex1Time", 5), + ("rex2Time", 6), + ("rex3Time", 7), + ("rex4Time", 8), + ] { + genesis.config.extra_fields.insert_value(field.to_string(), ts).unwrap(); + } + schedule_valid_rex5(&mut genesis, 9); + genesis.config.extra_fields.insert_value("rex6Time".to_string(), 10).unwrap(); + genesis.config.extra_fields.insert_value("rex6MinRotationDelay".to_string(), 7200).unwrap(); + let spec = ChainSpec::from_genesis(genesis); + + let expected = [ + MegaHardfork::MiniRex, + MegaHardfork::MiniRex1, + MegaHardfork::MiniRex2, + MegaHardfork::Rex, + MegaHardfork::Rex1, + MegaHardfork::Rex2, + MegaHardfork::Rex3, + MegaHardfork::Rex4, + MegaHardfork::Rex5, + MegaHardfork::Rex6, + ]; + let mega_order: Vec<&str> = spec + .hardforks + .forks_iter() + .map(|(hardfork, _)| hardfork.name()) + .filter(|name| expected.iter().any(|mega| mega.name() == *name)) + .collect(); + assert_eq!( + mega_order, + expected.iter().map(|mega| mega.name()).collect::>(), + "MegaETH forks must iterate in canonical activation order — fix the \ + `into_vec` literal, not this list" + ); + } + #[test] fn test_extract_from_json() { let genesis_info = r#" diff --git a/crates/stateless-core/src/db.rs b/crates/stateless-core/src/db.rs index f75f2626..8df3769f 100644 --- a/crates/stateless-core/src/db.rs +++ b/crates/stateless-core/src/db.rs @@ -29,6 +29,21 @@ pub struct BlockMeta { pub post_withdrawals_root: B256, } +impl BlockMeta { + /// Projects an RPC header into the meta of the block it seals — a header's roots are that + /// block's post-state. A missing `withdrawals_root` defaults to zero, the tip-observation + /// policy both binaries use; callers that must instead *reject* such headers (e.g. anchor + /// initialization from an operator-supplied hash) build the meta explicitly. + pub fn from_header(header: &alloy_rpc_types_eth::Header) -> Self { + Self { + block_number: header.number, + block_hash: header.hash, + post_state_root: header.state_root, + post_withdrawals_root: header.withdrawals_root.unwrap_or_default(), + } + } +} + /// Errors returned by persistence trait methods. /// /// This is the single typed error at the library/binary boundary: every @@ -107,7 +122,9 @@ pub trait ContractStore: Send + Sync { /// pipeline's [`ReorgResolver`](crate::pipeline::ReorgResolver) seam, which each scenario supplies. /// History-owning stores additionally implement /// [`DivergenceLookups`](crate::pipeline::DivergenceLookups) so the pipeline can bisect them. -pub trait ChainStore: ContractStore { +/// Deliberately independent of [`ContractStore`]: a chain-cursor store (e.g. an embedder whose +/// bytecode integrity is enforced at ingest) need not stub contract persistence. +pub trait ChainStore: Send + Sync { fn get_canonical_tip(&self) -> StoreResult>; fn get_anchor(&self) -> StoreResult>; fn advance_chain(&self, blocks: &[BlockMeta]) -> StoreResult<()>; diff --git a/crates/stateless-core/src/evm_database.rs b/crates/stateless-core/src/evm_database.rs index abdf2f3e..5927c3d7 100644 --- a/crates/stateless-core/src/evm_database.rs +++ b/crates/stateless-core/src/evm_database.rs @@ -5,6 +5,7 @@ //! validation. use std::{ + collections::BTreeMap, format, string::{String, ToString}, vec::Vec, @@ -225,8 +226,16 @@ impl WitnessExternalEnv { salt_witness: &SaltWitness, block_number: BlockNumber, ) -> Result { - let bucket_capacities = salt_witness - .kvs + Self::from_metadata_kvs(&salt_witness.kvs, block_number) + } + + /// Shared constructor body: scans the metadata key range of a witness's `kvs` map and + /// collects the bucket capacities (both witness types expose the same map layout). + fn from_metadata_kvs( + kvs: &BTreeMap>, + block_number: BlockNumber, + ) -> Result { + let bucket_capacities = kvs .range(METADATA_KEYS_RANGE) .map(|(key, value)| Self::parse_metadata_entry(key, value)) .collect::, _>>()?; @@ -260,13 +269,7 @@ impl WitnessExternalEnv { light_witness: &LightWitness, block_number: BlockNumber, ) -> Result { - let bucket_capacities = light_witness - .kvs - .range(METADATA_KEYS_RANGE) - .map(|(key, value)| Self::parse_metadata_entry(key, value)) - .collect::, _>>()?; - - Ok(Self { block_number, bucket_capacities }) + Self::from_metadata_kvs(&light_witness.kvs, block_number) } } diff --git a/crates/stateless-core/src/executor.rs b/crates/stateless-core/src/executor.rs index 71a6e9a3..96ef17cd 100644 --- a/crates/stateless-core/src/executor.rs +++ b/crates/stateless-core/src/executor.rs @@ -61,7 +61,7 @@ use revm::{ State, states::{BundleAccount, StateBuilder, bundle_state::BundleRetention}, }, - primitives::{B256, KECCAK_EMPTY, U256}, + primitives::{B256, KECCAK_EMPTY, U256, eip4844::BLOB_BASE_FEE_UPDATE_FRACTION_CANCUN}, state::Bytecode, }; use salt::{EphemeralSaltState, SaltValue, SaltWitness, StateRoot, StateUpdates, Witness}; @@ -69,7 +69,7 @@ use thiserror::Error; use tracing::debug; use crate::{ - chain_spec::{BLOB_GASPRICE_UPDATE_FRACTION, ChainSpec}, + chain_spec::ChainSpec, data_types::{Account, PlainKey, PlainValue}, evm_database::{WitnessDatabase, WitnessDatabaseError, WitnessExternalEnv}, withdrawals::{self, ADDRESS_L2_TO_L1_MESSAGE_PASSER, MptWitness}, @@ -263,10 +263,6 @@ impl ValidationOptions { /// - Chain configuration with appropriate spec ID for the block number /// - Block environment with gas limits, timestamps, and fee parameters /// - Blob gas pricing if excess blob gas is present in the header -/// -/// Creates an EVM environment from a block header and chain specification. -/// -/// This function sets up the configuration and block environment needed for EVM execution. pub fn create_evm_env( header: &alloy_consensus::Header, chain_spec: &ChainSpec, @@ -286,7 +282,8 @@ pub fn create_evm_env( }; if let Some(excess_blob_gas) = header.excess_blob_gas { - block_env.set_blob_excess_gas_and_price(excess_blob_gas, BLOB_GASPRICE_UPDATE_FRACTION); + block_env + .set_blob_excess_gas_and_price(excess_blob_gas, BLOB_BASE_FEE_UPDATE_FRACTION_CANCUN); } EvmEnv::new(cfg_env, block_env) diff --git a/crates/stateless-core/src/light_witness.rs b/crates/stateless-core/src/light_witness.rs index 035e200f..973a03e6 100644 --- a/crates/stateless-core/src/light_witness.rs +++ b/crates/stateless-core/src/light_witness.rs @@ -293,13 +293,6 @@ impl StateReader for LightWitnessExecutor { } } -impl LightWitnessExecutor { - /// Get the underlying kvs map - pub fn kvs(&self) -> &BTreeMap> { - &self.light_witness.kvs - } -} - #[cfg(test)] mod tests { // `std` is the `alloc` alias in no_std builds, where the prelude carries diff --git a/crates/stateless-core/src/pipeline/fetcher.rs b/crates/stateless-core/src/pipeline/fetcher.rs index b66393b7..6d0582ce 100644 --- a/crates/stateless-core/src/pipeline/fetcher.rs +++ b/crates/stateless-core/src/pipeline/fetcher.rs @@ -14,17 +14,17 @@ use tracing::{Instrument, debug, error, info, info_span, warn}; use crate::pipeline::{config::PipelineConfig, traits::BlockFetcher}; /// Invariant: every block in `[base_block, next_block)` is in exactly one of -/// `in_flight_blocks`, `sent`, or `failed`. All mutations go through the methods below. +/// `task_to_block` (in flight), `sent`, or `failed`. All mutations go through the methods below. struct FetcherState { /// Lowest block not yet sent downstream. base_block: u64, /// Next block to spawn fresh. next_block: u64, tasks: JoinSet<(u64, Result)>, - /// Task id → block, for panic recovery (`JoinError` only carries the id). + /// Task id → block, for panic recovery (`JoinError` only carries the id) and + /// block-in-flight lookups (`recover_gaps` scans the values; the set is bounded by + /// `max_in_flight`, so a linear scan is negligible next to the awaited fetches). task_to_block: HashMap, - /// Mirror of `task_to_block.values()` for O(1) block-in-flight lookup. - in_flight_blocks: HashSet, /// Successful blocks, waiting for `base_block` to catch up. sent: HashSet, /// Blocks awaiting retry. The RPC client retries transient errors internally, so failures @@ -42,7 +42,6 @@ impl FetcherState { next_block: start_block, tasks: JoinSet::new(), task_to_block: HashMap::new(), - in_flight_blocks: HashSet::new(), sent: HashSet::new(), failed: HashSet::new(), } @@ -75,7 +74,6 @@ impl FetcherState { let handle = self.tasks.spawn(async move { (bn, fetcher.fetch(bn).await) }.instrument(span)); self.task_to_block.insert(handle.id(), bn); - self.in_flight_blocks.insert(bn); } fn spawn_next(&mut self, fetcher: &Arc) { @@ -93,21 +91,18 @@ impl FetcherState { fn on_success(&mut self, id: Id, bn: u64) { self.task_to_block.remove(&id); - self.in_flight_blocks.remove(&bn); self.sent.insert(bn); } fn on_failure(&mut self, id: Id, bn: u64) { self.task_to_block.remove(&id); - self.in_flight_blocks.remove(&bn); self.failed.insert(bn); } /// Re-enqueues the panicked task's block. Returns `None` if the id is unknown - /// (shouldn't happen — would leak the block from `in_flight_blocks`). + /// (shouldn't happen — `recover_gaps` would pick the block up). fn on_panic(&mut self, id: Id) -> Option { let bn = self.task_to_block.remove(&id)?; - self.in_flight_blocks.remove(&bn); self.failed.insert(bn); Some(bn) } @@ -129,8 +124,8 @@ impl FetcherState { let mut recovered = 0; for bn in self.base_block..self.next_block { if !self.sent.contains(&bn) && - !self.in_flight_blocks.contains(&bn) && - !self.failed.contains(&bn) + !self.failed.contains(&bn) && + !self.task_to_block.values().any(|&in_flight| in_flight == bn) { self.failed.insert(bn); recovered += 1; @@ -344,12 +339,11 @@ mod tests { async fn on_failure_re_enqueues_for_immediate_retry() { let mut state = FetcherState::::new(100); let id = fresh_task_id(&mut state.tasks, 100).await; - state.in_flight_blocks.insert(100); state.task_to_block.insert(id, 100); state.on_failure(id, 100); assert!(state.failed.contains(&100)); - assert!(!state.in_flight_blocks.contains(&100)); + assert!(!state.task_to_block.contains_key(&id)); assert_eq!(state.pop_failed(), Some(100)); assert!(!state.failed.contains(&100)); @@ -382,7 +376,8 @@ mod tests { let mut state = FetcherState::::new(100); state.next_block = 103; state.sent.insert(100); - state.in_flight_blocks.insert(101); + let id = fresh_task_id(&mut state.tasks, 101).await; + state.task_to_block.insert(id, 101); state.failed.insert(102); assert_eq!(state.recover_gaps(), 0); @@ -395,13 +390,11 @@ mod tests { async fn on_panic_re_enqueues_known_task() { let mut state = FetcherState::::new(100); let id = fresh_task_id(&mut state.tasks, 100).await; - state.in_flight_blocks.insert(100); state.task_to_block.insert(id, 100); let bn = state.on_panic(id); assert_eq!(bn, Some(100)); assert!(state.failed.contains(&100)); - assert!(!state.in_flight_blocks.contains(&100)); assert!(!state.task_to_block.contains_key(&id)); } diff --git a/crates/stateless-core/src/pipeline/mod.rs b/crates/stateless-core/src/pipeline/mod.rs index d7d67c24..1603713d 100644 --- a/crates/stateless-core/src/pipeline/mod.rs +++ b/crates/stateless-core/src/pipeline/mod.rs @@ -105,7 +105,7 @@ where fetcher_shutdown.cancel(); await_handles(fetcher_handle, worker_handles, config.await_handles_timeout).await; - let transient_reason: String = match outcome { + match outcome { Ok(PipelineOutcome::Shutdown) => { info!("Shutting down"); return Ok(()); @@ -129,27 +129,16 @@ where } Ok(PipelineOutcome::Retry(msg)) => { warn!(reason = %msg, "Cycle ended with retry signal"); - msg } Err(e) => { // Any `Err` at this level is unexpected (every intentional transient/fatal // case returns `Ok(PipelineOutcome::..)`). Log and fall into the same // stale-detect + sleep + continue recovery path as `Retry`. error!(error = %e, "Cycle ended with unexpected error"); - e.to_string() } - }; + } - if handle_transient_restart( - transient_reason, - &*fetcher, - &*store, - &*hooks, - &config, - &shutdown, - ) - .await? - { + if handle_transient_restart(&*fetcher, &*store, &*hooks, &config, &shutdown).await? { return Ok(()); } } @@ -162,7 +151,6 @@ where /// returns `Ok(false)` so the outer loop `continue`s. Propagates `Err` only for store / /// hook failures the caller can't meaningfully recover from. async fn handle_transient_restart( - _reason: String, fetcher: &F, store: &S, hooks: &H, diff --git a/crates/stateless-core/src/pipeline/tests.rs b/crates/stateless-core/src/pipeline/tests.rs index a1511c5e..deaeca59 100644 --- a/crates/stateless-core/src/pipeline/tests.rs +++ b/crates/stateless-core/src/pipeline/tests.rs @@ -2,7 +2,6 @@ use std::{sync::Arc, time::Duration}; use alloy_primitives::{B256, BlockHash, BlockNumber, map::HashMap}; use eyre::{Result, anyhow}; -use revm::state::Bytecode; use tokio_util::sync::CancellationToken; use super::{ @@ -82,15 +81,6 @@ impl MockStore { } } -impl crate::ContractStore for MockStore { - fn get_contracts(&self, _: &[B256]) -> StoreResult<(HashMap, Vec)> { - Ok((HashMap::default(), vec![])) - } - fn add_contracts(&self, _: &[(B256, Bytecode)]) -> StoreResult<()> { - Ok(()) - } -} - impl ChainStore for MockStore { fn get_canonical_tip(&self) -> StoreResult> { Ok(self.chain.lock().unwrap().values().next_back().cloned()) diff --git a/crates/stateless-r2/src/fetch.rs b/crates/stateless-r2/src/fetch.rs index 55494887..8f7c90ff 100644 --- a/crates/stateless-r2/src/fetch.rs +++ b/crates/stateless-r2/src/fetch.rs @@ -190,6 +190,42 @@ pub struct RetryPacing { pub max: Duration, } +impl RetryPacing { + /// Starts executing this pacing, from `initial`. + pub fn schedule(&self) -> BackoffSchedule { + BackoffSchedule { + current_ms: self.initial.as_millis() as u64, + max_ms: self.max.as_millis() as u64, + } + } +} + +/// Stepping state for a [`RetryPacing`] — the arithmetic every retry loop that paces this +/// way needs, in one place. +/// +/// Owns the three invariants those loops rely on: up to 50% random jitter per sleep (keeps +/// parallel clients from retrying in lockstep through a shared outage), the `max` cap, and +/// a 1 ms floor so a zero-duration pacing cannot turn a retry loop into a busy-loop. +/// +/// Deadline handling stays with the caller: whether an overrunning sleep is clamped or +/// gives up is policy, not arithmetic. +#[derive(Clone, Copy, Debug)] +pub struct BackoffSchedule { + current_ms: u64, + max_ms: u64, +} + +impl BackoffSchedule { + /// Returns the next sleep in milliseconds (jittered, capped, floored at 1 ms) and + /// advances the doubling state. + pub fn next_sleep_ms(&mut self) -> u64 { + let jitter_ms = fastrand::u64(0..=self.current_ms / 2); + let sleep_ms = (self.current_ms + jitter_ms).min(self.max_ms).max(1); + self.current_ms = (self.current_ms * 2).min(self.max_ms); + sleep_ms + } +} + /// A successfully fetched object body plus the time the fetch spent queued on the /// concurrency cap. /// @@ -791,8 +827,7 @@ impl R2ObjectFetcher { on_retry: impl Fn(), ) -> Result { let key = keys::block_object_key(number, hash); - let max_backoff_ms = self.pacing.max.as_millis() as u64; - let mut backoff_ms = self.pacing.initial.as_millis() as u64; + let mut backoff = self.pacing.schedule(); let mut attempt = 0usize; let mut queue_wait = Duration::ZERO; @@ -838,10 +873,7 @@ impl R2ObjectFetcher { if !e.is_retryable() || attempt >= max_attempts { return Err(e); } - // Jittered doubling; `.max(1)` keeps a zero-duration policy from - // busy-looping. - let jitter_ms = fastrand::u64(0..=backoff_ms / 2); - let sleep_ms = (backoff_ms + jitter_ms).min(max_backoff_ms).max(1); + let sleep_ms = backoff.next_sleep_ms(); if deadline .is_some_and(|d| Instant::now() + Duration::from_millis(sleep_ms) >= d) { @@ -853,7 +885,6 @@ impl R2ObjectFetcher { "R2 witness GET failed, backing off", ); tokio::time::sleep(Duration::from_millis(sleep_ms)).await; - backoff_ms = (backoff_ms * 2).min(max_backoff_ms); } } } From 46cc35290f0ec3d2d2912bf38a11d03e43cbd398 Mon Sep 17 00:00:00 2001 From: "liquan.eth" Date: Thu, 3 Sep 2026 19:07:23 +0800 Subject: [PATCH 2/9] refactor(core): remove the unused EIP-3155 trace writer Deletes the `writer` parameter threaded through `validate_block`, `validate_block_deriving_updates`, `replay_block` and `verify_and_replay`, and with it EIP-3155 trace output. A feature deletion, not a cleanup. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01CHVpMMX9N69sUNbuKgpBVY --- bin/stateless-validator/src/chain_sync.rs | 1 - crates/stateless-core/src/executor.rs | 78 ++++------------------- 2 files changed, 11 insertions(+), 68 deletions(-) diff --git a/bin/stateless-validator/src/chain_sync.rs b/bin/stateless-validator/src/chain_sync.rs index b91436fb..7010e2b0 100644 --- a/bin/stateless-validator/src/chain_sync.rs +++ b/bin/stateless-validator/src/chain_sync.rs @@ -234,7 +234,6 @@ impl BlockProcessor for ValidatorProcessor { task.salt_witness, task.mpt_witness, &contracts, - None, ) }) .await diff --git a/crates/stateless-core/src/executor.rs b/crates/stateless-core/src/executor.rs index 96ef17cd..604d7c6e 100644 --- a/crates/stateless-core/src/executor.rs +++ b/crates/stateless-core/src/executor.rs @@ -30,9 +30,9 @@ //! The module integrates with the Salt witness system for state reconstruction //! and uses Revm for transaction execution. -use std::{boxed::Box, collections::BTreeMap, fmt::Debug, vec::Vec}; #[cfg(feature = "std")] -use std::{io::Write, time::Instant}; +use std::time::Instant; +use std::{boxed::Box, collections::BTreeMap, fmt::Debug, vec::Vec}; use alloy_consensus::{TxReceipt, proofs::calculate_receipt_root, transaction::Recovered}; use alloy_eips::eip2718::Encodable2718; @@ -52,8 +52,6 @@ use mega_evm::{ }; use op_alloy_consensus::OpTxEnvelope; use op_alloy_rpc_types::Transaction as OpTransaction; -#[cfg(feature = "std")] -use revm::inspector::inspectors::TracerEip3155; use revm::{ DatabaseRef, context::{BlockEnv, CfgEnv}, @@ -464,7 +462,6 @@ pub fn replay_block( block: &B, db: &DB, env_oracle: ENV, - #[cfg(feature = "std")] trace_writer: Option>, ) -> Result<(HashMap, BlockExecutionOutput), ValidationError> where B: BlockInput, @@ -489,28 +486,9 @@ where let BlockExecutionEnv { evm_env, executor_factory, ctx: execution_context } = create_block_execution_env(chain_spec, header, env_oracle); - // Plain execution path, shared by the non-tracer std branch and the no_std build. - // Extracted as a closure so the body lives in one place — any future change to the - // non-tracer path only needs to be made here. - let run_plain = |state: &mut _, ctx, env| { - let executor = executor_factory.create_executor(state, ctx, env); - execute_transactions(executor, block.txs_recovered()) - }; - - #[cfg(feature = "std")] - let (receipts_root, logs_bloom, gas_used) = if let Some(writer) = trace_writer { - let executor = executor_factory.create_executor_with_inspector( - &mut state, - execution_context, - evm_env, - TracerEip3155::new(writer), - ); - execute_transactions(executor, block.txs_recovered())? - } else { - run_plain(&mut state, execution_context, evm_env)? - }; - #[cfg(not(feature = "std"))] - let (receipts_root, logs_bloom, gas_used) = run_plain(&mut state, execution_context, evm_env)?; + let executor = executor_factory.create_executor(&mut state, execution_context, evm_env); + let (receipts_root, logs_bloom, gas_used) = + execute_transactions(executor, block.txs_recovered())?; // Merge transitions into bundle_state state.merge_transitions(BundleRetention::PlainState); @@ -739,7 +717,6 @@ fn verify_and_replay( block: &B, salt_witness: SaltWitness, contracts: &HashMap, - #[cfg(feature = "std")] writer: Option>, ) -> Result { let header = block.consensus_header(); @@ -757,14 +734,7 @@ fn verify_and_replay( // Replay block transactions let ((accounts, output), block_replay_time) = timed(|| { let witness_db = WitnessDatabase { header, witness: &witness, contracts }; - replay_block( - chain_spec, - block, - &witness_db, - ext_env, - #[cfg(feature = "std")] - writer, - ) + replay_block(chain_spec, block, &witness_db, ext_env) })?; let stats = ValidationStats { @@ -794,8 +764,6 @@ fn verify_and_replay( /// * `salt_witness` - The salt witness data needed for state validation /// * `mpt_witness` - The MPT witness data for withdrawal verification /// * `contracts` - Contract bytecode cache for transaction execution -/// * `writer` - Optional writer for EIP-3155 trace output. When provided, enables step-by-step EVM -/// execution tracing in EIP-3155 format. /// /// # Returns /// @@ -807,7 +775,6 @@ pub fn validate_block( salt_witness: SaltWitness, mpt_witness: MptWitness, contracts: &HashMap, - #[cfg(feature = "std")] writer: Option>, ) -> Result { // A block carrying only transaction hashes can't be replayed — fail fast before paying // the witness proof verification. `replay_block` re-checks for direct callers. @@ -817,14 +784,8 @@ pub fn validate_block( let header = block.consensus_header(); // Verify the witness proof and replay the block's transactions over it - let VerifiedReplay { witness, accounts, output, mut stats } = verify_and_replay( - chain_spec, - block, - salt_witness, - contracts, - #[cfg(feature = "std")] - writer, - )?; + let VerifiedReplay { witness, accounts, output, mut stats } = + verify_and_replay(chain_spec, block, salt_witness, contracts)?; // Extract and hash storage updates (only changed values) let withdrawal_storage = withdrawal_storage(&accounts); @@ -884,7 +845,6 @@ pub fn validate_block_deriving_updates( mpt_witness: MptWitness, contracts: &HashMap, options: ValidationOptions, - #[cfg(feature = "std")] writer: Option>, ) -> Result<(StateUpdates, ValidationStats), ValidationError> { // A block carrying only transaction hashes can't be replayed — fail fast before paying // the witness proof verification. `replay_block` re-checks for direct callers. @@ -913,14 +873,8 @@ pub fn validate_block_deriving_updates( } // Verify the witness proof and replay the block's transactions over it - let VerifiedReplay { witness, accounts, output, mut stats } = verify_and_replay( - chain_spec, - block, - salt_witness, - contracts, - #[cfg(feature = "std")] - writer, - )?; + let VerifiedReplay { witness, accounts, output, mut stats } = + verify_and_replay(chain_spec, block, salt_witness, contracts)?; // Check the header's claims (withdrawals root, receipts root, logs bloom, gas used) // before the more expensive state-update derivation. @@ -962,8 +916,6 @@ mod tests { fx.mpt_witness(&hash), &fx.contracts, options, - #[cfg(feature = "std")] - None, ) } @@ -974,15 +926,7 @@ mod tests { salt_witness: SaltWitness, hash: B256, ) -> Result { - validate_block( - &chain_spec(), - block, - salt_witness, - fx.mpt_witness(&hash), - &fx.contracts, - #[cfg(feature = "std")] - None, - ) + validate_block(&chain_spec(), block, salt_witness, fx.mpt_witness(&hash), &fx.contracts) } /// Empty-witness external env for tests that only exercise environment assembly. From e62e3e14aef654c4e6b9c3a8fe6f4ca3eb19f2a9 Mon Sep 17 00:00:00 2001 From: "liquan.eth" Date: Thu, 3 Sep 2026 19:10:35 +0800 Subject: [PATCH 3/9] perf: encode each transaction once in verify_block_integrity, read gas_used from the executor Both are definitional rewrites of security-critical derivations, argued in the PR body: `trie_hash()` is `keccak256(encoded_2718())`, and mega-evm's `gas_used` is the last receipt's cumulative gas. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01CHVpMMX9N69sUNbuKgpBVY --- crates/stateless-common/src/rpc_client.rs | 17 +++++--- crates/stateless-core/src/executor.rs | 47 ++++++++++++++++++++++- 2 files changed, 57 insertions(+), 7 deletions(-) diff --git a/crates/stateless-common/src/rpc_client.rs b/crates/stateless-common/src/rpc_client.rs index c6e272c8..8971c0d1 100644 --- a/crates/stateless-common/src/rpc_client.rs +++ b/crates/stateless-common/src/rpc_client.rs @@ -1599,13 +1599,19 @@ fn verify_block_integrity(block: &Block) -> Result<()> { // Verify transaction hashes and transactions root if let BlockTransactions::Full(ref transactions) = block.transactions { + // Encode each envelope exactly once: keccak of the encoding is the tx hash, and the + // same bytes feed the ordered trie for the transactions-root check. + let mut encoded_txs: Vec> = Vec::with_capacity(transactions.len()); for tx in transactions { - let tx_envelope = tx.inner.clone().into_inner(); + let tx_envelope = tx.inner.inner.inner(); + let mut encoded = Vec::with_capacity(tx_envelope.encode_2718_len()); + tx_envelope.encode_2718(&mut encoded); + let computed_hash = alloy_primitives::keccak256(&encoded); ensure!( - tx_envelope.trie_hash() == *tx_envelope.hash(), + computed_hash == *tx_envelope.hash(), "Transaction hash mismatch: expected {:?}, computed {:?}", tx_envelope.hash(), - tx_envelope.trie_hash() + computed_hash ); let recovered = tx_envelope @@ -1618,10 +1624,11 @@ fn verify_block_integrity(block: &Block) -> Result<()> { tx.from(), recovered ); + encoded_txs.push(encoded); } - let computed_tx_root = ordered_trie_root_with_encoder(transactions, |tx, buf| { - tx.inner.clone().into_inner().encode_2718(buf) + let computed_tx_root = ordered_trie_root_with_encoder(&encoded_txs, |tx_bytes, buf| { + buf.extend_from_slice(tx_bytes) }); ensure!( computed_tx_root == block.header.transactions_root, diff --git a/crates/stateless-core/src/executor.rs b/crates/stateless-core/src/executor.rs index 604d7c6e..85f3d1f2 100644 --- a/crates/stateless-core/src/executor.rs +++ b/crates/stateless-core/src/executor.rs @@ -546,8 +546,18 @@ where let logs_bloom = execution_result.receipts.iter().fold(Bloom::ZERO, |acc, receipt| acc | receipt.bloom()); - // Gas used is the cumulative gas used of the last receipt - let gas_used = execution_result.receipts.last().map(|r| r.cumulative_gas_used()).unwrap_or(0); + // `BlockExecutionResult::gas_used` is defined by mega-evm's `finish()` as + // `receipts.last().cumulative_gas_used()` — the exact expression the header check + // read before switching to this field, so the two cannot disagree regardless of how + // system transactions are accounted. The assertion pins that upstream definition: a + // future mega-evm that accounts gas outside the receipt chain fails loudly in every + // debug/test run instead of silently changing the header check. + let gas_used = execution_result.gas_used; + debug_assert_eq!( + gas_used, + execution_result.receipts.last().map(|r| r.cumulative_gas_used()).unwrap_or(0), + "mega-evm gas_used no longer equals the last receipt's cumulative gas" + ); let receipts_root = calculate_receipt_root(&execution_result.receipts); @@ -1105,6 +1115,39 @@ mod tests { /// returned updates must reproduce the header's state root when fed through the SALT trie /// update, locking its equivalence with the `validate_block` path the helpers were /// extracted from. + /// The header check reads `BlockExecutionResult::gas_used`; mega-evm defines that field + /// as the last receipt's cumulative gas, the expression the check read before. A + /// `debug_assert_eq!` at the derivation site pins the two against each other on every + /// replay; this pins the surviving value against what a real mainnet header claims, so a + /// mega-evm that starts accounting gas outside the receipt chain fails here rather than + /// as a consensus divergence in production. + #[test] + fn replayed_gas_used_matches_the_mainnet_header() { + let _logging = init_test_logging("stateless_core"); + let fx = TestFixtures::mainnet_shared(); + let paired = fx.paired_blocks(); + assert!(!paired.is_empty(), "no paired mainnet fixtures in test_data/mainnet"); + for (number, hash) in paired { + let block = &fx.blocks[&hash]; + let salt_witness = fx.salt_witnesses[&hash].clone(); + let header = block.consensus_header(); + let ext_env = WitnessExternalEnv::new(&salt_witness, header.number) + .expect("witness carries bucket metadata"); + let witness = Witness::from(salt_witness); + witness + .verify() + .unwrap_or_else(|e| panic!("witness verification failed for {number}: {e:?}")); + let witness_db = + WitnessDatabase { header, witness: &witness, contracts: &fx.contracts }; + let (_, output) = replay_block(&chain_spec(), block, &witness_db, ext_env) + .unwrap_or_else(|e| panic!("replay failed for {number} ({hash}): {e:?}")); + assert_eq!( + output.gas_used, block.header.gas_used, + "replayed gas_used disagrees with the header for {number} ({hash})" + ); + } + } + #[test] fn validate_block_deriving_updates_mainnet_fixtures() { let _logging = init_test_logging("stateless_core"); From 432e550ca1575cecd3daba8c69bd39945b32467e Mon Sep 17 00:00:00 2001 From: "liquan.eth" Date: Thu, 3 Sep 2026 19:19:58 +0800 Subject: [PATCH 4/9] perf: move blocking work off the async runtime Block integrity verification becomes the RPC retry loop's finalize step and runs on the blocking pool; the advancer's pre-advance hooks and store commit run there too. Panic semantics are preserved and pinned by tests. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01CHVpMMX9N69sUNbuKgpBVY --- AGENTS.md | 3 +- README.md | 4 +- crates/stateless-common/src/rpc_client.rs | 143 +++++++++++++++--- .../stateless-core/src/pipeline/advancer.rs | 40 +++-- crates/stateless-core/src/pipeline/mod.rs | 4 +- crates/stateless-core/src/pipeline/tests.rs | 72 +++++++-- 6 files changed, 219 insertions(+), 47 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index b68e91e1..5ba8cac8 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -55,6 +55,7 @@ Both binaries share a generic three-stage pipeline defined in `stateless-core::p 1. **Fetch** — `block_fetcher` streams blocks + witnesses from a `BlockFetcher` via a bounded in-flight window (concurrency capped by `fetcher_max_in_flight`). 2. **Process** — N workers run `BlockProcessor::process` (validator: EVM execution; trace server: pass-through). 3. **Advance** — `chain_advancer` reorders out-of-order results, verifies parent-hash continuity, detects reorgs, and persists via `ChainStore::advance_chain`. + The hooks' `pre_advance` and that commit are synchronous multi-millisecond disk work, so they run on the blocking pool rather than on an async worker (the trace server shares its runtime with the RPC handlers); a panic in either still unwinds the advancer unchanged, via `try_into_panic` → `resume_unwind`. The outer loop (`run_pipeline`) handles reorg rollback + restart, stale-data anchor reset, and transient vs fatal error classification. On a detected reorg, the rollback floor comes from a pluggable `ReorgResolver`: the core `BisectResolver` walks local history via `find_divergence_point`, while an embedder with an externally-supplied floor (e.g. the mega-reth FullNode) provides its own. @@ -155,7 +156,7 @@ The validator splits the two the same way: `--r2-max-concurrent-requests` caps R Client-side routing, budgets, and fallback match the S3 target, but edge behavior is zone configuration: **a cache rule making these objects cacheable must set 404s to bypass cache**, or a pre-upload frontier miss gets pinned for the negative-cache TTL (stalling the validator's tip-following in its fallback-less R2 mode) and a cached 404 can false-fire the below-band `kind="missing"` bucket-integrity alarm. The bucket is the same store the public gateway reads and can lead the generator at the frontier (uploader and generator RPC server publish from different files), so frontier hits are real; the frontier band is a small near-tip window (`R2_FRONTIER_WINDOW`, 32 blocks of uploader-lag grace on either side of the local tip — deliberately far narrower than the 4096-block routing window, so a stale catching-up tip cannot silence holes above it), hits there are labeled `witness_r2_frontier` (vs `witness_r2` past the band), the speculative frontier probe runs on an eighth of the remaining stage (vs half for blocks R2 must hold, so degraded R2 cannot burn half of every near-tip request's budget), and a `missing` classifies by band: in-band is the expected probe-ahead outcome (excluded from the alarm), below-band feeds `debug_trace_r2_witness_errors_total{kind="missing"}` (the bucket-integrity alarm, still covering recent-but-below-tip holes), and above-band — only reachable behind a stale catching-up tip — lands on its own `kind="missing_above_tip"` series, visible without flooding the alarm on every catch-up. Any witness-chain RPC attempt under a deadline is capped at the tightest of three bounds — half the full witness stage (`RpcClientConfig::witness_per_attempt_timeout`, derived from `--witness-timeout`), the global `--rpc-per-attempt-timeout-ms` (an explicitly stricter operator setting is honored, never loosened), and — only while the round still has an untried provider to rotate to — half of what the call still has as the attempt starts (recomputed after any concurrency-permit wait, so neither an old-block-clamped stage, a post-R2 remainder, nor a long permit queue defeats the reserve). -The round's last hop, and every hop of a single-provider chain, takes the remainder whole under the ceiling instead: rotation stays protected without structurally condemning a slow-but-honest transfer, and the witness decode runs outside the attempt window (bounded by the deadline alone), so CPU-bound decode neither burns the reserve nor reads as a provider stall while a corrupt payload still rotates as the provider's error; deadline-less chain-sync fetches keep the general 20s cap so a slower-than-cap transfer still completes. +The round's last hop, and every hop of a single-provider chain, takes the remainder whole under the ceiling instead: rotation stays protected without structurally condemning a slow-but-honest transfer, and the witness decode — like `eth_getBlock`'s integrity verification, which runs on the blocking pool — runs outside the attempt window (bounded by the deadline alone), so CPU-bound work neither burns the reserve nor reads as a provider stall while a corrupt payload still rotates as the provider's error; deadline-less chain-sync fetches keep the general 20s cap so a slower-than-cap transfer still completes. When a logical upstream call gives up on its deadline it logs one WARN naming the `phase` it died in (`before_attempt` / `permit_wait_clamped` / `attempt_clamped` / `before_backoff`) with `provider` / `round` / `permit_wait_ms` / `attempt_ms`, and the abandoned attempt is recorded as `outcome="deadline_clamped"` rather than dropped; best-effort internal probes (the throttled upstream tip seed) demote that give-up log to debug while the deadline metric still fires, so a probe whose failure is already degraded cannot page as a user-visible incident. Permit wait is timed separately (`debug_trace_upstream_permit_wait_seconds{method}`) and the acquire is clamped to the deadline (phase `permit_wait_clamped`, cut-short wait still sampled), so queueing behind our own `--witness-max-concurrent-requests` stays distinguishable from endpoint slowness and a saturated queue cannot block a call past its budget unobserved. The background chain-sync prefetch routes by freshness against the last observed remote head: frontier-fresh blocks give the generator a short exclusive grace (its "witness not found" means "not generated yet" — fallbacks are fed by the same pipeline and cannot be ahead) before falling back to the full endpoint chain, while deep catch-up blocks — and any block classified against a stale head observation (older than the grace, as during a long catch-up stretch when the tip is not re-polled) — use the full chain from the first attempt. diff --git a/README.md b/README.md index 04792fcc..e1a35a99 100644 --- a/README.md +++ b/README.md @@ -209,7 +209,7 @@ A count of zero, a non-numeric or blank one, one set alongside the S3 endpoint ( The configured target is published as the constant-1 gauge `debug_trace_r2_target_info{target}` (`r2_target_info` on the validator), so the target-less R2 series can be attributed to one target or the other during a rollout. Whether the domain actually delivered HTTP/2 is a separate question — selection is pure ALPN, so a misconfigured zone degrades to HTTP/1.1 with the h2 tuning inert — and is answered by `..._r2_negotiated_http_version_info{version}` plus a one-time warning naming the protocol that was negotiated. Any single witness-chain RPC attempt under a deadline is additionally capped at the tightest of: half the witness stage budget, the global per-attempt timeout, and — only while the round still has an untried provider to rotate to — half of what the call still has as the attempt starts (recomputed after any permit wait). -The round's last hop, and every hop of a single-provider chain, takes the remainder whole under the ceiling instead, so a stalled endpoint (or a saturated concurrency permit — waits are deadline-bounded too) can never consume the stage while a rotation is still worth reserving for, and a slow-but-honest transfer is never structurally condemned; the witness decode runs outside the attempt window, bounded by the request deadline alone. +The round's last hop, and every hop of a single-provider chain, takes the remainder whole under the ceiling instead, so a stalled endpoint (or a saturated concurrency permit — waits are deadline-bounded too) can never consume the stage while a rotation is still worth reserving for, and a slow-but-honest transfer is never structurally condemned; the witness decode — and `eth_getBlock`'s integrity verification, which runs on the blocking pool — runs outside the attempt window, bounded by the request deadline alone. `--r2-max-concurrent-requests` caps in-flight GETs separately from `--witness-max-concurrent-requests` — the RPC cap sizes a shared gateway, R2 tolerates far more. **Admission and response-size knobs** (each also settable via its `DEBUG_TRACE_SERVER_*` env var): @@ -320,7 +320,7 @@ Both binaries share a generic three-stage pipeline defined in `stateless-core`: Reorders out-of-order results (BTreeMap) Verifies parent-hash continuity Detects reorgs → rollback + restart - Persists via ChainStore::advance_chain() + Persists via ChainStore::advance_chain() (on the blocking pool) Outer loop (run_pipeline): Reorg → ReorgResolver decides floor → rollback → restart pipeline diff --git a/crates/stateless-common/src/rpc_client.rs b/crates/stateless-common/src/rpc_client.rs index 8971c0d1..a54cc905 100644 --- a/crates/stateless-common/src/rpc_client.rs +++ b/crates/stateless-common/src/rpc_client.rs @@ -452,15 +452,24 @@ impl RpcClient { }) } + /// The data provider this call's first round starts at: rotated per call so healthy + /// endpoints share load evenly. Within a round the order is fixed (start → start+1 → …). + /// + /// Safety: the constructor guarantees at least one data provider. The atomic op is + /// skipped for a single provider — pointless contention otherwise. + fn next_data_rr_start(&self) -> usize { + let n = self.data_providers.len(); + if n > 1 { self.data_rr_counter.fetch_add(1, Ordering::Relaxed) % n } else { 0 } + } + /// Deadline-aware counterpart of [`Self::call`]. /// /// With `deadline = Some(..)` the retry loop returns [`RpcDeadlineExceeded`] once the /// deadline passes, clamping each inter-round sleep so it doesn't overshoot. With /// `None` this is equivalent to [`Self::call`] and never returns `Err`. /// - /// Each call performs rounds of "try every data provider once in round-robin order". - /// The starting provider rotates per call via an atomic counter so healthy endpoints - /// share load evenly; within a round the order is fixed (start → start+1 → …). + /// Each call performs rounds of "try every data provider once in round-robin order", + /// starting at [`Self::next_data_rr_start`]. async fn call_with_deadline( &self, method: RpcMethod, @@ -481,18 +490,13 @@ impl RpcClient { best_effort: bool, f: impl Fn(RootProvider) -> BoxFuture>, ) -> std::result::Result { - // Safety: constructor guarantees at least one data provider. - let n = self.data_providers.len(); - // Skip the atomic op when there's a single provider — avoids pointless contention. - let rr_start = - if n > 1 { self.data_rr_counter.fetch_add(1, Ordering::Relaxed) % n } else { 0 }; round_robin_with_backoff( &self.data_providers, &self.data_provider_labels, &self.data_concurrency, &self.config.rpc_retry, AttemptCap::Fixed(self.config.per_attempt_timeout), - rr_start, + self.next_data_rr_start(), method, self.config.metrics.as_ref(), deadline, @@ -554,6 +558,13 @@ impl RpcClient { } /// Deadline-aware counterpart of [`Self::get_block`]. + /// + /// Verification is the retry loop's *finalize* step, not part of the attempt window: it + /// is per-transaction ECDSA recovery and re-encoding over a whole block — CPU-bound work + /// that would otherwise both burn the rotation reserve and read as a provider stall, + /// classifying a healthy endpoint serving a large block as stalled (see + /// [`round_robin_with_backoff`]). A verification failure still counts as that provider's + /// error, so a tampered block rotates exactly as a transport failure does. pub async fn get_block_with_deadline( &self, block_id: BlockId, @@ -561,15 +572,29 @@ impl RpcClient { deadline: Option, ) -> std::result::Result, RpcDeadlineExceeded> { let verify = !self.config.skip_block_verification; - self.call_with_deadline(RpcMethod::EthGetBlock, deadline, move |provider| { - Box::pin(async move { - let block = do_get_block_unchecked(&provider, block_id, full_txs).await?; - if verify { - verify_block_integrity(&block)?; - } - Ok(block) - }) - }) + round_robin_with_backoff( + &self.data_providers, + &self.data_provider_labels, + &self.data_concurrency, + &self.config.rpc_retry, + AttemptCap::Fixed(self.config.per_attempt_timeout), + self.next_data_rr_start(), + RpcMethod::EthGetBlock, + self.config.metrics.as_ref(), + deadline, + false, + move |provider, _provider_label| { + Box::pin(async move { do_get_block_unchecked(&provider, block_id, full_txs).await }) + }, + move |block, _provider_label| { + Box::pin(async move { + if !verify { + return Ok(block); + } + verify_block_on_blocking_pool(block).await + }) + }, + ) .await } @@ -1573,6 +1598,21 @@ async fn decode_witness_wire( Ok(result) } +/// [`RpcClient::get_block_with_deadline`]'s finalize half: [`verify_block_integrity`] on the +/// blocking pool (per-transaction ECDSA recovery plus a re-encode of every envelope is +/// CPU-bound over a full block), handing the block back untouched on success. +/// +/// A failure here is an integrity failure from this provider — the retry loop records it as +/// that provider's `Error` and rotates, exactly like a transport error. +async fn verify_block_on_blocking_pool(block: Block) -> Result> { + tokio::task::spawn_blocking(move || -> Result> { + verify_block_integrity(&block)?; + Ok(block) + }) + .await + .context("block verification task panicked")? +} + /// Verifies structural integrity of a block fetched from RPC. /// /// Checks: @@ -1660,7 +1700,7 @@ mod tests { find_divergence_point, pipeline::{BlockFetcher, DivergenceLookups}, }; - use stateless_test_utils::mock_rpc::{header_stub, parse_hex_u64, serve}; + use stateless_test_utils::mock_rpc::{consistent_header, header_stub, parse_hex_u64, serve}; use tokio_util::sync::CancellationToken; use super::*; @@ -1979,6 +2019,71 @@ mod tests { hb.stop().unwrap(); } + /// A block with no transactions, so [`verify_block_integrity`] reduces to its header-hash + /// check and the fixture needs no signed transactions. + fn block_stub(header: alloy_rpc_types_eth::Header) -> Block { + Block { + header, + uncles: Vec::new(), + transactions: alloy_rpc_types_eth::BlockTransactions::Hashes(Vec::new()), + withdrawals: None, + } + } + + /// Block verification is the retry loop's finalize step rather than part of the attempt + /// window (so a large block cannot read as a provider stall), and this pins the property + /// that move must not cost: an integrity failure is still that provider's error, so the + /// call rotates to the next provider instead of surfacing the bad block. + #[tokio::test] + async fn block_verification_failure_rotates_to_the_next_provider() { + let bad_hits = Arc::new(AtomicUsize::new(0)); + let (bad_handle, bad_url) = serve(Arc::clone(&bad_hits), |m| { + m.register_method("eth_getBlockByNumber", |_params, hits, _| { + hits.fetch_add(1, Ordering::Relaxed); + // A hash the header does not actually hash to — the first integrity check. + Ok::<_, ErrorObjectOwned>(block_stub(header_stub(7, BlockHash::from([9u8; 32])))) + }) + .unwrap(); + }) + .await; + let good_hits = Arc::new(AtomicUsize::new(0)); + let (good_handle, good_url) = serve(Arc::clone(&good_hits), |m| { + m.register_method("eth_getBlockByNumber", |_params, hits, _| { + hits.fetch_add(1, Ordering::Relaxed); + Ok::<_, ErrorObjectOwned>(block_stub(consistent_header(7))) + }) + .unwrap(); + }) + .await; + + let config = RpcClientConfig { + rpc_retry: BackoffPolicy::new(Duration::from_millis(1), Duration::from_millis(2)), + ..Default::default() + }; + let client = RpcClient::new_with_config( + &[bad_url.as_str(), good_url.as_str()], + &[good_url.as_str()], + config, + None, + ) + .unwrap(); + + let block = client + .get_block_with_deadline(BlockId::number(7), false, None) + .await + .expect("None deadline cannot time out"); + assert_eq!( + block.header.hash, + consistent_header(7).hash, + "the verified block must come from the second provider" + ); + assert_eq!(bad_hits.load(Ordering::Relaxed), 1, "the tampering provider must be tried"); + assert_eq!(good_hits.load(Ordering::Relaxed), 1, "rotation must reach the good provider"); + + bad_handle.stop().unwrap(); + good_handle.stop().unwrap(); + } + /// After every provider in a round fails, the helper sleeps and starts a new round. /// Observable via two providers that each fail one request then recover: the call must /// succeed after at least one full round of failures. diff --git a/crates/stateless-core/src/pipeline/advancer.rs b/crates/stateless-core/src/pipeline/advancer.rs index 360da60c..0ec5c01f 100644 --- a/crates/stateless-core/src/pipeline/advancer.rs +++ b/crates/stateless-core/src/pipeline/advancer.rs @@ -1,6 +1,6 @@ //! Advancer stage: reorders processed blocks, detects reorgs, persists progress. -use std::collections::BTreeMap; +use std::{collections::BTreeMap, sync::Arc}; use eyre::Result; use tokio_util::sync::CancellationToken; @@ -78,8 +78,8 @@ where /// and advances the canonical chain. pub(crate) async fn chain_advancer( fetcher: &F, - store: &S, - hooks: &H, + store: Arc, + hooks: Arc, resolver: &R, result_rx: kanal::Receiver>, initial_tip: BlockMeta, @@ -87,7 +87,7 @@ pub(crate) async fn chain_advancer( ) -> Result where F: BlockFetcher, - S: ChainStore, + S: ChainStore + 'static, H: PipelineHooks, R: ReorgResolver, { @@ -103,7 +103,6 @@ where // Reused across iterations to avoid per-iteration allocations; typical batch // size is small (<= `concurrent_workers`) and stable. let mut batch: Vec = Vec::new(); - let mut metas: Vec = Vec::new(); loop { let item = tokio::select! { @@ -127,7 +126,6 @@ where buffer.insert(item.block_number(), item); batch.clear(); - metas.clear(); while let Some(item) = buffer.remove(&next_expected) { if item.parent_hash() != current_tip.block_hash { debug!( @@ -139,7 +137,7 @@ where // Strategy is scenario-supplied (see `ReorgResolver`); `Floor` rolls back, // `Fatal`/`Retry` end the cycle. - let rollback_to = match resolver.resolve(fetcher, store, persisted_tip).await? { + let rollback_to = match resolver.resolve(fetcher, &store, persisted_tip).await? { ReorgResolution::Floor(floor) => { debug!(block = next_expected, floor, "Resolved reorg floor"); floor @@ -179,17 +177,37 @@ where } current_tip = item.to_block_meta(); next_expected += 1; - metas.push(current_tip.clone()); batch.push(item); } if !batch.is_empty() { - hooks.pre_advance(&batch)?; - store.advance_chain(&metas)?; + // The hooks' pre-advance persistence and the store commit are synchronous, + // potentially multi-ms disk work — run them off the async runtime so they can't + // stall other tasks (the trace server shares this runtime with its RPC handlers). + let advance_store = store.clone(); + let advance_hooks = hooks.clone(); + let owned_batch = std::mem::take(&mut batch); + let advanced = tokio::task::spawn_blocking(move || { + let metas: Vec = + owned_batch.iter().map(|item| item.to_block_meta()).collect(); + advance_hooks.pre_advance(&owned_batch)?; + advance_store.advance_chain(&metas)?; + Ok::<_, eyre::Report>(owned_batch) + }) + .await; + batch = match advanced { + Ok(buf) => buf?, + // A panic in the store/hooks must propagate unchanged, exactly as it did + // when these calls ran inline on this task. + Err(join_err) => match join_err.try_into_panic() { + Ok(payload) => std::panic::resume_unwind(payload), + Err(join_err) => return Err(join_err.into()), + }, + }; persisted_tip = current_tip.block_number; debug!( tip = current_tip.block_number, - advanced = metas.len(), + advanced = batch.len(), buffered = buffer.len(), "Chain advanced" ); diff --git a/crates/stateless-core/src/pipeline/mod.rs b/crates/stateless-core/src/pipeline/mod.rs index 1603713d..092b2262 100644 --- a/crates/stateless-core/src/pipeline/mod.rs +++ b/crates/stateless-core/src/pipeline/mod.rs @@ -93,8 +93,8 @@ where let outcome = chain_advancer( &*fetcher, - &*store, - &*hooks, + store.clone(), + hooks.clone(), &resolver, result_rx, initial_tip, diff --git a/crates/stateless-core/src/pipeline/tests.rs b/crates/stateless-core/src/pipeline/tests.rs index deaeca59..e10767a3 100644 --- a/crates/stateless-core/src/pipeline/tests.rs +++ b/crates/stateless-core/src/pipeline/tests.rs @@ -175,7 +175,7 @@ async fn run_advancer( tip: BlockMeta, rpc_hashes: HashMap, blocks: Vec, -) -> (Result, MockStore) { +) -> (Result, Arc) { run_advancer_with_resolver(tip, rpc_hashes, blocks, &BisectResolver).await } @@ -184,10 +184,10 @@ async fn run_advancer_with_resolver>( rpc_hashes: HashMap, blocks: Vec, resolver: &R, -) -> (Result, MockStore) { - let store = MockStore::new(tip.clone()); +) -> (Result, Arc) { + let store = Arc::new(MockStore::new(tip.clone())); let fetcher = MockFetcher { hashes: rpc_hashes }; - let hooks = NoopHooks; + let hooks = Arc::new(NoopHooks); let (tx, rx) = kanal::bounded(16); { @@ -210,7 +210,8 @@ async fn run_advancer_with_resolver>( } let result = - chain_advancer(&fetcher, &store, &hooks, resolver, rx, tip, CancellationToken::new()).await; + chain_advancer(&fetcher, store.clone(), hooks, resolver, rx, tip, CancellationToken::new()) + .await; (result, store) } @@ -274,9 +275,9 @@ impl PipelineHooks for BadBlockHooks { /// but specialized — generalizing the original would cascade through ~6 helper types for /// a single test case. async fn run_bad_block_advancer(tip: BlockMeta, blocks: Vec) -> Result { - let store = MockStore::new(tip.clone()); + let store = Arc::new(MockStore::new(tip.clone())); let fetcher = MockFetcher { hashes: HashMap::default() }; - let hooks = BadBlockHooks; + let hooks = Arc::new(BadBlockHooks); let (tx, rx) = kanal::bounded::< std::result::Result, ErrorAction)>, >(16); @@ -288,8 +289,55 @@ async fn run_bad_block_advancer(tip: BlockMeta, blocks: Vec) -> Result } } - chain_advancer(&fetcher, &store, &hooks, &BisectResolver, rx, tip, CancellationToken::new()) - .await + chain_advancer(&fetcher, store, hooks, &BisectResolver, rx, tip, CancellationToken::new()).await +} + +/// Hooks whose `pre_advance` panics — the advance step now runs on the blocking pool, and a +/// panic there arrives as a `JoinError` rather than unwinding this task on its own. +struct PanickingHooks; +impl PipelineHooks for PanickingHooks { + type Output = MockBlock; + + fn pre_advance(&self, _items: &[Self::Output]) -> Result<()> { + panic!("pre_advance exploded"); + } +} + +/// Moving the advance step onto the blocking pool must not turn a store/hooks panic into an +/// ordinary `Err` — a corrupted persistence layer has to keep taking the process down exactly +/// as it did when these calls ran inline. Pins the `try_into_panic` → `resume_unwind` branch. +#[tokio::test] +async fn test_chain_advancer_propagates_hook_panics() { + let tip = make_tip(10); + let store = Arc::new(MockStore::new(tip.clone())); + let fetcher = MockFetcher { hashes: HashMap::default() }; + let hooks = Arc::new(PanickingHooks); + let parent = tip.block_hash; + let (tx, rx) = kanal::bounded::< + std::result::Result, ErrorAction)>, + >(16); + tx.to_async().send(Ok(make_block(11, parent))).await.unwrap(); + + // Run it in its own task so the unwind is observable: a `JoinError` that `is_panic()` + // means the advancer panicked, an `Ok(Err(..))` would mean it swallowed the panic into + // the error channel. + let joined = tokio::spawn(async move { + chain_advancer(&fetcher, store, hooks, &BisectResolver, rx, tip, CancellationToken::new()) + .await + }) + .await; + let join_err = match joined { + Err(e) => e, + Ok(outcome) => panic!("a panicking hook must unwind the advancer, got {outcome:?}"), + }; + assert!(join_err.is_panic(), "the advancer must fail by panic, not by JoinError::Cancelled"); + let payload = join_err.into_panic(); + let message = payload + .downcast_ref::<&str>() + .map(|s| (*s).to_string()) + .or_else(|| payload.downcast_ref::().cloned()) + .unwrap_or_default(); + assert!(message.contains("pre_advance exploded"), "panic payload preserved: {message}"); } /// Covers the `verify_continuity` → Fatal branch in `chain_advancer` (`advancer.rs:109`). @@ -344,9 +392,9 @@ async fn test_chain_advancer_transient_error_returns_retry_outcome() { #[tokio::test] async fn test_chain_advancer_shutdown() { let tip = make_tip(10); - let store = MockStore::new(tip.clone()); + let store = Arc::new(MockStore::new(tip.clone())); let fetcher = MockFetcher { hashes: HashMap::default() }; - let hooks = NoopHooks; + let hooks = Arc::new(NoopHooks); let (_tx, rx) = kanal::bounded::< std::result::Result, ErrorAction)>, >(16); @@ -354,7 +402,7 @@ async fn test_chain_advancer_shutdown() { shutdown.cancel(); let outcome = - chain_advancer(&fetcher, &store, &hooks, &BisectResolver, rx, tip, shutdown).await.unwrap(); + chain_advancer(&fetcher, store, hooks, &BisectResolver, rx, tip, shutdown).await.unwrap(); assert!(matches!(outcome, PipelineOutcome::Shutdown)); } From 50741dc320e7cb83f3c285bfca014e37d53ea666 Mon Sep 17 00:00:00 2001 From: "liquan.eth" Date: Thu, 3 Sep 2026 19:21:19 +0800 Subject: [PATCH 5/9] perf(core): drop the per-state-read heap allocations in the witness path `WitnessDatabase` reads go through salt's `find()` with in-place decoding and stack-allocated key encodings; `LightWitness::metadata` stops cloning the `SaltValue` it only reads. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01CHVpMMX9N69sUNbuKgpBVY --- crates/stateless-core/src/data_types.rs | 45 +++++++++++++++++++--- crates/stateless-core/src/evm_database.rs | 28 ++++++++------ crates/stateless-core/src/light_witness.rs | 2 +- 3 files changed, 58 insertions(+), 17 deletions(-) diff --git a/crates/stateless-core/src/data_types.rs b/crates/stateless-core/src/data_types.rs index 47905aa3..4a859424 100644 --- a/crates/stateless-core/src/data_types.rs +++ b/crates/stateless-core/src/data_types.rs @@ -26,7 +26,7 @@ use std::{collections::BTreeMap, vec::Vec}; pub use alloy_primitives::Bytes; -use alloy_primitives::{Address, B256, U256}; +use alloy_primitives::{Address, B256, FixedBytes, U256}; use revm::primitives::KECCAK_EMPTY; use salt::{SaltKey, SaltValue}; @@ -67,14 +67,33 @@ impl PlainKey { /// - Unknown: preserved raw bytes from decode pub fn encode(&self) -> Vec { match self { - PlainKey::Account(addr) => addr.as_slice().to_vec(), - PlainKey::Storage(addr, slot) => { - addr.concat_const::(*slot).as_slice().to_vec() - } + PlainKey::Account(addr) => Self::account_key_bytes(addr).to_vec(), + PlainKey::Storage(addr, slot) => Self::storage_key_bytes(*addr, *slot).to_vec(), PlainKey::Unknown(data) => data.clone(), } } + /// Encoding of an account key — the raw address bytes. + /// + /// Same bytes as `PlainKey::Account(address).encode()` without the heap allocation, + /// for per-state-read hot paths. + #[inline] + pub(crate) fn account_key_bytes(address: &Address) -> &[u8] { + address.as_slice() + } + + /// Stack-allocated encoding of a storage-slot key — address (20) ++ slot (32). + /// + /// Same bytes as `PlainKey::Storage(address, slot).encode()` without the heap + /// allocation, for per-state-read hot paths. + #[inline] + pub(crate) fn storage_key_bytes( + address: Address, + slot: B256, + ) -> FixedBytes { + address.concat_const::(slot) + } + /// Decodes a byte slice into a PlainKey. /// /// Returns `PlainKey::Unknown` if the buffer length is neither 20 (account) @@ -239,6 +258,22 @@ mod tests { entries.into_iter().enumerate().map(|(i, v)| (SaltKey::from((0u32, i as u64)), v)).collect() } + /// The allocation-free encoders the witness read path uses must produce exactly what + /// `encode()` produces — `encode()` delegates to them today, and this pins that contract + /// against a future edit that stops delegating and lets the two drift into looking up + /// different keys. + #[test] + fn stack_key_encodings_match_encode() { + let addr = Address::from([0x11; 20]); + assert_eq!(PlainKey::account_key_bytes(&addr), PlainKey::Account(addr).encode().as_slice()); + for slot in [B256::ZERO, B256::from([0xAB; 32]), B256::from(U256::from(7))] { + assert_eq!( + PlainKey::storage_key_bytes(addr, slot).as_slice(), + PlainKey::Storage(addr, slot).encode().as_slice(), + ); + } + } + #[test] fn test_plain_key_round_trip() { let addr = Address::from([0xAB; 20]); diff --git a/crates/stateless-core/src/evm_database.rs b/crates/stateless-core/src/evm_database.rs index 5927c3d7..1618a1ec 100644 --- a/crates/stateless-core/src/evm_database.rs +++ b/crates/stateless-core/src/evm_database.rs @@ -8,7 +8,6 @@ use std::{ collections::BTreeMap, format, string::{String, ToString}, - vec::Vec, }; use alloy_consensus::Header; @@ -79,10 +78,16 @@ where W: StateReader, W::Error: core::fmt::Display, { - /// Get value from witness for the given plain key - fn plain_value(&self, plain_key: &[u8]) -> Result>, WitnessDatabaseError> { + /// Get the witness entry for the given plain key. + /// + /// Returns the whole `SaltValue` (an inline array) instead of going through salt's + /// `plain_value()`, which heap-allocates a copy of the value bytes on every read — + /// this lookup runs once per unique account/slot touched during block replay. + /// Callers decode from `.value()` in place. + fn find(&self, plain_key: &[u8]) -> Result, WitnessDatabaseError> { EphemeralSaltState::new(self.witness) - .plain_value(plain_key) + .find(plain_key) + .map(|found| found.map(|(_, salt_value)| salt_value)) .map_err(|e| WitnessDatabaseError(e.to_string())) } } @@ -98,9 +103,9 @@ where fn basic_ref(&self, address: Address) -> Result, Self::Error> { trace!(?address, "basic_ref"); - let raw_value = self.plain_value(&PlainKey::Account(address).encode())?; + let salt_value = self.find(PlainKey::account_key_bytes(&address))?; - match raw_value.and_then(|v| match PlainValue::decode(&v) { + match salt_value.and_then(|v| match PlainValue::decode(v.value()) { PlainValue::Account(acc) => Some(acc), _ => None, }) { @@ -134,10 +139,11 @@ where fn storage_ref(&self, address: Address, index: U256) -> Result { trace!(?address, index = %format_args!("{:#x}", index), "storage_ref"); - let raw_value = self.plain_value(&PlainKey::Storage(address, index.into()).encode())?; + let salt_value = + self.find(PlainKey::storage_key_bytes(address, index.into()).as_slice())?; - Ok(raw_value - .and_then(|v| match PlainValue::decode(&v) { + Ok(salt_value + .and_then(|v| match PlainValue::decode(v.value()) { PlainValue::Storage(value) => Some(value), _ => None, }) @@ -298,11 +304,11 @@ impl SaltEnv for WitnessExternalEnv { } fn bucket_id_for_account(account: Address) -> BucketId { - hasher::bucket_id(&PlainKey::Account(account).encode()) + hasher::bucket_id(PlainKey::account_key_bytes(&account)) } fn bucket_id_for_slot(address: Address, key: U256) -> BucketId { - hasher::bucket_id(&PlainKey::Storage(address, key.into()).encode()) + hasher::bucket_id(PlainKey::storage_key_bytes(address, key.into()).as_slice()) } } diff --git a/crates/stateless-core/src/light_witness.rs b/crates/stateless-core/src/light_witness.rs index 973a03e6..244b450e 100644 --- a/crates/stateless-core/src/light_witness.rs +++ b/crates/stateless-core/src/light_witness.rs @@ -180,7 +180,7 @@ impl StateReader for LightWitness { fn metadata(&self, bucket_id: BucketId) -> Result { let metadata_key = bucket_metadata_key(bucket_id); match self.kvs.get(&metadata_key) { - Some(Some(salt_value)) => BucketMeta::try_from(salt_value.clone()) + Some(Some(salt_value)) => BucketMeta::try_from(salt_value) .map_err(|_| LightWitnessError { message: "Failed to decode metadata" }), // A well-formed witness never maps a metadata key to a deletion, // but witness bytes are network input (and the light decode From ae09a4f9a24cae2325b6b766e6d4c6b6691c0c38 Mon Sep 17 00:00:00 2001 From: "liquan.eth" Date: Thu, 3 Sep 2026 19:27:46 +0800 Subject: [PATCH 6/9] fix: allow a witness-less RpcClient and type the unsatisfiable-range failure R2 mode stops handing the data endpoints to the client as placeholder witness endpoints, and an unsatisfiable witness provider range fails one request as `WitnessFetchError::NoProviderInRange` instead of panicking the process. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01CHVpMMX9N69sUNbuKgpBVY --- bin/debug-trace-server/src/data_provider.rs | 38 ++++++- bin/stateless-validator/src/app.rs | 7 +- crates/stateless-common/src/lib.rs | 2 +- crates/stateless-common/src/rpc_client.rs | 117 ++++++++++++++------ 4 files changed, 127 insertions(+), 37 deletions(-) diff --git a/bin/debug-trace-server/src/data_provider.rs b/bin/debug-trace-server/src/data_provider.rs index 23ed6da7..1ca34c71 100644 --- a/bin/debug-trace-server/src/data_provider.rs +++ b/bin/debug-trace-server/src/data_provider.rs @@ -45,7 +45,9 @@ use futures::{FutureExt, future::Shared}; use op_alloy_rpc_types::Transaction; use quick_cache::sync::Cache; use revm::state::Bytecode; -use stateless_common::{CodeFetchError, RpcClient, RpcDeadlineExceeded, WitnessSizeBreakdown}; +use stateless_common::{ + CodeFetchError, RpcClient, RpcDeadlineExceeded, WitnessFetchError, WitnessSizeBreakdown, +}; use stateless_core::{ ContractStore, LightWitness, StoreResult, db::StoreError, withdrawals::MptWitness, }; @@ -277,6 +279,18 @@ impl From for DataProviderError { } } +impl From for DataProviderError { + fn from(e: WitnessFetchError) -> Self { + match e { + // Only a blown deadline is a timeout. A range failure is a wiring bug in this + // process — routing it to `Timeout { Witness }` would fire the `deadline_witness` + // alarm, which must mean "an upstream witness fetch ran out of budget". + WitnessFetchError::Deadline(d) => d.into(), + WitnessFetchError::NoProviderInRange { .. } => eyre::eyre!("{e}").into(), + } + } +} + impl From for DataProviderError { fn from(e: CodeFetchError) -> Self { match e { @@ -3068,4 +3082,26 @@ mod tests { .into(); assert!(matches!(block_err, DataProviderError::Timeout { stage: TimeoutStage::Block, .. })); } + + /// A witness fetch whose provider range is unsatisfiable is a wiring bug, not a blown + /// budget: it must land on `Internal`, never on `Timeout { Witness }`. That bucket feeds + /// the `deadline_witness` error reason, whose whole value is meaning "an upstream witness + /// fetch ran out of time" — a wiring bug landing there would page for the wrong incident. + /// The deadline variant still classifies by method, exactly as before. + #[test] + fn witness_range_failure_is_internal_not_a_witness_timeout() { + let range_err: DataProviderError = + WitnessFetchError::NoProviderInRange { skip: 2, configured: 1 }.into(); + assert!(matches!(range_err, DataProviderError::Internal(_)), "got {range_err:?}"); + + let deadline_err: DataProviderError = WitnessFetchError::Deadline(RpcDeadlineExceeded { + method: stateless_common::RpcMethod::MegaGetBlockWitness, + elapsed: Duration::from_secs(3), + }) + .into(); + assert!(matches!( + deadline_err, + DataProviderError::Timeout { stage: TimeoutStage::Witness, .. } + )); + } } diff --git a/bin/stateless-validator/src/app.rs b/bin/stateless-validator/src/app.rs index d5347f58..712bda0d 100644 --- a/bin/stateless-validator/src/app.rs +++ b/bin/stateless-validator/src/app.rs @@ -331,8 +331,6 @@ pub async fn run() -> Result<()> { ..rpc_defaults } .with_metrics(Arc::new(metrics::ValidatorMetrics)); - // In R2 mode the RpcClient's witness providers are never used, but its constructor requires - // a non-empty list — hand it the data endpoints as a placeholder. let data_apis: Vec<&str> = args.rpc_endpoint.iter().map(String::as_str).collect(); let r2_witness = match args.witness_source { WitnessSource::Rpc => { @@ -361,8 +359,11 @@ pub async fn run() -> Result<()> { } }; + // In R2 mode the client carries no witness providers: witnesses come straight from R2, and + // a witness RPC call that slipped through fails structurally instead of quietly asking the + // data endpoints for `mega_getBlockWitness`. let witness_apis: Vec<&str> = if r2_witness.is_some() { - data_apis.clone() + Vec::new() } else { args.witness_endpoint.iter().map(String::as_str).collect() }; diff --git a/crates/stateless-common/src/lib.rs b/crates/stateless-common/src/lib.rs index aeed85f8..0ee6fbdc 100644 --- a/crates/stateless-common/src/lib.rs +++ b/crates/stateless-common/src/lib.rs @@ -4,7 +4,7 @@ pub use metrics::{RpcMethod, RpcMetrics}; pub mod rpc_client; pub use rpc_client::{ BackoffPolicy, CodeFetchError, RpcClient, RpcClientConfig, RpcDeadlineExceeded, - SetValidatedBlocksResponse, WitnessRequestKeys, + SetValidatedBlocksResponse, WitnessFetchError, WitnessRequestKeys, }; pub mod witness_encoding; pub use witness_encoding::{ diff --git a/crates/stateless-common/src/rpc_client.rs b/crates/stateless-common/src/rpc_client.rs index a54cc905..78ce374c 100644 --- a/crates/stateless-common/src/rpc_client.rs +++ b/crates/stateless-common/src/rpc_client.rs @@ -254,6 +254,24 @@ pub struct SetValidatedBlocksResponse { pub last_validated_block: (U64, B256), } +/// Error returned by the witness fetches that take a caller-computed provider range. +/// +/// `NoProviderInRange` is a wiring failure, not a transport one: either the caller's `skip` +/// selected past the configured witness endpoints, or the client carries none at all (a +/// deployment that sources witnesses elsewhere — the validator's R2 mode). Both binaries +/// reject an empty witness configuration at startup, so it stays unreachable in production; +/// it is a typed error rather than an `assert!` so a routing bug fails one request instead of +/// the process. +#[derive(Debug, thiserror::Error)] +pub enum WitnessFetchError { + #[error( + "witness fetch selected providers {skip}.. of {configured} configured — no witness provider in range" + )] + NoProviderInRange { skip: usize, configured: usize }, + #[error(transparent)] + Deadline(#[from] RpcDeadlineExceeded), +} + /// Errors returned by [`RpcClient::get_codes`] / [`RpcClient::get_codes_with_deadline`]. /// /// - `VerificationFailure` is deterministic (upstream returned bytecode whose keccak does not match @@ -320,7 +338,10 @@ impl RpcClient { /// # Arguments /// * `data_apis` - HTTP URLs of the standard JSON-RPC endpoints for blocks and contract data /// (tried in order, non-empty) - /// * `witness_apis` - HTTP URLs of the witness RPC endpoints (tried in order, non-empty) + /// * `witness_apis` - HTTP URLs of the witness RPC endpoints (tried in order). May be empty + /// when the deployment sources witnesses elsewhere (the validator's R2 witness mode); a + /// witness call on such a client returns [`WitnessFetchError::NoProviderInRange`] instead of + /// silently retrying against the wrong endpoints /// * `config` - Configuration controlling verification, retry, and concurrency behavior /// * `report_api` - Optional HTTP URL of the endpoint for reporting validated blocks pub fn new_with_config( @@ -332,9 +353,6 @@ impl RpcClient { if data_apis.is_empty() { return Err(eyre!("At least one data API URL must be provided")); } - if witness_apis.is_empty() { - return Err(eyre!("At least one witness API URL must be provided")); - } // One shared HTTP client for every provider (connection pools are keyed per host), so // the connect-phase bound applies uniformly to data, witness, and report endpoints. @@ -726,7 +744,8 @@ impl RpcClient { decode_witness_response, "Witness decoded", ) - .await?; + .await + .map_err(deadline_only)?; if let Some(ref metrics) = self.config.metrics { metrics.on_witness_fetch(WitnessSizeBreakdown::new(&witness.0, &witness.1)); @@ -757,7 +776,9 @@ impl RpcClient { hash: B256, deadline: Option, ) -> std::result::Result<(LightWitness, MptWitness), RpcDeadlineExceeded> { - self.get_witness_light_with_deadline_from(0, number, hash, deadline).await + self.get_witness_light_with_deadline_from(0, number, hash, deadline) + .await + .map_err(deadline_only) } /// Like [`Self::get_witness_light_with_deadline`], but skips the first `skip` witness @@ -766,15 +787,16 @@ impl RpcClient { /// position in the full configured witness endpoint list, and the shared witness /// concurrency cap still applies. /// - /// # Panics - /// Panics if `skip >= witness_provider_count()` — at least one provider must remain. + /// Returns [`WitnessFetchError::NoProviderInRange`] when `skip` selects past the + /// configured witness endpoints — a routing bug fails this one request rather than the + /// process. pub async fn get_witness_light_with_deadline_from( &self, skip: usize, number: u64, hash: B256, deadline: Option, - ) -> std::result::Result<(LightWitness, MptWitness), RpcDeadlineExceeded> { + ) -> std::result::Result<(LightWitness, MptWitness), WitnessFetchError> { self.witness_round_robin( skip..self.witness_providers.len(), number, @@ -809,7 +831,7 @@ impl RpcClient { "Witness light-decoded", ) .await - .expect("None deadline cannot time out") + .expect("pinned 0..1 range and a None deadline cannot fail") } /// Shared `mega_getBlockWitness` retry loop: primary-failover rounds (always start from @@ -821,8 +843,8 @@ impl RpcClient { /// the logged endpoint labels stay aligned with the full configured list because each /// label bakes in its original index (see [`endpoint_label`]). /// - /// # Panics - /// Panics if `providers` is empty or out of bounds — at least one provider must remain. + /// An empty or out-of-bounds `providers` range is a wiring failure, surfaced as + /// [`WitnessFetchError::NoProviderInRange`] rather than a panic. // A `warn`-level span (not the usual `info`) so it stays enabled at the default `warn` log // filter: the generic retry loop's per-attempt failure logs then inherit `block_number`, // which they cannot see otherwise, so an endpoint stall/error is traceable to its block. @@ -835,12 +857,11 @@ impl RpcClient { deadline: Option, decode: fn(&str) -> std::result::Result, trace_msg: &'static str, - ) -> std::result::Result { - assert!( - !providers.is_empty() && providers.end <= self.witness_providers.len(), - "witness provider range ({providers:?}) must select at least one of {} providers", - self.witness_providers.len() - ); + ) -> std::result::Result { + let configured = self.witness_providers.len(); + if providers.is_empty() || providers.end > configured { + return Err(WitnessFetchError::NoProviderInRange { skip: providers.start, configured }); + } // Deadline-bound witness attempts run under the reserve-half policy: the tightest of // the configured ceiling, the general per-attempt timeout, and — recomputed at each // attempt, after any permit wait — half of what the call still has, so neither a @@ -876,6 +897,7 @@ impl RpcClient { }, ) .await + .map_err(WitnessFetchError::Deadline) } /// Reports a range of validated blocks via the dedicated report endpoint. @@ -1613,6 +1635,19 @@ async fn verify_block_on_blocking_pool(block: Block) -> Result RpcDeadlineExceeded { + match e { + WitnessFetchError::Deadline(d) => d, + WitnessFetchError::NoProviderInRange { skip, configured } => unreachable!( + "full-range witness fetch on a client with no witness providers \ + (skip={skip}, configured={configured})" + ), + } +} + /// Verifies structural integrity of a block fetched from RPC. /// /// Checks: @@ -1833,11 +1868,12 @@ mod tests { .to_string() .contains("At least one data API") ); - assert!( - RpcClient::new(&[LOCALHOST_A], &[]) - .unwrap_err() - .to_string() - .contains("At least one witness API") + // An empty witness list is a legal configuration (the validator's R2 witness mode); + // a witness call on such a client returns `NoProviderInRange` instead. + assert_eq!( + RpcClient::new(&[LOCALHOST_A], &[]).unwrap().witness_provider_count(), + 0, + "an empty witness list must construct" ); for endpoints in [&[LOCALHOST_B][..], &[LOCALHOST_B, "http://localhost:8547"]] { @@ -2116,6 +2152,32 @@ mod tests { hb.stop().unwrap(); } + /// A `skip` past the configured witness endpoints — or a client built with none at all — + /// is a wiring failure, and must fail this one request rather than take the process down. + #[tokio::test] + async fn witness_fetch_out_of_range_returns_a_typed_error() { + let client = RpcClient::new(&[LOCALHOST_A], &[LOCALHOST_B]).unwrap(); + let err = client + .get_witness_light_with_deadline_from(1, 7, B256::ZERO, None) + .await + .expect_err("skip == provider count leaves no provider"); + assert!( + matches!(err, WitnessFetchError::NoProviderInRange { skip: 1, configured: 1 }), + "unexpected error: {err:?}" + ); + + // Same variant covers the no-witness-providers deployment (R2 mode). + let witnessless = RpcClient::new(&[LOCALHOST_A], &[]).unwrap(); + let err = witnessless + .get_witness_light_with_deadline_from(0, 7, B256::ZERO, None) + .await + .expect_err("a client with no witness providers cannot fetch a witness"); + assert!( + matches!(err, WitnessFetchError::NoProviderInRange { skip: 0, configured: 0 }), + "unexpected error: {err:?}" + ); + } + /// `get_witness` pins `rr_start = 0`, so every round visits the primary first and only /// falls through to the backup on failure. We can't easily make the primary succeed in /// a unit test (a valid witness payload needs real cryptographic proof material), but @@ -2269,15 +2331,6 @@ mod tests { hc.stop().unwrap(); } - /// Skipping every configured witness provider is a caller bug and must panic loudly - /// instead of silently retrying over an empty provider set. - #[tokio::test] - #[should_panic(expected = "must select at least one")] - async fn test_witness_fetch_skip_of_all_providers_panics() { - let client = RpcClient::new(&[LOCALHOST_A], &[LOCALHOST_B]).unwrap(); - let _ = client.get_witness_light_with_deadline_from(1, 1, BlockHash::ZERO, None).await; - } - /// Serves `mega_getBlockWitness` returning a stub that decodes-fails, while recording /// the provider's label to a shared `order` log on each hit. Used to verify call routing. async fn start_ordered_witness_rpc( From cccd964ef2cc5596affa26161d9b45b290366026 Mon Sep 17 00:00:00 2001 From: "liquan.eth" Date: Sun, 20 Sep 2026 15:17:59 +0800 Subject: [PATCH 7/9] refactor(core): state the key-encoder rule by ownership, and make its test bite MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The encoder docs justified themselves as "for per-state-read hot paths", which asks the next contributor an unmeasurable question; the decidable one is whether the caller needs to own the key. encode() is now documented as the owned form derived from the borrowing primaries. stack_key_encodings_match_encode compared each encoder against encode(), which delegates to it — true by construction, so it killed no mutation. It now asserts the literal byte layout, which catches a swapped concatenation (verified: the mutation fails the new test and passes the old one). Co-Authored-By: Claude Opus 5 (1M context) --- crates/stateless-core/src/data_types.rs | 42 +++++++++++++++---------- 1 file changed, 25 insertions(+), 17 deletions(-) diff --git a/crates/stateless-core/src/data_types.rs b/crates/stateless-core/src/data_types.rs index 4a859424..511613a6 100644 --- a/crates/stateless-core/src/data_types.rs +++ b/crates/stateless-core/src/data_types.rs @@ -59,7 +59,10 @@ pub enum PlainKey { } impl PlainKey { - /// Encodes the key into a byte vector. + /// Encodes the key into a byte vector — the owned form, for callers that must *keep* the + /// key (a map key, a stored key). A lookup that only borrows the bytes for the duration of + /// the call should use [`Self::account_key_bytes`] / [`Self::storage_key_bytes`], which + /// this delegates to; the choice is about ownership, not about how hot the path is. /// /// # Returns /// - Account: 20-byte address @@ -75,8 +78,8 @@ impl PlainKey { /// Encoding of an account key — the raw address bytes. /// - /// Same bytes as `PlainKey::Account(address).encode()` without the heap allocation, - /// for per-state-read hot paths. + /// The primary definition of this encoding; [`Self::encode`] is the owned form derived + /// from it, so the two cannot disagree about what key a lookup targets. #[inline] pub(crate) fn account_key_bytes(address: &Address) -> &[u8] { address.as_slice() @@ -84,8 +87,8 @@ impl PlainKey { /// Stack-allocated encoding of a storage-slot key — address (20) ++ slot (32). /// - /// Same bytes as `PlainKey::Storage(address, slot).encode()` without the heap - /// allocation, for per-state-read hot paths. + /// The primary definition of this encoding; [`Self::encode`] is the owned form derived + /// from it, so the two cannot disagree about what key a lookup targets. #[inline] pub(crate) fn storage_key_bytes( address: Address, @@ -258,20 +261,25 @@ mod tests { entries.into_iter().enumerate().map(|(i, v)| (SaltKey::from((0u32, i as u64)), v)).collect() } - /// The allocation-free encoders the witness read path uses must produce exactly what - /// `encode()` produces — `encode()` delegates to them today, and this pins that contract - /// against a future edit that stops delegating and lets the two drift into looking up - /// different keys. + /// Both forms of each key encoding must produce one specific byte layout: the witness read + /// path uses the borrowing form and `encode()` the owning one, and they have to name the + /// same SALT entry. Each is asserted against the literal bytes rather than against the + /// other, so the pin still holds if `encode()` ever stops delegating — comparing the two + /// would pass by construction while they delegate, and prove nothing. #[test] fn stack_key_encodings_match_encode() { - let addr = Address::from([0x11; 20]); - assert_eq!(PlainKey::account_key_bytes(&addr), PlainKey::Account(addr).encode().as_slice()); - for slot in [B256::ZERO, B256::from([0xAB; 32]), B256::from(U256::from(7))] { - assert_eq!( - PlainKey::storage_key_bytes(addr, slot).as_slice(), - PlainKey::Storage(addr, slot).encode().as_slice(), - ); - } + let addr = Address::from([0x11; ACCOUNT_ADDRESS_LEN]); + let slot = B256::from([0xAB; SLOT_KEY_LEN]); + + let expected_account = [0x11u8; ACCOUNT_ADDRESS_LEN]; + assert_eq!(PlainKey::account_key_bytes(&addr), expected_account); + assert_eq!(PlainKey::Account(addr).encode(), expected_account); + + // Address first, then slot — a swapped concatenation would look up a different entry. + let mut expected_storage = [0xABu8; STORAGE_SLOT_KEY_LEN]; + expected_storage[..ACCOUNT_ADDRESS_LEN].fill(0x11); + assert_eq!(PlainKey::storage_key_bytes(addr, slot).as_slice(), expected_storage); + assert_eq!(PlainKey::Storage(addr, slot).encode(), expected_storage); } #[test] From a418e8e037fd069d94589b07d69ab1e118e3f3ab Mon Sep 17 00:00:00 2001 From: "liquan.eth" Date: Sun, 20 Sep 2026 16:29:17 +0800 Subject: [PATCH 8/9] refactor(common): check the witness skip where it enters, not in the rotation MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The range guard sat in the private `witness_round_robin`, which widened the error type of all three of its callers even though two build their range from the provider count and cannot be out of range. That widening is what forced the `deadline_only` adapter, its `unreachable!`, two `map_err` sites and a reworded `expect`. The check now lives in `get_witness_light_with_deadline_from`, the only method taking a caller-computed value. `witness_round_robin` keeps returning `RpcDeadlineExceeded`, and its precondition is a `debug_assert!` beside the one `round_robin_with_backoff` already uses. Public signatures are unchanged. Half the old guard was dead — `end > configured` is unreachable from every call site — and it misreported that case as a bad `skip`. With the range now built from the provider count at each site, the one remaining failure is `skip >= configured`, which the error's fields describe exactly. Co-Authored-By: Claude Opus 5 (1M context) --- bin/debug-trace-server/src/data_provider.rs | 8 +- crates/stateless-common/src/rpc_client.rs | 82 ++++++++++----------- 2 files changed, 43 insertions(+), 47 deletions(-) diff --git a/bin/debug-trace-server/src/data_provider.rs b/bin/debug-trace-server/src/data_provider.rs index ca218595..90c93887 100644 --- a/bin/debug-trace-server/src/data_provider.rs +++ b/bin/debug-trace-server/src/data_provider.rs @@ -3168,11 +3168,9 @@ mod tests { assert!(matches!(block_err, DataProviderError::Timeout { stage: TimeoutStage::Block, .. })); } - /// A witness fetch whose provider range is unsatisfiable is a wiring bug, not a blown - /// budget: it must land on `Internal`, never on `Timeout { Witness }`. That bucket feeds - /// the `deadline_witness` error reason, whose whole value is meaning "an upstream witness - /// fetch ran out of time" — a wiring bug landing there would page for the wrong incident. - /// The deadline variant still classifies by method, exactly as before. + /// A range failure is a wiring bug: it must land on `Internal`, never on the + /// `deadline_witness` alarm's `Timeout { Witness }` (rationale at the `From` impl). The + /// deadline arm is asserted too — it is the delegation that keeps that alarm working. #[test] fn witness_range_failure_is_internal_not_a_witness_timeout() { let range_err: DataProviderError = diff --git a/crates/stateless-common/src/rpc_client.rs b/crates/stateless-common/src/rpc_client.rs index 6e1a7d47..c67f4a4a 100644 --- a/crates/stateless-common/src/rpc_client.rs +++ b/crates/stateless-common/src/rpc_client.rs @@ -231,12 +231,10 @@ pub struct SetValidatedBlocksResponse { /// past the configured witness endpoints. The constructor rejects an empty endpoint list, so /// reaching this means a routing bug in the skip computation rather than a misconfiguration. /// It is a typed error rather than an `assert!` so such a bug fails one request instead of the -/// process. +/// process — which is the whole reason this type exists. #[derive(Debug, thiserror::Error)] pub enum WitnessFetchError { - #[error( - "witness fetch selected providers {skip}.. of {configured} configured — no witness provider in range" - )] + #[error("witness fetch skip={skip} leaves none of {configured} configured providers")] NoProviderInRange { skip: usize, configured: usize }, #[error(transparent)] Deadline(#[from] RpcDeadlineExceeded), @@ -308,8 +306,7 @@ impl RpcClient { /// # Arguments /// * `data_apis` - HTTP URLs of the standard JSON-RPC endpoints for blocks and contract data /// (tried in order, non-empty) - /// * `witness_apis` - HTTP URLs of the witness RPC endpoints (tried in order, non-empty). R2 - /// does not replace them: it is tried first and these remain the fallback + /// * `witness_apis` - HTTP URLs of the witness RPC endpoints (tried in order, non-empty) /// * `config` - Configuration controlling verification, retry, and concurrency behavior /// * `report_api` - Optional HTTP URL of the endpoint for reporting validated blocks pub fn new_with_config( @@ -731,8 +728,7 @@ impl RpcClient { decode_witness_response, "Witness decoded", ) - .await - .map_err(deadline_only)?; + .await?; if let Some(ref metrics) = self.config.metrics { metrics.on_witness_fetch(WitnessSizeBreakdown::new(&witness.0, &witness.1)); @@ -763,9 +759,15 @@ impl RpcClient { hash: B256, deadline: Option, ) -> std::result::Result<(LightWitness, MptWitness), RpcDeadlineExceeded> { - self.get_witness_light_with_deadline_from(0, number, hash, deadline) - .await - .map_err(deadline_only) + self.witness_round_robin( + 0..self.witness_providers.len(), + number, + hash, + deadline, + decode_witness_response_light, + "Witness light-decoded", + ) + .await } /// Like [`Self::get_witness_light_with_deadline`], but skips the first `skip` witness @@ -784,15 +786,23 @@ impl RpcClient { hash: B256, deadline: Option, ) -> std::result::Result<(LightWitness, MptWitness), WitnessFetchError> { - self.witness_round_robin( - skip..self.witness_providers.len(), - number, - hash, - deadline, - decode_witness_response_light, - "Witness light-decoded", - ) - .await + // Checked here, where the caller-computed value enters, rather than deeper in the + // rotation: every other witness fetch builds its range from the provider count and + // cannot be out of range, so this is the only place the check has anything to do. + let configured = self.witness_providers.len(); + if skip >= configured { + return Err(WitnessFetchError::NoProviderInRange { skip, configured }); + } + Ok(self + .witness_round_robin( + skip..configured, + number, + hash, + deadline, + decode_witness_response_light, + "Witness light-decoded", + ) + .await?) } /// Like [`Self::get_witness_light`], but consults only the FIRST witness provider — @@ -818,7 +828,7 @@ impl RpcClient { "Witness light-decoded", ) .await - .expect("pinned 0..1 range and a None deadline cannot fail") + .expect("None deadline cannot time out") } /// Shared `mega_getBlockWitness` retry loop: primary-failover rounds (always start from @@ -830,8 +840,9 @@ impl RpcClient { /// the logged endpoint labels stay aligned with the full configured list because each /// label bakes in its original index (see [`endpoint_label`]). /// - /// An empty or out-of-bounds `providers` range is a wiring failure, surfaced as - /// [`WitnessFetchError::NoProviderInRange`] rather than a panic. + /// Every caller builds `providers` from the configured provider count, so the range is + /// non-empty and in bounds by construction; the one caller-supplied value (`skip`) is + /// checked in [`Self::get_witness_light_with_deadline_from`] before it gets here. // A `warn`-level span (not the usual `info`) so it stays enabled at the default `warn` log // filter: the generic retry loop's per-attempt failure logs then inherit `block_number`, // which they cannot see otherwise, so an endpoint stall/error is traceable to its block. @@ -844,11 +855,12 @@ impl RpcClient { deadline: Option, decode: fn(&str) -> std::result::Result, trace_msg: &'static str, - ) -> std::result::Result { - let configured = self.witness_providers.len(); - if providers.is_empty() || providers.end > configured { - return Err(WitnessFetchError::NoProviderInRange { skip: providers.start, configured }); - } + ) -> std::result::Result { + debug_assert!( + !providers.is_empty() && providers.end <= self.witness_providers.len(), + "witness provider range ({providers:?}) must select at least one of {} providers", + self.witness_providers.len() + ); // Deadline-bound witness attempts run under the reserve-half policy: the tightest of // the configured ceiling, the general per-attempt timeout, and — recomputed at each // attempt, after any permit wait — half of what the call still has, so neither a @@ -884,7 +896,6 @@ impl RpcClient { }, ) .await - .map_err(WitnessFetchError::Deadline) } /// Reports a range of validated blocks via the dedicated report endpoint. @@ -1628,19 +1639,6 @@ async fn verify_block_on_blocking_pool(block: Block) -> Result RpcDeadlineExceeded { - match e { - WitnessFetchError::Deadline(d) => d, - WitnessFetchError::NoProviderInRange { skip, configured } => unreachable!( - "full-range witness fetch on a client with no witness providers \ - (skip={skip}, configured={configured})" - ), - } -} - /// Verifies structural integrity of a block fetched from RPC. /// /// Checks: From 90ae6a86d9e6fdbd2850b8601600f857e648846b Mon Sep 17 00:00:00 2001 From: "liquan.eth" Date: Wed, 23 Sep 2026 11:59:58 +0800 Subject: [PATCH 9/9] fix(trace-server): log a witness range failure as a routing bug, not a deadline fetch_witness's error arm predates WitnessFetchError: it logged every failure as "Witness fetch deadline exceeded" without the error, which was true while the only error was RpcDeadlineExceeded. NoProviderInRange now reaches the same arm, so a routing bug would surface as a witness timeout with a budget it never spent. Match the variant before logging: Deadline keeps the existing WARN, NoProviderInRange gets its own error! with skip and configured as fields. The returned error and the source metrics are unchanged. Addresses the Codex review on data_provider.rs (PRRT_kwDOPRLqcM6kKZdZ). --- bin/debug-trace-server/src/data_provider.rs | 44 ++++++++++++++------- 1 file changed, 30 insertions(+), 14 deletions(-) diff --git a/bin/debug-trace-server/src/data_provider.rs b/bin/debug-trace-server/src/data_provider.rs index 90c93887..80031ef6 100644 --- a/bin/debug-trace-server/src/data_provider.rs +++ b/bin/debug-trace-server/src/data_provider.rs @@ -53,7 +53,7 @@ use stateless_core::{ ContractStore, LightWitness, StoreResult, db::StoreError, withdrawals::MptWitness, }; use stateless_db::ContractCache; -use tracing::{debug, instrument, trace, warn}; +use tracing::{debug, error, instrument, trace, warn}; use crate::{ block_data_cache::BlockDataCache, @@ -1345,19 +1345,35 @@ async fn fetch_witness( } Err(e) => { metrics.record_request(false, start.elapsed().as_secs_f64()); - let budget = deadline.saturating_duration_since(start); - // Attribution-grade context for the next timeout incident: the effective stage - // budget, the route taken, and whether the old-block clamp applied — the client - // only ever sees the generic `-32001` message. - warn!( - block_number, - block_hash = %block_hash, - source, - old_block = is_old_block(db_tip, block_number), - budget_ms = budget.as_millis() as u64, - elapsed_ms = start.elapsed().as_millis() as u64, - "Witness fetch deadline exceeded", - ); + match &e { + WitnessFetchError::Deadline(_) => { + let budget = deadline.saturating_duration_since(start); + // Attribution-grade context for the next timeout incident: the effective + // stage budget, the route taken, and whether the old-block clamp applied — + // the client only ever sees the generic `-32001` message. + warn!( + block_number, + block_hash = %block_hash, + source, + old_block = is_old_block(db_tip, block_number), + budget_ms = budget.as_millis() as u64, + elapsed_ms = start.elapsed().as_millis() as u64, + "Witness fetch deadline exceeded", + ); + } + // A routing bug, not a timeout: no upstream attempt ran, so the deadline + // warning and its budget fields would misattribute it. + WitnessFetchError::NoProviderInRange { skip, configured } => { + error!( + block_number, + block_hash = %block_hash, + source, + skip, + configured, + "Witness route selected no configured provider", + ); + } + } Err(e.into()) } }