Fix trade execution correctness and release 4.0.11

This commit is contained in:
0xfnzero
2026-05-21 17:50:06 +08:00
parent 0cf0064e80
commit 8ed2aef931
10 changed files with 186 additions and 132 deletions
+21 -1
View File
@@ -3,7 +3,9 @@ use once_cell::sync::Lazy;
use smallvec::SmallVec;
use solana_compute_budget_interface::ComputeBudgetInstruction;
use solana_sdk::instruction::Instruction;
use std::sync::Arc;
use std::{hash::Hash, sync::Arc};
const MAX_COMPUTE_BUDGET_CACHE_SIZE: usize = 4_096;
/// Cache key containing all parameters for compute budget instructions
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
@@ -17,6 +19,22 @@ struct ComputeBudgetCacheKey {
static COMPUTE_BUDGET_CACHE: Lazy<DashMap<ComputeBudgetCacheKey, Arc<SmallVec<[Instruction; 2]>>>> =
Lazy::new(|| DashMap::new());
#[inline]
fn prune_cache<K, V>(cache: &DashMap<K, V>, max_size: usize)
where
K: Eq + Hash + Clone,
{
let len = cache.len();
if len <= max_size {
return;
}
let remove_count = (len - max_size).max(max_size / 16).min(len);
let keys: Vec<K> = cache.iter().take(remove_count).map(|entry| entry.key().clone()).collect();
for key in keys {
cache.remove(&key);
}
}
/// Extend `instructions` with compute budget instructions; on cache hit extends from cached Arc (no SmallVec clone).
#[inline(always)]
pub fn extend_compute_budget_instructions(
@@ -41,6 +59,7 @@ pub fn extend_compute_budget_instructions(
let arc = Arc::new(insts);
instructions.extend(arc.iter().cloned());
COMPUTE_BUDGET_CACHE.insert(cache_key, arc);
prune_cache(&COMPUTE_BUDGET_CACHE, MAX_COMPUTE_BUDGET_CACHE_SIZE);
}
/// Returns compute budget instructions (allocates on cache hit; prefer `extend_compute_budget_instructions` on hot path).
@@ -59,5 +78,6 @@ pub fn compute_budget_instructions(unit_price: u64, unit_limit: u32) -> SmallVec
}
let arc = Arc::new(insts.clone());
COMPUTE_BUDGET_CACHE.insert(cache_key, arc);
prune_cache(&COMPUTE_BUDGET_CACHE, MAX_COMPUTE_BUDGET_CACHE_SIZE);
insts
}
+6 -5
View File
@@ -1,3 +1,4 @@
use anyhow::anyhow;
use solana_hash::Hash;
use solana_message::AddressLookupTableAccount;
use solana_sdk::{
@@ -93,19 +94,19 @@ fn build_versioned_transaction(
// 使用预分配的交易构建器以降低延迟
let mut builder = acquire_builder();
let versioned_msg = builder.build_zero_alloc(
let build_result = builder.build_zero_alloc(
&payer.pubkey(),
&full_instructions,
address_lookup_table_account,
blockhash,
);
release_builder(builder);
let versioned_msg = build_result?;
let msg_bytes = versioned_msg.serialize();
let signature = payer.as_ref().try_sign_message(&msg_bytes).expect("sign failed");
let signature =
payer.as_ref().try_sign_message(&msg_bytes).map_err(|e| anyhow!("sign failed: {e}"))?;
let tx = VersionedTransaction { signatures: vec![signature], message: versioned_msg };
// 归还构建器到池
release_builder(builder);
Ok(tx)
}
+25 -7
View File
@@ -471,8 +471,31 @@ impl ResultCollector {
}
}
/// Fast submit mode for callers that do not wait for on-chain confirmation.
/// Return as soon as one route accepts, a landed failure consumes the nonce, all routes finish,
/// or a short submit window expires. Slow HTTP routes continue in worker tasks but no longer
/// block post-buy monitoring / sell scheduling.
async fn wait_for_first_submitted(
&self,
timeout: Duration,
) -> Option<(bool, Vec<Signature>, Option<anyhow::Error>, Vec<SwqosSubmitTiming>)> {
let start = Instant::now();
let poll_interval = Duration::from_millis(1);
loop {
if self.success_flag.load(Ordering::Acquire)
|| self.landed_failed_flag.load(Ordering::Acquire)
|| self.completed_count.load(Ordering::Acquire) >= self.total_tasks
|| start.elapsed() >= timeout
{
return self.get_first();
}
tokio::time::sleep(poll_interval).await;
}
}
/// 等待全部任务完成(不等待链上确认),然后收集并返回所有签名。用于「多路提交」时返回多笔签名。
/// 轮询间隔 2ms,避免 50ms 间隔在最后一笔返回时多等几十 ms 拉高 submit 耗时。
#[allow(dead_code)]
async fn wait_for_all_submitted(
&self,
timeout_secs: u64,
@@ -627,9 +650,6 @@ pub async fn execute_parallel(
// Task preparation completed: one shared context (clone once per batch), then minimal per-task data.
let channel_count = selected_task_configs.len().max(1);
let collector = Arc::new(ResultCollector::new(channel_count));
// 上限最多 5s:单路卡死时才会等满;若所有通道在窗口内回完,主循环会提前结束(不睡满秒数)。
let submit_timeout_secs: u64 =
(3u64 + (channel_count.min(16) as u64).div_ceil(2).min(6)).clamp(3, 5);
let shared = Arc::new(SwqosSharedContext {
payer,
instructions,
@@ -714,12 +734,10 @@ pub async fn execute_parallel(
// All jobs enqueued (no spawn on hot path)
if !wait_transaction_confirmed {
// submit_timeout_secs 为「等齐各 SWQOS HTTP 应答」的上限;收齐后会立刻进入短 settle,不会睡满整段秒数。
// `wait_for_all_submitted` 仅在未收齐时追加长 grace,避免过早 drain 丢晚到签名。
let ret = collector.wait_for_all_submitted(submit_timeout_secs).await.unwrap_or((
let ret = collector.wait_for_first_submitted(Duration::from_millis(500)).await.unwrap_or((
false,
vec![],
Some(anyhow!("No SWQOS result within grace window (primary {}s)", submit_timeout_secs)),
Some(anyhow!("No SWQOS result within fast submit window")),
vec![],
));
let (success, signatures, last_error, submit_timings) = ret;
+3 -2
View File
@@ -32,12 +32,13 @@ pub struct PumpFunParams {
/// `Some(PDA(["creator-vault", fee_sharing_config]))` when pump-fees `SharingConfig` is **Active**; set by `from_mint_by_rpc` / [`refresh_fee_sharing_creator_vault_from_rpc`](Self::refresh_fee_sharing_creator_vault_from_rpc).
pub fee_sharing_creator_vault_if_active: Option<Pubkey>,
/// SPL Token or Token-2022 program id owning the **mint** (from gRPC / parser / cache).
/// **`Pubkey::default()`**ix 构建时使用 SDK 默认 **Token-2022**(与多数 Pump.fun 新发一致);显式传入 Legacy 或 Token-2022 id 可覆盖该默认值
/// **`Pubkey::default()`**ix 构建时使用 SDK 默认 **Token-2022**(与多数 Pump.fun 新发一致)。
/// 显式传入 Legacy 或 Token-2022 id 时严格按该值组装,不再用 mint 字符串后缀猜测。
pub token_program: Pubkey,
/// Whether to close token account when selling, only effective during sell operations
pub close_token_account_when_sell: Option<bool>,
/// Fee recipient for buy/sell account #2. Set from sol-parser-sdk (`tradeEvent.feeRecipient` / 同笔 create_v2+buy 回填的 `observed_fee_recipient`);热路径不查 RPC。
/// `Pubkey::default()` 时按 mayhem 从静态池随机(与 npm 静态池一致,可能落后于主网 Global)
/// `Pubkey::default()` 时只能使用 SDK 静态 fallback,可能落后于主网 Global;交易热路径应优先传入 gRPC / parser 观测值
pub fee_recipient: Pubkey,
/// Quote mint for v2 instructions (default: `So11111111111111111111111111111111111111112` for SOL-paired).
/// For USDC-paired coins, set to `EPjFWdd5AufqSSqeM2qN1xzybapC8G4wEGGkZwyTDt1v`.
+5 -5
View File
@@ -17,6 +17,7 @@ const PARALLEL_SENDER_COUNT: usize = 18;
/// 启动时预填充数量,必须 >= PARALLEL_SENDER_COUNT,否则 18 路并发 build 会触发分配或争抢
const TX_BUILDER_POOL_PREFILL: usize = 64;
use anyhow::Result;
use crossbeam_queue::ArrayQueue;
use once_cell::sync::Lazy;
use solana_message::AddressLookupTableAccount;
@@ -82,7 +83,7 @@ impl PreallocatedTxBuilder {
instructions: &[Instruction],
address_lookup_table_account: Option<&AddressLookupTableAccount>,
recent_blockhash: Hash,
) -> VersionedMessage {
) -> Result<VersionedMessage> {
self.reset();
self.instructions.extend_from_slice(instructions);
@@ -92,14 +93,13 @@ impl PreallocatedTxBuilder {
&self.instructions,
std::slice::from_ref(alt),
recent_blockhash,
)
.expect("v0 message compile failed");
VersionedMessage::V0(message)
)?;
Ok(VersionedMessage::V0(message))
} else {
// ✅ 没有查找表,使用 Legacy 消息(兼容所有 RPC
let message =
Message::new_with_blockhash(&self.instructions, Some(payer), &recent_blockhash);
VersionedMessage::Legacy(message)
Ok(VersionedMessage::Legacy(message))
}
}
}