use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::sync::Arc; use std::time::Instant; use crossbeam_queue::SegQueue; use solana_sdk::pubkey::Pubkey; use crate::common::AnyResult; use crate::streaming::common::BackpressureStrategy; use crate::streaming::common::{ MetricsEventType, MetricsManager, StreamClientConfig as ClientConfig, }; use crate::streaming::event_parser::common::filter::EventTypeFilter; use crate::streaming::event_parser::core::account_event_parser::AccountEventParser; use crate::streaming::event_parser::core::common_event_parser::CommonEventParser; use crate::streaming::event_parser::core::traits::get_high_perf_clock; use crate::streaming::event_parser::EventParser; use crate::streaming::event_parser::{ core::traits::UnifiedEvent, protocols::mutil::parser::MutilEventParser, Protocol, }; use crate::streaming::grpc::{BackpressureConfig, EventPretty}; use crate::streaming::shred::TransactionWithSlot; use once_cell::sync::OnceCell; /// High-performance Event processor using SegQueue for all strategies pub struct EventProcessor { pub(crate) metrics_manager: MetricsManager, pub(crate) config: ClientConfig, pub(crate) parser_cache: OnceCell>, pub(crate) protocols: Vec, pub(crate) event_type_filter: Option, pub(crate) callback: Option) + Send + Sync>>, pub(crate) backpressure_config: BackpressureConfig, /// High-performance lockfree queue for gRPC events pub(crate) grpc_queue: Arc)>>, /// High-performance lockfree queue for shred events pub(crate) shred_queue: Arc)>>, /// Fast O(1) counter for Drop strategy (avoids expensive SegQueue::len()) pub(crate) grpc_pending_count: Arc, pub(crate) shred_pending_count: Arc, /// Processing thread control pub(crate) processing_shutdown: Arc, } impl EventProcessor { /// Create a new high-performance event processor pub fn new(metrics_manager: MetricsManager, config: ClientConfig) -> Self { let backpressure_config = config.backpressure.clone(); let grpc_queue = Arc::new(SegQueue::new()); let shred_queue = Arc::new(SegQueue::new()); let grpc_pending_count = Arc::new(AtomicUsize::new(0)); let shred_pending_count = Arc::new(AtomicUsize::new(0)); let processing_shutdown = Arc::new(AtomicBool::new(false)); Self { metrics_manager, config, parser_cache: OnceCell::new(), protocols: vec![], event_type_filter: None, backpressure_config, callback: None, grpc_queue, shred_queue, grpc_pending_count, shred_pending_count, processing_shutdown, } } pub fn set_protocols_and_event_type_filter( &mut self, protocols: Vec, event_type_filter: Option, backpressure_config: BackpressureConfig, callback: Option) + Send + Sync>>, ) { self.protocols = protocols; self.event_type_filter = event_type_filter; // Check if Block processing thread should be started (before moving backpressure_config) let should_start_block_processing = true; // matches!(backpressure_config.strategy, BackpressureStrategy::Block); self.backpressure_config = backpressure_config; self.callback = callback; // Use stored values to initialize parser_cache let protocols_ref = &self.protocols; let event_type_filter_ref = self.event_type_filter.as_ref(); self.parser_cache.get_or_init(|| { Arc::new(MutilEventParser::new(protocols_ref.clone(), event_type_filter_ref.cloned())) }); // Start Block processing thread if using Block strategy if should_start_block_processing { self.start_block_processing_thread(); } } pub fn get_parser(&self) -> Arc { self.parser_cache.get().unwrap().clone() } /// Create adapter callback fn create_adapter_callback(&self) -> Arc) + Send + Sync> { let callback = self.callback.clone().unwrap(); let metrics_manager = self.metrics_manager.clone(); Arc::new(move |event: Box| { let processing_time_us = event.program_handle_time_consuming_us() as f64; callback(event); metrics_manager.update_metrics(MetricsEventType::Transaction, 1, processing_time_us); }) } pub async fn process_grpc_event_transaction_with_metrics( &self, event_pretty: EventPretty, bot_wallet: Option, ) -> AnyResult<()> { self.apply_backpressure_control(event_pretty, bot_wallet).await } /// Apply backpressure control strategy async fn apply_backpressure_control( &self, event_pretty: EventPretty, bot_wallet: Option, ) -> AnyResult<()> { match self.backpressure_config.strategy { BackpressureStrategy::Block => { // Block strategy: async wait if queue is full (backpressure control) loop { let current_pending = self.grpc_pending_count.load(Ordering::Relaxed); if current_pending < self.backpressure_config.permits { self.grpc_queue.push((event_pretty, bot_wallet)); self.grpc_pending_count.fetch_add(1, Ordering::Relaxed); break; } // Async yield to avoid blocking gRPC data source tokio::task::yield_now().await; } Ok(()) } BackpressureStrategy::Drop => { // Drop strategy: Use O(1) atomic counter instead of expensive O(n) len() // If pending count >= permits, DROP the event immediately let current_pending = self.grpc_pending_count.load(Ordering::Relaxed); if current_pending >= self.backpressure_config.permits { self.metrics_manager.increment_dropped_events(); Ok(()) } else { self.grpc_pending_count.fetch_add(1, Ordering::Relaxed); let processor = self.clone(); tokio::spawn(async move { match processor .process_grpc_event_transaction(event_pretty, bot_wallet) .await { Ok(_) => { processor.grpc_pending_count.fetch_sub(1, Ordering::Relaxed); } Err(e) => { log::error!("Error in async gRPC processing: {}", e); } } }); Ok(()) } } } } async fn process_grpc_event_transaction( &self, event_pretty: EventPretty, bot_wallet: Option, ) -> AnyResult<()> { if self.callback.is_none() { return Ok(()); } match event_pretty { EventPretty::Account(account_pretty) => { self.metrics_manager.add_account_process_count(); let account_event = AccountEventParser::parse_account_event( &self.protocols, account_pretty, self.event_type_filter.as_ref(), ); if let Some(event) = account_event { let processing_time_us = event.program_handle_time_consuming_us() as f64; self.invoke_callback(event); self.update_metrics(MetricsEventType::Account, 1, processing_time_us); } } EventPretty::Transaction(transaction_pretty) => { self.metrics_manager.add_tx_process_count(); let slot = transaction_pretty.slot; let signature = transaction_pretty.signature; let tx = transaction_pretty.tx; let block_time = transaction_pretty.block_time; let program_received_time_us = transaction_pretty.program_received_time_us; let transaction_index = transaction_pretty.transaction_index; // Use cache to get parser let parser = self.get_parser(); let adapter_callback = self.create_adapter_callback(); parser .parse_transaction_owned( tx, signature, Some(slot), block_time, program_received_time_us, bot_wallet, transaction_index, adapter_callback, ) .await?; } EventPretty::BlockMeta(block_meta_pretty) => { self.metrics_manager.add_block_meta_process_count(); let block_time_ms = block_meta_pretty .block_time .map(|ts| ts.seconds * 1000 + ts.nanos as i64 / 1_000_000) .unwrap_or_else(|| chrono::Utc::now().timestamp_millis()); let block_meta_event = CommonEventParser::generate_block_meta_event( block_meta_pretty.slot, &block_meta_pretty.block_hash, block_time_ms, block_meta_pretty.program_received_time_us, ); let processing_time_us = block_meta_event.program_handle_time_consuming_us() as f64; self.invoke_callback(block_meta_event); self.update_metrics(MetricsEventType::BlockMeta, 1, processing_time_us); } } Ok(()) } pub fn invoke_callback(&self, event: Box) { if let Some(callback) = self.callback.as_ref() { callback(event); } } /// Process a single transaction immediately pub async fn process_shred_transaction_immediate( &self, transaction_with_slot: TransactionWithSlot, bot_wallet: Option, ) -> AnyResult<()> { self.process_shred_transaction(transaction_with_slot, bot_wallet).await } /// Process shred transaction with backpressure control and performance monitoring pub async fn process_shred_transaction_with_metrics( &self, transaction_with_slot: TransactionWithSlot, bot_wallet: Option, ) -> AnyResult<()> { // Backpressure control logic self.apply_shred_backpressure_control(transaction_with_slot, bot_wallet).await } /// Apply shred backpressure control strategy async fn apply_shred_backpressure_control( &self, transaction_with_slot: TransactionWithSlot, bot_wallet: Option, ) -> AnyResult<()> { match self.backpressure_config.strategy { BackpressureStrategy::Block => { // Block strategy: async wait if queue is full (backpressure control) loop { let current_pending = self.shred_pending_count.load(Ordering::Relaxed); if current_pending < self.backpressure_config.permits { self.shred_queue.push((transaction_with_slot, bot_wallet)); self.shred_pending_count.fetch_add(1, Ordering::Relaxed); break; } // Async yield to avoid blocking shred data source tokio::task::yield_now().await; } Ok(()) } BackpressureStrategy::Drop => { // Drop strategy: Use O(1) atomic counter instead of expensive O(n) len() let current_pending = self.shred_pending_count.load(Ordering::Relaxed); if current_pending >= self.backpressure_config.permits { self.metrics_manager.increment_dropped_events(); Ok(()) } else { self.shred_pending_count.fetch_add(1, Ordering::Relaxed); let processor = self.clone(); tokio::spawn(async move { match processor .process_shred_transaction(transaction_with_slot, bot_wallet) .await { Ok(_) => { processor.shred_pending_count.fetch_sub(1, Ordering::Relaxed); } Err(e) => { log::error!("Error in async shred processing: {}", e); } } }); Ok(()) } } } } pub async fn process_shred_transaction( &self, transaction_with_slot: TransactionWithSlot, bot_wallet: Option, ) -> AnyResult<()> { if self.callback.is_none() { return Ok(()); } self.metrics_manager.add_tx_process_count(); let tx = transaction_with_slot.transaction; let slot = transaction_with_slot.slot; let signature = tx.signatures[0]; let program_received_time_us = transaction_with_slot.program_received_time_us; // Use cache to get parser let parser = self.get_parser(); let adapter_callback = self.create_adapter_callback(); parser .parse_versioned_transaction_owned( tx, signature, Some(slot), None, program_received_time_us, bot_wallet, None, &[], adapter_callback, ) .await?; Ok(()) } fn update_metrics(&self, ty: MetricsEventType, count: u64, time_us: f64) { self.metrics_manager.update_metrics(ty, count, time_us); } /// Start dedicated processing threads for all strategies fn start_block_processing_thread(&self) { // Reset shutdown flag self.processing_shutdown.store(false, Ordering::Relaxed); let grpc_queue = Arc::clone(&self.grpc_queue); let shred_queue = Arc::clone(&self.shred_queue); let grpc_pending_count = Arc::clone(&self.grpc_pending_count); let shred_pending_count = Arc::clone(&self.shred_pending_count); let shutdown_flag = Arc::clone(&self.processing_shutdown); let shutdown_flag_clone = Arc::clone(&self.processing_shutdown); let processor = self.clone(); let processor_clone = self.clone(); // 1. 专用线程 + 2. Busy-wait + 4. 无锁处理 std::thread::spawn(move || { // 创建blocking runtime for async processing let rt = tokio::runtime::Builder::new_current_thread().enable_all().build().unwrap(); while !shutdown_flag.load(Ordering::Relaxed) { if let Some((event_pretty, bot_wallet)) = grpc_queue.pop() { // Decrement pending counter when consuming from queue grpc_pending_count.fetch_sub(1, Ordering::Relaxed); // Process event in blocking runtime if let Err(e) = rt.block_on( processor.process_grpc_event_transaction(event_pretty, bot_wallet), ) { println!("Error processing gRPC event: {}", e); } } else { // 2. 优化忙等待: 使用轻量级休眠减少CPU占用 std::thread::yield_now(); } } }); // Shred处理也使用相同的低延迟优化 std::thread::spawn(move || { let rt = tokio::runtime::Builder::new_current_thread().enable_all().build().unwrap(); while !shutdown_flag_clone.load(Ordering::Relaxed) { if let Some((transaction_with_slot, bot_wallet)) = shred_queue.pop() { // Decrement pending counter when consuming from queue shred_pending_count.fetch_sub(1, Ordering::Relaxed); // Process transaction in blocking runtime if let Err(e) = rt.block_on( processor_clone .process_shred_transaction(transaction_with_slot, bot_wallet), ) { log::error!("Error processing shred transaction: {}", e); } } else { // 优化忙等待: 使用轻量级休眠减少CPU占用 std::thread::yield_now(); } } }); } /// Stop processing threads pub fn stop_processing(&self) { self.processing_shutdown.store(true, Ordering::Relaxed); } } // Implement Clone trait to support sharing between modules impl Clone for EventProcessor { fn clone(&self) -> Self { Self { metrics_manager: self.metrics_manager.clone(), config: self.config.clone(), parser_cache: self.parser_cache.clone(), protocols: self.protocols.clone(), event_type_filter: self.event_type_filter.clone(), backpressure_config: self.backpressure_config.clone(), callback: self.callback.clone(), grpc_queue: self.grpc_queue.clone(), shred_queue: self.shred_queue.clone(), grpc_pending_count: self.grpc_pending_count.clone(), shred_pending_count: self.shred_pending_count.clone(), processing_shutdown: self.processing_shutdown.clone(), } } }