diff --git a/Cargo.lock b/Cargo.lock index 2ec67af9..605c6c05 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5694,6 +5694,7 @@ dependencies = [ "salt", "serde", "stateless-core", + "stateless-r2", "stateless-test-utils", "thiserror 2.0.17", "tokio", @@ -5784,11 +5785,13 @@ dependencies = [ name = "stateless-test-utils" version = "2.0.16" dependencies = [ + "alloy-consensus", "alloy-genesis", "alloy-primitives", "alloy-rpc-types-eth", "bincode 2.0.1", "eyre", + "jsonrpsee", "op-alloy-rpc-types", "revm", "salt", diff --git a/bin/debug-trace-server/Cargo.toml b/bin/debug-trace-server/Cargo.toml index aa3d704f..a1e01edf 100644 --- a/bin/debug-trace-server/Cargo.toml +++ b/bin/debug-trace-server/Cargo.toml @@ -78,4 +78,4 @@ tokio = { workspace = true, features = ["macros", "rt-multi-thread", "time"] } zstd.workspace = true # stateless -stateless-test-utils = { path = "../../crates/stateless-test-utils" } +stateless-test-utils = { path = "../../crates/stateless-test-utils", features = ["mock-rpc"] } diff --git a/bin/debug-trace-server/src/data_provider.rs b/bin/debug-trace-server/src/data_provider.rs index 7e8b03c0..23ed6da7 100644 --- a/bin/debug-trace-server/src/data_provider.rs +++ b/bin/debug-trace-server/src/data_provider.rs @@ -1511,8 +1511,11 @@ impl ContractStore for NoopContractStore { pub(crate) mod test_support { use std::sync::atomic::{AtomicUsize, Ordering}; - use jsonrpsee::{RpcModule, server::ServerHandle, types::ErrorObjectOwned}; - use stateless_test_utils::fixtures::TestFixtures; + use jsonrpsee::{server::ServerHandle, types::ErrorObjectOwned}; + use stateless_test_utils::{ + fixtures::TestFixtures, + mock_rpc::{consistent_header, serve}, + }; use super::*; @@ -1555,13 +1558,6 @@ pub(crate) mod test_support { (number, hash, wire, LightWitness::from(salt)) } - /// Minimal self-consistent RPC `Header` for `number`: `hash` is the real `hash_slow()` - /// of the inner header, so `verify_hash = true` fetches accept it. - pub(crate) fn consistent_header(number: u64) -> alloy_rpc_types_eth::Header { - let inner = alloy_consensus::Header { number, ..Default::default() }; - alloy_rpc_types_eth::Header { hash: inner.hash_slow(), inner, ..Default::default() } - } - /// Per-method call counters for [`start_mock_rpc`], so round-trip-shape tests can /// assert which upstream calls a path made (and which it skipped). #[derive(Default)] @@ -1574,14 +1570,6 @@ pub(crate) mod test_support { pub(crate) header_by_hash: AtomicUsize, } - /// Boots a jsonrpsee server for `module` on an ephemeral port. - async fn serve(module: RpcModule) -> (ServerHandle, String) { - let server = - jsonrpsee::server::ServerBuilder::default().build("127.0.0.1:0").await.unwrap(); - let url = format!("http://{}", server.local_addr().unwrap()); - (server.start(module), url) - } - /// The simulated generator's "witness not generated yet" miss, shared by every mock /// witness source so the tests model one upstream wire shape. fn witness_not_generated() -> ErrorObjectOwned { @@ -1594,33 +1582,34 @@ pub(crate) mod test_support { /// call per method in [`MockRpcHits`]. pub(crate) async fn start_mock_rpc(tip: u64) -> (ServerHandle, String, Arc) { let hits = Arc::new(MockRpcHits::default()); - let mut module = RpcModule::new(hits.clone()); - module - .register_method("eth_getHeaderByNumber", move |params, hits, _| { - hits.header_by_number.fetch_add(1, Ordering::Relaxed); - // Parse as the client's own wire type so the mock can't drift from it. - let (tag,): (BlockNumberOrTag,) = params.parse().unwrap(); - let n = tag.as_number().unwrap_or(tip); - Ok::<_, ErrorObjectOwned>(consistent_header(n)) - }) - .unwrap(); - module - .register_method("eth_getHeaderByHash", move |params, hits, _| { - hits.header_by_hash.fetch_add(1, Ordering::Relaxed); - // Echo the requested hash (the client cross-checks it against the - // request); the content is a `tip` header, enough for number discovery. - let (hash,): (B256,) = params.parse().unwrap(); - let header = alloy_rpc_types_eth::Header { hash, ..consistent_header(tip) }; - Ok::<_, ErrorObjectOwned>(header) - }) - .unwrap(); - module - .register_method("eth_blockNumber", move |_params, hits, _| { - hits.block_number.fetch_add(1, Ordering::Relaxed); - Ok::<_, ErrorObjectOwned>(format!("{tip:#x}")) - }) - .unwrap(); - let (handle, url) = serve(module).await; + let (handle, url) = serve(hits.clone(), |module| { + module + .register_method("eth_getHeaderByNumber", move |params, hits, _| { + hits.header_by_number.fetch_add(1, Ordering::Relaxed); + // Parse as the client's own wire type so the mock can't drift from it. + let (tag,): (BlockNumberOrTag,) = params.parse().unwrap(); + let n = tag.as_number().unwrap_or(tip); + Ok::<_, ErrorObjectOwned>(consistent_header(n)) + }) + .unwrap(); + module + .register_method("eth_getHeaderByHash", move |params, hits, _| { + hits.header_by_hash.fetch_add(1, Ordering::Relaxed); + // Echo the requested hash (the client cross-checks it against the + // request); the content is a `tip` header, enough for number discovery. + let (hash,): (B256,) = params.parse().unwrap(); + let header = alloy_rpc_types_eth::Header { hash, ..consistent_header(tip) }; + Ok::<_, ErrorObjectOwned>(header) + }) + .unwrap(); + module + .register_method("eth_blockNumber", move |_params, hits, _| { + hits.block_number.fetch_add(1, Ordering::Relaxed); + Ok::<_, ErrorObjectOwned>(format!("{tip:#x}")) + }) + .unwrap(); + }) + .await; (handle, url, hits) } @@ -1632,17 +1621,18 @@ pub(crate) mod test_support { wire: Option, ) -> (ServerHandle, String, Arc) { let hits = Arc::new(AtomicUsize::new(0)); - let mut module = RpcModule::new((hits.clone(), wire)); - module - .register_method("mega_getBlockWitness", move |_p, (hits, wire), _| { - let call = hits.fetch_add(1, Ordering::Relaxed); - match wire { - Some(wire) if call >= misses_before_serve => Ok(wire.clone()), - _ => Err(witness_not_generated()), - } - }) - .unwrap(); - let (handle, url) = serve(module).await; + let (handle, url) = serve((hits.clone(), wire), |module| { + module + .register_method("mega_getBlockWitness", move |_p, (hits, wire), _| { + let call = hits.fetch_add(1, Ordering::Relaxed); + match wire { + Some(wire) if call >= misses_before_serve => Ok(wire.clone()), + _ => Err(witness_not_generated()), + } + }) + .unwrap(); + }) + .await; (handle, url, hits) } @@ -1655,18 +1645,19 @@ pub(crate) mod test_support { block: Block, wire: Option, ) -> (ServerHandle, String) { - let mut module = RpcModule::new((block, wire)); - module - .register_method("eth_getBlockByHash", |_p, (block, _), _| { - Ok::<_, ErrorObjectOwned>(block.clone()) - }) - .unwrap(); - module - .register_method("mega_getBlockWitness", |_p, (_, wire), _| { - wire.clone().ok_or_else(witness_not_generated) - }) - .unwrap(); - serve(module).await + serve((block, wire), |module| { + module + .register_method("eth_getBlockByHash", |_p, (block, _), _| { + Ok::<_, ErrorObjectOwned>(block.clone()) + }) + .unwrap(); + module + .register_method("mega_getBlockWitness", |_p, (_, wire), _| { + wire.clone().ok_or_else(witness_not_generated) + }) + .unwrap(); + }) + .await } /// Serves `wire` after `delay` on every call — a healthy-but-slow witness endpoint. @@ -1675,18 +1666,17 @@ pub(crate) mod test_support { wire: String, ) -> (ServerHandle, String, Arc) { let hits = Arc::new(AtomicUsize::new(0)); - let mut module = RpcModule::new((hits.clone(), wire)); - module - .register_async_method("mega_getBlockWitness", move |_p, ctx, _| async move { - ctx.0.fetch_add(1, Ordering::Relaxed); - tokio::time::sleep(delay).await; - Ok::<_, ErrorObjectOwned>(ctx.1.clone()) - }) - .unwrap(); - let server = - jsonrpsee::server::ServerBuilder::default().build("127.0.0.1:0").await.unwrap(); - let url = format!("http://{}", server.local_addr().unwrap()); - (server.start(module), url, hits) + let (handle, url) = serve((hits.clone(), wire), |module| { + module + .register_async_method("mega_getBlockWitness", move |_p, ctx, _| async move { + ctx.0.fetch_add(1, Ordering::Relaxed); + tokio::time::sleep(delay).await; + Ok::<_, ErrorObjectOwned>(ctx.1.clone()) + }) + .unwrap(); + }) + .await; + (handle, url, hits) } } @@ -1698,12 +1688,11 @@ mod tests { use stateless_test_utils::{ fixtures::TestFixtures, mock_r2::{mock_r2, mock_r2_held}, + mock_rpc::consistent_header, }; use super::{ - test_support::{ - block_and_witness_rpc, consistent_header, scripted_witness_rpc, start_mock_rpc, - }, + test_support::{block_and_witness_rpc, scripted_witness_rpc, start_mock_rpc}, *, }; use crate::server_db::test_support::StubBlockStore; diff --git a/bin/debug-trace-server/src/main.rs b/bin/debug-trace-server/src/main.rs index ecba45ea..8b6712d2 100644 --- a/bin/debug-trace-server/src/main.rs +++ b/bin/debug-trace-server/src/main.rs @@ -57,8 +57,8 @@ use clap::Parser; use eyre::Result; use jsonrpsee::server::{Server, ServerConfig, middleware::rpc::RpcServiceBuilder}; use stateless_common::{ - R2CountFlag, R2Flag, R2Flags, R2Target, R2TuningFlag, RedactedSecret, RpcClient, - RpcClientConfig, logging::LogArgs, validate_r2_flags, + R2CountFlag, R2Flag, R2Flags, R2Target, R2TuningFlag, R2WitnessTransport, RedactedSecret, + RpcClient, RpcClientConfig, logging::LogArgs, validate_r2_flags, }; use stateless_core::{ BisectResolver, ChainStore, ContractStore, DivergenceLookups, PipelineConfig, @@ -920,27 +920,28 @@ async fn main() -> Result<()> { }, ); let cf_access = access.is_some(); - let source = R2WitnessSource::new_custom_domain( + let transport = R2WitnessTransport::new_custom_domain( domain, access, r2_timeouts, rpc_retry, args.r2_max_concurrent_requests, connections, + metrics::record_r2_negotiated_version, )?; - metrics::record_r2_target(source.target_label()); - metrics::record_r2_connections(source.connections()); + metrics::record_r2_target(transport.target_label()); + metrics::record_r2_connections(transport.connections()); info!( - domain = %source.origin(), + domain = %transport.origin(), cf_access, - connections = source.connections(), + connections = transport.connections(), "Historical witness source: R2 (custom domain), RPC chain as fallback" ); - Some(source) + Some(R2WitnessSource::new(transport)) } R2Target::S3 => { let take = |v: &Option| v.clone().expect("S3 target"); - let source = R2WitnessSource::new( + let transport = R2WitnessTransport::new( args.r2_endpoint.as_deref().expect("S3 target"), take(&args.r2_bucket), take(&args.r2_access_key_id), @@ -949,13 +950,13 @@ async fn main() -> Result<()> { rpc_retry, args.r2_max_concurrent_requests, )?; - metrics::record_r2_target(source.target_label()); + metrics::record_r2_target(transport.target_label()); info!( - endpoint = %source.origin(), + endpoint = %transport.origin(), bucket = args.r2_bucket.as_deref().unwrap_or_default(), "Historical witness source: R2 (direct S3), RPC chain as fallback" ); - Some(source) + Some(R2WitnessSource::new(transport)) } }; let r2_witness_source = r2_source.map(Arc::new); diff --git a/bin/debug-trace-server/src/r2_witness.rs b/bin/debug-trace-server/src/r2_witness.rs index 5060dfec..1440b5a9 100644 --- a/bin/debug-trace-server/src/r2_witness.rs +++ b/bin/debug-trace-server/src/r2_witness.rs @@ -4,25 +4,24 @@ //! GETs or unsigned GETs through a Cloudflare custom domain, per construction — and decodes //! it with the **light** decoder — the trace server never verifies the witness proof, so the //! full decode's per-point elliptic-curve work would buy nothing (see -//! `stateless_core::light_witness`). The transport core is `stateless-r2`'s -//! [`R2ObjectFetcher`], shared with the validator's pipeline reader. +//! `stateless_core::light_witness`). The failure taxonomy and the transport wrapper are +//! [`stateless_common::r2_witness`], shared with the validator's adapter; the transport +//! core below that is `stateless-r2`'s [`R2ObjectFetcher`]. //! //! This adapter is request-serving, which shapes it differently from the validator's: //! every fetch runs under the caller's witness-stage deadline, failures surface immediately //! with **no pacing pause** (the caller's next move is the RPC fallback chain, not a blind //! re-enqueue), and the retry budget is small — a throttled R2 should hand over to the RPC //! chain quickly instead of burning the witness budget on backoff sleeps. +//! +//! [`R2ObjectFetcher`]: stateless_r2::fetch::R2ObjectFetcher use std::time::Instant; use alloy_primitives::B256; -use stateless_common::{BackoffPolicy, WitnessDecodingError, decode_witness_payload_light}; +pub use stateless_common::R2WitnessError; +use stateless_common::{R2WitnessTransport, decode_on_blocking_pool, decode_witness_payload_light}; use stateless_core::{LightWitness, withdrawals::MptWitness}; -use stateless_r2::{ - fetch::{CfAccessCredentials, FetchTimeouts, R2GetError, R2ObjectFetcher, RetryPacing}, - keys, -}; -use tokio::task::JoinError; use tracing::trace; use crate::metrics; @@ -38,140 +37,18 @@ const MAX_ATTEMPTS: usize = 3; /// below-band `kind="missing"` bucket-integrity alarm. pub(crate) const KIND_MISSING_ABOVE_TIP: &str = "missing_above_tip"; -/// Failure outcome of an R2 witness fetch. -#[derive(Debug, thiserror::Error)] -pub enum R2WitnessError { - /// The GET failed (absent object, transport, throttle, unexpected status, or out of - /// deadline while queued). - #[error(transparent)] - Get(#[from] R2GetError), - /// The object was fetched but its bytes did not decode to a witness tuple — a corrupt - /// witness in R2. Deterministic; not retried. - #[error("R2 witness for block {number} (key {key}) failed to decode: {source}")] - Decode { number: u64, key: String, source: WitnessDecodingError }, - /// The decode outran what was left of the caller's deadline — an oversized or - /// pathological object. The caller falls back with its reserved share of the stage; - /// the blocking decode itself cannot be cancelled and finishes in the background. - #[error("R2 witness decode for block {number} (key {key}) outran the deadline")] - DecodeTimeout { number: u64, key: String }, - /// The decode task panicked. This is a bug in our own decoder, not a problem with the - /// data in R2, so it is kept out of [`Self::Decode`]. - #[error("R2 witness decode task for block {number} (key {key}) panicked: {source}")] - DecodePanicked { number: u64, key: String, source: JoinError }, -} - -impl R2WitnessError { - /// Every label [`Self::kind`] can produce, for metrics pre-registration. - pub const KINDS: &'static [&'static str] = &[ - "missing", - "transport", - "throttled", - "status", - "connect", - "deadline", - "decode", - "decode_timeout", - "decode_panicked", - ]; - - /// Stable lowercase label for this variant — the `kind` label on the R2 witness error - /// counter. Every value returned here must appear in [`Self::KINDS`]. - pub const fn kind(&self) -> &'static str { - match self { - Self::Get(e) => e.kind(), - Self::Decode { .. } => "decode", - Self::DecodeTimeout { .. } => "decode_timeout", - Self::DecodePanicked { .. } => "decode_panicked", - } - } - - /// Whether the object was absent from the bucket — the one failure the frontier probe - /// treats as expected rather than alarming. - pub(crate) const fn is_missing(&self) -> bool { - matches!(self, Self::Get(R2GetError::Missing { .. })) - } -} - /// Fetches and light-decodes witnesses straight from an R2 bucket. -/// The fetcher's `Debug` redacts the credentials. +/// The transport's `Debug` redacts the credentials. #[derive(Debug)] pub struct R2WitnessSource { - fetcher: R2ObjectFetcher, -} - -/// The fetcher's pacing view of a `BackoffPolicy` — the adapter-layer conversion that keeps -/// `stateless-r2` free of a dependency on this workspace's backoff type. -fn pacing(backoff: &BackoffPolicy) -> RetryPacing { - RetryPacing { initial: backoff.initial, max: backoff.max } + transport: R2WitnessTransport, } impl R2WitnessSource { - /// The configured target's origin, for startup logging (see [`R2ObjectFetcher::origin`]). - pub fn origin(&self) -> &str { - self.fetcher.origin() - } - - /// The configured target's metric label (see [`R2ObjectFetcher::target_label`]). - pub const fn target_label(&self) -> &'static str { - self.fetcher.target_label() - } - - /// How many HTTP/2 connections the transport spreads its GETs over, for startup logging - /// (see [`R2ObjectFetcher::connections`]). - pub fn connections(&self) -> usize { - self.fetcher.connections() - } - - /// Builds a source from an R2 endpoint origin, bucket, and bucket-scoped S3 credentials. - /// - /// `timeouts` bounds each individual GET end-to-end and in its connect phase (further - /// clamped by the caller's deadline), `retry_backoff` paces the retries, and - /// `max_concurrent_requests` caps in-flight GETs (`None` = unlimited; see the - /// `--r2-max-concurrent-requests` flag for why it is separate from the RPC cap). - pub fn new( - endpoint: &str, - bucket: String, - access_key_id: String, - secret_access_key: String, - timeouts: FetchTimeouts, - retry_backoff: BackoffPolicy, - max_concurrent_requests: Option, - ) -> eyre::Result { - let fetcher = R2ObjectFetcher::new( - endpoint, - bucket, - access_key_id, - secret_access_key, - timeouts, - pacing(&retry_backoff), - max_concurrent_requests, - ) - .map_err(|e| eyre::eyre!(e))?; - Ok(Self { fetcher }) - } - - /// Builds a source that fetches unsigned through a Cloudflare custom domain fronting - /// the bucket (h2-multiplexed, edge-cacheable), with optional Cloudflare Access - /// service-token headers. The remaining parameters mean what they mean on [`Self::new`]. - pub fn new_custom_domain( - domain: &str, - access: Option, - timeouts: FetchTimeouts, - retry_backoff: BackoffPolicy, - max_concurrent_requests: Option, - connections: usize, - ) -> eyre::Result { - let fetcher = R2ObjectFetcher::new_custom_domain( - domain, - access, - timeouts, - pacing(&retry_backoff), - max_concurrent_requests, - connections, - ) - .map(|fetcher| fetcher.on_version_observed(metrics::record_r2_negotiated_version)) - .map_err(|e| eyre::eyre!(e))?; - Ok(Self { fetcher }) + /// Wraps an already-built transport. Construction (and the startup logging that reads + /// the configured target off it) lives at the wiring site, which owns the flags. + pub const fn new(transport: R2WitnessTransport) -> Self { + Self { transport } } /// Fetches and light-decodes the witness for `(number, hash)` under `deadline`. @@ -182,7 +59,8 @@ impl R2WitnessSource { deadline: Instant, ) -> Result<(LightWitness, MptWitness), R2WitnessError> { let fetched = self - .fetcher + .transport + .fetcher() .get_block_object( number, hash, @@ -196,36 +74,15 @@ impl R2WitnessSource { // reported on its own series instead of being subtracted the way the validator's // throughput pipeline does. metrics::record_r2_witness_queue_wait(fetched.queue_wait.as_secs_f64()); - decode_light_with_deadline(fetched.bytes, number, hash, deadline).await - } -} - -/// Light-decodes `bytes` on the blocking pool, bounded by the same `deadline` as the GET — -/// an oversized or pathological object must not eat the RPC fallback's share of the stage. -/// On timeout the blocking task is abandoned (it cannot be cancelled) and finishes in the -/// background. -async fn decode_light_with_deadline( - bytes: bytes::Bytes, - number: u64, - hash: B256, - deadline: Instant, -) -> Result<(LightWitness, MptWitness), R2WitnessError> { - let key = || keys::block_object_key(number, hash); - // A GET that lands right at the deadline gets no decode at all — nothing would wait - // for it. - if Instant::now() >= deadline { - return Err(R2WitnessError::DecodeTimeout { number, key: key() }); - } - // zstd + bincode over a multi-MB witness is CPU-bound; keep it off the runtime. - let decode = tokio::task::spawn_blocking(move || decode_witness_payload_light(&bytes)); - match tokio::time::timeout_at(deadline.into(), decode).await { - Ok(Ok(Ok(witness))) => { - trace!(number, "R2 witness fetched and light-decoded"); - Ok(witness) - } - Ok(Ok(Err(source))) => Err(R2WitnessError::Decode { number, key: key(), source }), - Ok(Err(source)) => Err(R2WitnessError::DecodePanicked { number, key: key(), source }), - Err(_) => Err(R2WitnessError::DecodeTimeout { number, key: key() }), + // The decode is bounded by the same deadline as the GET: an oversized or + // pathological object must not eat the RPC fallback's share of the stage. + let witness = + decode_on_blocking_pool(fetched.bytes, number, hash, Some(deadline), |bytes| { + decode_witness_payload_light(bytes) + }) + .await?; + trace!(number, "R2 witness fetched and light-decoded"); + Ok(witness) } } @@ -233,11 +90,15 @@ async fn decode_light_with_deadline( pub(crate) mod test_support { use std::time::Duration; + use stateless_common::BackoffPolicy; + use stateless_r2::fetch::FetchTimeouts; + use super::*; - /// Test source pointed at a mock endpoint, with millisecond retry pacing. - pub(crate) fn source(endpoint: &str) -> R2WitnessSource { - R2WitnessSource::new( + /// Millisecond-scale retry pacing and a 5s per-GET budget, so the mock-backed tests run + /// fast (production passes the policy built from the `--rpc-*-backoff-ms` flags). + pub(crate) fn test_transport(endpoint: &str) -> R2WitnessTransport { + R2WitnessTransport::new( endpoint, "witness-test".to_string(), "ak".to_string(), @@ -251,12 +112,18 @@ pub(crate) mod test_support { ) .unwrap() } + + /// Test source pointed at a mock endpoint. + pub(crate) fn source(endpoint: &str) -> R2WitnessSource { + R2WitnessSource::new(test_transport(endpoint)) + } } #[cfg(test)] mod tests { use std::{sync::atomic::Ordering, time::Duration}; + use stateless_r2::fetch::R2GetError; use stateless_test_utils::{fixtures::TestFixtures, mock_r2::mock_r2}; use super::{test_support::source, *}; @@ -276,19 +143,23 @@ mod tests { let (domain, _, heads) = stateless_test_utils::mock_r2::mock_r2_capturing(vec![(200, payload)]).await; - let source = R2WitnessSource::new_custom_domain( + let transport = R2WitnessTransport::new_custom_domain( &domain, None, - FetchTimeouts { + stateless_r2::fetch::FetchTimeouts { per_attempt: Duration::from_secs(5), connect: stateless_r2::fetch::DEFAULT_CONNECT_TIMEOUT, }, - BackoffPolicy::new(Duration::from_millis(5), Duration::from_millis(20)), + stateless_common::BackoffPolicy::new( + Duration::from_millis(5), + Duration::from_millis(20), + ), None, 1, + metrics::record_r2_negotiated_version, ) .unwrap(); - source + R2WitnessSource::new(transport) .get_witness_light(1, B256::ZERO, deadline()) .await .expect("valid object must fetch and decode"); @@ -344,13 +215,6 @@ mod tests { assert_eq!(hits.load(Ordering::SeqCst), MAX_ATTEMPTS); } - /// Every fetch-level kind must appear in this adapter's pre-registered [`KINDS`] — a new - /// `R2GetError` kind escaping metric pre-registration would drift silently otherwise. - #[test] - fn kinds_cover_all_fetch_kinds() { - assert!(R2GetError::KINDS.iter().all(|k| R2WitnessError::KINDS.contains(k))); - } - #[tokio::test] async fn undecodable_body_surfaces_decode_error() { let (endpoint, _) = mock_r2(vec![(200, "not a zstd witness")]).await; @@ -361,7 +225,7 @@ mod tests { /// A decode that outruns the deadline is abandoned so the caller falls back on its /// reserved budget share, instead of the object holding the witness stage hostage. - /// Driven through the extracted decode step with the deadline already gone — the + /// Driven through the shared decode step with the deadline already gone — the /// GET-succeeds-then-decode-overruns timing cannot be scripted deterministically. #[tokio::test] async fn decode_past_the_deadline_surfaces_decode_timeout() { @@ -370,10 +234,12 @@ mod tests { let (_, payload) = stateless_common::encode_witness_payload(&salt_witness, &mpt_witness) .expect("fixture witness must encode"); - let err = decode_light_with_deadline(payload.into(), 1, B256::ZERO, Instant::now()) - .await - .expect_err("an already-elapsed deadline must abandon the decode"); - assert!(matches!(err, R2WitnessError::DecodeTimeout { .. }), "{err}"); - assert_eq!(err.kind(), "decode_timeout"); + let err = decode_on_blocking_pool(payload, 1, B256::ZERO, Some(Instant::now()), |bytes| { + decode_witness_payload_light(bytes) + }) + .await + .expect_err("an already-elapsed deadline must abandon the decode") + .to_string(); + assert!(err.contains("outran the deadline"), "{err}"); } } diff --git a/bin/stateless-validator/Cargo.toml b/bin/stateless-validator/Cargo.toml index b89a315e..51beca81 100644 --- a/bin/stateless-validator/Cargo.toml +++ b/bin/stateless-validator/Cargo.toml @@ -51,7 +51,7 @@ jsonrpsee-types.workspace = true tempfile.workspace = true # stateless -stateless-test-utils = { path = "../../crates/stateless-test-utils" } +stateless-test-utils = { path = "../../crates/stateless-test-utils", features = ["mock-rpc"] } [features] test-bucket-resize = ["salt/test-bucket-resize", "stateless-core/test-bucket-resize"] diff --git a/bin/stateless-validator/src/app.rs b/bin/stateless-validator/src/app.rs index 29faec83..90a2ac3e 100644 --- a/bin/stateless-validator/src/app.rs +++ b/bin/stateless-validator/src/app.rs @@ -8,8 +8,8 @@ use alloy_rpc_types_eth::BlockId; use clap::{Parser, ValueEnum}; use eyre::Result; use stateless_common::{ - BackoffPolicy, R2CountFlag, R2Flag, R2Flags, R2Target, RedactedSecret, RpcClient, - RpcClientConfig, logging::LogArgs, validate_r2_flags, + BackoffPolicy, R2CountFlag, R2Flag, R2Flags, R2Target, R2WitnessTransport, RedactedSecret, + RpcClient, RpcClientConfig, logging::LogArgs, validate_r2_flags, }; use stateless_core::{ChainStore, ContractStore, chain_spec::ChainSpec, db::BlockMeta}; use stateless_db::ContractCache; @@ -356,8 +356,8 @@ pub async fn run() -> Result<()> { .r2_connect_timeout_ms .map_or(stateless_r2::fetch::DEFAULT_CONNECT_TIMEOUT, Duration::from_millis), }; - let client = build_r2_client(&args, timeouts, rpc_config.rpc_retry.clone())?; - Some(Arc::new(client)) + let transport = build_r2_transport(&args, timeouts, rpc_config.rpc_retry.clone())?; + Some(Arc::new(R2WitnessClient::new(transport))) } }; @@ -457,17 +457,17 @@ fn override_ms(ms: Option, default: Duration) -> Duration { ms.map(Duration::from_millis).unwrap_or(default) } -/// Builds the R2 witness client for `--witness-source r2`: the custom-domain target when +/// Builds the R2 witness transport for `--witness-source r2`: the custom-domain target when /// `--r2-custom-domain` is set, the SigV4-signed S3 target otherwise. /// /// Which target wins is already settled by the [`validate_r2_flags`] call below, so the arms /// read the one that was chosen — a set-but-empty flag belonging to the *other* target is /// rejected there rather than reaching a constructor. -fn build_r2_client( +fn build_r2_transport( args: &CommandLineArgs, timeouts: stateless_r2::fetch::FetchTimeouts, retry: BackoffPolicy, -) -> Result { +) -> Result { // `--witness-max-concurrent-requests` capped R2 GETs too before the caps were split. // Refuse the pre-split spelling by name rather than leave R2 uncapped: this mode has no // RPC fallback, so an uncapped fetcher aims its whole in-flight window at the bucket. @@ -485,7 +485,7 @@ fn build_r2_client( // Every coherence rule lives in the shared validator, so the reads below rest on an // invariant that was actually checked: no empty values, exactly one target, and an Access // pair that is either whole or absent. - let client = match validate_r2_flags(&r2_flags(args))? { + let transport = match validate_r2_flags(&r2_flags(args))? { R2Target::None => { return Err(eyre::eyre!( "--witness-source r2 needs an R2 target: configure --r2-custom-domain, or \ @@ -502,27 +502,28 @@ fn build_r2_client( }, ); let cf_access = access.is_some(); - let client = R2WitnessClient::new_custom_domain( + let transport = R2WitnessTransport::new_custom_domain( domain, access, timeouts, retry, args.r2_max_concurrent_requests, connections, + metrics::record_r2_negotiated_version, )?; - metrics::record_r2_connections(client.connections()); + metrics::record_r2_connections(transport.connections()); info!( - domain = %client.origin(), + domain = %transport.origin(), cf_access, - connections = client.connections(), - max_concurrent_requests = ?client.max_concurrent_requests(), + connections = transport.connections(), + max_concurrent_requests = ?transport.max_concurrent_requests(), "Witness source: R2 (custom domain)" ); - client + transport } R2Target::S3 => { let take = |v: &Option| v.clone().expect("S3 target"); - let client = R2WitnessClient::new( + let transport = R2WitnessTransport::new( args.r2_endpoint.as_deref().expect("S3 target"), take(&args.r2_bucket), take(&args.r2_access_key_id), @@ -532,16 +533,16 @@ fn build_r2_client( args.r2_max_concurrent_requests, )?; info!( - endpoint = %client.origin(), + endpoint = %transport.origin(), bucket = args.r2_bucket.as_deref().unwrap_or_default(), - max_concurrent_requests = ?client.max_concurrent_requests(), + max_concurrent_requests = ?transport.max_concurrent_requests(), "Witness source: R2 (direct S3)" ); - client + transport } }; - metrics::record_r2_target(client.target_label()); - Ok(client) + metrics::record_r2_target(transport.target_label()); + Ok(transport) } /// This binary's `--r2-*` flags, in the spellings its operators use. @@ -595,7 +596,7 @@ mod tests { "secret", ]; - /// Argv for `--witness-source r2` with the given target flags, so [`build_r2_client`] — + /// Argv for `--witness-source r2` with the given target flags, so [`build_r2_transport`] — /// the seam every R2-mode rule is gated behind — runs the rules from the path production /// takes. fn parse_r2_with_target(target: &[&str], extra: &[&str]) -> CommandLineArgs { @@ -615,20 +616,20 @@ mod tests { parse_r2_with_target(CUSTOM_DOMAIN_TARGET, extra) } - fn build(args: &CommandLineArgs) -> Result { + fn build(args: &CommandLineArgs) -> Result { let timeouts = stateless_r2::fetch::FetchTimeouts { per_attempt: Duration::from_secs(1), connect: Duration::from_secs(1), }; let retry = BackoffPolicy { initial: Duration::from_millis(1), max: Duration::from_millis(1) }; - build_r2_client(args, timeouts, retry) + build_r2_transport(args, timeouts, retry) } /// Carrying the pre-split spelling of the R2 concurrency cap into `--witness-source r2` /// must fail by name rather than leave R2 uncapped: that mode has no RPC fallback, so an /// uncapped fetcher aims its whole in-flight window at the bucket. Outside r2 mode the - /// rule is unreachable by construction — `build_r2_client` is only called from the + /// rule is unreachable by construction — `build_r2_transport` is only called from the /// `WitnessSource::R2` arm, the same call-site gating as every other R2 rule. #[test] fn r2_mode_refuses_the_pre_split_concurrency_spelling() { @@ -662,11 +663,11 @@ mod tests { /// The migration guard fires on the flags alone, so its test above would still pass with /// the constructors wired to the old field. This is the assertion that observes which cap - /// actually reaches the client — with both spellings set it must be the R2 one, not the - /// RPC one — on both target arms, plus the uncapped default, so a revert of either arm's - /// wiring fails here by value. + /// actually reaches the transport the fetcher runs on — with both spellings set it must be + /// the R2 one, not the RPC one — on both target arms, plus the uncapped default, so a + /// revert of either arm's wiring fails here by value. #[test] - fn the_r2_cap_not_the_rpc_one_reaches_the_client() { + fn the_r2_cap_not_the_rpc_one_reaches_the_transport() { let _guard = stateless_test_utils::env::env_lock(); for target in [CUSTOM_DOMAIN_TARGET, S3_TARGET] { diff --git a/bin/stateless-validator/src/r2_witness.rs b/bin/stateless-validator/src/r2_witness.rs index fdec1685..ed7dd04a 100644 --- a/bin/stateless-validator/src/r2_witness.rs +++ b/bin/stateless-validator/src/r2_witness.rs @@ -2,12 +2,13 @@ //! //! Fetches the primary witness object straight from the R2 bucket — via SigV4-signed S3 GETs //! or unsigned GETs through a Cloudflare custom domain, per construction — and returns -//! the same `(SaltWitness, MptWitness)` tuple the RPC path yields. The transport core is -//! [`R2ObjectFetcher`] from `stateless-r2`, shared with the debug-trace-server's historical -//! witness source; this adapter owns what is validator-specific: the **full** payload decode -//! (proof verification needs the elliptic-curve points the light decode skips), the validator -//! metrics, and the surfaced-failure pacing the pipeline fetcher relies on. The object body is -//! `zstd(bincode-legacy((SaltWitness, MptWitness)))`, which +//! the same `(SaltWitness, MptWitness)` tuple the RPC path yields. The failure taxonomy and +//! the transport wrapper are [`stateless_common::r2_witness`], shared with the +//! debug-trace-server's adapter; the transport core below that is `stateless-r2`'s +//! [`R2ObjectFetcher`]. This adapter owns what is validator-specific: the **full** payload +//! decode (proof verification needs the elliptic-curve points the light decode skips), the +//! validator metrics, and the surfaced-failure pacing the pipeline fetcher relies on. The +//! object body is `zstd(bincode-legacy((SaltWitness, MptWitness)))`, which //! [`stateless_common::decode_witness_payload`] inverts exactly. //! //! Operator note on missing objects: the pipeline retries a `Missing` witness indefinitely @@ -18,20 +19,18 @@ //! for the same block, and use the object key from the error's log line to check/backfill //! the bucket. On the custom-domain target, "appears once the uploader wins" additionally //! assumes the edge does not cache 404s — see the `--r2-custom-domain` flag docs. +//! +//! [`R2ObjectFetcher`]: stateless_r2::fetch::R2ObjectFetcher use std::time::{Duration, Instant}; use alloy_primitives::B256; use salt::SaltWitness; +pub use stateless_common::R2WitnessError; use stateless_common::{ - BackoffPolicy, WitnessDecodingError, WitnessSizeBreakdown, decode_witness_payload, + R2WitnessTransport, WitnessSizeBreakdown, decode_on_blocking_pool, decode_witness_payload, }; use stateless_core::withdrawals::MptWitness; -use stateless_r2::{ - fetch::{CfAccessCredentials, FetchTimeouts, R2GetError, R2ObjectFetcher, RetryPacing}, - keys, -}; -use tokio::task::JoinError; use tracing::trace; use crate::metrics; @@ -49,149 +48,19 @@ const DETERMINISTIC_FAILURE_THROTTLE: Duration = /// unboundedly, so there is no operator flag to mirror. const MAX_ATTEMPTS: usize = 9; -/// Failure outcome of an R2 witness fetch. -#[derive(Debug, thiserror::Error)] -pub enum R2WitnessError { - /// The GET failed (absent object, transport, throttle, or unexpected status — - /// see [`R2GetError`], and the module docs for the `Missing` operator note). - #[error(transparent)] - Get(#[from] R2GetError), - /// The object was fetched but its bytes did not decode to a `(SaltWitness, MptWitness)` tuple - /// — a corrupt witness in R2. Deterministic; not retried. - #[error("R2 witness for block {number} (key {key}) failed to decode: {source}")] - Decode { number: u64, key: String, source: WitnessDecodingError }, - /// The decode task panicked. This is a bug in our own decoder, not a problem with the data in - /// R2, so it is kept out of [`Self::Decode`]. - #[error("R2 witness decode task for block {number} (key {key}) panicked: {source}")] - DecodePanicked { number: u64, key: String, source: JoinError }, -} - -impl R2WitnessError { - /// Every label [`Self::kind`] can produce, for metrics pre-registration - /// (`crate::metrics::init_metrics` zero-inits the error counter per kind). - pub const KINDS: &'static [&'static str] = &[ - "missing", - "transport", - "throttled", - "status", - "connect", - "deadline", - "decode", - "decode_panicked", - ]; - - /// Stable lowercase label for this variant — the `kind` label on the R2 witness error - /// counter. Every value returned here must appear in [`Self::KINDS`]. - pub const fn kind(&self) -> &'static str { - match self { - Self::Get(e) => e.kind(), - Self::Decode { .. } => "decode", - Self::DecodePanicked { .. } => "decode_panicked", - } - } - - /// Whether an immediate retry against the same endpoint could plausibly succeed (transport - /// blips, 429, 5xx). Every other variant is deterministic and is surfaced without retrying. - const fn is_retryable(&self) -> bool { - matches!(self, Self::Get(e) if e.is_retryable()) - } -} - /// Fetches witness objects straight from an R2 bucket — SigV4-signed over the S3 API, or /// unsigned through a Cloudflare custom domain, per construction. -/// The fetcher's `Debug` redacts the credentials. +/// The transport's `Debug` redacts the credentials. #[derive(Debug)] pub struct R2WitnessClient { - fetcher: R2ObjectFetcher, - /// The configured in-flight GET cap, retained here because the fetcher decomposes it into - /// per-connection permits and cannot report the configured value back. - max_concurrent_requests: Option, -} - -/// The fetcher's pacing view of a `BackoffPolicy` — the adapter-layer conversion that keeps -/// `stateless-r2` free of a dependency on this workspace's backoff type. -fn pacing(backoff: &BackoffPolicy) -> RetryPacing { - RetryPacing { initial: backoff.initial, max: backoff.max } + transport: R2WitnessTransport, } impl R2WitnessClient { - /// The configured target's origin, for startup logging (see [`R2ObjectFetcher::origin`]). - pub fn origin(&self) -> &str { - self.fetcher.origin() - } - - /// The configured target's metric label (see [`R2ObjectFetcher::target_label`]). - pub const fn target_label(&self) -> &'static str { - self.fetcher.target_label() - } - - /// How many HTTP/2 connections the transport spreads its GETs over, for startup logging - /// (see [`R2ObjectFetcher::connections`]). - pub fn connections(&self) -> usize { - self.fetcher.connections() - } - - /// The configured cap on in-flight GETs (`None` = unlimited; see [`Self::new`] for the - /// exact semantics), for startup logging. - pub fn max_concurrent_requests(&self) -> Option { - self.max_concurrent_requests - } - - /// Builds a client from an R2 endpoint origin, bucket, and bucket-scoped S3 credentials. - /// - /// `timeouts` bounds each individual GET (end-to-end and connect). `retry_backoff` paces the - /// retries of retryable failures — first sleep `initial`, doubling up to `max`, each with - /// up to 50% jitter — and is the same policy the RPC path builds from - /// `--rpc-initial-backoff-ms` / `--rpc-max-backoff-ms`, so one pair of flags tunes both - /// paths. `max_concurrent_requests` caps the number of GETs in flight at once (`None` = - /// unlimited, `Some(0)` clamps to 1 — same semantics as the RPC witness semaphore; in R2 - /// mode this client is the only enforcement of `--r2-max-concurrent-requests`). Fails - /// if the endpoint is not a bare `scheme://host[:port]` origin or the HTTP client cannot be - /// built. - pub fn new( - endpoint: &str, - bucket: String, - access_key_id: String, - secret_access_key: String, - timeouts: FetchTimeouts, - retry_backoff: BackoffPolicy, - max_concurrent_requests: Option, - ) -> eyre::Result { - let fetcher = R2ObjectFetcher::new( - endpoint, - bucket, - access_key_id, - secret_access_key, - timeouts, - pacing(&retry_backoff), - max_concurrent_requests, - ) - .map_err(|e| eyre::eyre!(e))?; - Ok(Self { fetcher, max_concurrent_requests }) - } - - /// Builds a client that fetches unsigned through a Cloudflare custom domain fronting the - /// bucket (h2-multiplexed, edge-cacheable), with optional Cloudflare Access service-token - /// headers. The remaining parameters mean what they mean on [`Self::new`]. - pub fn new_custom_domain( - domain: &str, - access: Option, - timeouts: FetchTimeouts, - retry_backoff: BackoffPolicy, - max_concurrent_requests: Option, - connections: usize, - ) -> eyre::Result { - let fetcher = R2ObjectFetcher::new_custom_domain( - domain, - access, - timeouts, - pacing(&retry_backoff), - max_concurrent_requests, - connections, - ) - .map(|fetcher| fetcher.on_version_observed(metrics::record_r2_negotiated_version)) - .map_err(|e| eyre::eyre!(e))?; - Ok(Self { fetcher, max_concurrent_requests }) + /// Wraps an already-built transport. Construction (and the startup logging that reads + /// the configured target off it) lives at the wiring site, which owns the flags. + pub const fn new(transport: R2WitnessTransport) -> Self { + Self { transport } } /// Fetches and decodes the witness for `(number, hash)` from R2. @@ -215,7 +84,7 @@ impl R2WitnessClient { // fetch cycle would restart its ramp at `initial`, re-bursting GETs into the // same brownout the exhausted ramp just backed away from. let pause = if e.is_retryable() { - self.fetcher.pacing().max + self.transport.fetcher().pacing().max } else { DETERMINISTIC_FAILURE_THROTTLE }; @@ -232,27 +101,26 @@ impl R2WitnessClient { ) -> Result<(SaltWitness, MptWitness), R2WitnessError> { let started = Instant::now(); let fetched = self - .fetcher + .transport + .fetcher() .get_block_object(number, hash, MAX_ATTEMPTS, None, metrics::on_r2_witness_retry) .await?; let (bytes, queue_wait) = (fetched.bytes, fetched.queue_wait); - // zstd + bincode over a multi-MB witness is CPU-bound; keep it off the runtime. - let key = || keys::block_object_key(number, hash); - match tokio::task::spawn_blocking(move || decode_witness_payload(&bytes)).await { - Ok(Ok(witness)) => { - trace!(number, "R2 witness fetched and decoded"); - // Queue wait on the self-imposed concurrency cap is subtracted: folded in, it - // would masquerade as R2 slowness. - metrics::on_r2_witness_fetch_success( - started.elapsed().saturating_sub(queue_wait).as_secs_f64(), - WitnessSizeBreakdown::new(&witness.0, &witness.1), - ); - Ok(witness) - } - Ok(Err(source)) => Err(R2WitnessError::Decode { number, key: key(), source }), - Err(source) => Err(R2WitnessError::DecodePanicked { number, key: key(), source }), - } + // No deadline: the pipeline fetcher has no per-block budget to protect, so a slow + // decode must finish rather than be abandoned and re-fetched. + let witness = decode_on_blocking_pool(bytes, number, hash, None, |bytes| { + decode_witness_payload(bytes) + }) + .await?; + trace!(number, "R2 witness fetched and decoded"); + // Queue wait on the self-imposed concurrency cap is subtracted: folded in, it would + // masquerade as R2 slowness. + metrics::on_r2_witness_fetch_success( + started.elapsed().saturating_sub(queue_wait).as_secs_f64(), + WitnessSizeBreakdown::new(&witness.0, &witness.1), + ); + Ok(witness) } } @@ -260,6 +128,11 @@ impl R2WitnessClient { mod tests { use std::{str::FromStr, sync::atomic::Ordering}; + use stateless_common::BackoffPolicy; + use stateless_r2::{ + fetch::{FetchTimeouts, R2GetError}, + keys, + }; use stateless_test_utils::{fixtures::TestFixtures, mock_r2::mock_r2}; use super::*; @@ -284,53 +157,27 @@ mod tests { BackoffPolicy::new(Duration::from_millis(5), Duration::from_millis(20)) } - fn test_timeouts() -> FetchTimeouts { - FetchTimeouts { - per_attempt: Duration::from_secs(5), - connect: stateless_r2::fetch::DEFAULT_CONNECT_TIMEOUT, - } - } - fn client_with_backoff(endpoint: &str, retry_backoff: BackoffPolicy) -> R2WitnessClient { - R2WitnessClient::new( + let transport = R2WitnessTransport::new( endpoint, "witness-test".to_string(), "ak".to_string(), "sk".to_string(), - test_timeouts(), + FetchTimeouts { + per_attempt: Duration::from_secs(5), + connect: stateless_r2::fetch::DEFAULT_CONNECT_TIMEOUT, + }, retry_backoff, None, ) - .unwrap() + .unwrap(); + R2WitnessClient::new(transport) } async fn fetch(endpoint: &str) -> Result<(SaltWitness, MptWitness), R2WitnessError> { client_with_backoff(endpoint, test_backoff()).get_witness(1, B256::ZERO).await } - /// Construction errors from the shared fetcher must surface through the eyre conversion. - #[test] - fn rejects_endpoint_with_path() { - let err = R2WitnessClient::new( - "https://acc.r2.cloudflarestorage.com/witness-mainnet", - "witness-mainnet".to_string(), - "ak".to_string(), - "sk".to_string(), - test_timeouts(), - test_backoff(), - None, - ) - .unwrap_err(); - assert!(err.to_string().contains("Invalid R2 endpoint")); - } - - /// Every fetch-level kind must appear in this adapter's pre-registered [`KINDS`] — a new - /// `R2GetError` kind escaping metric pre-registration would drift silently otherwise. - #[test] - fn kinds_cover_all_fetch_kinds() { - assert!(R2GetError::KINDS.iter().all(|k| R2WitnessError::KINDS.contains(k))); - } - /// The only test of the success path (fetch → `spawn_blocking` decode): a fixture witness /// encoded with the uploader's `encode_witness_payload` must round-trip to the original /// tuple. @@ -360,15 +207,20 @@ mod tests { let (domain, _, heads) = stateless_test_utils::mock_r2::mock_r2_capturing(vec![(200, payload)]).await; - let client = R2WitnessClient::new_custom_domain( + let transport = R2WitnessTransport::new_custom_domain( &domain, None, - test_timeouts(), + FetchTimeouts { + per_attempt: Duration::from_secs(5), + connect: stateless_r2::fetch::DEFAULT_CONNECT_TIMEOUT, + }, test_backoff(), None, 1, + metrics::record_r2_negotiated_version, ) .unwrap(); + let client = R2WitnessClient::new(transport); let (decoded_salt, _) = client.get_witness(1, B256::ZERO).await.expect("valid object must fetch and decode"); assert_eq!(decoded_salt, salt_witness); @@ -377,22 +229,6 @@ mod tests { assert!(!head.contains("authorization:"), "custom-domain GET must be unsigned: {head}"); } - /// Construction errors from the shared fetcher's custom-domain arm surface through the - /// same eyre conversion as the S3 arm. - #[test] - fn custom_domain_rejects_origin_with_path() { - let err = R2WitnessClient::new_custom_domain( - "https://witness.example.com/witness-mainnet", - None, - test_timeouts(), - test_backoff(), - None, - 1, - ) - .unwrap_err(); - assert!(err.to_string().contains("Invalid R2 custom domain")); - } - #[tokio::test] async fn undecodable_body_surfaces_decode_without_retry() { let (endpoint, hits) = mock_r2(vec![(200, "not a zstd witness")]).await; diff --git a/bin/stateless-validator/tests/integration.rs b/bin/stateless-validator/tests/integration.rs index 2d195302..c69347f1 100644 --- a/bin/stateless-validator/tests/integration.rs +++ b/bin/stateless-validator/tests/integration.rs @@ -11,10 +11,7 @@ use std::{ use alloy_primitives::{B256, BlockHash}; use alloy_rpc_types_eth::Block; use clap::Parser; -use jsonrpsee::{ - RpcModule, - server::{ServerBuilder, ServerConfigBuilder}, -}; +use jsonrpsee::server::ServerConfigBuilder; use jsonrpsee_types::error::{ CALL_EXECUTION_FAILED_CODE, ErrorObject, ErrorObjectOwned, INVALID_PARAMS_CODE, }; @@ -24,7 +21,11 @@ use stateless_core::{ pipeline::run_pipeline, withdrawals::MptWitness, }; use stateless_db::ContractCache; -use stateless_test_utils::{fixtures::TestFixtures, logging::init_test_logging}; +use stateless_test_utils::{ + fixtures::TestFixtures, + logging::init_test_logging, + mock_rpc::{parse_hex_u64, serve_with_config}, +}; use stateless_validator::{ CommandLineArgs, VALIDATOR_DB_FILENAME, ValidatorDB, ValidatorFetcher, ValidatorHooks, ValidatorProcessor, load_or_create_chain_spec, run_with_signals, @@ -378,157 +379,159 @@ fn setup_test_db(fx: &TestFixtures) -> eyre::Result<(Arc, tempfile: async fn setup_mock_rpc_server( state: MockServerState, ) -> (jsonrpsee::server::ServerHandle, String) { - let mut module = RpcModule::new(state); - - 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 = u64::from_str_radix(&hex_number[2..], 16).unwrap_or(0); - - 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 cfg = ServerConfigBuilder::default().max_response_body_size(MAX_RESPONSE_BODY_SIZE).build(); + 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"), + ) + })?; + + 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}")) })?; - 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)) - }) - .unwrap(); - - module - .register_method("eth_blockNumber", |_params, ctx, _| { - let (&max_num, _) = ctx.fixtures.block_numbers.last_key_value().unwrap(); - Ok::(format!("0x{max_num:x}")) - }) - .unwrap(); - - module - .register_method("eth_getHeaderByNumber", |params, ctx, _| { - let (hex_number,): (String,) = params.parse().unwrap(); - let block_number = u64::from_str_radix(&hex_number[2..], 16).unwrap_or(0); - - 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 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()) - }) - .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()) - }) - .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 code = ctx.fixtures.contracts.get(&hash).cloned().unwrap_or_default(); - Ok::<_, ErrorObject<'static>>(code.original_bytes()) - }) - .unwrap(); - - 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 block_hash = BlockHash::from(keys.block_hash.0); - - let salt_witness = - ctx.fixtures.salt_witnesses.get(&block_hash).cloned().ok_or_else(|| { + Ok::<_, ErrorObject<'static>>(shape_block(block, full_block)) + }) + .unwrap(); + + module + .register_method("eth_blockNumber", |_params, ctx, _| { + let (&max_num, _) = ctx.fixtures.block_numbers.last_key_value().unwrap(); + Ok::(format!("0x{max_num:x}")) + }) + .unwrap(); + + 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()) + }) + .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()) + }) + .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 code = ctx.fixtures.contracts.get(&hash).cloned().unwrap_or_default(); + Ok::<_, ErrorObject<'static>>(code.original_bytes()) + }) + .unwrap(); + + 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 block_hash = BlockHash::from(keys.block_hash.0); + + let salt_witness = + ctx.fixtures.salt_witnesses.get(&block_hash).cloned().ok_or_else(|| { + make_rpc_error( + CALL_EXECUTION_FAILED_CODE, + format!("Witness for block {block_hash} not found"), + ) + })?; + + let mpt_witness = ctx.mpt_witnesses.get(&block_hash).cloned().ok_or_else(|| { make_rpc_error( CALL_EXECUTION_FAILED_CODE, format!("Witness for block {block_hash} not found"), ) })?; - let mpt_witness = ctx.mpt_witnesses.get(&block_hash).cloned().ok_or_else(|| { - make_rpc_error( - CALL_EXECUTION_FAILED_CODE, - format!("Witness for block {block_hash} not found"), - ) - })?; - - let encoded = encode_witness_response(&salt_witness, &mpt_witness).map_err(|e| { - make_rpc_error(CALL_EXECUTION_FAILED_CODE, format!("Failed to encode witness: {e}")) - })?; - - Ok::<_, ErrorObject<'static>>(encoded) - }) - .unwrap(); - - module - .register_method("mega_setValidatedBlocks", |params, ctx, _| { - use std::sync::atomic::Ordering; - let (first_block, last_block): ((u64, String), (u64, String)) = params.parse().unwrap(); - if ctx - .reject_reports - .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |n| n.checked_sub(1)) - .is_ok() - { - return Err(make_rpc_error( - CALL_EXECUTION_FAILED_CODE, - "transient report failure (scripted)".to_string(), - )); - } - ctx.validated_reports.lock().unwrap().push((first_block.0, last_block.0)); - let last_hash: BlockHash = last_block.1.parse().unwrap(); - Ok::(serde_json::json!({ - "accepted": true, - "lastValidatedBlock": [last_block.0, last_hash] - })) - }) - .unwrap(); - - let cfg = ServerConfigBuilder::default().max_response_body_size(MAX_RESPONSE_BODY_SIZE).build(); - let server = ServerBuilder::default().set_config(cfg).build("0.0.0.0:0").await.unwrap(); - let url = format!("http://{}", server.local_addr().unwrap()); - (server.start(module), url) + let encoded = + encode_witness_response(&salt_witness, &mpt_witness).map_err(|e| { + make_rpc_error( + CALL_EXECUTION_FAILED_CODE, + format!("Failed to encode witness: {e}"), + ) + })?; + + Ok::<_, ErrorObject<'static>>(encoded) + }) + .unwrap(); + + module + .register_method("mega_setValidatedBlocks", |params, ctx, _| { + use std::sync::atomic::Ordering; + let (first_block, last_block): ((u64, String), (u64, String)) = + params.parse().unwrap(); + if ctx + .reject_reports + .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |n| n.checked_sub(1)) + .is_ok() + { + return Err(make_rpc_error( + CALL_EXECUTION_FAILED_CODE, + "transient report failure (scripted)".to_string(), + )); + } + ctx.validated_reports.lock().unwrap().push((first_block.0, last_block.0)); + let last_hash: BlockHash = last_block.1.parse().unwrap(); + Ok::(serde_json::json!({ + "accepted": true, + "lastValidatedBlock": [last_block.0, last_hash] + })) + }) + .unwrap(); + }) + .await } /// Synthetic data integration test: validates consecutive blocks via the streaming pipeline. diff --git a/crates/stateless-common/Cargo.toml b/crates/stateless-common/Cargo.toml index 7a88942d..16a6130e 100644 --- a/crates/stateless-common/Cargo.toml +++ b/crates/stateless-common/Cargo.toml @@ -28,6 +28,7 @@ revm.workspace = true # stateless stateless-core = { path = "../stateless-core" } +stateless-r2 = { path = "../stateless-r2" } # misc base64.workspace = true @@ -50,7 +51,7 @@ zstd.workspace = true [dev-dependencies] jsonrpsee.workspace = true kanal.workspace = true -stateless-test-utils = { path = "../stateless-test-utils" } +stateless-test-utils = { path = "../stateless-test-utils", features = ["mock-rpc"] } tokio-util.workspace = true [[bench]] diff --git a/crates/stateless-common/src/lib.rs b/crates/stateless-common/src/lib.rs index 5df7ef92..aeed85f8 100644 --- a/crates/stateless-common/src/lib.rs +++ b/crates/stateless-common/src/lib.rs @@ -15,6 +15,8 @@ pub use witness_encoding::{ }; pub mod r2_args; pub use r2_args::{R2CountFlag, R2Flag, R2Flags, R2Target, R2TuningFlag, validate_r2_flags}; +pub mod r2_witness; +pub use r2_witness::{R2WitnessError, R2WitnessTransport, decode_on_blocking_pool}; pub mod secret; pub use secret::RedactedSecret; pub mod witness_size; diff --git a/crates/stateless-common/src/r2_witness.rs b/crates/stateless-common/src/r2_witness.rs new file mode 100644 index 00000000..157ad2c7 --- /dev/null +++ b/crates/stateless-common/src/r2_witness.rs @@ -0,0 +1,299 @@ +//! Shared core of the two binaries' direct-from-R2 witness adapters. +//! +//! Each binary reads witness objects straight from the R2 bucket through +//! [`R2ObjectFetcher`], but decodes and paces them differently: the trace server +//! light-decodes under a request deadline with no failure pauses, the validator +//! full-decodes with surfaced-failure pacing for its pipeline fetcher. What lives here is +//! the part that is identical by construction — the failure taxonomy with its metric +//! labels, and the transport wrapper (construction, target accessors) — so the two +//! adapters cannot drift apart on it. + +use std::time::Instant; + +use alloy_primitives::B256; +use stateless_r2::{ + fetch::{CfAccessCredentials, FetchTimeouts, R2GetError, R2ObjectFetcher, RetryPacing}, + keys, +}; +use tokio::task::JoinError; + +use crate::{BackoffPolicy, WitnessDecodingError}; + +/// Failure outcome of an R2 witness fetch, shared by both binaries' adapters. +/// +/// A binary whose fetches pass no deadline never produces [`Self::DecodeTimeout`] (or the +/// fetch-level `deadline` kind); its pre-registered series for those kinds stay at zero. +#[derive(Debug, thiserror::Error)] +pub enum R2WitnessError { + /// The GET failed (absent object, transport, throttle, unexpected status, or out of + /// deadline while queued) — see [`R2GetError`]. + #[error(transparent)] + Get(#[from] R2GetError), + /// The object was fetched but its bytes did not decode to a witness tuple — a corrupt + /// witness in R2. Deterministic; not retried. + #[error("R2 witness for block {number} (key {key}) failed to decode: {source}")] + Decode { number: u64, key: String, source: WitnessDecodingError }, + /// The decode outran what was left of the caller's deadline — an oversized or + /// pathological object. The blocking decode itself cannot be cancelled and finishes in + /// the background. + #[error("R2 witness decode for block {number} (key {key}) outran the deadline")] + DecodeTimeout { number: u64, key: String }, + /// The decode task panicked. This is a bug in our own decoder, not a problem with the + /// data in R2, so it is kept out of [`Self::Decode`]. + #[error("R2 witness decode task for block {number} (key {key}) panicked: {source}")] + DecodePanicked { number: u64, key: String, source: JoinError }, +} + +impl R2WitnessError { + /// Every label [`Self::kind`] can produce, for metrics pre-registration. + pub const KINDS: &'static [&'static str] = &[ + "missing", + "transport", + "throttled", + "status", + "connect", + "deadline", + "decode", + "decode_timeout", + "decode_panicked", + ]; + + /// Stable lowercase label for this variant — the `kind` label on the R2 witness error + /// counter. Every value returned here must appear in [`Self::KINDS`]. + pub const fn kind(&self) -> &'static str { + match self { + Self::Get(e) => e.kind(), + Self::Decode { .. } => "decode", + Self::DecodeTimeout { .. } => "decode_timeout", + Self::DecodePanicked { .. } => "decode_panicked", + } + } + + /// Whether the object was absent from the bucket — the one failure a caller probing + /// ahead of the uploader treats as expected rather than alarming. + pub const fn is_missing(&self) -> bool { + matches!(self, Self::Get(R2GetError::Missing { .. })) + } + + /// Whether an immediate retry against the same endpoint could plausibly succeed + /// (transport blips, 429, 5xx). Every other variant is deterministic and is surfaced + /// without retrying. + pub const fn is_retryable(&self) -> bool { + matches!(self, Self::Get(e) if e.is_retryable()) + } +} + +/// Decodes a fetched witness object with `decode` on the blocking pool — zstd + bincode over +/// a multi-MB witness is CPU-bound and must stay off the runtime — mapping both failure modes +/// onto [`R2WitnessError`]. +/// +/// `deadline` is the caller's budget for the decode: `Some` bounds it (an oversized or +/// pathological object must not eat what the caller reserved for its fallback, and an already +/// elapsed deadline skips the decode entirely — nothing would wait for it), while `None` lets +/// it run to completion. A decode abandoned on the deadline cannot be cancelled and finishes +/// in the background. +pub async fn decode_on_blocking_pool( + bytes: impl AsRef<[u8]> + Send + 'static, + number: u64, + hash: B256, + deadline: Option, + decode: impl FnOnce(&[u8]) -> Result + Send + 'static, +) -> Result { + let key = || keys::block_object_key(number, hash); + if deadline.is_some_and(|d| Instant::now() >= d) { + return Err(R2WitnessError::DecodeTimeout { number, key: key() }); + } + let task = tokio::task::spawn_blocking(move || decode(bytes.as_ref())); + let joined = match deadline { + Some(d) => match tokio::time::timeout_at(d.into(), task).await { + Ok(joined) => joined, + Err(_) => return Err(R2WitnessError::DecodeTimeout { number, key: key() }), + }, + None => task.await, + }; + match joined { + Ok(Ok(decoded)) => Ok(decoded), + Ok(Err(source)) => Err(R2WitnessError::Decode { number, key: key(), source }), + Err(source) => Err(R2WitnessError::DecodePanicked { number, key: key(), source }), + } +} + +/// The shared transport of the two R2 witness adapters: an [`R2ObjectFetcher`] plus the +/// construction and target accessors both binaries would otherwise duplicate verbatim. +/// The fetcher's `Debug` redacts the credentials. +#[derive(Debug)] +pub struct R2WitnessTransport { + fetcher: R2ObjectFetcher, + /// The configured in-flight GET cap, retained here because the fetcher decomposes it + /// into per-connection permits and cannot report the configured value back. + max_concurrent_requests: Option, +} + +/// The fetcher's pacing view of a [`BackoffPolicy`] — the adapter-layer conversion that +/// keeps `stateless-r2` free of a dependency on this workspace's backoff type. +fn pacing(backoff: &BackoffPolicy) -> RetryPacing { + RetryPacing { initial: backoff.initial, max: backoff.max } +} + +impl R2WitnessTransport { + /// Builds a transport from an R2 endpoint origin, bucket, and bucket-scoped S3 + /// credentials. + /// + /// `timeouts` bounds each individual GET (end-to-end and connect), `retry_backoff` + /// paces the fetcher's internal retries, and `max_concurrent_requests` caps in-flight + /// GETs (`None` = unlimited, `Some(0)` clamps to 1). Fails if the endpoint is not a + /// bare `scheme://host[:port]` origin or the HTTP client cannot be built. + pub fn new( + endpoint: &str, + bucket: String, + access_key_id: String, + secret_access_key: String, + timeouts: FetchTimeouts, + retry_backoff: BackoffPolicy, + max_concurrent_requests: Option, + ) -> eyre::Result { + let fetcher = R2ObjectFetcher::new( + endpoint, + bucket, + access_key_id, + secret_access_key, + timeouts, + pacing(&retry_backoff), + max_concurrent_requests, + ) + .map_err(|e| eyre::eyre!(e))?; + Ok(Self { fetcher, max_concurrent_requests }) + } + + /// Builds a transport that fetches unsigned through a Cloudflare custom domain + /// fronting the bucket (h2-multiplexed, edge-cacheable), with optional Cloudflare + /// Access service-token headers. `on_version_observed` receives the negotiated HTTP + /// version's metric label once known (each binary passes its own recorder); the + /// remaining parameters mean what they mean on [`Self::new`]. + pub fn new_custom_domain( + domain: &str, + access: Option, + timeouts: FetchTimeouts, + retry_backoff: BackoffPolicy, + max_concurrent_requests: Option, + connections: usize, + on_version_observed: impl Fn(&'static str) + Send + Sync + 'static, + ) -> eyre::Result { + let fetcher = R2ObjectFetcher::new_custom_domain( + domain, + access, + timeouts, + pacing(&retry_backoff), + max_concurrent_requests, + connections, + ) + .map(|fetcher| fetcher.on_version_observed(on_version_observed)) + .map_err(|e| eyre::eyre!(e))?; + Ok(Self { fetcher, max_concurrent_requests }) + } + + /// The underlying fetcher, for the adapter's own GETs and pacing reads. + pub fn fetcher(&self) -> &R2ObjectFetcher { + &self.fetcher + } + + /// The configured target's origin, for startup logging (see + /// [`R2ObjectFetcher::origin`]). + pub fn origin(&self) -> &str { + self.fetcher.origin() + } + + /// The configured target's metric label (see [`R2ObjectFetcher::target_label`]). + pub const fn target_label(&self) -> &'static str { + self.fetcher.target_label() + } + + /// How many HTTP/2 connections the transport spreads its GETs over, for startup + /// logging (see [`R2ObjectFetcher::connections`]). + pub fn connections(&self) -> usize { + self.fetcher.connections() + } + + /// The configured cap on in-flight GETs (`None` = unlimited; see [`Self::new`]), for + /// startup logging. + pub fn max_concurrent_requests(&self) -> Option { + self.max_concurrent_requests + } +} + +#[cfg(test)] +mod tests { + use std::time::Duration; + + use super::*; + + fn test_timeouts() -> FetchTimeouts { + FetchTimeouts { + per_attempt: Duration::from_secs(5), + connect: stateless_r2::fetch::DEFAULT_CONNECT_TIMEOUT, + } + } + + fn test_backoff() -> BackoffPolicy { + BackoffPolicy::new(Duration::from_millis(5), Duration::from_millis(20)) + } + + /// Every fetch-level kind must appear in the pre-registered [`R2WitnessError::KINDS`] + /// — a new [`R2GetError`] kind escaping metric pre-registration would drift silently + /// otherwise. One copy here guards both binaries' pre-registration loops. + #[test] + fn kinds_cover_all_fetch_kinds() { + assert!(R2GetError::KINDS.iter().all(|k| R2WitnessError::KINDS.contains(k))); + } + + /// Construction errors from the underlying fetcher must surface through the eyre + /// conversion, on both targets — one copy here covers both binaries' adapters. + #[test] + fn construction_rejects_a_target_carrying_a_path() { + let s3 = R2WitnessTransport::new( + "https://acc.r2.cloudflarestorage.com/witness-mainnet", + "witness-mainnet".to_string(), + "ak".to_string(), + "sk".to_string(), + test_timeouts(), + test_backoff(), + None, + ) + .unwrap_err(); + assert!(s3.to_string().contains("Invalid R2 endpoint"), "{s3}"); + + let custom_domain = R2WitnessTransport::new_custom_domain( + "https://witness.example.com/witness-mainnet", + None, + test_timeouts(), + test_backoff(), + None, + 1, + |_| {}, + ) + .unwrap_err(); + assert!(custom_domain.to_string().contains("Invalid R2 custom domain"), "{custom_domain}"); + } + + /// A decode whose deadline is already gone is abandoned before it starts — nothing + /// would wait for it — while a deadline-less decode runs to completion. + #[tokio::test] + async fn decode_respects_an_elapsed_deadline_and_runs_without_one() { + let elapsed = |_: &[u8]| -> Result { + panic!("an elapsed deadline must not start the decode") + }; + let err = + decode_on_blocking_pool(vec![0_u8; 8], 1, B256::ZERO, Some(Instant::now()), elapsed) + .await + .expect_err("an already-elapsed deadline must abandon the decode"); + assert!(matches!(err, R2WitnessError::DecodeTimeout { .. }), "{err}"); + assert_eq!(err.kind(), "decode_timeout"); + + let decoded = decode_on_blocking_pool(vec![7_u8; 4], 1, B256::ZERO, None, |bytes| { + Ok::(bytes.len()) + }) + .await + .expect("a deadline-less decode must run"); + assert_eq!(decoded, 4); + } +} diff --git a/crates/stateless-common/src/rpc_client.rs b/crates/stateless-common/src/rpc_client.rs index dadf0fbf..dfe2f38e 100644 --- a/crates/stateless-common/src/rpc_client.rs +++ b/crates/stateless-common/src/rpc_client.rs @@ -1643,17 +1643,14 @@ mod tests { }; use alloy_primitives::BlockHash; - use jsonrpsee::{ - RpcModule, - server::{ServerBuilder, ServerHandle}, - types::ErrorObjectOwned, - }; + use jsonrpsee::{server::ServerHandle, types::ErrorObjectOwned}; use stateless_core::{ PipelineConfig, block_fetcher, db::{BlockMeta, StoreResult}, find_divergence_point, pipeline::{BlockFetcher, DivergenceLookups}, }; + use stateless_test_utils::mock_rpc::{header_stub, parse_hex_u64, serve}; use tokio_util::sync::CancellationToken; use super::*; @@ -1693,31 +1690,6 @@ mod tests { ErrorObjectOwned::owned::<()>(-32000, msg, None) } - fn parse_hex_u64(s: &str) -> u64 { - u64::from_str_radix(s.strip_prefix("0x").unwrap_or(s), 16).unwrap() - } - - /// Minimal valid [`Header`] with `hash`/`number` populated; all other fields default. - fn header_stub(number: u64, hash: BlockHash) -> Header { - Header { - hash, - inner: alloy_consensus::Header { number, ..Default::default() }, - ..Default::default() - } - } - - /// Starts a jsonrpsee server bound to a random port and registers methods via `register`. - async fn serve( - ctx: Ctx, - register: impl FnOnce(&mut RpcModule), - ) -> (ServerHandle, String) { - let mut module = RpcModule::new(ctx); - register(&mut module); - let server = ServerBuilder::default().build("127.0.0.1:0").await.unwrap(); - let url = format!("http://{}", server.local_addr().unwrap()); - (server.start(module), url) - } - /// Serves `eth_getHeaderByNumber` with headers derived from `hashes`. async fn start_mock_rpc(hashes: HashMap) -> (ServerHandle, String) { serve(hashes, |m| { diff --git a/crates/stateless-test-utils/Cargo.toml b/crates/stateless-test-utils/Cargo.toml index a518aac9..e3127526 100644 --- a/crates/stateless-test-utils/Cargo.toml +++ b/crates/stateless-test-utils/Cargo.toml @@ -10,6 +10,7 @@ exclude.workspace = true [dependencies] # alloy +alloy-consensus.workspace = true alloy-genesis.workspace = true alloy-primitives.workspace = true alloy-rpc-types-eth.workspace = true @@ -23,6 +24,10 @@ op-alloy-rpc-types.workspace = true # misc bincode.workspace = true eyre.workspace = true +# Optional so that `mock_rpc`'s server stack is compiled only by the crates that mock an +# upstream node: this crate is a dev-dependency of `stateless-core` and `stateless-r2`, which +# have no other jsonrpsee edge and must keep building lean on their own. +jsonrpsee = { workspace = true, optional = true } # `serde` declared here (not inherited from sibling crates' feature unification): the fixtures # deserialize `Bytecode` from JSON, and this crate must also build standalone as a dev-dependency # of crates that pull in no other revm user (e.g. `stateless-r2`). @@ -32,3 +37,8 @@ serde_json.workspace = true tokio = { workspace = true, features = ["io-util", "net", "time"] } tracing = { workspace = true, features = ["std"] } tracing-subscriber.workspace = true + +[features] +# Opt-in so the jsonrpsee server stack stays out of the build for dev-dependents that mock no +# upstream node (`stateless-core`, `stateless-r2`). +mock-rpc = ["dep:jsonrpsee"] diff --git a/crates/stateless-test-utils/src/lib.rs b/crates/stateless-test-utils/src/lib.rs index 1e7daa7d..6907279e 100644 --- a/crates/stateless-test-utils/src/lib.rs +++ b/crates/stateless-test-utils/src/lib.rs @@ -2,3 +2,5 @@ pub mod env; pub mod fixtures; pub mod logging; pub mod mock_r2; +#[cfg(feature = "mock-rpc")] +pub mod mock_rpc; diff --git a/crates/stateless-test-utils/src/mock_rpc.rs b/crates/stateless-test-utils/src/mock_rpc.rs new file mode 100644 index 00000000..f7ce0d75 --- /dev/null +++ b/crates/stateless-test-utils/src/mock_rpc.rs @@ -0,0 +1,61 @@ +//! Shared jsonrpsee mock-upstream scaffolding. +//! +//! Every crate that talks to an upstream node fakes one in its tests; before this module +//! each hand-rolled the same server bootstrap and header shapes. Specialized mocks +//! (scripted witness sources, counting endpoints) stay with the tests that script them — +//! what lives here is only the scaffold they all share. + +use alloy_primitives::BlockHash; +use alloy_rpc_types_eth::Header; +use jsonrpsee::{ + RpcModule, + server::{ServerBuilder, ServerConfig, ServerHandle}, +}; + +/// Starts a jsonrpsee server on an ephemeral loopback port with methods registered via +/// `register`, returning the handle (dropping it stops the server) and the http URL. +pub async fn serve( + ctx: Ctx, + register: impl FnOnce(&mut RpcModule), +) -> (ServerHandle, String) { + serve_with_config(ServerConfig::default(), ctx, register).await +} + +/// [`serve`] with an explicit [`ServerConfig`] — for mocks whose fixture responses outgrow +/// jsonrpsee's default response-size cap. +pub async fn serve_with_config( + config: ServerConfig, + ctx: Ctx, + register: impl FnOnce(&mut RpcModule), +) -> (ServerHandle, String) { + let mut module = RpcModule::new(ctx); + register(&mut module); + let server = ServerBuilder::default().set_config(config).build("127.0.0.1:0").await.unwrap(); + let url = format!("http://{}", server.local_addr().unwrap()); + (server.start(module), url) +} + +/// Minimal RPC [`Header`] for `number` carrying the given `hash`; every other field is +/// default. For mocks that must serve prescribed hashes (e.g. divergence chains, where two +/// chains differ only by hash) — only `verify_hash = false` fetch paths accept it, since +/// the hash is not the header's real one. +pub fn header_stub(number: u64, hash: BlockHash) -> Header { + Header { + hash, + inner: alloy_consensus::Header { number, ..Default::default() }, + ..Default::default() + } +} + +/// Minimal self-consistent RPC [`Header`] for `number`: `hash` is the inner header's real +/// `hash_slow()`, so `verify_hash = true` fetch paths accept it too. +pub fn consistent_header(number: u64) -> Header { + let inner = alloy_consensus::Header { number, ..Default::default() }; + Header { hash: inner.hash_slow(), inner, ..Default::default() } +} + +/// Parses a `0x`-prefixed (or bare) hex string as `u64` — the wire shape of numeric +/// JSON-RPC block-number params. +pub fn parse_hex_u64(s: &str) -> u64 { + u64::from_str_radix(s.strip_prefix("0x").unwrap_or(s), 16).unwrap() +}