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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion bin/debug-trace-server/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"] }
153 changes: 71 additions & 82 deletions bin/debug-trace-server/src/data_provider.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::*;

Expand Down Expand Up @@ -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)]
Expand All @@ -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<C: Send + Sync + 'static>(module: RpcModule<C>) -> (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 {
Expand All @@ -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<MockRpcHits>) {
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)
}

Expand All @@ -1632,17 +1621,18 @@ pub(crate) mod test_support {
wire: Option<String>,
) -> (ServerHandle, String, Arc<AtomicUsize>) {
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)
}

Expand All @@ -1655,18 +1645,19 @@ pub(crate) mod test_support {
block: Block<Transaction>,
wire: Option<String>,
) -> (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.
Expand All @@ -1675,18 +1666,17 @@ pub(crate) mod test_support {
wire: String,
) -> (ServerHandle, String, Arc<AtomicUsize>) {
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)
}
}

Expand All @@ -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;
Expand Down
25 changes: 13 additions & 12 deletions bin/debug-trace-server/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<String>| 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),
Expand All @@ -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);
Expand Down
Loading
Loading