newton-tx-executor 0.7.3

Durable allowlisted transaction executor for Newton submissions
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};

/// Per-chain replacement fee policy.
#[derive(Debug, Clone, Copy)]
pub struct FeePolicy {
    /// Extra percentage applied to the market priority-fee quote.
    pub gas_bump_percent: u32,
    /// Maximum replacement or cancellation fee per gas, in wei.
    pub max_fee_per_gas_ceiling: u128,
}

/// Signed, deterministic EIP-2718 bytes.
#[derive(Debug, Clone)]
pub struct SignedAttempt {
    /// Explicit nonce encoded into the signed transaction.
    pub nonce: u64,
    /// Maximum gas units the transaction may consume.
    pub gas_limit: u64,
    /// EIP-1559 maximum fee per gas, in wei.
    pub max_fee_per_gas: u128,
    /// EIP-1559 priority fee per gas, in wei.
    pub max_priority_fee_per_gas: u128,
    /// Exact signed EIP-2718 bytes to persist before broadcast.
    pub raw_transaction: Bytes,
    /// Keccak-256 hash of `raw_transaction`.
    pub transaction_hash: B256,
}

/// One provider broadcast outcome.
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct BroadcastResult {
    /// Stable opaque identifier for the RPC endpoint that received the bytes.
    ///
    /// This must never contain credentials from the configured URL.
    pub provider: String,
    /// Accepted transaction hash or a bounded provider error.
    pub result: Result<B256, String>,
}

/// Receipt/finality observation.
#[derive(Debug, Clone)]
pub enum TransactionObservation {
    /// No receipt yet.
    Pending,
    /// At least one provider returned no receipt while another provider failed.
    ///
    /// This supports a safe same-nonce replacement, but cannot prove that a
    /// previously observed receipt was reorged away.
    PendingInconclusive {
        /// Providers that successfully returned no receipt.
        pending_providers: usize,
        /// Providers that could not provide receipt evidence.
        failed_providers: usize,
    },
    /// Successful mined receipt.
    Mined {
        /// Stable opaque identifier for the RPC provider that supplied this receipt.
        provider: String,
        /// Block containing the receipt.
        block_number: u64,
        /// Hash used to detect a pre-finality reorg.
        block_hash: B256,
        /// Confirmation depth at observation time.
        confirmations: u64,
        /// JSON-encoded receipt retained for recovery and audit.
        receipt: Vec<u8>,
    },
    /// Receipt consumed the nonce but reverted.
    Reverted {
        /// Stable opaque identifier for the RPC provider that supplied this receipt.
        provider: String,
        /// Block containing the reverted receipt.
        block_number: u64,
        /// Hash used to detect a pre-finality reorg.
        block_hash: B256,
        /// Confirmation depth at observation time.
        confirmations: u64,
        /// JSON-encoded receipt retained for recovery and audit.
        receipt: Vec<u8>,
    },
}

/// Amount of provider evidence requested for one receipt observation.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ObservationScope {
    /// Query one provider. A missing receipt means only "keep waiting" and must
    /// not be used to declare a reorg or make a replacement decision.
    Economical,
    /// Query all configured providers until a receipt is found, retaining
    /// partial pending evidence when some providers cannot answer. Use this
    /// before consequential state transitions.
    Exhaustive,
}

/// Verified on-chain projection for one requested item.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EffectStatus {
    /// Both stored hashes match the immutable request.
    Verified,
    /// Neither expected effect is present yet.
    Missing,
    /// The expected task exists but its response does not.
    TaskOnly,
    /// A non-zero stored hash differs from the immutable request.
    Conflict,
}

/// Per-item chain evidence.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct EffectObservation {
    /// Durable submission whose expected hashes were inspected.
    pub submission_id: SubmissionId,
    /// Relationship between stored and expected contract hashes.
    pub status: EffectStatus,
}

/// Authoritative classification for a transporter intent.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TransportEffectStatus {
    /// The immutable effect is present on the destination chain.
    Verified,
    /// The effect is not yet visible and may be retried idempotently.
    Missing,
    /// Destination state conflicts with the immutable intent.
    Conflict,
}

