use std::sync::Arc; use anyhow::{Context, Result}; use bridge_connector_common::result::BridgeSdkError; use near_sdk::json_types::U64; use serde_json::Value; use tracing::{info, warn}; use near_bridge_client::TransactionOptions; use near_jsonrpc_client::{JsonRpcClient, errors::JsonRpcError}; use near_primitives::{hash::CryptoHash, types::AccountId}; use near_rpc_client::NearRpcError; use solana_client::rpc_request::RpcResponseErrorData; use solana_rpc_client_api::{client_error::ErrorKind, request::RpcError}; use solana_sdk::{instruction::InstructionError, pubkey::Pubkey, transaction::TransactionError}; use omni_connector::OmniConnector; use omni_types::{ChainKind, FastTransfer, OmniAddress, TransferId, near_events::OmniBridgeEvent}; use crate::{ config, utils, utils::pending_transactions::PendingTransaction, workers::{PAUSED_ERROR, RetryableEvent}, }; use super::{EventAction, Transfer}; #[derive(Debug, serde::Serialize, serde::Deserialize)] pub struct UnverifiedTrasfer { pub tx_hash: CryptoHash, pub signer: AccountId, pub specific_errors: Option>, pub original_key: String, pub original_event: Value, } #[derive(Debug, serde::Deserialize)] enum UTXOChainMsg { MaxGasFee(U64), } pub async fn process_transfer_event( config: &config::Config, redis_connection_manager: &mut redis::aio::ConnectionManager, key: String, omni_connector: Arc, signer: AccountId, transfer: Transfer, near_nonce: Arc, ) -> Result { let transfer_message = match transfer { Transfer::Near { ref transfer_message, } => transfer_message.clone(), Transfer::Utxo { ref new_transfer_id, .. } => { let Ok(new_transfer_id) = new_transfer_id.try_into() else { warn!("Failed to build TransferId from: {new_transfer_id:?}"); return Ok(EventAction::Retry); }; let Ok(transfer_message) = omni_connector .near_get_transfer_message(new_transfer_id) .await else { warn!("Failed to get transfer message for UTXO transfer: {new_transfer_id:?}"); return Ok(EventAction::Retry); }; transfer_message } _ => { anyhow::bail!("Expected Transfer::Near or Transfer::Utxo variant, got: {transfer:?}"); } }; info!("Trying to process TransferMessage on NEAR"); match omni_connector .is_transfer_finalised( Some(transfer_message.get_origin_chain()), transfer_message.get_destination_chain(), transfer_message.destination_nonce, ) .await { Ok(true) => anyhow::bail!("Transfer is already finalised: {transfer_message:?}"), Ok(false) => {} Err(err) => { warn!("Failed to check if transfer is finalised: {err:?}"); return Ok(EventAction::Retry); } } if config.is_bridge_api_enabled() && !config .near .sign_without_checking_fee .as_ref() .is_some_and(|list| list.contains(&transfer_message.sender)) { let Ok(needed_fee) = utils::bridge_api::TransferFee::get_transfer_fee( config, &transfer_message.sender, &transfer_message.recipient, &transfer_message.token, ) .await else { warn!("Failed to get transfer fee for transfer: {transfer_message:?}"); return Ok(EventAction::Retry); }; if let Some(event_action) = needed_fee .check_fee( config, redis_connection_manager, &transfer_message, transfer_message.get_transfer_id(), &transfer_message.fee, ) .await { return Ok(event_action); } } let nonce = near_nonce .reserve_nonce() .await .context("Failed to reserve nonce for near transaction")?; match omni_connector .near_sign_transfer( TransferId { origin_chain: transfer_message.sender.get_chain(), origin_nonce: transfer_message.origin_nonce, }, Some(signer.clone()), Some(transfer_message.fee.clone()), TransactionOptions { nonce: Some(nonce), wait_until: near_primitives::views::TxExecutionStatus::Included, wait_final_outcome_timeout_sec: None, }, ) .await { Ok(tx_hash) => { let Ok(serialized_event) = serde_json::to_value(&transfer) else { warn!("Failed to serialize transfer: {transfer:?}"); return Ok(EventAction::Remove); }; utils::redis::add_event( config, redis_connection_manager, utils::redis::EVENTS, tx_hash.to_string(), RetryableEvent::new(UnverifiedTrasfer { tx_hash, signer, specific_errors: Some(vec![ "Signature request has already been submitted. Please try again later." .to_string(), "Signature request has timed out.".to_string(), "Attached deposit is lower than required".to_string(), "Exceeded the prepaid gas.".to_string(), ]), original_key: key, original_event: serialized_event, }), ) .await; info!("Signed transfer: {tx_hash:?}"); Ok(EventAction::Remove) } Err(err) => { if let BridgeSdkError::NearRpcError(near_rpc_error) = err { match near_rpc_error { NearRpcError::NonceError | NearRpcError::FinalizationError | NearRpcError::RpcBroadcastTxAsyncError(_) | NearRpcError::RpcQueryError(JsonRpcError::TransportError(_)) | NearRpcError::RpcTransactionError(JsonRpcError::TransportError(_)) => { warn!( "Failed to sign transfer ({}), retrying: {near_rpc_error:?}", transfer_message.origin_nonce ); return Ok(EventAction::Retry); } _ => { anyhow::bail!( "Failed to sign transfer ({}): {near_rpc_error:?}", transfer_message.origin_nonce ); } }; } anyhow::bail!( "Failed to sign transfer ({}): {err:?}", transfer_message.origin_nonce ); } } } pub async fn process_transfer_to_utxo_event( config: &config::Config, redis_connection_manager: &mut redis::aio::ConnectionManager, key: String, omni_connector: Arc, signer: AccountId, transfer: Transfer, near_nonce: Arc, ) -> Result { let Transfer::Near { ref transfer_message, } = transfer else { anyhow::bail!("Expected NearTransferWithTimestamp, got: {transfer:?}"); }; info!("Trying to process UtxoTransferMessage on NEAR"); let Some(recipient) = transfer_message.recipient.get_utxo_address() else { anyhow::bail!( "Expected UTXO recipient address, got: {:?}", transfer_message.recipient ); }; let nonce = near_nonce .reserve_nonce() .await .context("Failed to reserve nonce for near transaction")?; match omni_connector .near_submit_btc_transfer( transfer_message.recipient.get_chain(), recipient, transfer_message.amount.0 - transfer_message.fee.fee.0, None, TransferId { origin_chain: transfer_message.sender.get_chain(), origin_nonce: transfer_message.origin_nonce, }, TransactionOptions { nonce: Some(nonce), wait_until: near_primitives::views::TxExecutionStatus::Included, wait_final_outcome_timeout_sec: None, }, serde_json::from_str::(&transfer_message.msg) .map(|msg| match msg { UTXOChainMsg::MaxGasFee(max_fee) => max_fee.0, }) .ok(), ) .await { Ok(tx_hash) => { info!( "Submitted {:?} transfer: {tx_hash:?}", transfer_message.recipient.get_chain() ); let Ok(serialized_event) = serde_json::to_value(&transfer) else { warn!("Failed to serialize transfer: {transfer:?}"); return Ok(EventAction::Remove); }; utils::redis::add_event( config, redis_connection_manager, utils::redis::EVENTS, tx_hash.to_string(), RetryableEvent::new(UnverifiedTrasfer { tx_hash, signer, specific_errors: Some(vec![ "not exist".to_string(), "Previous btc tx has not been signed".to_string(), ]), original_key: key, original_event: serialized_event, }), ) .await; Ok(EventAction::Remove) } Err(err) => { if let BridgeSdkError::NearRpcError(near_rpc_error) = err { match near_rpc_error { NearRpcError::NonceError | NearRpcError::FinalizationError | NearRpcError::RpcBroadcastTxAsyncError(_) | NearRpcError::RpcQueryError(JsonRpcError::TransportError(_)) | NearRpcError::RpcTransactionError(JsonRpcError::TransportError(_)) => { warn!( "Failed to submit {:?} transfer ({}), retrying: {near_rpc_error:?}", transfer_message.recipient.get_chain(), transfer_message.origin_nonce ); return Ok(EventAction::Retry); } _ => { anyhow::bail!( "Failed to submit {:?} transfer ({}): {near_rpc_error:?}", transfer_message.recipient.get_chain(), transfer_message.origin_nonce ); } }; } else if let BridgeSdkError::InsufficientUTXOBalance = err { warn!( "Insufficient UTXO balance for {:?} transfer ({}), retrying", transfer_message.recipient.get_chain(), transfer_message.origin_nonce ); return Ok(EventAction::Retry); } else if let BridgeSdkError::InsufficientUTXOGasFee(err) = err { warn!( "Gas fee is too large for {:?} transfer ({}): {err}, retrying", transfer_message.recipient.get_chain(), transfer_message.origin_nonce ); return Ok(EventAction::Retry); } else if let BridgeSdkError::UtxoClientError(ref msg) = err { if msg == "Failed to estimate fee_rate" { warn!( "Failed to estimate fee_rate for {:?} transfer ({}), retrying", transfer_message.recipient.get_chain(), transfer_message.origin_nonce ); return Ok(EventAction::Retry); } } anyhow::bail!( "Failed to submit {:?} transfer ({}): {err:?}", transfer_message.recipient.get_chain(), transfer_message.origin_nonce ); } } } pub async fn process_sign_transfer_event( config: &config::Config, redis_connection_manager: &mut redis::aio::ConnectionManager, omni_connector: Arc, signer: AccountId, omni_bridge_event: OmniBridgeEvent, evm_nonces: Arc, ) -> Result { let OmniBridgeEvent::SignTransferEvent { message_payload, .. } = &omni_bridge_event else { anyhow::bail!("Expected SignTransferEvent, got: {omni_bridge_event:?}"); }; info!("Trying to process SignTransferEvent log on NEAR"); if message_payload.fee_recipient != Some(signer) { anyhow::bail!("Fee recipient mismatch"); } match omni_connector .is_transfer_finalised( None, message_payload.recipient.get_chain(), message_payload.destination_nonce, ) .await { Ok(true) => anyhow::bail!( "Transfer is already finalised: {:?}", message_payload.transfer_id ), Ok(false) => {} Err(err) => { warn!("Failed to check if transfer is finalised: {err:?}"); return Ok(EventAction::Retry); } } if config.is_bridge_api_enabled() { let transfer_message = match omni_connector .near_get_transfer_message(message_payload.transfer_id) .await { Ok(transfer_message) => transfer_message, Err(err) => { if err.to_string().contains("The transfer does not exist") { anyhow::bail!( "Transfer does not exist: {:?} (probably fee is 0 or transfer was already finalized)", message_payload.transfer_id ); } warn!( "Failed to get transfer message: {:?}", message_payload.transfer_id ); return Ok(EventAction::Retry); } }; let Ok(needed_fee) = utils::bridge_api::TransferFee::get_transfer_fee( config, &transfer_message.sender, &transfer_message.recipient, &transfer_message.token, ) .await else { warn!("Failed to get transfer fee for transfer: {transfer_message:?}"); return Ok(EventAction::Retry); }; if let Some(event_action) = needed_fee .check_fee( config, redis_connection_manager, &transfer_message, transfer_message.get_transfer_id(), &transfer_message.fee, ) .await { return Ok(event_action); } } let chain_kind = message_payload.recipient.get_chain(); let (fin_transfer_args, evm_nonce) = match chain_kind { ChainKind::Near => { anyhow::bail!("Near to Near transfers are not supported yet"); } ChainKind::Eth | ChainKind::Base | ChainKind::Arb | ChainKind::Bnb | ChainKind::Pol => { let nonce = evm_nonces .reserve_nonce(chain_kind) .await .context("Failed to reserve nonce for evm transaction")?; ( omni_connector::FinTransferArgs::EvmFinTransfer { chain_kind, event: omni_bridge_event.clone(), tx_nonce: Some(nonce.into()), }, Some(nonce), ) } ChainKind::Sol => { let OmniAddress::Sol(token) = message_payload.token_address.clone() else { anyhow::bail!( "Expected Sol token address, got: {:?}", message_payload.token_address ); }; ( omni_connector::FinTransferArgs::SolanaFinTransfer { event: omni_bridge_event.clone(), solana_token: Pubkey::new_from_array(token.0), }, None, ) } ChainKind::Btc | ChainKind::Zcash => { anyhow::bail!("Finishing BTC/ZEC transfers is not supported"); } }; match omni_connector.fin_transfer(fin_transfer_args).await { Ok(tx_hash) => { info!("Finalized deposit: {tx_hash}"); if let Some(nonce) = evm_nonce { if config.is_fee_bumping_enabled(chain_kind) { if let Err(err) = store_pending_transaction( config, redis_connection_manager, chain_kind, &tx_hash, nonce, omni_bridge_event, ) .await { warn!("Failed to store pending transaction {tx_hash}: {err:?}"); } } } Ok(EventAction::Remove) } Err(err) => { if let BridgeSdkError::EvmGasEstimateError(err) = err { let Some(evm) = (match chain_kind { ChainKind::Eth => &config.eth, ChainKind::Base => &config.base, ChainKind::Arb => &config.arb, ChainKind::Bnb => &config.bnb, ChainKind::Pol => &config.pol, ChainKind::Near | ChainKind::Sol | ChainKind::Btc | ChainKind::Zcash => { anyhow::bail!( "Failed to finalize deposit (unexpected: failed to get evm config): {err}" ); } }) else { anyhow::bail!( "Failed to finalize deposit (unexpected: config for {chain_kind:?} is not accessible): {err}" ); }; if evm .error_selectors_to_remove .iter() .any(|selector| err.contains(selector)) { anyhow::bail!( "Failed to finalize deposit: {err}. Found selector from the list of non-retryable errors in the config" ); } warn!("Failed to finalize deposit, retrying: {err}"); return Ok(EventAction::Retry); } if let BridgeSdkError::SolanaRpcError(ref client_error) = err { if let ErrorKind::RpcError(RpcError::RpcResponseError { data: RpcResponseErrorData::SendTransactionPreflightFailure(ref result), .. }) = client_error.kind { if let Some(TransactionError::InstructionError( _, InstructionError::Custom(error_code), )) = result.err { if error_code == PAUSED_ERROR { warn!("Solana bridge is paused"); return Ok(EventAction::Retry); } anyhow::bail!("Failed to finalize deposit: {err}"); } } } Ok(EventAction::Retry) } } } pub async fn process_unverified_transfer_event( config: &config::Config, redis_connection_manager: &mut redis::aio::ConnectionManager, jsonrpc_client: JsonRpcClient, unverified_event: UnverifiedTrasfer, ) { utils::redis::remove_event( config, redis_connection_manager, utils::redis::EVENTS, unverified_event.tx_hash.to_string(), ) .await; if !utils::near::is_tx_successful( &jsonrpc_client, unverified_event.tx_hash, unverified_event.signer, unverified_event.specific_errors, ) .await { utils::redis::add_event( config, redis_connection_manager, utils::redis::EVENTS, unverified_event.original_key, RetryableEvent::new(unverified_event.original_event), ) .await; } } pub async fn initiate_fast_transfer( fast_connector: Arc, transfer: Transfer, near_omni_nonce: Arc, ) -> Result { let Ok(near_bridge_client) = fast_connector.near_bridge_client() else { anyhow::bail!("Near bridge client is not configured"); }; let Transfer::Fast { block_number, tx_hash, token, amount, transfer_id, recipient, fee, msg, storage_deposit_amount, safe_confirmations, } = transfer.clone() else { anyhow::bail!("Expected FastTransferEvent, got: {transfer:?}"); }; // TODO: Fast transfer to other chain increases origin nonce by one, so regular relayer won't // be able to finalize it with a normal sign transfer. We need to catch and sign // `FastTransferEvent`. This will be possible once bridge-indexer will track these events // Related PR: https://github.com/Near-One/bridge-indexer-rs/pull/195 if recipient.get_chain() != ChainKind::Near { anyhow::bail!( "Fast transfer is supported only for transfers to NEAR for now, got: {:?}", recipient.get_chain() ); } info!("Trying to initiate FastTransfer on NEAR"); let Ok(token_id) = utils::storage::get_token_id(&fast_connector, transfer_id.origin_chain, &token).await else { warn!("Failed to get token id for transfer: {transfer_id:?}"); return Ok(EventAction::Retry); }; let fast_transfer = FastTransfer { transfer_id: transfer_id.into(), token_id: token_id.clone(), amount, fee: fee.clone(), recipient: recipient.clone(), msg: msg.clone(), }; match fast_connector.near_is_transfer_finalised(transfer_id).await { Ok(true) => anyhow::bail!("Transfer is already finalised: {transfer:?}"), Ok(false) => {} Err(err) => { warn!("Failed to check if transfer is finalised: {err:?}"); return Ok(EventAction::Retry); } } match fast_connector .near_get_fast_transfer_status(fast_transfer.id()) .await { Ok(Some(_)) => anyhow::bail!("Fast transfer is already finalised: {transfer:?}"), Ok(None) => {} Err(err) => { warn!("Failed to check if fast transfer is finalised: {err:?}"); return Ok(EventAction::Retry); } } let Ok(last_finalized_block_number) = fast_connector .evm_get_last_block_number(transfer_id.origin_chain) .await else { warn!("Failed to get last finalized block number for EVM chain"); return Ok(EventAction::Retry); }; let current_confirmations = last_finalized_block_number.saturating_sub(block_number); if current_confirmations < safe_confirmations { warn!( "Fast transfer block number ({block_number}) is not finalized yet, waiting for more confirmations. Current confirmations: {current_confirmations}", ); return Ok(EventAction::Retry); } let mut nonce = Some( near_omni_nonce .reserve_nonce() .await .context("Failed to reserve nonce for near transaction")?, ); let Ok(required_balance) = near_bridge_client .get_required_balance_for_fast_fin_transfer() .await else { warn!("Failed to get required balance for fast transfer"); return Ok(EventAction::Retry); }; match near_bridge_client .deposit_storage_if_required( required_balance + storage_deposit_amount.unwrap_or(0.into()).0, TransactionOptions { nonce, wait_until: near_primitives::views::TxExecutionStatus::Final, wait_final_outcome_timeout_sec: None, }, ) .await { Ok(true) => { nonce = Some( near_omni_nonce .reserve_nonce() .await .context("Failed to reserve nonce for near transaction")?, ); } Ok(false) => {} Err(err) => { warn!("Failed to deposit storage for fast transfer: {err:?}"); return Ok(EventAction::Retry); } } match fast_connector .near_fast_transfer( transfer_id.origin_chain, tx_hash, storage_deposit_amount.map(|storage_deposit_amount| storage_deposit_amount.0), TransactionOptions { nonce, wait_until: near_primitives::views::TxExecutionStatus::Included, wait_final_outcome_timeout_sec: None, }, ) .await { Ok(tx_hash) => { info!("Fast transfer initiated successfully: {tx_hash:?}"); Ok(EventAction::Remove) } Err(err) => { if let BridgeSdkError::NearRpcError(near_rpc_error) = err { match near_rpc_error { NearRpcError::NonceError | NearRpcError::FinalizationError | NearRpcError::RpcBroadcastTxAsyncError(_) | NearRpcError::RpcQueryError(JsonRpcError::TransportError(_)) | NearRpcError::RpcTransactionError(JsonRpcError::TransportError(_)) => { warn!("Failed to initiate fast transfer, retrying: {near_rpc_error:?}"); return Ok(EventAction::Retry); } _ => { anyhow::bail!("Failed to initiate fast transfer: {near_rpc_error:?}"); } }; } anyhow::bail!("Failed to initiate fast transfer: {err:?}"); } } } async fn store_pending_transaction( config: &config::Config, redis_connection_manager: &mut redis::aio::ConnectionManager, chain_kind: ChainKind, tx_hash: &str, nonce: u64, omni_bridge_event: OmniBridgeEvent, ) -> Result<()> { let pending_tx = PendingTransaction::new(tx_hash.to_string(), nonce, chain_kind, omni_bridge_event); utils::redis::zadd( config, redis_connection_manager, &utils::pending_transactions::get_pending_tx_key(chain_kind), nonce, pending_tx, ) .await; info!("Stored pending transaction {tx_hash} (nonce: {nonce}) for {chain_kind:?}"); Ok(()) }