fix(swqos): Speedlanding QUIC TLS, quinn runtime, and SWQoS fallback

- Enable quinn feature runtime-tokio so QUIC clients (Speedlanding/Soyas) initialize instead of failing with "no async runtime found".
- Align Speedlanding with the official client: use solana_tls_utils new_dummy_x509_certificate and SkipServerVerification; use fixed TLS SNI "speed-landing" (not the PoP hostname).
- Raise per-SWQoS init timeout to 15s; always log init failures/timeouts; if every SWQoS client fails, fall back to Default RPC and print active channel labels.
- Improve SWQoS submission visibility (Default RPC path and related logging); adjust Soyas TLS credential handling where applicable.

Made-with: Cursor
This commit is contained in:
0xfnzero
2026-04-16 20:11:41 +08:00
parent d711346c55
commit 06ed710869
5 changed files with 146 additions and 130 deletions
+2 -1
View File
@@ -109,7 +109,8 @@ parking_lot = "0.12"
arc-swap = "1.7"
sha2 = "0.10"
tonic-prost = "0.14.2"
quinn = { version = "0.11", default-features = false, features = ["rustls"] }
# 须含 runtime-tokio,否则 quinn::Endpoint::client 报 no async runtime foundQUICSpeedlanding/Soyas)无法初始化
quinn = { version = "0.11", default-features = false, features = ["rustls", "runtime-tokio"] }
rcgen = "0.13"
uuid = "1.11"
+58 -2
View File
@@ -136,8 +136,8 @@ impl TradingInfrastructure {
}
common::seed::start_rent_updater(rpc.clone());
// Create SWQOS clients with blacklist checking(单节点超时 5s,避免某一家卡死整段初始化
const SWQOS_CLIENT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
// Create SWQOS clients with blacklist checkingQUIC 握手可能较慢,单节点超时 15s
const SWQOS_CLIENT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(15);
let mut swqos_clients: Vec<Arc<SwqosClient>> = vec![];
for swqos in &config.swqos_configs {
if swqos.is_blacklisted() {
@@ -159,6 +159,11 @@ impl TradingInfrastructure {
{
Ok(Ok(swqos_client)) => swqos_clients.push(swqos_client),
Ok(Err(err)) => {
eprintln!(
"⚠️ SWQOS {:?} 初始化失败: {}(已从列表中排除)",
swqos.swqos_type(),
err
);
if sdk_log::sdk_log_enabled() {
warn!(
target: "sol_trade_sdk",
@@ -168,6 +173,11 @@ impl TradingInfrastructure {
}
}
Err(_) => {
eprintln!(
"⚠️ SWQOS {:?} 初始化超时({}s),已跳过",
swqos.swqos_type(),
SWQOS_CLIENT_TIMEOUT.as_secs()
);
if sdk_log::sdk_log_enabled() {
warn!(
target: "sol_trade_sdk",
@@ -180,6 +190,52 @@ impl TradingInfrastructure {
}
}
// 若全部失败、被黑名单跳过或仅配置了不可用通道,至少保留一条 Rpc Default,否则 execute_parallel 会因 swqos_clients 为空直接报错。
if swqos_clients.is_empty() {
eprintln!(
"⚠️ 无任何 SWQOS 客户端初始化成功,将回退为普通 RPC 发送: {}",
config.rpc_url
);
if sdk_log::sdk_log_enabled() {
warn!(
target: "sol_trade_sdk",
"no SWQOS clients initialized; falling back to Rpc Default ({})",
config.rpc_url
);
}
match SwqosConfig::get_swqos_client(
config.rpc_url.clone(),
config.commitment.clone(),
SwqosConfig::Default(config.rpc_url.clone()),
config.mev_protection,
)
.await
{
Ok(c) => swqos_clients.push(c),
Err(e) => {
if sdk_log::sdk_log_enabled() {
warn!(
target: "sol_trade_sdk",
"fallback Rpc Default client failed: {}",
e
);
}
}
}
}
if !swqos_clients.is_empty() {
let labels: Vec<&str> = swqos_clients
.iter()
.map(|c| c.get_swqos_type().as_str())
.collect();
eprintln!(
"️ SWQOS 通道已就绪: {} 条 → [{}]",
swqos_clients.len(),
labels.join(", ")
);
}
let swqos_count = swqos_clients.len();
let (max_sender_concurrency, effective_core_ids) = {
let num_cores = core_affinity::get_core_ids().map(|c| c.len()).unwrap_or(0);
+4 -1
View File
@@ -7,7 +7,7 @@ use solana_transaction_status::UiTransactionEncoding;
use crate::swqos::SwqosClientTrait;
use crate::{
common::SolanaRpcClient,
common::{sdk_log, SolanaRpcClient},
swqos::{common::poll_transaction_confirmation, SwqosType, TradeType},
};
use anyhow::Result;
@@ -25,6 +25,7 @@ impl SwqosClientTrait for SolRpcClient {
transaction: &VersionedTransaction,
wait_confirmation: bool,
) -> Result<()> {
let submit_start = Instant::now();
let signature = self
.rpc_client
.send_transaction_with_config(
@@ -39,6 +40,8 @@ impl SwqosClientTrait for SolRpcClient {
)
.await?;
sdk_log::log_swqos_submitted("Default", trade_type, submit_start.elapsed());
let start_time = Instant::now();
match poll_transaction_confirmation(&self.rpc_client, signature, wait_confirmation).await {
Ok(_) => (),
+21 -14
View File
@@ -6,7 +6,10 @@ use quinn::{
TransportConfig,
};
use rand::seq::IndexedRandom as _;
use rcgen::{CertificateParams, KeyPair as RcgenKeyPair};
use rustls::pki_types::{CertificateDer, PrivateKeyDer, PrivatePkcs8KeyDer};
use solana_client::rpc_client::SerializableTransaction;
use solana_sdk::signer::Signer;
use solana_sdk::{signature::Keypair, transaction::VersionedTransaction};
use std::time::Instant;
use std::{
@@ -69,18 +72,17 @@ impl rustls::client::danger::ServerCertVerifier for SkipServerVerification {
}
}
// Generate dummy self-signed certificate using rcgen
fn generate_self_signed_cert(_keypair: &Keypair) -> Result<(rustls::pki_types::CertificateDer<'static>, rustls::pki_types::PrivateKeyDer<'static>)> {
// Generate a new key pair for the certificate
let key_pair = rcgen::KeyPair::generate()?;
let params = rcgen::CertificateParams::default();
let cert = params.self_signed(&key_pair)?;
let cert_der = rustls::pki_types::CertificateDer::from(cert.der().to_vec());
let key_der = rustls::pki_types::PrivateKeyDer::try_from(key_pair.serialize_der())
.map_err(|e| anyhow::anyhow!("Failed to create private key: {:?}", e))?;
/// TLS 客户端证书:ECDSA P-256 + CN=钱包公钥(与 Speedlanding / Astralane QUIC 策略一致)。
fn generate_client_tls_credentials(keypair: &Keypair) -> Result<(CertificateDer<'static>, PrivateKeyDer<'static>)> {
let tls_key = RcgenKeyPair::generate_for(&rcgen::PKCS_ECDSA_P256_SHA256)?;
let mut cert_params = CertificateParams::new(vec![])?;
cert_params.distinguished_name.push(
rcgen::DnType::CommonName,
rcgen::DnValue::Utf8String(keypair.pubkey().to_string()),
);
let cert = cert_params.self_signed(&tls_key)?;
let cert_der = CertificateDer::from(cert.der().to_vec());
let key_der = PrivateKeyDer::Pkcs8(PrivatePkcs8KeyDer::from(tls_key.serialize_der()));
Ok((cert_der, key_der))
}
@@ -101,8 +103,13 @@ pub struct SoyasClient {
impl SoyasClient {
pub async fn new(rpc_url: String, endpoint_string: String, api_key: String) -> Result<Self> {
let rpc_client = SolanaRpcClient::new(rpc_url);
let keypair = Keypair::from_base58_string(&api_key);
let (cert, key) = generate_self_signed_cert(&keypair)?;
let keypair = Keypair::try_from_base58_string(api_key.trim()).map_err(|e| {
anyhow::anyhow!(
"Soyas api_token 无法解析为 Solana keypair base58QUIC mTLS 用): {}",
e
)
})?;
let (cert, key) = generate_client_tls_credentials(&keypair)?;
let mut crypto = rustls::ClientConfig::builder()
.dangerous()
.with_custom_certificate_verifier(SkipServerVerification::new())
+61 -112
View File
@@ -6,7 +6,9 @@ use quinn::{
TransportConfig,
};
use rand::seq::IndexedRandom as _;
use solana_sdk::signer::Signer;
use solana_sdk::{signature::Keypair, transaction::VersionedTransaction};
use solana_tls_utils::{new_dummy_x509_certificate, SkipServerVerification};
use std::time::Instant;
use std::{
net::{SocketAddr, ToSocketAddrs as _},
@@ -25,69 +27,9 @@ use crate::{
swqos::{SwqosType, TradeType},
};
// Skip server verification implementation
#[derive(Debug)]
struct SkipServerVerification;
impl SkipServerVerification {
fn new() -> Arc<Self> {
Arc::new(Self)
}
}
impl rustls::client::danger::ServerCertVerifier for SkipServerVerification {
fn verify_server_cert(
&self,
_end_entity: &rustls::pki_types::CertificateDer<'_>,
_intermediates: &[rustls::pki_types::CertificateDer<'_>],
_server_name: &rustls::pki_types::ServerName<'_>,
_ocsp_response: &[u8],
_now: rustls::pki_types::UnixTime,
) -> Result<rustls::client::danger::ServerCertVerified, rustls::Error> {
Ok(rustls::client::danger::ServerCertVerified::assertion())
}
fn verify_tls12_signature(
&self,
_message: &[u8],
_cert: &rustls::pki_types::CertificateDer<'_>,
_dss: &rustls::DigitallySignedStruct,
) -> Result<rustls::client::danger::HandshakeSignatureValid, rustls::Error> {
Ok(rustls::client::danger::HandshakeSignatureValid::assertion())
}
fn verify_tls13_signature(
&self,
_message: &[u8],
_cert: &rustls::pki_types::CertificateDer<'_>,
_dss: &rustls::DigitallySignedStruct,
) -> Result<rustls::client::danger::HandshakeSignatureValid, rustls::Error> {
Ok(rustls::client::danger::HandshakeSignatureValid::assertion())
}
fn supported_verify_schemes(&self) -> Vec<rustls::SignatureScheme> {
vec![rustls::SignatureScheme::ECDSA_NISTP256_SHA256]
}
}
// Generate dummy self-signed certificate using rcgen
fn generate_self_signed_cert(_keypair: &Keypair) -> Result<(rustls::pki_types::CertificateDer<'static>, rustls::pki_types::PrivateKeyDer<'static>)> {
// Generate a new key pair for the certificate
let key_pair = rcgen::KeyPair::generate()?;
let params = rcgen::CertificateParams::default();
let cert = params.self_signed(&key_pair)?;
let cert_der = rustls::pki_types::CertificateDer::from(cert.der().to_vec());
let key_der = rustls::pki_types::PrivateKeyDer::try_from(key_pair.serialize_der())
.map_err(|e| anyhow::anyhow!("Failed to create private key: {:?}", e))?;
Ok((cert_der, key_der))
}
const ALPN_TPU_PROTOCOL_ID: &[u8] = b"solana-tpu";
/// Fallback SNI when endpoint is IP or cannot extract host (keeps legacy behavior).
const SPEED_SERVER_FALLBACK: &str = "speed-landing";
/// QUIC TLS SNI:与 Speedlanding 官方客户端一致,固定为 `speed-landing`(勿用 PoP 主机名,否则易握手失败)。
const SPEED_SERVER: &str = "speed-landing";
const KEEP_ALIVE_INTERVAL: Duration = Duration::from_secs(25);
const MAX_IDLE_TIMEOUT: Duration = Duration::from_secs(5 * 60);
const CONNECT_TIMEOUT: Duration = Duration::from_secs(5);
@@ -98,34 +40,21 @@ pub struct SpeedlandingClient {
endpoint: Endpoint,
client_config: ClientConfig,
addr: SocketAddr,
/// TLS SNI: host from endpoint URL so server presents the right cert (e.g. nyc.speedlanding.trade).
server_name: String,
connection: ArcSwap<Connection>,
reconnect: Mutex<()>,
}
impl SpeedlandingClient {
/// Extract TLS SNI (host) from endpoint URL. Uses fallback "speed-landing" for IP or when host cannot be determined.
fn server_name_from_endpoint(endpoint: &str) -> String {
let without_scheme = endpoint
.strip_prefix("https://")
.or_else(|| endpoint.strip_prefix("http://"))
.unwrap_or(endpoint);
let host = without_scheme.split(':').next().unwrap_or("").trim();
if host.is_empty() {
return SPEED_SERVER_FALLBACK.to_string();
}
if !host.chars().any(|c| c.is_ascii_alphabetic()) {
return SPEED_SERVER_FALLBACK.to_string();
}
host.to_string()
}
pub async fn new(rpc_url: String, endpoint_string: String, api_key: String) -> Result<Self> {
let rpc_client = SolanaRpcClient::new(rpc_url);
let server_name = Self::server_name_from_endpoint(&endpoint_string);
let keypair = Keypair::from_base58_string(&api_key);
let (cert, key) = generate_self_signed_cert(&keypair)?;
// Speedlanding QUIC:与官方一致使用 `solana_tls_utils::new_dummy_x509_certificate`Ed25519 dummy cert+ SNI `speed-landing`。
let keypair = Keypair::try_from_base58_string(api_key.trim()).map_err(|e| {
anyhow::anyhow!(
"Speedlanding api_token 无法解析为 Solana keypair base58(用于 mTLS);请确认粘贴的是机器人提供的密钥而非其它字符串: {}",
e
)
})?;
let (cert, key) = new_dummy_x509_certificate(&keypair);
let mut crypto = rustls::ClientConfig::builder()
.dangerous()
.with_custom_certificate_verifier(SkipServerVerification::new())
@@ -148,18 +77,23 @@ impl SpeedlandingClient {
.to_socket_addrs()?
.next()
.ok_or_else(|| anyhow::anyhow!("Address not resolved"))?;
let connecting = endpoint.connect(addr, &server_name)?;
let connecting = endpoint.connect(addr, SPEED_SERVER)?;
let connection = timeout(CONNECT_TIMEOUT, connecting)
.await
.context("Speedlanding QUIC connect timeout")?
.context("Speedlanding QUIC handshake failed")?;
.with_context(|| {
format!(
"Speedlanding QUIC handshake failed(请确认:1) 机器人登记的身份与钱包公钥 {} 一致 2) 本机 UDP 可访问 {} 3) region 与 PoP 匹配)",
keypair.pubkey(),
endpoint_string
)
})?;
Ok(Self {
rpc_client: Arc::new(rpc_client),
endpoint,
client_config,
addr,
server_name,
connection: ArcSwap::from_pointee(connection),
reconnect: Mutex::new(()),
})
@@ -181,12 +115,17 @@ impl SpeedlandingClient {
let connecting = self.endpoint.connect_with(
self.client_config.clone(),
self.addr,
self.server_name.as_str(),
SPEED_SERVER,
)?;
let connection = timeout(CONNECT_TIMEOUT, connecting)
.await
.context("Speedlanding QUIC reconnect timeout")?
.context("Speedlanding QUIC re-handshake failed")?;
.with_context(|| {
format!(
"Speedlanding QUIC re-handshake failed(对端 {} SNI {}",
self.addr, SPEED_SERVER
)
})?;
self.connection.store(Arc::new(connection));
return Ok(self.connection.load_full());
}
@@ -219,51 +158,61 @@ impl SwqosClientTrait for SpeedlandingClient {
Ok(Err(_)) | Err(_) => true,
};
if need_retry {
if crate::common::sdk_log::sdk_log_enabled() {
eprintln!(" [Speedlanding] {} submission failed after {:?}, reconnecting", trade_type, start_time.elapsed());
}
eprintln!(
" [Speedlanding] {} QUIC 首次发送失败 {:?},正在重试",
trade_type,
start_time.elapsed()
);
let connection = self.ensure_connected().await?;
send_result =
timeout(SEND_TIMEOUT, Self::try_send_bytes(&connection, &*buf_guard)).await;
}
match send_result.context("Speedlanding QUIC send timeout") {
Ok(Ok(())) => {
if crate::common::sdk_log::sdk_log_enabled() {
crate::common::sdk_log::log_swqos_submitted("Speedlanding", trade_type, start_time.elapsed());
}
// 提交结果与「详细耗时/SDK 开关」无关,便于确认当前通道确实在执行
crate::common::sdk_log::log_swqos_submitted("Speedlanding", trade_type, start_time.elapsed());
}
Ok(Err(e)) => {
if crate::common::sdk_log::sdk_log_enabled() {
crate::common::sdk_log::log_swqos_submission_failed("Speedlanding", trade_type, start_time.elapsed(), &e);
}
crate::common::sdk_log::log_swqos_submission_failed(
"Speedlanding",
trade_type,
start_time.elapsed(),
&e,
);
return Err(e.into());
}
Err(e) => {
if crate::common::sdk_log::sdk_log_enabled() {
crate::common::sdk_log::log_swqos_submission_failed("Speedlanding", trade_type, start_time.elapsed(), "timeout");
}
crate::common::sdk_log::log_swqos_submission_failed(
"Speedlanding",
trade_type,
start_time.elapsed(),
"timeout",
);
return Err(e.into());
}
}
match poll_transaction_confirmation(&self.rpc_client, signature, wait_confirmation).await {
Ok(_) => (),
Err(e) => {
if crate::common::sdk_log::sdk_log_enabled() {
println!(" signature: {:?}", signature);
println!(
" [{:width$}] {} confirmation failed: {:?}",
"Speedlanding",
trade_type,
start_time.elapsed(),
width = crate::common::sdk_log::SWQOS_LABEL_WIDTH
);
}
println!(" signature: {:?}", signature);
crate::common::sdk_log::log_swqos_submission_failed(
"Speedlanding",
trade_type,
start_time.elapsed(),
&e,
);
return Err(e);
}
}
if wait_confirmation && crate::common::sdk_log::sdk_log_enabled() {
if wait_confirmation {
println!(" signature: {:?}", signature);
println!(" [{:width$}] {} confirmed: {:?}", "Speedlanding", trade_type, start_time.elapsed(), width = crate::common::sdk_log::SWQOS_LABEL_WIDTH);
println!(
" [{:width$}] {} confirmed: {:?}",
"Speedlanding",
trade_type,
start_time.elapsed(),
width = crate::common::sdk_log::SWQOS_LABEL_WIDTH
);
}
Ok(())
}