From 45ab6eb06131f1376a055e009a6d0106963394d3 Mon Sep 17 00:00:00 2001 From: 0xfnzero <0xfnzero@users.noreply.github.com> Date: Mon, 25 May 2026 00:41:25 +0800 Subject: [PATCH] Fix grpc subscription cleanup and API compatibility --- src/streaming/grpc/pool.rs | 84 ++++++++++++++-------- src/streaming/yellowstone_grpc.rs | 93 ++++++++++++++++++++++--- src/streaming/yellowstone_sub_system.rs | 2 +- 3 files changed, 137 insertions(+), 42 deletions(-) diff --git a/src/streaming/grpc/pool.rs b/src/streaming/grpc/pool.rs index b0e16e8..3726369 100644 --- a/src/streaming/grpc/pool.rs +++ b/src/streaming/grpc/pool.rs @@ -1,6 +1,5 @@ use super::types::{AccountPretty, BlockMetaPretty, TransactionPretty}; use crate::streaming::event_parser::common::high_performance_clock::get_high_perf_clock; -use log::warn; use solana_sdk::{pubkey::Pubkey, signature::Signature}; use std::collections::VecDeque; use std::ops::DerefMut; @@ -96,23 +95,20 @@ impl PooledAccountPretty { /// 从 gRPC 更新重置数据 pub fn reset_from_update(&mut self, account_update: SubscribeUpdateAccount) -> bool { let Some(account_info) = account_update.account else { - warn!("grpc account update missing nested account, skipping"); + log::debug!("drop account update without account payload"); return false; }; self.account.slot = account_update.slot; self.account.signature = if let Some(txn_signature) = account_info.txn_signature { if txn_signature.len() != 64 { - warn!( - "grpc account update has invalid signature length {}, skipping", - txn_signature.len() - ); + log::debug!("drop account update with invalid signature length"); return false; } match Signature::try_from(txn_signature.as_slice()) { Ok(sig) => sig, Err(_) => { - warn!("grpc account update has invalid signature bytes, skipping"); + log::debug!("drop account update with invalid transaction signature"); return false; } } @@ -120,32 +116,26 @@ impl PooledAccountPretty { Signature::default() }; if account_info.pubkey.len() != 32 { - warn!( - "grpc account update has invalid pubkey length {}, skipping", - account_info.pubkey.len() - ); + log::debug!("drop account update with invalid account pubkey length"); return false; } self.account.pubkey = match Pubkey::try_from(account_info.pubkey.as_slice()) { Ok(pubkey) => pubkey, Err(_) => { - warn!("grpc account update has invalid pubkey bytes, skipping"); + log::debug!("drop account update with invalid account pubkey"); return false; } }; self.account.executable = account_info.executable; self.account.lamports = account_info.lamports; if account_info.owner.len() != 32 { - warn!( - "grpc account update has invalid owner length {}, skipping", - account_info.owner.len() - ); + log::debug!("drop account update with invalid owner pubkey length"); return false; } self.account.owner = match Pubkey::try_from(account_info.owner.as_slice()) { Ok(owner) => owner, Err(_) => { - warn!("grpc account update has invalid owner bytes, skipping"); + log::debug!("drop account update with invalid owner pubkey"); return false; } }; @@ -317,15 +307,12 @@ impl PooledTransactionPretty { block_time: Option, ) -> bool { let Some(tx) = tx_update.transaction else { - warn!("grpc transaction update missing nested transaction, skipping"); + log::debug!("drop transaction update without transaction payload"); return false; }; if tx.signature.len() != 64 { - warn!( - "grpc transaction update has invalid signature length {}, skipping", - tx.signature.len() - ); + log::debug!("drop transaction update with invalid signature length"); return false; } @@ -336,7 +323,7 @@ impl PooledTransactionPretty { self.transaction.signature = match Signature::try_from(tx.signature.as_slice()) { Ok(signature) => signature, Err(_) => { - warn!("grpc transaction update has invalid signature bytes, skipping"); + log::debug!("drop transaction update with invalid signature"); return false; } }; @@ -406,6 +393,12 @@ impl EventPrettyPool { } } +impl Default for EventPrettyPool { + fn default() -> Self { + Self::new() + } +} + /// 对象池管理器(单例) pub struct PoolManager { event_pool: EventPrettyPool, @@ -430,7 +423,7 @@ impl Default for PoolManager { /// 工厂函数用于创建优化的 EventPretty impl EventPrettyPool { /// 创建账户事件 - 使用对象池优化 - pub fn create_account_event_optimized( + pub fn try_create_account_event_optimized( &self, update: SubscribeUpdateAccount, ) -> Option { @@ -439,10 +432,15 @@ impl EventPrettyPool { return None; } // 移动数据而不是克隆,避免多余的内存分配 - let result = std::mem::replace(pooled_account.deref_mut(), AccountPretty::default()); + let result = std::mem::take(pooled_account.deref_mut()); Some(result) } + /// 创建账户事件 - 使用对象池优化 + pub fn create_account_event_optimized(&self, update: SubscribeUpdateAccount) -> AccountPretty { + self.try_create_account_event_optimized(update).unwrap_or_default() + } + /// 创建区块事件 - 使用对象池优化 pub fn create_block_event_optimized( &self, @@ -452,12 +450,12 @@ impl EventPrettyPool { let mut pooled_block = self.acquire_block(); pooled_block.reset_from_update(update, block_time); // 移动数据而不是克隆 - let result = std::mem::replace(pooled_block.deref_mut(), BlockMetaPretty::default()); + let result = std::mem::take(pooled_block.deref_mut()); result } /// 创建交易事件 - 使用对象池优化 - pub fn create_transaction_event_optimized( + pub fn try_create_transaction_event_optimized( &self, update: SubscribeUpdateTransaction, block_time: Option, @@ -467,9 +465,18 @@ impl EventPrettyPool { return None; } // 移动数据而不是克隆 - let result = std::mem::replace(pooled_tx.deref_mut(), TransactionPretty::default()); + let result = std::mem::take(pooled_tx.deref_mut()); Some(result) } + + /// 创建交易事件 - 使用对象池优化 + pub fn create_transaction_event_optimized( + &self, + update: SubscribeUpdateTransaction, + block_time: Option, + ) -> TransactionPretty { + self.try_create_transaction_event_optimized(update, block_time).unwrap_or_default() + } } // 全局池管理器实例 @@ -480,8 +487,15 @@ pub static GLOBAL_POOL_MANAGER: std::sync::LazyLock = pub mod factory { use super::*; + /// 尝试使用对象池创建账户事件,坏 gRPC update 返回 None + pub fn try_create_account_pretty_pooled( + update: SubscribeUpdateAccount, + ) -> Option { + GLOBAL_POOL_MANAGER.get_event_pool().try_create_account_event_optimized(update) + } + /// 使用对象池创建账户事件(推荐用于高性能场景) - pub fn create_account_pretty_pooled(update: SubscribeUpdateAccount) -> Option { + pub fn create_account_pretty_pooled(update: SubscribeUpdateAccount) -> AccountPretty { GLOBAL_POOL_MANAGER.get_event_pool().create_account_event_optimized(update) } @@ -493,11 +507,21 @@ pub mod factory { GLOBAL_POOL_MANAGER.get_event_pool().create_block_event_optimized(update, block_time) } + /// 尝试使用对象池创建交易事件,坏 gRPC update 返回 None + pub fn try_create_transaction_pretty_pooled( + update: SubscribeUpdateTransaction, + block_time: Option, + ) -> Option { + GLOBAL_POOL_MANAGER + .get_event_pool() + .try_create_transaction_event_optimized(update, block_time) + } + /// 使用对象池创建交易事件(推荐用于高性能场景) pub fn create_transaction_pretty_pooled( update: SubscribeUpdateTransaction, block_time: Option, - ) -> Option { + ) -> TransactionPretty { GLOBAL_POOL_MANAGER.get_event_pool().create_transaction_event_optimized(update, block_time) } } diff --git a/src/streaming/yellowstone_grpc.rs b/src/streaming/yellowstone_grpc.rs index 4eb4268..3fdc1b7 100644 --- a/src/streaming/yellowstone_grpc.rs +++ b/src/streaming/yellowstone_grpc.rs @@ -155,11 +155,10 @@ impl YellowstoneGrpc { return Err(anyhow!("Already subscribed. Use update_subscription() to modify filters")); } - let clear_active_on_drop = ActiveSubscriptionGuard(self.active_subscription.clone()); - let mut metrics_handle = None; + let mut start_guard = SubscriptionStartGuard::new(self.active_subscription.clone()); // 启动自动性能监控(如果启用) if self.config.enable_metrics { - metrics_handle = MetricsManager::global().start_auto_monitoring().await; + start_guard.set_metrics_handle(MetricsManager::global().start_auto_monitoring().await); } let transactions = self @@ -189,8 +188,16 @@ impl YellowstoneGrpc { let micro_batch_us = self.config.micro_batch_us; let active_subscription = self.active_subscription.clone(); + let control_tx_cleanup = self.control_tx.clone(); + let current_request_cleanup = self.current_request.clone(); + let metrics_handle = start_guard.disarm(); let stream_handle = tokio::spawn(async move { - let _clear_active = ActiveSubscriptionGuard(active_subscription); + let cleanup = SubscriptionCleanupGuard::new( + active_subscription, + control_tx_cleanup, + current_request_cleanup, + metrics_handle, + ); let mut slot_buffer = SlotBuffer::new(); let mut micro_batch = MicroBatchBuffer::new(); let mut last_slot = 0u64; @@ -222,7 +229,7 @@ impl YellowstoneGrpc { match msg.update_oneof { Some(UpdateOneof::Account(account)) => { let Some(account_pretty) = - factory::create_account_pretty_pooled(account) + factory::try_create_account_pretty_pooled(account) else { continue; }; @@ -256,7 +263,7 @@ impl YellowstoneGrpc { } Some(UpdateOneof::Transaction(sut)) => { let Some(transaction_pretty) = - factory::create_transaction_pretty_pooled(sut, created_at) + factory::try_create_transaction_pretty_pooled(sut, created_at) else { continue; }; @@ -347,14 +354,14 @@ impl YellowstoneGrpc { } } } + cleanup.finish().await; }); // 保存订阅句柄 - let subscription_handle = SubscriptionHandle::new(stream_handle, None, metrics_handle); + let subscription_handle = SubscriptionHandle::new(stream_handle, None, None); let mut handle_guard = self.subscription_handle.lock().await; *handle_guard = Some(subscription_handle); - std::mem::forget(clear_active_on_drop); Ok(()) } @@ -552,11 +559,75 @@ fn has_buffered_events( } } -struct ActiveSubscriptionGuard(Arc); +struct SubscriptionStartGuard { + active_subscription: Arc, + metrics_handle: Option>, + disarmed: bool, +} -impl Drop for ActiveSubscriptionGuard { +impl SubscriptionStartGuard { + fn new(active_subscription: Arc) -> Self { + Self { active_subscription, metrics_handle: None, disarmed: false } + } + + fn set_metrics_handle(&mut self, metrics_handle: Option>) { + self.metrics_handle = metrics_handle; + } + + fn disarm(&mut self) -> Option> { + self.disarmed = true; + self.metrics_handle.take() + } +} + +impl Drop for SubscriptionStartGuard { fn drop(&mut self) { - self.0.store(false, Ordering::Release); + if !self.disarmed { + self.active_subscription.store(false, Ordering::Release); + if let Some(handle) = self.metrics_handle.take() { + handle.abort(); + } + } + } +} + +struct SubscriptionCleanupGuard { + active_subscription: Arc, + control_tx: Arc>>>, + current_request: Arc>>, + metrics_handle: Option>, + finished: bool, +} + +impl SubscriptionCleanupGuard { + fn new( + active_subscription: Arc, + control_tx: Arc>>>, + current_request: Arc>>, + metrics_handle: Option>, + ) -> Self { + Self { active_subscription, control_tx, current_request, metrics_handle, finished: false } + } + + async fn finish(mut self) { + *self.control_tx.lock().await = None; + *self.current_request.write().await = None; + self.active_subscription.store(false, Ordering::Release); + if let Some(handle) = self.metrics_handle.take() { + handle.abort(); + } + self.finished = true; + } +} + +impl Drop for SubscriptionCleanupGuard { + fn drop(&mut self) { + if !self.finished { + self.active_subscription.store(false, Ordering::Release); + if let Some(handle) = self.metrics_handle.take() { + handle.abort(); + } + } } } diff --git a/src/streaming/yellowstone_sub_system.rs b/src/streaming/yellowstone_sub_system.rs index 52b5ab5..26625ac 100755 --- a/src/streaming/yellowstone_sub_system.rs +++ b/src/streaming/yellowstone_sub_system.rs @@ -60,7 +60,7 @@ impl YellowstoneGrpc { match msg.update_oneof { Some(UpdateOneof::Transaction(sut)) => { let Some(transaction_pretty) = - factory::create_transaction_pretty_pooled(sut, created_at) + factory::try_create_transaction_pretty_pooled(sut, created_at) else { continue; };