diff --git a/Package.swift b/Package.swift index 794f6ec..0555e2c 100644 --- a/Package.swift +++ b/Package.swift @@ -5,7 +5,7 @@ import PackageDescription import Foundation let tag = "v0.6.0" -let checksum = "ac07f72265cde752e9f7951383e8d1a01f3568676d156ee592a715bc5a38d32a" +let checksum = "7cb8c8c49221d991f7cbe71c73dcad8e7ceb4d8ad410d68d7a781eb70e04fbdf" let url = "https://github.com/synonymdev/bitkit-core/releases/download/\(tag)/BitkitCore.xcframework.zip" let localBinary = ProcessInfo.processInfo.environment["BITKIT_CORE_LOCAL"] == "1" diff --git a/bindings/ios/bitkitcore.swift b/bindings/ios/bitkitcore.swift index 5b7f48d..65c4e4f 100644 --- a/bindings/ios/bitkitcore.swift +++ b/bindings/ios/bitkitcore.swift @@ -2576,7 +2576,7 @@ public protocol UsdtWalletProtocol: AnyObject, Sendable { func history() throws -> [UsdtTransfer] - func quoteTransfer(recipient: String, amount: UInt64) async throws -> UsdtQuote + func quoteTransfer(recipient: String, amount: UInt64, destination: UsdtDestination) async throws -> UsdtQuote func receiveAddress() -> String @@ -2712,13 +2712,13 @@ open func history()throws -> [UsdtTransfer] { }) } -open func quoteTransfer(recipient: String, amount: UInt64)async throws -> UsdtQuote { +open func quoteTransfer(recipient: String, amount: UInt64, destination: UsdtDestination)async throws -> UsdtQuote { return try await uniffiRustCallAsync( rustFutureFunc: { uniffi_bitkitcore_fn_method_usdtwallet_quote_transfer( self.uniffiClonePointer(), - FfiConverterString.lower(recipient),FfiConverterUInt64.lower(amount) + FfiConverterString.lower(recipient),FfiConverterUInt64.lower(amount),FfiConverterTypeUsdtDestination_lower(destination) ) }, pollFunc: ffi_bitkitcore_rust_future_poll_rust_buffer, @@ -16046,16 +16046,20 @@ public func FfiConverterTypeUsdtPaymentRequest_lower(_ value: UsdtPaymentRequest public struct UsdtQuote { public var id: String public var recipient: String + public var destination: UsdtDestination public var amount: UInt64 + public var receivedAmount: UInt64 public var maximumFee: UInt64 public var expiresAt: UInt64 // Default memberwise initializers are never public by default, so we // declare one manually. - public init(id: String, recipient: String, amount: UInt64, maximumFee: UInt64, expiresAt: UInt64) { + public init(id: String, recipient: String, destination: UsdtDestination, amount: UInt64, receivedAmount: UInt64, maximumFee: UInt64, expiresAt: UInt64) { self.id = id self.recipient = recipient + self.destination = destination self.amount = amount + self.receivedAmount = receivedAmount self.maximumFee = maximumFee self.expiresAt = expiresAt } @@ -16074,9 +16078,15 @@ extension UsdtQuote: Equatable, Hashable { if lhs.recipient != rhs.recipient { return false } + if lhs.destination != rhs.destination { + return false + } if lhs.amount != rhs.amount { return false } + if lhs.receivedAmount != rhs.receivedAmount { + return false + } if lhs.maximumFee != rhs.maximumFee { return false } @@ -16089,7 +16099,9 @@ extension UsdtQuote: Equatable, Hashable { public func hash(into hasher: inout Hasher) { hasher.combine(id) hasher.combine(recipient) + hasher.combine(destination) hasher.combine(amount) + hasher.combine(receivedAmount) hasher.combine(maximumFee) hasher.combine(expiresAt) } @@ -16108,7 +16120,9 @@ public struct FfiConverterTypeUsdtQuote: FfiConverterRustBuffer { try UsdtQuote( id: FfiConverterString.read(from: &buf), recipient: FfiConverterString.read(from: &buf), + destination: FfiConverterTypeUsdtDestination.read(from: &buf), amount: FfiConverterUInt64.read(from: &buf), + receivedAmount: FfiConverterUInt64.read(from: &buf), maximumFee: FfiConverterUInt64.read(from: &buf), expiresAt: FfiConverterUInt64.read(from: &buf) ) @@ -16117,7 +16131,9 @@ public struct FfiConverterTypeUsdtQuote: FfiConverterRustBuffer { public static func write(_ value: UsdtQuote, into buf: inout [UInt8]) { FfiConverterString.write(value.id, into: &buf) FfiConverterString.write(value.recipient, into: &buf) + FfiConverterTypeUsdtDestination.write(value.destination, into: &buf) FfiConverterUInt64.write(value.amount, into: &buf) + FfiConverterUInt64.write(value.receivedAmount, into: &buf) FfiConverterUInt64.write(value.maximumFee, into: &buf) FfiConverterUInt64.write(value.expiresAt, into: &buf) } @@ -16146,7 +16162,9 @@ public struct UsdtTransfer { */ public var txHash: String? public var userOperationHash: String? + public var bridgeGuid: String? public var recipient: String + public var destination: UsdtDestination public var amount: UInt64 public var receivedAmount: UInt64 public var fee: UInt64? @@ -16159,11 +16177,13 @@ public struct UsdtTransfer { public init(id: String, /** * Source transaction hash, absent until execution is observed. - */txHash: String?, userOperationHash: String?, recipient: String, amount: UInt64, receivedAmount: UInt64, fee: UInt64?, isIncoming: Bool, status: UsdtTransferStatus, timestamp: UInt64) { + */txHash: String?, userOperationHash: String?, bridgeGuid: String?, recipient: String, destination: UsdtDestination, amount: UInt64, receivedAmount: UInt64, fee: UInt64?, isIncoming: Bool, status: UsdtTransferStatus, timestamp: UInt64) { self.id = id self.txHash = txHash self.userOperationHash = userOperationHash + self.bridgeGuid = bridgeGuid self.recipient = recipient + self.destination = destination self.amount = amount self.receivedAmount = receivedAmount self.fee = fee @@ -16189,9 +16209,15 @@ extension UsdtTransfer: Equatable, Hashable { if lhs.userOperationHash != rhs.userOperationHash { return false } + if lhs.bridgeGuid != rhs.bridgeGuid { + return false + } if lhs.recipient != rhs.recipient { return false } + if lhs.destination != rhs.destination { + return false + } if lhs.amount != rhs.amount { return false } @@ -16217,7 +16243,9 @@ extension UsdtTransfer: Equatable, Hashable { hasher.combine(id) hasher.combine(txHash) hasher.combine(userOperationHash) + hasher.combine(bridgeGuid) hasher.combine(recipient) + hasher.combine(destination) hasher.combine(amount) hasher.combine(receivedAmount) hasher.combine(fee) @@ -16241,7 +16269,9 @@ public struct FfiConverterTypeUsdtTransfer: FfiConverterRustBuffer { id: FfiConverterString.read(from: &buf), txHash: FfiConverterOptionString.read(from: &buf), userOperationHash: FfiConverterOptionString.read(from: &buf), + bridgeGuid: FfiConverterOptionString.read(from: &buf), recipient: FfiConverterString.read(from: &buf), + destination: FfiConverterTypeUsdtDestination.read(from: &buf), amount: FfiConverterUInt64.read(from: &buf), receivedAmount: FfiConverterUInt64.read(from: &buf), fee: FfiConverterOptionUInt64.read(from: &buf), @@ -16255,7 +16285,9 @@ public struct FfiConverterTypeUsdtTransfer: FfiConverterRustBuffer { FfiConverterString.write(value.id, into: &buf) FfiConverterOptionString.write(value.txHash, into: &buf) FfiConverterOptionString.write(value.userOperationHash, into: &buf) + FfiConverterOptionString.write(value.bridgeGuid, into: &buf) FfiConverterString.write(value.recipient, into: &buf) + FfiConverterTypeUsdtDestination.write(value.destination, into: &buf) FfiConverterUInt64.write(value.amount, into: &buf) FfiConverterUInt64.write(value.receivedAmount, into: &buf) FfiConverterOptionUInt64.write(value.fee, into: &buf) @@ -23405,6 +23437,99 @@ extension UrPayload: Codable {} +// Note that we don't yet support `indirect` for enums. +// See https://github.com/mozilla/uniffi-rs/issues/396 for further discussion. + +public enum UsdtDestination { + + case stable + case ethereum + case arbitrum + case polygon + case plasma +} + + +#if compiler(>=6) +extension UsdtDestination: Sendable {} +#endif + +#if swift(>=5.8) +@_documentation(visibility: private) +#endif +public struct FfiConverterTypeUsdtDestination: FfiConverterRustBuffer { + typealias SwiftType = UsdtDestination + + public static func read(from buf: inout (data: Data, offset: Data.Index)) throws -> UsdtDestination { + let variant: Int32 = try readInt(&buf) + switch variant { + + case 1: return .stable + + case 2: return .ethereum + + case 3: return .arbitrum + + case 4: return .polygon + + case 5: return .plasma + + default: throw UniffiInternalError.unexpectedEnumCase + } + } + + public static func write(_ value: UsdtDestination, into buf: inout [UInt8]) { + switch value { + + + case .stable: + writeInt(&buf, Int32(1)) + + + case .ethereum: + writeInt(&buf, Int32(2)) + + + case .arbitrum: + writeInt(&buf, Int32(3)) + + + case .polygon: + writeInt(&buf, Int32(4)) + + + case .plasma: + writeInt(&buf, Int32(5)) + + } + } +} + + +#if swift(>=5.8) +@_documentation(visibility: private) +#endif +public func FfiConverterTypeUsdtDestination_lift(_ buf: RustBuffer) throws -> UsdtDestination { + return try FfiConverterTypeUsdtDestination.lift(buf) +} + +#if swift(>=5.8) +@_documentation(visibility: private) +#endif +public func FfiConverterTypeUsdtDestination_lower(_ value: UsdtDestination) -> RustBuffer { + return FfiConverterTypeUsdtDestination.lower(value) +} + + +extension UsdtDestination: Equatable, Hashable {} + +extension UsdtDestination: Codable {} + + + + + + public enum UsdtError: Swift.Error { @@ -23594,6 +23719,18 @@ public enum UsdtTransferStatus { * Source payment failed or was proven not to have executed. */ case failed + /** + * Source payment executed; destination delivery is pending. + */ + case bridging + /** + * Delivery is blocked or its message could not be recovered; it may still complete. + */ + case bridgeNeedsAttention + /** + * Delivery was permanently stopped. This does not imply a refund of source funds or fees. + */ + case bridgeFailed /** * Another operation consumed the payment nonce. */ @@ -23621,7 +23758,13 @@ public struct FfiConverterTypeUsdtTransferStatus: FfiConverterRustBuffer { case 3: return .failed - case 4: return .replaced + case 4: return .bridging + + case 5: return .bridgeNeedsAttention + + case 6: return .bridgeFailed + + case 7: return .replaced default: throw UniffiInternalError.unexpectedEnumCase } @@ -23643,9 +23786,21 @@ public struct FfiConverterTypeUsdtTransferStatus: FfiConverterRustBuffer { writeInt(&buf, Int32(3)) - case .replaced: + case .bridging: writeInt(&buf, Int32(4)) + + case .bridgeNeedsAttention: + writeInt(&buf, Int32(5)) + + + case .bridgeFailed: + writeInt(&buf, Int32(6)) + + + case .replaced: + writeInt(&buf, Int32(7)) + } } } @@ -29768,7 +29923,7 @@ private let initializationResult: InitializationResult = { if (uniffi_bitkitcore_checksum_method_usdtwallet_history() != 4617) { return InitializationResult.apiChecksumMismatch } - if (uniffi_bitkitcore_checksum_method_usdtwallet_quote_transfer() != 56645) { + if (uniffi_bitkitcore_checksum_method_usdtwallet_quote_transfer() != 3732) { return InitializationResult.apiChecksumMismatch } if (uniffi_bitkitcore_checksum_method_usdtwallet_receive_address() != 540) { diff --git a/bindings/ios/bitkitcoreFFI.h b/bindings/ios/bitkitcoreFFI.h index 63b655a..7ebda6a 100644 --- a/bindings/ios/bitkitcoreFFI.h +++ b/bindings/ios/bitkitcoreFFI.h @@ -692,7 +692,7 @@ RustBuffer uniffi_bitkitcore_fn_method_usdtwallet_history(void*_Nonnull ptr, Rus #endif #ifndef UNIFFI_FFIDEF_UNIFFI_BITKITCORE_FN_METHOD_USDTWALLET_QUOTE_TRANSFER #define UNIFFI_FFIDEF_UNIFFI_BITKITCORE_FN_METHOD_USDTWALLET_QUOTE_TRANSFER -uint64_t uniffi_bitkitcore_fn_method_usdtwallet_quote_transfer(void*_Nonnull ptr, RustBuffer recipient, uint64_t amount +uint64_t uniffi_bitkitcore_fn_method_usdtwallet_quote_transfer(void*_Nonnull ptr, RustBuffer recipient, uint64_t amount, RustBuffer destination ); #endif #ifndef UNIFFI_FFIDEF_UNIFFI_BITKITCORE_FN_METHOD_USDTWALLET_RECEIVE_ADDRESS diff --git a/src/lib.rs b/src/lib.rs index 035fb1d..985a20e 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -91,8 +91,9 @@ pub use modules::onchain; pub use modules::scanner::{DecodingError, LnurlPayData, Scanner}; pub use modules::seedqr::{decode_compact_seed_qr, decode_standard_seed_qr, SeedQrError}; pub use modules::usdt::{ - usdt_address, usdt_format_amount, usdt_parse_amount, usdt_parse_payment_request, UsdtError, - UsdtPaymentRequest, UsdtQuote, UsdtTransfer, UsdtTransferStatus, UsdtWallet, + usdt_address, usdt_format_amount, usdt_parse_amount, usdt_parse_payment_request, + UsdtDestination, UsdtError, UsdtPaymentRequest, UsdtQuote, UsdtTransfer, UsdtTransferStatus, + UsdtWallet, }; use bip39::Mnemonic; diff --git a/src/modules/usdt/README.md b/src/modules/usdt/README.md index d52a9c5..ca18ae2 100644 --- a/src/modules/usdt/README.md +++ b/src/modules/usdt/README.md @@ -14,7 +14,7 @@ Owned mnemonic/passphrase/seed buffers are zeroized and signing keys are erased ## Quotes and fees -`quote_transfer` takes a raw recipient and a positive atomic amount; it never receives signing credentials. Local quotes last at most 120 seconds and newly quoted paymaster terms must expire within 15 minutes. `send` validates the owner, nonces, balance, gas estimates, the current slow gas-price recommendation and deadlines before signing the stored plan. Quotes use the fast gas-price recommendation; a modest price increase does not invalidate a quote that still covers the current slow recommendation. Changes beyond the approved bounds require a new review; signing cannot raise the approved fee. +`quote_transfer` takes a raw recipient, positive atomic amount and destination; it never receives signing credentials. Local quotes last at most 120 seconds and newly quoted paymaster terms must expire within 15 minutes. `send` validates the owner, nonces, balance, gas estimates, the current slow gas-price recommendation and deadlines before signing the stored plan. Quotes use the fast gas-price recommendation; a modest price increase does not invalidate a quote that still covers the current slow recommendation. Changes beyond the approved bounds require a new review; signing cannot raise the approved fee. The pinned ERC-20 paymaster collects USDT. Its finite approval includes a 5% margin; the displayed maximum fee comes from signed gas limits and paymaster terms, not the allowance. Call/pre-verification estimates receive 10% execution/L1-data headroom; the charged pre-verification margin is included in the maximum. A residual paymaster allowance can remain and is reset to a finite amount on the next payment. @@ -30,7 +30,7 @@ A matching event in a canonical receipt settles the payment. Discovery logs alon Seed restoration recovers deposits and outgoing activity from genesis, including transfers before delegation and sends through another wallet. Supported direct EntryPoint calls and paymaster modes recover payment/fee attribution; unknown wrappers or payment modes preserve raw token transfers instead of guessing their intent. Failed payments retain attempted amounts but have no delivered amount. Transaction hashes are absent until execution is observed; callers derive explorer links from the source transaction hash rather than storing a second copy of it. -`sync_history` returns `true` when caught up and `false` when more work remains. It uses adaptive log ranges and a 20-second soft budget between persisted receipts; an in-flight receipt may finish later. Learned range limits survive budget exits. A single-block log overflow falls back to that block's individual receipts. Complete receipt enrichment and fallback scans are retained by canonical block hash within the revisit window; log-only incoming observations remain eligible for later enrichment. Zero/self transfers are discarded before enrichment. Network failures preserve completed work and never silently skip a block. Callers should defer catch-up during payment review/submission so historical work does not delay sends. +`sync_history` returns `true` when caught up and `false` when more work remains. It uses adaptive log ranges and a 20-second soft budget between persisted receipts; an in-flight receipt may finish later. Learned range limits survive budget exits. A single-block log overflow falls back to that block's individual receipts. Complete receipt enrichment and fallback scans are retained by canonical block hash within the revisit window; log-only incoming observations and successful bridge receipts missing tracking or fee evidence remain eligible for later enrichment. Incomplete bridge metadata does not keep an executed source payment pending. Zero/self transfers are discarded before enrichment. Network failures preserve completed work and never silently skip a block. Callers should defer catch-up during payment review/submission so historical work does not delay sends. `check_recent_execution` checks one recent direct Arbitrum payment with a five-second request budget. It requires the expected operation outcome and token transfer in a matching canonical receipt and does not scan history, rebroadcast, expire payments or reconcile nonces. It can confirm execution at the current L2 tip; this is provisional sequencer execution, not parent-chain finality. Native send screens may call it approximately once per second during a short foreground window, with cancellation and rate-limit backoff between checks. Missing evidence leaves Pending intact. Normal recovery handles older payments outside its 64-block lookup window. @@ -40,17 +40,31 @@ Payment outcomes and expiry decisions trust the configured chain RPC. A maliciou Storage is wallet-specific and must have one owning `UsdtWallet` object. Drop it before deleting its database during an explicit wallet wipe. Async exports use UniFFI's Tokio adapter; Kotlin cancellation can drop the polled future, whereas the current Swift bindings may finish an in-flight call after task cancellation. Callers must check cancellation between calls. Sends are not detached onto the global runtime used by stateless exports. -## Transport +## Transport and cross-network APIs Both chain and bundler endpoints must be controlled, credential-free HTTPS URLs; HTTP is accepted only on loopback for fixtures. Provider keys belong on the server. Chain/bundler calls share an 80/minute budget with a burst of 20. Responses are bounded to 2 MiB, except protocol-projected receipts up to 16 MiB. The companion service documents provider requirements, receipt projection and deployment limits. +The outbound bridge API supports Ethereum (30101), Polygon (30109), Plasma (30383) and Stable (30396), alongside direct Arbitrum transfers. Native release flows expose Arbitrum only; bridge routes require explicit service enablement and destination acceptance. Plain deposits on another chain are not automatically forwarded. Recipient validation rejects the destination token and the pinned EntryPoint, paymaster and Simple7702 delegate addresses on every destination. Direct Arbitrum sends also reject the source OFT and helper. + +Bridge quotes include 10% native messaging-fee headroom and 20% token-conversion headroom, both within the displayed maximum USDT fee. Before signing or rebroadcasting, the stored native fee, helper liquidity and token approval are checked against current requirements without raising approved limits. The service reports OFT/helper execution reverts as sanitized RPC code `3`, which core maps to `UnsupportedRoute` during initial quoting. A reverted fee recheck for an already reviewed bridge quote requires a fresh quote (`QuoteExpired`); insufficient helper liquidity remains `UnsupportedRoute`. Provider outages remain retryable network errors. Delivery checks process up to three transfers concurrently outside the send lock, with a ten-second request budget, even when source recovery fails; failed lookups retain the last known status, while an explicit `INFLIGHT` or `CONFIRMING` update clears a previous needs-attention state. + +Bridges use the pinned OFT and TransactionValueHelper with zero account ETH, a finite USDT approval covering principal/fee, and atomic helper-allowance revocation. The deployed helper requires native liquidity and retains behaviors noted in its OpenZeppelin audit; its verified runtime is not the audit-remediated implementation. Source success means bridging, not delivered. + +`Pending` means source execution is unresolved; `Failed` means the source payment failed or was proved unexecuted; `Replaced` means another operation consumed its nonce. `Bridging` means source execution succeeded and destination delivery is unresolved. For bridges, `Confirmed` means delivery was reported. `BridgeNeedsAttention` covers retryable delivery problems or missing message evidence; without a GUID, no delivery lookup is possible. `BridgeFailed` means LayerZero reports a burned or skipped message: delivery polling stops, while source transaction, GUID, amount and fees remain visible. Neither bridge status implies a refund. Terminal delivery states survive restart and source-history rescans for the same transaction and GUID. + +LayerZero status must match the operation GUID/pathway before confirmation; blocked delivery remains visible and never triggers an automatic paid retry. + +RPC providers see queried addresses. Delivery checks use `bitkit_getBridgeMessages([sourceTransactionHash])` on the existing chain-service endpoint. The service queries LayerZero Scan without forwarding device headers, projects only message identity/pathway/status fields, and applies its shared request and response limits. LayerZero sees the service IP and the transaction hash; the service still sees the requesting device. Manually opening LayerZero Scan from transaction details connects the browser directly. No delivery requests are made for Arbitrum-only transfers. + +Unavailable, unmatched or unknown delivery responses preserve the last status. Those lookups and retryable problems (`FAILED`, `BLOCKED`, `PAYLOAD_STORED`) wait at least one minute before another automatic attempt in the same wallet session. `APPLICATION_BURNED` and `APPLICATION_SKIPPED` stop polling as `BridgeFailed`; `DELIVERED` stops polling as `Confirmed`. No status triggers an automatic paid retry or refund. + ## Validation and bindings Run `cargo test --locked --lib modules::usdt`; CI runs these deterministic tests. They cover independent signing/address vectors, fee bounds, uncertain submission, nonce recovery and restored history. Fixtures use public test credentials. -For the ignored deployed-contract test, start a fresh Arbitrum Anvil fork on port 18545 and `tests/usdt-fork/provider.mjs` on 18546 after installing its pinned dependencies. Run `cargo test deployed_contracts_collect_usdt_fees_without_account_eth -- --ignored`. The fixture requires Anvil, sets local balances/signing terms and executes deployed contracts; it does not establish real provider pricing. +For the ignored deployed-contract test, start a fresh Arbitrum Anvil fork on port 18545 and `tests/usdt-fork/provider.mjs` on 18546 after installing its pinned dependencies. Run `cargo test deployed_contracts_collect_usdt_fees_and_revert_failed_bridges_atomically -- --ignored`. The fixture requires Anvil, sets local balances/signing terms and executes deployed contracts; it does not establish real provider pricing or destination delivery. -To include the service, start it with `NODE_ENV=test ARBITRUM_RPC_URL=http://127.0.0.1:18545 LOCAL_PROVIDER_URL=http://127.0.0.1:18546`, then pass `USDT_FORK_RPC_URL=http://127.0.0.1:3100/v1/usdt/chain-rpc` and `USDT_FORK_BUNDLER_URL=http://127.0.0.1:3100/v1/usdt/rpc` to the ignored test. +To include the service, start it with `USDT_BRIDGE_NETWORKS=ethereum,polygon,plasma,stable NODE_ENV=test ARBITRUM_RPC_URL=http://127.0.0.1:18545 LOCAL_PROVIDER_URL=http://127.0.0.1:18546`, then pass `USDT_FORK_RPC_URL=http://127.0.0.1:3100/v1/usdt/chain-rpc` and `USDT_FORK_BUNDLER_URL=http://127.0.0.1:3100/v1/usdt/rpc` to the ignored test. Build iOS and Android sequentially with the repository scripts; Android temporarily edits the manifest/example. Generated bindings and native artifacts must use the same source. App configuration and local package overrides belong in each native repository's USDT documentation. @@ -60,5 +74,12 @@ Build iOS and Android sequentially with the repository scripts; Android temporar - [Alto EIP-7702 request validation](https://github.com/pimlicolabs/alto/blob/96529592b67a69be23c013359cbc9990657af64a/src/rpc/rpcHandler.ts) - [Simple7702Account](https://github.com/eth-infinitism/account-abstraction/blob/releases/v0.8/contracts/accounts/Simple7702Account.sol) - [Pimlico supported tokens](https://docs.pimlico.io/references/paymaster/erc20-paymaster/supported-tokens) +- [Pimlico paymaster deployments](https://docs.pimlico.io/references/paymaster/erc20-paymaster/contract-addresses) - [Pimlico pricing](https://www.pimlico.io/pricing) - [Pimlico public endpoint limits](https://docs.pimlico.io/references/bundler/public-endpoint) +- [USDT0 documentation](https://docs.usdt0.to/) +- [Transaction helper audit](https://www.openzeppelin.com/news/usdt0-transaction-helper-audit) +- [Verified deployed helper](https://arbitrum.blockscout.com/api/v2/smart-contracts/0xa90f03c856d01f698e7071b393387cd75a8a319a) +- [LayerZero message statuses](https://docs.layerzero.network/v2/tools/layerzeroscan/mainnet/messages/get-messagesstatus) + +Destination token addresses follow the official USDT0 ecosystem listings for [Polygon](https://usdt0.to/ecosystem/polygon), [Plasma](https://usdt0.to/ecosystem/plasma) and [Stable](https://usdt0.to/ecosystem/stable). diff --git a/src/modules/usdt/history.rs b/src/modules/usdt/history.rs index c2da32f..464ac11 100644 --- a/src/modules/usdt/history.rs +++ b/src/modules/usdt/history.rs @@ -1,9 +1,9 @@ use super::{ account::{SimpleAccount, ENTRY_POINT}, amount::token_amount, - transaction::{entry_point_event, event_data, EntryPoint, Erc20}, - types::TOKEN, - UsdtError, UsdtTransfer, UsdtTransferStatus, UsdtWallet, + transaction::{entry_point_event, event_data, BridgeHelper, EntryPoint, Erc20}, + types::{BRIDGE_HELPER, OFT, TOKEN}, + UsdtDestination, UsdtError, UsdtTransfer, UsdtTransferStatus, UsdtWallet, }; use alloy_primitives::{Address, Bytes, B256, U256}; use alloy_sol_types::{SolCall, SolEvent}; @@ -164,6 +164,7 @@ impl UsdtWallet { if self.store.history_progress()? != Some(number) { self.store.save_history_progress(number)?; } + let mut complete = true; for hash in &block.transactions { let id = format!("{hash:#x}"); if self.store.has_history_receipt(&id, &block_hash, true)? { @@ -173,13 +174,16 @@ impl UsdtWallet { return Ok(false); } let receipt = self.rpc.block_receipt(*hash, block.hash, number).await?; - self.save_receipt_history(&id, number, &block, &receipt, true) + complete &= self + .save_receipt_history(&id, number, &block, &receipt, true) .await?; } if self.rpc.block(number).await?.hash != block.hash { return Err(UsdtError::NetworkUnavailable); } - self.store.complete_history_block(number, &block_hash)?; + if complete { + self.store.complete_history_block(number, &block_hash)?; + } Ok(true) } @@ -190,17 +194,25 @@ impl UsdtWallet { block: &super::rpc::Block, receipt: &Value, complete_receipt: bool, - ) -> Result<(), UsdtError> { + ) -> Result { // Finish and persist a receipt before yielding the work budget. let timestamp = block.timestamp()?; let transfers = self.receipt_history(hash, timestamp, receipt).await?; + // Missing bridge evidence remains eligible for enrichment without keeping the source pending. + let complete_receipt = complete_receipt + && transfers.iter().all(|transfer| { + transfer.destination == UsdtDestination::Arbitrum + || transfer.status == UsdtTransferStatus::Failed + || (transfer.bridge_guid.is_some() && transfer.fee.is_some()) + }); self.store.save_history_receipt( &transfers, hash, number, &format!("{:#x}", block.hash), complete_receipt, - ) + )?; + Ok(complete_receipt) } async fn history_logs( @@ -278,7 +290,9 @@ impl UsdtWallet { id: format!("{hash}:{index}"), tx_hash: Some(hash.into()), user_operation_hash: None, + bridge_guid: None, recipient: event.to.to_checksum(None), + destination: UsdtDestination::Arbitrum, amount: token_amount(event.value)?, received_amount: token_amount(event.value)?, fee: None, @@ -322,8 +336,8 @@ impl UsdtWallet { continue; } let operation_hash = format!("{:#x}", event.userOpHash); - let (recipient, amount) = if let Some(saved) = saved { - (saved.recipient, saved.amount) + let (recipient, amount, destination) = if let Some(saved) = saved { + (saved.recipient, saved.amount, saved.destination) } else { let Some(op) = batch.as_ref().and_then(|batch| { batch @@ -336,16 +350,20 @@ impl UsdtWallet { if !super::paymaster::supported_payment(&op.paymasterAndData) { continue; } - let Some((recipient, amount)) = decode_payment(&op.callData, self.address) else { + let Some((recipient, amount, destination)) = + decode_payment(&op.callData, self.address) + else { continue; }; - (recipient.to_checksum(None), amount) + (recipient.to_checksum(None), amount, destination) }; let mut transfer = UsdtTransfer { id: operation_hash.clone(), tx_hash: Some(hash.into()), user_operation_hash: Some(operation_hash), + bridge_guid: None, recipient, + destination, amount, received_amount: amount, fee: None, @@ -369,19 +387,38 @@ impl UsdtWallet { } } -fn decode_payment(data: &[u8], sender: Address) -> Option<(Address, u64)> { - let calls = decode_calls(data).ok()?; +fn decode_payment(data: &[u8], sender: Address) -> Option<(Address, u64, UsdtDestination)> { let mut payment = None; - for (target, data) in calls { - if target != TOKEN { - return None; - } - if let Ok(call) = Erc20::transferCall::abi_decode(&data) { - if payment.is_some() || call.amount.is_zero() || call.recipient == sender { + for (target, data) in decode_calls(data).ok()? { + let next = if target == TOKEN { + if let Ok(call) = Erc20::transferCall::abi_decode(&data) { + if call.amount.is_zero() || call.recipient == sender { + return None; + } + ( + call.recipient, + token_amount(call.amount).ok()?, + UsdtDestination::Arbitrum, + ) + } else if Erc20::approveCall::abi_decode(&data).is_ok() { + continue; + } else { + return None; + } + } else if target == BRIDGE_HELPER { + let call = BridgeHelper::sendCall::abi_decode(&data).ok()?; + if call.oft != OFT { return None; } - payment = Some((call.recipient, token_amount(call.amount).ok()?)); - } else if Erc20::approveCall::abi_decode(&data).is_err() { + ( + Address::from_word(call.param.to), + token_amount(call.param.amountLD).ok()?, + UsdtDestination::from_endpoint(call.param.dstEid)?, + ) + } else { + return None; + }; + if payment.replace(next).is_some() { return None; } } diff --git a/src/modules/usdt/rpc.rs b/src/modules/usdt/rpc.rs index 1762cc1..d11ca3e 100644 --- a/src/modules/usdt/rpc.rs +++ b/src/modules/usdt/rpc.rs @@ -1,4 +1,4 @@ -use super::UsdtError; +use super::{UsdtError, UsdtTransfer, UsdtTransferStatus}; use alloy_primitives::{Address, Bytes, B256, U256}; use alloy_sol_types::SolCall; use serde::{de::DeserializeOwned, Deserialize}; @@ -139,17 +139,26 @@ impl Rpc { if matches!(error.code, -32002 | -32603) { return Err(UsdtError::NetworkUnavailable); } + if method == "eth_call" + && error.code == 3 + && serde_json::from_value::
(params[0]["to"].clone()) + .is_ok_and(|to| [super::types::OFT, super::types::BRIDGE_HELPER].contains(&to)) + { + return Err(UsdtError::UnsupportedRoute); + } if matches!( method, "eth_chainId" | "eth_blockNumber" | "eth_getCode" + | "eth_getBalance" | "eth_getTransactionCount" | "eth_call" | "eth_getLogs" | "eth_getBlockByNumber" | "eth_getTransactionReceipt" | "eth_getTransactionByHash" + | "bitkit_getBridgeMessages" ) { return Err(UsdtError::NetworkUnavailable); } @@ -171,6 +180,50 @@ impl Rpc { *next = scheduled + REQUEST_INTERVAL; } + pub async fn bridge_status( + &self, + transfer: &UsdtTransfer, + ) -> Result { + let (Some(guid), Some(tx_hash)) = + (transfer.bridge_guid.as_deref(), transfer.tx_hash.as_deref()) + else { + return Err(UsdtError::NetworkUnavailable); + }; + let hash: B256 = tx_hash.parse().map_err(|_| UsdtError::InvalidResponse)?; + let response: Value = self.call("bitkit_getBridgeMessages", json!([hash])).await?; + let messages = response["data"] + .as_array() + .ok_or(UsdtError::InvalidResponse)?; + let message = messages.iter().find(|message| { + message["guid"] + .as_str() + .is_some_and(|value| value.eq_ignore_ascii_case(guid)) + && message["pathway"]["srcEid"].as_u64() == Some(30110) + && message["pathway"]["dstEid"].as_u64() + == transfer.destination.endpoint().map(u64::from) + && message["pathway"]["sender"]["address"] + .as_str() + .is_some_and(|a| a.eq_ignore_ascii_case(&super::types::OFT.to_string())) + && message["source"]["tx"]["txHash"] + .as_str() + .is_some_and(|h| h.eq_ignore_ascii_case(tx_hash)) + }); + Ok(match message.and_then(|m| m["status"]["name"].as_str()) { + Some("DELIVERED") => UsdtTransferStatus::Confirmed, + Some("INFLIGHT" | "CONFIRMING") => UsdtTransferStatus::Bridging, + Some("FAILED" | "BLOCKED" | "PAYLOAD_STORED") => { + UsdtTransferStatus::BridgeNeedsAttention + } + Some("APPLICATION_BURNED" | "APPLICATION_SKIPPED") => UsdtTransferStatus::BridgeFailed, + _ => return Err(UsdtError::NetworkUnavailable), + }) + } + + pub async fn balance(&self, address: Address) -> Result { + self.call("eth_getBalance", json!([address, "pending"])) + .await + } + pub async fn block(&self, number: u64) -> Result { self.call::>("eth_getBlockByNumber", json!([U256::from(number), false])) .await? @@ -319,6 +372,41 @@ mod tests { } } + #[tokio::test] + async fn bridge_reverts_are_distinct_from_rpc_failures() { + use super::super::types::{BRIDGE_HELPER, OFT, TOKEN}; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + for (target, code, expected_route) in [ + (OFT, 3, true), + (BRIDGE_HELPER, 3, true), + (TOKEN, 3, false), + (BRIDGE_HELPER, -32602, false), + (BRIDGE_HELPER, -32002, false), + ] { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("http://{}", listener.local_addr().unwrap()); + let server = tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.unwrap(); + let mut request = [0; 4096]; + assert!(socket.read(&mut request).await.unwrap() > 0); + let body = json!({"error":{"code":code,"message":"Provider rejected the request"}}) + .to_string(); + socket.write_all(format!("HTTP/1.1 200 OK\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}", body.len()).as_bytes()).await.unwrap(); + }); + let error = Rpc::new(url, 42161) + .unwrap() + .call::("eth_call", json!([{"to": target, "data":"0x"}, "latest"])) + .await + .unwrap_err(); + if expected_route { + assert!(matches!(error, UsdtError::UnsupportedRoute)); + } else { + assert!(matches!(error, UsdtError::NetworkUnavailable)); + } + server.await.unwrap(); + } + } + #[tokio::test(start_paused = true)] async fn chain_and_bundler_share_bursts_and_sustained_budget() { let chain = Rpc::new("https://chain.example".into(), 42161).unwrap(); diff --git a/src/modules/usdt/store.rs b/src/modules/usdt/store.rs index d361566..310520b 100644 --- a/src/modules/usdt/store.rs +++ b/src/modules/usdt/store.rs @@ -286,6 +286,26 @@ impl Store { let mut transfer = transfer.clone(); if let Some(saved) = existing { transfer.id = saved.id; + if saved + .tx_hash + .as_deref() + .zip(transfer.tx_hash.as_deref()) + .is_some_and(|(a, b)| a.eq_ignore_ascii_case(b)) + && saved + .bridge_guid + .as_ref() + .zip(transfer.bridge_guid.as_ref()) + .is_some_and(|(a, b)| a.eq_ignore_ascii_case(b)) + && transfer.status == UsdtTransferStatus::Bridging + && matches!( + saved.status, + UsdtTransferStatus::Confirmed + | UsdtTransferStatus::BridgeNeedsAttention + | UsdtTransferStatus::BridgeFailed + ) + { + transfer.status = saved.status; + } write_transfer(tx, &transfer)?; } else { tx.execute( @@ -297,6 +317,15 @@ impl Store { Ok(()) } + pub fn awaiting_delivery(&self) -> Result, UsdtError> { + let connection = self.connection()?; + let mut statement = connection.prepare( + "SELECT data FROM usdt_transfers WHERE json_extract(data, '$.status') IN ('Bridging','BridgeNeedsAttention') AND json_extract(data, '$.bridge_guid') IS NOT NULL ORDER BY json_extract(data, '$.timestamp') DESC, id", + )?; + let rows = statement.query_map([], |row| row.get::<_, String>(0))?; + rows.map(|row| decode(&row?)).collect() + } + pub fn pending_operation(&self) -> Result, UsdtError> { let saved: Option<(String, String)> = self .connection()? diff --git a/src/modules/usdt/tests.rs b/src/modules/usdt/tests.rs index ae20583..5829e41 100644 --- a/src/modules/usdt/tests.rs +++ b/src/modules/usdt/tests.rs @@ -241,7 +241,14 @@ struct ChainState { block_timestamps: std::collections::BTreeMap, reject_broadcast: bool, delay_gas_estimate: bool, + bridge_messages: std::collections::HashMap, + bridge_requests: Vec, + bridge_delay: std::time::Duration, paymaster: alloy_primitives::Address, + helper_balance: alloy_primitives::U256, + native_message_fee: u64, + helper_token_fee: u64, + quote_revert: Option, history_input: Option, history_target: Option, receipt_logs: Option>, @@ -296,7 +303,14 @@ impl MockChain { block_timestamps: Default::default(), reject_broadcast: false, delay_gas_estimate: false, + bridge_messages: Default::default(), + bridge_requests: vec![], + bridge_delay: std::time::Duration::ZERO, paymaster: paymaster::PAYMASTER, + helper_balance: U256::from(1_000_000_000_000_000u64), + native_message_fee: 10_000_000_000, + helper_token_fee: 300_000, + quote_revert: None, history_input: None, history_target: Some(account::ENTRY_POINT), receipt_logs: None, @@ -324,6 +338,7 @@ impl MockChain { })); let server_state = state.clone(); let task = tokio::spawn(async move { + let mut requests = tokio::task::JoinSet::new(); while let Ok((mut socket, _)) = listener.accept().await { let mut request = Vec::new(); let header_end = loop { @@ -390,6 +405,15 @@ impl MockChain { if delay { tokio::time::sleep(std::time::Duration::from_secs(6)).await; } + let bridge_delay = if method == "bitkit_getBridgeMessages" { + let mut state = server_state.lock().unwrap(); + state + .bridge_requests + .push(body["params"][0].as_str().unwrap().into()); + state.bridge_delay + } else { + std::time::Duration::ZERO + }; let mut response = server_state.lock().unwrap().respond(&body); if body["method"] == "eth_getLogs" { if let Some(logs) = response["result"].as_array_mut() { @@ -409,7 +433,11 @@ impl MockChain { } let response = response.to_string(); let response=format!("HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{response}",response.len()); - let _ = socket.write_all(response.as_bytes()).await; + while requests.try_join_next().is_some() {} + requests.spawn(async move { + tokio::time::sleep(bridge_delay).await; + let _ = socket.write_all(response.as_bytes()).await; + }); } }); Self { url, state, task } @@ -451,6 +479,11 @@ impl ChainState { } } let result = match body["method"].as_str().unwrap() { + "bitkit_getBridgeMessages" => self + .bridge_messages + .get(&body["params"][0].as_str().unwrap().to_ascii_lowercase()) + .cloned() + .unwrap_or_else(|| json!({"data":[]})), "eth_chainId" => json!(U256::from(self.chain)), "eth_blockNumber" => json!(U256::from(self.tip)), "eth_getBlockByNumber" => { @@ -482,9 +515,41 @@ impl ChainState { } "eth_getCode" => json!(self.account_code), "eth_getTransactionCount" => json!(U256::from(self.authorization_nonce)), + "eth_getBalance" => json!(self.helper_balance), "eth_call" => { let data: Bytes = serde_json::from_value(body["params"][0]["data"].clone()).unwrap(); + use transaction::{BridgeHelper, MessagingFee, OFTLimit, OFTReceipt, Oft}; + let target: alloy_primitives::Address = + serde_json::from_value(body["params"][0]["to"].clone()).unwrap(); + if self.quote_revert == Some(target) + && (data.starts_with(&Oft::quoteSendCall::SELECTOR) + || data.starts_with(&BridgeHelper::quoteSendCall::SELECTOR)) + { + return json!({"jsonrpc":"2.0","id":1,"error":{"code":3,"message":"Bridge contract call reverted"}}); + } + if data.starts_with(&Oft::tokenCall::SELECTOR) { + assert!([types::OFT, types::BRIDGE_HELPER].contains(&target)); + } + if [ + Oft::peersCall::SELECTOR, + Oft::quoteOFTCall::SELECTOR, + Oft::quoteSendCall::SELECTOR, + ] + .iter() + .any(|selector| data.starts_with(selector)) + { + assert_eq!(target, types::OFT); + } + if [ + BridgeHelper::maxGasCall::SELECTOR, + BridgeHelper::quoteSendCall::SELECTOR, + ] + .iter() + .any(|selector| data.starts_with(selector)) + { + assert_eq!(target, types::BRIDGE_HELPER); + } let encoded = if data.starts_with(&transaction::EntryPoint::getNonceCall::SELECTOR) { let block = serde_json::from_value::(body["params"][1].clone()).ok(); @@ -496,6 +561,37 @@ impl ChainState { self.nonce }; U256::from(nonce).abi_encode() + } else if data.starts_with(&Oft::tokenCall::SELECTOR) { + types::TOKEN.abi_encode() + } else if data.starts_with(&Oft::peersCall::SELECTOR) { + alloy_primitives::B256::repeat_byte(1).abi_encode() + } else if data.starts_with(&Oft::quoteOFTCall::SELECTOR) { + let param = Oft::quoteOFTCall::abi_decode(&data).unwrap().param; + ( + OFTLimit { + minAmountLD: U256::from(1), + maxAmountLD: U256::MAX, + }, + Vec::::new(), + OFTReceipt { + amountSentLD: param.amountLD, + amountReceivedLD: param.amountLD, + }, + ) + .abi_encode_params() + } else if data.starts_with(&Oft::quoteSendCall::SELECTOR) { + MessagingFee { + nativeFee: U256::from(self.native_message_fee), + lzTokenFee: U256::ZERO, + } + .abi_encode() + } else if data.starts_with(&BridgeHelper::maxGasCall::SELECTOR) { + U256::from(1_000_000_000_000_000u64).abi_encode() + } else if data.starts_with(&BridgeHelper::quoteSendCall::SELECTOR) { + let param = BridgeHelper::quoteSendCall::abi_decode(&data) + .unwrap() + .param; + (param.amountLD + U256::from(self.helper_token_fee)).abi_encode() } else { self.balance.abi_encode() }; @@ -773,11 +869,11 @@ async fn signed_operation_survives_uncertain_broadcast_and_restart() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let next = wallet - .quote_transfer(RECIPIENT.into(), 2_000_000) + .quote_transfer(RECIPIENT.into(), 2_000_000, UsdtDestination::Arbitrum) .await .unwrap(); chain.state.lock().unwrap().reject_broadcast = true; @@ -810,7 +906,9 @@ async fn signed_operation_survives_uncertain_broadcast_and_restart() { // A pending payment is rejected locally even if the provider is now misconfigured. chain.state.lock().unwrap().chain = 1; assert!(matches!( - wallet.quote_transfer(RECIPIENT.into(), 2_000_000).await, + wallet + .quote_transfer(RECIPIENT.into(), 2_000_000, UsdtDestination::Arbitrum) + .await, Err(UsdtError::PendingTransfer) )); assert!(matches!( @@ -849,7 +947,7 @@ async fn fluctuating_gas_estimates_preserve_the_reviewed_operation() { chain.state.lock().unwrap().pre_verification_estimates = [80_000, 80_100, 80_200, 80_300].into(); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let reviewed = wallet.store.quote("e.id).unwrap().plan.operation; @@ -889,12 +987,14 @@ async fn wrong_network_owner_nonce_balance_and_paymaster_cannot_sign() { types::TOKEN.to_checksum(None) ); assert!(matches!( - wallet.quote_transfer(request, 1_000_000).await, + wallet + .quote_transfer(request, 1_000_000, UsdtDestination::Arbitrum) + .await, Err(UsdtError::InvalidAddress) )); } let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); assert!(matches!( @@ -936,7 +1036,9 @@ async fn wrong_network_owner_nonce_balance_and_paymaster_cannot_sign() { state.paymaster = Address::repeat_byte(1); } assert!(matches!( - wallet.quote_transfer(RECIPIENT.into(), 1_000_000).await, + wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) + .await, Err(UsdtError::UnsupportedRoute) )); assert!(chain.state.lock().unwrap().operations.is_empty()); @@ -953,7 +1055,7 @@ async fn expired_unmined_operation_releases_nonce_for_a_new_approval() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); wallet @@ -973,7 +1075,9 @@ async fn expired_unmined_operation_releases_nonce_for_a_new_approval() { UsdtTransferStatus::Pending ); assert!(matches!( - wallet.quote_transfer(RECIPIENT.into(), 1_000_000).await, + wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) + .await, Err(UsdtError::PendingTransfer) )); { @@ -995,7 +1099,7 @@ async fn expired_unmined_operation_releases_nonce_for_a_new_approval() { UsdtTransferStatus::Failed ); let next = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); assert_eq!( @@ -1009,13 +1113,112 @@ async fn expired_unmined_operation_releases_nonce_for_a_new_approval() { } } +#[tokio::test] +async fn bridge_payment_bounds_token_fees_and_revokes_helper_approval() { + use alloy_primitives::U256; + use alloy_sol_types::SolCall; + let chain = MockChain::start().await; + let dir = tempfile::tempdir().unwrap(); + let wallet = chain.wallet(&dir); + let quote = wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Polygon) + .await + .unwrap(); + let plan = wallet.store.quote("e.id).unwrap().plan; + let calls = history::decode_calls(&plan.operation.call_data).unwrap(); + assert_eq!(calls.len(), 4); + let paymaster = transaction::Erc20::approveCall::abi_decode(&calls[0].1).unwrap(); + assert_eq!(calls[0].0, types::TOKEN); + assert_eq!(paymaster.spender, paymaster::PAYMASTER); + let approval = transaction::Erc20::approveCall::abi_decode(&calls[1].1).unwrap(); + assert_eq!(calls[1].0, types::TOKEN); + assert_eq!(approval.spender, types::BRIDGE_HELPER); + assert_eq!(approval.amount, U256::from(1_360_001)); + let send = transaction::BridgeHelper::sendCall::abi_decode(&calls[2].1).unwrap(); + assert_eq!(calls[2].0, types::BRIDGE_HELPER); + assert_eq!(send.oft, types::OFT); + assert_eq!(send.param.dstEid, 30109); + assert_eq!( + send.param.to, + RECIPIENT + .parse::() + .unwrap() + .into_word() + ); + assert_eq!(send.param.minAmountLD, U256::from(1_000_000)); + assert_eq!(send.fee.nativeFee, U256::from(11_000_000_001u64)); + let revoke = transaction::Erc20::approveCall::abi_decode(&calls[3].1).unwrap(); + assert_eq!(revoke.spender, types::BRIDGE_HELPER); + assert!(revoke.amount.is_zero()); + let bridge_fee = approval.amount.to::() - quote.amount; + assert_eq!(bridge_fee, 360_001); + assert!(paymaster.amount > U256::from(quote.maximum_fee - bridge_fee)); + assert_eq!(quote.received_amount, 1_000_000); + for target in [types::OFT, types::BRIDGE_HELPER] { + chain.state.lock().unwrap().quote_revert = Some(target); + assert!(matches!( + wallet + .send(quote.id.clone(), TEST_PHRASE.into(), None) + .await, + Err(UsdtError::QuoteExpired) + )); + assert!(matches!( + wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Polygon) + .await, + Err(UsdtError::UnsupportedRoute) + )); + } + chain.state.lock().unwrap().quote_revert = None; + chain.state.lock().unwrap().native_message_fee = 12_000_000_000; + assert!(matches!( + wallet + .send(quote.id.clone(), TEST_PHRASE.into(), None) + .await, + Err(UsdtError::QuoteExpired) + )); + chain.state.lock().unwrap().native_message_fee = 10_500_000_000; + chain.state.lock().unwrap().helper_token_fee = 400_000; + assert!(matches!( + wallet + .send(quote.id.clone(), TEST_PHRASE.into(), None) + .await, + Err(UsdtError::QuoteExpired) + )); + chain.state.lock().unwrap().helper_token_fee = 300_000; + chain.state.lock().unwrap().helper_balance = U256::ZERO; + assert!(matches!( + wallet + .send(quote.id.clone(), TEST_PHRASE.into(), None) + .await, + Err(UsdtError::UnsupportedRoute) + )); + assert!(chain.state.lock().unwrap().operations.is_empty()); + assert!(matches!( + wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Polygon) + .await, + Err(UsdtError::UnsupportedRoute) + )); + chain.state.lock().unwrap().helper_balance = U256::from(1_000_000_000_000_000u64); + let sent = wallet + .send(quote.id, TEST_PHRASE.into(), None) + .await + .unwrap(); + chain.state.lock().unwrap().helper_balance = U256::ZERO; + chain.state.lock().unwrap().tip += 3; + wallet.refresh_transfers().await.unwrap(); + assert_eq!(chain.state.lock().unwrap().operations.len(), 1); + assert!(wallet.store.pending_plan(&sent.id).unwrap().is_some()); +} + #[tokio::test] async fn seed_restore_recovers_mined_payments_without_local_submission_data() { let chain = MockChain::start().await; let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -1043,65 +1246,93 @@ async fn seed_restore_recovers_mined_payments_without_local_submission_data() { } #[tokio::test] -async fn bundled_operations_cannot_contribute_another_payments_fee() { +async fn bundled_operations_cannot_contribute_another_payments_bridge_status_or_fee() { use alloy_primitives::{B256, U256}; use alloy_sol_types::SolEvent; use serde_json::json; - let chain = MockChain::start().await; - let dir = tempfile::tempdir().unwrap(); - let wallet = chain.wallet(&dir); - let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) - .await - .unwrap(); - let mut transfer = wallet - .send(quote.id, TEST_PHRASE.into(), None) - .await - .unwrap(); - let own_hash = transfer - .user_operation_hash - .as_ref() - .unwrap() - .parse() - .unwrap(); - let other_hash = B256::repeat_byte(2); - let event = |hash, success| { - transaction::EntryPoint::UserOperationEvent { - userOpHash: hash, - sender: wallet.address, - paymaster: paymaster::PAYMASTER, - nonce: U256::ZERO, - success, - actualGasCost: U256::from(1), - actualGasUsed: U256::from(1), - } - .encode_log_data() - }; - let log = |address, data: alloy_primitives::LogData| json!({"address":address,"topics":data.topics(),"data":data.data}); - let receipt = json!({"logs":[ - log(paymaster::PAYMASTER, transaction::Paymaster::UserOperationSponsored { userOpHash:other_hash, user:wallet.address, paymasterMode:1, token:types::TOKEN, tokenAmountPaid:U256::from(500), exchangeRate:U256::from(1) }.encode_log_data()), - log(account::ENTRY_POINT, event(other_hash, true)), - log(paymaster::PAYMASTER, transaction::Paymaster::UserOperationSponsored { userOpHash:own_hash, user:wallet.address, paymasterMode:1, token:types::TOKEN, tokenAmountPaid:U256::from(123), exchangeRate:U256::from(1) }.encode_log_data()), - log(types::TOKEN, transaction::Erc20::Transfer { from: wallet.address, to: RECIPIENT.parse().unwrap(), value: U256::from(transfer.amount) }.encode_log_data()), - log(account::ENTRY_POINT, event(own_hash, true)), - ]}); - assert_eq!( - transaction::operation_logs(&receipt, own_hash).unwrap(), - &receipt["logs"].as_array().unwrap()[2..] - ); - let mut failed = receipt.clone(); - failed["logs"][4] = log(account::ENTRY_POINT, event(own_hash, false)); - wallet.settle(&mut transfer, &failed).unwrap(); - assert_eq!(transfer.status, UsdtTransferStatus::Failed); - assert_eq!(transfer.fee, Some(123)); - wallet.settle(&mut transfer, &receipt).unwrap(); - assert_eq!(transfer.status, UsdtTransferStatus::Confirmed); - assert_eq!(transfer.fee, Some(123)); + for (destination, status, fee, fee_with_bridge_log) in [ + ( + UsdtDestination::Arbitrum, + UsdtTransferStatus::Confirmed, + Some(123), + Some(123), + ), + ( + UsdtDestination::Polygon, + UsdtTransferStatus::BridgeNeedsAttention, + None, + Some(623), + ), + ] { + let chain = MockChain::start().await; + let dir = tempfile::tempdir().unwrap(); + let wallet = chain.wallet(&dir); + let quote = wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, destination) + .await + .unwrap(); + let mut transfer = wallet + .send(quote.id, TEST_PHRASE.into(), None) + .await + .unwrap(); + let own_hash = transfer + .user_operation_hash + .as_ref() + .unwrap() + .parse() + .unwrap(); + let other_hash = B256::repeat_byte(2); + let guid = B256::repeat_byte(3); + let event = |hash, success| { + transaction::EntryPoint::UserOperationEvent { + userOpHash: hash, + sender: wallet.address, + paymaster: paymaster::PAYMASTER, + nonce: U256::ZERO, + success, + actualGasCost: U256::from(1), + actualGasUsed: U256::from(1), + } + .encode_log_data() + }; + let log = |address, data: alloy_primitives::LogData| json!({"address":address,"topics":data.topics(),"data":data.data}); + let receipt = json!({"logs":[ + log(types::OFT, transaction::Oft::OFTSent { guid, dstEid:30109, fromAddress:types::BRIDGE_HELPER, amountSentLD:U256::from(1_000_000), amountReceivedLD:U256::from(1_000_000) }.encode_log_data()), + log(types::BRIDGE_HELPER, transaction::BridgeHelper::LogSend { sender:wallet.address, oft:types::OFT, amountLD:U256::from(1_000_000), nativeFee:U256::from(100), feeInToken:U256::from(500), totalAmount:U256::from(1_000_500) }.encode_log_data()), + log(account::ENTRY_POINT, event(other_hash, true)), + log(types::TOKEN, transaction::Erc20::Transfer { from:wallet.address, to:RECIPIENT.parse().unwrap(), value:U256::from(1_000_000) }.encode_log_data()), + log(paymaster::PAYMASTER, transaction::Paymaster::UserOperationSponsored { userOpHash:own_hash, user:wallet.address, paymasterMode:1, token:types::TOKEN, tokenAmountPaid:U256::from(123), exchangeRate:U256::from(1) }.encode_log_data()), + log(account::ENTRY_POINT, event(own_hash, true)), + ]}); + assert_eq!( + transaction::operation_logs(&receipt, own_hash).unwrap(), + &receipt["logs"].as_array().unwrap()[3..] + ); + let mut failed = receipt.clone(); + failed["logs"][5] = log(account::ENTRY_POINT, event(own_hash, false)); + wallet.settle(&mut transfer, &failed).unwrap(); + assert_eq!(transfer.status, UsdtTransferStatus::Failed); + assert_eq!(transfer.fee, Some(123)); + assert_eq!(transfer.bridge_guid, None); + wallet.settle(&mut transfer, &receipt).unwrap(); + assert_eq!(transfer.status, status); + assert_eq!(transfer.bridge_guid, None); + assert_eq!(transfer.fee, fee); + + let mut receipt = receipt; + let bridge_log = receipt["logs"][1].clone(); + receipt["logs"] + .as_array_mut() + .unwrap() + .insert(5, bridge_log); + wallet.settle(&mut transfer, &receipt).unwrap(); + assert_eq!(transfer.fee, fee_with_bridge_log); + } } #[tokio::test] #[ignore = "requires a fresh local Arbitrum fork and tests/usdt-fork/provider.mjs"] -async fn deployed_contracts_collect_usdt_fees_without_account_eth() { +async fn deployed_contracts_collect_usdt_fees_and_revert_failed_bridges_atomically() { use alloy_primitives::U256; use serde_json::json; let rpc = rpc::Rpc::new("http://127.0.0.1:18545".into(), types::CHAIN_ID).unwrap(); @@ -1111,16 +1342,11 @@ async fn deployed_contracts_collect_usdt_fees_without_account_eth() { let wallet = UsdtWallet::new( usdt_address(TEST_PHRASE.into(), None).unwrap(), dir.path().join("usdt.sqlite").to_string_lossy().into(), - std::env::var("USDT_FORK_RPC_URL").unwrap_or_else(|_| "http://127.0.0.1:18545".into()), + std::env::var("USDT_FORK_RPC_URL").unwrap_or_else(|_| "http://127.0.0.1:18546".into()), std::env::var("USDT_FORK_BUNDLER_URL").unwrap_or_else(|_| "http://127.0.0.1:18546".into()), ) .unwrap(); - assert_eq!( - rpc.call::("eth_getBalance", json!([wallet.address, "pending"])) - .await - .unwrap(), - U256::ZERO - ); + assert_eq!(rpc.balance(wallet.address).await.unwrap(), U256::ZERO); let initial = wallet.balance().await.unwrap(); // Only locally mined transactions belong to this fixture's history. wallet @@ -1128,7 +1354,7 @@ async fn deployed_contracts_collect_usdt_fees_without_account_eth() { .complete_history(wallet.block_number().await.unwrap()) .unwrap(); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -1141,12 +1367,86 @@ async fn deployed_contracts_collect_usdt_fees_without_account_eth() { let fee = transfer.fee.unwrap(); assert!(fee > 0 && fee <= quote.maximum_fee); assert_eq!(wallet.balance().await.unwrap(), initial - 1_000_000 - fee); + + let bridge = wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Polygon) + .await + .unwrap(); + let helper_balance = rpc.balance(types::BRIDGE_HELPER).await.unwrap(); + let _: serde_json::Value = rpc + .call("anvil_setBalance", json!([types::BRIDGE_HELPER, "0x0"])) + .await + .unwrap(); + assert!(matches!( + wallet + .send(bridge.id.clone(), TEST_PHRASE.into(), None) + .await, + Err(UsdtError::UnsupportedRoute) + )); + let _: serde_json::Value = rpc + .call( + "anvil_setBalance", + json!([types::BRIDGE_HELPER, helper_balance]), + ) + .await + .unwrap(); + // The helper can lose liquidity after preflight but before execution. + let fixture = rpc::Rpc::new("http://127.0.0.1:18546".into(), types::CHAIN_ID).unwrap(); + let _: bool = fixture + .call("test_drainHelperBeforeNextBroadcast", json!([])) + .await + .unwrap(); + let before = wallet.balance().await.unwrap(); + let sent = wallet + .send(bridge.id, TEST_PHRASE.into(), None) + .await + .unwrap(); + let history = wallet.refresh_transfers().await.unwrap(); + let failed = history.iter().find(|t| t.id == sent.id).unwrap(); + assert_eq!(failed.status, UsdtTransferStatus::Failed); + assert!(failed.bridge_guid.is_none()); + assert!(before - wallet.balance().await.unwrap() <= bridge.maximum_fee); + let _: serde_json::Value = rpc + .call( + "anvil_setBalance", + json!([types::BRIDGE_HELPER, helper_balance]), + ) + .await + .unwrap(); + + let bridge = wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Polygon) + .await + .unwrap(); + let before = wallet.balance().await.unwrap(); + let sent = wallet + .send(bridge.id, TEST_PHRASE.into(), None) + .await + .unwrap(); + let history = wallet.refresh_transfers().await.unwrap(); + let bridged = history.iter().find(|t| t.id == sent.id).unwrap(); + assert_eq!(bridged.status, UsdtTransferStatus::Bridging); + assert!(bridged.bridge_guid.is_some()); + assert!(bridged.fee.unwrap() > 0); + assert!(bridged.fee.unwrap() <= bridge.maximum_fee); + sync_history_to_tip(&wallet).await; assert_eq!( - rpc.call::("eth_getBalance", json!([wallet.address, "pending"])) - .await - .unwrap(), - U256::ZERO + wallet.balance().await.unwrap(), + before - 1_000_000 - bridged.fee.unwrap() ); + alloy_sol_types::sol! { function allowance(address owner, address spender) view returns (uint256); } + let allowance = rpc + .contract( + types::TOKEN, + allowanceCall { + owner: wallet.address, + spender: types::BRIDGE_HELPER, + }, + ) + .await + .unwrap(); + assert!(allowance.is_zero()); + assert_eq!(rpc.balance(wallet.address).await.unwrap(), U256::ZERO); } #[tokio::test] @@ -1157,7 +1457,7 @@ async fn history_preserves_receipts_with_external_account_call_shapes() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); wallet @@ -1230,7 +1530,7 @@ async fn wrapped_history_preserves_signed_payments_and_restores_token_transfers( let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -1328,7 +1628,9 @@ fn stored_activity_is_complete_and_sorted_newest_first() { id: format!("receipt-{index}"), tx_hash: Some(format!("tx-{index}")), user_operation_hash: None, + bridge_guid: None, recipient: RECIPIENT.into(), + destination: UsdtDestination::Arbitrum, amount: 1, received_amount: 1, fee: None, @@ -1352,14 +1654,20 @@ async fn nonce_advance_with_delayed_logs_remains_pending_and_history_reconciles( let chain = MockChain::start().await; let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); - for recipient in [wallet.receive_address(), types::TOKEN.to_checksum(None)] { + for (recipient, destination) in [ + (wallet.receive_address(), UsdtDestination::Arbitrum), + (wallet.receive_address(), UsdtDestination::Ethereum), + (types::TOKEN.to_checksum(None), UsdtDestination::Arbitrum), + ] { assert!(matches!( - wallet.quote_transfer(recipient, 1_000_000).await, + wallet + .quote_transfer(recipient, 1_000_000, destination) + .await, Err(UsdtError::InvalidAddress) )); } let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -1383,7 +1691,9 @@ async fn nonce_advance_with_delayed_logs_remains_pending_and_history_reconciles( ); assert!(wallet.store.pending_plan(&sent.id).unwrap().is_some()); assert!(matches!( - wallet.quote_transfer(RECIPIENT.into(), 1_000_000).await, + wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) + .await, Err(UsdtError::PendingTransfer) )); { @@ -1406,7 +1716,7 @@ async fn interrupted_history_resumes_without_repeating_completed_work() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); wallet @@ -1464,7 +1774,7 @@ async fn replacement_after_expiry_recovers_pending_send_after_restart() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); chain.state.lock().unwrap().reject_broadcast = true; @@ -1516,7 +1826,7 @@ async fn replacement_after_expiry_recovers_pending_send_after_restart() { ); assert!(wallet.store.pending_plan(&sent.id).unwrap().is_none()); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); assert_eq!( @@ -1618,7 +1928,7 @@ async fn quote_expiring_during_validation_is_not_signed() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let mut data = wallet.store.quote("e.id).unwrap(); @@ -1673,6 +1983,143 @@ async fn invalid_chain_data_and_stored_json_have_distinct_errors() { )); } +#[tokio::test] +async fn stalled_bridge_status_checks_leave_time_for_source_recovery_and_sending() { + use tokio::time::Duration; + let chain = MockChain::start().await; + let directory = tempfile::tempdir().unwrap(); + let wallet = chain.wallet(&directory); + chain.state.lock().unwrap().bridge_delay = Duration::from_secs(25); + let bridges: Vec<_> = (0..5) + .map(|index| UsdtTransfer { + id: format!("bridge-{index}"), + tx_hash: Some(format!("{:#x}", alloy_primitives::B256::repeat_byte(index))), + user_operation_hash: None, + bridge_guid: Some(format!( + "{:#x}", + alloy_primitives::B256::repeat_byte(index + 10) + )), + recipient: RECIPIENT.into(), + destination: UsdtDestination::Polygon, + amount: 1_000_000, + received_amount: 1_000_000, + fee: Some(20), + is_incoming: false, + status: UsdtTransferStatus::Bridging, + timestamp: 1, + }) + .collect(); + for bridge in &bridges { + chain.state.lock().unwrap().bridge_messages.insert( + bridge.tx_hash.clone().unwrap(), + serde_json::json!({"data":[{ + "guid":bridge.bridge_guid, + "pathway":{"srcEid":30110,"dstEid":30109,"sender":{"address":types::OFT}}, + "source":{"tx":{"txHash":bridge.tx_hash}},"status":{"name":"DELIVERED"} + }]}), + ); + } + wallet + .store + .save_history_receipt(&bridges, "bridges", 1000, "block", true) + .unwrap(); + // A signed operation with no nonce consumption must still resolve once expired. + let quote = wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) + .await + .unwrap(); + chain.state.lock().unwrap().reject_broadcast = true; + let sent = wallet + .send(quote.id, TEST_PHRASE.into(), None) + .await + .unwrap(); + { + let mut state = chain.state.lock().unwrap(); + state.tip += 3; + state.timestamp += alloy_primitives::U256::from(601); + } + // Completion includes bridge polling and the shared chain/bundler request budget. + let history = tokio::time::timeout(Duration::from_secs(25), wallet.refresh_transfers()) + .await + .unwrap() + .unwrap(); + assert_eq!(chain.state.lock().unwrap().bridge_requests.len(), 3); + assert_eq!( + history + .iter() + .find(|transfer| transfer.id == sent.id) + .unwrap() + .status, + UsdtTransferStatus::Failed + ); + assert_eq!( + history + .iter() + .filter(|transfer| transfer.status == UsdtTransferStatus::Bridging) + .count(), + 5 + ); + chain.state.lock().unwrap().reject_broadcast = false; + let quote = wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) + .await + .unwrap(); + // Replenish the shared RPC burst so this checks polling, not rate-limit waiting. + tokio::time::sleep(Duration::from_secs(15)).await; + // Destination polling must not hold the mutation lock needed to send. + let refresh_wallet = wallet.clone(); + let refresh = tokio::spawn(async move { refresh_wallet.refresh_transfers().await }); + tokio::time::timeout(Duration::from_secs(3), async { + while chain.state.lock().unwrap().bridge_requests.len() < 5 { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .unwrap(); + let result = tokio::time::timeout( + Duration::from_secs(8), + wallet.send(quote.id, TEST_PHRASE.into(), None), + ) + .await; + assert_eq!(result.unwrap().unwrap().status, UsdtTransferStatus::Pending); + assert!(!refresh.is_finished()); + refresh.await.unwrap().unwrap(); + assert_eq!( + chain + .state + .lock() + .unwrap() + .bridge_requests + .iter() + .collect::>() + .len(), + 5 + ); + // Failed lookups wait a minute, then recover without delaying source reconciliation. + let attempts = chain.state.lock().unwrap().bridge_requests.len(); + wallet.refresh_transfers().await.unwrap(); + assert_eq!(chain.state.lock().unwrap().bridge_requests.len(), attempts); + tokio::time::pause(); + tokio::time::advance(Duration::from_secs(60)).await; + tokio::time::resume(); + chain.state.lock().unwrap().bridge_delay = Duration::from_millis(1200); + chain.state.lock().unwrap().tip += 3; + for _ in 0..2 { + { + let mut state = chain.state.lock().unwrap(); + state.fail_block_read_at = Some(state.block_reads + 1); + } + tokio::time::timeout(Duration::from_secs(25), wallet.refresh_transfers()) + .await + .unwrap() + .unwrap_err(); + } + let history = wallet.history().unwrap(); + assert!(bridges.iter().all(|bridge| history.iter().any( + |transfer| transfer.id == bridge.id && transfer.status == UsdtTransferStatus::Confirmed + ))); +} + #[tokio::test] async fn each_payment_authorizes_the_current_nonce_after_delegation() { use alloy_primitives::{Bytes, U256}; @@ -1680,7 +2127,7 @@ async fn each_payment_authorizes_the_current_nonce_after_delegation() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); wallet @@ -1701,7 +2148,7 @@ async fn each_payment_authorizes_the_current_nonce_after_delegation() { UsdtTransferStatus::Confirmed ); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); wallet @@ -1724,7 +2171,7 @@ async fn consumed_authorization_preserves_pending_payment_until_signed_expiry() let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let quote_data = wallet.store.quote("e.id).unwrap(); @@ -1738,7 +2185,7 @@ async fn consumed_authorization_preserves_pending_payment_until_signed_expiry() )); assert!(chain.state.lock().unwrap().operations.is_empty()); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); chain.state.lock().unwrap().reject_broadcast = true; @@ -1772,7 +2219,9 @@ async fn consumed_authorization_preserves_pending_payment_until_signed_expiry() assert_eq!(retained.operation.eip7702_auth, signed.eip7702_auth); assert_eq!(retained.operation.signature, signed.signature); assert!(matches!( - restored.quote_transfer(RECIPIENT.into(), 1).await, + restored + .quote_transfer(RECIPIENT.into(), 1, UsdtDestination::Arbitrum) + .await, Err(UsdtError::PendingTransfer) )); chain.state.lock().unwrap().timestamp = U256::from(retained.expires_at + 1); @@ -1780,7 +2229,10 @@ async fn consumed_authorization_preserves_pending_payment_until_signed_expiry() assert_eq!(history[0].status, UsdtTransferStatus::Failed); assert_eq!(history[0].fee, Some(0)); assert!(restored.store.pending_plan(&sent.id).unwrap().is_none()); - let quote = restored.quote_transfer(RECIPIENT.into(), 1).await.unwrap(); + let quote = restored + .quote_transfer(RECIPIENT.into(), 1, UsdtDestination::Arbitrum) + .await + .unwrap(); assert_eq!( restored .store @@ -1802,7 +2254,9 @@ async fn foreign_delegation_and_unbounded_paymaster_terms_cannot_authorize_payme let wallet = chain.wallet(&dir); chain.state.lock().unwrap().account_code = Bytes::from_static(&[0xef, 0x01, 0x00, 1]); assert!(matches!( - wallet.quote_transfer(RECIPIENT.into(), 1).await, + wallet + .quote_transfer(RECIPIENT.into(), 1, UsdtDestination::Arbitrum) + .await, Err(UsdtError::UnsupportedDelegation) )); assert_eq!(wallet.balance().await.unwrap(), 10_000_000); @@ -1812,12 +2266,17 @@ async fn foreign_delegation_and_unbounded_paymaster_terms_cannot_authorize_payme for expiry in [0, timestamp + 901] { chain.state.lock().unwrap().paymaster_valid_until = Some(expiry); assert!(matches!( - wallet.quote_transfer(RECIPIENT.into(), 1).await, + wallet + .quote_transfer(RECIPIENT.into(), 1, UsdtDestination::Arbitrum) + .await, Err(UsdtError::InvalidResponse) )); } chain.state.lock().unwrap().paymaster_valid_until = Some(timestamp + 900); - let quote = wallet.quote_transfer(RECIPIENT.into(), 1).await.unwrap(); + let quote = wallet + .quote_transfer(RECIPIENT.into(), 1, UsdtDestination::Arbitrum) + .await + .unwrap(); chain.state.lock().unwrap().account_code = Bytes::from( [ &[0xef, 0x01, 0x00][..], @@ -1840,7 +2299,10 @@ async fn seed_restore_includes_external_token_sends_without_duplicate_operation_ let chain = MockChain::start().await; let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); - let quote = wallet.quote_transfer(RECIPIENT.into(), 77).await.unwrap(); + let quote = wallet + .quote_transfer(RECIPIENT.into(), 77, UsdtDestination::Arbitrum) + .await + .unwrap(); wallet .send(quote.id, TEST_PHRASE.into(), None) .await @@ -1885,7 +2347,7 @@ async fn gas_price_changes_respect_the_approved_fee() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let plan = wallet.store.quote("e.id).unwrap().plan; @@ -1912,7 +2374,7 @@ async fn consumed_nonce_recovery_requires_complete_receipts_and_resumes_after_re let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -1950,7 +2412,7 @@ async fn consumed_nonce_recovery_requires_complete_receipts_and_resumes_after_re 0 ); restored - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); } @@ -1962,7 +2424,7 @@ async fn consuming_block_receipts_recover_a_payment_hidden_from_log_queries() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -1995,7 +2457,7 @@ async fn consumed_nonce_recovery_preserves_malformed_operation_evidence() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -2045,7 +2507,7 @@ async fn dense_block_history_recovers_large_receipts_without_skipping_after_rest let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -2130,7 +2592,7 @@ async fn first_submission_precheck_releases_an_operation_that_was_never_sent() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); chain @@ -2154,7 +2616,7 @@ async fn first_submission_precheck_releases_an_operation_that_was_never_sent() { assert_eq!(failed.received_amount, 0); assert_eq!(failed.fee, Some(0)); wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); } @@ -2168,7 +2630,7 @@ async fn unknown_paymaster_history_preserves_principal_fee_and_refund() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); wallet @@ -2244,7 +2706,7 @@ async fn settlement_requires_matching_canonical_receipts() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -2300,6 +2762,261 @@ async fn settlement_requires_matching_canonical_receipts() { ); } +#[tokio::test] +async fn bridge_settlement_recovers_guid_fees_and_preserves_delivery_on_rescan() { + use alloy_primitives::{B256, U256}; + use alloy_sol_types::SolEvent; + use serde_json::json; + let chain = MockChain::start().await; + let dir = tempfile::tempdir().unwrap(); + let wallet = chain.wallet(&dir); + let guid = B256::repeat_byte(0xab); + let message = json!({"guid":guid,"pathway":{"srcEid":30110,"dstEid":30109,"sender":{"address":types::OFT}},"source":{"tx":{"txHash":B256::repeat_byte(7)}},"status":{"name":"DELIVERED"}}); + let quote = wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Polygon) + .await + .unwrap(); + let sent = wallet + .send(quote.id, TEST_PHRASE.into(), None) + .await + .unwrap(); + { + let mut state = chain.state.lock().unwrap(); + state.mined = true; + state.tip += 3; + let mut logs = state.event_logs(); + let log = |address, event: alloy_primitives::LogData, index| json!({"address":address,"topics":event.topics(),"data":event.data,"logIndex":U256::from(index)}); + logs.insert( + 1, + log( + types::OFT, + transaction::Oft::OFTSent { + guid, + dstEid: 30109, + fromAddress: types::BRIDGE_HELPER, + amountSentLD: U256::from(1_000_000), + amountReceivedLD: U256::from(999_999), + } + .encode_log_data(), + 1u64, + ), + ); + logs.insert( + 2, + log( + types::BRIDGE_HELPER, + transaction::BridgeHelper::LogSend { + sender: wallet.address, + oft: types::OFT, + amountLD: U256::from(1_000_000), + nativeFee: U256::from(10_000_000_000u64), + feeInToken: U256::from(300_000), + totalAmount: U256::from(1_300_000), + } + .encode_log_data(), + 2u64, + ), + ); + logs[3]["logIndex"] = json!("0x3"); + state.receipt_logs = Some(logs); + } + sync_history_to_tip(&wallet).await; + let pending = wallet.history().unwrap().remove(0); + assert_eq!(pending.status, UsdtTransferStatus::Bridging); + assert_eq!(pending.bridge_guid, Some(format!("{guid:#x}"))); + assert_eq!(pending.received_amount, 999_999); + assert_eq!(pending.fee, Some(300_123)); + assert!(wallet.store.pending_plan(&sent.id).unwrap().is_none()); + for (pointer, value, expected) in [ + ("/guid", json!(B256::ZERO), None), + ("/pathway/srcEid", json!(30101), None), + ("/pathway/dstEid", json!(30101), None), + ("/pathway/sender/address", json!(RECIPIENT), None), + ("/source/tx/txHash", json!(B256::ZERO), None), + ("/status/name", json!("NEW_PROVIDER_STATUS"), None), + ( + "/status/name", + json!("FAILED"), + Some(UsdtTransferStatus::BridgeNeedsAttention), + ), + ( + "/status/name", + json!("BLOCKED"), + Some(UsdtTransferStatus::BridgeNeedsAttention), + ), + ( + "/status/name", + json!("PAYLOAD_STORED"), + Some(UsdtTransferStatus::BridgeNeedsAttention), + ), + ( + "/status/name", + json!("INFLIGHT"), + Some(UsdtTransferStatus::Bridging), + ), + ( + "/status/name", + json!("CONFIRMING"), + Some(UsdtTransferStatus::Bridging), + ), + ( + "/status/name", + json!("APPLICATION_BURNED"), + Some(UsdtTransferStatus::BridgeFailed), + ), + ( + "/status/name", + json!("APPLICATION_SKIPPED"), + Some(UsdtTransferStatus::BridgeFailed), + ), + ] { + let mut changed = message.clone(); + *changed.pointer_mut(pointer).unwrap() = value; + chain + .state + .lock() + .unwrap() + .bridge_messages + .insert(pending.tx_hash.clone().unwrap(), json!({"data":[changed]})); + match expected { + Some(status) => assert_eq!(wallet.rpc.bridge_status(&pending).await.unwrap(), status), + None => assert!(matches!( + wallet.rpc.bridge_status(&pending).await, + Err(UsdtError::NetworkUnavailable) + )), + } + } + let mut retryable = message.clone(); + retryable["status"]["name"] = json!("FAILED"); + chain.state.lock().unwrap().bridge_messages.insert( + pending.tx_hash.clone().unwrap(), + json!({"data":[retryable]}), + ); + assert_eq!( + wallet.refresh_transfers().await.unwrap()[0].status, + UsdtTransferStatus::BridgeNeedsAttention + ); + let requests = chain.state.lock().unwrap().bridge_requests.len(); + wallet.refresh_transfers().await.unwrap(); + assert_eq!(chain.state.lock().unwrap().bridge_requests.len(), requests); + tokio::time::pause(); + tokio::time::advance(std::time::Duration::from_secs(60)).await; + tokio::time::resume(); + chain.state.lock().unwrap().bridge_messages.insert( + pending.tx_hash.clone().unwrap(), + json!({"data":[message.clone()]}), + ); + chain.state.lock().unwrap().chain = 1; + let mut delivered = wallet.refresh_transfers().await.unwrap().remove(0); + assert_eq!(delivered.status, UsdtTransferStatus::Confirmed); + delivered.bridge_guid = delivered.bridge_guid.map(|guid| guid.to_uppercase()); + delivered.tx_hash = delivered.tx_hash.map(|hash| hash.to_uppercase()); + wallet.store.update_transfer(&delivered).unwrap(); + for status in [ + UsdtTransferStatus::Confirmed, + UsdtTransferStatus::BridgeFailed, + ] { + if status == UsdtTransferStatus::BridgeFailed { + delivered.status = UsdtTransferStatus::Bridging; + wallet.store.update_transfer(&delivered).unwrap(); + let mut stopped = message.clone(); + stopped["status"]["name"] = json!("APPLICATION_BURNED"); + chain + .state + .lock() + .unwrap() + .bridge_messages + .insert(pending.tx_hash.clone().unwrap(), json!({"data":[stopped]})); + assert_eq!(wallet.refresh_transfers().await.unwrap()[0].status, status); + } + let reads = { + let mut state = chain.state.lock().unwrap(); + state.chain = 42161; + let current_hash = state.block_hash(20000); + state.block_hashes.insert( + 20000, + if current_hash == B256::repeat_byte(0xac) { + B256::repeat_byte(0xab) + } else { + B256::repeat_byte(0xac) + }, + ); + state.receipt_reads + }; + sync_history_to_tip(&wallet).await; + assert_eq!(chain.state.lock().unwrap().receipt_reads, reads + 1); + assert_eq!(wallet.history().unwrap()[0].status, status); + } + drop(wallet); + let wallet = chain.wallet(&dir); + let requests = chain.state.lock().unwrap().bridge_requests.len(); + let failed = wallet.refresh_transfers().await.unwrap().remove(0); + assert_eq!(failed.status, UsdtTransferStatus::BridgeFailed); + assert_eq!(failed.fee, Some(300_123)); + assert!(failed.tx_hash.is_some() && failed.bridge_guid.is_some()); + assert_eq!(chain.state.lock().unwrap().bridge_requests.len(), requests); + let replacement_guid = B256::repeat_byte(0xad); + { + let mut state = chain.state.lock().unwrap(); + state.block_hashes.insert(20000, B256::repeat_byte(0xae)); + state.receipt_logs.as_mut().unwrap()[1]["topics"][1] = json!(replacement_guid); + } + sync_history_to_tip(&wallet).await; + let replacement = wallet.history().unwrap().remove(0); + assert_eq!(replacement.id, delivered.id); + assert_eq!( + replacement.bridge_guid, + Some(format!("{replacement_guid:#x}")) + ); + assert_eq!(replacement.status, UsdtTransferStatus::Bridging); + drop(wallet); + let restored_dir = tempfile::tempdir().unwrap(); + let restored = chain.wallet(&restored_dir); + sync_history_to_tip(&restored).await; + let recovered = restored.history().unwrap().remove(0); + assert_eq!(recovered.destination, UsdtDestination::Polygon); + assert_eq!( + recovered.bridge_guid, + Some(format!("{replacement_guid:#x}")) + ); + assert_eq!(recovered.fee, Some(300_123)); +} + +#[tokio::test] +async fn infrastructure_and_token_addresses_are_not_payment_recipients() { + let chain = MockChain::start().await; + let dir = tempfile::tempdir().unwrap(); + let wallet = chain.wallet(&dir); + for destination in [ + UsdtDestination::Arbitrum, + UsdtDestination::Ethereum, + UsdtDestination::Polygon, + UsdtDestination::Plasma, + UsdtDestination::Stable, + ] { + let mut recipients = vec![ + destination.token(), + account::ENTRY_POINT, + account::DELEGATE, + paymaster::PAYMASTER, + ]; + if destination == UsdtDestination::Arbitrum { + recipients.extend([types::OFT, types::BRIDGE_HELPER]); + } + for recipient in recipients { + assert!(matches!( + wallet + .quote_transfer(recipient.to_checksum(None), 1_000_000, destination) + .await, + Err(UsdtError::InvalidAddress) + )); + } + if let Some(eid) = destination.endpoint() { + assert_eq!(UsdtDestination::from_endpoint(eid), Some(destination)); + } + } +} + #[tokio::test] async fn settlement_uses_the_canonical_operation_outcome() { use alloy_primitives::{B256, U256}; @@ -2310,7 +3027,7 @@ async fn settlement_uses_the_canonical_operation_outcome() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -2370,7 +3087,7 @@ async fn unrecognized_operations_preserve_raw_debits_and_refunds() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); wallet @@ -2455,7 +3172,7 @@ async fn recent_execution_requires_the_expected_token_transfer() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -2513,7 +3230,7 @@ async fn execution_check_preserves_unmined_payments_and_throttling() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -2548,7 +3265,7 @@ async fn interrupted_nonce_recovery_restarts_on_a_changed_block() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -2590,7 +3307,7 @@ async fn history_reuses_receipts_only_while_their_block_remains_canonical() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); let sent = wallet @@ -2644,7 +3361,7 @@ async fn incoming_log_progress_does_not_hide_later_receipt_evidence() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); wallet @@ -2680,6 +3397,96 @@ async fn incoming_log_progress_does_not_hide_later_receipt_evidence() { assert_eq!(chain.state.lock().unwrap().receipt_reads, 1); } +#[tokio::test] +async fn bridge_history_retries_incomplete_receipt_enrichment() { + use alloy_primitives::{B256, U256}; + use alloy_sol_types::SolEvent; + use serde_json::json; + + for (dense, missing) in [ + (false, types::OFT), + (true, types::OFT), + (false, paymaster::PAYMASTER), + ] { + let chain = MockChain::start().await; + let dir = tempfile::tempdir().unwrap(); + let wallet = chain.wallet(&dir); + let quote = wallet + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Polygon) + .await + .unwrap(); + let sent = wallet + .send(quote.id, TEST_PHRASE.into(), None) + .await + .unwrap(); + let complete = { + let mut state = chain.state.lock().unwrap(); + state.mined = true; + state.tip += 3; + if dense { + state.oversized_block = Some(20000); + } + let mut logs = state.event_logs(); + let helper = transaction::BridgeHelper::LogSend { + sender: wallet.address, + oft: types::OFT, + amountLD: U256::from(1_000_000), + totalAmount: U256::from(1_000_500), + feeInToken: U256::from(500), + nativeFee: U256::from(1), + } + .encode_log_data(); + let oft = transaction::Oft::OFTSent { + guid: B256::repeat_byte(0x42), + dstEid: UsdtDestination::Polygon.endpoint().unwrap(), + fromAddress: types::BRIDGE_HELPER, + amountSentLD: U256::from(1_000_000), + amountReceivedLD: U256::from(1_000_000), + } + .encode_log_data(); + for (address, data) in [(types::BRIDGE_HELPER, helper), (types::OFT, oft)] { + logs.insert( + 0, + json!({"address":address,"topics":data.topics(),"data":data.data, + "transactionHash":B256::repeat_byte(7),"blockNumber":"0x4e20"}), + ); + } + for (index, log) in logs.iter_mut().enumerate() { + log["logIndex"] = json!(format!("0x{index:x}")); + } + state.receipt_logs = Some( + logs.iter() + .filter(|log| { + serde_json::from_value::(log["address"].clone()) + .unwrap() + != missing + }) + .cloned() + .collect(), + ); + logs + }; + sync_history_to_tip(&wallet).await; + let partial = wallet.history().unwrap().remove(0); + assert_eq!(partial.id, sent.id); + assert!(partial.bridge_guid.is_none() || partial.fee.is_none()); + assert_ne!(partial.status, UsdtTransferStatus::Pending); + assert!(wallet.store.pending_plan(&sent.id).unwrap().is_none()); + drop(wallet); + chain.state.lock().unwrap().receipt_logs = Some(complete); + let wallet = chain.wallet(&dir); + sync_history_to_tip(&wallet).await; + let enriched = wallet.history().unwrap().remove(0); + assert_eq!(enriched.id, sent.id); + assert_eq!( + enriched.bridge_guid, + Some(format!("{:#x}", B256::repeat_byte(0x42))) + ); + assert_eq!(enriched.fee, Some(623)); + assert_eq!(enriched.status, UsdtTransferStatus::Bridging); + } +} + #[tokio::test] async fn pending_payment_recovers_after_chain_head_retreat_and_restart() { for hide_logs in [false, true] { @@ -2688,7 +3495,7 @@ async fn pending_payment_recovers_after_chain_head_retreat_and_restart() { let wallet = chain.wallet(&dir); chain.state.lock().unwrap().tip = 20010; let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); chain.state.lock().unwrap().tip = 19998; @@ -2732,7 +3539,7 @@ async fn failed_restored_payment_preserves_attempted_recipient_and_fee() { let dir = tempfile::tempdir().unwrap(); let wallet = chain.wallet(&dir); let quote = wallet - .quote_transfer(RECIPIENT.into(), 1_000_000) + .quote_transfer(RECIPIENT.into(), 1_000_000, UsdtDestination::Arbitrum) .await .unwrap(); wallet diff --git a/src/modules/usdt/transaction.rs b/src/modules/usdt/transaction.rs index 8b4d3d6..5fe4d58 100644 --- a/src/modules/usdt/transaction.rs +++ b/src/modules/usdt/transaction.rs @@ -16,6 +16,47 @@ sol! { bytes paymasterAndData; bytes signature; } + #[derive(Debug)] + struct SendParam { + uint32 dstEid; + bytes32 to; + uint256 amountLD; + uint256 minAmountLD; + bytes extraOptions; + bytes composeMsg; + bytes oftCmd; + } + #[derive(Debug)] + struct MessagingFee { + uint256 nativeFee; + uint256 lzTokenFee; + } + struct OFTLimit { + uint256 minAmountLD; + uint256 maxAmountLD; + } + struct OFTFeeDetail { + int256 feeAmountLD; + string description; + } + struct OFTReceipt { + uint256 amountSentLD; + uint256 amountReceivedLD; + } + interface Oft { + function token() external view returns (address); + function peers(uint32 eid) external view returns (bytes32); + function quoteOFT(SendParam param) external view returns (OFTLimit limit, OFTFeeDetail[] fees, OFTReceipt receipt); + function quoteSend(SendParam param, bool payInLzToken) external view returns (MessagingFee fee); + event OFTSent(bytes32 indexed guid, uint32 dstEid, address indexed fromAddress, uint256 amountSentLD, uint256 amountReceivedLD); + } + interface BridgeHelper { + function token() external view returns (address); + function maxGas() external view returns (uint256); + function quoteSend(SendParam param, MessagingFee fee) external view returns (uint256 totalAmount); + function send(address oft, SendParam param, MessagingFee fee) external payable; + event LogSend(address indexed sender, address indexed oft, uint256 amountLD, uint256 nativeFee, uint256 feeInToken, uint256 totalAmount); + } interface EntryPoint { function getNonce(address sender, uint192 key) view returns (uint256); function handleOps(PackedOperation[] ops, address beneficiary); diff --git a/src/modules/usdt/types.rs b/src/modules/usdt/types.rs index 391dfc6..ccdafb0 100644 --- a/src/modules/usdt/types.rs +++ b/src/modules/usdt/types.rs @@ -3,6 +3,45 @@ use serde::{Deserialize, Serialize}; pub(super) const CHAIN_ID: u64 = 42161; pub(super) const TOKEN: Address = address!("Fd086bC7CD5C481DCC9C85ebE478A1C0b69FCbb9"); +pub(super) const OFT: Address = address!("14E4A1B13bf7F943c8ff7C51fb60FA964A298D92"); +pub(super) const BRIDGE_HELPER: Address = address!("a90f03c856D01F698E7071B393387cd75a8a319A"); + +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize, uniffi::Enum)] +pub enum UsdtDestination { + Stable, + Ethereum, + Arbitrum, + Polygon, + Plasma, +} + +impl UsdtDestination { + pub(super) fn token(self) -> Address { + match self { + Self::Arbitrum => TOKEN, + Self::Ethereum => address!("dAC17F958D2ee523a2206206994597C13D831ec7"), + Self::Polygon => address!("c2132D05D31c914a87C6611C10748AEb04B58e8F"), + Self::Plasma => address!("B8CE59FC3717ada4C02eaDF9682A9e934F625ebb"), + Self::Stable => address!("779Ded0c9e1022225f8E0630b35a9b54bE713736"), + } + } + + pub(super) fn from_endpoint(eid: u32) -> Option { + [Self::Ethereum, Self::Polygon, Self::Plasma, Self::Stable] + .into_iter() + .find(|d| d.endpoint() == Some(eid)) + } + + pub(super) fn endpoint(self) -> Option { + match self { + Self::Stable => Some(30396), + Self::Ethereum => Some(30101), + Self::Arbitrum => None, + Self::Polygon => Some(30109), + Self::Plasma => Some(30383), + } + } +} #[derive(Clone, Debug, PartialEq, Eq, uniffi::Record)] pub struct UsdtPaymentRequest { @@ -16,7 +55,9 @@ pub struct UsdtPaymentRequest { pub struct UsdtQuote { pub id: String, pub recipient: String, + pub destination: UsdtDestination, pub amount: u64, + pub received_amount: u64, pub maximum_fee: u64, pub expires_at: u64, } @@ -29,6 +70,12 @@ pub enum UsdtTransferStatus { Confirmed, /// Source payment failed or was proven not to have executed. Failed, + /// Source payment executed; destination delivery is pending. + Bridging, + /// Delivery is blocked or its message could not be recovered; it may still complete. + BridgeNeedsAttention, + /// Delivery was permanently stopped. This does not imply a refund of source funds or fees. + BridgeFailed, /// Another operation consumed the payment nonce. Replaced, } @@ -39,7 +86,9 @@ pub struct UsdtTransfer { /// Source transaction hash, absent until execution is observed. pub tx_hash: Option, pub user_operation_hash: Option, + pub bridge_guid: Option, pub recipient: String, + pub destination: UsdtDestination, pub amount: u64, pub received_amount: u64, pub fee: Option, diff --git a/src/modules/usdt/wallet.rs b/src/modules/usdt/wallet.rs index dca9588..56b6378 100644 --- a/src/modules/usdt/wallet.rs +++ b/src/modules/usdt/wallet.rs @@ -1,26 +1,34 @@ use super::{ account::{validate_delegation, ENTRY_POINT}, - amount::token_amount, + amount::{token_amount, with_margin}, keys::{derive_owner_key, parse_address}, paymaster::{Pimlico, PAYMASTER}, rpc::Rpc, store::{QuoteData, Store}, - transaction::{entry_point_event, event_data, EntryPoint, Erc20, Paymaster, Plan}, - types::{CHAIN_ID, TOKEN}, + transaction::{ + entry_point_event, event_data, BridgeHelper, EntryPoint, Erc20, Oft, Paymaster, Plan, + SendParam, + }, + types::{BRIDGE_HELPER, CHAIN_ID, OFT, TOKEN}, user_operation::Authorization, - UsdtError, UsdtQuote, UsdtTransfer, UsdtTransferStatus, + UsdtDestination, UsdtError, UsdtQuote, UsdtTransfer, UsdtTransferStatus, }; use alloy_primitives::{Address, Bytes, B256, U256}; use alloy_sol_types::{SolCall, SolEvent}; use serde_json::{json, Value}; -use std::sync::{atomic::AtomicU64, Arc}; -use tokio::sync::Mutex; +use std::collections::HashMap; +use std::sync::{ + atomic::{AtomicU64, AtomicUsize, Ordering}, + Arc, +}; +use tokio::{sync::Mutex, time::Instant}; const RECENT_EXECUTION_BLOCKS: u64 = 64; const RECENT_EXECUTION_BUDGET: std::time::Duration = std::time::Duration::from_secs(5); const NONCE_RECOVERY_BUDGET: std::time::Duration = std::time::Duration::from_secs(20); const EXPIRY_SEARCH_BLOCKS: u64 = 4096; const QUOTE_LIFETIME_SECONDS: u64 = 120; +const BRIDGE_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(60); #[derive(uniffi::Object)] pub struct UsdtWallet { @@ -29,6 +37,8 @@ pub struct UsdtWallet { pub(super) paymaster: Pimlico, pub(super) store: Store, operation: Mutex<()>, + bridge_poll_offset: AtomicUsize, + bridge_retry_after: Mutex>, pub(super) history_range_limit: AtomicU64, } @@ -57,6 +67,8 @@ impl UsdtWallet { paymaster, store, operation: Mutex::new(()), + bridge_poll_offset: AtomicUsize::new(0), + bridge_retry_after: Mutex::new(HashMap::new()), history_range_limit: AtomicU64::new(super::history::MAX_LOG_RANGE), })) } @@ -82,39 +94,30 @@ impl UsdtWallet { &self, recipient: String, amount: u64, + destination: UsdtDestination, ) -> Result { if amount == 0 { return Err(UsdtError::InvalidAmount); } let recipient = parse_address(recipient.trim())?; - if [ - self.address, - TOKEN, - ENTRY_POINT, - PAYMASTER, - super::account::DELEGATE, - ] - .contains(&recipient) + if recipient == self.address + || recipient == destination.token() + || [ENTRY_POINT, PAYMASTER, super::account::DELEGATE].contains(&recipient) + || (destination == UsdtDestination::Arbitrum + && [OFT, BRIDGE_HELPER].contains(&recipient)) { return Err(UsdtError::InvalidAddress); } self.store.require_no_pending()?; self.rpc.verify_chain().await?; self.require_balance(amount, 0).await?; - let calls = vec![( - TOKEN, - Erc20::transferCall { - recipient, - amount: U256::from(amount), - } - .abi_encode() - .into(), - )]; + let (calls, received_amount, bridge_fee) = + self.transfer_calls(recipient, amount, destination).await?; let nonce = self.nonce("latest").await?; let authorization = self.authorization().await?; let created_block = self.block_number().await?; let timestamp = self.block_timestamp(created_block).await?; - let (operation, maximum_fee, operation_expires_at) = self + let (operation, gas_fee, operation_expires_at) = self .paymaster .prepare(self.address, nonce, authorization, &calls, timestamp) .await?; @@ -123,11 +126,16 @@ impl UsdtWallet { .saturating_sub(timestamp) .min(QUOTE_LIFETIME_SECONDS), ); + let maximum_fee = gas_fee + .checked_add(bridge_fee) + .ok_or(UsdtError::InvalidAmount)?; self.require_balance(amount, maximum_fee).await?; let quote = UsdtQuote { id: uuid::Uuid::new_v4().to_string(), recipient: recipient.to_checksum(None), + destination, amount, + received_amount, maximum_fee, expires_at, }; @@ -172,6 +180,7 @@ impl UsdtWallet { } self.require_balance(data.quote.amount, data.quote.maximum_fee) .await?; + self.validate_bridge(&data.plan).await?; self.paymaster.validate_gas(&data.plan.operation).await?; let block = self.block_number().await?; if self.block_timestamp(block).await?.saturating_add(5) >= data.plan.expires_at { @@ -192,9 +201,11 @@ impl UsdtWallet { id: quote_id, tx_hash: None, user_operation_hash: Some(format!("{hash:#x}")), + bridge_guid: None, recipient: data.quote.recipient, + destination: data.quote.destination, amount: data.quote.amount, - received_amount: data.quote.amount, + received_amount: data.quote.received_amount, fee: None, is_incoming: false, status: UsdtTransferStatus::Pending, @@ -228,6 +239,9 @@ impl UsdtWallet { let Some(mut transfer) = self.store.transfer(&id)? else { return Ok(None); }; + if transfer.destination != UsdtDestination::Arbitrum { + return Ok(Some(transfer)); + } let Some(plan) = self.store.pending_plan(&id)? else { return Ok(Some(transfer)); }; @@ -279,15 +293,23 @@ impl UsdtWallet { /// Reconciles pending execution using chain proofs and may rebroadcast the identical signed operation. pub async fn refresh_transfers(&self) -> Result, UsdtError> { + let pending = self.refresh_pending_transfers().await; + self.refresh_bridges(&self.store.awaiting_delivery()?) + .await?; + pending?; + self.history() + } +} + +impl UsdtWallet { + async fn refresh_pending_transfers(&self) -> Result<(), UsdtError> { let _guard = self.operation.lock().await; if let Some((mut transfer, plan)) = self.store.pending_operation()? { self.recover_pending(&mut transfer, &plan).await?; } - self.history() + Ok(()) } -} -impl UsdtWallet { async fn recover_pending( &self, transfer: &mut UsdtTransfer, @@ -324,7 +346,7 @@ impl UsdtWallet { if self.block_timestamp(confirmed_tip).await? > plan.expires_at { transfer.mark_unexecuted(UsdtTransferStatus::Failed); self.store.update_transfer(transfer)?; - } else { + } else if self.validate_bridge(plan).await.is_ok() { let _ = self.broadcast(plan, hash).await; } return Ok(()); @@ -423,6 +445,72 @@ impl UsdtWallet { self.store.update_transfer(transfer) } + async fn refresh_bridges(&self, transfers: &[UsdtTransfer]) -> Result<(), UsdtError> { + let mut retry_after = self.bridge_retry_after.lock().await; + retry_after.retain(|_, deadline| *deadline > Instant::now()); + let mut bridges: Vec<_> = transfers + .iter() + .filter(|transfer| !retry_after.contains_key(&transfer.id)) + .collect(); + drop(retry_after); + if bridges.is_empty() { + return Ok(()); + } + let batch = [0, 1, 2]; + let offset = self + .bridge_poll_offset + .fetch_add(batch.len(), Ordering::Relaxed) + % bridges.len(); + bridges.rotate_left(offset); + let bridges = &bridges; + let check = |index: usize| async move { + let transfer = bridges.get(index).copied()?; + Some(( + transfer, + tokio::time::timeout( + std::time::Duration::from_secs(10), + self.rpc.bridge_status(transfer), + ) + .await, + )) + }; + let (first, second, third) = + tokio::join!(check(batch[0]), check(batch[1]), check(batch[2])); + for (previous, result) in [first, second, third].into_iter().flatten() { + if !matches!(result, Ok(Ok(status)) if status != UsdtTransferStatus::BridgeNeedsAttention) + { + self.bridge_retry_after + .lock() + .await + .insert(previous.id.clone(), Instant::now() + BRIDGE_RETRY_DELAY); + } + match result { + Ok(Ok(status)) if status != previous.status => { + let _guard = self.operation.lock().await; + let Some(mut current) = self.store.transfer(&previous.id)? else { + continue; + }; + if current + .tx_hash + .as_deref() + .zip(previous.tx_hash.as_deref()) + .is_some_and(|(a, b)| a.eq_ignore_ascii_case(b)) + && current.bridge_guid == previous.bridge_guid + && current.status == previous.status + { + current.status = status; + self.store.update_transfer(¤t)?; + } + } + Ok(Ok(_)) => {} + _ => log::warn!( + "USDT bridge delivery lookup unavailable; retaining last known status" + ), + } + } + Ok(()) + } + pub(super) async fn block_number(&self) -> Result { u64::try_from(self.rpc.call::("eth_blockNumber", json!([])).await?) .map_err(|_| UsdtError::InvalidResponse) @@ -572,6 +660,7 @@ impl UsdtWallet { } let mut transfer_proven = false; let mut gas_fee = None; + let mut bridge_fee = None; for log in logs { let address: Address = serde_json::from_value(log["address"].clone())?; let data = event_data(log)?; @@ -593,19 +682,235 @@ impl UsdtWallet { } } } + if event.success && address == BRIDGE_HELPER { + if let Ok(event) = BridgeHelper::LogSend::decode_log_data(&data) { + if event.sender == self.address + && event.oft == OFT + && event.amountLD == U256::from(transfer.amount) + { + bridge_fee = Some(token_amount(event.feeInToken)?); + } + } + } + if event.success && address == OFT { + if let Ok(event) = Oft::OFTSent::decode_log_data(&data) { + if event.fromAddress == BRIDGE_HELPER + && transfer.destination.endpoint() == Some(event.dstEid) + && event.amountSentLD == U256::from(transfer.amount) + { + transfer.bridge_guid = Some(format!("{:#x}", event.guid)); + transfer.received_amount = token_amount(event.amountReceivedLD)?; + } + } + } } - if event.success && !transfer_proven { + if event.success && transfer.destination == UsdtDestination::Arbitrum && !transfer_proven { return Err(UsdtError::InvalidResponse); } - transfer.fee = gas_fee; + transfer.fee = if event.success && transfer.destination != UsdtDestination::Arbitrum { + gas_fee + .zip(bridge_fee) + .and_then(|(gas, bridge)| gas.checked_add(bridge)) + } else { + gas_fee + }; transfer.status = if !event.success { transfer.received_amount = 0; UsdtTransferStatus::Failed - } else { + } else if transfer.destination == UsdtDestination::Arbitrum { UsdtTransferStatus::Confirmed + } else if transfer.bridge_guid.is_some() { + UsdtTransferStatus::Bridging + } else { + UsdtTransferStatus::BridgeNeedsAttention }; Ok(()) } + async fn validate_bridge(&self, plan: &Plan) -> Result<(), UsdtError> { + let calls = super::history::decode_calls(&plan.operation.call_data)?; + let Some((_, data)) = calls.iter().find(|(target, _)| *target == BRIDGE_HELPER) else { + return Ok(()); + }; + let send = + BridgeHelper::sendCall::abi_decode(data).map_err(|_| UsdtError::InvalidResponse)?; + let requote = |error| match error { + UsdtError::UnsupportedRoute => UsdtError::QuoteExpired, + error => error, + }; + let required = self + .rpc + .contract( + OFT, + Oft::quoteSendCall { + param: send.param.clone(), + payInLzToken: false, + }, + ) + .await + .map_err(requote)?; + if !required.lzTokenFee.is_zero() || required.nativeFee > send.fee.nativeFee { + return Err(UsdtError::QuoteExpired); + } + if self.rpc.balance(BRIDGE_HELPER).await? < send.fee.nativeFee { + return Err(UsdtError::UnsupportedRoute); + } + let total = self + .rpc + .contract( + BRIDGE_HELPER, + BridgeHelper::quoteSendCall { + param: send.param, + fee: send.fee, + }, + ) + .await + .map_err(requote)?; + let allowance = calls + .iter() + .filter(|(target, _)| *target == TOKEN) + .filter_map(|(_, data)| Erc20::approveCall::abi_decode(data).ok()) + .find(|call| call.spender == BRIDGE_HELPER) + .ok_or(UsdtError::InvalidResponse)?; + if total > allowance.amount { + return Err(UsdtError::QuoteExpired); + } + Ok(()) + } + + async fn transfer_calls( + &self, + recipient: Address, + amount: u64, + destination: UsdtDestination, + ) -> Result<(Vec<(Address, Bytes)>, u64, u64), UsdtError> { + let Some(eid) = destination.endpoint() else { + return Ok(( + vec![( + TOKEN, + Erc20::transferCall { + recipient, + amount: U256::from(amount), + } + .abi_encode() + .into(), + )], + amount, + 0, + )); + }; + let token = self.rpc.contract(OFT, Oft::tokenCall {}).await?; + let helper_token = self + .rpc + .contract(BRIDGE_HELPER, BridgeHelper::tokenCall {}) + .await?; + let peer = self.rpc.contract(OFT, Oft::peersCall { eid }).await?; + if token != TOKEN || helper_token != TOKEN || peer.is_zero() { + return Err(UsdtError::UnsupportedRoute); + } + if recipient.into_word() == peer { + return Err(UsdtError::InvalidAddress); + } + let mut param = SendParam { + dstEid: eid, + to: recipient.into_word(), + amountLD: U256::from(amount), + minAmountLD: U256::ZERO, + extraOptions: Bytes::new(), + composeMsg: Bytes::new(), + oftCmd: Bytes::new(), + }; + let oft = self + .rpc + .contract( + OFT, + Oft::quoteOFTCall { + param: param.clone(), + }, + ) + .await?; + if U256::from(amount) < oft.limit.minAmountLD + || U256::from(amount) > oft.limit.maxAmountLD + || oft.receipt.amountSentLD != U256::from(amount) + || oft.receipt.amountReceivedLD.is_zero() + || oft.receipt.amountReceivedLD > U256::from(amount) + { + return Err(UsdtError::InvalidAmount); + } + param.minAmountLD = oft.receipt.amountReceivedLD; + let mut fee = self + .rpc + .contract( + OFT, + Oft::quoteSendCall { + param: param.clone(), + payInLzToken: false, + }, + ) + .await?; + // Native headroom is quoted into the approved USDT maximum. + fee.nativeFee = with_margin(fee.nativeFee, 10)?; + let maximum_native = self + .rpc + .contract(BRIDGE_HELPER, BridgeHelper::maxGasCall {}) + .await?; + if !fee.lzTokenFee.is_zero() + || fee.nativeFee > maximum_native + || self.rpc.balance(BRIDGE_HELPER).await? < fee.nativeFee + { + return Err(UsdtError::UnsupportedRoute); + } + let total = self + .rpc + .contract( + BRIDGE_HELPER, + BridgeHelper::quoteSendCall { + param: param.clone(), + fee: fee.clone(), + }, + ) + .await?; + let token_fee = total + .checked_sub(U256::from(amount)) + .ok_or(UsdtError::InvalidResponse)?; + let token_fee = with_margin(token_fee, 20)?; + let approval = U256::from(amount) + .checked_add(token_fee) + .ok_or(UsdtError::InvalidResponse)?; + Ok(( + vec![ + ( + TOKEN, + Erc20::approveCall { + spender: BRIDGE_HELPER, + amount: approval, + } + .abi_encode() + .into(), + ), + ( + BRIDGE_HELPER, + BridgeHelper::sendCall { + oft: OFT, + param, + fee, + } + .abi_encode() + .into(), + ), + ( + TOKEN, + Erc20::approveCall { + spender: BRIDGE_HELPER, + amount: U256::ZERO, + } + .abi_encode() + .into(), + ), + ], + token_amount(oft.receipt.amountReceivedLD)?, + token_amount(token_fee)?, + )) + } } pub(super) fn now() -> u64 { diff --git a/tests/usdt-fork/provider.mjs b/tests/usdt-fork/provider.mjs index 65ad95b..689fbd8 100644 --- a/tests/usdt-fork/provider.mjs +++ b/tests/usdt-fork/provider.mjs @@ -109,7 +109,12 @@ function pack(op) { signature: op.signature, }; } +let drainHelperBeforeBroadcast = false; async function dispatch(method, params) { + if (method === 'test_drainHelperBeforeNextBroadcast') { + drainHelperBeforeBroadcast = true; + return true; + } if (method === 'eth_estimateUserOperationGas' || method === 'eth_sendUserOperation') { const op = params[0]; assert.equal(op.factory, '0x7702'); @@ -161,6 +166,10 @@ async function dispatch(method, params) { return { paymaster: pmAddress, paymasterData: concat([unsigned, signature]), ...limits }; } if (method === 'eth_sendUserOperation') { + if (drainHelperBeforeBroadcast) { + drainHelperBeforeBroadcast = false; + await rpc.send('anvil_setBalance', ['0xa90f03c856d01f698e7071b393387cd75a8a319a', '0x0']); + } const op = pack(params[0]); const hash = await rpc.send('eth_call', [ { to: entryAddress, data: entry.interface.encodeFunctionData('getUserOpHash', [op]) }, @@ -192,6 +201,7 @@ async function dispatch(method, params) { ); return hash; } + if (method === 'bitkit_getBridgeMessages') return { data: [] }; return rpc.send(method, params); } const server = createServer(async (request, response) => {