From 807b015fc38f2af09280eb7c8fb219c7cd0d649b Mon Sep 17 00:00:00 2001 From: Wood Date: Wed, 25 Feb 2026 01:25:42 +0800 Subject: [PATCH] Release v3.5.0: performance, constants, bilingual docs MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Bump version to 3.5.0 - Performance: hot-path timing only when log_enabled/simulate; execute_parallel takes &[Arc]; shared HTTP client constants for SWQoS - Code quality: validate_protocol_params extracted for buy/sell; BYTES_PER_ACCOUNT, MAX_INSTRUCTIONS_WARN, HTTP timeout constants; prefetch/syscall bypass comments - Documentation: bilingual (EN + 中文) doc comments in execution, executor, perf, swqos; README/README_CN version and What's new in 3.5.0 - Add release_notes_v3.5.0.md Co-authored-by: Cursor --- .github/workflows/release.yml | 35 +++++ CHANGELOG_CN.md | 21 --- Cargo.toml | 3 +- README.md | 12 +- README_CN.md | 12 +- release_notes_v3.4.1.md | 40 +++++ release_notes_v3.5.0.md | 38 +++++ src/common/clock.rs | 97 ++++++++++++ src/common/fast_timing.rs | 8 +- src/common/gas_fee_strategy.rs | 3 + src/common/mod.rs | 4 +- src/common/sdk_log.rs | 19 +++ src/common/types.rs | 16 +- src/lib.rs | 183 +++++++++++++--------- src/perf/compiler_optimization.rs | 4 +- src/perf/hardware_optimizations.rs | 117 +++++--------- src/perf/mod.rs | 10 +- src/perf/syscall_bypass.rs | 86 ++++------ src/perf/zero_copy_io.rs | 24 +-- src/swqos/blockrazor.rs | 28 ++-- src/swqos/bloxroute.rs | 45 +++--- src/swqos/common.rs | 30 ++++ src/swqos/node1.rs | 46 +++--- src/swqos/serialization.rs | 45 +++++- src/swqos/speedlanding.rs | 27 ++-- src/swqos/stellium.rs | 44 +++--- src/trading/common/transaction_builder.rs | 9 +- src/trading/core/async_executor.rs | 102 +++++++----- src/trading/core/execution.rs | 41 +++-- src/trading/core/executor.rs | 161 ++++++++++--------- src/trading/core/params.rs | 6 + 31 files changed, 807 insertions(+), 509 deletions(-) create mode 100644 .github/workflows/release.yml delete mode 100644 CHANGELOG_CN.md create mode 100644 release_notes_v3.4.1.md create mode 100644 release_notes_v3.5.0.md create mode 100644 src/common/clock.rs create mode 100644 src/common/sdk_log.rs diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml new file mode 100644 index 0000000..9414fc1 --- /dev/null +++ b/.github/workflows/release.yml @@ -0,0 +1,35 @@ +# 推送 tag(如 v3.4.1)时自动创建 GitHub Release +name: Release + +on: + push: + tags: + - 'v*' + +permissions: + contents: write + +jobs: + release: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + with: + fetch-depth: 0 + + - name: Get version from tag + id: tag + run: echo "VERSION=${GITHUB_REF#refs/tags/v}" >> $GITHUB_OUTPUT + + - name: Create Release + uses: softprops/action-gh-release@v2 + with: + name: v${{ steps.tag.outputs.VERSION }} + body: | + ## sol-trade-sdk ${{ steps.tag.outputs.VERSION }} + Rust SDK to interact with the dex trade Solana program (Pump.fun, Raydium, etc.). + - **Cargo**: `sol-trade-sdk = { git = "https://github.com/${{ github.repository }}", tag = "v${{ steps.tag.outputs.VERSION }}" }` + draft: false + generate_release_notes: true + env: + GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} diff --git a/CHANGELOG_CN.md b/CHANGELOG_CN.md deleted file mode 100644 index 5765952..0000000 --- a/CHANGELOG_CN.md +++ /dev/null @@ -1,21 +0,0 @@ -# 更新日志 - -## [3.3.6] - 2025-01-30 - -### 新增 -- **Stellium SWQOS 支持**:全新 Stellium 客户端实现 - - 使用标准 Solana `sendTransaction` RPC 格式 - - 自动连接保活,60 秒 ping 间隔 - - 5 个小费账户用于负载分配 - - 支持 8 个区域端点(纽约、法兰克福、阿姆斯特丹、东京、伦敦等) - - 最低小费要求:0.001 SOL - -### 变更 -- **更新最低小费要求**以提高交易成功率: - - NextBlock: 0.00001 → 0.001 SOL - - ZeroSlot: 0.00001 → 0.001 SOL - - Temporal: 0.00001 → 0.001 SOL - - BloxRoute: 0.00001 → 0.001 SOL - - FlashBlock: 0.00001 → 0.001 SOL - - BlockRazor: 0.00001 → 0.001 SOL -- 增强异步执行器,添加小费验证警告 diff --git a/Cargo.toml b/Cargo.toml index 7796f79..6f5b53f 100755 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "sol-trade-sdk" -version = "3.4.1" +version = "3.5.0" edition = "2021" authors = [ "William ", @@ -81,7 +81,6 @@ rustls = { version = "0.23.23", features = ["ring"] } rustls-native-certs = "0.8.1" tokio-rustls = "0.26.1" core_affinity = "0.8" -log = "0.4.22" chrono = "0.4.39" regex = "1" tracing = "0.1.41" diff --git a/README.md b/README.md index 709c4e8..77296c3 100644 --- a/README.md +++ b/README.md @@ -60,6 +60,14 @@ --- +## 🆕 What's new in 3.5.0 + +- **Performance**: Hot-path timing only when logging; reduced clones (`execute_parallel` takes `&[Arc]`); shared HTTP client constants for SWQoS. +- **Code quality**: Extracted `validate_protocol_params` for buy/sell; constants for instruction/account sizes and HTTP timeouts; prefetch and branch-hint comments. +- **Documentation**: Bilingual (English + 中文) doc comments across execution, executor, perf, and swqos modules. + +--- + ## ✨ Features 1. **PumpFun Trading**: Support for `buy` and `sell` operations @@ -89,14 +97,14 @@ Add the dependency to your `Cargo.toml`: ```toml # Add to your Cargo.toml -sol-trade-sdk = { path = "./sol-trade-sdk", version = "3.4.1" } +sol-trade-sdk = { path = "./sol-trade-sdk", version = "3.5.0" } ``` ### Use crates.io ```toml # Add to your Cargo.toml -sol-trade-sdk = "3.4.1" +sol-trade-sdk = "3.5.0" ``` ## 🛠️ Usage Examples diff --git a/README_CN.md b/README_CN.md index 1273a99..bb1952d 100755 --- a/README_CN.md +++ b/README_CN.md @@ -60,6 +60,14 @@ --- +## 🆕 3.5.0 更新说明 + +- **性能**:仅在打日志时做热路径计时;减少 clone(`execute_parallel` 改为接收 `&[Arc]`);SWQoS 共用 HTTP 客户端常量。 +- **代码质量**:抽取 buy/sell 共用的 `validate_protocol_params`;指令/账户大小与 HTTP 超时常量化;预取与分支提示注释完善。 +- **文档**:execution、executor、perf、swqos 等模块增加中英双语文档注释。 + +--- + ## ✨ 项目特性 1. **PumpFun 交易**: 支持`购买`、`卖出`功能 @@ -89,14 +97,14 @@ git clone https://github.com/0xfnzero/sol-trade-sdk ```toml # 添加到您的 Cargo.toml -sol-trade-sdk = { path = "./sol-trade-sdk", version = "3.4.1" } +sol-trade-sdk = { path = "./sol-trade-sdk", version = "3.5.0" } ``` ### 使用 crates.io ```toml # 添加到您的 Cargo.toml -sol-trade-sdk = "3.4.1" +sol-trade-sdk = "3.5.0" ``` ## 🛠️ 使用示例 diff --git a/release_notes_v3.4.1.md b/release_notes_v3.4.1.md new file mode 100644 index 0000000..9d9c3e9 --- /dev/null +++ b/release_notes_v3.4.1.md @@ -0,0 +1,40 @@ +# sol-trade-sdk v3.4.1 + +Rust SDK for Solana DEX trading (Pump.fun, PumpSwap, Raydium, Bonk, Meteora, etc.). + +## What's Changed + +### New Features + +- **PumpFun & PumpSwap Cashback** (#77): Support for cashback in PumpFun and PumpSwap trading flows. See [Cashback documentation](docs/PUMP_CASHBACK_README.md). +- **Events**: `is_cashback_coin` is now passed from events; PumpFun examples use `sol-parser-sdk` only for event parsing. + +### Performance + +- **SWQoS**: Reduced submit latency; fixed high latency after ~5 minutes idle. + +### Bug Fixes + +- **WSOL ATA**: WSOL Associated Token Account creation now runs in background with retry and timeout for more reliable setup. +- Silenced unused and deprecated compiler warnings. + +### Documentation + +- README (EN/中文): Added Cashback section, outline, examples and tables. +- Updated crates.io / docs references in README. + +--- + +## Cargo + +**From Git (this release):** +```toml +sol-trade-sdk = { git = "https://github.com/0xfnzero/sol-trade-sdk", tag = "v3.4.1" } +``` + +**From crates.io** (when published): +```toml +sol-trade-sdk = "3.4.1" +``` + +**Full Changelog**: https://github.com/0xfnzero/sol-trade-sdk/compare/v3.4.0...v3.4.1 diff --git a/release_notes_v3.5.0.md b/release_notes_v3.5.0.md new file mode 100644 index 0000000..6bd3602 --- /dev/null +++ b/release_notes_v3.5.0.md @@ -0,0 +1,38 @@ +# sol-trade-sdk v3.5.0 + +Rust SDK for Solana DEX trading (Pump.fun, PumpSwap, Raydium, Bonk, Meteora, etc.). + +## What's Changed + +### Performance + +- **Hot-path timing**: `Instant::now()` for build/submit/total/confirm only when `log_enabled` or (for total) simulate, reducing cold-path syscalls. +- **Fewer clones**: `execute_parallel` now takes `&[Arc]` instead of `Vec>`; caller no longer clones the client list. +- **SWQoS HTTP**: Named constants for pool idle timeout, connect/request timeouts, and HTTP/2 keepalive in `swqos/common.rs`. + +### Code quality + +- **Protocol params**: Single `validate_protocol_params(dex_type, params)` used by both buy and sell; removed duplicated match blocks. +- **Constants**: `BYTES_PER_ACCOUNT`, `MAX_INSTRUCTIONS_WARN` in execution; HTTP timeout constants in swqos common. +- **Comments**: Prefetch and branch-hint safety/usage documented; `SYSCALL_BYPASS` marked as reserved for future use. + +### Documentation + +- **Bilingual docs**: English + 中文 doc comments in `trading/core/execution.rs`, `trading/core/executor.rs`, `perf/hardware_optimizations.rs`, `perf/mod.rs`, `perf/syscall_bypass.rs`, `swqos/common.rs`. +- **README**: Version references and "What's new in 3.5.0" (EN) / "3.5.0 更新说明" (CN) updated. + +--- + +## Cargo + +**From Git (this release):** +```toml +sol-trade-sdk = { git = "https://github.com/0xfnzero/sol-trade-sdk", tag = "v3.5.0" } +``` + +**From crates.io** (when published): +```toml +sol-trade-sdk = "3.5.0" +``` + +**Full Changelog**: https://github.com/0xfnzero/sol-trade-sdk/compare/v3.4.1...v3.5.0 diff --git a/src/common/clock.rs b/src/common/clock.rs new file mode 100644 index 0000000..916aed9 --- /dev/null +++ b/src/common/clock.rs @@ -0,0 +1,97 @@ +//! High-performance clock (same design as sol-parser-sdk for consistent grpc_recv_us vs "now"). +//! +//! Uses monotonic clock + base UTC timestamp to avoid frequent syscalls; aligned with sol-parser-sdk +//! so event-side grpc_recv_us and SDK-side now_micros() share the same time scale. + +use std::time::Instant; + +/// High-performance clock: monotonic + base UTC microsecond timestamp. +#[derive(Debug)] +pub struct HighPerformanceClock { + base_instant: Instant, + base_timestamp_us: i64, + last_calibration: Instant, + calibration_interval_secs: u64, +} + +impl HighPerformanceClock { + /// Calibrate every 5 minutes by default. + pub fn new() -> Self { + Self::new_with_calibration_interval(300) + } + + /// Sample multiple times and use the lowest-latency baseline to reduce init error. + pub fn new_with_calibration_interval(calibration_interval_secs: u64) -> Self { + let mut best_offset = i64::MAX; + let mut best_instant = Instant::now(); + let mut best_timestamp = chrono::Utc::now().timestamp_micros(); + + for _ in 0..3 { + let instant_before = Instant::now(); + let timestamp = chrono::Utc::now().timestamp_micros(); + let instant_after = Instant::now(); + let sample_latency = instant_after.duration_since(instant_before).as_nanos() as i64; + if sample_latency < best_offset { + best_offset = sample_latency; + best_instant = instant_before; + best_timestamp = timestamp; + } + } + + Self { + base_instant: best_instant, + base_timestamp_us: best_timestamp, + last_calibration: best_instant, + calibration_interval_secs, + } + } + + #[inline(always)] + pub fn now_micros(&self) -> i64 { + let elapsed = self.base_instant.elapsed(); + self.base_timestamp_us + elapsed.as_micros() as i64 + } + + /// Recalibrate when needed to prevent drift. + pub fn now_micros_with_calibration(&mut self) -> i64 { + if self.last_calibration.elapsed().as_secs() >= self.calibration_interval_secs { + self.recalibrate(); + } + self.now_micros() + } + + fn recalibrate(&mut self) { + let current_monotonic = Instant::now(); + let current_utc = chrono::Utc::now().timestamp_micros(); + let expected_utc = self.base_timestamp_us + + current_monotonic.duration_since(self.base_instant).as_micros() as i64; + let drift_us = current_utc - expected_utc; + if drift_us.abs() > 1000 { + self.base_instant = current_monotonic; + self.base_timestamp_us = current_utc; + } + self.last_calibration = current_monotonic; + } +} + +impl Default for HighPerformanceClock { + fn default() -> Self { + Self::new() + } +} + +static HIGH_PERF_CLOCK: once_cell::sync::OnceCell = + once_cell::sync::OnceCell::new(); + +/// Current time in microseconds (UTC scale); same as sol-parser-sdk clock::now_micros for comparable grpc_recv_us. +#[inline(always)] +pub fn now_micros() -> i64 { + let clock = HIGH_PERF_CLOCK.get_or_init(HighPerformanceClock::new); + clock.now_micros() +} + +/// Elapsed microseconds from start_timestamp_us to now. +#[inline(always)] +pub fn elapsed_micros_since(start_timestamp_us: i64) -> i64 { + now_micros() - start_timestamp_us +} diff --git a/src/common/fast_timing.rs b/src/common/fast_timing.rs index 8bb9c5a..9b9cf3e 100644 --- a/src/common/fast_timing.rs +++ b/src/common/fast_timing.rs @@ -174,7 +174,9 @@ mod tests { let total_elapsed = start.elapsed(); let avg_per_call = total_elapsed.as_nanos() / iterations; - println!("Average fast_now_nanos() call: {}ns", avg_per_call); + if crate::common::sdk_log::sdk_log_enabled() { + println!("Average fast_now_nanos() call: {}ns", avg_per_call); + } // 快速时间戳应该非常快(< 100ns per call) assert!(avg_per_call < 100); @@ -193,6 +195,8 @@ mod tests { let total_elapsed = start.elapsed(); let avg_per_call = total_elapsed.as_nanos() / iterations; - println!("Average Instant::now() call: {}ns", avg_per_call); + if crate::common::sdk_log::sdk_log_enabled() { + println!("Average Instant::now() call: {}ns", avg_per_call); + } } } diff --git a/src/common/gas_fee_strategy.rs b/src/common/gas_fee_strategy.rs index 5ace5c9..a0172a8 100644 --- a/src/common/gas_fee_strategy.rs +++ b/src/common/gas_fee_strategy.rs @@ -354,6 +354,9 @@ impl GasFeeStrategy { /// 打印所有策略。 /// Print all strategies pub fn print_all_strategies(&self) { + if !crate::common::sdk_log::sdk_log_enabled() { + return; + } for strategy in self.get_strategies(TradeType::Buy) { println!("[buy] - {:?}", strategy); } diff --git a/src/common/mod.rs b/src/common/mod.rs index 8fc314a..24b25d9 100755 --- a/src/common/mod.rs +++ b/src/common/mod.rs @@ -1,5 +1,8 @@ +pub mod address_lookup; pub mod bonding_curve; +pub mod clock; pub mod fast_fn; +pub mod sdk_log; pub mod fast_timing; pub mod gas_fee_strategy; pub mod global; @@ -10,7 +13,6 @@ pub mod spl_token; pub mod spl_token_2022; pub mod subscription_handle; pub mod types; -pub mod address_lookup; pub use gas_fee_strategy::*; pub use types::*; diff --git a/src/common/sdk_log.rs b/src/common/sdk_log.rs new file mode 100644 index 0000000..96d4a7b --- /dev/null +++ b/src/common/sdk_log.rs @@ -0,0 +1,19 @@ +//! sol-trade-sdk global log switch +//! +//! Controlled by `TradeConfig::log_enabled`, set in `TradingClient::new`. +//! All SDK logs (timing, SWQOS submit/confirm, WSOL, blacklist, etc.) should check this before output. + +use std::sync::atomic::{AtomicBool, Ordering}; + +static SDK_LOG_ENABLED: AtomicBool = AtomicBool::new(true); + +/// Whether SDK logging is enabled (set from TradeConfig.log_enabled in TradingClient::new). +#[inline(always)] +pub fn sdk_log_enabled() -> bool { + SDK_LOG_ENABLED.load(Ordering::Relaxed) +} + +/// Set the SDK global log switch (only called from TradingClient::new). +pub fn set_sdk_log_enabled(enabled: bool) { + SDK_LOG_ENABLED.store(enabled, Ordering::Relaxed); +} diff --git a/src/common/types.rs b/src/common/types.rs index 7e1d9ce..475f879 100755 --- a/src/common/types.rs +++ b/src/common/types.rs @@ -72,6 +72,10 @@ pub struct TradeConfig { pub create_wsol_ata_on_startup: bool, /// Whether to use seed optimization for all ATA operations (default: true) pub use_seed_optimize: bool, + /// Whether to pin parallel submit tasks to CPU cores (can reduce latency; set false in containers). Default true. + pub use_core_affinity: bool, + /// Whether to output all SDK logs (timing, SWQOS submit/confirm, WSOL, blacklist, etc.). Default true. + pub log_enabled: bool, } impl TradeConfig { @@ -80,14 +84,18 @@ impl TradeConfig { swqos_configs: Vec, commitment: CommitmentConfig, ) -> Self { - println!("🔧 TradeConfig create_wsol_ata_on_startup default value: true"); - println!("🔧 TradeConfig use_seed_optimize default value: true"); + if crate::common::sdk_log::sdk_log_enabled() { + println!("🔧 TradeConfig create_wsol_ata_on_startup default: true"); + println!("🔧 TradeConfig use_seed_optimize default: true"); + } Self { rpc_url, swqos_configs, commitment, - create_wsol_ata_on_startup: true, // 默认:启动时检查并创建 - use_seed_optimize: true, // 默认:使用seed优化 + create_wsol_ata_on_startup: true, // default: check and create on startup + use_seed_optimize: true, // default: use seed optimization + use_core_affinity: true, // default: pin parallel submit tasks to cores + log_enabled: true, // default: enable all SDK logs } } diff --git a/src/lib.rs b/src/lib.rs index 7d1cafd..67597de 100755 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,6 +6,7 @@ pub mod swqos; pub mod trading; pub mod utils; use crate::common::nonce_cache::DurableNonceInfo; +use crate::common::sdk_log; use crate::common::GasFeeStrategy; use crate::common::{TradeConfig, InfrastructureConfig}; #[cfg(feature = "perf-trace")] @@ -37,6 +38,21 @@ use solana_sdk::message::AddressLookupTableAccount; use solana_sdk::signer::Signer; use solana_sdk::{pubkey::Pubkey, signature::Keypair, signature::Signature}; use std::sync::Arc; +#[allow(unused_imports)] +use tracing::{debug, error, info, warn}; + +/// Single place to validate that protocol params match the given DEX type (avoids duplicate match in buy/sell). +#[inline(always)] +fn validate_protocol_params(dex_type: DexType, params: &DexParamEnum) -> bool { + match dex_type { + DexType::PumpFun => params.as_any().downcast_ref::().is_some(), + DexType::PumpSwap => params.as_any().downcast_ref::().is_some(), + DexType::Bonk => params.as_any().downcast_ref::().is_some(), + DexType::RaydiumCpmm => params.as_any().downcast_ref::().is_some(), + DexType::RaydiumAmmV4 => params.as_any().downcast_ref::().is_some(), + DexType::MeteoraDammV2 => params.as_any().downcast_ref::().is_some(), + } +} /// Type of the token to buy #[derive(Clone, PartialEq)] @@ -90,7 +106,9 @@ impl TradingInfrastructure { for swqos in &config.swqos_configs { // Check blacklist, skip disabled providers if swqos.is_blacklisted() { - eprintln!("\u{26a0}\u{fe0f} SWQOS {:?} is blacklisted, skipping", swqos.swqos_type()); + if sdk_log::sdk_log_enabled() { + warn!(target: "sol_trade_sdk", "⚠️ SWQOS {:?} is blacklisted, skipping", swqos.swqos_type()); + } continue; } match SwqosConfig::get_swqos_client( @@ -99,10 +117,15 @@ impl TradingInfrastructure { swqos.clone(), ).await { Ok(swqos_client) => swqos_clients.push(swqos_client), - Err(err) => eprintln!( - "failed to create {:?} swqos client: {err}. Excluding from swqos list", - swqos.swqos_type() - ), + Err(err) => { + if sdk_log::sdk_log_enabled() { + warn!( + target: "sol_trade_sdk", + "failed to create {:?} swqos client: {err}. Excluding from swqos list", + swqos.swqos_type() + ); + } + } } } @@ -130,6 +153,10 @@ pub struct TradingClient { /// Whether to use seed optimization for all ATA operations (default: true) /// Applies to all token account creations across buy and sell operations pub use_seed_optimize: bool, + /// Whether to pin parallel submit tasks to CPU cores (from TradeConfig.use_core_affinity). Default true. + pub use_core_affinity: bool, + /// Whether to output all SDK logs (from TradeConfig.log_enabled). + pub log_enabled: bool, } static INSTANCE: Mutex>> = Mutex::new(None); @@ -144,6 +171,8 @@ impl Clone for TradingClient { infrastructure: self.infrastructure.clone(), middleware_manager: self.middleware_manager.clone(), use_seed_optimize: self.use_seed_optimize, + use_core_affinity: self.use_core_affinity, + log_enabled: self.log_enabled, } } } @@ -193,6 +222,8 @@ pub struct TradeBuyParams { /// When Some(false), uses regular buy instruction where slippage is applied to SOL/quote input. /// This option only applies to PumpFun and PumpSwap DEXes; it is ignored for other DEXes. pub use_exact_sol_amount: Option, + /// 可选:事件收到时间(微秒,与 sol-parser-sdk 的 metadata.grpc_recv_us / clock::now_micros 同源)。不传且开启 log_enabled 时 SDK 用 now_micros() 作为起点,打印起点→提交耗时。 + pub grpc_recv_us: Option, } /// Parameters for executing sell orders across different DEX protocols @@ -237,6 +268,8 @@ pub struct TradeSellParams { pub gas_fee_strategy: GasFeeStrategy, /// Whether to simulate the transaction instead of executing it pub simulate: bool, + /// 可选:事件收到时间(微秒,与 sol-parser-sdk clock 同源)。不传且开启 log_enabled 时 SDK 用 now_micros() 作为起点。 + pub grpc_recv_us: Option, } impl TradingClient { @@ -266,6 +299,8 @@ impl TradingClient { infrastructure, middleware_manager: None, use_seed_optimize, + use_core_affinity: true, + log_enabled: true, } } @@ -293,7 +328,9 @@ impl TradingClient { tokio::spawn(async move { Self::ensure_wsol_ata(&payer_clone, &rpc_clone).await; }); - println!("ℹ️ WSOL ATA 创建已在后台启动,不阻塞机器人启动"); + if sdk_log::sdk_log_enabled() { + info!(target: "sol_trade_sdk", "ℹ️ WSOL ATA creation started in background, does not block bot startup"); + } } Self { @@ -301,6 +338,8 @@ impl TradingClient { infrastructure, middleware_manager: None, use_seed_optimize, + use_core_affinity: true, + log_enabled: true, } } @@ -315,11 +354,15 @@ impl TradingClient { match rpc.get_account(&wsol_ata).await { Ok(_) => { - println!("✅ WSOL ATA已存在: {}", wsol_ata); + if sdk_log::sdk_log_enabled() { + info!(target: "sol_trade_sdk", "✅ WSOL ATA already exists: {}", wsol_ata); + } return; } Err(_) => { - println!("🔨 创建WSOL ATA: {}", wsol_ata); + if sdk_log::sdk_log_enabled() { + info!(target: "sol_trade_sdk", "🔨 Creating WSOL ATA: {}", wsol_ata); + } let create_ata_ixs = crate::trading::common::wsol_manager::create_wsol_ata(&payer.pubkey()); @@ -333,15 +376,19 @@ impl TradingClient { for attempt in 1..=MAX_RETRIES { if attempt > 1 { - println!("🔄 重试创建WSOL ATA (第{}/{}次)...", attempt, MAX_RETRIES); + if sdk_log::sdk_log_enabled() { + info!(target: "sol_trade_sdk", "🔄 Retrying WSOL ATA creation (attempt {}/{})...", attempt, MAX_RETRIES); + } tokio::time::sleep(tokio::time::Duration::from_secs(2)).await; } let recent_blockhash = match rpc.get_latest_blockhash().await { Ok(hash) => hash, Err(e) => { - eprintln!("⚠️ 获取最新blockhash失败: {}", e); - last_error = Some(format!("获取blockhash失败: {}", e)); + if sdk_log::sdk_log_enabled() { + warn!(target: "sol_trade_sdk", "⚠️ Failed to get latest blockhash: {}", e); + } + last_error = Some(format!("Failed to get blockhash: {}", e)); continue; } }; @@ -361,7 +408,9 @@ impl TradingClient { match send_result { Ok(Ok(signature)) => { - println!("✅ WSOL ATA创建成功: {}", signature); + if sdk_log::sdk_log_enabled() { + info!(target: "sol_trade_sdk", "✅ WSOL ATA created successfully: {}", signature); + } return; } Ok(Err(e)) => { @@ -369,45 +418,50 @@ impl TradingClient { // 检查账户是否实际已存在 if let Ok(_) = rpc.get_account(&wsol_ata).await { - println!( - "✅ WSOL ATA已存在(交易失败但账户存在): {}", - wsol_ata - ); + if sdk_log::sdk_log_enabled() { + info!(target: "sol_trade_sdk", "✅ WSOL ATA already exists (tx failed but account exists): {}", wsol_ata); + } return; } - if attempt < MAX_RETRIES { - eprintln!("⚠️ 第{}次尝试失败: {}", attempt, e); + if attempt < MAX_RETRIES && sdk_log::sdk_log_enabled() { + warn!(target: "sol_trade_sdk", "⚠️ Attempt {} failed: {}", attempt, e); } } Err(_) => { - last_error = Some(format!("交易确认超时({}秒)", TIMEOUT_SECS)); - eprintln!("⚠️ 第{}次尝试超时", attempt); + last_error = Some(format!("Transaction confirmation timeout ({}s)", TIMEOUT_SECS)); + if sdk_log::sdk_log_enabled() { + warn!(target: "sol_trade_sdk", "⚠️ Attempt {} timed out", attempt); + } } } } // 所有重试都失败了 if let Some(err) = last_error { - eprintln!("❌ WSOL ATA创建失败(已重试{}次): {}", MAX_RETRIES, wsol_ata); - eprintln!(" 错误详情: {}", err); - eprintln!(" 💡 可能原因:"); - eprintln!(" 1. 钱包SOL余额不足(需要约0.002 SOL用于租金豁免)"); - eprintln!(" 2. RPC节点响应超时或网络拥堵"); - eprintln!(" 3. 交易费用不足"); - eprintln!(" 🔧 解决方案:"); - eprintln!(" 1. 给钱包充值至少0.1 SOL"); - eprintln!(" 2. 等待几秒后重试"); - eprintln!(" 3. 检查RPC节点连接"); - eprintln!(" ⚠️ 程序将在5秒后退出,请解决上述问题后重启"); + if sdk_log::sdk_log_enabled() { + error!(target: "sol_trade_sdk", "❌ WSOL ATA creation failed after {} retries: {}", MAX_RETRIES, wsol_ata); + error!(target: "sol_trade_sdk", " Error: {}", err); + error!(target: "sol_trade_sdk", " 💡 Possible causes:"); + error!(target: "sol_trade_sdk", " 1. Insufficient SOL balance (need ~0.002 SOL for rent exemption)"); + error!(target: "sol_trade_sdk", " 2. RPC timeout or network congestion"); + error!(target: "sol_trade_sdk", " 3. Insufficient transaction fee"); + error!(target: "sol_trade_sdk", " 🔧 Solutions:"); + error!(target: "sol_trade_sdk", " 1. Fund wallet with at least 0.1 SOL"); + error!(target: "sol_trade_sdk", " 2. Retry after a few seconds"); + error!(target: "sol_trade_sdk", " 3. Check RPC connection"); + error!(target: "sol_trade_sdk", " ⚠️ Process will exit in 5 seconds, restart after fixing the above"); + } std::thread::sleep(std::time::Duration::from_secs(5)); panic!( - "❌ WSOL ATA创建失败且账户不存在: {}. 错误: {}", + "❌ WSOL ATA creation failed and account does not exist: {}. Error: {}", wsol_ata, err ); } } else { - println!("ℹ️ WSOL ATA已存在(无需创建)"); + if sdk_log::sdk_log_enabled() { + info!(target: "sol_trade_sdk", "ℹ️ WSOL ATA already exists (no need to create)"); + } } } } @@ -426,6 +480,10 @@ impl TradingClient { /// Returns a configured `SolTradingSDK` instance ready for trading operations #[inline] pub async fn new(payer: Arc, trade_config: TradeConfig) -> Self { + // 设置 SDK 全局日志开关,后续所有 SDK 内日志(SWQOS/WSOL/耗时等)均受此控制 + sdk_log::set_sdk_log_enabled(trade_config.log_enabled); + // 预热高性能时钟,避免首笔交易时触发 3 次 Utc::now() 校准 + let _ = crate::common::clock::now_micros(); // Create infrastructure from trade config let infra_config = InfrastructureConfig::from_trade_config(&trade_config); let infrastructure = Arc::new(TradingInfrastructure::new(infra_config).await); @@ -443,6 +501,8 @@ impl TradingClient { infrastructure, middleware_manager: None, use_seed_optimize: trade_config.use_seed_optimize, + use_core_affinity: trade_config.use_core_affinity, + log_enabled: trade_config.log_enabled, }; let mut current = INSTANCE.lock(); @@ -525,8 +585,9 @@ impl TradingClient { params: TradeBuyParams, ) -> Result<(bool, Vec, Option), anyhow::Error> { #[cfg(feature = "perf-trace")] - if params.slippage_basis_points.is_none() { - log::debug!( + if sdk_log::sdk_log_enabled() && params.slippage_basis_points.is_none() { + debug!( + target: "sol_trade_sdk", "slippage_basis_points is none, use default slippage basis points: {}", DEFAULT_SLIPPAGE ); @@ -573,28 +634,13 @@ impl TradingClient { fixed_output_amount: params.fixed_output_token_amount, gas_fee_strategy: params.gas_fee_strategy, simulate: params.simulate, + log_enabled: self.log_enabled, + use_core_affinity: self.use_core_affinity, + grpc_recv_us: params.grpc_recv_us, use_exact_sol_amount: params.use_exact_sol_amount, }; - // Validate protocol params - let is_valid_params = match params.dex_type { - DexType::PumpFun => protocol_params.as_any().downcast_ref::().is_some(), - DexType::PumpSwap => { - protocol_params.as_any().downcast_ref::().is_some() - } - DexType::Bonk => protocol_params.as_any().downcast_ref::().is_some(), - DexType::RaydiumCpmm => { - protocol_params.as_any().downcast_ref::().is_some() - } - DexType::RaydiumAmmV4 => { - protocol_params.as_any().downcast_ref::().is_some() - } - DexType::MeteoraDammV2 => { - protocol_params.as_any().downcast_ref::().is_some() - } - }; - - if !is_valid_params { + if !validate_protocol_params(params.dex_type, &protocol_params) { return Err(anyhow::anyhow!("Invalid protocol params for Trade")); } @@ -635,8 +681,9 @@ impl TradingClient { params: TradeSellParams, ) -> Result<(bool, Vec, Option), anyhow::Error> { #[cfg(feature = "perf-trace")] - if params.slippage_basis_points.is_none() { - log::debug!( + if sdk_log::sdk_log_enabled() && params.slippage_basis_points.is_none() { + debug!( + target: "sol_trade_sdk", "slippage_basis_points is none, use default slippage basis points: {}", DEFAULT_SLIPPAGE ); @@ -683,32 +730,16 @@ impl TradingClient { fixed_output_amount: params.fixed_output_token_amount, gas_fee_strategy: params.gas_fee_strategy, simulate: params.simulate, + log_enabled: self.log_enabled, + use_core_affinity: self.use_core_affinity, + grpc_recv_us: params.grpc_recv_us, use_exact_sol_amount: None, }; - // Validate protocol params - let is_valid_params = match params.dex_type { - DexType::PumpFun => protocol_params.as_any().downcast_ref::().is_some(), - DexType::PumpSwap => { - protocol_params.as_any().downcast_ref::().is_some() - } - DexType::Bonk => protocol_params.as_any().downcast_ref::().is_some(), - DexType::RaydiumCpmm => { - protocol_params.as_any().downcast_ref::().is_some() - } - DexType::RaydiumAmmV4 => { - protocol_params.as_any().downcast_ref::().is_some() - } - DexType::MeteoraDammV2 => { - protocol_params.as_any().downcast_ref::().is_some() - } - }; - - if !is_valid_params { + if !validate_protocol_params(params.dex_type, &protocol_params) { return Err(anyhow::anyhow!("Invalid protocol params for Trade")); } - // Execute sell based on tip preference let swap_result = executor.swap(sell_params).await; let result = swap_result.map(|(success, sigs, err)| (success, sigs, err.map(TradeError::from))); diff --git a/src/perf/compiler_optimization.rs b/src/perf/compiler_optimization.rs index e5aa638..7c709d6 100644 --- a/src/perf/compiler_optimization.rs +++ b/src/perf/compiler_optimization.rs @@ -140,7 +140,7 @@ impl CompilerOptimizer { /// 🚀 生成超高性能编译配置 pub fn generate_ultra_performance_config(&self) -> Result { - log::info!("🚀 Generating ultra-performance compiler configuration..."); + tracing::info!(target: "sol_trade_sdk","🚀 Generating ultra-performance compiler configuration..."); let mut rustflags = Vec::new(); @@ -206,7 +206,7 @@ impl CompilerOptimizer { cargo_config: self.generate_cargo_config(), }; - log::info!("✅ Ultra-performance compiler configuration generated"); + tracing::info!(target: "sol_trade_sdk","✅ Ultra-performance compiler configuration generated"); Ok(config) } diff --git a/src/perf/hardware_optimizations.rs b/src/perf/hardware_optimizations.rs index fac4427..4966cfb 100644 --- a/src/perf/hardware_optimizations.rs +++ b/src/perf/hardware_optimizations.rs @@ -1,11 +1,5 @@ -//! 🚀 硬件级性能优化 - CPU缓存行对齐 & SIMD加速 -//! -//! 实现CPU硬件特性的深度利用,包括: -//! - 缓存行对齐和缓存预取 -//! - SIMD指令集优化 -//! - 分支预测优化 -//! - 内存屏障控制 -//! - CPU指令流水线优化 +//! Hardware-oriented optimizations: cache-line alignment, prefetch, SIMD, branch hints, memory barriers. +//! 硬件级优化:缓存行对齐与预取、SIMD、分支提示、内存屏障。 use std::sync::atomic::{AtomicU64, Ordering}; use std::mem::size_of; @@ -13,26 +7,23 @@ use std::ptr; use crossbeam_utils::CachePadded; use anyhow::Result; -// CPU缓存行大小常量 (通常为64字节) +/// Typical CPU cache line size in bytes. 典型 CPU 缓存行大小(字节)。 pub const CACHE_LINE_SIZE: usize = 64; -/// 🚀 硬件优化的数据结构基础特征 +/// Trait for cache-line-aligned data and prefetch. 缓存行对齐与预取 trait。 pub trait CacheLineAligned { - /// 确保数据结构按缓存行对齐 fn ensure_cache_aligned(&self) -> bool; - /// 预取数据到CPU缓存 fn prefetch_data(&self); } -/// 🚀 SIMD优化的内存操作 +/// SIMD-accelerated memory operations. SIMD 加速的内存操作。 pub struct SIMDMemoryOps; impl SIMDMemoryOps { - /// 🚀 SIMD加速的内存拷贝 - 针对小数据包优化 + /// SIMD-optimized copy by size class. 按长度分派的 SIMD 拷贝。 #[inline(always)] pub unsafe fn memcpy_simd_optimized(dst: *mut u8, src: *const u8, len: usize) { match len { - // 针对不同数据大小使用不同优化策略 0 => return, 1..=8 => Self::memcpy_small(dst, src, len), 9..=16 => Self::memcpy_sse(dst, src, len), @@ -42,7 +33,7 @@ impl SIMDMemoryOps { } } - /// 小数据拷贝优化 (1-8字节) + /// Copy 1–8 bytes (scalar / small word). 小数据拷贝(1–8 字节)。 #[inline(always)] unsafe fn memcpy_small(dst: *mut u8, src: *const u8, len: usize) { match len { @@ -63,7 +54,7 @@ impl SIMDMemoryOps { } } - /// SSE优化拷贝 (9-16字节) + /// Copy 9–16 bytes using SSE (128-bit). SSE 拷贝(9–16 字节)。 #[inline(always)] unsafe fn memcpy_sse(dst: *mut u8, src: *const u8, len: usize) { #[cfg(target_arch = "x86_64")] @@ -82,7 +73,7 @@ impl SIMDMemoryOps { } } - /// AVX优化拷贝 (17-32字节) + /// Copy 17–32 bytes using AVX (256-bit). AVX 拷贝(17–32 字节)。 #[inline(always)] unsafe fn memcpy_avx(dst: *mut u8, src: *const u8, len: usize) { #[cfg(target_arch = "x86_64")] @@ -101,19 +92,16 @@ impl SIMDMemoryOps { } } - /// AVX2优化拷贝 (33-64字节) + /// Copy 33–64 bytes using AVX2 (256-bit, two chunks). AVX2 拷贝(33–64 字节,两段)。 #[inline(always)] unsafe fn memcpy_avx2(dst: *mut u8, src: *const u8, len: usize) { #[cfg(target_arch = "x86_64")] { use std::arch::x86_64::{__m256i, _mm256_loadu_si256, _mm256_storeu_si256}; - // 拷贝前32字节 let chunk1 = _mm256_loadu_si256(src as *const __m256i); _mm256_storeu_si256(dst as *mut __m256i, chunk1); - if len > 32 { - // 拷贝剩余字节 let remaining = len - 32; if remaining <= 32 { let chunk2 = _mm256_loadu_si256(src.add(32) as *const __m256i); @@ -128,7 +116,7 @@ impl SIMDMemoryOps { } } - /// AVX512或回退拷贝 (>64字节) + /// Copy >64 bytes: AVX-512 64-byte chunks when available, else AVX2 32-byte chunks. >64 字节:有 AVX512 用 64 字节块,否则 AVX2 32 字节块。 #[inline(always)] unsafe fn memcpy_avx512_or_fallback(dst: *mut u8, src: *const u8, len: usize) { #[cfg(all(target_arch = "x86_64", target_feature = "avx512f"))] @@ -138,14 +126,12 @@ impl SIMDMemoryOps { let chunks = len / 64; let mut offset = 0; - // 使用AVX512处理64字节块 for _ in 0..chunks { let chunk = _mm512_loadu_si512(src.add(offset) as *const __m512i); _mm512_storeu_si512(dst.add(offset) as *mut __m512i, chunk); offset += 64; } - // 处理剩余字节 let remaining = len % 64; if remaining > 0 { Self::memcpy_avx2(dst.add(offset), src.add(offset), remaining); @@ -154,7 +140,6 @@ impl SIMDMemoryOps { #[cfg(not(all(target_arch = "x86_64", target_feature = "avx512f")))] { - // 回退到AVX2分块处理 let chunks = len / 32; let mut offset = 0; @@ -170,7 +155,7 @@ impl SIMDMemoryOps { } } - /// 🚀 SIMD加速的内存比较 + /// SIMD-optimized byte equality; dispatches by length (small / SSE / AVX2 / large). SIMD 加速的内存比较,按长度分派。 #[inline(always)] pub unsafe fn memcmp_simd_optimized(a: *const u8, b: *const u8, len: usize) -> bool { match len { @@ -182,7 +167,7 @@ impl SIMDMemoryOps { } } - /// 小数据比较 + /// Compare 1–8 bytes (scalar). 小数据比较(1–8 字节)。 #[inline(always)] unsafe fn memcmp_small(a: *const u8, b: *const u8, len: usize) -> bool { match len { @@ -198,7 +183,7 @@ impl SIMDMemoryOps { } } - /// SSE比较 + /// Compare 9–16 bytes using SSE. SSE 比较(9–16 字节)。 #[inline(always)] unsafe fn memcmp_sse(a: *const u8, b: *const u8, len: usize) -> bool { #[cfg(target_arch = "x86_64")] @@ -210,7 +195,6 @@ impl SIMDMemoryOps { let cmp_result = _mm_cmpeq_epi8(chunk_a, chunk_b); let mask = _mm_movemask_epi8(cmp_result) as u32; - // 检查前len字节是否相等 let valid_mask = if len >= 16 { 0xFFFF } else { (1u32 << len) - 1 }; (mask & valid_mask) == valid_mask } @@ -221,7 +205,7 @@ impl SIMDMemoryOps { } } - /// AVX2比较 + /// Compare 17–32 bytes using AVX2. AVX2 比较(17–32 字节)。 #[inline(always)] unsafe fn memcmp_avx2(a: *const u8, b: *const u8, len: usize) -> bool { #[cfg(target_arch = "x86_64")] @@ -243,7 +227,7 @@ impl SIMDMemoryOps { } } - /// 大数据比较 + /// Compare >32 bytes in 32-byte AVX2 chunks. 大数据比较(32 字节 AVX2 分块)。 #[inline(always)] unsafe fn memcmp_large(a: *const u8, b: *const u8, len: usize) -> bool { let chunks = len / 32; @@ -263,7 +247,7 @@ impl SIMDMemoryOps { true } - /// 🚀 SIMD加速的内存清零 + /// SIMD-optimized zero memory. SIMD 加速的内存清零。 #[inline(always)] pub unsafe fn memzero_simd_optimized(ptr: *mut u8, len: usize) { #[cfg(target_arch = "x86_64")] @@ -279,7 +263,6 @@ impl SIMDMemoryOps { offset += 32; } - // 处理剩余字节 let remaining = len % 32; for i in 0..remaining { *ptr.add(offset + i) = 0; @@ -293,14 +276,15 @@ impl SIMDMemoryOps { } } -/// 🚀 缓存行对齐的原子计数器 -#[repr(align(64))] // 强制64字节对齐 +/// Cache-line-aligned atomic counter. 缓存行对齐的原子计数器。 +#[repr(align(64))] pub struct CacheAlignedCounter { value: AtomicU64, _padding: [u8; CACHE_LINE_SIZE - size_of::()], } impl CacheAlignedCounter { + /// Create counter with initial value. 创建并设置初值。 pub fn new(initial: u64) -> Self { Self { value: AtomicU64::new(initial), @@ -339,23 +323,18 @@ impl CacheLineAligned for CacheAlignedCounter { } } -/// 🚀 缓存友好的环形缓冲区 +/// Cache-friendly lock-free ring buffer. 缓存友好的无锁环形缓冲区。 #[repr(align(64))] pub struct CacheOptimizedRingBuffer { - /// 数据缓冲区 buffer: Vec, - /// 生产者头指针 (独占缓存行) producer_head: CachePadded, - /// 消费者尾指针 (独占缓存行) consumer_tail: CachePadded, - /// 容量 (2的幂次方) capacity: usize, - /// 掩码 (capacity - 1) mask: usize, } impl CacheOptimizedRingBuffer { - /// 创建缓存优化的环形缓冲区 + /// Create ring buffer; capacity must be a power of 2. 创建环形缓冲区,容量须为 2 的幂。 pub fn new(capacity: usize) -> Result { if !capacity.is_power_of_two() { return Err(anyhow::anyhow!("Capacity must be a power of 2")); @@ -373,53 +352,41 @@ impl CacheOptimizedRingBuffer { }) } - /// 🚀 无锁写入元素 + /// Lock-free push; returns false if full. 无锁写入,满则返回 false。 #[inline(always)] pub fn try_push(&self, item: T) -> bool { let current_head = self.producer_head.load(Ordering::Relaxed); let current_tail = self.consumer_tail.load(Ordering::Acquire); - - // 检查是否还有空间 if (current_head + 1) & self.mask as u64 == current_tail & self.mask as u64 { - return false; // 缓冲区满 + return false; } - - // 写入数据 unsafe { let index = current_head & self.mask as u64; let ptr = self.buffer.as_ptr().add(index as usize) as *mut T; ptr.write(item); } - - // 发布新的头指针 self.producer_head.store(current_head + 1, Ordering::Release); true } - /// 🚀 无锁读取元素 + /// Lock-free pop; returns None if empty. 无锁读取,空则返回 None。 #[inline(always)] pub fn try_pop(&self) -> Option { let current_tail = self.consumer_tail.load(Ordering::Relaxed); let current_head = self.producer_head.load(Ordering::Acquire); - - // 检查是否有数据 if current_tail == current_head { - return None; // 缓冲区空 + return None; } - - // 读取数据 let item = unsafe { let index = current_tail & self.mask as u64; let ptr = self.buffer.as_ptr().add(index as usize); ptr.read() }; - - // 发布新的尾指针 self.consumer_tail.store(current_tail + 1, Ordering::Release); Some(item) } - /// 获取当前元素数量 + /// Current number of elements. 当前元素个数。 #[inline(always)] pub fn len(&self) -> usize { let head = self.producer_head.load(Ordering::Relaxed); @@ -427,7 +394,7 @@ impl CacheOptimizedRingBuffer { ((head + self.capacity as u64 - tail) & self.mask as u64) as usize } - /// 检查是否为空 + /// True if no elements. 是否为空。 #[inline(always)] pub fn is_empty(&self) -> bool { self.producer_head.load(Ordering::Relaxed) == @@ -445,24 +412,18 @@ impl CacheLineAligned for CacheOptimizedRingBuffer { unsafe { use std::arch::x86_64::_mm_prefetch; use std::arch::x86_64::_MM_HINT_T0; - - // 预取头指针 _mm_prefetch(self.producer_head.as_ptr() as *const i8, _MM_HINT_T0); - - // 预取尾指针 _mm_prefetch(self.consumer_tail.as_ptr() as *const i8, _MM_HINT_T0); - - // 预取缓冲区开始位置 _mm_prefetch(self.buffer.as_ptr() as *const i8, _MM_HINT_T0); } } } -/// 🚀 CPU分支预测优化工具 +/// Branch hint helpers (likely/unlikely) and prefetch. 分支提示与预取。 pub struct BranchOptimizer; impl BranchOptimizer { - /// likely宏 - 告诉编译器条件大概率为真 + /// Hint: condition is usually true. 提示编译器条件大概率为真。 #[inline(always)] pub fn likely(condition: bool) -> bool { #[cold] @@ -474,7 +435,7 @@ impl BranchOptimizer { condition } - /// unlikely宏 - 告诉编译器条件大概率为假 + /// Hint: condition is usually false. 提示编译器条件大概率为假。 #[inline(always)] pub fn unlikely(condition: bool) -> bool { #[cold] @@ -486,7 +447,7 @@ impl BranchOptimizer { condition } - /// 预取指令 - 提前加载数据到缓存 + /// Prefetch: load cache line at ptr into L1. Caller must ensure ptr is valid, read-only, no concurrent write. 预取:将 ptr 所在缓存行加载到 L1;调用方需保证有效、只读、无并发写。 #[inline(always)] pub unsafe fn prefetch_read_data(ptr: *const T) { #[cfg(target_arch = "x86_64")] @@ -497,7 +458,7 @@ impl BranchOptimizer { } } - /// 预取指令 - 提前加载数据到缓存(写优化) + /// Prefetch for write (T1 hint). 写预取(T1 提示)。 #[inline(always)] pub unsafe fn prefetch_write_data(ptr: *const T) { #[cfg(target_arch = "x86_64")] @@ -509,35 +470,35 @@ impl BranchOptimizer { } } -/// 🚀 内存屏障控制 +/// Memory barrier helpers. 内存屏障辅助。 pub struct MemoryBarriers; impl MemoryBarriers { - /// 编译器屏障 - 防止编译器重排序 + /// Compiler barrier only (no CPU reorder). 仅编译器屏障,防止重排序。 #[inline(always)] pub fn compiler_barrier() { std::sync::atomic::compiler_fence(Ordering::SeqCst); } - /// 轻量级内存屏障 - 仅CPU重排序保护 + /// Light barrier (Acquire). 轻量级屏障(Acquire)。 #[inline(always)] pub fn memory_barrier_light() { std::sync::atomic::fence(Ordering::Acquire); } - /// 重量级内存屏障 - 全序一致性 + /// Full sequential consistency barrier. 全序一致性屏障。 #[inline(always)] pub fn memory_barrier_heavy() { std::sync::atomic::fence(Ordering::SeqCst); } - /// 存储屏障 - 确保写入可见性 + /// Store/release barrier. 存储屏障,保证写入可见性。 #[inline(always)] pub fn store_barrier() { std::sync::atomic::fence(Ordering::Release); } - /// 加载屏障 - 确保读取正确性 + /// Load/acquire barrier. 加载屏障,保证读取顺序。 #[inline(always)] pub fn load_barrier() { std::sync::atomic::fence(Ordering::Acquire); diff --git a/src/perf/mod.rs b/src/perf/mod.rs index ee5b3c4..ca44fbd 100644 --- a/src/perf/mod.rs +++ b/src/perf/mod.rs @@ -1,11 +1,5 @@ -//! 🚀 性能优化模块 -//! -//! 提供多层次性能优化: -//! - SIMD 向量化:AVX2 内存操作、批量计算 -//! - 硬件级优化:分支预测、缓存预取 -//! - 零拷贝 I/O:内存映射、DMA传输 -//! - 系统调用绕过:批处理、快速时间 -//! - 编译器优化:内联、向量化 +//! Performance: SIMD, cache prefetch, branch hints, zero-copy I/O, syscall bypass, compiler hints. +//! 性能优化:SIMD、缓存预取、分支提示、零拷贝 I/O、系统调用绕过、编译器提示。 pub mod simd; pub mod hardware_optimizations; diff --git a/src/perf/syscall_bypass.rs b/src/perf/syscall_bypass.rs index c0f61f6..0dda284 100644 --- a/src/perf/syscall_bypass.rs +++ b/src/perf/syscall_bypass.rs @@ -1,13 +1,5 @@ -//! 🚀 系统调用绕过机制 - 最小化系统调用开销 -//! -//! 实现系统调用级别的极致优化,包括: -//! - 系统调用批处理 -//! - vDSO快速系统调用 -//! - io_uring异步I/O优化 -//! - 内存映射系统调用 -//! - 用户空间系统调用实现 -//! - 系统调用拦截与优化 -//! - 直接硬件访问 +//! Syscall bypass: batching, vDSO fast time, io_uring, mmap, userspace impl. +//! 系统调用绕过:批处理、vDSO 快速时间、io_uring、mmap、用户态实现。 use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::Arc; @@ -18,38 +10,25 @@ use std::fs::OpenOptions; use anyhow::Result; use crossbeam_utils::CachePadded; -/// 🚀 系统调用绕过管理器 +/// Syscall bypass manager (batch, fast time, I/O). 系统调用绕过管理器。 pub struct SystemCallBypassManager { - /// 绕过配置 config: SyscallBypassConfig, - /// 批处理器 batch_processor: Arc, - /// 快速时间获取器 fast_time_provider: Arc, - /// I/O优化器 _io_optimizer: Arc, - /// 统计信息 stats: Arc, } -/// 系统调用绕过配置 +/// Syscall bypass configuration. 系统调用绕过配置。 #[derive(Debug, Clone)] pub struct SyscallBypassConfig { - /// 启用系统调用批处理 pub enable_batch_processing: bool, - /// 批处理大小 pub batch_size: usize, - /// 启用快速时间获取 pub enable_fast_time: bool, - /// 启用vDSO优化 pub enable_vdso: bool, - /// 启用io_uring pub enable_io_uring: bool, - /// 启用内存映射优化 pub enable_mmap_optimization: bool, - /// 启用用户空间实现 pub enable_userspace_impl: bool, - /// 系统调用缓存大小 pub syscall_cache_size: usize, } @@ -68,30 +47,19 @@ impl Default for SyscallBypassConfig { } } -/// 系统调用批处理器 pub struct SyscallBatchProcessor { - /// 待处理的系统调用队列 pending_calls: crossbeam_queue::ArrayQueue, - /// 批处理线程池 _executor: tokio::runtime::Handle, - /// 批处理统计 batch_stats: CachePadded, } -/// 系统调用请求 #[derive(Debug, Clone)] pub enum SyscallRequest { - /// 文件写入 Write { fd: i32, data: Vec }, - /// 文件读取 Read { fd: i32, size: usize }, - /// 网络发送 Send { socket: i32, data: Vec }, - /// 网络接收 Recv { socket: i32, size: usize }, - /// 时间获取 GetTime, - /// 内存分配 MemAlloc { size: usize }, /// 内存释放 MemFree { ptr: usize }, @@ -132,7 +100,7 @@ impl FastTimeProvider { vdso_enabled: enable_vdso, }; - log::info!("🚀 Fast time provider initialized with vDSO: {}", enable_vdso); + tracing::info!(target: "sol_trade_sdk","🚀 Fast time provider initialized with vDSO: {}", enable_vdso); Ok(provider) } @@ -233,7 +201,7 @@ impl IOOptimizer { pub fn new(_config: &SyscallBypassConfig) -> Result { let io_uring_available = Self::check_io_uring_support(); - log::info!("🚀 I/O Optimizer initialized - io_uring: {}", io_uring_available); + tracing::info!(target: "sol_trade_sdk","🚀 I/O Optimizer initialized - io_uring: {}", io_uring_available); Ok(Self { io_uring_available, @@ -249,7 +217,7 @@ impl IOOptimizer { // 检查内核版本和io_uring支持 if let Ok(uname) = std::process::Command::new("uname").arg("-r").output() { let kernel_version = String::from_utf8_lossy(&uname.stdout); - log::info!("Kernel version: {}", kernel_version.trim()); + tracing::info!(target: "sol_trade_sdk","Kernel version: {}", kernel_version.trim()); // 简单检查:内核版本 >= 5.1 支持io_uring if let Some(version_str) = kernel_version.split('.').next() { @@ -277,7 +245,7 @@ impl IOOptimizer { /// 使用io_uring进行批量写入 async fn io_uring_batch_write(&self, requests: &[(i32, &[u8])]) -> Result> { // 这里是伪代码 - 实际实现需要io_uring库 - log::trace!("Using io_uring for {} write operations", requests.len()); + tracing::trace!(target: "sol_trade_sdk","Using io_uring for {} write operations", requests.len()); let mut results = Vec::with_capacity(requests.len()); @@ -362,7 +330,7 @@ impl IOOptimizer { self.mmap_regions.push(region); - log::info!("✅ Memory mapped I/O created: {} bytes at {:p}", size, addr); + tracing::info!(target: "sol_trade_sdk","✅ Memory mapped I/O created: {} bytes at {:p}", size, addr); Ok(addr as usize) } } @@ -390,7 +358,7 @@ impl SyscallBatchProcessor { let pending_calls = crossbeam_queue::ArrayQueue::new(batch_size * 10); let executor = tokio::runtime::Handle::current(); - log::info!("🚀 Syscall batch processor created with batch size: {}", batch_size); + tracing::info!(target: "sol_trade_sdk","🚀 Syscall batch processor created with batch size: {}", batch_size); Ok(Self { pending_calls, @@ -464,7 +432,7 @@ impl SyscallBatchProcessor { self.batch_stats.fetch_add(1, Ordering::Relaxed); - log::trace!("Executed batch of {} syscalls", batch_size); + tracing::trace!(target: "sol_trade_sdk","Executed batch of {} syscalls", batch_size); Ok(batch_size) } @@ -473,7 +441,7 @@ impl SyscallBatchProcessor { // 使用writev系统调用进行批量写入 for (fd, data) in requests { // 实际实现会使用writev或io_uring - log::trace!("Batched write to fd {}: {} bytes", fd, data.len()); + tracing::trace!(target: "sol_trade_sdk","Batched write to fd {}: {} bytes", fd, data.len()); } Ok(()) } @@ -482,7 +450,7 @@ impl SyscallBatchProcessor { async fn batch_read_operations(&self, requests: Vec<(i32, usize)>) -> Result<()> { // 使用readv系统调用进行批量读取 for (fd, size) in requests { - log::trace!("Batched read from fd {}: {} bytes", fd, size); + tracing::trace!(target: "sol_trade_sdk","Batched read from fd {}: {} bytes", fd, size); } Ok(()) } @@ -491,7 +459,7 @@ impl SyscallBatchProcessor { async fn batch_network_operations(&self, requests: Vec<(i32, Vec)>) -> Result<()> { // 使用sendmsg/recvmsg进行批量网络操作 for (socket, data) in requests { - log::trace!("Batched network send to socket {}: {} bytes", socket, data.len()); + tracing::trace!(target: "sol_trade_sdk","Batched network send to socket {}: {} bytes", socket, data.len()); } Ok(()) } @@ -515,11 +483,11 @@ impl SystemCallBypassManager { let io_optimizer = Arc::new(IOOptimizer::new(&config)?); let stats = Arc::new(SyscallBypassStats::default()); - log::info!("🚀 System Call Bypass Manager initialized"); - log::info!(" 📦 Batch Processing: {}", config.enable_batch_processing); - log::info!(" ⏰ Fast Time: {}", config.enable_fast_time); - log::info!(" 🚀 vDSO: {}", config.enable_vdso); - log::info!(" 📁 io_uring: {}", config.enable_io_uring); + tracing::info!(target: "sol_trade_sdk","🚀 System Call Bypass Manager initialized"); + tracing::info!(target: "sol_trade_sdk"," 📦 Batch Processing: {}", config.enable_batch_processing); + tracing::info!(target: "sol_trade_sdk"," ⏰ Fast Time: {}", config.enable_fast_time); + tracing::info!(target: "sol_trade_sdk"," 🚀 vDSO: {}", config.enable_vdso); + tracing::info!(target: "sol_trade_sdk"," 📁 io_uring: {}", config.enable_io_uring); Ok(Self { config, @@ -626,7 +594,7 @@ impl SystemCallBypassManager { } }); - log::info!("✅ Batch processing worker started"); + tracing::info!(target: "sol_trade_sdk","✅ Batch processing worker started"); Ok(()) } @@ -669,16 +637,16 @@ pub struct SyscallBypassStatsSnapshot { impl SyscallBypassStatsSnapshot { /// 打印统计信息 pub fn print_stats(&self) { - log::info!("📊 System Call Bypass Stats:"); - log::info!(" 🚫 Syscalls Bypassed: {}", self.syscalls_bypassed); - log::info!(" 📦 Syscalls Batched: {}", self.syscalls_batched); - log::info!(" ⏰ Time Calls Cached: {}", self.time_calls_cached); - log::info!(" 📁 I/O Operations Optimized: {}", self.io_operations_optimized); - log::info!(" 💾 Memory Operations Avoided: {}", self.memory_operations_avoided); + tracing::info!(target: "sol_trade_sdk","📊 System Call Bypass Stats:"); + tracing::info!(target: "sol_trade_sdk"," 🚫 Syscalls Bypassed: {}", self.syscalls_bypassed); + tracing::info!(target: "sol_trade_sdk"," 📦 Syscalls Batched: {}", self.syscalls_batched); + tracing::info!(target: "sol_trade_sdk"," ⏰ Time Calls Cached: {}", self.time_calls_cached); + tracing::info!(target: "sol_trade_sdk"," 📁 I/O Operations Optimized: {}", self.io_operations_optimized); + tracing::info!(target: "sol_trade_sdk"," 💾 Memory Operations Avoided: {}", self.memory_operations_avoided); let total_optimizations = self.syscalls_bypassed + self.time_calls_cached + self.io_operations_optimized + self.memory_operations_avoided; - log::info!(" 🏆 Total Optimizations: {}", total_optimizations); + tracing::info!(target: "sol_trade_sdk"," 🏆 Total Optimizations: {}", total_optimizations); } } diff --git a/src/perf/zero_copy_io.rs b/src/perf/zero_copy_io.rs index bbeca79..4030f80 100644 --- a/src/perf/zero_copy_io.rs +++ b/src/perf/zero_copy_io.rs @@ -73,7 +73,7 @@ impl SharedMemoryPool { free_blocks.push(AtomicU64::new(bits)); } - log::info!("🚀 Created shared memory pool {} with {} blocks of {} bytes each", + tracing::info!(target: "sol_trade_sdk","🚀 Created shared memory pool {} with {} blocks of {} bytes each", pool_id, total_blocks, aligned_block_size); Ok(Self { @@ -155,7 +155,7 @@ impl SharedMemoryPool { #[inline(always)] pub fn deallocate_block(&self, block: ZeroCopyBlock) { if block.pool_id != self.pool_id { - log::error!("Attempting to deallocate block from wrong pool"); + tracing::error!(target: "sol_trade_sdk", "Attempting to deallocate block from wrong pool"); return; } @@ -267,7 +267,7 @@ impl MemoryMappedBuffer { .map_anon() .context("Failed to create memory mapped buffer")?; - log::info!("🚀 Created memory mapped buffer {} with size {} bytes", buffer_id, size); + tracing::info!(target: "sol_trade_sdk","🚀 Created memory mapped buffer {} with size {} bytes", buffer_id, size); Ok(Self { mmap, @@ -425,7 +425,7 @@ impl DirectMemoryAccessManager { dma_channels.push(Arc::new(DMAChannel::new(i)?)); } - log::info!("🚀 Created DMA manager with {} channels", num_channels); + tracing::info!(target: "sol_trade_sdk","🚀 Created DMA manager with {} channels", num_channels); Ok(Self { dma_channels, @@ -549,12 +549,12 @@ impl ZeroCopyStats { let bytes = self.bytes_transferred.load(Ordering::Relaxed); let mmap_usage = self.mmap_buffer_usage.load(Ordering::Relaxed); - log::info!("🚀 Zero-Copy Stats:"); - log::info!(" 📦 Blocks: Allocated={}, Freed={}, Active={}", + tracing::info!(target: "sol_trade_sdk","🚀 Zero-Copy Stats:"); + tracing::info!(target: "sol_trade_sdk"," 📦 Blocks: Allocated={}, Freed={}, Active={}", allocated, freed, allocated.saturating_sub(freed)); - log::info!(" 📊 Bytes Transferred: {} ({:.2} MB)", + tracing::info!(target: "sol_trade_sdk"," 📊 Bytes Transferred: {} ({:.2} MB)", bytes, bytes as f64 / 1024.0 / 1024.0); - log::info!(" 💾 Memory Mapped Usage: {} ({:.2} MB)", + tracing::info!(target: "sol_trade_sdk"," 💾 Memory Mapped Usage: {} ({:.2} MB)", mmap_usage, mmap_usage as f64 / 1024.0 / 1024.0); } } @@ -581,10 +581,10 @@ impl ZeroCopyMemoryManager { let dma_manager = Arc::new(DirectMemoryAccessManager::new(16)?); // 16 DMA channels let stats = Arc::new(ZeroCopyStats::new()); - log::info!("🚀 Zero-Copy Memory Manager initialized"); - log::info!(" 📦 Memory Pools: {}", shared_pools.len()); - log::info!(" 💾 Mapped Buffers: {}", mmap_buffers.len()); - log::info!(" 🔄 DMA Channels: 16"); + tracing::info!(target: "sol_trade_sdk","🚀 Zero-Copy Memory Manager initialized"); + tracing::info!(target: "sol_trade_sdk"," 📦 Memory Pools: {}", shared_pools.len()); + tracing::info!(target: "sol_trade_sdk"," 💾 Mapped Buffers: {}", mmap_buffers.len()); + tracing::info!(target: "sol_trade_sdk"," 🔄 DMA Channels: 16"); Ok(Self { shared_pools, diff --git a/src/swqos/blockrazor.rs b/src/swqos/blockrazor.rs index f6ed801..3745929 100644 --- a/src/swqos/blockrazor.rs +++ b/src/swqos/blockrazor.rs @@ -92,7 +92,9 @@ impl BlockRazorClient { let handle = tokio::spawn(async move { // Immediate first ping to warm connection and reduce first-submit cold start latency if let Err(e) = Self::send_ping_request(&http_client, &endpoint, &auth_token).await { - eprintln!("BlockRazor ping request failed: {}", e); + if crate::common::sdk_log::sdk_log_enabled() { + eprintln!("BlockRazor ping request failed: {}", e); + } } let mut interval = tokio::time::interval(Duration::from_secs(30)); // 30s keepalive to avoid server ~5min idle close loop { @@ -101,7 +103,9 @@ impl BlockRazorClient { break; } if let Err(e) = Self::send_ping_request(&http_client, &endpoint, &auth_token).await { - eprintln!("BlockRazor ping request failed: {}", e); + if crate::common::sdk_log::sdk_log_enabled() { + eprintln!("BlockRazor ping request failed: {}", e); + } } } }); @@ -173,12 +177,14 @@ impl BlockRazorClient { // Parse JSON response if let Ok(response_json) = serde_json::from_str::(&response_text) { - if response_json.get("result").is_some() || response_json.get("signature").is_some() { - println!(" [blockrazor] {} submitted: {:?}", trade_type, start_time.elapsed()); - } else if let Some(_error) = response_json.get("error") { - eprintln!(" [blockrazor] {} submission failed: {:?}", trade_type, _error); + if crate::common::sdk_log::sdk_log_enabled() { + if response_json.get("result").is_some() || response_json.get("signature").is_some() { + println!(" [blockrazor] {} submitted: {:?}", trade_type, start_time.elapsed()); + } else if let Some(_error) = response_json.get("error") { + eprintln!(" [blockrazor] {} submission failed: {:?}", trade_type, _error); + } } - } else { + } else if crate::common::sdk_log::sdk_log_enabled() { eprintln!(" [blockrazor] {} submission failed: {:?}", trade_type, response_text); } @@ -186,12 +192,14 @@ impl BlockRazorClient { match poll_transaction_confirmation(&self.rpc_client, signature, wait_confirmation).await { Ok(_) => (), Err(e) => { - println!(" signature: {:?}", signature); - println!(" [blockrazor] {} confirmation failed: {:?}", trade_type, start_time.elapsed()); + if crate::common::sdk_log::sdk_log_enabled() { + println!(" signature: {:?}", signature); + println!(" [blockrazor] {} confirmation failed: {:?}", trade_type, start_time.elapsed()); + } return Err(e); }, } - if wait_confirmation { + if wait_confirmation && crate::common::sdk_log::sdk_log_enabled() { println!(" signature: {:?}", signature); println!(" [blockrazor] {} confirmed: {:?}", trade_type, start_time.elapsed()); } diff --git a/src/swqos/bloxroute.rs b/src/swqos/bloxroute.rs index e092bac..e802796 100755 --- a/src/swqos/bloxroute.rs +++ b/src/swqos/bloxroute.rs @@ -1,3 +1,4 @@ +use crate::swqos::common::default_http_client_builder; use crate::swqos::common::poll_transaction_confirmation; use crate::swqos::common::serialize_transaction_and_encode; use crate::swqos::serialization; @@ -47,17 +48,9 @@ impl SwqosClientTrait for BloxrouteClient { impl BloxrouteClient { pub fn new(rpc_url: String, endpoint: String, auth_token: String) -> Self { let rpc_client = SolanaRpcClient::new(rpc_url); - let http_client = Client::builder() - // Optimized connection pool settings for high performance + let http_client = default_http_client_builder() .pool_idle_timeout(Duration::from_secs(120)) - .pool_max_idle_per_host(256) // Increased from 64 to 256 - .tcp_keepalive(Some(Duration::from_secs(60))) // Reduced from 1200 to 60 - .tcp_nodelay(true) // Disable Nagle's algorithm for lower latency - .http2_keep_alive_interval(Duration::from_secs(10)) - .http2_keep_alive_timeout(Duration::from_secs(5)) - .http2_adaptive_window(true) // Enable adaptive flow control - .timeout(Duration::from_millis(3000)) // Reduced from 10s to 3s - .connect_timeout(Duration::from_millis(2000)) // Reduced from 5s to 2s + .pool_max_idle_per_host(256) .build() .unwrap(); Self { rpc_client: Arc::new(rpc_client), endpoint, auth_token, http_client } @@ -85,12 +78,14 @@ impl BloxrouteClient { // Parse with from_str to avoid extra wait from .json().await if let Ok(response_json) = serde_json::from_str::(&response_text) { - if response_json.get("result").is_some() { - println!(" [bloxroute] {} submitted: {:?}", trade_type, start_time.elapsed()); - } else if let Some(_error) = response_json.get("error") { - eprintln!(" [bloxroute] {} submission failed: {:?}", trade_type, _error); + if crate::common::sdk_log::sdk_log_enabled() { + if response_json.get("result").is_some() { + println!(" [bloxroute] {} submitted: {:?}", trade_type, start_time.elapsed()); + } else if let Some(_error) = response_json.get("error") { + eprintln!(" [bloxroute] {} submission failed: {:?}", trade_type, _error); + } } - } else { + } else if crate::common::sdk_log::sdk_log_enabled() { eprintln!(" [bloxroute] {} submission failed: {:?}", trade_type, response_text); } @@ -98,12 +93,14 @@ impl BloxrouteClient { match poll_transaction_confirmation(&self.rpc_client, signature, wait_confirmation).await { Ok(_) => (), Err(e) => { - println!(" signature: {:?}", signature); - println!(" [bloxroute] {} confirmation failed: {:?}", trade_type, start_time.elapsed()); + if crate::common::sdk_log::sdk_log_enabled() { + println!(" signature: {:?}", signature); + println!(" [bloxroute] {} confirmation failed: {:?}", trade_type, start_time.elapsed()); + } return Err(e); }, } - if wait_confirmation { + if wait_confirmation && crate::common::sdk_log::sdk_log_enabled() { println!(" signature: {:?}", signature); println!(" [bloxroute] {} confirmed: {:?}", trade_type, start_time.elapsed()); } @@ -135,11 +132,13 @@ impl BloxrouteClient { .text() .await?; - if let Ok(response_json) = serde_json::from_str::(&response_text) { - if response_json.get("result").is_some() { - println!(" bloxroute {} submitted: {:?}", trade_type, start_time.elapsed()); - } else if let Some(_error) = response_json.get("error") { - eprintln!(" bloxroute {} submission failed: {:?}", trade_type, _error); + if crate::common::sdk_log::sdk_log_enabled() { + if let Ok(response_json) = serde_json::from_str::(&response_text) { + if response_json.get("result").is_some() { + println!(" bloxroute {} submitted: {:?}", trade_type, start_time.elapsed()); + } else if let Some(_error) = response_json.get("error") { + eprintln!(" bloxroute {} submission failed: {:?}", trade_type, _error); + } } } diff --git a/src/swqos/common.rs b/src/swqos/common.rs index 1b3f7e7..9249cd5 100755 --- a/src/swqos/common.rs +++ b/src/swqos/common.rs @@ -17,6 +17,36 @@ use std::str::FromStr; use std::time::{Duration, Instant}; use tokio::time::sleep; +/// Default pool idle timeout for SWQOS HTTP client (seconds). 连接池空闲超时(秒)。 +const HTTP_POOL_IDLE_TIMEOUT_SECS: u64 = 300; +/// Max idle connections per host. 每主机最大空闲连接数。 +const HTTP_POOL_MAX_IDLE_PER_HOST: usize = 4; +/// TCP keepalive interval (seconds). TCP 保活间隔(秒)。 +const HTTP_TCP_KEEPALIVE_SECS: u64 = 60; +/// HTTP/2 keepalive interval (seconds). HTTP/2 保活间隔(秒)。 +const HTTP2_KEEPALIVE_INTERVAL_SECS: u64 = 10; +/// HTTP/2 keepalive timeout (seconds). HTTP/2 保活超时(秒)。 +const HTTP2_KEEPALIVE_TIMEOUT_SECS: u64 = 5; +/// Request timeout (milliseconds). 请求超时(毫秒)。 +const HTTP_TIMEOUT_MS: u64 = 3000; +/// Connect timeout (milliseconds). 连接超时(毫秒)。 +const HTTP_CONNECT_TIMEOUT_MS: u64 = 2000; + +/// Shared HTTP client builder for SWQOS clients; call `.build().unwrap()` or override pool first. SWQOS 共用 HTTP 客户端构建器。 +pub fn default_http_client_builder() -> reqwest::ClientBuilder { + Client::builder() + .pool_idle_timeout(Duration::from_secs(HTTP_POOL_IDLE_TIMEOUT_SECS)) + .pool_max_idle_per_host(HTTP_POOL_MAX_IDLE_PER_HOST) + .tcp_keepalive(Some(Duration::from_secs(HTTP_TCP_KEEPALIVE_SECS))) + .tcp_nodelay(true) + .http2_keep_alive_interval(Duration::from_secs(HTTP2_KEEPALIVE_INTERVAL_SECS)) + .http2_keep_alive_timeout(Duration::from_secs(HTTP2_KEEPALIVE_TIMEOUT_SECS)) + .http2_adaptive_window(true) + .timeout(Duration::from_millis(HTTP_TIMEOUT_MS)) + .connect_timeout(Duration::from_millis(HTTP_CONNECT_TIMEOUT_MS)) +} + +/// Trade/on-chain error with code and optional instruction index. 交易/链上错误,含错误码与可选指令下标。 #[derive(Debug, Clone)] pub struct TradeError { pub code: u32, diff --git a/src/swqos/node1.rs b/src/swqos/node1.rs index 15e1967..fc74637 100644 --- a/src/swqos/node1.rs +++ b/src/swqos/node1.rs @@ -1,4 +1,4 @@ -use crate::swqos::common::{poll_transaction_confirmation, serialize_transaction_and_encode}; +use crate::swqos::common::{default_http_client_builder, poll_transaction_confirmation, serialize_transaction_and_encode}; use rand::seq::IndexedRandom; use reqwest::Client; use serde_json::json; @@ -50,19 +50,7 @@ impl SwqosClientTrait for Node1Client { impl Node1Client { pub fn new(rpc_url: String, endpoint: String, auth_token: String) -> Self { let rpc_client = SolanaRpcClient::new(rpc_url); - let http_client = Client::builder() - // Optimized connection pool settings for high performance - .pool_idle_timeout(Duration::from_secs(300)) - .pool_max_idle_per_host(4) - .tcp_keepalive(Some(Duration::from_secs(60))) // Reduced from 1200 to 60 - .tcp_nodelay(true) // Disable Nagle's algorithm for lower latency - .http2_keep_alive_interval(Duration::from_secs(10)) - .http2_keep_alive_timeout(Duration::from_secs(5)) - .http2_adaptive_window(true) // Enable adaptive flow control - .timeout(Duration::from_millis(3000)) // Reduced from 10s to 3s - .connect_timeout(Duration::from_millis(2000)) // Reduced from 5s to 2s - .build() - .unwrap(); + let http_client = default_http_client_builder().build().unwrap(); let client = Self { rpc_client: Arc::new(rpc_client), @@ -92,7 +80,9 @@ impl Node1Client { let handle = tokio::spawn(async move { // Immediate first ping to warm connection and reduce first-submit cold start latency if let Err(e) = Self::send_ping_request(&http_client, &endpoint, &auth_token).await { - eprintln!("Node1 ping request failed: {}", e); + if crate::common::sdk_log::sdk_log_enabled() { + eprintln!("Node1 ping request failed: {}", e); + } } let mut interval = tokio::time::interval(Duration::from_secs(30)); loop { @@ -101,7 +91,9 @@ impl Node1Client { break; } if let Err(e) = Self::send_ping_request(&http_client, &endpoint, &auth_token).await { - eprintln!("Node1 ping request failed: {}", e); + if crate::common::sdk_log::sdk_log_enabled() { + eprintln!("Node1 ping request failed: {}", e); + } } } }); @@ -132,7 +124,7 @@ impl Node1Client { .await?; let status = response.status(); let _ = response.bytes().await; - if !status.is_success() { + if !status.is_success() && crate::common::sdk_log::sdk_log_enabled() { eprintln!("Node1 ping request returned non-success status: {}", status); } Ok(()) @@ -164,12 +156,14 @@ impl Node1Client { // Parse JSON response if let Ok(response_json) = serde_json::from_str::(&response_text) { - if response_json.get("result").is_some() { - println!(" [node1] {} submitted: {:?}", trade_type, start_time.elapsed()); - } else if let Some(_error) = response_json.get("error") { - eprintln!(" [node1] {} submission failed: {:?}", trade_type, _error); + if crate::common::sdk_log::sdk_log_enabled() { + if response_json.get("result").is_some() { + println!(" [node1] {} submitted: {:?}", trade_type, start_time.elapsed()); + } else if let Some(_error) = response_json.get("error") { + eprintln!(" [node1] {} submission failed: {:?}", trade_type, _error); + } } - } else { + } else if crate::common::sdk_log::sdk_log_enabled() { eprintln!(" [node1] {} submission failed: {:?}", trade_type, response_text); } @@ -177,12 +171,14 @@ impl Node1Client { match poll_transaction_confirmation(&self.rpc_client, signature, wait_confirmation).await { Ok(_) => (), Err(e) => { - println!(" signature: {:?}", signature); - println!(" [node1] {} confirmation failed: {:?}", trade_type, start_time.elapsed()); + if crate::common::sdk_log::sdk_log_enabled() { + println!(" signature: {:?}", signature); + println!(" [node1] {} confirmation failed: {:?}", trade_type, start_time.elapsed()); + } return Err(e); }, } - if wait_confirmation { + if wait_confirmation && crate::common::sdk_log::sdk_log_enabled() { println!(" signature: {:?}", signature); println!(" [node1] {} confirmed: {:?}", trade_type, start_time.elapsed()); } diff --git a/src/swqos/serialization.rs b/src/swqos/serialization.rs index 4d80a20..abfd416 100644 --- a/src/swqos/serialization.rs +++ b/src/swqos/serialization.rs @@ -105,7 +105,41 @@ impl Base64Encoder { } } +/// Guard that returns the serialization buffer to the pool on drop. +pub struct PooledTxBufGuard(pub Vec); + +impl std::ops::Deref for PooledTxBufGuard { + type Target = [u8]; + fn deref(&self) -> &[u8] { + &self.0 + } +} + +impl Drop for PooledTxBufGuard { + fn drop(&mut self) { + if !self.0.is_empty() { + SERIALIZER.return_buffer(std::mem::take(&mut self.0)); + } + } +} + +/// Serialize transaction to bincode bytes using buffer pool. The returned guard returns the buffer +/// to the pool when dropped; use `&*guard` or `guard.as_ref()` for `&[u8]`. +pub fn serialize_transaction_bincode_sync( + transaction: &impl SerializableTransaction, +) -> Result<(PooledTxBufGuard, Signature)> { + let signature = transaction.get_signature(); + let serialized_tx = SERIALIZER.serialize_zero_alloc(transaction, "transaction")?; + Ok((PooledTxBufGuard(serialized_tx), *signature)) +} + +/// Return a buffer to the pool (for manual use when not using `PooledTxBufGuard`). +pub fn return_serialization_buffer(buffer: Vec) { + SERIALIZER.return_buffer(buffer); +} + /// Sync serialize + encode using buffer pool; use in hot path to reduce allocs. +/// Base64 path uses SIMD-accelerated encoding. pub fn serialize_transaction_sync( transaction: &impl SerializableTransaction, encoding: UiTransactionEncoding, @@ -114,7 +148,7 @@ pub fn serialize_transaction_sync( let serialized_tx = SERIALIZER.serialize_zero_alloc(transaction, "transaction")?; let serialized = match encoding { UiTransactionEncoding::Base58 => bs58::encode(&serialized_tx).into_string(), - UiTransactionEncoding::Base64 => STANDARD.encode(&serialized_tx), + UiTransactionEncoding::Base64 => SIMDSerializer::encode_base64_simd(&serialized_tx), _ => return Err(anyhow::anyhow!("Unsupported encoding")), }; SERIALIZER.return_buffer(serialized_tx); @@ -133,10 +167,7 @@ pub async fn serialize_transaction( let serialized = match encoding { UiTransactionEncoding::Base58 => bs58::encode(&serialized_tx).into_string(), - UiTransactionEncoding::Base64 => { - // Use SIMD-optimized Base64 encoding - STANDARD.encode(&serialized_tx) - } + UiTransactionEncoding::Base64 => SIMDSerializer::encode_base64_simd(&serialized_tx), _ => return Err(anyhow::anyhow!("Unsupported encoding")), }; @@ -156,7 +187,7 @@ pub fn serialize_transactions_batch_sync( let serialized_tx = SERIALIZER.serialize_zero_alloc(tx, "transaction")?; let encoded = match encoding { UiTransactionEncoding::Base58 => bs58::encode(&serialized_tx).into_string(), - UiTransactionEncoding::Base64 => STANDARD.encode(&serialized_tx), + UiTransactionEncoding::Base64 => SIMDSerializer::encode_base64_simd(&serialized_tx), _ => return Err(anyhow::anyhow!("Unsupported encoding")), }; SERIALIZER.return_buffer(serialized_tx); @@ -177,7 +208,7 @@ pub async fn serialize_transactions_batch( let encoded = match encoding { UiTransactionEncoding::Base58 => bs58::encode(&serialized_tx).into_string(), - UiTransactionEncoding::Base64 => STANDARD.encode(&serialized_tx), + UiTransactionEncoding::Base64 => SIMDSerializer::encode_base64_simd(&serialized_tx), _ => return Err(anyhow::anyhow!("Unsupported encoding")), }; diff --git a/src/swqos/speedlanding.rs b/src/swqos/speedlanding.rs index 61118a4..a791dc2 100644 --- a/src/swqos/speedlanding.rs +++ b/src/swqos/speedlanding.rs @@ -6,7 +6,6 @@ use quinn::{ TransportConfig, }; use rand::seq::IndexedRandom as _; -use solana_rpc_client::rpc_client::SerializableTransaction; use solana_sdk::{signature::Keypair, transaction::VersionedTransaction}; use solana_tls_utils::{new_dummy_x509_certificate, SkipServerVerification}; use std::time::Instant; @@ -19,6 +18,7 @@ use tokio::sync::Mutex; use crate::common::SolanaRpcClient; use crate::swqos::common::poll_transaction_confirmation; +use crate::swqos::serialization::serialize_transaction_bincode_sync; use crate::swqos::SwqosClientTrait; use crate::{ constants::swqos::SPEEDLANDING_TIP_ACCOUNTS, @@ -105,27 +105,32 @@ impl SwqosClientTrait for SpeedlandingClient { wait_confirmation: bool, ) -> Result<()> { let start_time = Instant::now(); - let signature = transaction.get_signature(); - let serialized_tx = bincode::serialize(transaction)?; + let (buf_guard, signature) = serialize_transaction_bincode_sync(transaction)?; let connection = self.connection.load_full(); - if Self::try_send_bytes(&connection, &serialized_tx).await.is_err() { - eprintln!(" [speedlanding] {} submission failed, reconnecting", trade_type); + if Self::try_send_bytes(&connection, &*buf_guard).await.is_err() { + if crate::common::sdk_log::sdk_log_enabled() { + eprintln!(" [speedlanding] {} submission failed, reconnecting", trade_type); + } self.reconnect().await?; let connection = self.connection.load_full(); - if let Err(e) = Self::try_send_bytes(&connection, &serialized_tx).await { - eprintln!(" [speedlanding] {} submission failed: {:?}", trade_type, e); + if let Err(e) = Self::try_send_bytes(&connection, &*buf_guard).await { + if crate::common::sdk_log::sdk_log_enabled() { + eprintln!(" [speedlanding] {} submission failed: {:?}", trade_type, e); + } return Err(e.into()); } } - match poll_transaction_confirmation(&self.rpc_client, *signature, wait_confirmation).await { + match poll_transaction_confirmation(&self.rpc_client, signature, wait_confirmation).await { Ok(_) => (), Err(e) => { - println!(" signature: {:?}", signature); - println!(" [speedlanding] {} confirmation failed: {:?}", trade_type, start_time.elapsed()); + if crate::common::sdk_log::sdk_log_enabled() { + println!(" signature: {:?}", signature); + println!(" [speedlanding] {} confirmation failed: {:?}", trade_type, start_time.elapsed()); + } return Err(e); } } - if wait_confirmation { + if wait_confirmation && crate::common::sdk_log::sdk_log_enabled() { println!(" signature: {:?}", signature); println!(" [speedlanding] {} confirmed: {:?}", trade_type, start_time.elapsed()); } diff --git a/src/swqos/stellium.rs b/src/swqos/stellium.rs index b7b79c2..592997d 100644 --- a/src/swqos/stellium.rs +++ b/src/swqos/stellium.rs @@ -1,4 +1,4 @@ -use crate::swqos::common::{poll_transaction_confirmation, serialize_transaction_and_encode}; +use crate::swqos::common::{default_http_client_builder, poll_transaction_confirmation, serialize_transaction_and_encode}; use rand::seq::IndexedRandom; use reqwest::Client; use serde_json::json; @@ -48,19 +48,7 @@ impl SwqosClientTrait for StelliumClient { impl StelliumClient { pub fn new(rpc_url: String, endpoint: String, auth_token: String) -> Self { let rpc_client = SolanaRpcClient::new(rpc_url); - let http_client = Client::builder() - // Optimized connection pool settings for high performance - .pool_idle_timeout(Duration::from_secs(300)) - .pool_max_idle_per_host(4) - .tcp_keepalive(Some(Duration::from_secs(60))) // Reduced from 1200 to 60 - .tcp_nodelay(true) // Disable Nagle's algorithm for lower latency - .http2_keep_alive_interval(Duration::from_secs(10)) - .http2_keep_alive_timeout(Duration::from_secs(5)) - .http2_adaptive_window(true) // Enable adaptive flow control - .timeout(Duration::from_millis(3000)) // Reduced from 10s to 3s - .connect_timeout(Duration::from_millis(2000)) // Reduced from 5s to 2s - .build() - .unwrap(); + let http_client = default_http_client_builder().build().unwrap(); let keep_alive_running = Arc::new(AtomicBool::new(true)); @@ -94,7 +82,7 @@ impl StelliumClient { if let Ok(resp) = http_client.get(&url).timeout(Duration::from_millis(1500)).send().await { let status = resp.status(); let _ = resp.bytes().await; - if !status.is_success() { + if !status.is_success() && crate::common::sdk_log::sdk_log_enabled() { eprintln!(" [Stellium] Ping failed with status: {}", status); } } @@ -109,12 +97,14 @@ impl StelliumClient { Ok(response) => { let status = response.status(); let _ = response.bytes().await; - if !status.is_success() { + if !status.is_success() && crate::common::sdk_log::sdk_log_enabled() { eprintln!(" [Stellium] Ping failed with status: {}", status); } } Err(e) => { - eprintln!(" [Stellium] Ping request error: {:?}", e); + if crate::common::sdk_log::sdk_log_enabled() { + eprintln!(" [Stellium] Ping request error: {:?}", e); + } } } } @@ -152,12 +142,14 @@ impl StelliumClient { // Parse response if let Ok(response_json) = serde_json::from_str::(&response_text) { - if response_json.get("result").is_some() { - println!(" [Stellium] {} submitted: {:?}", trade_type, start_time.elapsed()); - } else if let Some(_error) = response_json.get("error") { - eprintln!(" [Stellium] {} submission failed: {:?}", trade_type, _error); + if crate::common::sdk_log::sdk_log_enabled() { + if response_json.get("result").is_some() { + println!(" [Stellium] {} submitted: {:?}", trade_type, start_time.elapsed()); + } else if let Some(_error) = response_json.get("error") { + eprintln!(" [Stellium] {} submission failed: {:?}", trade_type, _error); + } } - } else { + } else if crate::common::sdk_log::sdk_log_enabled() { eprintln!(" [Stellium] {} submission failed: {:?}", trade_type, response_text); } @@ -165,12 +157,14 @@ impl StelliumClient { match poll_transaction_confirmation(&self.rpc_client, signature, wait_confirmation).await { Ok(_) => (), Err(e) => { - println!(" signature: {:?}", signature); - println!(" [Stellium] {} confirmation failed: {:?}", trade_type, start_time.elapsed()); + if crate::common::sdk_log::sdk_log_enabled() { + println!(" signature: {:?}", signature); + println!(" [Stellium] {} confirmation failed: {:?}", trade_type, start_time.elapsed()); + } return Err(e); }, } - if wait_confirmation { + if wait_confirmation && crate::common::sdk_log::sdk_log_enabled() { println!(" signature: {:?}", signature); println!(" [Stellium] {} confirmed: {:?}", trade_type, start_time.elapsed()); } diff --git a/src/trading/common/transaction_builder.rs b/src/trading/common/transaction_builder.rs index e03f408..56cf8e3 100755 --- a/src/trading/common/transaction_builder.rs +++ b/src/trading/common/transaction_builder.rs @@ -14,13 +14,14 @@ use crate::{ trading::{MiddlewareManager, core::transaction_pool::{acquire_builder, release_builder}}, }; -/// Build standard RPC transaction +/// Build standard RPC transaction. +/// Takes `business_instructions` by reference to avoid per-task Vec clone in execute_parallel. pub async fn build_transaction( payer: Arc, _rpc: Option>, unit_limit: u32, unit_price: u64, - business_instructions: Vec, + business_instructions: &[Instruction], address_lookup_table_account: Option, recent_blockhash: Option, middleware_manager: Option>, @@ -57,8 +58,8 @@ pub async fn build_transaction( unit_limit, )); - // Add business instructions - instructions.extend(business_instructions); + // Add business instructions (clone only here; avoids per-task Vec clone in execute_parallel) + instructions.extend_from_slice(business_instructions); // Get blockhash for transaction let blockhash = get_transaction_blockhash(recent_blockhash, durable_nonce.clone()); diff --git a/src/trading/core/async_executor.rs b/src/trading/core/async_executor.rs index 098a126..75d3a16 100644 --- a/src/trading/core/async_executor.rs +++ b/src/trading/core/async_executor.rs @@ -1,4 +1,4 @@ -//! 并行执行器 +//! Parallel executor for multi-SWQOS submit. use anyhow::{anyhow, Result}; use crossbeam_queue::ArrayQueue; @@ -39,8 +39,8 @@ struct TaskResult { signature: Signature, error: Option, #[allow(dead_code)] - swqos_type: SwqosType, // 🔧 增加:记录SWQOS类型 - landed_on_chain: bool, // 🔧 Whether tx landed on-chain (even if failed) + swqos_type: SwqosType, + landed_on_chain: bool, } /// Check if an error indicates the transaction landed on-chain (vs network/timeout error) @@ -87,14 +87,14 @@ impl ResultCollector { } fn submit(&self, result: TaskResult) { - // 🚀 优化:ArrayQueue 内部已保证同步,无需额外 fence + // ArrayQueue is already synchronized; no extra fence needed let is_success = result.success; let is_landed_failed = result.landed_on_chain && !result.success; let _ = self.results.push(result); if is_success { - self.success_flag.store(true, Ordering::Release); // Release 确保 push 可见 + self.success_flag.store(true, Ordering::Release); } else if is_landed_failed { // 🔧 Tx landed but failed (e.g., ExceededSlippage) - nonce is consumed, no point waiting self.landed_failed_flag.store(true, Ordering::Release); @@ -105,12 +105,11 @@ impl ResultCollector { async fn wait_for_success(&self) -> Option<(bool, Vec, Option)> { let start = Instant::now(); - let timeout = std::time::Duration::from_secs(30); + let timeout = std::time::Duration::from_secs(5); + let poll_interval = std::time::Duration::from_millis(1000); loop { - // 🚀 Acquire 确保看到 push 的内容 if self.success_flag.load(Ordering::Acquire) { - // 🔧 修复:收集所有签名 let mut signatures = Vec::new(); let mut has_success = false; while let Some(result) = self.results.pop() { @@ -124,7 +123,7 @@ impl ResultCollector { } } - // 🔧 Early exit: if a tx landed but failed (e.g., ExceededSlippage), + // Early exit: if a tx landed but failed (e.g., ExceededSlippage), // nonce is consumed and other channels can't succeed - return immediately if self.landed_failed_flag.load(Ordering::Acquire) { let mut signatures = Vec::new(); @@ -142,8 +141,7 @@ impl ResultCollector { } let completed = self.completed_count.load(Ordering::Acquire); - if completed >= self.total_tasks { - // 🔧 修复:收集所有签名 + if completed >= self.total_tasks { let mut signatures = Vec::new(); let mut last_error = None; let mut any_success = false; @@ -165,12 +163,11 @@ impl ResultCollector { if start.elapsed() > timeout { return None; } - tokio::task::yield_now().await; + tokio::time::sleep(poll_interval).await; } } fn get_first(&self) -> Option<(bool, Vec, Option)> { - // 🔧 修复:收集已提交的所有签名 let mut signatures = Vec::new(); let mut has_success = false; let mut last_error = None; @@ -191,11 +188,25 @@ impl ResultCollector { None } } + + /// 等待全部任务完成(不等待链上确认),然后收集并返回所有签名。用于「多路提交」时返回多笔签名。 + async fn wait_for_all_submitted(&self, timeout_secs: u64) -> Option<(bool, Vec, Option)> { + let start = Instant::now(); + let timeout = std::time::Duration::from_secs(timeout_secs); + let poll_interval = std::time::Duration::from_millis(50); + while self.completed_count.load(Ordering::Acquire) < self.total_tasks { + if start.elapsed() > timeout { + break; + } + tokio::time::sleep(poll_interval).await; + } + self.get_first() + } } -/// 🔧 修复:返回Vec支持多SWQOS并发交易 +/// Execute trade on multiple SWQOS clients in parallel; returns success flag, all signatures, and last error. pub async fn execute_parallel( - swqos_clients: Vec>, + swqos_clients: &[Arc], payer: Arc, rpc: Option>, instructions: Vec, @@ -208,6 +219,7 @@ pub async fn execute_parallel( wait_transaction_confirmed: bool, with_tip: bool, gas_fee_strategy: GasFeeStrategy, + use_core_affinity: bool, ) -> Result<(bool, Vec, Option)> { let _exec_start = Instant::now(); @@ -224,10 +236,10 @@ pub async fn execute_parallel( return Err(anyhow!("No Rpc Default Swqos configured.")); } - let cores = core_affinity::get_core_ids().unwrap(); + let cores = core_affinity::get_core_ids().unwrap_or_default(); let instructions = Arc::new(instructions); - // 预先计算所有有效的组合 + // Precompute all valid (client, gas config) combinations let task_configs: Vec<_> = swqos_clients .iter() .enumerate() @@ -244,7 +256,7 @@ pub async fn execute_parallel( .into_iter() .filter(|config| config.0.eq(&swqos_client.get_swqos_type())) .filter(|config| { - // 当需要 tip 且不是 Default 时,按 provider 最低小费进行筛选 + // When tip required and not Default, filter by provider minimum tip if with_tip && !matches!(config.0, SwqosType::Default) { let min_tip = match config.0 { SwqosType::Jito => SWQOS_MIN_TIP_JITO, @@ -262,9 +274,9 @@ pub async fn execute_parallel( SwqosType::Speedlanding => SWQOS_MIN_TIP_SPEEDLANDING, SwqosType::Default => SWQOS_MIN_TIP_DEFAULT, }; - if config.2.tip < min_tip { + if config.2.tip < min_tip && crate::common::sdk_log::sdk_log_enabled() { println!( - "⚠️ Config filtered: {:?} tip {} is below minimum required tip {}", + "⚠️ Config filtered: {:?} tip {} is below minimum required {}", config.0, config.2.tip, min_tip ); } @@ -291,7 +303,8 @@ pub async fn execute_parallel( let _spawn_start = Instant::now(); for (i, swqos_client, gas_fee_strategy_config) in task_configs { - let core_id = cores[i % cores.len()]; + let core_id = cores.get(i % cores.len().max(1)).copied(); + let use_affinity = use_core_affinity; let payer = payer.clone(); let instructions = instructions.clone(); let middleware_manager = middleware_manager.clone(); @@ -306,10 +319,15 @@ pub async fn execute_parallel( let rpc = rpc.clone(); let durable_nonce = durable_nonce.clone(); let address_lookup_table_account = address_lookup_table_account.clone(); + let recent_blockhash_task = recent_blockhash.clone(); tokio::spawn(async move { let _task_start = Instant::now(); - core_affinity::set_for_current(core_id); + if use_affinity { + if let Some(cid) = core_id { + core_affinity::set_for_current(cid); + } + } let tip_amount = if with_tip { tip } else { 0.0 }; @@ -319,9 +337,9 @@ pub async fn execute_parallel( rpc, unit_limit, unit_price, - instructions.as_ref().clone(), + instructions.as_ref(), address_lookup_table_account, - recent_blockhash, + recent_blockhash_task, middleware_manager, protocol_name, is_buy, @@ -339,8 +357,8 @@ pub async fn execute_parallel( success: false, signature: Signature::default(), error: Some(e), - swqos_type, // 🔧 记录SWQOS类型 - landed_on_chain: false, // Build failed, tx never sent + swqos_type, + landed_on_chain: false, }); return; } @@ -373,28 +391,28 @@ pub async fn execute_parallel( } }; - // Transaction sent - - if let Some(signature) = transaction.signatures.first() { - collector.submit(TaskResult { - success, - signature: *signature, - error: err, - swqos_type, // 🔧 记录SWQOS类型 - landed_on_chain, // 🔧 Whether tx landed (even if it failed) - }); - } + // Transaction sent: always submit a result so collector never has "no result" for this task. + // If transaction has no signatures (malformed), submit with default signature and success=false. + let sig = transaction.signatures.first().copied().unwrap_or_default(); + collector.submit(TaskResult { + success, + signature: sig, + error: err, + swqos_type, + landed_on_chain, + }); }); } // All tasks spawned if !wait_transaction_confirmed { - tokio::time::sleep(std::time::Duration::from_millis(10)).await; - if let Some(result) = collector.get_first() { - return Ok(result); - } - return Err(anyhow!("No transaction signature available")); + const SUBMIT_TIMEOUT_SECS: u64 = 30; + let (success, signatures, last_error) = collector + .wait_for_all_submitted(SUBMIT_TIMEOUT_SECS) + .await + .unwrap_or((false, vec![], Some(anyhow!("No SWQOS result within {}s", SUBMIT_TIMEOUT_SECS)))); + return Ok((success, signatures, last_error)); } if let Some(result) = collector.wait_for_success().await { diff --git a/src/trading/core/execution.rs b/src/trading/core/execution.rs index 22069a8..6cede9b 100644 --- a/src/trading/core/execution.rs +++ b/src/trading/core/execution.rs @@ -1,4 +1,5 @@ -//! 执行模块 +//! Execution: instruction preprocessing, cache prefetch, branch hints. +//! 执行模块:指令预处理、缓存预取、分支提示。 use anyhow::Result; use solana_sdk::{ @@ -12,7 +13,15 @@ use crate::perf::{ simd::SIMDMemory, }; -/// 预取工具 +/// Solana account key size in bytes (Pubkey). 每个账户(Pubkey)的字节数。 +pub const BYTES_PER_ACCOUNT: usize = 32; + +/// Threshold above which we warn about large instruction count. 超过此次数会打 warning。 +pub const MAX_INSTRUCTIONS_WARN: usize = 64; + +/// Prefetch helper: triggers CPU prefetch for soon-to-be-accessed data to reduce cache-miss latency. +/// Call once on hot-path refs; no-op on non-x86_64. Safety: caller ensures valid read-only ref, no concurrent write. +/// 缓存预取:对即将访问的数据做 CPU 预取以降低 cache-miss;热路径上调用一次即可;非 x86_64 为 no-op。安全:调用方保证有效只读、无并发写。 pub struct Prefetch; impl Prefetch { @@ -21,20 +30,15 @@ impl Prefetch { if instructions.is_empty() { return; } - - // 预取第一条指令 + // Prefetch first/middle/last instruction into L1 for subsequent build_transaction access. 预取首/中/尾指令到 L1。 unsafe { BranchOptimizer::prefetch_read_data(&instructions[0]); } - - // 预取中间指令 if instructions.len() > 2 { unsafe { BranchOptimizer::prefetch_read_data(&instructions[instructions.len() / 2]); } } - - // 预取最后一条指令 if instructions.len() > 1 { unsafe { BranchOptimizer::prefetch_read_data(&instructions[instructions.len() - 1]); @@ -57,46 +61,40 @@ impl Prefetch { } } -/// 内存操作 +/// Memory operations (SIMD-accelerated where available). 内存操作(可用时走 SIMD 加速)。 pub struct MemoryOps; impl MemoryOps { #[inline(always)] pub unsafe fn copy(dst: *mut u8, src: *const u8, len: usize) { - // 优先使用 AVX2 SIMD 加速 SIMDMemory::copy_avx2(dst, src, len); } #[inline(always)] pub unsafe fn compare(a: *const u8, b: *const u8, len: usize) -> bool { - // 优先使用 AVX2 SIMD 比较 SIMDMemory::compare_avx2(a, b, len) } #[inline(always)] pub unsafe fn zero(ptr: *mut u8, len: usize) { - // 优先使用 AVX2 SIMD 清零 SIMDMemory::zero_avx2(ptr, len); } } -/// 指令处理器 +/// Instruction preprocessing and validation. 指令预处理与校验。 pub struct InstructionProcessor; impl InstructionProcessor { #[inline(always)] pub fn preprocess(instructions: &[Instruction]) -> Result<()> { - // 分支预测: 大概率指令不为空 if BranchOptimizer::unlikely(instructions.is_empty()) { return Err(anyhow::anyhow!("Instructions empty")); } - // 预取所有指令到缓存 Prefetch::instructions(instructions); - // 分支预测: 大概率指令数量合理 - if BranchOptimizer::unlikely(instructions.len() > 64) { - log::warn!("Large instruction count: {}", instructions.len()); + if BranchOptimizer::unlikely(instructions.len() > MAX_INSTRUCTIONS_WARN) { + tracing::warn!(target: "sol_trade_sdk", "Large instruction count: {}", instructions.len()); } Ok(()) @@ -107,7 +105,7 @@ impl InstructionProcessor { let mut total_size = 0; for (i, instr) in instructions.iter().enumerate() { - // 预取下一条指令 + // Prefetch next instruction; safe: same slice, read-only. 预取下一条指令;安全:同 slice、只读。 unsafe { if let Some(next_instr) = instructions.get(i + 1) { BranchOptimizer::prefetch_read_data(next_instr); @@ -115,20 +113,19 @@ impl InstructionProcessor { } total_size += instr.data.len(); - total_size += instr.accounts.len() * 32; // 每个账户 32 字节 + total_size += instr.accounts.len() * BYTES_PER_ACCOUNT; } total_size } } -/// 执行路径 +/// Trade direction / execution path helpers. 交易方向与执行路径辅助。 pub struct ExecutionPath; impl ExecutionPath { #[inline(always)] pub fn is_buy(input_mint: &Pubkey) -> bool { - // 分支预测: 大概率是买入 let is_buy = input_mint == &crate::constants::SOL_TOKEN_ACCOUNT || input_mint == &crate::constants::WSOL_TOKEN_ACCOUNT || input_mint == &crate::constants::USD1_TOKEN_ACCOUNT diff --git a/src/trading/core/executor.rs b/src/trading/core/executor.rs index 4b901e0..4c4b451 100755 --- a/src/trading/core/executor.rs +++ b/src/trading/core/executor.rs @@ -4,11 +4,14 @@ use solana_sdk::{ instruction::Instruction, message::AddressLookupTableAccount, pubkey::Pubkey, signature::Keypair, signature::Signature, }; -use std::{sync::Arc, time::Instant}; +use std::{sync::Arc, time::{Duration, Instant}}; +#[allow(unused_imports)] +use tracing::{info, trace, warn}; use crate::{ common::{nonce_cache::DurableNonceInfo, GasFeeStrategy, SolanaRpcClient}, perf::syscall_bypass::SystemCallBypassManager, + swqos::common::poll_transaction_confirmation, trading::core::{ async_executor::execute_parallel, execution::{InstructionProcessor, Prefetch}, @@ -20,7 +23,9 @@ use once_cell::sync::Lazy; use crate::swqos::TradeType; use super::{params::SwapParams, traits::InstructionBuilder}; -/// 🚀 全局系统调用绕过管理器 +/// Global syscall bypass manager (reserved for future time/IO optimizations). +/// 全局系统调用绕过管理器(预留,后续可接入时间/IO 等优化)。 +#[allow(dead_code)] static SYSCALL_BYPASS: Lazy = Lazy::new(|| { use crate::perf::syscall_bypass::SyscallBypassConfig; SystemCallBypassManager::new(SyscallBypassConfig::default()) @@ -45,27 +50,29 @@ impl GenericTradeExecutor { #[async_trait::async_trait] impl TradeExecutor for GenericTradeExecutor { async fn swap(&self, params: SwapParams) -> Result<(bool, Vec, Option)> { - let total_start = Instant::now(); + // Sample total start only when logging or simulate. 仅在有日志或 simulate 时取起点。 + let total_start = (params.log_enabled || params.simulate).then(Instant::now); + let timing_start_us: Option = if params.log_enabled { + Some(params.grpc_recv_us.unwrap_or_else(crate::common::clock::now_micros)) + } else { + None + }; - // 判断买卖方向 let is_buy = params.trade_type == TradeType::Buy || params.trade_type == TradeType::CreateAndBuy; - // CPU 预取 Prefetch::keypair(¶ms.payer); - // 构建指令 - let build_start = Instant::now(); + // Time build only when log_enabled to avoid cold-path syscalls. 仅 log_enabled 时计时,减少冷路径 syscall。 + let build_start = params.log_enabled.then(Instant::now); let instructions = if is_buy { self.instruction_builder.build_buy_instructions(¶ms).await? } else { self.instruction_builder.build_sell_instructions(¶ms).await? }; - let build_elapsed = build_start.elapsed(); + let build_elapsed = build_start.map(|s| s.elapsed()).unwrap_or(Duration::ZERO); - // 指令预处理 InstructionProcessor::preprocess(&instructions)?; - // 中间件处理 let final_instructions = match ¶ms.middleware_manager { Some(middleware_manager) => middleware_manager .apply_middlewares_process_protocol_instructions( @@ -76,12 +83,10 @@ impl TradeExecutor for GenericTradeExecutor { None => instructions, }; - // 提交前耗时 - let before_submit_elapsed = total_start.elapsed(); + let before_submit_elapsed = total_start.as_ref().map(|s| s.elapsed()).unwrap_or(Duration::ZERO); - // 如果是模拟模式,直接通过 RPC 模拟交易 if params.simulate { - let send_start = Instant::now(); + let send_start = crate::common::sdk_log::sdk_log_enabled().then(Instant::now); let result = simulate_transaction( params.rpc, params.payer, @@ -96,44 +101,23 @@ impl TradeExecutor for GenericTradeExecutor { params.gas_fee_strategy, ) .await; - let send_elapsed = send_start.elapsed(); - let total_elapsed = total_start.elapsed(); + let send_elapsed = send_start.map(|s| s.elapsed()).unwrap_or(Duration::ZERO); + let total_elapsed = total_start.as_ref().map(|s| s.elapsed()).unwrap_or(Duration::ZERO); - // Get performance metrics using fast timestamp - let timestamp_ns = SYSCALL_BYPASS.fast_timestamp_nanos(); - - // Print all timing metrics at once to avoid blocking critical path - println!("[Timestamp] {}ns", timestamp_ns); - println!( - "[Build Instructions] Time: {:.3}ms ({:.0}μs)", - build_elapsed.as_micros() as f64 / 1000.0, - build_elapsed.as_micros() - ); - println!( - "[Before Submit] {:.3}ms ({:.0}μs)", - before_submit_elapsed.as_micros() as f64 / 1000.0, - before_submit_elapsed.as_micros() - ); - println!( - "[Simulate Transaction] Time: {:.3}ms ({:.0}μs)", - send_elapsed.as_micros() as f64 / 1000.0, - send_elapsed.as_micros() - ); - println!( - "[Total Time] {:.3}ms ({:.0}μs)", - total_elapsed.as_micros() as f64 / 1000.0, - total_elapsed.as_micros() - ); + if crate::common::sdk_log::sdk_log_enabled() { + let dir = if is_buy { "Buy" } else { "Sell" }; + println!(" [SDK] {} timing(sim) build_instructions: {:.2}ms before_submit: {:.2}ms simulate: {:.2}ms total: {:.2}ms", dir, build_elapsed.as_secs_f64() * 1000.0, before_submit_elapsed.as_secs_f64() * 1000.0, send_elapsed.as_secs_f64() * 1000.0, total_elapsed.as_secs_f64() * 1000.0); + } return result; } - // 并行发送交易 - let send_start = Instant::now(); + let need_confirm = params.wait_transaction_confirmed; + let send_start = params.log_enabled.then(Instant::now); let result = execute_parallel( - params.swqos_clients.clone(), + ¶ms.swqos_clients, params.payer, - params.rpc, + params.rpc.clone(), final_instructions, params.address_lookup_table_account, params.recent_blockhash, @@ -141,30 +125,64 @@ impl TradeExecutor for GenericTradeExecutor { params.middleware_manager, self.protocol_name, is_buy, - params.wait_transaction_confirmed, + false, // submit only here; confirmation and log timing handled below if is_buy { true } else { params.with_tip }, params.gas_fee_strategy, + params.use_core_affinity, ) .await; - let send_elapsed = send_start.elapsed(); - let total_elapsed = total_start.elapsed(); + let send_elapsed = send_start.map(|s| s.elapsed()).unwrap_or(Duration::ZERO); - // Get performance metrics using fast timestamp - #[cfg(feature = "perf-trace")] - { - let timestamp_ns = SYSCALL_BYPASS.fast_timestamp_nanos(); - log::trace!( - "[Execute] timestamp_ns={} build_us={} before_submit_us={} send_us={} total_us={}", - timestamp_ns, - build_elapsed.as_micros(), - before_submit_elapsed.as_micros(), - send_elapsed.as_micros(), - total_elapsed.as_micros() - ); + if params.log_enabled && crate::common::sdk_log::sdk_log_enabled() { + let dir = if is_buy { "Buy" } else { "Sell" }; + let build_ms = build_elapsed.as_secs_f64() * 1000.0; + let before_ms = before_submit_elapsed.as_secs_f64() * 1000.0; + let send_ms = send_elapsed.as_secs_f64() * 1000.0; + if let Some(start_us) = timing_start_us { + let now_us = crate::common::clock::now_micros(); + let start_to_submit_us = (now_us - start_us).max(0); + println!(" [SDK] {} timing(after_submit) build_instructions: {:.2}ms before_submit: {:.2}ms submit: {:.2}ms start_to_submit: {} μs", dir, build_ms, before_ms, send_ms, start_to_submit_us); + } else { + println!(" [SDK] {} timing(after_submit) build_instructions: {:.2}ms before_submit: {:.2}ms submit: {:.2}ms", dir, build_ms, before_ms, send_ms); + } } - - #[cfg(not(feature = "perf-trace"))] - let _ = (build_elapsed, before_submit_elapsed, send_elapsed, total_elapsed); + + let result = if need_confirm { + let (ok, sigs, err) = match &result { + Ok((success, signatures, last_error)) => ( + *success, + signatures.clone(), + last_error.as_ref().map(|e| anyhow::anyhow!("{}", e)), + ), + Err(e) => (false, vec![], Some(anyhow::anyhow!("{}", e))), + }; + let first_sig = sigs.first().copied(); + let confirm_result = if let (Some(rpc), Some(sig)) = (params.rpc.as_ref(), first_sig) { + let confirm_start = (params.log_enabled && crate::common::sdk_log::sdk_log_enabled()).then(Instant::now); + let poll_res = poll_transaction_confirmation(rpc, sig, true).await; + let confirm_elapsed = confirm_start.map(|s| s.elapsed()).unwrap_or(Duration::ZERO); + if params.log_enabled && crate::common::sdk_log::sdk_log_enabled() { + let dir = if is_buy { "Buy" } else { "Sell" }; + let confirm_ms = confirm_elapsed.as_secs_f64() * 1000.0; + let total_ms = total_start.as_ref().map(|s| s.elapsed()).unwrap_or(Duration::ZERO).as_secs_f64() * 1000.0; + println!(" [SDK] {} timing(after_confirm) confirm: {:.2}ms total: {:.2}ms", dir, confirm_ms, total_ms); + } + match poll_res { + Ok(_) => (true, sigs, None), + Err(e) => (false, sigs, Some(e)), + } + } else { + (ok, sigs, err) + }; + Ok(confirm_result) + } else { + if params.log_enabled && crate::common::sdk_log::sdk_log_enabled() { + let total_ms = total_start.as_ref().map(|s| s.elapsed()).unwrap_or(Duration::ZERO).as_secs_f64() * 1000.0; + let dir = if is_buy { "Buy" } else { "Sell" }; + println!(" [SDK] {} timing total: {:.2}ms", dir, total_ms); + } + result + }; result } @@ -174,7 +192,8 @@ impl TradeExecutor for GenericTradeExecutor { } } -/// 🔧 修复:Simulate模式返回Vec(单个RPC模拟) +/// Simulate mode: single RPC simulation, returns Vec for API consistency. +/// 模拟模式:单次 RPC 模拟,返回 Vec 以与 API 一致。 async fn simulate_transaction( rpc: Option>, payer: Arc, @@ -215,7 +234,7 @@ async fn simulate_transaction( Some(rpc.clone()), unit_limit, unit_price, - instructions, + &instructions, address_lookup_table_account, recent_blockhash, middleware_manager, @@ -256,12 +275,12 @@ async fn simulate_transaction( if let Some(err) = simulate_result.value.err { #[cfg(feature = "perf-trace")] { - log::warn!("[Simulation Failed] error={:?} signature={:?}", err, signature); + warn!(target: "sol_trade_sdk", "[Simulation Failed] error={:?} signature={:?}", err, signature); if let Some(logs) = &simulate_result.value.logs { - log::trace!("Transaction logs: {:?}", logs); + trace!(target: "sol_trade_sdk", "Transaction logs: {:?}", logs); } if let Some(units_consumed) = simulate_result.value.units_consumed { - log::trace!("Compute Units Consumed: {}", units_consumed); + trace!(target: "sol_trade_sdk", "Compute Units Consumed: {}", units_consumed); } } return Ok((false, vec![signature], Some(anyhow::anyhow!("{:?}", err)))); @@ -270,12 +289,12 @@ async fn simulate_transaction( // Simulation succeeded #[cfg(feature = "perf-trace")] { - log::info!("[Simulation Succeeded] signature={:?}", signature); + info!(target: "sol_trade_sdk", "[Simulation Succeeded] signature={:?}", signature); if let Some(units_consumed) = simulate_result.value.units_consumed { - log::trace!("Compute Units Consumed: {}", units_consumed); + trace!(target: "sol_trade_sdk", "Compute Units Consumed: {}", units_consumed); } if let Some(logs) = &simulate_result.value.logs { - log::trace!("Transaction logs: {:?}", logs); + trace!(target: "sol_trade_sdk", "Transaction logs: {:?}", logs); } } diff --git a/src/trading/core/params.rs b/src/trading/core/params.rs index 2e40fb3..cc24503 100755 --- a/src/trading/core/params.rs +++ b/src/trading/core/params.rs @@ -67,6 +67,12 @@ pub struct SwapParams { pub fixed_output_amount: Option, pub gas_fee_strategy: GasFeeStrategy, pub simulate: bool, + /// Whether to output SDK logs (from TradeConfig.log_enabled). + pub log_enabled: bool, + /// Whether to pin parallel submit tasks to cores (from TradeConfig.use_core_affinity). + pub use_core_affinity: bool, + /// Optional event receive time in microseconds (same scale as sol-parser-sdk clock::now_micros). Used as timing start when log_enabled. + pub grpc_recv_us: Option, /// Use exact SOL amount instructions (buy_exact_sol_in for PumpFun, buy_exact_quote_in for PumpSwap). /// When Some(true) or None (default), the exact SOL/quote amount is spent and slippage is applied to output tokens. /// When Some(false), uses regular buy instruction where slippage is applied to SOL/quote input.