Release v3.5.0: performance, constants, bilingual docs
- Bump version to 3.5.0 - Performance: hot-path timing only when log_enabled/simulate; execute_parallel takes &[Arc<SwqosClient>]; 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 <cursoragent@cursor.com>
This commit is contained in:
@@ -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 }}
|
||||
@@ -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
|
||||
- 增强异步执行器,添加小费验证警告
|
||||
+1
-2
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "sol-trade-sdk"
|
||||
version = "3.4.1"
|
||||
version = "3.5.0"
|
||||
edition = "2021"
|
||||
authors = [
|
||||
"William <byteblock6@gmail.com>",
|
||||
@@ -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"
|
||||
|
||||
@@ -60,6 +60,14 @@
|
||||
|
||||
---
|
||||
|
||||
## 🆕 What's new in 3.5.0
|
||||
|
||||
- **Performance**: Hot-path timing only when logging; reduced clones (`execute_parallel` takes `&[Arc<SwqosClient>]`); 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
|
||||
|
||||
+10
-2
@@ -60,6 +60,14 @@
|
||||
|
||||
---
|
||||
|
||||
## 🆕 3.5.0 更新说明
|
||||
|
||||
- **性能**:仅在打日志时做热路径计时;减少 clone(`execute_parallel` 改为接收 `&[Arc<SwqosClient>]`);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"
|
||||
```
|
||||
|
||||
## 🛠️ 使用示例
|
||||
|
||||
@@ -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
|
||||
@@ -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<SwqosClient>]` instead of `Vec<Arc<SwqosClient>>`; 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
|
||||
@@ -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<HighPerformanceClock> =
|
||||
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
|
||||
}
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
+3
-1
@@ -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::*;
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
+12
-4
@@ -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<SwqosConfig>,
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+107
-76
@@ -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::<PumpFunParams>().is_some(),
|
||||
DexType::PumpSwap => params.as_any().downcast_ref::<PumpSwapParams>().is_some(),
|
||||
DexType::Bonk => params.as_any().downcast_ref::<BonkParams>().is_some(),
|
||||
DexType::RaydiumCpmm => params.as_any().downcast_ref::<RaydiumCpmmParams>().is_some(),
|
||||
DexType::RaydiumAmmV4 => params.as_any().downcast_ref::<RaydiumAmmV4Params>().is_some(),
|
||||
DexType::MeteoraDammV2 => params.as_any().downcast_ref::<MeteoraDammV2Params>().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<Option<Arc<TradingClient>>> = 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<bool>,
|
||||
/// 可选:事件收到时间(微秒,与 sol-parser-sdk 的 metadata.grpc_recv_us / clock::now_micros 同源)。不传且开启 log_enabled 时 SDK 用 now_micros() 作为起点,打印起点→提交耗时。
|
||||
pub grpc_recv_us: Option<i64>,
|
||||
}
|
||||
|
||||
/// 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<i64>,
|
||||
}
|
||||
|
||||
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<Keypair>, 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<Signature>, Option<TradeError>), 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::<PumpFunParams>().is_some(),
|
||||
DexType::PumpSwap => {
|
||||
protocol_params.as_any().downcast_ref::<PumpSwapParams>().is_some()
|
||||
}
|
||||
DexType::Bonk => protocol_params.as_any().downcast_ref::<BonkParams>().is_some(),
|
||||
DexType::RaydiumCpmm => {
|
||||
protocol_params.as_any().downcast_ref::<RaydiumCpmmParams>().is_some()
|
||||
}
|
||||
DexType::RaydiumAmmV4 => {
|
||||
protocol_params.as_any().downcast_ref::<RaydiumAmmV4Params>().is_some()
|
||||
}
|
||||
DexType::MeteoraDammV2 => {
|
||||
protocol_params.as_any().downcast_ref::<MeteoraDammV2Params>().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<Signature>, Option<TradeError>), 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::<PumpFunParams>().is_some(),
|
||||
DexType::PumpSwap => {
|
||||
protocol_params.as_any().downcast_ref::<PumpSwapParams>().is_some()
|
||||
}
|
||||
DexType::Bonk => protocol_params.as_any().downcast_ref::<BonkParams>().is_some(),
|
||||
DexType::RaydiumCpmm => {
|
||||
protocol_params.as_any().downcast_ref::<RaydiumCpmmParams>().is_some()
|
||||
}
|
||||
DexType::RaydiumAmmV4 => {
|
||||
protocol_params.as_any().downcast_ref::<RaydiumAmmV4Params>().is_some()
|
||||
}
|
||||
DexType::MeteoraDammV2 => {
|
||||
protocol_params.as_any().downcast_ref::<MeteoraDammV2Params>().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)));
|
||||
|
||||
@@ -140,7 +140,7 @@ impl CompilerOptimizer {
|
||||
|
||||
/// 🚀 生成超高性能编译配置
|
||||
pub fn generate_ultra_performance_config(&self) -> Result<CompilerConfig> {
|
||||
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)
|
||||
}
|
||||
|
||||
|
||||
@@ -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::<AtomicU64>()],
|
||||
}
|
||||
|
||||
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<T> {
|
||||
/// 数据缓冲区
|
||||
buffer: Vec<T>,
|
||||
/// 生产者头指针 (独占缓存行)
|
||||
producer_head: CachePadded<AtomicU64>,
|
||||
/// 消费者尾指针 (独占缓存行)
|
||||
consumer_tail: CachePadded<AtomicU64>,
|
||||
/// 容量 (2的幂次方)
|
||||
capacity: usize,
|
||||
/// 掩码 (capacity - 1)
|
||||
mask: usize,
|
||||
}
|
||||
|
||||
impl<T: Copy + Default> CacheOptimizedRingBuffer<T> {
|
||||
/// 创建缓存优化的环形缓冲区
|
||||
/// Create ring buffer; capacity must be a power of 2. 创建环形缓冲区,容量须为 2 的幂。
|
||||
pub fn new(capacity: usize) -> Result<Self> {
|
||||
if !capacity.is_power_of_two() {
|
||||
return Err(anyhow::anyhow!("Capacity must be a power of 2"));
|
||||
@@ -373,53 +352,41 @@ impl<T: Copy + Default> CacheOptimizedRingBuffer<T> {
|
||||
})
|
||||
}
|
||||
|
||||
/// 🚀 无锁写入元素
|
||||
/// 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<T> {
|
||||
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<T: Copy + Default> CacheOptimizedRingBuffer<T> {
|
||||
((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<T> CacheLineAligned for CacheOptimizedRingBuffer<T> {
|
||||
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<T>(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<T>(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);
|
||||
|
||||
+2
-8
@@ -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;
|
||||
|
||||
+27
-59
@@ -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<SyscallBatchProcessor>,
|
||||
/// 快速时间获取器
|
||||
fast_time_provider: Arc<FastTimeProvider>,
|
||||
/// I/O优化器
|
||||
_io_optimizer: Arc<IOOptimizer>,
|
||||
/// 统计信息
|
||||
stats: Arc<SyscallBypassStats>,
|
||||
}
|
||||
|
||||
/// 系统调用绕过配置
|
||||
/// 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<SyscallRequest>,
|
||||
/// 批处理线程池
|
||||
_executor: tokio::runtime::Handle,
|
||||
/// 批处理统计
|
||||
batch_stats: CachePadded<AtomicU64>,
|
||||
}
|
||||
|
||||
/// 系统调用请求
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum SyscallRequest {
|
||||
/// 文件写入
|
||||
Write { fd: i32, data: Vec<u8> },
|
||||
/// 文件读取
|
||||
Read { fd: i32, size: usize },
|
||||
/// 网络发送
|
||||
Send { socket: i32, data: Vec<u8> },
|
||||
/// 网络接收
|
||||
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<Self> {
|
||||
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<Vec<usize>> {
|
||||
// 这里是伪代码 - 实际实现需要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<u8>)>) -> 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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+12
-12
@@ -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,
|
||||
|
||||
+18
-10
@@ -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::<serde_json::Value>(&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());
|
||||
}
|
||||
|
||||
+22
-23
@@ -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::<serde_json::Value>(&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::<serde_json::Value>(&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::<serde_json::Value>(&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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
+21
-25
@@ -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::<serde_json::Value>(&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());
|
||||
}
|
||||
|
||||
@@ -105,7 +105,41 @@ impl Base64Encoder {
|
||||
}
|
||||
}
|
||||
|
||||
/// Guard that returns the serialization buffer to the pool on drop.
|
||||
pub struct PooledTxBufGuard(pub Vec<u8>);
|
||||
|
||||
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<u8>) {
|
||||
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")),
|
||||
};
|
||||
|
||||
|
||||
+16
-11
@@ -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());
|
||||
}
|
||||
|
||||
+19
-25
@@ -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::<serde_json::Value>(&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());
|
||||
}
|
||||
|
||||
@@ -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<Keypair>,
|
||||
_rpc: Option<Arc<SolanaRpcClient>>,
|
||||
unit_limit: u32,
|
||||
unit_price: u64,
|
||||
business_instructions: Vec<Instruction>,
|
||||
business_instructions: &[Instruction],
|
||||
address_lookup_table_account: Option<AddressLookupTableAccount>,
|
||||
recent_blockhash: Option<Hash>,
|
||||
middleware_manager: Option<Arc<MiddlewareManager>>,
|
||||
@@ -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());
|
||||
|
||||
@@ -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<anyhow::Error>,
|
||||
#[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<Signature>, Option<anyhow::Error>)> {
|
||||
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<Signature>, Option<anyhow::Error>)> {
|
||||
// 🔧 修复:收集已提交的所有签名
|
||||
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<Signature>, Option<anyhow::Error>)> {
|
||||
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<Signature>支持多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<Arc<SwqosClient>>,
|
||||
swqos_clients: &[Arc<SwqosClient>],
|
||||
payer: Arc<Keypair>,
|
||||
rpc: Option<Arc<SolanaRpcClient>>,
|
||||
instructions: Vec<Instruction>,
|
||||
@@ -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<Signature>, Option<anyhow::Error>)> {
|
||||
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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<SystemCallBypassManager> = 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<Signature>, Option<anyhow::Error>)> {
|
||||
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<i64> = 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<Signature>(单个RPC模拟)
|
||||
/// Simulate mode: single RPC simulation, returns Vec<Signature> for API consistency.
|
||||
/// 模拟模式:单次 RPC 模拟,返回 Vec<Signature> 以与 API 一致。
|
||||
async fn simulate_transaction(
|
||||
rpc: Option<Arc<SolanaRpcClient>>,
|
||||
payer: Arc<Keypair>,
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -67,6 +67,12 @@ pub struct SwapParams {
|
||||
pub fixed_output_amount: Option<u64>,
|
||||
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<i64>,
|
||||
/// 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.
|
||||
|
||||
Reference in New Issue
Block a user