Merge pull request #103 from HelvetiCrypt/wait-for-all-submits

Add SwapParams.wait_for_all_submits to restore all-routes signature collection
This commit is contained in:
Wood
2026-05-28 02:46:39 +08:00
committed by GitHub
24 changed files with 84 additions and 7 deletions
+1
View File
@@ -179,6 +179,7 @@ async fn pumpfun_copy_trade_with_grpc(
)),
address_lookup_table_account,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: false,
close_input_token_ata: false,
create_mint_ata: true,
+2
View File
@@ -171,6 +171,7 @@ async fn bonk_copy_trade_with_grpc(trade_info: BonkTradeEvent) -> AnyResult<()>
)),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: true,
close_input_token_ata: false,
create_mint_ata: true,
@@ -222,6 +223,7 @@ async fn bonk_copy_trade_with_grpc(trade_info: BonkTradeEvent) -> AnyResult<()>
)),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
with_tip: false,
durable_nonce: None,
create_output_token_ata: false,
+2
View File
@@ -139,6 +139,7 @@ async fn bonk_sniper_trade_with_shreds(trade_info: BonkTradeEvent) -> AnyResult<
)),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: true,
close_input_token_ata: true,
create_mint_ata: true,
@@ -183,6 +184,7 @@ async fn bonk_sniper_trade_with_shreds(trade_info: BonkTradeEvent) -> AnyResult<
)),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: true,
close_output_token_ata: true,
close_mint_token_ata: false,
+10
View File
@@ -628,6 +628,7 @@ async fn handle_buy_pumpfun(
extension_params: DexParamEnum::PumpFun(param),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: false,
close_input_token_ata: false,
create_mint_ata: create_mint_ata,
@@ -684,6 +685,7 @@ async fn handle_buy_pumpswap(
extension_params: DexParamEnum::PumpSwap(param),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: true,
close_input_token_ata: false,
create_mint_ata: create_mint_ata,
@@ -740,6 +742,7 @@ async fn handle_buy_bonk(
extension_params: DexParamEnum::Bonk(param),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: true,
close_input_token_ata: false,
create_mint_ata: create_mint_ata,
@@ -800,6 +803,7 @@ async fn handle_buy_raydium_v4(
extension_params: DexParamEnum::RaydiumAmmV4(param),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: true,
close_input_token_ata: false,
create_mint_ata: create_mint_ata,
@@ -861,6 +865,7 @@ async fn handle_buy_raydium_cpmm(
extension_params: DexParamEnum::RaydiumCpmm(param),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: true,
close_input_token_ata: false,
create_mint_ata: create_mint_ata,
@@ -1031,6 +1036,7 @@ async fn handle_sell_pumpfun(
extension_params: DexParamEnum::PumpFun(param),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: true,
close_output_token_ata: false,
close_mint_token_ata: false,
@@ -1090,6 +1096,7 @@ async fn handle_sell_pumpswap(
extension_params: DexParamEnum::PumpSwap(param),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: true,
close_output_token_ata: false,
close_mint_token_ata: false,
@@ -1149,6 +1156,7 @@ async fn handle_sell_bonk(
extension_params: DexParamEnum::Bonk(param),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: true,
close_output_token_ata: false,
close_mint_token_ata: false,
@@ -1211,6 +1219,7 @@ async fn handle_sell_raydium_v4(
extension_params: DexParamEnum::RaydiumAmmV4(param),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: true,
close_output_token_ata: false,
close_mint_token_ata: false,
@@ -1274,6 +1283,7 @@ async fn handle_sell_raydium_cpmm(
extension_params: DexParamEnum::RaydiumCpmm(param),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: true,
close_output_token_ata: false,
close_mint_token_ata: false,
@@ -43,6 +43,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: false, //if input token is SOL/WSOL,set to true,if input token is USDC,set to false.
close_input_token_ata: false, //if input token is SOL/WSOL,set to true,if input token is USDC,set to false.
create_mint_ata: true,
@@ -84,6 +85,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: false, //if output token is SOL/WSOL,set to true,if output token is USDC,set to false.
close_output_token_ata: false, //if output token is SOL/WSOL,set to true,if output token is USDC,set to false.
close_mint_token_ata: false,
+1
View File
@@ -104,6 +104,7 @@ async fn test_middleware() -> AnyResult<()> {
),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: true,
close_input_token_ata: true,
create_mint_ata: true,
+1
View File
@@ -175,6 +175,7 @@ async fn pumpfun_copy_trade_with_grpc(
)),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: false,
close_input_token_ata: false,
create_mint_ata: true,
@@ -169,6 +169,7 @@ async fn pumpfun_copy_trade(e: sol_parser_sdk::core::events::PumpFunTradeEvent)
)),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: false,
close_input_token_ata: false,
create_mint_ata: true,
@@ -220,6 +221,7 @@ async fn pumpfun_copy_trade(e: sol_parser_sdk::core::events::PumpFunTradeEvent)
)),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: false,
close_output_token_ata: false,
close_mint_token_ata: false,
@@ -158,6 +158,7 @@ async fn pumpfun_sniper_trade(e: sol_parser_sdk::core::events::PumpFunTradeEvent
)),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: true,
close_input_token_ata: true,
create_mint_ata: true,
@@ -203,6 +204,7 @@ async fn pumpfun_sniper_trade(e: sol_parser_sdk::core::events::PumpFunTradeEvent
)),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: true,
close_output_token_ata: true,
close_mint_token_ata: false,
@@ -42,6 +42,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: true,
close_input_token_ata: true,
create_mint_ata: true,
@@ -81,6 +82,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: true,
close_output_token_ata: true,
close_mint_token_ata: false,
+2
View File
@@ -236,6 +236,7 @@ async fn pumpswap_trade_with_grpc(
extension_params: DexParamEnum::PumpSwap(params.clone()),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: is_sol,
close_input_token_ata: is_sol,
create_mint_ata: true,
@@ -277,6 +278,7 @@ async fn pumpswap_trade_with_grpc(
extension_params: DexParamEnum::PumpSwap(params.clone()),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: is_sol,
close_output_token_ata: is_sol,
close_mint_token_ata: false,
@@ -157,6 +157,7 @@ async fn raydium_amm_v4_copy_trade_with_grpc(trade_info: RaydiumAmmV4SwapEvent)
extension_params: DexParamEnum::RaydiumAmmV4(params),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: is_wsol,
close_input_token_ata: is_wsol,
create_mint_ata: true,
@@ -199,6 +200,7 @@ async fn raydium_amm_v4_copy_trade_with_grpc(trade_info: RaydiumAmmV4SwapEvent)
extension_params: DexParamEnum::RaydiumAmmV4(params),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: is_wsol,
close_output_token_ata: is_wsol,
close_mint_token_ata: false,
@@ -161,6 +161,7 @@ async fn raydium_cpmm_copy_trade_with_grpc(trade_info: RaydiumCpmmSwapEvent) ->
extension_params: DexParamEnum::RaydiumCpmm(buy_params),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: is_wsol,
close_input_token_ata: is_wsol,
create_mint_ata: true,
@@ -201,6 +202,7 @@ async fn raydium_cpmm_copy_trade_with_grpc(trade_info: RaydiumCpmmSwapEvent) ->
extension_params: DexParamEnum::RaydiumCpmm(sell_params),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: is_wsol,
close_output_token_ata: is_wsol,
close_mint_token_ata: false,
+2
View File
@@ -42,6 +42,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_input_token_ata: true,
close_input_token_ata: true,
create_mint_ata: true,
@@ -84,6 +85,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
),
address_lookup_table_account: None,
wait_tx_confirmed: true,
wait_for_all_submits: false,
create_output_token_ata: true,
close_output_token_ata: true,
close_mint_token_ata: false,
+12
View File
@@ -362,6 +362,11 @@ pub struct TradeBuyParams {
pub address_lookup_table_account: Option<AddressLookupTableAccount>,
/// Whether to wait for transaction confirmation before returning
pub wait_tx_confirmed: bool,
/// Fast-submit only (`wait_tx_confirmed = false`): when true, wait for every
/// SWQOS route's HTTP submit response so all submitted signatures are
/// returned. Set to true when confirming externally against a pinned
/// durable nonce; defaults to false. See `SwapParams.wait_for_all_submits`.
pub wait_for_all_submits: bool,
/// Whether to create input token associated token account
pub create_input_token_ata: bool,
/// Whether to close input token associated token account after trade
@@ -414,6 +419,11 @@ pub struct TradeSellParams {
pub address_lookup_table_account: Option<AddressLookupTableAccount>,
/// Whether to wait for transaction confirmation before returning
pub wait_tx_confirmed: bool,
/// Fast-submit only (`wait_tx_confirmed = false`): when true, wait for every
/// SWQOS route's HTTP submit response so all submitted signatures are
/// returned. Set to true when confirming externally against a pinned
/// durable nonce; defaults to false. See `SwapParams.wait_for_all_submits`.
pub wait_for_all_submits: bool,
/// Whether to create output token associated token account
pub create_output_token_ata: bool,
/// Whether to close output token associated token account after trade
@@ -868,6 +878,7 @@ impl TradingClient {
gas_fee_strategy: params.gas_fee_strategy,
simulate: params.simulate,
log_enabled: self.log_enabled,
wait_for_all_submits: params.wait_for_all_submits,
use_dedicated_sender_threads: self.use_dedicated_sender_threads,
sender_thread_cores: self.sender_thread_cores.clone(),
max_sender_concurrency: self.max_sender_concurrency,
@@ -983,6 +994,7 @@ impl TradingClient {
gas_fee_strategy: params.gas_fee_strategy,
simulate: params.simulate,
log_enabled: self.log_enabled,
wait_for_all_submits: params.wait_for_all_submits,
use_dedicated_sender_threads: self.use_dedicated_sender_threads,
sender_thread_cores: self.sender_thread_cores.clone(),
max_sender_concurrency: self.max_sender_concurrency,
+1
View File
@@ -382,6 +382,7 @@ mod tests {
gas_fee_strategy: GasFeeStrategy::new(),
simulate: true,
log_enabled: false,
wait_for_all_submits: false,
use_dedicated_sender_threads: false,
sender_thread_cores: None,
max_sender_concurrency: 0,
+1
View File
@@ -364,6 +364,7 @@ mod tests {
gas_fee_strategy: GasFeeStrategy::new(),
simulate: true,
log_enabled: false,
wait_for_all_submits: false,
use_dedicated_sender_threads: false,
sender_thread_cores: None,
max_sender_concurrency: 0,
+1
View File
@@ -919,6 +919,7 @@ mod tests {
gas_fee_strategy: GasFeeStrategy::new(),
simulate: false,
log_enabled: false,
wait_for_all_submits: false,
use_dedicated_sender_threads: false,
sender_thread_cores: None,
max_sender_concurrency: 0,
+1
View File
@@ -589,6 +589,7 @@ mod tests {
gas_fee_strategy: GasFeeStrategy::new(),
simulate: true,
log_enabled: false,
wait_for_all_submits: false,
use_dedicated_sender_threads: false,
sender_thread_cores: None,
max_sender_concurrency: 0,
+1
View File
@@ -382,6 +382,7 @@ mod tests {
gas_fee_strategy: GasFeeStrategy::new(),
simulate: true,
log_enabled: false,
wait_for_all_submits: false,
use_dedicated_sender_threads: false,
sender_thread_cores: None,
max_sender_concurrency: 0,
+1
View File
@@ -394,6 +394,7 @@ mod tests {
gas_fee_strategy: GasFeeStrategy::new(),
simulate: true,
log_enabled: false,
wait_for_all_submits: false,
use_dedicated_sender_threads: false,
sender_thread_cores: None,
max_sender_concurrency: 0,
+20 -7
View File
@@ -496,7 +496,8 @@ impl ResultCollector {
/// 等待全部任务完成(不等待链上确认),然后收集并返回所有签名。用于「多路提交」时返回多笔签名。
/// 轮询间隔 2ms,避免 50ms 间隔在最后一笔返回时多等几十 ms 拉高 submit 耗时。
#[allow(dead_code)]
/// Re-enabled via `SwapParams.wait_for_all_submits` for callers that confirm
/// externally against a pinned durable nonce and need every submitted sig.
async fn wait_for_all_submitted(
&self,
timeout_secs: u64,
@@ -600,6 +601,7 @@ pub async fn execute_parallel(
protocol_name: &'static str,
is_buy: bool,
wait_transaction_confirmed: bool,
wait_for_all_submits: bool,
with_tip: bool,
gas_fee_strategy: GasFeeStrategy,
use_dedicated_sender_threads: bool,
@@ -735,12 +737,23 @@ pub async fn execute_parallel(
// All jobs enqueued (no spawn on hot path)
if !wait_transaction_confirmed {
let ret = collector.wait_for_first_submitted(FAST_SUBMIT_RESULT_TIMEOUT).await.unwrap_or((
false,
vec![],
Some(anyhow!("No SWQOS result within submit result window")),
vec![],
));
let ret = if wait_for_all_submits {
collector.wait_for_all_submitted(FAST_SUBMIT_RESULT_TIMEOUT.as_secs()).await.unwrap_or(
(
false,
vec![],
Some(anyhow!("No SWQOS result within submit result window")),
vec![],
),
)
} else {
collector.wait_for_first_submitted(FAST_SUBMIT_RESULT_TIMEOUT).await.unwrap_or((
false,
vec![],
Some(anyhow!("No SWQOS result within submit result window")),
vec![],
))
};
let (success, signatures, last_error, submit_timings) = ret;
return Ok((success, signatures, last_error, submit_timings));
}
+5
View File
@@ -157,6 +157,10 @@ impl TradeExecutor for GenericTradeExecutor {
}
let need_confirm = params.wait_tx_confirmed;
// When the caller confirms externally (need_confirm = false) and opts in
// via SwapParams.wait_for_all_submits, return every route's signature so
// pinned-nonce confirmation can poll all of them.
let wait_for_all_submits = !need_confirm && params.wait_for_all_submits;
let sender_config = params.sender_concurrency_config();
let result = execute_parallel(
params.swqos_clients.as_slice(),
@@ -169,6 +173,7 @@ impl TradeExecutor for GenericTradeExecutor {
self.protocol_name,
is_buy,
false, // submit only here; confirmation and log timing handled below
wait_for_all_submits,
if is_buy { true } else { params.with_tip },
params.gas_fee_strategy,
params.use_dedicated_sender_threads,
+8
View File
@@ -82,6 +82,14 @@ pub struct SwapParams {
pub simulate: bool,
/// Whether to output SDK logs (from TradeConfig.log_enabled).
pub log_enabled: bool,
/// Fast-submit only (`wait_tx_confirmed = false`): when true, wait for every
/// SWQOS route's HTTP submit response so all submitted signatures are
/// returned. Defaults to false (post-4.0.11 behaviour: return after the
/// first route accepts). Set to true when the caller does its own on-chain
/// confirmation against a pinned durable nonce — only one route's tx can
/// land, but the caller cannot know which in advance, so it needs every
/// signature to poll via `getSignatureStatuses`.
pub wait_for_all_submits: bool,
/// Use dedicated sender threads (internal; set via client.with_dedicated_sender_threads()).
pub use_dedicated_sender_threads: bool,
/// Core indices for dedicated sender threads (from TradeConfig.sender_thread_cores). Arc avoids cloning the Vec on hot path.