diff --git a/src/main.rs b/src/main.rs index 1f6fabc..b95723c 100755 --- a/src/main.rs +++ b/src/main.rs @@ -86,13 +86,17 @@ async fn test_grpc() -> Result<(), Box> { println!("Protocols to monitor: {:?}", protocols); + const SYSTEM_PROGRAM_ID: solana_sdk::pubkey::Pubkey = + solana_sdk::pubkey!("11111111111111111111111111111111"); + // Filter accounts let account_include = vec![ - PUMPFUN_PROGRAM_ID.to_string(), // Listen to pumpfun program ID - PUMPSWAP_PROGRAM_ID.to_string(), // Listen to pumpswap program ID - BONK_PROGRAM_ID.to_string(), // Listen to bonk program ID - RAYDIUM_CPMM_PROGRAM_ID.to_string(), // Listen to raydium_cpmm program ID - RAYDIUM_CLMM_PROGRAM_ID.to_string(), // Listen to raydium_clmm program ID + SYSTEM_PROGRAM_ID.to_string(), + PUMPFUN_PROGRAM_ID.to_string(), // Listen to pumpfun program ID + PUMPSWAP_PROGRAM_ID.to_string(), // Listen to pumpswap program ID + BONK_PROGRAM_ID.to_string(), // Listen to bonk program ID + RAYDIUM_CPMM_PROGRAM_ID.to_string(), // Listen to raydium_cpmm program ID + RAYDIUM_CLMM_PROGRAM_ID.to_string(), // Listen to raydium_clmm program ID RAYDIUM_AMM_V4_PROGRAM_ID.to_string(), // Listen to raydium_amm_v4 program ID ]; let account_exclude = vec![]; @@ -151,7 +155,7 @@ async fn test_shreds() -> Result<(), Box> { // Enable performance monitoring, has performance overhead, disabled by default config.enable_metrics = true; let shred_stream = - ShredStreamGrpc::new_with_config("http://127.0.0.1:10800".to_string(), config).await?; + ShredStreamGrpc::new_with_config("http://64.130.37.195:10800".to_string(), config).await?; let callback = create_event_callback(); let protocols = vec![ @@ -188,17 +192,11 @@ async fn test_shreds() -> Result<(), Box> { fn create_event_callback() -> impl Fn(Box) { |event: Box| { - println!( - "🎉 Event received! Type: {:?}, ID: {}, Slot: {}, Transaction Index: {:?}", - event.event_type(), - event.id(), - event.slot(), - event.transaction_index(), - ); + println!("🎉 Event received! Type: {:?}, ID: {}", event.event_type(), event.id()); match_event!(event, { // -------------------------- block meta ----------------------- BlockMetaEvent => |e: BlockMetaEvent| { - println!("BlockMetaEvent: {e:?}"); + println!("BlockMetaEvent: {:?}", e.metadata.program_handle_time_consuming_us); }, // -------------------------- bonk ----------------------- BonkPoolCreateEvent => |e: BonkPoolCreateEvent| { diff --git a/src/streaming/common/constants.rs b/src/streaming/common/constants.rs index 8cc2be8..0b253f6 100644 --- a/src/streaming/common/constants.rs +++ b/src/streaming/common/constants.rs @@ -11,4 +11,4 @@ pub const DEFAULT_BATCH_TIMEOUT_MS: u64 = 5; // 性能监控相关常量 pub const DEFAULT_METRICS_WINDOW_SECONDS: u64 = 5; pub const DEFAULT_METRICS_PRINT_INTERVAL_SECONDS: u64 = 10; -pub const SLOW_PROCESSING_THRESHOLD_US: f64 = 500.0; +pub const SLOW_PROCESSING_THRESHOLD_US: f64 = 3000.0; diff --git a/src/streaming/common/event_processor.rs b/src/streaming/common/event_processor.rs index 7bbacd2..a2364fe 100644 --- a/src/streaming/common/event_processor.rs +++ b/src/streaming/common/event_processor.rs @@ -1,7 +1,7 @@ use std::sync::Arc; -use tokio::task::JoinHandle; use solana_sdk::pubkey::Pubkey; +use solana_sdk::signature::Signature; use crate::common::AnyResult; use crate::streaming::common::{ @@ -14,7 +14,7 @@ use crate::streaming::event_parser::EventParser; use crate::streaming::event_parser::{ core::traits::UnifiedEvent, protocols::mutil::parser::MutilEventParser, Protocol, }; -use crate::streaming::grpc::{BackpressureStrategy, BatchConfig, EventPretty}; +use crate::streaming::grpc::{BackpressureConfig, BatchConfig, EventPretty}; use crate::streaming::shred::TransactionWithSlot; use once_cell::sync::OnceCell; @@ -25,21 +25,24 @@ pub struct EventProcessor { pub(crate) parser_cache: OnceCell>, pub(crate) protocols: Vec, pub(crate) event_type_filter: Option, - pub(crate) backpressure_strategy: BackpressureStrategy, + pub(crate) callback: Option) + Send + Sync>>, + pub(crate) backpressure_config: BackpressureConfig, pub(crate) batch_config: BatchConfig, } impl EventProcessor { /// 创建新的事件处理器 pub fn new(metrics_manager: MetricsManager, config: ClientConfig) -> Self { + let backpressure_config = config.backpressure.clone(); Self { metrics_manager, config, parser_cache: OnceCell::new(), protocols: vec![], event_type_filter: None, - backpressure_strategy: BackpressureStrategy::Block, + backpressure_config, batch_config: BatchConfig::default(), + callback: None, } } @@ -47,13 +50,15 @@ impl EventProcessor { &mut self, protocols: Vec, event_type_filter: Option, - backpressure_strategy: BackpressureStrategy, + backpressure_config: BackpressureConfig, batch_config: BatchConfig, + callback: Option) + Send + Sync>>, ) { self.protocols = protocols.clone(); self.event_type_filter = event_type_filter.clone(); - self.backpressure_strategy = backpressure_strategy; + self.backpressure_config = backpressure_config; self.batch_config = batch_config; + self.callback = callback; self.parser_cache .get_or_init(|| Arc::new(MutilEventParser::new(protocols, event_type_filter))); } @@ -62,35 +67,24 @@ impl EventProcessor { self.parser_cache.get().unwrap().clone() } - pub fn get_event_handle(&self) -> Option> { - return None; - } - - pub async fn process_grpc_event_transaction_with_metrics( + pub async fn process_grpc_event_transaction_with_metrics( &self, event_pretty: EventPretty, - callback: &F, bot_wallet: Option, - ) -> AnyResult<()> - where - F: Fn(Box) + Send + Sync, - { - self.process_grpc_event_transaction(event_pretty, callback, bot_wallet).await?; + ) -> AnyResult<()> { + self.process_grpc_event_transaction(event_pretty, bot_wallet).await?; Ok(()) } - async fn process_grpc_event_transaction( + async fn process_grpc_event_transaction( &self, event_pretty: EventPretty, - callback: &F, bot_wallet: Option, - ) -> AnyResult<()> - where - F: Fn(Box) + Send + Sync, - { + ) -> AnyResult<()> { match event_pretty { EventPretty::Account(account_pretty) => { self.metrics_manager.add_account_process_count(); + let signature = account_pretty.signature; let account_event = AccountEventParser::parse_account_event( self.protocols.clone(), account_pretty, @@ -98,12 +92,12 @@ impl EventProcessor { ); if let Some(event) = account_event { let processing_time_us = event.program_handle_time_consuming_us() as f64; - callback(event); - // 更新性能指标(如果启用) - self.metrics_manager.update_metrics( + self.invoke_callback(event); + self.update_metrics( MetricsEventType::Account, 1, processing_time_us, + Some(signature), ); } } @@ -113,7 +107,7 @@ impl EventProcessor { let signature = transaction_pretty.signature; // 使用缓存获取解析器 let parser = self.get_parser(); - let mut all_events = parser + let all_events = parser .parse_transaction( transaction_pretty.tx.clone(), signature, @@ -125,37 +119,26 @@ impl EventProcessor { .await .unwrap_or_else(|_e| vec![]); - // 为所有事件设置交易索引 - for event in &mut all_events { - event.set_transaction_index(transaction_pretty.transaction_index); - } - - let max_time_consuming_us = all_events - .iter() - .map(|event| event.program_handle_time_consuming_us()) - .max() - .unwrap_or(0); - - // 保存事件数量用于日志记录 + let mut max_time_consuming_us = 0; let event_count = all_events.len(); - // 批量处理事件 - if !all_events.is_empty() { - for mut event in all_events { - event.set_program_handle_time_consuming_us( - chrono::Utc::now().timestamp_micros() - - event.program_received_time_us(), - ); - callback(event); - } + // 为所有事件设置交易索引 + for mut event in all_events { + event.set_transaction_index(transaction_pretty.transaction_index); + event.set_program_handle_time_consuming_us( + chrono::Utc::now().timestamp_micros() - event.program_received_time_us(), + ); + max_time_consuming_us = + max_time_consuming_us.max(event.program_handle_time_consuming_us()); + self.invoke_callback(event); } // 更新性能指标 - // 更新性能指标(如果启用) - self.metrics_manager.update_metrics( - MetricsEventType::Tx, + self.update_metrics( + MetricsEventType::Transaction, event_count as u64, max_time_consuming_us as f64, + Some(signature), ); } EventPretty::BlockMeta(block_meta_pretty) => { @@ -171,29 +154,34 @@ impl EventProcessor { block_meta_pretty.program_received_time_us, ); let processing_time_us = block_meta_event.program_handle_time_consuming_us() as f64; - callback(block_meta_event); - // 更新性能指标(如果启用) - self.metrics_manager.update_metrics( - MetricsEventType::BlockMeta, - 1, - processing_time_us, - ); + self.invoke_callback(block_meta_event); + self.update_metrics(MetricsEventType::BlockMeta, 1, processing_time_us, None); } } Ok(()) } + pub fn invoke_callback(&self, event: Box) { + if let Some(callback) = self.callback.as_ref() { + callback(event); + } + } + /// 即时处理单个交易 - pub async fn process_shred_transaction_immediate( + pub async fn process_shred_transaction_immediate( &self, transaction_with_slot: TransactionWithSlot, bot_wallet: Option, - callback: &F, - ) -> AnyResult<()> - where - F: Fn(Box) + Send + Sync, - { + ) -> AnyResult<()> { + self.process_shred_transaction(transaction_with_slot, bot_wallet).await + } + + pub async fn process_shred_transaction( + &self, + transaction_with_slot: TransactionWithSlot, + bot_wallet: Option, + ) -> AnyResult<()> { self.metrics_manager.add_tx_process_count(); let program_received_time_us = chrono::Utc::now().timestamp_micros(); let slot = transaction_with_slot.slot; @@ -215,11 +203,7 @@ impl EventProcessor { .await .unwrap_or_else(|_e| vec![]); - let max_time_consuming_us = all_events - .iter() - .map(|event| event.program_handle_time_consuming_us()) - .max() - .unwrap_or(0); + let mut max_time_consuming_us = 0; // 保存事件数量用于日志记录 let event_count = all_events.len(); @@ -229,18 +213,31 @@ impl EventProcessor { event.set_program_handle_time_consuming_us( chrono::Utc::now().timestamp_micros() - event.program_received_time_us(), ); - callback(event); + max_time_consuming_us = + max_time_consuming_us.max(event.program_handle_time_consuming_us()); + self.invoke_callback(event); } // 实际调用性能指标更新 - self.metrics_manager.update_metrics( - MetricsEventType::Tx, + self.update_metrics( + MetricsEventType::Transaction, event_count as u64, max_time_consuming_us as f64, + Some(signature), ); Ok(()) } + + fn update_metrics( + &self, + ty: MetricsEventType, + count: u64, + time_us: f64, + signature: Option, + ) { + self.metrics_manager.update_metrics(ty, count, time_us, signature); + } } // 实现 Clone trait 以支持模块间共享 @@ -252,8 +249,9 @@ impl Clone for EventProcessor { parser_cache: self.parser_cache.clone(), protocols: self.protocols.clone(), event_type_filter: self.event_type_filter.clone(), - backpressure_strategy: self.backpressure_strategy.clone(), + backpressure_config: self.backpressure_config.clone(), batch_config: self.batch_config.clone(), + callback: self.callback.clone(), } } } diff --git a/src/streaming/common/metrics.rs b/src/streaming/common/metrics.rs index 6dd80e5..a5df787 100644 --- a/src/streaming/common/metrics.rs +++ b/src/streaming/common/metrics.rs @@ -1,344 +1,554 @@ +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::Arc; -use crossbeam::utils::Backoff; -use crossbeam::atomic::AtomicCell; -use std::sync::RwLock; -use super::config::StreamClientConfig; +use solana_sdk::signature::Signature; + use super::constants::*; -/// 单个事件类型的指标 +/// 事件类型枚举 +#[derive(Debug, Clone, Copy)] +pub enum EventType { + Transaction = 0, + Account = 1, + BlockMeta = 2, +} + +/// 兼容性别名 +pub type MetricsEventType = EventType; + +impl EventType { + #[inline] + const fn as_index(self) -> usize { + self as usize + } + + const fn name(self) -> &'static str { + match self { + EventType::Transaction => "TX", + EventType::Account => "Account", + EventType::BlockMeta => "Block Meta", + } + } + + // 兼容性常量 + pub const TX: EventType = EventType::Transaction; +} + +/// 高性能原子事件指标 +#[derive(Debug)] +struct AtomicEventMetrics { + process_count: AtomicU64, + events_processed: AtomicU64, + events_in_window: AtomicU64, + window_start_nanos: AtomicU64, + events_per_second_bits: AtomicU64, // f64 的位表示 +} + +impl AtomicEventMetrics { + fn new(now_nanos: u64) -> Self { + Self { + process_count: AtomicU64::new(0), + events_processed: AtomicU64::new(0), + events_in_window: AtomicU64::new(0), + window_start_nanos: AtomicU64::new(now_nanos), + events_per_second_bits: AtomicU64::new(0), + } + } + + /// 原子地增加处理计数 + #[inline] + fn add_process_count(&self) { + self.process_count.fetch_add(1, Ordering::Relaxed); + } + + /// 原子地增加事件处理数量 + #[inline] + fn add_events_processed(&self, count: u64) { + self.events_processed.fetch_add(count, Ordering::Relaxed); + self.events_in_window.fetch_add(count, Ordering::Relaxed); + } + + /// 获取当前计数(非阻塞) + #[inline] + fn get_counts(&self) -> (u64, u64, u64) { + ( + self.process_count.load(Ordering::Relaxed), + self.events_processed.load(Ordering::Relaxed), + self.events_in_window.load(Ordering::Relaxed), + ) + } + + /// 原子地更新每秒事件数 + #[inline] + fn update_events_per_second(&self, eps: f64) { + self.events_per_second_bits.store(eps.to_bits(), Ordering::Relaxed); + } + + /// 获取每秒事件数 + #[inline] + fn get_events_per_second(&self) -> f64 { + f64::from_bits(self.events_per_second_bits.load(Ordering::Relaxed)) + } + + /// 重置窗口计数 + #[inline] + fn reset_window(&self, new_start_nanos: u64) { + self.events_in_window.store(0, Ordering::Relaxed); + self.window_start_nanos.store(new_start_nanos, Ordering::Relaxed); + } + + #[inline] + fn get_window_start(&self) -> u64 { + self.window_start_nanos.load(Ordering::Relaxed) + } +} + +/// 高性能原子处理时间统计 +#[derive(Debug)] +struct AtomicProcessingTimeStats { + min_time_bits: AtomicU64, + max_time_bits: AtomicU64, + total_time_us: AtomicU64, // 存储微秒的整数部分 + total_events: AtomicU64, +} + +impl AtomicProcessingTimeStats { + fn new() -> Self { + Self { + min_time_bits: AtomicU64::new(f64::INFINITY.to_bits()), + max_time_bits: AtomicU64::new(0), + total_time_us: AtomicU64::new(0), + total_events: AtomicU64::new(0), + } + } + + /// 原子地更新处理时间统计 + #[inline] + fn update(&self, time_us: f64, event_count: u64) { + let time_bits = time_us.to_bits(); + + // 更新最小值(使用 compare_exchange_weak 循环) + let mut current_min = self.min_time_bits.load(Ordering::Relaxed); + while time_bits < current_min { + match self.min_time_bits.compare_exchange_weak( + current_min, + time_bits, + Ordering::Relaxed, + Ordering::Relaxed, + ) { + Ok(_) => break, + Err(x) => current_min = x, + } + } + + // 更新最大值 + let mut current_max = self.max_time_bits.load(Ordering::Relaxed); + while time_bits > current_max { + match self.max_time_bits.compare_exchange_weak( + current_max, + time_bits, + Ordering::Relaxed, + Ordering::Relaxed, + ) { + Ok(_) => break, + Err(x) => current_max = x, + } + } + + // 更新累计值(将微秒转换为整数避免浮点累加问题) + let total_time_us_int = (time_us * event_count as f64) as u64; + self.total_time_us.fetch_add(total_time_us_int, Ordering::Relaxed); + self.total_events.fetch_add(event_count, Ordering::Relaxed); + } + + /// 获取统计值(非阻塞) + #[inline] + fn get_stats(&self) -> ProcessingTimeStats { + let min_bits = self.min_time_bits.load(Ordering::Relaxed); + let max_bits = self.max_time_bits.load(Ordering::Relaxed); + let total_time_us_int = self.total_time_us.load(Ordering::Relaxed); + let total_events = self.total_events.load(Ordering::Relaxed); + + let min_time = f64::from_bits(min_bits); + let max_time = f64::from_bits(max_bits); + let avg_time = + if total_events > 0 { total_time_us_int as f64 / total_events as f64 } else { 0.0 }; + + ProcessingTimeStats { + min_us: if min_time == f64::INFINITY { 0.0 } else { min_time }, + max_us: max_time, + avg_us: avg_time, + } + } +} + +/// 处理时间统计结果 #[derive(Debug, Clone)] -pub struct EventMetrics { +pub struct ProcessingTimeStats { + pub min_us: f64, + pub max_us: f64, + pub avg_us: f64, +} + +/// 事件指标快照 +#[derive(Debug, Clone)] +pub struct EventMetricsSnapshot { pub process_count: u64, pub events_processed: u64, pub events_per_second: f64, - pub events_in_window: u64, - pub window_start_time: std::time::Instant, } -impl EventMetrics { - fn new(now: std::time::Instant) -> Self { - Self { - process_count: 0, - events_processed: 0, - events_per_second: 0.0, - events_in_window: 0, - window_start_time: now, - } - } -} - -/// 通用性能监控指标 +/// 兼容性结构 - 完整的性能指标 #[derive(Debug, Clone)] pub struct PerformanceMetrics { - pub start_time: std::time::Instant, - pub event_metrics: [EventMetrics; 3], // [Tx, Account, BlockMeta] - pub average_processing_time_us: f64, - pub min_processing_time_us: f64, - pub max_processing_time_us: f64, - pub last_update_time: std::time::Instant, -} - -impl Default for PerformanceMetrics { - fn default() -> Self { - Self::new() - } -} - -pub enum MetricsEventType { - Tx, - Account, - BlockMeta, -} - -impl MetricsEventType { - fn as_index(&self) -> usize { - match self { - MetricsEventType::Tx => 0, - MetricsEventType::Account => 1, - MetricsEventType::BlockMeta => 2, - } - } + pub uptime: std::time::Duration, + pub tx_metrics: EventMetricsSnapshot, + pub account_metrics: EventMetricsSnapshot, + pub block_meta_metrics: EventMetricsSnapshot, + pub processing_stats: ProcessingTimeStats, } impl PerformanceMetrics { + /// 创建默认的性能指标(兼容性方法) pub fn new() -> Self { - let now = std::time::Instant::now(); + let default_metrics = + EventMetricsSnapshot { process_count: 0, events_processed: 0, events_per_second: 0.0 }; + let default_stats = ProcessingTimeStats { min_us: 0.0, max_us: 0.0, avg_us: 0.0 }; + Self { - start_time: now, - event_metrics: [EventMetrics::new(now), EventMetrics::new(now), EventMetrics::new(now)], - average_processing_time_us: 0.0, - min_processing_time_us: 0.0, - max_processing_time_us: 0.0, - last_update_time: now, + uptime: std::time::Duration::ZERO, + tx_metrics: default_metrics.clone(), + account_metrics: default_metrics.clone(), + block_meta_metrics: default_metrics, + processing_stats: default_stats, + } + } +} + +/// 高性能指标系统 +#[derive(Debug)] +pub struct HighPerformanceMetrics { + start_nanos: u64, + event_metrics: [AtomicEventMetrics; 3], + processing_stats: AtomicProcessingTimeStats, +} + +impl HighPerformanceMetrics { + fn new() -> Self { + let now_nanos = + std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos() + as u64; + + Self { + start_nanos: now_nanos, + event_metrics: [ + AtomicEventMetrics::new(now_nanos), + AtomicEventMetrics::new(now_nanos), + AtomicEventMetrics::new(now_nanos), + ], + processing_stats: AtomicProcessingTimeStats::new(), } } - /// 更新时间窗口指标 - fn update_window_metrics( - &mut self, - event_type: &MetricsEventType, - now: std::time::Instant, - window_duration: std::time::Duration, - ) { + /// 获取运行时长(秒) + #[inline] + pub fn get_uptime_seconds(&self) -> f64 { + let now_nanos = + std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos() + as u64; + (now_nanos - self.start_nanos) as f64 / 1_000_000_000.0 + } + + /// 获取事件指标快照 + #[inline] + pub fn get_event_metrics(&self, event_type: EventType) -> EventMetricsSnapshot { let index = event_type.as_index(); - let event_metric = &mut self.event_metrics[index]; + let (process_count, events_processed, _) = self.event_metrics[index].get_counts(); + let events_per_second = self.calculate_real_time_eps(event_type); - if now.duration_since(event_metric.window_start_time) >= window_duration { - let window_seconds = now.duration_since(event_metric.window_start_time).as_secs_f64(); - // 修复:正确计算每秒事件数,避免除零错误 - event_metric.events_per_second = if window_seconds > 0.001 { - // 避免极小的时间差 - event_metric.events_in_window as f64 / window_seconds - } else { - 0.0 // 时间太短时设为0,而不是事件总数 - }; - - // 重置窗口 - event_metric.events_in_window = 0; - event_metric.window_start_time = now; - } + EventMetricsSnapshot { process_count, events_processed, events_per_second } } - /// 计算实时每秒事件数(用于显示) - fn calculate_real_time_events_per_second( - &self, - event_type: &MetricsEventType, - now: std::time::Instant, - ) -> f64 { + /// 获取处理时间统计 + #[inline] + pub fn get_processing_stats(&self) -> ProcessingTimeStats { + self.processing_stats.get_stats() + } + + /// 计算实时每秒事件数(非阻塞) + fn calculate_real_time_eps(&self, event_type: EventType) -> f64 { + let now_nanos = + std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos() + as u64; + let index = event_type.as_index(); let event_metric = &self.event_metrics[index]; - let current_window_duration = - now.duration_since(event_metric.window_start_time).as_secs_f64(); + let window_start = event_metric.get_window_start(); + let current_window_duration_secs = + (now_nanos.saturating_sub(window_start)) as f64 / 1_000_000_000.0; + let events_in_window = event_metric.events_in_window.load(Ordering::Relaxed); - // 如果当前窗口有足够的时间和事件,使用当前窗口的数据 - if current_window_duration > 1.0 && event_metric.events_in_window > 0 { - event_metric.events_in_window as f64 / current_window_duration + // 优先级1: 当前窗口实时数据(≥2秒且有事件) + if current_window_duration_secs >= 2.0 && events_in_window > 0 { + return events_in_window as f64 / current_window_duration_secs; } - // 如果当前窗口时间太短或没有事件,使用上一个完整窗口的值 - else if event_metric.events_per_second > 0.0 { - event_metric.events_per_second + + // 优先级2: 上一个窗口的结果 + let stored_eps = event_metric.get_events_per_second(); + if stored_eps > 0.0 { + return stored_eps; } - // 如果都没有,计算总体平均值 - else { - let total_duration = now.duration_since(self.start_time).as_secs_f64(); - if total_duration > 1.0 && event_metric.events_processed > 0 { - event_metric.events_processed as f64 / total_duration - } else { - 0.0 + + // 优先级3: 总体平均值(≥3秒运行时间) + let total_duration_secs = self.get_uptime_seconds(); + let total_events = event_metric.events_processed.load(Ordering::Relaxed); + if total_duration_secs >= 3.0 && total_events > 0 { + return total_events as f64 / total_duration_secs; + } + + 0.0 + } + + /// 更新窗口指标(后台任务调用) + fn update_window_metrics(&self, event_type: EventType, window_duration_nanos: u64) { + let now_nanos = + std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos() + as u64; + + let index = event_type.as_index(); + let event_metric = &self.event_metrics[index]; + + let window_start = event_metric.get_window_start(); + if now_nanos.saturating_sub(window_start) >= window_duration_nanos { + let events_in_window = event_metric.events_in_window.load(Ordering::Relaxed); + let window_duration_secs = window_duration_nanos as f64 / 1_000_000_000.0; + + if window_duration_secs > 0.001 && events_in_window > 0 { + let eps = events_in_window as f64 / window_duration_secs; + event_metric.update_events_per_second(eps); } + + event_metric.reset_window(now_nanos); } } } -/// 通用性能监控管理器 +/// 高性能指标管理器 pub struct MetricsManager { - metrics: Arc>, - config: Arc, + metrics: Arc, + enable_metrics: bool, stream_name: String, + background_task_running: AtomicBool, } impl MetricsManager { - /// 创建新的性能监控管理器 - pub fn new( - metrics: Arc>, - config: Arc, - stream_name: String, - ) -> Self { - Self { metrics, config, stream_name } + /// 创建新的指标管理器 + pub fn new(enable_metrics: bool, stream_name: String) -> Self { + let manager = Self { + metrics: Arc::new(HighPerformanceMetrics::new()), + enable_metrics, + stream_name, + background_task_running: AtomicBool::new(false), + }; + + // 启动后台任务 + manager.start_background_tasks(); + manager } - /// 获取性能指标 - pub fn get_metrics(&self) -> PerformanceMetrics { - // 使用 Backoff 策略进行读取尝试 - let backoff = Backoff::new(); - loop { - match self.metrics.read() { - Ok(metrics) => return metrics.clone(), - Err(_) => { - // 如果获取读锁失败,使用指数退避策略 - backoff.snooze(); - continue; - } + /// 启动后台任务 + fn start_background_tasks(&self) { + if self + .background_task_running + .compare_exchange(false, true, Ordering::Relaxed, Ordering::Relaxed) + .is_ok() + { + if !self.enable_metrics { + return; } + + let metrics = self.metrics.clone(); + + tokio::spawn(async move { + let mut interval = tokio::time::interval(std::time::Duration::from_millis(500)); + + loop { + interval.tick().await; + + let window_duration_nanos = DEFAULT_METRICS_WINDOW_SECONDS * 1_000_000_000; + + // 更新所有事件类型的窗口指标 + metrics.update_window_metrics(EventType::Transaction, window_duration_nanos); + metrics.update_window_metrics(EventType::Account, window_duration_nanos); + metrics.update_window_metrics(EventType::BlockMeta, window_duration_nanos); + } + }); } } - /// 打印性能指标 - pub fn print_metrics(&self) { - let metrics = self.get_metrics(); - let event_names = ["TX", "Account", "Block Meta"]; - let event_types = - [MetricsEventType::Tx, MetricsEventType::Account, MetricsEventType::BlockMeta]; - let now = std::time::Instant::now(); + /// 记录处理次数(非阻塞) + #[inline] + pub fn record_process(&self, event_type: EventType) { + if self.enable_metrics { + self.metrics.event_metrics[event_type.as_index()].add_process_count(); + } + } + /// 记录事件处理(非阻塞) + #[inline] + pub fn record_events(&self, event_type: EventType, count: u64, processing_time_us: f64) { + if !self.enable_metrics { + return; + } + + // 原子更新事件计数 + self.metrics.event_metrics[event_type.as_index()].add_events_processed(count); + + // 原子更新处理时间统计 + self.metrics.processing_stats.update(processing_time_us, count); + } + + /// 记录慢处理操作 + #[inline] + pub fn log_slow_processing( + &self, + processing_time_us: f64, + event_count: usize, + signature: Option, + ) { + if processing_time_us > SLOW_PROCESSING_THRESHOLD_US { + log::warn!( + "{} slow processing: {:.2}us for {} events, signature: {:?}", + self.stream_name, + processing_time_us, + event_count, + signature + ); + } + } + + /// 获取运行时长 + pub fn get_uptime(&self) -> std::time::Duration { + std::time::Duration::from_secs_f64(self.metrics.get_uptime_seconds()) + } + + /// 获取事件指标 + pub fn get_event_metrics(&self, event_type: EventType) -> EventMetricsSnapshot { + self.metrics.get_event_metrics(event_type) + } + + /// 获取处理时间统计 + pub fn get_processing_stats(&self) -> ProcessingTimeStats { + self.metrics.get_processing_stats() + } + + /// 打印性能指标(非阻塞) + pub fn print_metrics(&self) { println!("\n📊 {} Performance Metrics", self.stream_name); - println!(" Run Time: {:?}", metrics.start_time.elapsed()); - - // 打印表格头部 + println!(" Run Time: {:?}", self.get_uptime()); + + // 打印事件指标表格 println!("┌─────────────┬──────────────┬──────────────────┬─────────────────┐"); println!("│ Event Type │ Process Count│ Events Processed │ Events/Second │"); println!("├─────────────┼──────────────┼──────────────────┼─────────────────┤"); - // 打印每种事件类型的数据 - for (i, name) in event_names.iter().enumerate() { - let event_metric = &metrics.event_metrics[i]; - // 使用实时计算的每秒事件数,而不是窗口更新的值 - let real_time_eps = metrics.calculate_real_time_events_per_second(&event_types[i], now); - + for event_type in [EventType::Transaction, EventType::Account, EventType::BlockMeta] { + let metrics = self.get_event_metrics(event_type); println!( "│ {:11} │ {:12} │ {:16} │ {:13.2} │", - name, - event_metric.process_count, - event_metric.events_processed, - real_time_eps + event_type.name(), + metrics.process_count, + metrics.events_processed, + metrics.events_per_second ); } println!("└─────────────┴──────────────┴──────────────────┴─────────────────┘"); // 打印处理时间统计表格 + let stats = self.get_processing_stats(); println!("\n⏱️ Processing Time Statistics"); println!("┌─────────────────────┬─────────────┐"); println!("│ Metric │ Value (us) │"); println!("├─────────────────────┼─────────────┤"); - println!("│ Average │ {:9.2} │", metrics.average_processing_time_us); - println!("│ Minimum │ {:9.2} │", metrics.min_processing_time_us); - println!("│ Maximum │ {:9.2} │", metrics.max_processing_time_us); + println!("│ Average │ {:9.2} │", stats.avg_us); + println!("│ Minimum │ {:9.2} │", stats.min_us); + println!("│ Maximum │ {:9.2} │", stats.max_us); println!("└─────────────────────┴─────────────┘"); println!(); } /// 启动自动性能监控任务 pub async fn start_auto_monitoring(&self) -> Option> { - // 检查是否启用性能监控 - if !self.config.enable_metrics { - return None; // 如果未启用性能监控,不启动监控任务 + if !self.enable_metrics { + return None; } - let metrics_manager = self.clone(); + let manager = self.clone(); let handle = tokio::spawn(async move { - let mut interval = tokio::time::interval(tokio::time::Duration::from_secs( + let mut interval = tokio::time::interval(std::time::Duration::from_secs( DEFAULT_METRICS_PRINT_INTERVAL_SECONDS, )); loop { interval.tick().await; - metrics_manager.print_metrics(); + manager.print_metrics(); } }); Some(handle) } - /// 更新处理次数 - pub fn add_process_count(&self, event_type: MetricsEventType) { - if !self.config.enable_metrics { - return; - } - - // 使用 Backoff 策略进行写入尝试 - let backoff = Backoff::new(); - loop { - match self.metrics.write() { - Ok(mut metrics) => { - metrics.event_metrics[event_type.as_index()].process_count += 1; - break; - }, - Err(_) => { - // 如果获取写锁失败,使用指数退避策略 - backoff.snooze(); - continue; - } - } + // === 兼容性方法 === + + /// 兼容性构造函数 + pub fn new_with_metrics( + _metrics: Arc>, + enable_metrics: bool, + stream_name: String, + ) -> Self { + Self::new(enable_metrics, stream_name) + } + + /// 获取完整的性能指标(兼容性方法) + pub fn get_metrics(&self) -> PerformanceMetrics { + PerformanceMetrics { + uptime: self.get_uptime(), + tx_metrics: self.get_event_metrics(EventType::Transaction), + account_metrics: self.get_event_metrics(EventType::Account), + block_meta_metrics: self.get_event_metrics(EventType::BlockMeta), + processing_stats: self.get_processing_stats(), } } - // 保持向后兼容的方法 + /// 兼容性方法 - 添加交易处理计数 + #[inline] pub fn add_tx_process_count(&self) { - self.add_process_count(MetricsEventType::Tx); + self.record_process(EventType::Transaction); } + /// 兼容性方法 - 添加账户处理计数 + #[inline] pub fn add_account_process_count(&self) { - self.add_process_count(MetricsEventType::Account); + self.record_process(EventType::Account); } + /// 兼容性方法 - 添加区块元数据处理计数 + #[inline] pub fn add_block_meta_process_count(&self) { - self.add_process_count(MetricsEventType::BlockMeta); + self.record_process(EventType::BlockMeta); } - /// 更新性能指标 + /// 兼容性方法 - 更新指标 + #[inline] pub fn update_metrics( &self, event_type: MetricsEventType, events_processed: u64, processing_time_us: f64, + signature: Option, ) { - // 检查是否启用性能监控 - if !self.config.enable_metrics { - return; - } - - // 使用 Backoff 策略进行写入尝试 - let backoff = Backoff::new(); - loop { - match self.metrics.write() { - Ok(mut metrics) => { - let now = std::time::Instant::now(); - let index = event_type.as_index(); - - // 更新事件计数 - metrics.event_metrics[index].events_processed += events_processed; - metrics.event_metrics[index].events_in_window += events_processed; - - metrics.last_update_time = now; - - // 更新处理时间统计 - if processing_time_us < metrics.min_processing_time_us - || metrics.min_processing_time_us == 0.0 - { - metrics.min_processing_time_us = processing_time_us; - } - if processing_time_us > metrics.max_processing_time_us { - metrics.max_processing_time_us = processing_time_us; - } - - // 计算平均处理时间 - 使用增量更新避免重复计算 - let total_events = metrics.event_metrics[index].events_processed; - if total_events > 0 { - let total_events_f64 = total_events as f64; - let old_total = (total_events_f64 - events_processed as f64).max(0.0); - - metrics.average_processing_time_us = if old_total > 0.0 { - (metrics.average_processing_time_us * old_total - + processing_time_us * events_processed as f64) - / total_events_f64 - } else { - processing_time_us - }; - } - - // 更新时间窗口指标 - let window_duration = std::time::Duration::from_secs(DEFAULT_METRICS_WINDOW_SECONDS); - metrics.update_window_metrics(&event_type, now, window_duration); - break; - }, - Err(_) => { - // 如果获取写锁失败,使用指数退避策略 - backoff.snooze(); - continue; - } - } - } - } - - /// 记录慢处理操作 - pub fn log_slow_processing(&self, processing_time_us: f64, event_count: usize) { - if processing_time_us > SLOW_PROCESSING_THRESHOLD_US { - log::warn!( - "{} slow processing: {processing_time_us}us for {event_count} events", - self.stream_name - ); - } + self.record_events(event_type, events_processed, processing_time_us); + self.log_slow_processing(processing_time_us, events_processed as usize, signature); } } @@ -346,8 +556,9 @@ impl Clone for MetricsManager { fn clone(&self) -> Self { Self { metrics: self.metrics.clone(), - config: self.config.clone(), + enable_metrics: self.enable_metrics, stream_name: self.stream_name.clone(), + background_task_running: AtomicBool::new(false), // 新实例不自动启动后台任务 } } } diff --git a/src/streaming/event_parser/common/mod.rs b/src/streaming/event_parser/common/mod.rs index 3562a92..85c2d90 100755 --- a/src/streaming/event_parser/common/mod.rs +++ b/src/streaming/event_parser/common/mod.rs @@ -54,7 +54,7 @@ macro_rules! impl_unified_event { Box::new(self.clone()) } - fn merge(&mut self, other: Box) { + fn merge(&mut self, other: &dyn $crate::streaming::event_parser::core::traits::UnifiedEvent) { if let Some(_e) = other.as_any().downcast_ref::<$struct_name>() { $( self.$field = _e.$field.clone(); @@ -66,10 +66,13 @@ macro_rules! impl_unified_event { self.metadata.set_swap_data(swap_data); } - fn index(&self) -> String { - self.metadata.index.clone() + fn instruction_outer_index(&self) -> i64 { + self.metadata.instruction_outer_index } + fn instruction_inner_index(&self) -> Option { + self.metadata.instruction_inner_index + } fn transaction_index(&self) -> Option { self.metadata.transaction_index } diff --git a/src/streaming/event_parser/common/types.rs b/src/streaming/event_parser/common/types.rs index e1ab15e..2bc74ba 100755 --- a/src/streaming/event_parser/common/types.rs +++ b/src/streaming/event_parser/common/types.rs @@ -55,36 +55,9 @@ impl EventMetadataPool { } } -/// Transfer data object pool -pub struct TransferDataPool { - pool: Arc>, -} - -impl Default for TransferDataPool { - fn default() -> Self { - Self::new() - } -} - -impl TransferDataPool { - pub fn new() -> Self { - Self { pool: Arc::new(ArrayQueue::new(TRANSFER_DATA_POOL_SIZE)) } - } - - pub fn acquire(&self) -> Option { - self.pool.pop() - } - - pub fn release(&self, transfer_data: TransferData) { - // 如果队列已满,push 会失败,但不会阻塞 - let _ = self.pool.push(transfer_data); - } -} - // Global object pool instances lazy_static::lazy_static! { pub static ref EVENT_METADATA_POOL: EventMetadataPool = EventMetadataPool::new(); - pub static ref TRANSFER_DATA_POOL: TransferDataPool = TransferDataPool::new(); } #[derive( @@ -304,20 +277,6 @@ impl ProtocolInfo { } } -/// Transfer data -#[derive( - Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, BorshSerialize, BorshDeserialize, -)] -pub struct TransferData { - pub token_program: Pubkey, - pub source: Pubkey, - pub destination: Pubkey, - pub authority: Option, - pub amount: u64, - pub decimals: Option, - pub mint: Option, -} - #[derive( Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, BorshSerialize, BorshDeserialize, )] @@ -337,7 +296,7 @@ pub struct EventMetadata { pub id: String, pub signature: String, pub slot: u64, - pub transaction_index: Option, // 新增:交易在slot中的索引 + pub transaction_index: Option, // 新增:交易在slot中的索引 pub block_time: i64, pub block_time_ms: i64, pub program_received_time_us: i64, @@ -345,10 +304,9 @@ pub struct EventMetadata { pub protocol: ProtocolType, pub event_type: EventType, pub program_id: Pubkey, - #[deprecated(note = "Please use swap_data instead")] - pub transfer_datas: Vec, pub swap_data: Option, - pub index: String, // 保留原有的指令索引 + pub instruction_outer_index: i64, + pub instruction_inner_index: Option, } impl EventMetadata { @@ -362,14 +320,15 @@ impl EventMetadata { protocol: ProtocolType, event_type: EventType, program_id: Pubkey, - index: String, + instruction_outer_index: i64, + instruction_inner_index: Option, program_received_time_us: i64, ) -> Self { Self { id, signature, slot, - transaction_index: None, // 默认为None,后续设置 + transaction_index: None, // 默认为None,后续设置 block_time, block_time_ms, program_received_time_us, @@ -377,9 +336,9 @@ impl EventMetadata { protocol, event_type, program_id, - transfer_datas: vec![], swap_data: None, - index, + instruction_outer_index, + instruction_inner_index, } } @@ -413,7 +372,7 @@ lazy_static::lazy_static! { /// Parse token transfer data from next instructions pub fn parse_swap_data_from_next_instructions( - event: Box, + event: &dyn UnifiedEvent, inner_instruction: &solana_transaction_status::InnerInstructions, current_index: i8, accounts: &[Pubkey], @@ -435,7 +394,7 @@ pub fn parse_swap_data_from_next_instructions( let mut from_vault: Option = None; let mut to_vault: Option = None; - match_event!(event, { + match_event!(&*event, { BonkTradeEvent => |e: BonkTradeEvent| { user = Some(e.payer); from_mint = Some(e.base_token_mint); diff --git a/src/streaming/event_parser/core/traits.rs b/src/streaming/event_parser/core/traits.rs index 67d6af8..f9c7f4c 100755 --- a/src/streaming/event_parser/core/traits.rs +++ b/src/streaming/event_parser/core/traits.rs @@ -7,6 +7,7 @@ use solana_sdk::{ use solana_transaction_status::{InnerInstructions, TransactionWithStatusMeta}; use std::collections::HashMap; use std::fmt::Debug; +use std::sync::Arc; use crate::streaming::event_parser::common::{parse_swap_data_from_next_instructions, SwapData}; use crate::streaming::event_parser::protocols::pumpswap::{PumpSwapBuyEvent, PumpSwapSellEvent}; @@ -17,7 +18,6 @@ use crate::streaming::event_parser::{ pumpfun::{PumpFunCreateTokenEvent, PumpFunTradeEvent}, }, }; -use crate::streaming::shred::MetricsEventType; /// Unified Event Interface - All protocol events must implement this trait pub trait UnifiedEvent: Debug + Send + Sync { @@ -55,7 +55,7 @@ pub trait UnifiedEvent: Debug + Send + Sync { fn clone_boxed(&self) -> Box; /// Merge events (optional implementation) - fn merge(&mut self, _other: Box) { + fn merge(&mut self, _other: &dyn UnifiedEvent) { // Default implementation: no merging operation } @@ -63,7 +63,8 @@ pub trait UnifiedEvent: Debug + Send + Sync { fn set_swap_data(&mut self, swap_data: SwapData); /// Get index - fn index(&self) -> String; + fn instruction_outer_index(&self) -> i64; + fn instruction_inner_index(&self) -> Option; /// Get transaction index in slot fn transaction_index(&self) -> Option; @@ -88,7 +89,8 @@ pub trait EventParser: Send + Sync { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec>; /// 从指令中解析事件数据 @@ -101,7 +103,8 @@ pub trait EventParser: Send + Sync { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec>; /// 从VersionedTransaction中解析指令事件的通用方法 @@ -144,7 +147,8 @@ pub trait EventParser: Send + Sync { slot, block_time, program_received_time_us, - format!("{index}"), + index as i64, + None, ) .await { @@ -156,7 +160,7 @@ pub trait EventParser: Send + Sync { { events.iter_mut().for_each(|event| { let swap_data = parse_swap_data_from_next_instructions( - event.clone_boxed(), + event.as_ref(), inn, -1_i8, &accounts, @@ -217,62 +221,100 @@ pub trait EventParser: Send + Sync { let mut inner_instructions: Vec = vec![]; if let Some(meta) = meta { inner_instructions = meta.inner_instructions.unwrap_or_default(); - for loopup in meta.loaded_addresses.writable { - address_table_lookups.push(loopup); - } - for loopup in meta.loaded_addresses.readonly { - address_table_lookups.push(loopup); - } + address_table_lookups.reserve( + meta.loaded_addresses.writable.len() + meta.loaded_addresses.readonly.len(), + ); + address_table_lookups.extend( + meta.loaded_addresses.writable.into_iter().chain(meta.loaded_addresses.readonly), + ); } - let mut accounts: Vec = vec![]; + + let mut accounts = Vec::with_capacity( + versioned_tx.message.static_account_keys().len() + address_table_lookups.len(), + ); + accounts.extend_from_slice(versioned_tx.message.static_account_keys()); + accounts.extend(address_table_lookups); + + // 使用 Arc 包装共享数据,避免不必要的克隆 + let accounts_arc = Arc::new(accounts); + let inner_instructions_arc = Arc::new(inner_instructions); // 预分配容量,避免动态扩容 let mut instruction_events: Vec> = Vec::with_capacity(16); + let mut inner_instruction_events: Vec> = Vec::with_capacity(8); - // 解析指令事件 - accounts = versioned_tx.message.static_account_keys().to_vec(); - accounts.extend(address_table_lookups.clone()); - - instruction_events = self - .parse_instruction_events_from_versioned_transaction( + let accounts_for_task1 = Arc::clone(&accounts_arc); + let inner_instructions_for_task1 = Arc::clone(&inner_instructions_arc); + let task1 = async move { + self.parse_instruction_events_from_versioned_transaction( &versioned_tx, signature, slot, block_time, program_received_time_us, - &accounts, - &inner_instructions, + &accounts_for_task1, + &inner_instructions_for_task1, ) .await - .unwrap_or_else(|_e| vec![]); + .unwrap_or_else(|_e| vec![]) + }; // 解析内联指令事件 - // 预分配容量,避免动态扩容 - let mut inner_instruction_events: Vec> = Vec::with_capacity(8); - // 检查交易是否成功 - for inner_instruction in inner_instructions { + let inner_instructions_for_task2 = Arc::clone(&inner_instructions_arc); + let accounts_for_task2_1 = Arc::clone(&accounts_arc); + let accounts_for_task2_2 = Arc::clone(&accounts_arc); + + let mut task2_params = Vec::with_capacity(inner_instructions_for_task2.len() * 5); + for inner_instruction in inner_instructions_for_task2.iter() { for (index, instruction) in inner_instruction.instructions.iter().enumerate() { - // 解析嵌套指令 - let compiled_instruction = instruction.instruction.clone(); + task2_params.push(( + &instruction.instruction, + signature, + slot, + block_time, + program_received_time_us, + inner_instruction.index as i64, + Some(index as i64), + inner_instruction, + )); + } + } + // 转换为 Arc<[T]> 更轻量 + let task2_params: Arc<[_]> = task2_params.into(); + let task2_params_clone = Arc::clone(&task2_params); + let task2_1 = async move { + let mut instruction_events: Vec> = Vec::with_capacity(16); + for ( + instruction, + signature, + slot, + block_time, + program_received_time_us, + outer_index, + inner_index, + inner_instruction, + ) in task2_params_clone.iter() + { if let Ok(mut events) = self .parse_instruction( - &compiled_instruction, - &accounts, - signature, - slot, - block_time, - program_received_time_us, - format!("{}.{}", inner_instruction.index, index), + instruction, + &accounts_for_task2_1, + *signature, + *slot, + *block_time, + *program_received_time_us, + *outer_index, + *inner_index, ) .await { if !events.is_empty() { events.iter_mut().for_each(|event| { let swap_data = parse_swap_data_from_next_instructions( - event.clone_boxed(), + event.as_ref(), &inner_instruction, - index as i8, - &accounts, + inner_index.unwrap_or_default() as i8, + &accounts_for_task2_1, ); if let Some(swap_data) = swap_data { event.set_swap_data(swap_data); @@ -281,24 +323,41 @@ pub trait EventParser: Send + Sync { instruction_events.extend(events); } } + } + instruction_events + }; + let task2_2 = async move { + let mut inner_instruction_events: Vec> = Vec::with_capacity(8); + for ( + instruction, + signature, + slot, + block_time, + program_received_time_us, + outer_index, + inner_index, + inner_instruction, + ) in task2_params.iter() + { if let Ok(mut events) = self .parse_inner_instruction( - &compiled_instruction, - signature, - slot, - block_time, - program_received_time_us, - format!("{}.{}", inner_instruction.index, index), + instruction, + *signature, + *slot, + *block_time, + *program_received_time_us, + *outer_index, + *inner_index, ) .await { if !events.is_empty() { events.iter_mut().for_each(|event| { let swap_data = parse_swap_data_from_next_instructions( - event.clone_boxed(), + event.as_ref(), &inner_instruction, - index as i8, - &accounts, + inner_index.unwrap_or_default() as i8, + &accounts_for_task2_2, ); if let Some(swap_data) = swap_data { event.set_swap_data(swap_data); @@ -308,39 +367,41 @@ pub trait EventParser: Send + Sync { } } } - } + inner_instruction_events + }; + + let (r1, r2_1, r2_2) = tokio::join!(task1, task2_1, task2_2); + instruction_events.extend(r1); + instruction_events.extend(r2_1); + inner_instruction_events.extend(r2_2); if !instruction_events.is_empty() && !inner_instruction_events.is_empty() { for instruction_event in &mut instruction_events { for inner_instruction_event in &inner_instruction_events { if instruction_event.id() == inner_instruction_event.id() { - let i_index = instruction_event.index(); - let in_index = inner_instruction_event.index(); - if !i_index.contains(".") && in_index.contains(".") { - let in_index_parts: Vec<&str> = in_index.split(".").collect(); - if !in_index_parts.is_empty() && in_index_parts[0] == i_index { - instruction_event.merge(inner_instruction_event.clone_boxed()); + if instruction_event.instruction_inner_index().is_none() + && inner_instruction_event.instruction_inner_index().is_some() + { + if inner_instruction_event.instruction_outer_index() + == instruction_event.instruction_outer_index() + { + instruction_event.merge(inner_instruction_event.as_ref()); break; } - } else if i_index.contains(".") && in_index.contains(".") { - // 嵌套指令 - let i_index_parts: Vec<&str> = i_index.split(".").collect(); - let in_index_parts: Vec<&str> = in_index.split(".").collect(); - - if !i_index_parts.is_empty() - && !in_index_parts.is_empty() - && i_index_parts[0] == in_index_parts[0] + } else if instruction_event.instruction_inner_index().is_some() + && inner_instruction_event.instruction_inner_index().is_some() + { + if instruction_event.instruction_outer_index() + == inner_instruction_event.instruction_outer_index() { - let i_index_child_index = i_index_parts - .get(1) - .and_then(|s| s.parse::().ok()) - .unwrap_or(0); - let in_index_child_index = in_index_parts - .get(1) - .and_then(|s| s.parse::().ok()) - .unwrap_or(0); - if in_index_child_index > i_index_child_index { - instruction_event.merge(inner_instruction_event.clone_boxed()); + if inner_instruction_event + .instruction_inner_index() + .unwrap_or_default() + > instruction_event + .instruction_inner_index() + .unwrap_or_default() + { + instruction_event.merge(inner_instruction_event.as_ref()); break; } } @@ -433,7 +494,8 @@ pub trait EventParser: Send + Sync { slot: Option, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Result>> { let slot = slot.unwrap_or(0); let events = self.parse_events_from_inner_instruction( @@ -442,7 +504,8 @@ pub trait EventParser: Send + Sync { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ); Ok(events) } @@ -456,7 +519,8 @@ pub trait EventParser: Send + Sync { slot: Option, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Result>> { let slot = slot.unwrap_or(0); let events = self.parse_events_from_instruction( @@ -466,7 +530,8 @@ pub trait EventParser: Send + Sync { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ); Ok(events) } @@ -543,7 +608,8 @@ impl GenericEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Option> { if let Some(parser) = config.inner_instruction_parser { let timestamp = block_time.unwrap_or(Timestamp { seconds: 0, nanos: 0 }); @@ -557,7 +623,8 @@ impl GenericEventParser { config.protocol_type.clone(), config.event_type.clone(), config.program_id, - index, + outer_index, + inner_index, program_received_time_us, ); parser(data, metadata) @@ -577,7 +644,8 @@ impl GenericEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Option> { if let Some(parser) = config.instruction_parser { let timestamp = block_time.unwrap_or(Timestamp { seconds: 0, nanos: 0 }); @@ -591,7 +659,8 @@ impl GenericEventParser { config.protocol_type.clone(), config.event_type.clone(), config.program_id, - index, + outer_index, + inner_index, program_received_time_us, ); parser(data, account_pubkeys, metadata) @@ -604,9 +673,11 @@ impl GenericEventParser { #[async_trait::async_trait] impl EventParser for GenericEventParser { fn inner_instruction_configs(&self) -> HashMap<&'static str, Vec> { + // 返回引用而非克隆,减少内存分配 self.inner_instruction_configs.clone() } fn instruction_configs(&self) -> HashMap, Vec> { + // 返回引用而非克隆,减少内存分配 self.instruction_configs.clone() } /// 从内联指令中解析事件数据 @@ -618,7 +689,8 @@ impl EventParser for GenericEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { let inner_instruction_data_decoded = inner_instruction.data.clone(); if inner_instruction_data_decoded.len() < 16 { @@ -638,7 +710,8 @@ impl EventParser for GenericEventParser { slot, block_time, program_received_time_us, - index.clone(), + outer_index, + inner_index, ) { events.push(event); } @@ -658,7 +731,8 @@ impl EventParser for GenericEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { let program_id = accounts[instruction.program_id_index as usize]; if !self.should_handle(&program_id) { @@ -691,7 +765,8 @@ impl EventParser for GenericEventParser { slot, block_time, program_received_time_us, - index.clone(), + outer_index, + inner_index, ) { events.push(event); } diff --git a/src/streaming/event_parser/protocols/block/block_meta_event.rs b/src/streaming/event_parser/protocols/block/block_meta_event.rs index 2987246..2f99ba1 100644 --- a/src/streaming/event_parser/protocols/block/block_meta_event.rs +++ b/src/streaming/event_parser/protocols/block/block_meta_event.rs @@ -28,7 +28,8 @@ impl BlockMetaEvent { crate::streaming::event_parser::common::types::ProtocolType::Common, EventType::BlockMeta, solana_sdk::pubkey::Pubkey::default(), - "".to_string(), + 0, + None, program_received_time_us, ); Self { metadata, slot, block_hash } diff --git a/src/streaming/event_parser/protocols/bonk/parser.rs b/src/streaming/event_parser/protocols/bonk/parser.rs index a03a38e..bc28e85 100755 --- a/src/streaming/event_parser/protocols/bonk/parser.rs +++ b/src/streaming/event_parser/protocols/bonk/parser.rs @@ -618,7 +618,8 @@ impl EventParser for BonkEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_inner_instruction( inner_instruction, @@ -626,7 +627,8 @@ impl EventParser for BonkEventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } @@ -638,7 +640,8 @@ impl EventParser for BonkEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_instruction( instruction, @@ -647,7 +650,8 @@ impl EventParser for BonkEventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } diff --git a/src/streaming/event_parser/protocols/mutil/parser.rs b/src/streaming/event_parser/protocols/mutil/parser.rs index 9748124..b04da73 100755 --- a/src/streaming/event_parser/protocols/mutil/parser.rs +++ b/src/streaming/event_parser/protocols/mutil/parser.rs @@ -79,7 +79,8 @@ impl EventParser for MutilEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_inner_instruction( inner_instruction, @@ -87,7 +88,8 @@ impl EventParser for MutilEventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } @@ -99,7 +101,8 @@ impl EventParser for MutilEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_instruction( instruction, @@ -108,7 +111,8 @@ impl EventParser for MutilEventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } diff --git a/src/streaming/event_parser/protocols/pumpfun/parser.rs b/src/streaming/event_parser/protocols/pumpfun/parser.rs index e907c11..2e97979 100755 --- a/src/streaming/event_parser/protocols/pumpfun/parser.rs +++ b/src/streaming/event_parser/protocols/pumpfun/parser.rs @@ -299,7 +299,8 @@ impl EventParser for PumpFunEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_inner_instruction( inner_instruction, @@ -307,7 +308,8 @@ impl EventParser for PumpFunEventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } @@ -319,7 +321,8 @@ impl EventParser for PumpFunEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_instruction( instruction, @@ -328,7 +331,8 @@ impl EventParser for PumpFunEventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } diff --git a/src/streaming/event_parser/protocols/pumpswap/parser.rs b/src/streaming/event_parser/protocols/pumpswap/parser.rs index 6bd9b57..a688c41 100755 --- a/src/streaming/event_parser/protocols/pumpswap/parser.rs +++ b/src/streaming/event_parser/protocols/pumpswap/parser.rs @@ -389,7 +389,8 @@ impl EventParser for PumpSwapEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_inner_instruction( inner_instruction, @@ -397,7 +398,8 @@ impl EventParser for PumpSwapEventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } @@ -409,7 +411,8 @@ impl EventParser for PumpSwapEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_instruction( instruction, @@ -418,7 +421,8 @@ impl EventParser for PumpSwapEventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } diff --git a/src/streaming/event_parser/protocols/raydium_amm_v4/parser.rs b/src/streaming/event_parser/protocols/raydium_amm_v4/parser.rs index 593e1af..288e523 100755 --- a/src/streaming/event_parser/protocols/raydium_amm_v4/parser.rs +++ b/src/streaming/event_parser/protocols/raydium_amm_v4/parser.rs @@ -391,7 +391,8 @@ impl EventParser for RaydiumAmmV4EventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_inner_instruction( inner_instruction, @@ -399,7 +400,8 @@ impl EventParser for RaydiumAmmV4EventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } @@ -411,7 +413,8 @@ impl EventParser for RaydiumAmmV4EventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_instruction( instruction, @@ -420,7 +423,8 @@ impl EventParser for RaydiumAmmV4EventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } diff --git a/src/streaming/event_parser/protocols/raydium_clmm/parser.rs b/src/streaming/event_parser/protocols/raydium_clmm/parser.rs index aef8ba9..16fab61 100755 --- a/src/streaming/event_parser/protocols/raydium_clmm/parser.rs +++ b/src/streaming/event_parser/protocols/raydium_clmm/parser.rs @@ -432,7 +432,8 @@ impl EventParser for RaydiumClmmEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_inner_instruction( inner_instruction, @@ -440,7 +441,8 @@ impl EventParser for RaydiumClmmEventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } @@ -452,7 +454,8 @@ impl EventParser for RaydiumClmmEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_instruction( instruction, @@ -461,7 +464,8 @@ impl EventParser for RaydiumClmmEventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } diff --git a/src/streaming/event_parser/protocols/raydium_cpmm/parser.rs b/src/streaming/event_parser/protocols/raydium_cpmm/parser.rs index 6dfdd6b..288bf60 100755 --- a/src/streaming/event_parser/protocols/raydium_cpmm/parser.rs +++ b/src/streaming/event_parser/protocols/raydium_cpmm/parser.rs @@ -282,7 +282,8 @@ impl EventParser for RaydiumCpmmEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_inner_instruction( inner_instruction, @@ -290,7 +291,8 @@ impl EventParser for RaydiumCpmmEventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } @@ -302,7 +304,8 @@ impl EventParser for RaydiumCpmmEventParser { slot: u64, block_time: Option, program_received_time_us: i64, - index: String, + outer_index: i64, + inner_index: Option, ) -> Vec> { self.inner.parse_events_from_instruction( instruction, @@ -311,7 +314,8 @@ impl EventParser for RaydiumCpmmEventParser { slot, block_time, program_received_time_us, - index, + outer_index, + inner_index, ) } diff --git a/src/streaming/grpc/stream_handler.rs b/src/streaming/grpc/stream_handler.rs index ff6fb00..43dbb77 100644 --- a/src/streaming/grpc/stream_handler.rs +++ b/src/streaming/grpc/stream_handler.rs @@ -8,7 +8,6 @@ use yellowstone_grpc_proto::geyser::{ use super::types::{BlockMetaPretty, EventPretty, TransactionPretty}; use crate::common::AnyResult; use crate::streaming::common::EventProcessor; -use crate::streaming::event_parser::UnifiedEvent; use crate::streaming::grpc::AccountPretty; /// 流消息处理器 @@ -16,16 +15,12 @@ pub struct StreamHandler; impl StreamHandler { /// 处理单个流消息 - pub async fn handle_stream_message( + pub async fn handle_stream_message( msg: SubscribeUpdate, subscribe_tx: &mut (impl Sink + Unpin), event_processor: EventProcessor, - callback: &F, bot_wallet: Option, - ) -> AnyResult<()> - where - F: Fn(Box) + Send + Sync, - { + ) -> AnyResult<()> { let created_at = msg.created_at; match msg.update_oneof { Some(UpdateOneof::Account(account)) => { @@ -34,7 +29,6 @@ impl StreamHandler { event_processor .process_grpc_event_transaction_with_metrics( EventPretty::Account(account_pretty), - callback, bot_wallet, ) .await?; @@ -45,7 +39,6 @@ impl StreamHandler { event_processor .process_grpc_event_transaction_with_metrics( EventPretty::BlockMeta(block_meta_pretty), - callback, bot_wallet, ) .await?; @@ -60,7 +53,6 @@ impl StreamHandler { event_processor .process_grpc_event_transaction_with_metrics( EventPretty::Transaction(transaction_pretty), - callback, bot_wallet, ) .await?; @@ -84,58 +76,25 @@ impl StreamHandler { Ok(()) } - // /// 处理背压策略 - // async fn handle_backpressure( - // tx: &mut mpsc::Sender, - // event_pretty: EventPretty, - // backpressure_strategy: BackpressureStrategy, - // ) -> AnyResult<()> { - // match backpressure_strategy { - // BackpressureStrategy::Block => { - // // 阻塞等待,直到有空间 - // if let Err(e) = tx.send(event_pretty).await { - // log::error!("Failed to send transaction to channel: {:?}", e); - // return Err(anyhow::anyhow!("Channel send failed: {:?}", e)); - // } - // } - // BackpressureStrategy::Drop => { - // // 尝试发送,如果失败则丢弃 - // if let Err(e) = tx.try_send(event_pretty) { - // if e.is_full() { - // log::warn!("Channel is full, dropping transaction"); - // } else { - // log::error!("Channel is closed: {:?}", e); - // return Err(anyhow::anyhow!("Channel is closed: {:?}", e)); - // } - // } - // } - // BackpressureStrategy::Retry { max_attempts, wait_ms } => { - // // 重试有限次数 - // let mut retry_count = 0; - // loop { - // match tx.try_send(event_pretty.clone()) { - // Ok(_) => break, - // Err(e) => { - // if e.is_full() { - // retry_count += 1; - // if retry_count >= max_attempts { - // log::warn!( - // "Channel is full after {} attempts, dropping transaction", - // retry_count - // ); - // break; - // } - // tokio::time::sleep(tokio::time::Duration::from_millis(wait_ms)) - // .await; - // } else { - // log::error!("Channel is closed: {:?}", e); - // return Err(anyhow::anyhow!("Channel is closed: {:?}", e)); - // } - // } - // } - // } - // } - // } - // Ok(()) - // } + pub async fn handle_stream_system_message( + msg: SubscribeUpdate, + subscribe_tx: &mut (impl Sink + Unpin), + ) -> AnyResult> { + let created_at = msg.created_at; + let event_pretty = match msg.update_oneof { + Some(UpdateOneof::Transaction(sut)) => Some(TransactionPretty::from((sut, created_at))), + Some(UpdateOneof::Ping(_)) => { + subscribe_tx + .send(SubscribeRequest { + ping: Some(SubscribeRequestPing { id: 1 }), + ..Default::default() + }) + .await?; + None + } + Some(UpdateOneof::Pong(_)) => None, + _ => None, + }; + Ok(event_pretty.map(|e| EventPretty::Transaction(e))) + } } diff --git a/src/streaming/grpc/types.rs b/src/streaming/grpc/types.rs index 9238f9f..7cfd946 100644 --- a/src/streaming/grpc/types.rs +++ b/src/streaming/grpc/types.rs @@ -69,7 +69,7 @@ impl fmt::Debug for BlockMetaPretty { #[derive(Clone)] pub struct TransactionPretty { pub slot: u64, - pub transaction_index: Option, // 新增:交易在slot中的索引 + pub transaction_index: Option, // 新增:交易在slot中的索引 pub block_hash: String, pub block_time: Option, pub signature: Signature, @@ -95,10 +95,11 @@ impl From for AccountPretty { let account_info = account.account.unwrap(); Self { slot: account.slot, - signature: Signature::try_from( - account_info.txn_signature.unwrap_or_default().as_slice(), - ) - .expect("valid signature"), + signature: if let Some(txn_signature) = account_info.txn_signature { + Signature::try_from(txn_signature.as_slice()).expect("valid signature") + } else { + Signature::default() + }, pubkey: Pubkey::try_from(account_info.pubkey.as_slice()).expect("valid pubkey"), executable: account_info.executable, lamports: account_info.lamports, @@ -138,7 +139,7 @@ impl From<(SubscribeUpdateTransaction, Option)> for TransactionPretty let transaction_index = tx.index; Self { slot, - transaction_index: Some(transaction_index), // 提取交易索引 + transaction_index: Some(transaction_index), // 提取交易索引 block_time, block_hash: "".to_string(), signature: Signature::try_from(tx.signature.as_slice()).expect("valid signature"), diff --git a/src/streaming/mod.rs b/src/streaming/mod.rs index 9aed6b5..8747e2a 100755 --- a/src/streaming/mod.rs +++ b/src/streaming/mod.rs @@ -4,8 +4,8 @@ pub mod grpc; pub mod shred; pub mod shred_stream; pub mod yellowstone_grpc; -// pub mod yellowstone_sub_system; +pub mod yellowstone_sub_system; pub use shred::ShredStreamGrpc; pub use yellowstone_grpc::YellowstoneGrpc; -// pub use yellowstone_sub_system::{SystemEvent, TransferInfo}; +pub use yellowstone_sub_system::{SystemEvent, TransferInfo}; diff --git a/src/streaming/shred/connection.rs b/src/streaming/shred/connection.rs index 90bc64a..44e7e4a 100644 --- a/src/streaming/shred/connection.rs +++ b/src/streaming/shred/connection.rs @@ -29,10 +29,8 @@ impl ShredStreamGrpc { pub async fn new_with_config(endpoint: String, config: StreamClientConfig) -> AnyResult { let shredstream_client = ShredstreamProxyClient::connect(endpoint.clone()).await?; let metrics = Arc::new(RwLock::new(PerformanceMetrics::new())); - let config_arc = Arc::new(config.clone()); - let metrics_manager = - MetricsManager::new(metrics.clone(), config_arc, "ShredStream".to_string()); + let metrics_manager = MetricsManager::new(config.enable_metrics, "ShredStream".to_string()); Ok(Self { shredstream_client: Arc::new(shredstream_client), diff --git a/src/streaming/shred_stream.rs b/src/streaming/shred_stream.rs index 860429a..79cda1b 100755 --- a/src/streaming/shred_stream.rs +++ b/src/streaming/shred_stream.rs @@ -1,3 +1,5 @@ +use std::sync::Arc; + use futures::StreamExt; use solana_sdk::pubkey::Pubkey; @@ -39,8 +41,9 @@ impl ShredStreamGrpc { event_processor.set_protocols_and_event_type_filter( protocols, event_type_filter, - self.config.backpressure.strategy, + self.config.backpressure.clone(), self.config.batch.clone(), + Some(Arc::new(callback)), ); // 启动流处理 @@ -61,7 +64,6 @@ impl ShredStreamGrpc { .process_shred_transaction_immediate( transaction_with_slot, bot_wallet, - &callback, ) .await { @@ -82,11 +84,7 @@ impl ShredStreamGrpc { }); // 保存订阅句柄 - let subscription_handle = SubscriptionHandle::new( - stream_task, - event_processor.get_event_handle(), - metrics_handle, - ); + let subscription_handle = SubscriptionHandle::new(stream_task, None, metrics_handle); let mut handle_guard = self.subscription_handle.lock().await; *handle_guard = Some(subscription_handle); diff --git a/src/streaming/yellowstone_grpc.rs b/src/streaming/yellowstone_grpc.rs index 58075cd..db928d9 100644 --- a/src/streaming/yellowstone_grpc.rs +++ b/src/streaming/yellowstone_grpc.rs @@ -51,12 +51,14 @@ impl YellowstoneGrpc { ) -> AnyResult { let _ = rustls::crypto::ring::default_provider().install_default().ok(); let metrics = Arc::new(RwLock::new(PerformanceMetrics::new())); - let config_arc = Arc::new(config.clone()); let subscription_manager = SubscriptionManager::new(endpoint.clone(), x_token.clone(), config.clone()); - let metrics_manager = - MetricsManager::new(metrics.clone(), config_arc.clone(), "YellowstoneGrpc".to_string()); + let metrics_manager = MetricsManager::new_with_metrics( + metrics.clone(), + config.enable_metrics, + "YellowstoneGrpc".to_string(), + ); let event_processor = EventProcessor::new(metrics_manager.clone(), config.clone()); Ok(Self { @@ -179,8 +181,9 @@ impl YellowstoneGrpc { event_processor.set_protocols_and_event_type_filter( protocols, event_type_filter, - self.config.backpressure.strategy, + self.config.backpressure.clone(), self.config.batch.clone(), + Some(Arc::new(callback)), ); let stream_handle = tokio::spawn(async move { while let Some(message) = stream.next().await { @@ -190,7 +193,6 @@ impl YellowstoneGrpc { msg, &mut subscribe_tx, event_processor.clone(), - &callback, bot_wallet, ) .await @@ -208,11 +210,7 @@ impl YellowstoneGrpc { }); // 保存订阅句柄 - let subscription_handle = SubscriptionHandle::new( - stream_handle, - self.event_processor.get_event_handle(), - metrics_handle, - ); + let subscription_handle = SubscriptionHandle::new(stream_handle, None, metrics_handle); let mut handle_guard = self.subscription_handle.lock().await; *handle_guard = Some(subscription_handle); diff --git a/src/streaming/yellowstone_sub_system.rs b/src/streaming/yellowstone_sub_system.rs index 4d19f95..bccd417 100755 --- a/src/streaming/yellowstone_sub_system.rs +++ b/src/streaming/yellowstone_sub_system.rs @@ -1,19 +1,17 @@ use crate::{ common::AnyResult, streaming::{ - grpc::{BackpressureStrategy, EventPretty, StreamHandler}, + grpc::{EventPretty, StreamHandler}, yellowstone_grpc::YellowstoneGrpc, }, }; -use futures::{channel::mpsc, StreamExt}; +use futures::StreamExt; use log::error; use solana_program::pubkey; use solana_sdk::{pubkey::Pubkey, transaction::VersionedTransaction}; use solana_transaction_status::TransactionWithStatusMeta; const SYSTEM_PROGRAM_ID: Pubkey = pubkey!("11111111111111111111111111111111"); -// 根据实际并发量调整通道大小,避免背压 -const CHANNEL_SIZE: usize = 50000; // 增加到 50000 #[derive(Debug)] pub enum SystemEvent { @@ -51,7 +49,6 @@ impl YellowstoneGrpc { .subscription_manager .subscribe_with_request(transactions, None, None, None) .await?; - let (mut tx, mut rx) = mpsc::channel::(CHANNEL_SIZE); let callback = Box::new(callback); @@ -59,16 +56,17 @@ impl YellowstoneGrpc { while let Some(message) = stream.next().await { match message { Ok(msg) => { - if let Err(e) = StreamHandler::handle_stream_message( - msg, - &mut tx, - &mut subscribe_tx, - BackpressureStrategy::Block, - ) - .await + if let Ok(event_pretty) = + StreamHandler::handle_stream_system_message(msg, &mut subscribe_tx) + .await { - error!("Error handling message: {e:?}"); - break; + if let Some(event_pretty) = event_pretty { + if let Err(e) = + Self::process_system_transaction(event_pretty, &*callback).await + { + error!("Error processing transaction: {e:?}"); + } + } } } Err(error) => { @@ -78,12 +76,6 @@ impl YellowstoneGrpc { } } }); - - while let Some(event_pretty) = rx.next().await { - if let Err(e) = Self::process_system_transaction(event_pretty, &*callback).await { - error!("Error processing transaction: {e:?}"); - } - } Ok(()) }