diff --git a/Cargo.toml b/Cargo.toml index 4419359..72b9ae6 100755 --- a/Cargo.toml +++ b/Cargo.toml @@ -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 found,QUIC(Speedlanding/Soyas)无法初始化 +quinn = { version = "0.11", default-features = false, features = ["rustls", "runtime-tokio"] } rcgen = "0.13" uuid = "1.11" diff --git a/src/lib.rs b/src/lib.rs index a41ad8a..b6a141b 100755 --- a/src/lib.rs +++ b/src/lib.rs @@ -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 checking(QUIC 握手可能较慢,单节点超时 15s) + const SWQOS_CLIENT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(15); let mut swqos_clients: Vec> = 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); diff --git a/src/swqos/solana_rpc.rs b/src/swqos/solana_rpc.rs index 64cc92b..e2670fb 100755 --- a/src/swqos/solana_rpc.rs +++ b/src/swqos/solana_rpc.rs @@ -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(_) => (), diff --git a/src/swqos/soyas.rs b/src/swqos/soyas.rs index bb0eac2..b46e71e 100644 --- a/src/swqos/soyas.rs +++ b/src/swqos/soyas.rs @@ -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 { 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 base58(QUIC mTLS 用): {}", + e + ) + })?; + let (cert, key) = generate_client_tls_credentials(&keypair)?; let mut crypto = rustls::ClientConfig::builder() .dangerous() .with_custom_certificate_verifier(SkipServerVerification::new()) diff --git a/src/swqos/speedlanding.rs b/src/swqos/speedlanding.rs index 8ffc39a..c936093 100644 --- a/src/swqos/speedlanding.rs +++ b/src/swqos/speedlanding.rs @@ -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 { - 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 { - Ok(rustls::client::danger::ServerCertVerified::assertion()) - } - - fn verify_tls12_signature( - &self, - _message: &[u8], - _cert: &rustls::pki_types::CertificateDer<'_>, - _dss: &rustls::DigitallySignedStruct, - ) -> Result { - Ok(rustls::client::danger::HandshakeSignatureValid::assertion()) - } - - fn verify_tls13_signature( - &self, - _message: &[u8], - _cert: &rustls::pki_types::CertificateDer<'_>, - _dss: &rustls::DigitallySignedStruct, - ) -> Result { - Ok(rustls::client::danger::HandshakeSignatureValid::assertion()) - } - - fn supported_verify_schemes(&self) -> Vec { - 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, 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 { 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(()) }