From cc29851c67c3cdf95a1fa842c6ea1d41ddbf3a8e Mon Sep 17 00:00:00 2001 From: William Date: Fri, 3 Jan 2025 14:43:26 +0800 Subject: [PATCH] add jito --- README.md | 3 +- crates/pumpfun/Cargo.toml | 6 +- crates/pumpfun/src/error/mod.rs | 3 + crates/pumpfun/src/jito/mod.rs | 217 ++++++++++++++++++++++++ crates/pumpfun/src/lib.rs | 289 ++++++++++++++++++++++++++++---- 5 files changed, 483 insertions(+), 35 deletions(-) create mode 100644 crates/pumpfun/src/jito/mod.rs diff --git a/README.md b/README.md index 84d2e20..5a37a38 100644 --- a/README.md +++ b/README.md @@ -12,6 +12,7 @@ This repository is forked from [https://github.com/nhuxhr/pumpfun-rs](https://gi 4. Add `logs_subscribe` to subscribe the logs of the PumpFun program. 5. Add `logs_events` to define the event of the logs. 6. Add `logs_parser` to parse the logs. +7. Add `jito` to send transaction with Jito. ## Installation @@ -19,7 +20,7 @@ Add the following to your `Cargo.toml`: ```toml [dependencies] -mai3-pumpfun-sdk = "2.3.0" +mai3-pumpfun-sdk = "2.4.0" ``` ## Usage diff --git a/crates/pumpfun/Cargo.toml b/crates/pumpfun/Cargo.toml index e97ea26..90f7cf0 100644 --- a/crates/pumpfun/Cargo.toml +++ b/crates/pumpfun/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mai3-pumpfun-sdk" -version = "2.3.0" +version = "2.4.0" edition = "2021" authors = ["William "] repository = "https://github.com/MiracleAI-Labs/pumpfun-sdk" @@ -27,3 +27,7 @@ futures = "0.3.31" futures-util = "0.3.31" base64 = "0.22.1" bs58 = "0.5.1" +rand = "0.8.5" +bincode = "1.3.3" +reqwest = { version = "0.11.27", features = ["json"] } + diff --git a/crates/pumpfun/src/error/mod.rs b/crates/pumpfun/src/error/mod.rs index 09d1f12..315e0eb 100644 --- a/crates/pumpfun/src/error/mod.rs +++ b/crates/pumpfun/src/error/mod.rs @@ -52,6 +52,8 @@ pub enum ClientError { Parse(String, String), + Jito(String, String), + Join(String), Subscribe(String, String), @@ -84,6 +86,7 @@ impl std::fmt::Display for ClientError { Self::RateLimitExceeded => write!(f, "Rate limit exceeded"), Self::Solana(msg, details) => write!(f, "Solana error: {}, details: {}", msg, details), Self::Parse(msg, details) => write!(f, "Parse error: {}, details: {}", msg, details), + Self::Jito(msg, details) => write!(f, "Jito error: {}, details: {}", msg, details), Self::Join(msg) => write!(f, "Task join error: {}", msg), Self::Subscribe(msg, details) => write!(f, "Subscribe error: {}, details: {}", msg, details), Self::Send(msg, details) => write!(f, "Send error: {}, details: {}", msg, details), diff --git a/crates/pumpfun/src/jito/mod.rs b/crates/pumpfun/src/jito/mod.rs new file mode 100644 index 0000000..5f12227 --- /dev/null +++ b/crates/pumpfun/src/jito/mod.rs @@ -0,0 +1,217 @@ +use { + bincode, bs58, reqwest, serde::Deserialize, serde_json::{json, Value}, + std::{str::FromStr, time::Duration} +}; + +use anchor_client::solana_sdk::{ + commitment_config::CommitmentConfig, pubkey::Pubkey, signature::Signature, transaction::Transaction +}; + +use crate::error::ClientError::{self, *}; + +/// 常量定义 +pub const MAX_RETRIES: u8 = 3; +pub const RETRY_DELAY: Duration = Duration::from_millis(200); + +/// 交易配置 +#[derive(Debug, Clone)] +pub struct TransactionConfig { + pub skip_preflight: bool, + pub preflight_commitment: CommitmentConfig, + pub encoding: String, + pub last_n_blocks: u64, +} + +impl Default for TransactionConfig { + fn default() -> Self { + Self { + skip_preflight: true, // Jito 建议跳过预检 + preflight_commitment: CommitmentConfig::confirmed(), + encoding: "base58".to_string(), + last_n_blocks: 100, + } + } +} +#[derive(Clone)] +pub struct JitoClient { + endpoint: String, + client: reqwest::Client, + config: TransactionConfig, +} + +impl JitoClient { + pub fn new(endpoint: &str) -> Self { + Self { + endpoint: endpoint.to_string(), + client: reqwest::Client::new(), + config: TransactionConfig::default(), + } + } + + /// 获取随机的 tip account + pub async fn get_tip_account(&self) -> Result { + let response = self.send_request("getTipAccounts", json!([])).await?; + + if let Some(accounts) = response["result"].as_array() { + if accounts.is_empty() { + return Err(ClientError::Other( + "No JITO tip accounts found".to_string() + )); + } + + let random_index = rand::random::() % accounts.len(); + if let Some(account) = accounts.get(random_index) { + if let Some(address) = account.as_str() { + return Pubkey::from_str(address) + .map_err(|e| ClientError::Parse( + "Invalid tip account address".to_string(), + e.to_string() + )); + } + } + } + + Err(ClientError::Other("Failed to get Tip Account".to_string())) + } + + /// 估算优先费用 + pub async fn estimate_priority_fees( + &self, + account: &Pubkey, + ) -> Result { + let params = json!({ + "last_n_blocks": self.config.last_n_blocks, + "account": account.to_string(), + "api_version": 2 + }); + + let response = self.send_request("qn_estimatePriorityFees", params).await?; + + // 解析响应 + if let Some(result) = response.get("result") { + let estimate: PriorityFeeEstimate = serde_json::from_value(result.clone()) + .map_err(|e| ClientError::Parse( + "Failed to parse priority fee estimate".to_string(), + e.to_string() + ))?; + + Ok(estimate) + } else { + Err(ClientError::Parse( + "Invalid response format".to_string(), + "Missing result field".to_string() + )) + } + } + + /// 发送交易 + pub async fn send_transaction( + &self, + transaction: &Transaction, + ) -> Result { + let wire_transaction = bincode::serialize(transaction).map_err(|e| + ClientError::Parse( + "Transaction serialization failed".to_string(), + e.to_string() + ))?; + + let encoded_tx = bs58::encode(&wire_transaction).into_string(); + + for retry in 0..MAX_RETRIES { + match self.try_send_transaction(&encoded_tx).await { + Ok(signature) => { + return Ok(Signature::from_str(&signature).map_err(|e| + ClientError::Parse( + "Invalid signature".to_string(), + e.to_string() + ))?); + }, + Err(e) => { + println!("Retry {} failed: {:?}", retry, e); + if retry == MAX_RETRIES - 1 { + return Err(e); + } + tokio::time::sleep(RETRY_DELAY).await; + } + } + } + + Err(ClientError::Other("Max retries exceeded".to_string())) + } + + async fn try_send_transaction(&self, encoded_tx: &str) -> Result { + let params = json!([ + encoded_tx, + { + "skipPreflight": self.config.skip_preflight, + "preflightCommitment": self.config.preflight_commitment.commitment, + "encoding": self.config.encoding, + "maxRetries": MAX_RETRIES, + "minContextSlot": null + } + ]); + + let response = self.send_request( + "sendTransaction", + params + ).await?; + + response["result"] + .as_str() + .map(|s| s.to_string()) + .ok_or_else(|| ClientError::Parse( + "Invalid response format".to_string(), + "Missing result field".to_string() + )) + } + + async fn send_request(&self, method: &str, params: Value) -> Result { + let request_body = json!({ + "jsonrpc": "2.0", + "id": 1, + "method": method, + "params": params + }); + + let response = self.client + .post(&self.endpoint) + .header("Content-Type", "application/json") + .json(&request_body) + .send() + .await + .map_err(|e| ClientError::Solana( + "Request failed".to_string(), + e.to_string() + ))?; + + let response_data: Value = response.json().await + .map_err(|e| ClientError::Parse( + "Invalid JSON response".to_string(), + e.to_string() + ))?; + + if let Some(error) = response_data.get("error") { + return Err(ClientError::Solana( + "RPC error".to_string(), + error.to_string() + )); + } + + Ok(response_data) + } +} + +#[derive(Debug, Deserialize)] +pub struct PriorityFeeEstimate { + pub recommended: u64, + pub per_compute_unit: PriorityFeeLevel, + pub per_transaction: PriorityFeeLevel, +} + +#[derive(Debug, Deserialize)] +pub struct PriorityFeeLevel { + pub extreme: u64, // 95th percentile + pub high: u64, // 80th percentile + pub medium: u64, // 60th percentile + pub low: u64, // 40th percentile +} diff --git a/crates/pumpfun/src/lib.rs b/crates/pumpfun/src/lib.rs index 2c8ca3f..beb7096 100644 --- a/crates/pumpfun/src/lib.rs +++ b/crates/pumpfun/src/lib.rs @@ -5,6 +5,7 @@ pub mod constants; pub mod error; pub mod instruction; pub mod utils; +pub mod jito; use crate::error::ClientError::*; @@ -15,6 +16,10 @@ use anchor_client::{ pubkey::Pubkey, signature::{Keypair, Signature}, signer::Signer, + instruction::Instruction, + system_instruction, + compute_budget::ComputeBudgetInstruction, + transaction::Transaction, }, Client, Cluster, Program, }; @@ -22,10 +27,18 @@ use anchor_spl::associated_token::{ get_associated_token_address, spl_associated_token_account::instruction::create_associated_token_account, }; -use borsh::BorshDeserialize; -pub use pumpfun_cpi as cpi; -use solana_sdk::compute_budget::ComputeBudgetInstruction; + use std::sync::Arc; +use borsh::BorshDeserialize; +use std::time::Instant; +pub use pumpfun_cpi as cpi; + +use crate::jito::JitoClient; +use crate::error::ClientError; + +const DEFAULT_SLIPPAGE: u64 = 500; // 10% +const DEFAULT_COMPUTE_UNIT_LIMIT: u32 = 68_000; +const DEFAULT_COMPUTE_UNIT_PRICE: u64 = 400_000; /// Configuration for priority fee compute unit parameters #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -42,6 +55,8 @@ pub struct PumpFun { pub rpc: RpcClient, /// Keypair used to sign transactions pub payer: Arc, + /// Jito client instance + pub jito_client: Option, /// Anchor client instance pub client: Client>, /// Anchor program instance @@ -63,6 +78,7 @@ impl PumpFun { /// Returns a new PumpFun client instance configured with the provided parameters pub fn new( cluster: Cluster, + jito_url: String, payer: Arc, options: Option, ws: Option, @@ -74,6 +90,11 @@ impl PumpFun { cluster.url() }); + let mut jito_client = None; + if !jito_url.is_empty() { + jito_client = Some(JitoClient::new(&jito_url)); + } + // Create Anchor Client with optional commitment config let client: Client> = if let Some(options) = options { Client::new_with_options(cluster.clone(), payer.clone(), options) @@ -88,6 +109,7 @@ impl PumpFun { Self { rpc, payer, + jito_client, client, program, } @@ -321,6 +343,126 @@ impl PumpFun { Ok(signature) } + /// Buys tokens from a bonding curve with Jito + pub async fn buy_with_jito( + &self, + mint: &Pubkey, + amount_sol: u64, + slippage_basis_points: Option, + priority_fee: Option, + ) -> Result { + let start_time = Instant::now(); + + if self.jito_client.is_none() { + return Err(ClientError::Other( + "Jito client not found".to_string(), + )); + } + + // Get accounts and calculate buy amounts + let global_account = self.get_global_account()?; + + // 获取 bonding curve pda + let bonding_curve_pda = Self::get_bonding_curve_pda(mint).unwrap(); + // 获取 bonding curve account + let bonding_curve_account = self.get_bonding_curve_account(mint)?; + // 获取 buy amount + let buy_amount = bonding_curve_account + .get_buy_price(amount_sol) + .map_err(error::ClientError::BondingCurveError)?; + + let buy_amount_with_slippage = + utils::calculate_with_slippage_buy(amount_sol, slippage_basis_points.unwrap_or(500)); + + let mut unit_limit = DEFAULT_COMPUTE_UNIT_LIMIT; + let mut unit_price = DEFAULT_COMPUTE_UNIT_PRICE; + + // 准备所有指令 + let mut instructions: Vec = vec![]; + + // Add priority fee if provided + if let Some(fee) = priority_fee { + if let Some(limit) = fee.limit { + unit_limit = limit; + let limit_ix = ComputeBudgetInstruction::set_compute_unit_limit(limit); + instructions.push(limit_ix); + } + + if let Some(price) = fee.price { + unit_price = price; + let price_ix = ComputeBudgetInstruction::set_compute_unit_price(price); + instructions.push(price_ix); + } + } + + // 获取 jito client + let jito_client = self.jito_client.as_ref().unwrap(); + + // 获取优先费用估算 + let priority_fees = jito_client.estimate_priority_fees(&bonding_curve_pda).await?; + + // 计算每计算单元的优先费用(使用 Extreme 级别) + let priority_fee_per_cu = priority_fees.per_compute_unit.extreme; + + // 完整的单位转换过程 + let total_priority_fee_microlamports = priority_fee_per_cu as u128 * unit_limit as u128; + let total_priority_fee_lamports = total_priority_fee_microlamports / 1_000_000; + let total_priority_fee_sol = total_priority_fee_lamports as f64 / 1_000_000_000.0; + + println!("Priority fee details:"); + println!(" Per CU (microlamports): {}", priority_fee_per_cu); + println!(" Total (lamports): {}", total_priority_fee_lamports); + println!(" Total (SOL): {:.9}", total_priority_fee_sol); + + // 获取 tip account + let tip_account = jito_client.get_tip_account().await.unwrap(); + + // Create Associated Token Account if needed + let ata: Pubkey = get_associated_token_address(&self.payer.pubkey(), mint); + if self.rpc.get_account(&ata).is_err() { + instructions.push(create_associated_token_account( + &self.payer.pubkey(), + &self.payer.pubkey(), + mint, + &constants::accounts::TOKEN_PROGRAM, + )); + } + + // Add buy instruction + instructions.push(instruction::buy( + &self.payer.clone().as_ref(), + mint, + &global_account.fee_recipient, + cpi::instruction::Buy { + _amount: buy_amount, + _max_sol_cost: buy_amount_with_slippage, + }, + )); + + instructions.push( + system_instruction::transfer( + &self.payer.pubkey(), + &tip_account, + total_priority_fee_lamports as u64, + ), + ); + + // 创建并发送交易 + let recent_blockhash = self.rpc.get_latest_blockhash()?; + let transaction = Transaction::new_signed_with_payer( + &instructions, + Some(&self.payer.pubkey()), + &[&self.payer.clone()], + recent_blockhash, + ); + + // 通过 Jito 发送交易 + let signature = jito_client.send_transaction(&transaction).await.unwrap(); + println!("Total Jito buy operation time: {:?}ms", start_time.elapsed().as_millis()); + + Ok(signature) + } + /// Sells tokens back to the bonding curve in exchange for SOL /// /// # Arguments @@ -393,6 +535,117 @@ impl PumpFun { Ok(signature) } + /// Sells tokens back to the bonding curve in exchange for SOL with Jito + pub async fn sell_with_jito( + &self, + mint: &Pubkey, + amount_token: Option, + slippage_basis_points: Option, + priority_fee: Option, + ) -> Result { + let start_time = Instant::now(); + + if self.jito_client.is_none() { + return Err(ClientError::Other( + "Jito client not found".to_string(), + )); + } + + // Get accounts and calculate sell amounts + let ata: Pubkey = get_associated_token_address(&self.payer.pubkey(), mint); + let balance = self.rpc.get_token_account_balance(&ata).unwrap(); + let balance_u64: u64 = balance.amount.parse::().unwrap(); + let _amount = amount_token.unwrap_or(balance_u64); + let global_account = self.get_global_account()?; + let bonding_curve_pda = Self::get_bonding_curve_pda(mint).unwrap(); + let bonding_curve_account = self.get_bonding_curve_account(mint)?; + let min_sol_output = bonding_curve_account + .get_sell_price(_amount, global_account.fee_basis_points) + .map_err(error::ClientError::BondingCurveError)?; + let _min_sol_output = utils::calculate_with_slippage_sell( + min_sol_output, + slippage_basis_points.unwrap_or(500), + ); + + let mut unit_limit = DEFAULT_COMPUTE_UNIT_LIMIT; + let mut unit_price = DEFAULT_COMPUTE_UNIT_PRICE; + + // 准备所有指令 + let mut instructions: Vec = vec![]; + + // Add priority fee if provided + if let Some(fee) = priority_fee { + if let Some(limit) = fee.limit { + unit_limit = limit; + let limit_ix = ComputeBudgetInstruction::set_compute_unit_limit(limit); + instructions.push(limit_ix); + } + + if let Some(price) = fee.price { + unit_price = price; + let price_ix = ComputeBudgetInstruction::set_compute_unit_price(price); + instructions.push(price_ix); + } + } + + // 获取 jito client + let jito_client = self.jito_client.as_ref().unwrap(); + + // 获取优先费用估算 + let priority_fees = jito_client.estimate_priority_fees(&bonding_curve_pda).await?; + + // 计算每计算单元的优先费用(使用 Extreme 级别) + let priority_fee_per_cu = priority_fees.per_compute_unit.extreme; + + // 完整的单位转换过程 + let total_priority_fee_microlamports = priority_fee_per_cu as u128 * unit_limit as u128; + let total_priority_fee_lamports = total_priority_fee_microlamports / 1_000_000; + let total_priority_fee_sol = total_priority_fee_lamports as f64 / 1_000_000_000.0; + + println!("Priority fee details:"); + println!(" Per CU (microlamports): {}", priority_fee_per_cu); + println!(" Total (lamports): {}", total_priority_fee_lamports); + println!(" Total (SOL): {:.9}", total_priority_fee_sol); + + // 获取 tip account + let tip_account = jito_client.get_tip_account().await.unwrap(); + + // Add buy instruction + instructions.push(instruction::sell( + &self.payer.clone().as_ref(), + mint, + &global_account.fee_recipient, + cpi::instruction::Sell { + _amount, + _min_sol_output, + }, + )); + + // 添加 tip 指令 + instructions.push( + system_instruction::transfer( + &self.payer.pubkey(), + &tip_account, + total_priority_fee_lamports as u64, + ), + ); + + // 创建并发送交易 + let recent_blockhash = self.rpc.get_latest_blockhash()?; + let transaction = Transaction::new_signed_with_payer( + &instructions, + Some(&self.payer.pubkey()), + &[&self.payer.clone()], + recent_blockhash, + ); + + // 通过 Jito 发送交易 + let signature = jito_client.send_transaction(&transaction).await.unwrap(); + println!("Total Jito sell operation time: {:?}ms", start_time.elapsed().as_millis()); + + Ok(signature) + } + /// Gets the Program Derived Address (PDA) for the global state account /// /// # Returns @@ -521,34 +774,4 @@ mod tests { assert!(bonding_curve_pda.is_some()); assert!(metadata_pda != Pubkey::default()); } - - // #[tokio::test] - // async fn test_logs_subscription() { - // let ws_url = "wss://api.mainnet-beta.solana.com"; - // let program_address = "6EF8rrecthR5Dkzon8Nwu78hRvfCKubJ14M5uBEwF6P"; - // let commitment = CommitmentConfig::confirmed(); - - // let process_logs_callback = |signature: String, logs: Vec| { - // println!("Signature: {}", signature); - // for log in logs { - // println!("Log: {}", log); - // } - // }; - - // let subscription = start_subscription( - // ws_url, - // program_address, - // commitment, - // process_logs_callback, - // ) - // .await.unwrap(); - - // // 模拟运行5秒后关闭订阅 - // tokio::time::sleep(tokio::time::Duration::from_secs(5)).await; - - // (subscription.unsub_fn)(); // 调用取消逻辑 - // subscription.task.await.unwrap(); - - // println!("Subscription closed."); - // } }