use alloy::primitives::{Address, Bytes, B256, U256};
use async_trait::async_trait;
use newton_chainio::avs::errors::is_transient_rpc_error;
use newton_submission_service::ExecutableIntent;
use newton_task_submission::SubmissionId;
use std::time::Duration;
mod alloy_backend;
pub use alloy_backend::{AlloyBackend, AlloyTaskBackend};
#[derive(Debug, Clone, Copy)]
pub struct FeePolicy {
pub gas_bump_percent: u32,
pub max_fee_per_gas_ceiling: u128,
}
#[derive(Debug, Clone)]
pub struct SignedAttempt {
pub nonce: u64,
pub gas_limit: u64,
pub max_fee_per_gas: u128,
pub max_priority_fee_per_gas: u128,
pub raw_transaction: Bytes,
pub transaction_hash: B256,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct BroadcastResult {
pub provider: String,
pub result: Result<B256, String>,
}
#[derive(Debug, Clone)]
pub enum TransactionObservation {
Pending,
PendingInconclusive {
pending_providers: usize,
failed_providers: usize,
},
Mined {
provider: String,
block_number: u64,
block_hash: B256,
confirmations: u64,
receipt: Vec<u8>,
},
Reverted {
provider: String,
block_number: u64,
block_hash: B256,
confirmations: u64,
receipt: Vec<u8>,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ObservationScope {
Economical,
Exhaustive,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EffectStatus {
Verified,
Missing,
TaskOnly,
Conflict,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct EffectObservation {
pub submission_id: SubmissionId,
pub status: EffectStatus,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TransportEffectStatus {
Verified,
Missing,
Conflict,
}
#[async_trait]
pub trait ChainBackend: Send + Sync {
fn address(&self) -> Address;
fn finality_confirmations(&self) -> u64 {
1
}
async fn latest_transaction_count(&self) -> Result<u64, BackendError>;
async fn pending_transaction_count(&self) -> Result<u64, BackendError>;
async fn balance(&self) -> Result<U256, BackendError>;
async fn prepare(
&self,
intent: &ExecutableIntent,
nonce: u64,
previous_fees: Option<(u128, u128)>,
) -> Result<SignedAttempt, BackendError>;
async fn prepare_cancellation(&self, nonce: u64, fees: (u128, u128)) -> Result<SignedAttempt, BackendError>;
async fn broadcast(&self, raw_transaction: &Bytes) -> Vec<BroadcastResult>;
async fn observe(
&self,
transaction_hash: B256,
receipt_provider: Option<&str>,
scope: ObservationScope,
) -> Result<TransactionObservation, BackendError>;
async fn classify_onchain_effects(&self, intent: &ExecutableIntent)
-> Result<Vec<EffectObservation>, BackendError>;
async fn classify_transport_effect(
&self,
_intent: &ExecutableIntent,
) -> Result<TransportEffectStatus, BackendError> {
Err(BackendError::Policy(
"backend does not support transporter effect classification".to_string(),
))
}
}
#[derive(Debug, thiserror::Error)]
pub enum BackendError {
#[error("executor configuration error: {0}")]
Configuration(String),
#[error("executor policy violation: {0}")]
Policy(String),
#[error("transaction simulation failed: {0}")]
Simulation(String),
#[error("retryable transaction simulation failure: {0}")]
SimulationRetryable(String),
#[error("transaction signing failed: {0}")]
Signing(String),
#[error("transient RPC error: {0}")]
RpcTransient(String),
#[error("permanent RPC error: {0}")]
RpcPermanent(String),
#[error("RPC operation {operation} timed out after {timeout:?}")]
Timeout {
operation: &'static str,
timeout: Duration,
},
#[error("transaction decode error: {0}")]
Decode(String),
}
impl BackendError {
fn rpc(error: impl std::fmt::Display) -> Self {
let message = error.to_string();
if is_transient_rpc_error(&message) {
Self::RpcTransient(message)
} else {
Self::RpcPermanent(message)
}
}
pub const fn is_transient(&self) -> bool {
matches!(
self,
Self::RpcTransient(_) | Self::Timeout { .. } | Self::SimulationRetryable(_)
)
}
pub fn is_transport_destination_behind(&self) -> bool {
matches!(self, Self::Simulation(message) if message == "GlobalTableRootInFuture")
}
}