use super::{
BackendError, BroadcastResult, ChainBackend, EffectObservation, EffectStatus, FeePolicy, ObservationScope,
SignedAttempt, TransactionObservation, TransportEffectStatus,
};
use alloy::{
consensus::{SignableTransaction, TxEip1559},
eips::{eip2718::Encodable2718, eip2930::AccessList},
network::TxSigner,
primitives::{keccak256, Address, Bytes, TxKind, B256, U256},
providers::{DynProvider, Provider, ProviderBuilder},
rpc::types::TransactionRequest,
signers::{local::PrivateKeySigner, Signer},
sol_types::{SolCall, SolError},
};
use async_trait::async_trait;
use futures::{stream::FuturesUnordered, StreamExt};
use newton_chainio::{
avs::errors::{classify_batch_item_revert, classify_top_level_revert, is_contract_revert_error},
error::BatchItemError,
};
use newton_core::{
batch_task_manager::{
BatchTaskManager::{
batchCreateAndRespondToTasksCall, batchRespondToTasksCall, BatchPartialFailure, BatchTaskManagerInstance,
},
INewtonPolicy as BatchNewtonPolicy, INewtonProverTaskManager as BatchTaskTypes,
NewtonMessage as BatchNewtonMessage,
},
common::provider_answered,
ecdsa_operator_table_updater::ECDSAOperatorTableUpdater::{
calculateOperatorTableLeafCall, confirmGlobalTableRoot_0Call, getCertificateVerifierCall,
getGlobalTableRootByTimestampCall, updateOperatorTableCall, GlobalTableRootInFuture, InvalidGlobalTableRoot,
},
newton_prover_task_manager::{
INewtonPolicy as TaskNewtonPolicy, INewtonProverTaskManager as TaskTypes, NewtonMessage as TaskNewtonMessage,
NewtonProverTaskManager,
},
view_bn254_certificate_verifier::ViewBN254CertificateVerifier::{
isReferenceTimestampSetCall, latestReferenceTimestampCall, OperatorSet as VerifierOperatorSet,
},
};
use newton_metric::ChainRpcOutcome;
use newton_submission_service::ExecutableIntent;
use newton_task_submission::BatchIntentItem;
use std::{
future::Future,
sync::atomic::{AtomicUsize, Ordering},
time::Duration,
};
use tokio::sync::OnceCell;
const BASE_FEE_HEADROOM_NUM: u128 = 5;
const BASE_FEE_HEADROOM_DEN: u128 = 2;
const MAX_CLASSIFICATION_CONCURRENCY: usize = 8;
fn error_outcome(error: &(dyn std::error::Error + 'static)) -> ChainRpcOutcome {
if provider_answered(error) {
ChainRpcOutcome::Rejected
} else {
ChainRpcOutcome::Failure
}
}
fn provider_id(index: usize) -> String {
format!("rpc-{index}")
}
#[derive(Debug, Default)]
struct ProviderFailures {
messages: Vec<String>,
saw_transient: bool,
first_permanent: Option<BackendError>,
}
impl ProviderFailures {
fn record(&mut self, provider: &str, error: BackendError) {
self.messages.push(format!("{provider}: {error}"));
if error.is_transient() {
self.saw_transient = true;
} else {
self.first_permanent.get_or_insert(error);
}
}
fn into_error(self, operation: &str) -> BackendError {
if !self.saw_transient {
if let Some(error) = self.first_permanent {
return error;
}
}
BackendError::RpcTransient(format!(
"all providers failed {operation}: {}",
self.messages.join("; ")
))
}
}
fn finish_transaction_observation(
provider_count: usize,
pending_count: usize,
failures: ProviderFailures,
) -> Result<TransactionObservation, BackendError> {
if pending_count == provider_count {
Ok(TransactionObservation::Pending)
} else if pending_count > 0 {
Ok(TransactionObservation::PendingInconclusive {
pending_providers: pending_count,
failed_providers: failures.messages.len(),
})
} else {
Err(failures.into_error("transaction observation"))
}
}
fn finish_transport_effect_observation(
provider_count: usize,
missing_count: usize,
conflict_count: usize,
failures: ProviderFailures,
) -> Result<TransportEffectStatus, BackendError> {
if missing_count == 0 && conflict_count == 0 {
return Err(failures.into_error("transporter effect classification"));
}
if failures.messages.is_empty() {
if conflict_count == provider_count {
return Ok(TransportEffectStatus::Conflict);
}
return Ok(TransportEffectStatus::Missing);
}
Err(BackendError::RpcTransient(format!(
"transporter effect observation is inconclusive: {missing_count} provider(s) reported missing, \
{conflict_count} reported conflict; {}",
failures.messages.join("; ")
)))
}
#[derive(Debug)]
struct RpcProvider {
id: String,
client: DynProvider,
chain_verified: OnceCell<()>,
}
#[derive(Debug, Clone, Copy)]
enum BackendContract {
Task(Address),
Transporter(Address),
}
#[derive(Debug)]
pub struct AlloyBackend {
signer: PrivateKeySigner,
chain_id: u64,
contract: BackendContract,
providers: Vec<RpcProvider>,
fee_policy: FeePolicy,
finality_confirmations: u64,
rpc_timeout: Duration,
next_observation_provider: AtomicUsize,
}
pub type AlloyTaskBackend = AlloyBackend;
impl AlloyBackend {
pub fn new(
signer: PrivateKeySigner,
chain_id: u64,
batch_task_manager: Address,
rpc_urls: Vec<String>,
fee_policy: FeePolicy,
finality_confirmations: u64,
rpc_timeout: Duration,
) -> Result<Self, BackendError> {
if rpc_urls.is_empty() {
return Err(BackendError::Configuration(
"at least one RPC URL is required".to_string(),
));
}
let providers = rpc_urls
.into_iter()
.enumerate()
.map(|(index, url)| {
let id = provider_id(index);
let parsed = url
.parse()
.map_err(|error| BackendError::Configuration(format!("invalid RPC URL {id}: {error}")))?;
Ok(RpcProvider {
id,
client: ProviderBuilder::new().connect_http(parsed).erased(),
chain_verified: OnceCell::new(),
})
})
.collect::<Result<Vec<_>, BackendError>>()?;
Ok(Self {
signer,
chain_id,
contract: BackendContract::Task(batch_task_manager),
providers,
fee_policy,
finality_confirmations,
rpc_timeout,
next_observation_provider: AtomicUsize::new(0),
})
}
pub fn new_transporter(
signer: PrivateKeySigner,
chain_id: u64,
operator_table_updater: Address,
rpc_urls: Vec<String>,
fee_policy: FeePolicy,
finality_confirmations: u64,
rpc_timeout: Duration,
) -> Result<Self, BackendError> {
let mut backend = Self::new(
signer,
chain_id,
operator_table_updater,
rpc_urls,
fee_policy,
finality_confirmations,
rpc_timeout,
)?;
backend.contract = BackendContract::Transporter(operator_table_updater);
Ok(backend)
}
fn calldata(&self, intent: &ExecutableIntent) -> Result<Bytes, BackendError> {
let bytes = match intent {
ExecutableIntent::BatchCreateAndRespond {
contract_role, items, ..
} => {
self.validate_role(contract_role)?;
batchCreateAndRespondToTasksCall {
tasks: items.iter().map(|item| to_batch_task(&item.task)).collect(),
responses: items.iter().map(|item| to_batch_response(&item.response)).collect(),
signatureDataArray: items.iter().map(|item| item.signature_data.clone()).collect(),
attestationDataArray: items.iter().map(|item| item.attestation_data.clone()).collect(),
}
.abi_encode()
}
ExecutableIntent::BatchRespond {
contract_role, items, ..
} => {
self.validate_role(contract_role)?;
batchRespondToTasksCall {
tasks: items.iter().map(|item| to_batch_task(&item.task)).collect(),
responses: items.iter().map(|item| to_batch_response(&item.response)).collect(),
signatureDataArray: items.iter().map(|item| item.signature_data.clone()).collect(),
attestationDataArray: items.iter().map(|item| item.attestation_data.clone()).collect(),
}
.abi_encode()
}
ExecutableIntent::ConfirmGlobalTableRoot {
contract_role,
root,
reference_timestamp,
reference_block_number,
} => {
self.validate_transporter_role(contract_role)?;
confirmGlobalTableRoot_0Call {
globalTableRoot: *root,
referenceTimestamp: *reference_timestamp,
referenceBlockNumber: *reference_block_number,
}
.abi_encode()
}
ExecutableIntent::UpdateOperatorTable {
contract_role,
reference_timestamp,
root,
operator_set_index,
proof,
operator_table_bytes,
..
} => {
self.validate_transporter_role(contract_role)?;
updateOperatorTableCall {
referenceTimestamp: *reference_timestamp,
globalTableRoot: *root,
operatorSetIndex: *operator_set_index,
proof: proof.clone(),
operatorTableBytes: operator_table_bytes.clone(),
}
.abi_encode()
}
};
Ok(Bytes::from(bytes))
}
fn validate_role(&self, role: &str) -> Result<(), BackendError> {
if !matches!(self.contract, BackendContract::Task(_)) {
return Err(BackendError::Policy(
"transporter signer cannot execute task intent".to_string(),
));
}
if role != "batch_task_manager" {
return Err(BackendError::Policy(format!(
"task signer cannot use contract role {role}"
)));
}
Ok(())
}
fn validate_transporter_role(&self, role: &str) -> Result<(), BackendError> {
if !matches!(self.contract, BackendContract::Transporter(_)) {
return Err(BackendError::Policy(
"task signer cannot execute transporter intent".to_string(),
));
}
if role != "operator_table_updater" {
return Err(BackendError::Policy(format!(
"transporter signer cannot use contract role {role}"
)));
}
Ok(())
}
const fn target_address(&self) -> Address {
match self.contract {
BackendContract::Task(address) | BackendContract::Transporter(address) => address,
}
}
async fn rpc_call<T, E, F>(
&self,
provider_index: usize,
operation: &'static str,
future: F,
) -> Result<T, BackendError>
where
E: std::error::Error + 'static,
F: Future<Output = Result<T, E>>,
{
match tokio::time::timeout(self.rpc_timeout, future).await {
Ok(Ok(value)) => {
newton_metric::record_chain_rpc(self.chain_id, provider_index, operation, ChainRpcOutcome::Success);
Ok(value)
}
Ok(Err(error)) => {
newton_metric::record_chain_rpc(self.chain_id, provider_index, operation, error_outcome(&error));
Err(BackendError::rpc(error))
}
Err(_) => {
newton_metric::record_chain_rpc(self.chain_id, provider_index, operation, ChainRpcOutcome::Failure);
Err(BackendError::Timeout {
operation,
timeout: self.rpc_timeout,
})
}
}
}
async fn task_manager_address(
&self,
provider_index: usize,
provider: &DynProvider,
) -> Result<Address, BackendError> {
let BackendContract::Task(batch_task_manager) = self.contract else {
return Err(BackendError::Policy(
"transporter backend cannot inspect task-manager effects".to_string(),
));
};
let batch = BatchTaskManagerInstance::new(batch_task_manager, provider.clone());
self.rpc_call(provider_index, "task_manager", async {
batch.taskManager().call().await
})
.await
}
async fn inspect_item(
&self,
provider_index: usize,
provider: &DynProvider,
task_manager: Address,
item: &BatchIntentItem,
) -> Result<EffectObservation, BackendError> {
let manager = NewtonProverTaskManager::new(task_manager, provider.clone());
let task_hash = async { manager.taskHash(item.task.taskId).call().await };
let response_hash = async { manager.normalizedTaskResponseHash(item.task.taskId).call().await };
let (stored_task_hash, stored_response_hash) = self
.rpc_call(
provider_index,
"classify_effect",
futures::future::try_join(task_hash, response_hash),
)
.await?;
let task_matches = stored_task_hash == item.expected_task_hash;
let response_matches = stored_response_hash == item.expected_response_hash;
let status = if task_matches && response_matches {
EffectStatus::Verified
} else if stored_task_hash == B256::ZERO && stored_response_hash == B256::ZERO {
EffectStatus::Missing
} else if task_matches && stored_response_hash == B256::ZERO {
EffectStatus::TaskOnly
} else {
EffectStatus::Conflict
};
if status == EffectStatus::Conflict {
tracing::warn!(
task_id = %item.task.taskId,
stored_task_hash = %stored_task_hash,
expected_task_hash = %item.expected_task_hash,
stored_response_hash = %stored_response_hash,
expected_response_hash = %item.expected_response_hash,
"on-chain task effect conflicts with durable submission"
);
}
Ok(EffectObservation {
submission_id: item.submission_id,
status,
})
}
async fn ensure_chain(&self, provider_index: usize, provider: &RpcProvider) -> Result<(), BackendError> {
provider
.chain_verified
.get_or_try_init(|| async {
let observed_chain_id = self
.rpc_call(provider_index, "chain_id", provider.client.get_chain_id())
.await?;
if observed_chain_id != self.chain_id {
return Err(BackendError::Policy(format!(
"RPC {} chain id {observed_chain_id} does not match configured chain {}",
provider.id, self.chain_id
)));
}
Ok(())
})
.await
.map(|_| ())
}
async fn validate_transport_leaf(
&self,
provider_index: usize,
provider: &DynProvider,
intent: &ExecutableIntent,
) -> Result<(), BackendError> {
let ExecutableIntent::UpdateOperatorTable {
operator_table_bytes,
expected_leaf,
..
} = intent
else {
return Ok(());
};
self.validate_transporter_role(intent.contract_role())?;
let input = Bytes::from(
calculateOperatorTableLeafCall {
operatorTableBytes: operator_table_bytes.clone(),
}
.abi_encode(),
);
let request = TransactionRequest::default()
.to(self.target_address())
.input(input.into());
let output = self
.rpc_call(provider_index, "calculate_operator_table_leaf", async {
provider.call(request).await
})
.await?;
let observed = calculateOperatorTableLeafCall::abi_decode_returns(&output)
.map_err(|error| BackendError::Decode(error.to_string()))?;
if observed != *expected_leaf {
return Err(BackendError::Policy(format!(
"operator table leaf {observed} does not match persisted expected leaf {expected_leaf}"
)));
}
Ok(())
}
async fn market_fees(
&self,
provider_index: usize,
provider: &DynProvider,
) -> Result<Option<(u128, u128)>, BackendError> {
let fee_history = self
.rpc_call(
provider_index,
"fee_history",
provider.get_fee_history(10, Default::default(), &[90.0]),
)
.await?;
let Some(base_fee) = fee_history.latest_block_base_fee() else {
return Ok(None);
};
let priority = fee_history
.reward
.as_ref()
.and_then(|rewards| {
rewards
.iter()
.filter_map(|block_rewards| block_rewards.first().copied())
.max()
})
.unwrap_or(2_000_000_000)
.saturating_mul(u128::from(100u32.saturating_add(self.fee_policy.gas_bump_percent)))
/ 100;
Ok(Some((base_fee, priority)))
}
async fn escalated_fees(&self, previous: (u128, u128)) -> (u128, u128) {
for (provider_index, provider) in self.providers.iter().enumerate() {
if self.ensure_chain(provider_index, provider).await.is_err() {
continue;
}
if let Ok(Some(market)) = self.market_fees(provider_index, &provider.client).await {
return escalate_floor(
Some(market),
previous.0,
previous.1,
self.fee_policy.max_fee_per_gas_ceiling,
);
}
}
escalate_floor(None, previous.0, previous.1, self.fee_policy.max_fee_per_gas_ceiling)
}
async fn prepare_with_provider(
&self,
provider_index: usize,
provider: &DynProvider,
intent: &ExecutableIntent,
nonce: u64,
previous_fees: Option<(u128, u128)>,
) -> Result<SignedAttempt, BackendError> {
self.validate_transport_leaf(provider_index, provider, intent).await?;
let input = self.calldata(intent)?;
let request = TransactionRequest::default()
.from(self.address())
.to(self.target_address())
.input(input.clone().into())
.nonce(nonce)
.value(U256::ZERO);
match tokio::time::timeout(self.rpc_timeout, async { provider.call(request.clone()).await }).await {
Ok(Ok(_)) => {
newton_metric::record_chain_rpc(self.chain_id, provider_index, "simulate", ChainRpcOutcome::Success)
}
Ok(Err(error)) => {
newton_metric::record_chain_rpc(self.chain_id, provider_index, "simulate", error_outcome(&error));
let error = alloy::contract::Error::TransportError(error);
if let Some(revert_data) = error.as_revert_data() {
return Err(classify_simulation_revert(intent, &revert_data));
}
if is_contract_revert_error(&error) {
return Err(BackendError::Simulation(error.to_string()));
}
return Err(BackendError::rpc(error));
}
Err(_) => {
newton_metric::record_chain_rpc(self.chain_id, provider_index, "simulate", ChainRpcOutcome::Failure);
return Err(BackendError::Timeout {
operation: "simulate",
timeout: self.rpc_timeout,
});
}
}
let gas_limit = self
.rpc_call(provider_index, "estimate_gas", async {
provider.estimate_gas(request).await
})
.await?;
let (max_fee_per_gas, max_priority_fee_per_gas) = match previous_fees {
Some(previous) => self.escalated_fees(previous).await,
None => {
let estimate = self
.rpc_call(provider_index, "estimate_fees", provider.estimate_eip1559_fees())
.await?;
(estimate.max_fee_per_gas, estimate.max_priority_fee_per_gas)
}
};
self.sign_attempt(
nonce,
gas_limit,
max_fee_per_gas,
max_priority_fee_per_gas,
TxKind::Call(self.target_address()),
input,
)
.await
}
async fn sign_attempt(
&self,
nonce: u64,
gas_limit: u64,
max_fee_per_gas: u128,
max_priority_fee_per_gas: u128,
to: TxKind,
input: Bytes,
) -> Result<SignedAttempt, BackendError> {
if max_priority_fee_per_gas > max_fee_per_gas {
return Err(BackendError::Policy("priority fee exceeds max fee per gas".to_string()));
}
let mut transaction = TxEip1559 {
chain_id: self.chain_id,
nonce,
gas_limit,
max_fee_per_gas,
max_priority_fee_per_gas,
to,
value: U256::ZERO,
access_list: AccessList::default(),
input,
};
let signature = self
.signer
.sign_transaction(&mut transaction)
.await
.map_err(|error| BackendError::Signing(error.to_string()))?;
let envelope: alloy::consensus::TxEnvelope = transaction.into_signed(signature).into();
let raw_transaction = Bytes::from(envelope.encoded_2718());
Ok(SignedAttempt {
nonce,
gas_limit,
max_fee_per_gas,
max_priority_fee_per_gas,
transaction_hash: keccak256(&raw_transaction),
raw_transaction,
})
}
async fn transaction_count(&self, pending: bool) -> Result<u64, BackendError> {
let mut failures = ProviderFailures::default();
for (provider_index, provider) in self.providers.iter().enumerate() {
if let Err(error) = self.ensure_chain(provider_index, provider).await {
failures.record(&provider.id, error);
continue;
}
let call = provider.client.get_transaction_count(self.address());
let result = if pending {
self.rpc_call(provider_index, "pending_nonce", async { call.pending().await })
.await
} else {
self.rpc_call(provider_index, "latest_nonce", async { call.latest().await })
.await
};
match result {
Ok(count) => return Ok(count),
Err(error) => failures.record(&provider.id, error),
}
}
Err(failures.into_error("transaction-count read"))
}
async fn native_balance(&self) -> Result<U256, BackendError> {
let mut failures = ProviderFailures::default();
for (provider_index, provider) in self.providers.iter().enumerate() {
if let Err(error) = self.ensure_chain(provider_index, provider).await {
failures.record(&provider.id, error);
continue;
}
match self
.rpc_call(provider_index, "balance", async {
provider.client.get_balance(self.address()).latest().await
})
.await
{
Ok(balance) => return Ok(balance),
Err(error) => failures.record(&provider.id, error),
}
}
Err(failures.into_error("native-balance read"))
}
}
fn escalate_floor(
market: Option<(u128, u128)>,
previous_max_fee: u128,
previous_priority_fee: u128,
ceiling: u128,
) -> (u128, u128) {
let market_priority = market.map(|(_, priority)| priority).unwrap_or(0);
let priority = market_priority.max(previous_priority_fee.saturating_mul(110) / 100);
let previous_floor = previous_max_fee.saturating_mul(110) / 100;
let max_fee = match market {
Some((base_fee, _)) => previous_floor.max(
base_fee
.saturating_mul(BASE_FEE_HEADROOM_NUM)
.saturating_div(BASE_FEE_HEADROOM_DEN)
.saturating_add(priority),
),
None => previous_floor,
}
.min(ceiling);
(max_fee, priority.min(max_fee))
}
fn to_batch_task(value: &TaskTypes::Task) -> BatchTaskTypes::Task {
BatchTaskTypes::Task {
taskId: value.taskId,
policyClient: value.policyClient,
policyId: value.policyId,
policyRevision: value.policyRevision,
taskCreatedBlock: value.taskCreatedBlock,
quorumThresholdPercentage: value.quorumThresholdPercentage,
intent: to_batch_intent(&value.intent),
intentSignature: value.intentSignature.clone(),
policies: value.policies.iter().map(to_batch_policy_spec).collect(),
wasmArgs: value.wasmArgs.clone(),
quorumNumbers: value.quorumNumbers.clone(),
initializationTimestamp: value.initializationTimestamp,
}
}
fn to_batch_response(value: &TaskTypes::TaskResponse) -> BatchTaskTypes::TaskResponse {
BatchTaskTypes::TaskResponse {
taskId: value.taskId,
policyClient: value.policyClient,
policyId: value.policyId,
intent: to_batch_intent(&value.intent),
intentSignature: value.intentSignature.clone(),
allowed: value.allowed,
policyTaskData: value.policyTaskData.iter().map(to_batch_policy_task_data).collect(),
initializationTimestamp: value.initializationTimestamp,
}
}
fn to_batch_policy_task_data(value: &TaskNewtonMessage::PolicyTaskData) -> BatchNewtonMessage::PolicyTaskData {
BatchNewtonMessage::PolicyTaskData {
policyId: value.policyId,
policyAddress: value.policyAddress,
policy: value.policy.clone(),
policyData: value.policyData.iter().map(to_batch_policy_data).collect(),
}
}
fn to_batch_policy_data(value: &TaskNewtonMessage::PolicyData) -> BatchNewtonMessage::PolicyData {
BatchNewtonMessage::PolicyData {
wasmArgs: value.wasmArgs.clone(),
data: value.data.clone(),
expireBlock: value.expireBlock,
}
}
fn to_batch_intent(value: &TaskNewtonMessage::Intent) -> BatchNewtonMessage::Intent {
BatchNewtonMessage::Intent {
from: value.from,
to: value.to,
value: value.value,
data: value.data.clone(),
chainId: value.chainId,
functionSignature: value.functionSignature.clone(),
}
}
fn to_batch_policy_spec(
value: &newton_core::newton_prover_task_manager::INewtonPolicyClient::PolicySpec,
) -> newton_core::batch_task_manager::INewtonPolicyClient::PolicySpec {
newton_core::batch_task_manager::INewtonPolicyClient::PolicySpec {
policy: value.policy,
config: to_batch_policy_config(&value.config),
}
}
fn to_batch_policy_config(
value: &newton_core::newton_prover_task_manager::INewtonPolicy::PolicyConfig,
) -> newton_core::batch_task_manager::INewtonPolicy::PolicyConfig {
newton_core::batch_task_manager::INewtonPolicy::PolicyConfig {
policyParams: value.policyParams.clone(),
expireAfter: value.expireAfter,
}
}
fn classify_simulation_revert(intent: &ExecutableIntent, revert_data: &[u8]) -> BackendError {
if matches!(intent, ExecutableIntent::ConfirmGlobalTableRoot { .. })
&& revert_data.starts_with(&GlobalTableRootInFuture::SELECTOR)
{
return BackendError::Simulation("GlobalTableRootInFuture".to_string());
}
if matches!(intent, ExecutableIntent::UpdateOperatorTable { .. })
&& revert_data.starts_with(&InvalidGlobalTableRoot::SELECTOR)
{
return BackendError::Simulation("InvalidGlobalTableRoot".to_string());
}
let item_count = match intent {
ExecutableIntent::BatchCreateAndRespond { items, .. } | ExecutableIntent::BatchRespond { items, .. } => {
items.len()
}
ExecutableIntent::ConfirmGlobalTableRoot { .. } | ExecutableIntent::UpdateOperatorTable { .. } => 0,
};
if item_count == 0 {
return BackendError::Simulation(classify_top_level_revert(revert_data).to_string());
}
classify_simulation_revert_for_items(item_count, revert_data)
}
fn classify_simulation_revert_for_items(item_count: usize, revert_data: &[u8]) -> BackendError {
if item_count == 1 {
if let Ok(partial) = BatchPartialFailure::abi_decode(revert_data) {
if let Some(failure) = partial.failures.first() {
let classified = classify_batch_item_revert(&failure.reason);
if classified.is_retryable() {
return BackendError::SimulationRetryable(format_batch_item_error(&classified));
}
return BackendError::Simulation(format_batch_item_error(&classified));
}
}
}
BackendError::Simulation(classify_top_level_revert(revert_data).to_string())
}
fn format_batch_item_error(error: &BatchItemError) -> String {
match error {
BatchItemError::TaskAlreadyExists => "TaskAlreadyExists".to_string(),
BatchItemError::TaskAlreadyResponded => "TaskAlreadyResponded".to_string(),
BatchItemError::LikelyOutOfGas { gas_forwarded } => {
format!("LikelyOutOfGas(gas_forwarded={gas_forwarded:?})")
}
BatchItemError::InsufficientGasForItem { gas_left } => {
format!("InsufficientGasForItem(gas_left={gas_left:?})")
}
BatchItemError::ContractRevert { selector, name } => {
format!("ContractRevert({name}, 0x{selector})")
}
BatchItemError::Unknown { raw } => format!("UnknownRevert(0x{})", hex::encode(raw)),
}
}
#[async_trait]
impl ChainBackend for AlloyBackend {
fn address(&self) -> Address {
self.signer.address()
}
fn finality_confirmations(&self) -> u64 {
self.finality_confirmations.max(1)
}
async fn latest_transaction_count(&self) -> Result<u64, BackendError> {
self.transaction_count(false).await
}
async fn pending_transaction_count(&self) -> Result<u64, BackendError> {
self.transaction_count(true).await
}
async fn balance(&self) -> Result<U256, BackendError> {
self.native_balance().await
}
async fn prepare(
&self,
intent: &ExecutableIntent,
nonce: u64,
previous_fees: Option<(u128, u128)>,
) -> Result<SignedAttempt, BackendError> {
let mut failures = Vec::new();
let mut saw_transient = false;
let mut simulation_failure = None;
let mut permanent_failure = None;
for (provider_index, provider) in self.providers.iter().enumerate() {
if let Err(error) = self.ensure_chain(provider_index, provider).await {
failures.push(format!("{}: {error}", provider.id));
if error.is_transient() {
saw_transient = true;
} else {
permanent_failure.get_or_insert(error);
}
continue;
}
match self
.prepare_with_provider(provider_index, &provider.client, intent, nonce, previous_fees)
.await
{
Ok(attempt) => return Ok(attempt),
Err(error) if error.is_transient() => {
saw_transient = true;
failures.push(format!("{}: {error}", provider.id));
}
Err(error @ BackendError::Simulation(_)) => {
failures.push(format!("{}: {error}", provider.id));
simulation_failure.get_or_insert(error);
}
Err(error) => {
failures.push(format!("{}: {error}", provider.id));
permanent_failure.get_or_insert(error);
}
}
}
if !saw_transient {
if let Some(error) = simulation_failure {
return Err(error);
}
if let Some(error) = permanent_failure {
return Err(error);
}
}
Err(BackendError::RpcTransient(format!(
"all providers failed transaction preparation: {}",
failures.join("; ")
)))
}
async fn prepare_cancellation(&self, nonce: u64, fees: (u128, u128)) -> Result<SignedAttempt, BackendError> {
let (max_fee_per_gas, max_priority_fee_per_gas) = self.escalated_fees(fees).await;
self.sign_attempt(
nonce,
21_000,
max_fee_per_gas,
max_priority_fee_per_gas,
TxKind::Call(self.address()),
Bytes::new(),
)
.await
}
async fn broadcast(&self, raw_transaction: &Bytes) -> Vec<BroadcastResult> {
let expected_hash = keccak256(raw_transaction);
let mut results = Vec::with_capacity(self.providers.len());
for (provider_index, provider) in self.providers.iter().enumerate() {
let result = match self
.rpc_call(
provider_index,
"broadcast",
provider.client.send_raw_transaction(raw_transaction),
)
.await
{
Ok(pending) if *pending.tx_hash() == expected_hash => Ok(expected_hash),
Ok(pending) => Err(format!(
"provider returned transaction hash {}, expected {expected_hash}",
pending.tx_hash()
)),
Err(error) => Err(error.to_string()),
};
let accepted = result.is_ok();
results.push(BroadcastResult {
provider: provider.id.clone(),
result,
});
if accepted {
break;
}
}
results
}
async fn observe(
&self,
transaction_hash: B256,
receipt_provider: Option<&str>,
scope: ObservationScope,
) -> Result<TransactionObservation, BackendError> {
let pinned_index = receipt_provider.and_then(|id| self.providers.iter().position(|provider| provider.id == id));
if scope == ObservationScope::Economical {
let provider_index = pinned_index.unwrap_or_else(|| {
self.next_observation_provider.fetch_add(1, Ordering::Relaxed) % self.providers.len()
});
return self
.observe_with_provider(provider_index, &self.providers[provider_index], transaction_hash)
.await;
}
let provider_order = pinned_index
.into_iter()
.chain((0..self.providers.len()).filter(|index| Some(*index) != pinned_index));
let mut failures = ProviderFailures::default();
let mut pending_count = 0_usize;
for provider_index in provider_order {
let provider = &self.providers[provider_index];
match self
.observe_with_provider(provider_index, provider, transaction_hash)
.await
{
Ok(TransactionObservation::Pending) => pending_count += 1,
Ok(observation) => return Ok(observation),
Err(error) => failures.record(&provider.id, error),
}
}
finish_transaction_observation(self.providers.len(), pending_count, failures)
}
async fn classify_onchain_effects(
&self,
intent: &ExecutableIntent,
) -> Result<Vec<EffectObservation>, BackendError> {
let mut failures = ProviderFailures::default();
for (provider_index, provider) in self.providers.iter().enumerate() {
match self.classify_with_provider(provider_index, provider, intent).await {
Ok(observations) => return Ok(observations),
Err(error) => failures.record(&provider.id, error),
}
}
Err(failures.into_error("on-chain effect classification"))
}
async fn classify_transport_effect(
&self,
intent: &ExecutableIntent,
) -> Result<TransportEffectStatus, BackendError> {
let mut failures = ProviderFailures::default();
let mut missing_count = 0_usize;
let mut conflict_count = 0_usize;
for (provider_index, provider) in self.providers.iter().enumerate() {
match self
.classify_transport_with_provider(provider_index, provider, intent)
.await
{
Ok(TransportEffectStatus::Verified) => return Ok(TransportEffectStatus::Verified),
Ok(TransportEffectStatus::Missing) => missing_count += 1,
Ok(TransportEffectStatus::Conflict) => conflict_count += 1,
Err(error) => failures.record(&provider.id, error),
}
}
finish_transport_effect_observation(self.providers.len(), missing_count, conflict_count, failures)
}
}
impl AlloyBackend {
async fn observe_with_provider(
&self,
provider_index: usize,
provider: &RpcProvider,
transaction_hash: B256,
) -> Result<TransactionObservation, BackendError> {
self.ensure_chain(provider_index, provider).await?;
let provider_id = &provider.id;
let Some(receipt) = self
.rpc_call(
provider_index,
"receipt",
provider.client.get_transaction_receipt(transaction_hash),
)
.await?
else {
return Ok(TransactionObservation::Pending);
};
let block_number = receipt
.block_number
.ok_or_else(|| BackendError::RpcPermanent(format!("{provider_id}: mined receipt lacks block number")))?;
let block_hash = receipt
.block_hash
.ok_or_else(|| BackendError::RpcPermanent(format!("{provider_id}: mined receipt lacks block hash")))?;
let Some(canonical_block) = self
.rpc_call(provider_index, "receipt_canonical_block", async {
provider.client.get_block_by_number(block_number.into()).await
})
.await?
else {
return Ok(TransactionObservation::Pending);
};
if canonical_block.header.hash != block_hash {
return Ok(TransactionObservation::Pending);
}
let head = self
.rpc_call(provider_index, "block_number", provider.client.get_block_number())
.await?;
let confirmations = head.saturating_sub(block_number).saturating_add(1);
let status = receipt.status();
let encoded = serde_json::to_vec(&receipt).map_err(|error| BackendError::Decode(error.to_string()))?;
if status {
Ok(TransactionObservation::Mined {
provider: provider_id.clone(),
block_number,
block_hash,
confirmations,
receipt: encoded,
})
} else {
Ok(TransactionObservation::Reverted {
provider: provider_id.clone(),
block_number,
block_hash,
confirmations,
receipt: encoded,
})
}
}
async fn classify_with_provider(
&self,
provider_index: usize,
provider: &RpcProvider,
intent: &ExecutableIntent,
) -> Result<Vec<EffectObservation>, BackendError> {
self.ensure_chain(provider_index, provider).await?;
let items = match intent {
ExecutableIntent::BatchCreateAndRespond { items, .. } | ExecutableIntent::BatchRespond { items, .. } => {
items
}
ExecutableIntent::ConfirmGlobalTableRoot { .. } | ExecutableIntent::UpdateOperatorTable { .. } => {
return Err(BackendError::Policy(
"transporter intent cannot use task effect classification".to_string(),
));
}
};
let task_manager = self.task_manager_address(provider_index, &provider.client).await?;
let mut remaining = items.iter();
let mut pending = FuturesUnordered::new();
for item in remaining.by_ref().take(MAX_CLASSIFICATION_CONCURRENCY) {
pending.push(self.inspect_item(provider_index, &provider.client, task_manager, item));
}
let mut observations = Vec::with_capacity(items.len());
while let Some(observation) = pending.next().await {
observations.push(observation?);
if let Some(item) = remaining.next() {
pending.push(self.inspect_item(provider_index, &provider.client, task_manager, item));
}
}
Ok(observations)
}
async fn classify_transport_with_provider(
&self,
provider_index: usize,
provider: &RpcProvider,
intent: &ExecutableIntent,
) -> Result<TransportEffectStatus, BackendError> {
self.ensure_chain(provider_index, provider).await?;
self.validate_transporter_role(intent.contract_role())?;
match intent {
ExecutableIntent::ConfirmGlobalTableRoot {
root,
reference_timestamp,
..
} => {
let request = TransactionRequest::default().to(self.target_address()).input(
Bytes::from(
getGlobalTableRootByTimestampCall {
referenceTimestamp: *reference_timestamp,
}
.abi_encode(),
)
.into(),
);
let output = self
.rpc_call(provider_index, "classify_transport_root", async {
provider.client.call(request).await
})
.await?;
let observed = getGlobalTableRootByTimestampCall::abi_decode_returns(&output)
.map_err(|error| BackendError::Decode(error.to_string()))?;
if observed == *root {
Ok(TransportEffectStatus::Verified)
} else if observed == B256::ZERO {
Ok(TransportEffectStatus::Missing)
} else {
Ok(TransportEffectStatus::Conflict)
}
}
ExecutableIntent::UpdateOperatorTable {
reference_timestamp,
operator_table_bytes,
..
} => {
let (operator_set, curve_type) = decode_transport_table_identity(operator_table_bytes)?;
let verifier_request = TransactionRequest::default()
.to(self.target_address())
.input(Bytes::from(getCertificateVerifierCall { curveType: curve_type }.abi_encode()).into());
let verifier_output = self
.rpc_call(provider_index, "classify_transport_verifier", async {
provider.client.call(verifier_request).await
})
.await?;
let verifier = getCertificateVerifierCall::abi_decode_returns(&verifier_output)
.map_err(|error| BackendError::Decode(error.to_string()))?;
if verifier == Address::ZERO {
return Err(BackendError::Decode(
"operator-table updater returned a zero certificate verifier".to_string(),
));
}
let is_set_request = TransactionRequest::default().to(verifier).input(
Bytes::from(
isReferenceTimestampSetCall {
operatorSet: operator_set.clone(),
referenceTimestamp: *reference_timestamp,
}
.abi_encode(),
)
.into(),
);
let is_set_output = self
.rpc_call(provider_index, "classify_transport_table", async {
provider.client.call(is_set_request).await
})
.await?;
let is_set = isReferenceTimestampSetCall::abi_decode_returns(&is_set_output)
.map_err(|error| BackendError::Decode(error.to_string()))?;
if is_set {
return Ok(TransportEffectStatus::Verified);
}
let latest_request = TransactionRequest::default().to(verifier).input(
Bytes::from(
latestReferenceTimestampCall {
operatorSet: operator_set,
}
.abi_encode(),
)
.into(),
);
let latest_output = self
.rpc_call(provider_index, "classify_transport_table_latest", async {
provider.client.call(latest_request).await
})
.await?;
let latest = latestReferenceTimestampCall::abi_decode_returns(&latest_output)
.map_err(|error| BackendError::Decode(error.to_string()))?;
if latest > *reference_timestamp {
Ok(TransportEffectStatus::Conflict)
} else {
Ok(TransportEffectStatus::Missing)
}
}
ExecutableIntent::BatchCreateAndRespond { .. } | ExecutableIntent::BatchRespond { .. } => Err(
BackendError::Policy("task intent cannot use transporter effect classification".to_string()),
),
}
}
}
fn decode_transport_table_identity(operator_table_bytes: &Bytes) -> Result<(VerifierOperatorSet, u8), BackendError> {
const IDENTITY_PREFIX_BYTES: usize = 3 * 32;
if operator_table_bytes.len() < IDENTITY_PREFIX_BYTES {
return Err(BackendError::Decode(
"operator table is shorter than its ABI identity prefix".to_string(),
));
}
let avs_word = &operator_table_bytes[..32];
if avs_word[..12].iter().any(|byte| *byte != 0) {
return Err(BackendError::Decode(
"operator table contains a non-canonical AVS address".to_string(),
));
}
let id_word = &operator_table_bytes[32..64];
if id_word[..28].iter().any(|byte| *byte != 0) {
return Err(BackendError::Decode(
"operator table operator-set ID exceeds uint32".to_string(),
));
}
let curve_word = &operator_table_bytes[64..96];
if curve_word[..31].iter().any(|byte| *byte != 0) {
return Err(BackendError::Decode(
"operator table curve type exceeds uint8".to_string(),
));
}
let operator_set = VerifierOperatorSet {
avs: Address::from_slice(&avs_word[12..]),
id: u32::from_be_bytes(id_word[28..].try_into().expect("four-byte slice")),
};
Ok((operator_set, curve_word[31]))
}
#[cfg(test)]
mod tests {
use super::*;
use newton_core::batch_task_manager::IBatchTaskManager::FailedItem;
fn backend(max_fee: u128) -> AlloyTaskBackend {
AlloyTaskBackend::new(
PrivateKeySigner::random(),
31_337,
Address::repeat_byte(1),
vec!["http://127.0.0.1:8545".to_string()],
FeePolicy {
gas_bump_percent: 20,
max_fee_per_gas_ceiling: max_fee,
},
1,
Duration::from_secs(1),
)
.expect("backend")
}
fn transporter_backend() -> AlloyBackend {
AlloyBackend::new_transporter(
PrivateKeySigner::random(),
31_338,
Address::repeat_byte(2),
vec!["http://127.0.0.1:8545".to_string()],
FeePolicy {
gas_bump_percent: 20,
max_fee_per_gas_ceiling: 1_000,
},
1,
Duration::from_secs(1),
)
.expect("transporter backend")
}
#[test]
fn provider_ids_are_stable_configuration_slots() {
assert_eq!(provider_id(0), provider_id(0));
assert_ne!(provider_id(0), provider_id(1));
assert_eq!(provider_id(0), "rpc-0");
}
#[test]
fn transporter_backend_encodes_only_allowlisted_updater_calls() {
let backend = transporter_backend();
let confirm = ExecutableIntent::ConfirmGlobalTableRoot {
contract_role: "operator_table_updater".to_string(),
root: B256::repeat_byte(3),
reference_timestamp: 4,
reference_block_number: 5,
};
let update = ExecutableIntent::UpdateOperatorTable {
contract_role: "operator_table_updater".to_string(),
reference_timestamp: 4,
root: B256::repeat_byte(3),
operator_set_index: 0,
proof: Bytes::from(vec![6]),
operator_table_bytes: Bytes::from(vec![7]),
expected_leaf: B256::repeat_byte(8),
};
assert_eq!(
&backend.calldata(&confirm).expect("confirm calldata")[..4],
confirmGlobalTableRoot_0Call::SELECTOR
);
assert_eq!(
&backend.calldata(&update).expect("update calldata")[..4],
updateOperatorTableCall::SELECTOR
);
let wrong_role = ExecutableIntent::ConfirmGlobalTableRoot {
contract_role: "batch_task_manager".to_string(),
root: B256::ZERO,
reference_timestamp: 0,
reference_block_number: 0,
};
assert!(matches!(backend.calldata(&wrong_role), Err(BackendError::Policy(_))));
}
#[test]
fn task_and_transporter_backends_reject_each_others_intents() {
let transport_intent = ExecutableIntent::ConfirmGlobalTableRoot {
contract_role: "operator_table_updater".to_string(),
root: B256::ZERO,
reference_timestamp: 0,
reference_block_number: 0,
};
assert!(matches!(
backend(1_000).calldata(&transport_intent),
Err(BackendError::Policy(_))
));
let task_intent = ExecutableIntent::BatchRespond {
contract_role: "batch_task_manager".to_string(),
items: Vec::new(),
};
assert!(matches!(
transporter_backend().calldata(&task_intent),
Err(BackendError::Policy(_))
));
}
#[tokio::test]
async fn cancellation_preparation_is_deterministic_and_local() {
let backend = backend(1_000);
let first = backend.prepare_cancellation(7, (100, 10)).await.expect("first");
let second = backend.prepare_cancellation(7, (100, 10)).await.expect("second");
assert_eq!(first.raw_transaction, second.raw_transaction);
assert_eq!(first.transaction_hash, keccak256(&first.raw_transaction));
assert_eq!(first.nonce, 7);
}
#[tokio::test]
async fn cancellation_enforces_fee_ceiling_before_signing() {
let attempt = backend(99).prepare_cancellation(7, (100, 10)).await.expect("attempt");
assert_eq!(attempt.max_fee_per_gas, 99);
assert!(attempt.max_priority_fee_per_gas <= attempt.max_fee_per_gas);
}
#[test]
fn replacement_fees_track_market_and_preserve_the_node_floor() {
assert_eq!(escalate_floor(Some((1_000, 100)), 500, 50, u128::MAX), (2_600, 100));
assert_eq!(escalate_floor(Some((100, 40)), 500, 50, u128::MAX), (550, 55));
assert_eq!(escalate_floor(None, 500, 50, u128::MAX), (550, 55));
}
#[test]
fn rpc_classification_recognizes_transient_transport_failures() {
assert!(matches!(
BackendError::rpc("HTTP 503 Service Unavailable"),
BackendError::RpcTransient(_)
));
}
#[test]
fn provider_failure_aggregation_preserves_retryability() {
let mut permanent_only = ProviderFailures::default();
permanent_only.record("rpc-a", BackendError::Policy("wrong chain".to_string()));
assert!(matches!(
permanent_only.into_error("chain read"),
BackendError::Policy(_)
));
let mut mixed = ProviderFailures::default();
mixed.record("rpc-a", BackendError::Policy("wrong chain".to_string()));
mixed.record(
"rpc-b",
BackendError::Timeout {
operation: "chain_id",
timeout: Duration::from_secs(1),
},
);
let error = mixed.into_error("chain read");
assert!(matches!(&error, BackendError::RpcTransient(_)));
assert!(error.to_string().contains("rpc-a"));
assert!(error.to_string().contains("rpc-b"));
}
#[test]
fn partial_pending_evidence_remains_distinct_from_authoritative_pending() {
assert!(matches!(
finish_transaction_observation(2, 2, ProviderFailures::default()),
Ok(TransactionObservation::Pending)
));
let mut partial_failure = ProviderFailures::default();
partial_failure.record("rpc-1", BackendError::RpcPermanent("unauthorized".to_string()));
let uncertain = finish_transaction_observation(2, 1, partial_failure);
assert!(matches!(
uncertain,
Ok(TransactionObservation::PendingInconclusive {
pending_providers: 1,
failed_providers: 1,
})
));
}
#[test]
fn transporter_conflict_requires_complete_provider_agreement() {
assert!(matches!(
finish_transport_effect_observation(2, 0, 2, ProviderFailures::default()),
Ok(TransportEffectStatus::Conflict)
));
assert!(matches!(
finish_transport_effect_observation(2, 1, 1, ProviderFailures::default()),
Ok(TransportEffectStatus::Missing)
));
let mut incomplete = ProviderFailures::default();
incomplete.record(
"rpc-1",
BackendError::Timeout {
operation: "classify_transport_root",
timeout: Duration::from_secs(1),
},
);
assert!(matches!(
finish_transport_effect_observation(2, 0, 1, incomplete),
Err(BackendError::RpcTransient(_))
));
}
#[test]
fn single_item_batch_out_of_gas_is_retryable() {
let revert = BatchPartialFailure {
failures: vec![FailedItem {
index: U256::ZERO,
taskId: B256::ZERO,
reason: Bytes::new(),
}],
}
.abi_encode();
assert!(matches!(
classify_simulation_revert_for_items(1, &revert),
BackendError::SimulationRetryable(_)
));
}
#[test]
fn future_transporter_root_has_a_distinct_deferred_classification() {
let intent = ExecutableIntent::ConfirmGlobalTableRoot {
contract_role: "operator_table_updater".to_string(),
root: B256::repeat_byte(1),
reference_timestamp: 100,
reference_block_number: 10,
};
let error = classify_simulation_revert(&intent, &GlobalTableRootInFuture::SELECTOR);
assert!(error.is_transport_destination_behind());
}
#[test]
fn invalid_global_table_root_is_a_permanent_simulation_failure() {
let intent = ExecutableIntent::UpdateOperatorTable {
contract_role: "operator_table_updater".to_string(),
reference_timestamp: 100,
root: B256::repeat_byte(1),
operator_set_index: 0,
proof: Bytes::new(),
operator_table_bytes: Bytes::new(),
expected_leaf: B256::repeat_byte(2),
};
assert!(matches!(
classify_simulation_revert(&intent, &InvalidGlobalTableRoot::SELECTOR),
BackendError::Simulation(message) if message == "InvalidGlobalTableRoot"
));
}
#[test]
fn transporter_table_identity_decodes_abi_prefix() {
let avs = Address::repeat_byte(0x42);
let mut encoded = vec![0_u8; 96];
encoded[12..32].copy_from_slice(avs.as_slice());
encoded[60..64].copy_from_slice(&17_u32.to_be_bytes());
encoded[95] = 2;
let (operator_set, curve_type) =
decode_transport_table_identity(&Bytes::from(encoded)).expect("decode table identity");
assert_eq!(operator_set.avs, avs);
assert_eq!(operator_set.id, 17);
assert_eq!(curve_type, 2);
}
#[test]
fn transporter_table_identity_rejects_noncanonical_prefix() {
let mut encoded = vec![0_u8; 96];
encoded[32] = 1;
assert!(matches!(
decode_transport_table_identity(&Bytes::from(encoded)),
Err(BackendError::Decode(_))
));
}
#[test]
fn single_item_idempotent_revert_is_reconciled_not_retried_as_rpc() {
let revert = BatchPartialFailure {
failures: vec![FailedItem {
index: U256::ZERO,
taskId: B256::ZERO,
reason: Bytes::copy_from_slice(&newton_chainio::avs::errors::selectors::TASK_ALREADY_EXISTS),
}],
}
.abi_encode();
assert!(matches!(
classify_simulation_revert_for_items(1, &revert),
BackendError::Simulation(_)
));
}
}