/// Narrow chain boundary used by signer workers and tests.
#[async_trait]
pub trait ChainBackend: Send + Sync {
    /// EVM sender owned by this signer lane.
    fn address(&self) -> Address;
    /// Required confirmation depth for this chain/provider policy.
    fn finality_confirmations(&self) -> u64 {
        1
    }
    /// Mined transaction count. This is the only source for a new task nonce.
    async fn latest_transaction_count(&self) -> Result<u64, BackendError>;
    /// Pending transaction count, used only to detect unexpected in-flight activity.
    async fn pending_transaction_count(&self) -> Result<u64, BackendError>;
    /// Native balance available to pay the maximum transaction cost.
    async fn balance(&self) -> Result<U256, BackendError>;
    /// Validates, encodes, estimates, signs, and returns bytes without broadcasting.
    async fn prepare(
        &self,
        intent: &ExecutableIntent,
        nonce: u64,
        previous_fees: Option<(u128, u128)>,
    ) -> Result<SignedAttempt, BackendError>;
    /// Prepares an allowlisted same-nonce self-send without broadcasting.
    async fn prepare_cancellation(&self, nonce: u64, fees: (u128, u128)) -> Result<SignedAttempt, BackendError>;
    /// Broadcasts exact persisted bytes.
    async fn broadcast(&self, raw_transaction: &Bytes) -> Vec<BroadcastResult>;
    /// Reads current receipt/finality state.
    async fn observe(
        &self,
        transaction_hash: B256,
        receipt_provider: Option<&str>,
        scope: ObservationScope,
    ) -> Result<TransactionObservation, BackendError>;
    /// Classifies each requested item from authoritative task-manager hashes.
    async fn classify_onchain_effects(&self, intent: &ExecutableIntent)
        -> Result<Vec<EffectObservation>, BackendError>;
    /// Classifies one transporter effect from destination state.
    async fn classify_transport_effect(
        &self,
        _intent: &ExecutableIntent,
    ) -> Result<TransportEffectStatus, BackendError> {
        Err(BackendError::Policy(
            "backend does not support transporter effect classification".to_string(),
        ))
    }
}

/// Chain preparation or observation failure.
#[derive(Debug, thiserror::Error)]
pub enum BackendError {
    /// Invalid configuration.
    #[error("executor configuration error: {0}")]
    Configuration(String),
    /// Executor policy denied an intent.
    #[error("executor policy violation: {0}")]
    Policy(String),
    /// Read-only preflight failed.
    #[error("transaction simulation failed: {0}")]
    Simulation(String),
    /// Per-item batch simulation failed for a context-dependent gas reason.
    #[error("retryable transaction simulation failure: {0}")]
    SimulationRetryable(String),
    /// Signing failed.
    #[error("transaction signing failed: {0}")]
    Signing(String),
    /// Retryable RPC or transport failure.
    #[error("transient RPC error: {0}")]
    RpcTransient(String),
    /// Deterministic RPC protocol/configuration failure.
    #[error("permanent RPC error: {0}")]
    RpcPermanent(String),
    /// One provider exceeded the configured request deadline.
    #[error("RPC operation {operation} timed out after {timeout:?}")]
    Timeout {
        /// Bounded operation name.
        operation: &'static str,
        /// Applied request deadline.
        timeout: Duration,
    },
    /// Serialization failure.
    #[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)
        }
    }

    /// Whether retrying after releasing an unsigned reservation is safe.
    pub const fn is_transient(&self) -> bool {
        matches!(
            self,
            Self::RpcTransient(_) | Self::Timeout { .. } | Self::SimulationRetryable(_)
        )
    }

    /// Whether transporter preflight observed a destination head behind the
    /// finalized source timestamp. The planner normally prevents this revert;
    /// retaining the classification closes the provider/head race between the
    /// read-only plan and executor simulation.
    pub fn is_transport_destination_behind(&self) -> bool {
        matches!(self, Self::Simulation(message) if message == "GlobalTableRootInFuture")
    }
}