use std::num::NonZeroUsize;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use alloy_network::{Ethereum as AlloyEthereum, EthereumWallet, NetworkWallet, TransactionBuilder};
use alloy_primitives::{Address, Bytes, TxHash, U256};
use alloy_provider::fillers::{
BlobGasFiller, ChainIdFiller, FillProvider, GasFiller, JoinFill, NonceFiller, WalletFiller,
};
use alloy_provider::{
Identity, PendingTransactionError, Provider, ProviderBuilder, RootProvider, WalletProvider,
};
use alloy_rpc_client::RpcClient;
use alloy_rpc_types_eth::{BlockId, TransactionReceipt, TransactionRequest};
use alloy_transport::TransportError;
use alloy_transport::layers::{FallbackLayer, ThrottleLayer};
use alloy_transport_http::{Client, Http};
use r402_protocol::network::{ChainId, ChainProvider};
use tower::ServiceBuilder;
#[cfg(feature = "telemetry")]
use tracing::Instrument;
use url::Url;
use crate::chain::account::Eip155ChainReference;
use crate::chain::nonce::PendingNonceManager;
#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
pub enum Eip155ChainProviderError {
#[error("at least one signer must be provided")]
EmptyWallet,
#[error("at least one HTTP RPC endpoint is required")]
NoHttpEndpoint,
#[error("failed to build HTTP RPC client")]
HttpClient,
}
pub type InnerFiller = JoinFill<
GasFiller,
JoinFill<BlobGasFiller, JoinFill<NonceFiller<PendingNonceManager>, ChainIdFiller>>,
>;
pub type InnerProvider = FillProvider<
JoinFill<JoinFill<Identity, InnerFiller>, WalletFiller<EthereumWallet>>,
RootProvider,
>;
#[derive(Debug)]
pub struct Eip155ChainProvider {
chain: Eip155ChainReference,
eip1559: bool,
flashblocks: bool,
receipt_timeout_secs: u64,
inner: InnerProvider,
signer_addresses: Arc<[Address]>,
signer_cursor: Arc<AtomicUsize>,
nonce_manager: PendingNonceManager,
}
impl Eip155ChainProvider {
pub fn new(
chain: Eip155ChainReference,
wallet: EthereumWallet,
rpc_endpoints: &[(Url, Option<u32>)],
eip1559: bool,
flashblocks: bool,
receipt_timeout_secs: u64,
) -> Result<Self, Eip155ChainProviderError> {
let signer_addresses =
NetworkWallet::<AlloyEthereum>::signer_addresses(&wallet).collect::<Vec<_>>();
if signer_addresses.is_empty() {
return Err(Eip155ChainProviderError::EmptyWallet);
}
let signer_addresses: Arc<[Address]> = signer_addresses.into();
let signer_cursor = Arc::new(AtomicUsize::new(0));
let chain_id: ChainId = chain.into();
let client = Self::rpc_client(&chain_id, rpc_endpoints)?;
let nonce_manager = PendingNonceManager::default();
let filler = JoinFill::new(
GasFiller::default(),
JoinFill::new(
BlobGasFiller::default(),
JoinFill::new(
NonceFiller::new(nonce_manager.clone()),
ChainIdFiller::default(),
),
),
);
let inner: InnerProvider = ProviderBuilder::default()
.filler(filler)
.wallet(wallet)
.connect_client(client);
#[cfg(feature = "telemetry")]
tracing::info!(chain=%chain_id, signers=?signer_addresses, "Using EVM provider");
Ok(Self {
chain,
eip1559,
flashblocks,
receipt_timeout_secs,
inner,
signer_addresses,
signer_cursor,
nonce_manager,
})
}
pub fn rpc_client(
chain_id: &ChainId,
endpoints: &[(Url, Option<u32>)],
) -> Result<RpcClient, Eip155ChainProviderError> {
#[cfg(not(feature = "telemetry"))]
let _ = chain_id;
let http_client = Client::builder()
.http1_only()
.no_proxy()
.build()
.map_err(|_| Eip155ChainProviderError::HttpClient)?;
let transports = endpoints
.iter()
.filter_map(|(url, rate_limit)| {
let scheme = url.scheme();
let is_http = scheme == "http" || scheme == "https";
if !is_http {
return None;
}
#[cfg(feature = "telemetry")]
tracing::info!(chain=%chain_id, rpc_url=%url, rate_limit=?rate_limit, "Using HTTP transport");
let limit = rate_limit.unwrap_or(u32::MAX);
let service = ServiceBuilder::new()
.layer(ThrottleLayer::new(limit))
.service(Http::with_client(http_client.clone(), url.clone()));
Some(service)
})
.collect::<Vec<_>>();
let count =
NonZeroUsize::new(transports.len()).ok_or(Eip155ChainProviderError::NoHttpEndpoint)?;
let fallback = ServiceBuilder::new()
.layer(FallbackLayer::default().with_active_transport_count(count))
.service(transports);
Ok(RpcClient::new(fallback, false))
}
#[allow(
clippy::indexing_slicing,
reason = "bounds guaranteed by constructor requiring non-empty signers"
)]
fn next_signer_address(&self) -> Address {
debug_assert!(
!self.signer_addresses.is_empty(),
"signer_addresses must not be empty"
);
if self.signer_addresses.len() == 1 {
self.signer_addresses[0]
} else {
let next =
self.signer_cursor.fetch_add(1, Ordering::Relaxed) % self.signer_addresses.len();
self.signer_addresses[next]
}
}
}
#[derive(Debug, thiserror::Error)]
pub enum MetaTransactionSendError {
#[error(transparent)]
Transport(#[from] TransportError),
#[error(transparent)]
PendingTransaction(#[from] PendingTransactionError),
#[error("receipt wait failed for {hash}")]
ReceiptWait {
hash: TxHash,
#[source]
source: PendingTransactionError,
},
#[error("{0}")]
Custom(String),
}
#[derive(Debug, Clone)]
pub struct MetaTransaction {
pub to: Address,
pub calldata: Bytes,
pub confirmations: u64,
pub from: Option<Address>,
pub value: U256,
}
impl MetaTransaction {
#[must_use]
pub const fn new(to: Address, calldata: Bytes, confirmations: u64) -> Self {
Self {
to,
calldata,
confirmations,
from: None,
value: U256::ZERO,
}
}
#[must_use]
pub const fn with_from(mut self, from: Address) -> Self {
self.from = Some(from);
self
}
#[must_use]
pub fn with_data_suffix(mut self, suffix: &[u8]) -> Self {
self.calldata = crate::chain::append_data_suffix(self.calldata, suffix);
self
}
}
impl ChainProvider for Eip155ChainProvider {
fn signer_addresses(&self) -> Vec<String> {
self.inner
.signer_addresses()
.map(|a| a.to_string())
.collect()
}
fn chain_id(&self) -> ChainId {
self.chain.into()
}
}
pub trait Eip155MetaTransactionProvider {
type Error;
type Inner: Provider;
fn inner(&self) -> &Self::Inner;
fn chain(&self) -> &Eip155ChainReference;
fn send_transaction(
&self,
tx: MetaTransaction,
) -> impl Future<Output = Result<TransactionReceipt, Self::Error>> + Send;
fn send_raw_transaction(
&self,
encoded: &[u8],
confirmations: u64,
) -> impl Future<Output = Result<TransactionReceipt, Self::Error>> + Send;
}
impl<T: Eip155MetaTransactionProvider> Eip155MetaTransactionProvider for Arc<T> {
type Error = T::Error;
type Inner = T::Inner;
fn inner(&self) -> &Self::Inner {
(**self).inner()
}
fn chain(&self) -> &Eip155ChainReference {
(**self).chain()
}
fn send_transaction(
&self,
tx: MetaTransaction,
) -> impl Future<Output = Result<TransactionReceipt, Self::Error>> + Send {
(**self).send_transaction(tx)
}
fn send_raw_transaction(
&self,
encoded: &[u8],
confirmations: u64,
) -> impl Future<Output = Result<TransactionReceipt, Self::Error>> + Send {
(**self).send_raw_transaction(encoded, confirmations)
}
}
impl Eip155MetaTransactionProvider for Eip155ChainProvider {
type Error = MetaTransactionSendError;
type Inner = InnerProvider;
fn inner(&self) -> &Self::Inner {
&self.inner
}
fn chain(&self) -> &Eip155ChainReference {
&self.chain
}
async fn send_transaction(
&self,
tx: MetaTransaction,
) -> Result<TransactionReceipt, Self::Error> {
let from_address = match tx.from {
Some(pinned) => {
if !self.signer_addresses.contains(&pinned) {
return Err(MetaTransactionSendError::Custom(format!(
"requested signer {pinned} is not in the configured wallet"
)));
}
pinned
}
None => self.next_signer_address(),
};
let mut txr = TransactionRequest::default()
.with_to(tx.to)
.with_from(from_address)
.with_input(tx.calldata)
.with_value(tx.value);
if !self.eip1559 {
let provider = &self.inner;
let gas_fut = provider.get_gas_price();
#[cfg(feature = "telemetry")]
let gas: u128 = gas_fut
.instrument(tracing::info_span!("get_gas_price"))
.await?;
#[cfg(not(feature = "telemetry"))]
let gas: u128 = gas_fut.await?;
txr.set_gas_price(gas);
}
if txr.gas.is_none() {
let block_id = if self.flashblocks {
BlockId::latest()
} else {
BlockId::pending()
};
let gas_limit = self.inner.estimate_gas(txr.clone()).block(block_id).await?;
txr.set_gas_limit(gas_limit);
}
let pending_tx = match self.inner.send_transaction(txr).await {
Ok(pending) => pending,
Err(e) => {
self.nonce_manager.reset_nonce(from_address).await;
return Err(MetaTransactionSendError::Transport(e));
}
};
let timeout = std::time::Duration::from_secs(self.receipt_timeout_secs);
let tx_hash = *pending_tx.tx_hash();
let watcher = pending_tx
.with_required_confirmations(tx.confirmations)
.with_timeout(Some(timeout));
match watcher.get_receipt().await {
Ok(receipt) => Ok(receipt),
Err(e) => {
self.nonce_manager.reset_nonce(from_address).await;
Err(MetaTransactionSendError::ReceiptWait {
hash: tx_hash,
source: e,
})
}
}
}
async fn send_raw_transaction(
&self,
encoded: &[u8],
confirmations: u64,
) -> Result<TransactionReceipt, Self::Error> {
let pending_tx = self.inner.send_raw_transaction(encoded).await?;
let timeout = std::time::Duration::from_secs(self.receipt_timeout_secs);
let tx_hash = *pending_tx.tx_hash();
let watcher = pending_tx
.with_required_confirmations(confirmations)
.with_timeout(Some(timeout));
match watcher.get_receipt().await {
Ok(receipt) => Ok(receipt),
Err(e) => Err(MetaTransactionSendError::ReceiptWait {
hash: tx_hash,
source: e,
}),
}
}
}
#[cfg(test)]
#[allow(
clippy::expect_used,
clippy::unwrap_used,
reason = "test assertions on known-valid fixtures"
)]
mod tests {
use std::str::FromStr;
use std::sync::{Arc, Mutex};
use super::*;
const RAW_TX_HASH: TxHash = TxHash::repeat_byte(0x11);
fn chain_id() -> ChainId {
"eip155:8453".parse().expect("fixture chain id")
}
#[test]
fn rpc_client_rejects_empty_endpoints() {
let err = Eip155ChainProvider::rpc_client(&chain_id(), &[]).unwrap_err();
assert_eq!(
err,
Eip155ChainProviderError::NoHttpEndpoint,
"empty endpoint list must return NoHttpEndpoint, not panic"
);
}
#[test]
fn rpc_client_rejects_non_http_endpoints() {
let ws = Url::parse("ws://127.0.0.1:8545").expect("fixture ws url");
let err = Eip155ChainProvider::rpc_client(&chain_id(), &[(ws, None)]).unwrap_err();
assert_eq!(
err,
Eip155ChainProviderError::NoHttpEndpoint,
"non-HTTP endpoints must return NoHttpEndpoint after filtering"
);
}
#[test]
fn rpc_client_accepts_https_endpoint() {
let url = Url::parse("https://mainnet.base.org").expect("fixture rpc url");
Eip155ChainProvider::rpc_client(&chain_id(), &[(url, None)])
.expect("HTTPS endpoint must construct");
}
#[test]
fn meta_transaction_new_defaults_value_zero() {
let tx = MetaTransaction::new(Address::ZERO, Bytes::new(), 1);
assert_eq!(tx.value, U256::ZERO, "new() defaults value to zero");
assert_eq!(tx.from, None, "new() does not pin a signer");
let pinned = Address::repeat_byte(0x42);
let tx = tx.with_from(pinned);
assert_eq!(tx.from, Some(pinned), "with_from pins signer");
assert_eq!(tx.value, U256::ZERO, "with_from leaves value unchanged");
}
#[test]
fn with_data_suffix_preserves_value() {
let mut tx = MetaTransaction::new(Address::ZERO, Bytes::from_static(&[0x01]), 2);
tx.value = U256::from(9u64);
let tx = tx.with_data_suffix(&[0xaa]);
assert_eq!(tx.value, U256::from(9u64), "suffix must not clobber value");
assert_eq!(
tx.calldata.as_ref(),
&[0x01, 0xaa],
"suffix appends calldata"
);
assert_eq!(tx.confirmations, 2, "suffix must not clobber confirmations");
}
fn anvil_wallet() -> EthereumWallet {
let signer = alloy_signer_local::PrivateKeySigner::from_str(
"0xac0974bec39a17e36ba4a6b4d238ff944bacb478cbed5efcae784d7bf4f2ff80",
)
.expect("anvil key");
EthereumWallet::from(signer)
}
fn provider_http(rpc: &str, receipt_timeout_secs: u64) -> Eip155ChainProvider {
let url = Url::parse(rpc).expect("rpc url");
Eip155ChainProvider::new(
Eip155ChainReference::new(8453),
anvil_wallet(),
&[(url, None)],
true,
false,
receipt_timeout_secs,
)
.expect("provider")
}
fn mined_receipt_json(hash: TxHash) -> serde_json::Value {
serde_json::json!({
"transactionHash": hash,
"transactionIndex": "0x0",
"blockHash": "0x2222222222222222222222222222222222222222222222222222222222222222",
"blockNumber": "0x1",
"from": "0xf39Fd6e51aad88F6F4ce6aB8827279cffFb92266",
"to": "0x4020074e9dF2ce1deE5A9C1b5c3f541D02a10003",
"cumulativeGasUsed": "0x1",
"gasUsed": "0x1",
"effectiveGasPrice": "0x1",
"contractAddress": null,
"logs": [],
"logsBloom": format!("0x{}", "0".repeat(512)),
"status": "0x1",
"type": "0x2"
})
}
fn rpc_result(
id: &serde_json::Value,
result: &serde_json::Value,
) -> wiremock::ResponseTemplate {
wiremock::ResponseTemplate::new(200).set_body_json(serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"result": result
}))
}
fn rpc_error(id: &serde_json::Value, message: &str) -> wiremock::ResponseTemplate {
wiremock::ResponseTemplate::new(200).set_body_json(serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"error": { "code": -32000, "message": message }
}))
}
fn raw_tx_hex(body: &serde_json::Value) -> Option<String> {
body.get("params")
.and_then(serde_json::Value::as_array)
.and_then(|params| params.first())
.and_then(serde_json::Value::as_str)
.map(str::to_owned)
}
#[derive(Clone, Copy)]
enum RawRpcMode {
BroadcastFail,
ReceiptRpcFail,
ReceiptTimeout,
Mined,
}
#[derive(Clone)]
struct RawRpc {
mode: RawRpcMode,
captured: Arc<Mutex<Option<String>>>,
}
impl RawRpc {
fn remember_raw(&self, body: &serde_json::Value) {
let Some(hex) = raw_tx_hex(body) else {
return;
};
*self.captured.lock().expect("captured") = Some(hex);
}
}
impl wiremock::Respond for RawRpc {
fn respond(&self, request: &wiremock::Request) -> wiremock::ResponseTemplate {
let body: serde_json::Value = serde_json::from_slice(&request.body).unwrap_or_default();
let id = body
.get("id")
.cloned()
.unwrap_or_else(|| serde_json::json!(1));
let method = body
.get("method")
.and_then(serde_json::Value::as_str)
.unwrap_or("");
if method == "eth_sendRawTransaction" {
self.remember_raw(&body);
return match self.mode {
RawRpcMode::BroadcastFail => rpc_error(&id, "broadcast failed"),
_ => rpc_result(&id, &serde_json::json!(RAW_TX_HASH)),
};
}
if method == "eth_getTransactionReceipt" {
return match self.mode {
RawRpcMode::Mined => rpc_result(&id, &mined_receipt_json(RAW_TX_HASH)),
RawRpcMode::ReceiptRpcFail => rpc_error(&id, "receipt failed"),
_ => rpc_result(&id, &serde_json::Value::Null),
};
}
if method == "eth_blockNumber" {
return rpc_result(&id, &serde_json::json!("0x1"));
}
rpc_result(&id, &serde_json::json!("0x"))
}
}
async fn mount_raw_rpc(
server: &wiremock::MockServer,
mode: RawRpcMode,
) -> Arc<Mutex<Option<String>>> {
let captured = Arc::new(Mutex::new(None));
let script = RawRpc {
mode,
captured: Arc::clone(&captured),
};
wiremock::Mock::given(wiremock::matchers::method("POST"))
.respond_with(script)
.mount(server)
.await;
captured
}
#[tokio::test]
async fn send_raw_transaction_broadcast_error_is_transport() {
let server = wiremock::MockServer::start().await;
let captured = mount_raw_rpc(&server, RawRpcMode::BroadcastFail).await;
let provider = provider_http(&server.uri(), 1);
let encoded = [0x02, 0xaa, 0xbb];
let err = provider
.send_raw_transaction(&encoded, 1)
.await
.expect_err("broadcast failure must not succeed");
assert!(
matches!(err, MetaTransactionSendError::Transport(_)),
"broadcast RPC error must be Transport, got {err}"
);
let hex = captured.lock().expect("captured").clone();
assert_eq!(
hex.as_deref(),
Some("0x02aabb"),
"raw envelope must be forwarded as hex"
);
}
#[tokio::test]
async fn send_raw_transaction_waits_for_receipt() {
let server = wiremock::MockServer::start().await;
mount_raw_rpc(&server, RawRpcMode::Mined).await;
let provider = provider_http(&server.uri(), 5);
let receipt = provider
.send_raw_transaction(&[0x02], 1)
.await
.expect("mined raw tx");
assert_eq!(
receipt.transaction_hash, RAW_TX_HASH,
"receipt hash must match eth_sendRawTransaction result"
);
assert!(receipt.status(), "fixture receipt is successful");
}
#[tokio::test]
async fn send_raw_transaction_receipt_rpc_error_is_receipt_wait() {
let server = wiremock::MockServer::start().await;
mount_raw_rpc(&server, RawRpcMode::ReceiptRpcFail).await;
let provider = provider_http(&server.uri(), 1);
let err = provider
.send_raw_transaction(&[0x02], 1)
.await
.expect_err("receipt RPC failure must not succeed");
match err {
MetaTransactionSendError::ReceiptWait { hash, .. } => {
assert_eq!(
hash, RAW_TX_HASH,
"ReceiptWait must carry the broadcast hash"
);
}
other => panic!("expected ReceiptWait, got {other}"),
}
}
#[tokio::test]
async fn send_raw_transaction_receipt_timeout_is_receipt_wait() {
let server = wiremock::MockServer::start().await;
mount_raw_rpc(&server, RawRpcMode::ReceiptTimeout).await;
let provider = provider_http(&server.uri(), 1);
let err = provider
.send_raw_transaction(&[0x02], 1)
.await
.expect_err("missing receipt must time out");
assert!(
matches!(err, MetaTransactionSendError::ReceiptWait { .. }),
"timeout must be ReceiptWait, got {err}"
);
}
}