Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion machine-funding-server/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ license.workspace = true
rust-version.workspace = true

[dependencies]
alloy-consensus = { workspace = true, features = ["k256"] }
async-trait.workspace = true
axum.workspace = true
seismic-alloy-network.workspace = true
Expand Down Expand Up @@ -38,7 +39,6 @@ uuid.workspace = true
url.workspace = true

[dev-dependencies]
alloy-consensus = { workspace = true, features = ["k256"] }
http-body-util.workspace = true
tempfile = "3.21.0"
tower = { workspace = true, features = ["util"] }
213 changes: 198 additions & 15 deletions machine-funding-server/src/base_chain.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,8 @@ use crate::{
ServiceError,
},
};
use alloy_eips::eip2718::Encodable2718;
use alloy_consensus::{transaction::SignerRecoverable, Transaction, TxEnvelope};
use alloy_eips::eip2718::{Decodable2718, Encodable2718};
use alloy_network::{Ethereum, EthereumWallet, TransactionBuilder};
use alloy_primitives::{keccak256, Bytes, TxKind, B256, U256};
use alloy_provider::{Provider, RootProvider};
Expand All @@ -26,6 +27,8 @@ use tokio::time::{sleep, timeout};
use url::Url;

const ETH_TRANSFER_GAS_LIMIT: u64 = 21_000;
const MAX_ETH_TRANSFER_GAS_LIMIT: u64 = 100_000;
const GAS_HEADROOM_DIVISOR: u64 = 4;
const ERC20_TRANSFER_GAS_LIMIT: u64 = 120_000;
const FEE_HEADROOM_MULTIPLIER: u128 = 2;
const RECEIPT_POLL_INTERVAL: Duration = Duration::from_millis(500);
Expand Down Expand Up @@ -162,12 +165,7 @@ impl BaseChainDriver {
max_priority_fee_per_gas: u128,
) -> Result<TransactionRequest, ServiceError> {
let (to, value, data, gas) = match input.asset {
FundingAsset::BaseEth => (
input.recipient,
input.amount,
Bytes::new(),
ETH_TRANSFER_GAS_LIMIT,
),
FundingAsset::BaseEth => (input.recipient, input.amount, Bytes::new(), None),
FundingAsset::BaseErc20Usdc => (
self.config.token_address,
U256::ZERO,
Expand All @@ -177,7 +175,7 @@ impl BaseChainDriver {
}
.abi_encode()
.into(),
ERC20_TRANSFER_GAS_LIMIT,
Some(ERC20_TRANSFER_GAS_LIMIT),
),
FundingAsset::Susdc | FundingAsset::SusdcGas | FundingAsset::Erc20Usdc => {
return Err(deployment_identity_error());
Expand All @@ -188,7 +186,7 @@ impl BaseChainDriver {
to: Some(TxKind::Call(to)),
max_fee_per_gas: Some(max_fee_per_gas),
max_priority_fee_per_gas: Some(max_priority_fee_per_gas),
gas: Some(gas),
gas,
value: Some(value),
input: TransactionInput::from(data),
nonce: Some(nonce),
Expand All @@ -204,8 +202,13 @@ impl BaseChainDriver {
max_fee_per_gas: u128,
max_priority_fee_per_gas: u128,
) -> Result<PreparedTransaction, ServiceError> {
let request =
let mut request =
self.transaction_request(input, nonce, max_fee_per_gas, max_priority_fee_per_gas)?;
if input.asset == FundingAsset::BaseEth {
let estimate =
rpc_timeout(async { self.provider.estimate_gas(request.clone()).await }).await?;
request.gas = Some(eth_gas_limit(estimate)?);
}
let envelope = TransactionBuilder::<Ethereum>::build(request, &self.wallet)
.await
.map_err(chain_error)?;
Expand Down Expand Up @@ -257,6 +260,33 @@ impl ChainDriver for BaseChainDriver {
&self.operator_key
}

async fn can_retry_reverted(
&self,
input: &FundingInput,
transaction: &PreparedTransaction,
) -> Result<bool, ServiceError> {
self.validate_input(input)?;
if !legacy_eth_transfer(input, transaction, &self.config)? {
return Ok(false);
}
let hash = B256::from_str(&transaction.hash).map_err(chain_error)?;
let receipt = rpc_timeout(self.provider.get_transaction_receipt(hash)).await?;
let Some(receipt) = receipt else {
return Ok(false);
};
if receipt.transaction_hash != hash
|| receipt_status(Some(&receipt)) != Some(ChainResult::Reverted)
{
return Ok(false);
}
let latest = rpc_timeout(self.provider.get_block_number()).await?;
Ok(latest
>= receipt
.block_number
.unwrap()
.saturating_add(self.config.confirmations.saturating_sub(1)))
}

/// Only Base assets bound to this exact deployment are signed; a request
/// persisted under a rotated token or reserve is rejected before signing.
fn validate_input(&self, input: &FundingInput) -> Result<(), ServiceError> {
Expand Down Expand Up @@ -348,6 +378,37 @@ impl ChainDriver for BaseChainDriver {
}
}

fn eth_gas_limit(estimate: u64) -> Result<u64, ServiceError> {
let limit = estimate
.checked_add(estimate.div_ceil(GAS_HEADROOM_DIVISOR))
.filter(|limit| *limit <= MAX_ETH_TRANSFER_GAS_LIMIT && estimate >= ETH_TRANSFER_GAS_LIMIT)
.ok_or_else(|| {
chain_error("Base ETH transfer gas estimate exceeds the supported budget")
})?;
Ok(limit)
}

fn legacy_eth_transfer(
input: &FundingInput,
transaction: &PreparedTransaction,
config: &BaseConfig,
) -> Result<bool, ServiceError> {
if input.asset != FundingAsset::BaseEth {
return Ok(false);
}
let raw = hex::decode(transaction.serialized_transaction.trim_start_matches("0x"))
.map_err(chain_error)?;
let envelope = TxEnvelope::decode_2718(&mut raw.as_slice()).map_err(chain_error)?;
Ok(format!("{:#x}", keccak256(&raw)) == transaction.hash
&& envelope.gas_limit() == ETH_TRANSFER_GAS_LIMIT
&& envelope.chain_id() == Some(config.chain_id)
&& envelope.nonce() == transaction.nonce
&& envelope.to() == Some(input.recipient)
&& envelope.value() == input.amount
&& envelope.input().is_empty()
&& envelope.recover_signer().map_err(chain_error)? == config.reserve_address)
}

fn receipt_status(receipt: Option<&TransactionReceipt>) -> Option<ChainResult> {
let receipt = receipt?;
receipt.block_number?;
Expand All @@ -362,11 +423,132 @@ fn receipt_status(receipt: Option<&TransactionReceipt>) -> Option<ChainResult> {
#[cfg(test)]
mod tests {
use super::*;
use alloy_consensus::transaction::SignerRecoverable;
use alloy_consensus::{Transaction, TxEnvelope};
use alloy_eips::eip2718::Decodable2718;
use alloy_network::TxSigner;
use alloy_primitives::Address;
use axum::{extract::State, routing::post, Json, Router};
use serde_json::{json, Value};

async fn gas_rpc(Json(request): Json<Value>) -> Json<Value> {
assert_eq!(request["method"], "eth_estimateGas");
assert!(request["params"][0].get("gas").is_none());
Json(json!({"jsonrpc":"2.0", "id":request["id"], "result":"0x52e4"}))
}

async fn estimated_driver() -> (BaseChainDriver, tokio::task::JoinHandle<()>) {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let url = format!("http://{}", listener.local_addr().unwrap());
let task = tokio::spawn(async move {
axum::serve(listener, Router::new().route("/", post(gas_rpc)))
.await
.unwrap();
});
let mut driver = driver();
driver.provider = RootProvider::new_http(Url::parse(&url).unwrap());
(driver, task)
}

async fn recovery_rpc(State(receipt): State<Value>, Json(request): Json<Value>) -> Json<Value> {
let result = match request["method"].as_str().unwrap() {
"eth_getTransactionReceipt" => receipt,
"eth_blockNumber" => json!("0xc"),
other => panic!("unexpected RPC {other}"),
};
Json(json!({"jsonrpc":"2.0", "id":request["id"], "result":result}))
}

async fn signed_legacy(
driver: &BaseChainDriver,
input: &FundingInput,
gas: u64,
) -> PreparedTransaction {
let mut request = driver
.transaction_request(input, 7, 2_000_000_000, 1_000_000)
.unwrap();
request.gas = Some(gas);
let envelope = TransactionBuilder::<Ethereum>::build(request, &driver.wallet)
.await
.unwrap();
let raw = envelope.encoded_2718();
PreparedTransaction {
hash: format!("{:#x}", keccak256(&raw)),
nonce: 7,
serialized_transaction: format!("0x{}", hex::encode(raw)),
}
}

#[test]
fn native_gas_margin_is_bounded_and_never_truncated() {
assert_eq!(eth_gas_limit(21_000).unwrap(), 26_250);
assert_eq!(eth_gas_limit(21_220).unwrap(), 26_525);
assert_eq!(eth_gas_limit(80_000).unwrap(), MAX_ETH_TRANSFER_GAS_LIMIT);
for estimate in [0, 20_999, 80_001, u64::MAX] {
assert!(eth_gas_limit(estimate).is_err());
}
}

#[tokio::test]
async fn legacy_retry_requires_matching_signed_transfer_and_confirmed_revert() {
let mut driver = driver();
let input = input(FundingAsset::BaseEth, 10_000_000_000_000);
let old = signed_legacy(&driver, &input, ETH_TRANSFER_GAS_LIMIT).await;
assert!(legacy_eth_transfer(&input, &old, &driver.config).unwrap());
let new = signed_legacy(&driver, &input, 26_525).await;
assert!(!legacy_eth_transfer(&input, &new, &driver.config).unwrap());
let mut changed = input.clone();
changed.amount += U256::from(1);
assert!(!legacy_eth_transfer(&changed, &old, &driver.config).unwrap());
for (status, block, confirmations, expected) in [
(0, Some(12), 1, true),
(1, Some(12), 1, false),
(0, None, 1, false),
(0, Some(12), 2, false),
] {
let receipt = json!({
"blockHash": format!("0x{}", "1".repeat(64)),
"blockNumber": block.map(|n| format!("0x{n:x}")),
"transactionHash": old.hash,
"transactionIndex":"0x0", "type":"0x2", "contractAddress":null,
"cumulativeGasUsed":"0x5208", "gasUsed":"0x5208", "effectiveGasPrice":"0x1",
"logs":[], "logsBloom":format!("0x{}", "0".repeat(512)),
"status":format!("0x{status:x}"), "from":driver.config.reserve_address, "to":input.recipient
});
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
driver.provider = RootProvider::new_http(
Url::parse(&format!("http://{}", listener.local_addr().unwrap())).unwrap(),
);
driver.config.confirmations = confirmations;
let task = tokio::spawn(async move {
axum::serve(
listener,
Router::new()
.route("/", post(recovery_rpc))
.with_state(receipt),
)
.await
.unwrap();
});
assert_eq!(
driver.can_retry_reverted(&input, &old).await.unwrap(),
expected
);
task.abort();
}
}

#[tokio::test]
async fn delegated_wallet_drip_uses_estimated_gas_with_headroom() {
let (driver, task) = estimated_driver().await;
let input = input(FundingAsset::BaseEth, 10_000_000_000_000);
let prepared = driver
.sign_transaction(&input, 7, 2_000_000_000, 1_000_000)
.await
.unwrap();
let raw = hex::decode(prepared.serialized_transaction.trim_start_matches("0x")).unwrap();
let envelope = TxEnvelope::decode_2718(&mut raw.as_slice()).unwrap();
task.abort();
assert_eq!(envelope.gas_limit(), 26_525);
assert_eq!(envelope.value(), input.amount);
}

fn base_config(reserve: Address) -> BaseConfig {
BaseConfig {
Expand Down Expand Up @@ -414,7 +596,7 @@ mod tests {

#[tokio::test]
async fn eth_gas_drip_is_a_plain_value_transfer_signed_by_the_reserve() {
let driver = driver();
let (driver, task) = estimated_driver().await;
let gas = input(FundingAsset::BaseEth, 1_000_000_000_000_000);
let prepared = driver
.sign_transaction(&gas, 7, 2_000_000_000, 1_000_000)
Expand All @@ -429,7 +611,8 @@ mod tests {
assert_eq!(envelope.to(), Some(gas.recipient));
assert_eq!(envelope.value(), gas.amount);
assert!(envelope.input().is_empty());
assert_eq!(envelope.gas_limit(), ETH_TRANSFER_GAS_LIMIT);
assert_eq!(envelope.gas_limit(), 26_525);
task.abort();
assert_eq!(
envelope.recover_signer().unwrap(),
driver.config.reserve_address
Expand Down
7 changes: 7 additions & 0 deletions machine-funding-server/src/chain.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,13 @@ pub trait ChainDriver: Clone + Send + Sync + 'static {
Ok(())
}
async fn pending_nonce(&self) -> Result<u64, ServiceError>;
async fn can_retry_reverted(
&self,
_input: &FundingInput,
_transaction: &PreparedTransaction,
) -> Result<bool, ServiceError> {
Ok(false)
}
async fn prepare(
&self,
input: &FundingInput,
Expand Down
30 changes: 27 additions & 3 deletions machine-funding-server/src/service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -177,8 +177,10 @@ where
"idempotency_key was already used for a different request",
));
}
if let Some(response) = terminal_result(&existing, true)? {
return Ok(response);
if !base_eth_revert(&existing) {
if let Some(response) = terminal_result(&existing, true)? {
return Ok(response);
}
}

let lock_guard = self
Expand Down Expand Up @@ -297,7 +299,24 @@ where
"idempotency_key was already used for a different request",
));
}
if let Some(response) = terminal_result(&current, true)? {
let retry = if base_eth_revert(&current) {
let FundingRecord::Failed {
input, transaction, ..
} = &current
else {
unreachable!()
};
self.driver
.can_retry_reverted(&FundingInput::try_from(input)?, transaction)
.await?
} else {
false
};
if retry {
self.store
.retry_reverted(request_record_key, queue, &current)
.await?;
} else if let Some(response) = terminal_result(&current, true)? {
return Ok(response);
}

Expand Down Expand Up @@ -498,6 +517,11 @@ where
}
}

fn base_eth_revert(record: &FundingRecord) -> bool {
matches!(record, FundingRecord::Failed { input, code, .. }
if input.asset == FundingAsset::BaseEth && code == "transaction_reverted")
}

fn terminal_result(
record: &FundingRecord,
replayed: bool,
Expand Down
Loading
Loading