mod cache_store;
mod transport;
use std::{
fs,
path::{Path, PathBuf},
};
use alloy_provider::{
layers::CacheLayer,
transport::{
layers::{RateLimitRetryPolicy, RetryBackoffLayer},
RpcError, TransportError, TransportErrorKind,
},
DynProvider, Provider, ProviderBuilder,
};
use alloy_rpc_client::{ClientBuilder, RpcClient};
use clap::Parser;
use tracing::{debug, info, warn};
pub use self::cache_store::{ExternalEnvSnapshot, RpcCacheStore};
use self::{
cache_store::CacheFileEnvelope,
transport::{CachingTransport, ReplayTransport, TransportCache},
};
use super::{EvmeError, Result};
pub type OpProvider = DynProvider<op_alloy_network::Optimism>;
#[derive(Debug)]
pub struct BuildProviderOutput {
pub provider: OpProvider,
pub cache_store: RpcCacheStore,
pub chain_id: u64,
pub external_env: Option<ExternalEnvSnapshot>,
}
#[derive(Parser, Debug, Clone)]
#[command(next_help_heading = "RPC Options")]
pub struct RpcArgs {
#[arg(
long = "rpc",
visible_aliases = ["rpc-url"],
alias = "fork.rpc",
)]
pub rpc_url: Option<String>,
#[arg(
long = "rpc.capture-file",
value_parser = parse_non_empty_path,
requires = "rpc_url",
conflicts_with_all = ["replay_file", "cache_dir", "clear_cache", "no_cache_file", "cache_size"],
)]
pub capture_file: Option<PathBuf>,
#[arg(
long = "rpc.replay-file",
value_parser = parse_non_empty_path,
conflicts_with_all = ["rpc_url", "capture_file", "cache_dir", "clear_cache", "no_cache_file", "cache_size"],
)]
pub replay_file: Option<PathBuf>,
#[arg(id = "cache_size", long = "rpc.cache-size", default_value_t = 10_000)]
pub cache_size: u32,
#[arg(long = "rpc.cache-dir", value_parser = parse_non_empty_path)]
pub cache_dir: Option<PathBuf>,
#[arg(long = "rpc.no-cache-file")]
pub no_cache_file: bool,
#[arg(long = "rpc.clear-cache")]
pub clear_cache: bool,
#[arg(long = "rpc.max-retries", default_value_t = 5)]
pub max_retries: u32,
#[arg(long = "rpc.backoff-ms", default_value_t = 1_000)]
pub backoff_ms: u64,
#[arg(long = "rpc.rate-limit", default_value_t = 660)]
pub compute_units_per_sec: u64,
}
impl RpcArgs {
pub async fn build_provider(&self) -> Result<BuildProviderOutput> {
let rpc_url_str = self.rpc_url.as_deref().ok_or_else(|| {
EvmeError::InvalidInput("No RPC URL provided. Pass '--rpc <URL>'.".to_string())
})?;
let url: reqwest::Url = rpc_url_str.parse().map_err(|e| {
EvmeError::RpcError(format!("Invalid RPC URL '{}': {}", rpc_url_str, e))
})?;
let chain_id = self.resolve_chain_id(url.clone()).await?;
if self.cache_size == 0 {
let provider = build_bare_op_provider(self.build_retry_client(url));
info!(
rpc_url = %rpc_url_str,
max_retries = self.max_retries,
backoff_ms = self.backoff_ms,
"Built RPC provider (cache disabled)",
);
return Ok(BuildProviderOutput {
provider,
cache_store: RpcCacheStore::noop(),
chain_id,
external_env: None,
});
}
let cache_path = if self.no_cache_file {
None
} else {
Some(resolve_cache_path(self.cache_dir.as_deref(), chain_id)?)
};
let cache_layer = CacheLayer::new(self.cache_size);
let cache = cache_layer.cache();
let cache_store = match cache_path {
Some(path) => {
if self.clear_cache {
if let Err(e) = fs::remove_file(&path) {
if e.kind() != std::io::ErrorKind::NotFound {
return Err(EvmeError::RpcError(format!(
"Failed to clear RPC cache at {}: {e}",
path.display(),
)));
}
} else {
info!(path = %path.display(), "Cleared existing RPC cache");
}
}
if let Some(parent) = path.parent() {
if let Err(e) = fs::create_dir_all(parent) {
warn!(
path = %parent.display(),
error = %e,
"Failed to create cache directory; persist may fail",
);
}
}
if path.exists() {
if let Err(err) = cache.load_cache(path.clone()) {
warn!(
path = %path.display(),
error = %err,
"Failed to load RPC cache; starting empty",
);
}
}
RpcCacheStore::new(cache, path)
}
None => RpcCacheStore::noop(),
};
let client = self.build_retry_client(url);
let provider = ProviderBuilder::new()
.disable_recommended_fillers()
.layer(cache_layer)
.network::<op_alloy_network::Optimism>()
.connect_client(client);
info!(
rpc_url = %rpc_url_str,
cache_size = self.cache_size,
max_retries = self.max_retries,
backoff_ms = self.backoff_ms,
"Built RPC provider",
);
Ok(BuildProviderOutput {
provider: DynProvider::new(provider),
cache_store,
chain_id,
external_env: None,
})
}
pub async fn build_replay_provider(&self) -> Result<BuildProviderOutput> {
let path = self.replay_file.as_ref().expect("replay mode requires --rpc.replay-file");
let envelope = CacheFileEnvelope::load(path)?;
let chain_id = envelope.chain_id;
debug!(
path = %path.display(),
chain_id,
has_external_env = envelope.external_env.is_some(),
"Loaded replay envelope",
);
let transport_cache = TransportCache::from_value(&envelope.cache)?;
debug!(entries = transport_cache.len(), "Seeded transport cache from envelope");
let replay_client = ClientBuilder::default()
.transport(ReplayTransport::new(path.clone(), transport_cache), true);
let provider = ProviderBuilder::new()
.disable_recommended_fillers()
.network::<op_alloy_network::Optimism>()
.connect_client(replay_client);
info!(
path = %path.display(),
chain_id,
"Built RPC provider (replay from cache file)",
);
Ok(BuildProviderOutput {
provider: DynProvider::new(provider),
cache_store: RpcCacheStore::noop(),
chain_id,
external_env: envelope.external_env,
})
}
pub async fn build_capture_provider(&self) -> Result<BuildProviderOutput> {
let path = self.capture_file.as_ref().expect("capture mode requires --rpc.capture-file");
let rpc_url_str = self.rpc_url.as_ref().expect("capture mode requires --rpc");
let url: reqwest::Url = rpc_url_str.parse().map_err(|e| {
EvmeError::RpcError(format!("Invalid RPC URL '{}': {}", rpc_url_str, e))
})?;
let existing_envelope = if path.exists() {
let env = CacheFileEnvelope::load(path)?;
debug!(
path = %path.display(),
chain_id = env.chain_id,
has_external_env = env.external_env.is_some(),
"Found existing capture envelope; entries will merge after chain-id validation",
);
Some(env)
} else {
debug!(path = %path.display(), "No existing capture envelope; starting fresh");
None
};
let transport_cache = TransportCache::new();
let http = alloy_transport_http::Http::new(url.clone());
let caching = CachingTransport::new(http, transport_cache.clone());
let client = self.build_client(caching, &url);
let provider = ProviderBuilder::new()
.disable_recommended_fillers()
.network::<op_alloy_network::Optimism>()
.connect_client(client);
let chain_id = provider.get_chain_id().await.map_err(|e| {
EvmeError::RpcError(format!("Failed to fetch chain ID from '{}': {e}", rpc_url_str))
})?;
if let Some(ref env) = existing_envelope {
if env.chain_id != chain_id {
return Err(EvmeError::RpcError(format!(
"Chain ID mismatch: cache file {} contains chain {} but endpoint returned {}",
path.display(),
env.chain_id,
chain_id,
)));
}
transport_cache.merge(&env.cache)?;
debug!(
entries = transport_cache.len(),
"Merged existing envelope entries into capture cache",
);
}
info!(
rpc_url = %rpc_url_str,
path = %path.display(),
chain_id,
"Built RPC provider (capture to cache file)",
);
let prev_external_env = existing_envelope.and_then(|e| e.external_env);
Ok(BuildProviderOutput {
provider: DynProvider::new(provider),
cache_store: RpcCacheStore::new_envelope(transport_cache, path.clone(), chain_id),
chain_id,
external_env: prev_external_env,
})
}
fn build_retry_client(&self, url: reqwest::Url) -> RpcClient {
self.build_client(alloy_transport_http::Http::new(url.clone()), &url)
}
fn build_client<T: alloy_transport::IntoBoxTransport>(
&self,
transport: T,
url: &reqwest::Url,
) -> RpcClient {
let is_local =
url.host_str().is_some_and(|h| h == "localhost" || h == "127.0.0.1" || h == "::1");
if self.max_retries > 0 {
let policy = RateLimitRetryPolicy::default().or(|err: &TransportError| {
matches!(err, RpcError::Transport(TransportErrorKind::Custom(_)))
});
let retry = RetryBackoffLayer::new_with_policy(
self.max_retries,
self.backoff_ms,
self.compute_units_per_sec,
policy,
);
ClientBuilder::default().layer(retry).transport(transport, is_local)
} else {
ClientBuilder::default().transport(transport, is_local)
}
}
async fn resolve_chain_id(&self, url: reqwest::Url) -> Result<u64> {
let url_str = url.as_str().to_string();
let bare = build_bare_op_provider(self.build_retry_client(url));
bare.get_chain_id().await.map_err(|e| {
EvmeError::RpcError(format!("Failed to fetch chain ID from '{}': {}", url_str, e))
})
}
}
fn parse_non_empty_path(s: &str) -> std::result::Result<PathBuf, String> {
if s.trim().is_empty() {
Err("path must not be empty".to_string())
} else {
Ok(PathBuf::from(s))
}
}
fn build_bare_op_provider(client: RpcClient) -> OpProvider {
DynProvider::new(
ProviderBuilder::new()
.disable_recommended_fillers()
.network::<op_alloy_network::Optimism>()
.connect_client(client),
)
}
fn resolve_cache_path(user_cache_dir: Option<&Path>, chain_id: u64) -> Result<PathBuf> {
let dir = match user_cache_dir {
Some(dir) => dir.to_path_buf(),
None => dirs::cache_dir()
.ok_or_else(|| {
EvmeError::RpcError(
"Could not determine default cache directory; pass --rpc.cache-dir \
or --rpc.no-cache-file"
.to_string(),
)
})?
.join("mega-evme")
.join("rpc"),
};
Ok(dir.join(format!("rpc-cache-{chain_id}.json")))
}
#[cfg(test)]
mod tests {
use std::path::PathBuf;
use super::*;
#[test]
fn test_resolve_cache_path_with_explicit_dir() {
let dir = PathBuf::from("/some/user/dir");
let path = resolve_cache_path(Some(&dir), 4326).expect("resolve");
assert_eq!(path, PathBuf::from("/some/user/dir/rpc-cache-4326.json"));
}
#[test]
fn test_resolve_cache_path_default_is_under_platform_cache() {
let Some(expected_root) = dirs::cache_dir() else {
return;
};
let path = resolve_cache_path(None, 11_155_420).expect("resolve");
let expected = expected_root.join("mega-evme").join("rpc").join("rpc-cache-11155420.json");
assert_eq!(path, expected);
}
}