use crate::Provider;
use alloy_eips::BlockId;
use alloy_network::{Network, TransactionBuilder};
use alloy_primitives::{Address, Bytes};
use alloy_rpc_types_eth::TransactionInputKind;
use alloy_sol_types::{SolCall, SolError, SolValue};
use alloy_transport::TransportError;
#[cfg(not(target_family = "wasm"))]
use futures::future::BoxFuture as CcipFuture;
#[cfg(target_family = "wasm")]
use futures::future::LocalBoxFuture as CcipFuture;
use futures::{stream, StreamExt};
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::sync::Semaphore;
mod bounds;
#[cfg(all(feature = "ccip-read-http", not(all(target_os = "wasi", target_env = "p1"))))]
mod http;
#[cfg(all(feature = "ccip-read-http", not(all(target_os = "wasi", target_env = "p1"))))]
pub use http::HttpCcipReadGateway;
const BATCH_GATEWAY_SENTINEL: &str = "x-batch-gateway:true";
mod abi {
alloy_sol_types::sol! {
error OffchainLookup(
address sender,
string[] urls,
bytes callData,
bytes4 callbackFunction,
bytes extraData
);
error HttpError(uint16 status, string message);
struct BatchGatewayRequest {
address sender;
string[] urls;
bytes data;
}
function query(BatchGatewayRequest[] requests)
external
view
returns (bool[] failures, bytes[] responses);
}
}
#[derive(Clone, Debug)]
pub struct CcipReadConfig {
pub max_redirects: usize,
pub max_batch_size: usize,
pub max_concurrent_requests: usize,
pub max_total_requests: usize,
pub max_gateway_urls: usize,
pub max_revert_data_size: usize,
pub max_response_size: usize,
}
impl Default for CcipReadConfig {
fn default() -> Self {
Self {
max_redirects: 4,
max_batch_size: 50,
max_concurrent_requests: 4,
max_total_requests: 100,
max_gateway_urls: 8,
max_revert_data_size: 1_048_576,
max_response_size: 1_048_576,
}
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct CcipReadRequest {
pub sender: Address,
pub urls: Vec<String>,
pub data: Bytes,
}
#[derive(Clone, Debug, thiserror::Error)]
#[error("{message}")]
pub struct CcipReadGatewayError {
pub status: Option<u16>,
pub message: String,
}
impl CcipReadGatewayError {
pub fn new(message: impl Into<String>) -> Self {
Self { status: None, message: message.into() }
}
pub fn http(status: u16, message: impl Into<String>) -> Self {
Self { status: Some(status), message: message.into() }
}
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum CcipReadError {
#[error(transparent)]
Transport(#[from] TransportError),
#[error("CCIP Read requires an eth_call target")]
MissingTarget,
#[error("OffchainLookup sender {sender} does not match call target {target}")]
SenderMismatch {
sender: Address,
target: Address,
},
#[error("CCIP Read redirect limit of {0} exceeded")]
TooManyRedirects(usize),
#[error("invalid OffchainLookup revert: {0}")]
InvalidOffchainLookup(alloy_sol_types::Error),
#[error("CCIP Read gateway request failed: {0}")]
Gateway(#[from] CcipReadGatewayError),
#[error("invalid ENSIP-21 batch request: {0}")]
InvalidBatch(String),
#[error("invalid CCIP Read configuration: {0}")]
InvalidConfig(String),
#[error("CCIP Read resource limit exceeded: {0}")]
ResourceLimit(String),
}
#[cfg_attr(target_family = "wasm", async_trait::async_trait(?Send))]
#[cfg_attr(not(target_family = "wasm"), async_trait::async_trait)]
pub trait CcipReadGateway: Send + Sync {
async fn request(
&self,
request: &CcipReadRequest,
max_response_size: usize,
) -> Result<Bytes, CcipReadGatewayError>;
}
#[derive(Clone, Copy, Debug, Default)]
#[cfg(all(feature = "ccip-read-http", target_os = "wasi", target_env = "p1"))]
pub struct HttpCcipReadGateway;
#[cfg(all(feature = "ccip-read-http", target_os = "wasi", target_env = "p1"))]
#[async_trait::async_trait(?Send)]
impl CcipReadGateway for HttpCcipReadGateway {
async fn request(
&self,
_request: &CcipReadRequest,
_max_response_size: usize,
) -> Result<Bytes, CcipReadGatewayError> {
Err(CcipReadGatewayError::new(
"the default CCIP Read HTTP gateway is unavailable on WASI Preview 1",
))
}
}
#[derive(Clone, Debug)]
pub struct CcipReadClient<G> {
gateway: G,
config: CcipReadConfig,
}
#[cfg(feature = "ccip-read-http")]
impl Default for CcipReadClient<HttpCcipReadGateway> {
fn default() -> Self {
Self::new(HttpCcipReadGateway::default())
}
}
impl<G> CcipReadClient<G> {
pub fn new(gateway: G) -> Self {
Self { gateway, config: CcipReadConfig::default() }
}
pub const fn with_config(mut self, config: CcipReadConfig) -> Self {
self.config = config;
self
}
pub const fn config(&self) -> &CcipReadConfig {
&self.config
}
pub const fn gateway(&self) -> &G {
&self.gateway
}
}
impl<G: CcipReadGateway> CcipReadClient<G> {
pub async fn call<P, N>(
&self,
provider: &P,
transaction: N::TransactionRequest,
) -> Result<Bytes, CcipReadError>
where
P: Provider<N>,
N: Network,
{
self.call_at(provider, transaction, BlockId::latest()).await
}
pub async fn call_at<P, N>(
&self,
provider: &P,
mut transaction: N::TransactionRequest,
block: BlockId,
) -> Result<Bytes, CcipReadError>
where
P: Provider<N>,
N: Network,
{
if self.config.max_concurrent_requests == 0 {
return Err(CcipReadError::InvalidConfig(
"max_concurrent_requests must be greater than zero".into(),
));
}
let target = transaction.to().ok_or(CcipReadError::MissingTarget)?;
let context = BatchContext::new(&self.config);
let mut redirects = 0;
loop {
let error = match provider.call(transaction.clone()).block(block).await {
Ok(result) => return Ok(result),
Err(error) => error,
};
let Some(revert) = extract_offchain_lookup(&error, self.config.max_revert_data_size)?
else {
return Err(CcipReadError::Transport(error));
};
if redirects == self.config.max_redirects {
return Err(CcipReadError::TooManyRedirects(self.config.max_redirects));
}
redirects += 1;
bounds::offchain_lookup(&revert, &self.config)?;
let lookup = abi::OffchainLookup::abi_decode(&revert)
.map_err(CcipReadError::InvalidOffchainLookup)?;
if lookup.sender != target {
return Err(CcipReadError::SenderMismatch { sender: lookup.sender, target });
}
let request =
CcipReadRequest { sender: lookup.sender, urls: lookup.urls, data: lookup.callData };
let response = self.fetch(request, &context).await?;
let mut callback = lookup.callbackFunction.to_vec();
callback.extend_from_slice(&(response, lookup.extraData).abi_encode_params());
transaction.set_input_kind(callback, TransactionInputKind::Both);
}
}
fn fetch<'a>(
&'a self,
request: CcipReadRequest,
context: &'a BatchContext<'a>,
) -> CcipFuture<'a, Result<Bytes, CcipReadError>> {
Box::pin(async move {
if request.urls.len() > context.config.max_gateway_urls {
return Err(CcipReadError::ResourceLimit(format!(
"gateway URL count {} exceeds limit {}",
request.urls.len(),
context.config.max_gateway_urls
)));
}
if request.urls.iter().any(|url| url == BATCH_GATEWAY_SENTINEL) {
context.reserve(1)?;
return self.local_batch(request.data, context).await;
}
context.reserve(request.urls.len().max(1))?;
let _permit =
context.concurrency.acquire().await.expect("CCIP Read semaphore is never closed");
let response = self.gateway.request(&request, context.config.max_response_size).await?;
if response.len() > context.config.max_response_size {
return Err(CcipReadError::ResourceLimit(format!(
"gateway response is {} bytes; limit is {}",
response.len(),
context.config.max_response_size
)));
}
Ok(response)
})
}
fn local_batch<'a>(
&'a self,
data: Bytes,
context: &'a BatchContext<'a>,
) -> CcipFuture<'a, Result<Bytes, CcipReadError>> {
Box::pin(async move {
bounds::batch(&data, context.config)?;
let call = abi::queryCall::abi_decode(&data)
.map_err(|err| CcipReadError::InvalidBatch(err.to_string()))?;
let requests = call.requests.into_iter().map(|request| {
let request = CcipReadRequest {
sender: request.sender,
urls: request.urls,
data: request.data,
};
self.fetch(request, context)
});
let (failures, responses) = stream::iter(requests)
.buffered(context.config.max_concurrent_requests)
.map(|result| match result {
Ok(response) => (false, response),
Err(error) => (true, encode_batch_error(&error)),
})
.unzip::<_, _, Vec<_>, Vec<_>>()
.await;
let encoded: Bytes =
abi::queryCall::abi_encode_returns(&abi::queryReturn { failures, responses })
.into();
if encoded.len() > context.config.max_response_size {
return Err(CcipReadError::ResourceLimit(format!(
"batch gateway response is {} bytes; limit is {}",
encoded.len(),
context.config.max_response_size
)));
}
Ok(encoded)
})
}
}
struct BatchContext<'a> {
config: &'a CcipReadConfig,
total_requests: AtomicUsize,
concurrency: Semaphore,
}
impl<'a> BatchContext<'a> {
fn new(config: &'a CcipReadConfig) -> Self {
Self {
config,
total_requests: AtomicUsize::new(0),
concurrency: Semaphore::new(config.max_concurrent_requests),
}
}
fn reserve(&self, count: usize) -> Result<(), CcipReadError> {
let limit = self.config.max_total_requests;
let mut current = self.total_requests.load(Ordering::Relaxed);
loop {
let Some(next) = current.checked_add(count).filter(|next| *next <= limit) else {
return Err(CcipReadError::ResourceLimit(format!(
"total gateway request budget of {limit} exceeded"
)));
};
match self.total_requests.compare_exchange_weak(
current,
next,
Ordering::Relaxed,
Ordering::Relaxed,
) {
Ok(_) => return Ok(()),
Err(actual) => current = actual,
}
}
}
}
fn extract_offchain_lookup(
error: &TransportError,
max_revert_data_size: usize,
) -> Result<Option<Bytes>, CcipReadError> {
let Some(raw) = error.as_error_resp().and_then(|payload| payload.data.as_ref()) else {
return Ok(None);
};
let max_json_size = max_revert_data_size.saturating_mul(2).saturating_add(4_096);
if raw.get().len() > max_json_size {
return Ok(None);
}
let Ok(value) = serde_json::from_str(raw.get()) else {
return Ok(None);
};
let Some(data) = find_offchain_lookup(&value) else {
return Ok(None);
};
if data.len() > max_revert_data_size {
return Err(CcipReadError::ResourceLimit(format!(
"OffchainLookup revert data is {} bytes; limit is {max_revert_data_size}",
data.len()
)));
}
Ok(Some(data))
}
fn find_offchain_lookup(value: &serde_json::Value) -> Option<Bytes> {
match value {
serde_json::Value::String(value) => {
let data: Bytes = value.parse().ok()?;
data.starts_with(abi::OffchainLookup::SELECTOR.as_slice()).then_some(data)
}
serde_json::Value::Object(values) => values.values().find_map(find_offchain_lookup),
serde_json::Value::Array(values) => values.iter().find_map(find_offchain_lookup),
_ => None,
}
}
fn encode_batch_error(error: &CcipReadError) -> Bytes {
if let CcipReadError::Gateway(CcipReadGatewayError { status: Some(status), message }) = error {
return abi::HttpError { status: *status, message: message.clone() }.abi_encode().into();
}
alloy_sol_types::Revert::from(error.to_string()).abi_encode().into()
}
#[cfg(feature = "ccip-read-http")]
pub fn shared_http_ccip_read_client() -> &'static CcipReadClient<HttpCcipReadGateway> {
static CLIENT: std::sync::OnceLock<CcipReadClient<HttpCcipReadGateway>> =
std::sync::OnceLock::new();
CLIENT.get_or_init(CcipReadClient::default)
}
#[cfg(feature = "ccip-read-http")]
#[cfg_attr(target_family = "wasm", async_trait::async_trait(?Send))]
#[cfg_attr(not(target_family = "wasm"), async_trait::async_trait)]
pub trait ProviderCcipReadExt<N: Network>: Provider<N> {
async fn call_with_ccip_read(
&self,
transaction: N::TransactionRequest,
) -> Result<Bytes, CcipReadError>;
async fn call_with_ccip_read_at(
&self,
transaction: N::TransactionRequest,
block: BlockId,
) -> Result<Bytes, CcipReadError>;
}
#[cfg(feature = "ccip-read-http")]
#[cfg_attr(target_family = "wasm", async_trait::async_trait(?Send))]
#[cfg_attr(not(target_family = "wasm"), async_trait::async_trait)]
impl<P, N> ProviderCcipReadExt<N> for P
where
P: Provider<N>,
N: Network,
{
async fn call_with_ccip_read(
&self,
transaction: N::TransactionRequest,
) -> Result<Bytes, CcipReadError> {
shared_http_ccip_read_client().call(self, transaction).await
}
async fn call_with_ccip_read_at(
&self,
transaction: N::TransactionRequest,
block: BlockId,
) -> Result<Bytes, CcipReadError> {
shared_http_ccip_read_client().call_at(self, transaction, block).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ProviderBuilder;
use alloy_json_rpc::ErrorPayload;
use alloy_primitives::{address, bytes, fixed_bytes, U256};
use alloy_rpc_types_eth::{TransactionInput, TransactionRequest};
use alloy_transport::mock::Asserter;
use std::{
collections::VecDeque,
sync::{Arc, Mutex, PoisonError},
};
#[derive(Clone, Debug, Default)]
struct MockGateway {
responses: Arc<Mutex<VecDeque<Result<Bytes, CcipReadGatewayError>>>>,
requests: Arc<Mutex<Vec<CcipReadRequest>>>,
}
impl MockGateway {
fn with_responses(
responses: impl IntoIterator<Item = Result<Bytes, CcipReadGatewayError>>,
) -> Self {
Self {
responses: Arc::new(Mutex::new(responses.into_iter().collect())),
requests: Arc::default(),
}
}
fn requests(&self) -> Vec<CcipReadRequest> {
self.requests.lock().unwrap_or_else(PoisonError::into_inner).clone()
}
}
#[async_trait::async_trait]
impl CcipReadGateway for MockGateway {
async fn request(
&self,
request: &CcipReadRequest,
_max_response_size: usize,
) -> Result<Bytes, CcipReadGatewayError> {
self.requests.lock().unwrap_or_else(PoisonError::into_inner).push(request.clone());
self.responses
.lock()
.unwrap_or_else(PoisonError::into_inner)
.pop_front()
.unwrap_or_else(|| Err(CcipReadGatewayError::new("no mock response")))
}
}
#[derive(Clone, Copy, Debug)]
struct BatchMockGateway;
#[async_trait::async_trait]
impl CcipReadGateway for BatchMockGateway {
async fn request(
&self,
request: &CcipReadRequest,
_max_response_size: usize,
) -> Result<Bytes, CcipReadGatewayError> {
match request.data.as_ref() {
[1] => Ok(bytes!("aaaa")),
[2] => Err(CcipReadGatewayError::http(404, "not found")),
_ => Err(CcipReadGatewayError::new("unexpected request")),
}
}
}
fn revert_error(data: Bytes) -> ErrorPayload {
ErrorPayload::internal_error_with_message_and_obj(
"call failed".into(),
serde_json::value::to_raw_value(&data).unwrap(),
)
}
fn offchain_lookup(sender: Address, urls: Vec<String>, call_data: Bytes) -> Bytes {
abi::OffchainLookup {
sender,
urls,
callData: call_data,
callbackFunction: fixed_bytes!("12345678"),
extraData: bytes!("010203"),
}
.abi_encode()
.into()
}
fn word_at(data: &[u8], offset: usize) -> usize {
U256::from_be_slice(&data[offset..offset + 32]).to()
}
#[tokio::test]
async fn bounds_decoded_lookup_before_allocating_aliased_urls() {
let target = address!("1111111111111111111111111111111111111111");
let canonical =
offchain_lookup(target, vec!["a".repeat(1024), String::new()], Bytes::new());
let mut aliased = canonical.to_vec();
let array = 4 + word_at(&aliased, 4 + 32) + 32;
aliased.copy_within(array..array + 32, array + 32);
for (revert, limit, accepted) in [
(canonical.clone(), canonical.len(), true),
(Bytes::from(aliased.clone()), aliased.len(), false),
(Bytes::from(aliased), 4096, true),
] {
let asserter = Asserter::new();
asserter.push_failure(revert_error(revert));
asserter.push_success(&bytes!("feed"));
let provider = ProviderBuilder::new().connect_mocked_client(asserter);
let gateway = MockGateway::with_responses([Ok(bytes!("01"))]);
let client = CcipReadClient::new(gateway.clone())
.with_config(CcipReadConfig { max_revert_data_size: limit, ..Default::default() });
let result = client.call(&provider, TransactionRequest::default().to(target)).await;
if accepted {
assert_eq!(result.unwrap(), bytes!("feed"));
} else {
assert!(matches!(result, Err(CcipReadError::ResourceLimit(_))));
assert!(gateway.requests().is_empty());
}
}
}
#[tokio::test]
async fn bounds_decoded_batch_before_allocating_aliased_requests() {
let sender = Address::ZERO;
let request = |data| abi::BatchGatewayRequest {
sender,
urls: vec!["https://example.test".into()],
data,
};
let canonical =
abi::queryCall { requests: vec![request(vec![1; 1024].into()), request(Bytes::new())] }
.abi_encode();
let mut aliased = canonical.clone();
let array = 4 + word_at(&aliased, 4) + 32;
aliased.copy_within(array..array + 32, array + 32);
for (data, accepted) in [(canonical, true), (aliased, false)] {
let gateway = MockGateway::with_responses([Ok(bytes!("01")), Ok(bytes!("02"))]);
let client = CcipReadClient::new(gateway.clone()).with_config(CcipReadConfig {
max_revert_data_size: data.len(),
..Default::default()
});
let context = BatchContext::new(client.config());
let result = client
.fetch(
CcipReadRequest {
sender,
urls: vec![BATCH_GATEWAY_SENTINEL.into()],
data: data.into(),
},
&context,
)
.await;
if accepted {
let decoded = abi::queryCall::abi_decode_returns(&result.unwrap()).unwrap();
assert_eq!(decoded.failures, vec![false, false]);
} else {
assert!(matches!(result, Err(CcipReadError::ResourceLimit(_))));
assert!(gateway.requests().is_empty());
}
}
}
#[test]
fn checks_array_counts_offsets_and_lossy_strings_before_decoding() {
let config = CcipReadConfig::default();
let canonical = offchain_lookup(Address::ZERO, vec!["a".repeat(256)], Bytes::new());
let urls = 4 + word_at(&canonical, 4 + 32);
let mut excessive = canonical.to_vec();
excessive[urls..urls + 32].fill(0xff);
assert!(bounds::offchain_lookup(&excessive, &config).is_err());
let mut invalid_offset = canonical.to_vec();
invalid_offset[4 + 32..4 + 64].fill(0xff);
assert!(matches!(
bounds::offchain_lookup(&invalid_offset, &config),
Err(CcipReadError::InvalidOffchainLookup(_))
));
let mut invalid_utf8 = canonical.to_vec();
let string = urls + 32 + word_at(&canonical, urls + 32) + 32;
invalid_utf8[string..string + 256].fill(0xff);
let config = CcipReadConfig { max_revert_data_size: canonical.len(), ..config };
bounds::offchain_lookup(&canonical, &config).unwrap();
assert!(matches!(
bounds::offchain_lookup(&invalid_utf8, &config),
Err(CcipReadError::ResourceLimit(_))
));
let mut batch = abi::queryCall { requests: vec![] }.abi_encode();
let array = 4 + word_at(&batch, 4);
batch[array..array + 32]
.copy_from_slice(&U256::from(config.max_batch_size + 1).to_be_bytes::<32>());
assert!(matches!(bounds::batch(&batch, &config), Err(CcipReadError::ResourceLimit(_))));
}
#[tokio::test]
async fn follows_offchain_lookup_and_calls_callback() {
assert_eq!(abi::OffchainLookup::SELECTOR, [0x55, 0x6f, 0x18, 0x30]);
let target = address!("1111111111111111111111111111111111111111");
let call_data = bytes!("abcdef");
let urls = vec!["https://example.test/{sender}/{data}".to_string()];
let asserter = Asserter::new();
asserter.push_failure(revert_error(offchain_lookup(
target,
urls.clone(),
call_data.clone(),
)));
asserter.push_success(&bytes!("feed"));
let provider = ProviderBuilder::new().connect_mocked_client(asserter);
let gateway = MockGateway::with_responses([Ok(bytes!("deadbeef"))]);
let client = CcipReadClient::new(gateway.clone());
let result =
client.call(&provider, TransactionRequest::default().to(target)).await.unwrap();
assert_eq!(result, bytes!("feed"));
assert_eq!(
gateway.requests(),
vec![CcipReadRequest { sender: target, urls, data: call_data }]
);
}
#[tokio::test]
async fn callback_keeps_input_and_data_in_sync() {
let target = address!("1111111111111111111111111111111111111111");
let revert =
offchain_lookup(target, vec!["https://example.test/{data}".into()], bytes!("abcdef"));
let asserter = Asserter::new();
asserter.push_failure(revert_error(revert));
asserter.push_success(&bytes!("feed"));
let provider = ProviderBuilder::new().connect_mocked_client(asserter);
let gateway = MockGateway::with_responses([Ok(bytes!("deadbeef"))]);
let client = CcipReadClient::new(gateway);
let result = client
.call(
&provider,
TransactionRequest::default()
.to(target)
.input(TransactionInput::both(bytes!("00"))),
)
.await
.unwrap();
assert_eq!(result, bytes!("feed"));
}
#[tokio::test]
async fn rejects_sender_mismatch() {
let target = address!("1111111111111111111111111111111111111111");
let sender = address!("2222222222222222222222222222222222222222");
let revert =
offchain_lookup(sender, vec!["https://example.test/{data}".into()], Bytes::new());
let asserter = Asserter::new();
asserter.push_failure(revert_error(revert));
let provider = ProviderBuilder::new().connect_mocked_client(asserter);
let gateway = MockGateway::default();
let error = CcipReadClient::new(gateway.clone())
.call(&provider, TransactionRequest::default().to(target))
.await
.unwrap_err();
assert!(matches!(
error,
CcipReadError::SenderMismatch { sender: actual_sender, target: actual_target }
if actual_sender == sender && actual_target == target
));
assert!(gateway.requests().is_empty());
}
#[tokio::test]
async fn rejects_excessive_gateway_url_list() {
let target = address!("1111111111111111111111111111111111111111");
let revert =
offchain_lookup(target, vec!["https://example.test/{data}".into(); 9], Bytes::new());
let asserter = Asserter::new();
asserter.push_failure(revert_error(revert));
let provider = ProviderBuilder::new().connect_mocked_client(asserter);
let gateway = MockGateway::default();
let error = CcipReadClient::new(gateway.clone())
.call(&provider, TransactionRequest::default().to(target))
.await
.unwrap_err();
assert!(matches!(error, CcipReadError::ResourceLimit(message) if message.contains("URL")));
assert!(gateway.requests().is_empty());
}
#[tokio::test]
async fn enforces_response_limit_for_custom_gateways() {
let target = address!("1111111111111111111111111111111111111111");
let revert =
offchain_lookup(target, vec!["https://example.test/{data}".into()], Bytes::new());
let asserter = Asserter::new();
asserter.push_failure(revert_error(revert));
let provider = ProviderBuilder::new().connect_mocked_client(asserter);
let gateway = MockGateway::with_responses([Ok(bytes!("0102"))]);
let config = CcipReadConfig { max_response_size: 1, ..Default::default() };
let error = CcipReadClient::new(gateway)
.with_config(config)
.call(&provider, TransactionRequest::default().to(target))
.await
.unwrap_err();
assert!(
matches!(error, CcipReadError::ResourceLimit(message) if message.contains("response"))
);
}
#[tokio::test]
async fn rejects_oversized_revert_data_before_decoding() {
let target = address!("1111111111111111111111111111111111111111");
let revert =
offchain_lookup(target, vec!["https://example.test/{data}".into()], bytes!("01020304"));
let asserter = Asserter::new();
asserter.push_failure(revert_error(revert));
let provider = ProviderBuilder::new().connect_mocked_client(asserter);
let config = CcipReadConfig { max_revert_data_size: 4, ..Default::default() };
let error = CcipReadClient::new(MockGateway::default())
.with_config(config)
.call(&provider, TransactionRequest::default().to(target))
.await
.unwrap_err();
assert!(
matches!(error, CcipReadError::ResourceLimit(message) if message.contains("revert"))
);
}
#[tokio::test]
async fn preserves_oversized_non_ccip_rpc_errors() {
let target = address!("1111111111111111111111111111111111111111");
let asserter = Asserter::new();
asserter.push_failure(revert_error(vec![0u8; 5_000].into()));
let provider = ProviderBuilder::new().connect_mocked_client(asserter);
let config = CcipReadConfig { max_revert_data_size: 1, ..Default::default() };
let error = CcipReadClient::new(MockGateway::default())
.with_config(config)
.call(&provider, TransactionRequest::default().to(target))
.await
.unwrap_err();
assert!(matches!(error, CcipReadError::Transport(_)));
}
#[tokio::test]
async fn enforces_redirect_limit() {
let target = address!("1111111111111111111111111111111111111111");
let revert =
offchain_lookup(target, vec!["https://example.test/{data}".into()], Bytes::new());
let asserter = Asserter::new();
for _ in 0..3 {
asserter.push_failure(revert_error(revert.clone()));
}
let provider = ProviderBuilder::new().connect_mocked_client(asserter);
let gateway =
MockGateway::with_responses([Ok(bytes!("01")), Ok(bytes!("02")), Ok(bytes!("03"))]);
let config = CcipReadConfig { max_redirects: 2, ..Default::default() };
let error = CcipReadClient::new(gateway.clone())
.with_config(config)
.call(&provider, TransactionRequest::default().to(target))
.await
.unwrap_err();
assert!(matches!(error, CcipReadError::TooManyRedirects(2)));
assert_eq!(gateway.requests().len(), 2);
}
#[tokio::test]
async fn executes_batch_gateway_requests_in_original_order() {
assert_eq!(abi::queryCall::SELECTOR, [0xa7, 0x80, 0xba, 0xb6]);
let sender = address!("1111111111111111111111111111111111111111");
let batch = abi::queryCall {
requests: vec![
abi::BatchGatewayRequest {
sender,
urls: vec!["https://one.test".into()],
data: bytes!("01"),
},
abi::BatchGatewayRequest {
sender,
urls: vec!["https://two.test".into()],
data: bytes!("02"),
},
],
}
.abi_encode()
.into();
let client = CcipReadClient::new(BatchMockGateway);
let context = BatchContext::new(client.config());
let encoded = client
.fetch(
CcipReadRequest { sender, urls: vec![BATCH_GATEWAY_SENTINEL.into()], data: batch },
&context,
)
.await
.unwrap();
let decoded = abi::queryCall::abi_decode_returns(&encoded).unwrap();
assert_eq!(decoded.failures, vec![false, true]);
assert_eq!(decoded.responses[0], bytes!("aaaa"));
let http_error = abi::HttpError::abi_decode(&decoded.responses[1]).unwrap();
assert_eq!(http_error.status, 404);
assert_eq!(http_error.message, "not found");
}
#[tokio::test]
async fn enforces_response_limit_for_local_batch() {
let sender = address!("1111111111111111111111111111111111111111");
let batch = abi::queryCall {
requests: vec![abi::BatchGatewayRequest {
sender,
urls: vec!["https://one.test".into()],
data: bytes!("01"),
}],
}
.abi_encode()
.into();
let config = CcipReadConfig { max_response_size: 8, ..Default::default() };
let client = CcipReadClient::new(BatchMockGateway).with_config(config);
let context = BatchContext::new(client.config());
let error = client
.fetch(
CcipReadRequest { sender, urls: vec![BATCH_GATEWAY_SENTINEL.into()], data: batch },
&context,
)
.await
.unwrap_err();
assert!(matches!(
error,
CcipReadError::ResourceLimit(message) if message.contains("batch gateway response")
));
}
#[tokio::test]
async fn limits_nested_batch_recursion() {
let sender = address!("1111111111111111111111111111111111111111");
let mut data = Bytes::new();
for _ in 0..3 {
data = abi::queryCall {
requests: vec![abi::BatchGatewayRequest {
sender,
urls: vec![BATCH_GATEWAY_SENTINEL.into()],
data,
}],
}
.abi_encode()
.into();
}
let config = CcipReadConfig { max_total_requests: 2, ..Default::default() };
let client = CcipReadClient::new(MockGateway::default()).with_config(config);
let context = BatchContext::new(client.config());
let encoded = client
.fetch(
CcipReadRequest { sender, urls: vec![BATCH_GATEWAY_SENTINEL.into()], data },
&context,
)
.await
.unwrap();
let outer = abi::queryCall::abi_decode_returns(&encoded).unwrap();
assert_eq!(outer.failures, vec![false]);
let middle = abi::queryCall::abi_decode_returns(&outer.responses[0]).unwrap();
assert_eq!(middle.failures, vec![true]);
assert!(alloy_sol_types::Revert::abi_decode(&middle.responses[0])
.unwrap()
.reason
.contains("budget"));
}
#[tokio::test]
async fn rejects_zero_max_concurrent_requests_before_eth_call() {
let asserter = Asserter::new();
let provider = ProviderBuilder::new().connect_mocked_client(asserter);
let client = CcipReadClient::new(MockGateway::default())
.with_config(CcipReadConfig { max_concurrent_requests: 0, ..Default::default() });
let error = client
.call(
&provider,
TransactionRequest::default()
.to(address!("1111111111111111111111111111111111111111")),
)
.await
.unwrap_err();
assert!(
matches!(error, CcipReadError::InvalidConfig(message) if message.contains("max_concurrent_requests"))
);
}
}