From eddac7a6791a15fc39f53c443b102cefb929022f Mon Sep 17 00:00:00 2001 From: ysq Date: Sun, 7 Sep 2025 18:07:45 +0800 Subject: [PATCH] perf: optimize trade executor performance and reduce memory usage - Implement singleton pattern for TradeFactory - Replace TradeTimer with direct Instant measurements - Optimize PumpFun instruction building and parallel execution - Reduce unnecessary clones and improve memory allocation - Simplify address lookup table management --- src/common/address_lookup_cache.rs | 2 +- src/instruction/pumpfun.rs | 24 +++-- src/trading/common/address_lookup_manager.rs | 13 ++- src/trading/common/transaction_builder.rs | 12 ++- src/trading/core/executor.rs | 60 +++-------- src/trading/core/mod.rs | 3 +- src/trading/core/parallel.rs | 52 ++++----- src/trading/core/timer.rs | 45 -------- src/trading/factory.rs | 107 +++++++++---------- 9 files changed, 127 insertions(+), 191 deletions(-) delete mode 100755 src/trading/core/timer.rs diff --git a/src/common/address_lookup_cache.rs b/src/common/address_lookup_cache.rs index f2f0122..72d73ba 100755 --- a/src/common/address_lookup_cache.rs +++ b/src/common/address_lookup_cache.rs @@ -103,5 +103,5 @@ pub async fn get_address_lookup_table_account( lookup_table_address: &Pubkey, ) -> AddressLookupTableAccount { let cache = AddressLookupTableCache::get_instance(); - return cache.get_table_content(&lookup_table_address); + cache.get_table_content(lookup_table_address) } diff --git a/src/instruction/pumpfun.rs b/src/instruction/pumpfun.rs index 23f31ee..0f6b23d 100755 --- a/src/instruction/pumpfun.rs +++ b/src/instruction/pumpfun.rs @@ -45,7 +45,7 @@ impl InstructionBuilder for PumpFunInstructionBuilder { return Err(anyhow!("Amount cannot be zero")); } - let bonding_curve = protocol_params.bonding_curve.clone(); + let bonding_curve = &protocol_params.bonding_curve; let max_sol_cost = calculate_with_slippage_buy( params.sol_amount, @@ -53,12 +53,20 @@ impl InstructionBuilder for PumpFunInstructionBuilder { ); let creator_vault_pda = protocol_params.creator_vault; - let mut creator = Pubkey::default(); - if let Some(default_creator_ata) = get_creator_vault_pda(&creator) { - if default_creator_ata != creator_vault_pda { - creator = creator_vault_pda; + // Optimize creator lookup - avoid PDA calculation if not default + let creator = if creator_vault_pda == Pubkey::default() { + Pubkey::default() + } else { + // Fast check against cached default creator vault + static DEFAULT_CREATOR_VAULT: std::sync::LazyLock> = + std::sync::LazyLock::new(|| get_creator_vault_pda(&Pubkey::default())); + + if Some(creator_vault_pda) == *DEFAULT_CREATOR_VAULT { + Pubkey::default() + } else { + creator_vault_pda } - } + }; let buy_token_amount = get_buy_token_amount_from_sol_amount( bonding_curve.virtual_token_reserves as u128, @@ -68,7 +76,7 @@ impl InstructionBuilder for PumpFunInstructionBuilder { params.sol_amount, ); - let mut instructions = vec![]; + let mut instructions = Vec::with_capacity(2); // Create associated token account instructions.push(create_associated_token_account( @@ -99,7 +107,7 @@ impl InstructionBuilder for PumpFunInstructionBuilder { .downcast_ref::() .ok_or_else(|| anyhow!("Invalid protocol params for PumpFun"))?; - let bonding_curve = protocol_params.bonding_curve.clone(); + let bonding_curve = &protocol_params.bonding_curve; let token_amount = if let Some(amount) = params.token_amount { if amount == 0 { diff --git a/src/trading/common/address_lookup_manager.rs b/src/trading/common/address_lookup_manager.rs index 62a41e2..9d393f1 100755 --- a/src/trading/common/address_lookup_manager.rs +++ b/src/trading/common/address_lookup_manager.rs @@ -7,12 +7,11 @@ use crate::common::address_lookup_cache::get_address_lookup_table_account; pub async fn get_address_lookup_table_accounts( lookup_table_key: Option, ) -> Vec { - let mut address_lookup_table_accounts = vec![]; - - if let Some(lookup_table_key) = lookup_table_key { - let account = get_address_lookup_table_account(&lookup_table_key).await; - address_lookup_table_accounts.push(account); + match lookup_table_key { + Some(key) => { + let account = get_address_lookup_table_account(&key).await; + vec![account] + } + None => Vec::new(), } - - address_lookup_table_accounts } diff --git a/src/trading/common/transaction_builder.rs b/src/trading/common/transaction_builder.rs index c00ce33..88445e3 100755 --- a/src/trading/common/transaction_builder.rs +++ b/src/trading/common/transaction_builder.rs @@ -27,13 +27,13 @@ pub async fn build_transaction( recent_blockhash: Hash, data_size_limit: u32, middleware_manager: Option>, - protocol_name: String, + protocol_name: &str, is_buy: bool, with_tip: bool, tip_account: &Pubkey, tip_amount: f64, ) -> Result { - let mut instructions = vec![]; + let mut instructions = Vec::with_capacity(business_instructions.len() + 5); // 添加nonce指令 if is_buy { @@ -84,12 +84,16 @@ async fn build_versioned_transaction( address_lookup_table_accounts: Vec, blockhash: Hash, middleware_manager: Option>, - protocol_name: String, + protocol_name: &str, is_buy: bool, ) -> Result { let full_instructions = match middleware_manager { Some(middleware_manager) => middleware_manager - .apply_middlewares_process_full_instructions(instructions, protocol_name, is_buy)?, + .apply_middlewares_process_full_instructions( + instructions, + protocol_name.to_string(), + is_buy, + )?, None => instructions, }; let v0_message: v0::Message = v0::Message::try_compile( diff --git a/src/trading/core/executor.rs b/src/trading/core/executor.rs index e293abd..bf759e0 100755 --- a/src/trading/core/executor.rs +++ b/src/trading/core/executor.rs @@ -1,10 +1,9 @@ use anyhow::Result; -use std::sync::Arc; +use std::{sync::Arc, time::Instant}; use super::{ parallel::parallel_execute_with_tips, params::{BuyParams, SellParams}, - timer::TradeTimer, traits::{InstructionBuilder, TradeExecutor}, }; use crate::{swqos::SwqosClient, trading::middleware::MiddlewareManager}; @@ -38,26 +37,12 @@ impl TradeExecutor for GenericTradeExecutor { if data_size_limit == 0 { data_size_limit = MAX_LOADED_ACCOUNTS_DATA_SIZE_LIMIT; } - let timer = TradeTimer::new("Building buy transaction instructions"); - // Validate parameters - convert to BuyParams for validation - let buy_params = BuyParams { - rpc: params.rpc, - payer: params.payer.clone(), - mint: params.mint, - sol_amount: params.sol_amount, - slippage_basis_points: params.slippage_basis_points, - priority_fee: params.priority_fee.clone(), - lookup_table_key: params.lookup_table_key, - recent_blockhash: params.recent_blockhash, - data_size_limit: data_size_limit, - wait_transaction_confirmed: params.wait_transaction_confirmed, - protocol_params: params.protocol_params.clone(), - }; + let start = Instant::now(); - // Build instructions - let instructions = self.instruction_builder.build_buy_instructions(&buy_params).await?; - let final_instructions = match middleware_manager.clone() { + // Build instructions directly from params to avoid unnecessary cloning + let instructions = self.instruction_builder.build_buy_instructions(¶ms).await?; + let final_instructions = match &middleware_manager { Some(middleware_manager) => middleware_manager .apply_middlewares_process_protocol_instructions( instructions, @@ -67,19 +52,19 @@ impl TradeExecutor for GenericTradeExecutor { None => instructions, }; - timer.finish(); + println!("Building buy transaction instructions time cost: {:?}", start.elapsed()); // Execute transactions in parallel parallel_execute_with_tips( swqos_clients, params.payer, final_instructions, - params.priority_fee, + Arc::new(params.priority_fee), params.lookup_table_key, params.recent_blockhash, data_size_limit, middleware_manager, - self.protocol_name.to_string(), + self.protocol_name, true, params.wait_transaction_confirmed, true, @@ -95,26 +80,11 @@ impl TradeExecutor for GenericTradeExecutor { swqos_clients: Vec>, middleware_manager: Option>, ) -> Result<()> { - let timer = TradeTimer::new("Building sell transaction instructions"); + let start = Instant::now(); - // Convert to SellParams for instruction building - let sell_params = SellParams { - rpc: params.rpc, - payer: params.payer.clone(), - mint: params.mint, - token_amount: params.token_amount, - slippage_basis_points: params.slippage_basis_points, - priority_fee: params.priority_fee.clone(), - lookup_table_key: params.lookup_table_key, - recent_blockhash: params.recent_blockhash, - wait_transaction_confirmed: params.wait_transaction_confirmed, - protocol_params: params.protocol_params.clone(), - with_tip: params.with_tip, - }; - - // Build instructions - let instructions = self.instruction_builder.build_sell_instructions(&sell_params).await?; - let final_instructions = match middleware_manager.clone() { + // Build instructions directly from params to avoid unnecessary cloning + let instructions = self.instruction_builder.build_sell_instructions(¶ms).await?; + let final_instructions = match &middleware_manager { Some(middleware_manager) => middleware_manager .apply_middlewares_process_protocol_instructions( instructions, @@ -124,19 +94,19 @@ impl TradeExecutor for GenericTradeExecutor { None => instructions, }; - timer.finish(); + println!("Building sell transaction instructions time cost: {:?}", start.elapsed()); // Execute transactions in parallel parallel_execute_with_tips( swqos_clients, params.payer, final_instructions, - params.priority_fee, + Arc::new(params.priority_fee), params.lookup_table_key, params.recent_blockhash, 0, middleware_manager, - self.protocol_name.to_string(), + self.protocol_name, false, params.wait_transaction_confirmed, params.with_tip, diff --git a/src/trading/core/mod.rs b/src/trading/core/mod.rs index 398bf30..1a05b79 100755 --- a/src/trading/core/mod.rs +++ b/src/trading/core/mod.rs @@ -1,5 +1,4 @@ pub mod params; pub mod traits; pub mod executor; -pub mod parallel; -pub mod timer; \ No newline at end of file +pub mod parallel; \ No newline at end of file diff --git a/src/trading/core/parallel.rs b/src/trading/core/parallel.rs index e532f6b..12e4e05 100755 --- a/src/trading/core/parallel.rs +++ b/src/trading/core/parallel.rs @@ -1,14 +1,14 @@ use anyhow::{anyhow, Result}; use solana_hash::Hash; use solana_sdk::{instruction::Instruction, pubkey::Pubkey, signature::Keypair}; -use std::{str::FromStr, sync::Arc}; +use std::{str::FromStr, sync::Arc, time::Instant}; use tokio::sync::mpsc; use tokio::task::JoinHandle; use crate::{ common::PriorityFee, swqos::{SwqosClient, SwqosType, TradeType}, - trading::{common::build_transaction, core::timer::TradeTimer, MiddlewareManager}, + trading::{common::build_transaction, MiddlewareManager}, }; /// Generic function for parallel transaction execution @@ -16,26 +16,34 @@ pub async fn parallel_execute_with_tips( swqos_clients: Vec>, payer: Arc, instructions: Vec, - priority_fee: PriorityFee, + priority_fee: Arc, lookup_table_key: Option, recent_blockhash: Hash, data_size_limit: u32, middleware_manager: Option>, - protocol_name: String, + protocol_name: &'static str, is_buy: bool, wait_transaction_confirmed: bool, with_tip: bool, ) -> Result<()> { let cores = core_affinity::get_core_ids().unwrap(); - let mut handles: Vec>> = vec![]; - - if is_buy && swqos_clients.len() > priority_fee.buy_tip_fees.len() { + let mut handles: Vec>> = Vec::with_capacity(swqos_clients.len()); + if is_buy + && (swqos_clients.len() > priority_fee.buy_tip_fees.len() + || priority_fee.buy_tip_fees.is_empty()) + { return Err(anyhow!("Number of tip clients exceeds the configured buy tip fees")); } - if !is_buy && swqos_clients.len() > priority_fee.sell_tip_fees.len() { + if !is_buy + && !with_tip + && (swqos_clients.len() > priority_fee.sell_tip_fees.len() + || priority_fee.sell_tip_fees.is_empty()) + { return Err(anyhow!("Number of tip clients exceeds the configured sell tip fees")); } + let instructions = Arc::new(instructions); + for i in 0..swqos_clients.len() { let swqos_client = swqos_clients[i].clone(); if !with_tip && !matches!(swqos_client.get_swqos_type(), SwqosType::Default) { @@ -47,43 +55,36 @@ pub async fn parallel_execute_with_tips( let core_id = cores[i % cores.len()]; let middleware_manager = middleware_manager.clone(); - let protocol_name = protocol_name.clone(); let handle = tokio::spawn(async move { core_affinity::set_for_current(core_id); - let mut timer = TradeTimer::new(format!( - "Building transaction instructions: {:?}", - swqos_client.get_swqos_type() - )); + let swqos_type = swqos_client.get_swqos_type(); + let mut start = Instant::now(); - let tip_account = swqos_client.get_tip_account()?; - let tip_account = Arc::new(Pubkey::from_str(&tip_account).map_err(|e| anyhow!(e))?); - if priority_fee.buy_tip_fees.len() == 0 { - return Err(anyhow!("buy_tip_fees is empty")); - } + let tip_account_str = swqos_client.get_tip_account()?; + let tip_account = Arc::new(Pubkey::from_str(&tip_account_str).unwrap_or_default()); let tip_amount = priority_fee.buy_tip_fees[i]; let transaction = build_transaction( payer, &priority_fee, - instructions, + (*instructions).clone(), lookup_table_key, recent_blockhash, data_size_limit, middleware_manager, protocol_name, is_buy, - swqos_client.get_swqos_type() != SwqosType::Default, + swqos_type != SwqosType::Default, &tip_account, tip_amount, ) .await?; - timer.stage(format!( - "Submitting transaction instructions: {:?}", - swqos_client.get_swqos_type() - )); + println!("Building transaction instructions: {:?} {:?}", swqos_type, start.elapsed()); + + start = Instant::now(); swqos_client .send_transaction( @@ -92,7 +93,8 @@ pub async fn parallel_execute_with_tips( ) .await?; - timer.finish(); + println!("Submitting transaction instructions: {:?} {:?}", swqos_type, start.elapsed()); + Ok::<(), anyhow::Error>(()) }); diff --git a/src/trading/core/timer.rs b/src/trading/core/timer.rs deleted file mode 100755 index 3a81315..0000000 --- a/src/trading/core/timer.rs +++ /dev/null @@ -1,45 +0,0 @@ -use std::time::Instant; - -/// Trade time measurement tool -#[derive(Clone)] -pub struct TradeTimer { - start_time: Instant, - stage: String, -} - -impl TradeTimer { - /// Create a new timer - pub fn new(stage: impl Into) -> Self { - Self { start_time: Instant::now(), stage: stage.into() } - } - - /// Record current stage time and start a new stage - pub fn stage(&mut self, new_stage: impl Into) { - let elapsed = self.start_time.elapsed(); - println!(" {} time cost: {:?}", self.stage, elapsed); - - self.start_time = Instant::now(); - self.stage = new_stage.into(); - } - - /// Complete timing and output final time cost - pub fn finish(mut self) { - let elapsed = self.start_time.elapsed(); - println!(" {} time cost: {:?}", self.stage, elapsed); - self.stage.clear(); // Clear stage to avoid duplicate printing in Drop - } - - /// Get the elapsed time of current stage (without resetting the timer) - pub fn elapsed(&self) -> std::time::Duration { - self.start_time.elapsed() - } -} - -impl Drop for TradeTimer { - fn drop(&mut self) { - if !self.stage.is_empty() { - let elapsed = self.start_time.elapsed(); - println!(" {} time cost: {:?}", self.stage, elapsed); - } - } -} diff --git a/src/trading/factory.rs b/src/trading/factory.rs index 1050d96..0d5b06b 100755 --- a/src/trading/factory.rs +++ b/src/trading/factory.rs @@ -19,70 +19,69 @@ pub enum DexType { RaydiumAmmV4, } -impl std::fmt::Display for DexType { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - match self { - DexType::PumpFun => write!(f, "PumpFun"), - DexType::PumpSwap => write!(f, "PumpSwap"), - DexType::Bonk => write!(f, "Bonk"), - DexType::RaydiumCpmm => write!(f, "RaydiumCpmm"), - DexType::RaydiumAmmV4 => write!(f, "RaydiumAmmV4"), - } - } -} - -impl std::str::FromStr for DexType { - type Err = anyhow::Error; - - fn from_str(s: &str) -> Result { - match s.to_lowercase().as_str() { - "pumpfun" => Ok(DexType::PumpFun), - "pumpswap" => Ok(DexType::PumpSwap), - "bonk" => Ok(DexType::Bonk), - "raydiumcpmm" => Ok(DexType::RaydiumCpmm), - "raydiumammv4" => Ok(DexType::RaydiumAmmV4), - _ => Err(anyhow!("Unsupported protocol: {}", s)), - } - } -} - /// 交易工厂 - 用于创建不同协议的交易执行器 pub struct TradeFactory; impl TradeFactory { - /// 创建指定协议的交易执行器 + /// 创建指定协议的交易执行器(零开销单例) pub fn create_executor(dex_type: DexType) -> Arc { match dex_type { - DexType::PumpFun => { - let instruction_builder = Arc::new(PumpFunInstructionBuilder); - Arc::new(GenericTradeExecutor::new(instruction_builder, "PumpFun")) - } - DexType::PumpSwap => { - let instruction_builder = Arc::new(PumpSwapInstructionBuilder); - Arc::new(GenericTradeExecutor::new(instruction_builder, "PumpSwap")) - } - DexType::Bonk => { - let instruction_builder = Arc::new(BonkInstructionBuilder); - Arc::new(GenericTradeExecutor::new(instruction_builder, "Bonk")) - } - DexType::RaydiumCpmm => { - let instruction_builder = Arc::new(RaydiumCpmmInstructionBuilder); - Arc::new(GenericTradeExecutor::new(instruction_builder, "RaydiumCpmm")) - } - DexType::RaydiumAmmV4 => { - let instruction_builder = Arc::new(RaydiumAmmV4InstructionBuilder); - Arc::new(GenericTradeExecutor::new(instruction_builder, "RaydiumAmmV4")) - } + DexType::PumpFun => Self::pumpfun_executor(), + DexType::PumpSwap => Self::pumpswap_executor(), + DexType::Bonk => Self::bonk_executor(), + DexType::RaydiumCpmm => Self::raydium_cpmm_executor(), + DexType::RaydiumAmmV4 => Self::raydium_amm_v4_executor(), } } - /// 获取所有支持的协议 - pub fn supported_dex_types() -> Vec { - vec![DexType::PumpFun, DexType::PumpSwap, DexType::Bonk, DexType::RaydiumCpmm] + // Static instances created at compile time - zero runtime overhead + #[inline] + fn pumpfun_executor() -> Arc { + static INSTANCE: std::sync::LazyLock> = + std::sync::LazyLock::new(|| { + let instruction_builder = Arc::new(PumpFunInstructionBuilder); + Arc::new(GenericTradeExecutor::new(instruction_builder, "PumpFun")) + }); + INSTANCE.clone() } - /// 检查协议是否支持 - pub fn is_supported(dex_type: &DexType) -> bool { - Self::supported_dex_types().contains(dex_type) + #[inline] + fn pumpswap_executor() -> Arc { + static INSTANCE: std::sync::LazyLock> = + std::sync::LazyLock::new(|| { + let instruction_builder = Arc::new(PumpSwapInstructionBuilder); + Arc::new(GenericTradeExecutor::new(instruction_builder, "PumpSwap")) + }); + INSTANCE.clone() + } + + #[inline] + fn bonk_executor() -> Arc { + static INSTANCE: std::sync::LazyLock> = + std::sync::LazyLock::new(|| { + let instruction_builder = Arc::new(BonkInstructionBuilder); + Arc::new(GenericTradeExecutor::new(instruction_builder, "Bonk")) + }); + INSTANCE.clone() + } + + #[inline] + fn raydium_cpmm_executor() -> Arc { + static INSTANCE: std::sync::LazyLock> = + std::sync::LazyLock::new(|| { + let instruction_builder = Arc::new(RaydiumCpmmInstructionBuilder); + Arc::new(GenericTradeExecutor::new(instruction_builder, "RaydiumCpmm")) + }); + INSTANCE.clone() + } + + #[inline] + fn raydium_amm_v4_executor() -> Arc { + static INSTANCE: std::sync::LazyLock> = + std::sync::LazyLock::new(|| { + let instruction_builder = Arc::new(RaydiumAmmV4InstructionBuilder); + Arc::new(GenericTradeExecutor::new(instruction_builder, "RaydiumAmmV4")) + }); + INSTANCE.clone() } }