Release solana-streamer 1.5.2

This commit is contained in:
0xfnzero
2026-05-26 21:23:27 +08:00
parent 8b967f58d6
commit 02c8a6d168
8 changed files with 245 additions and 87 deletions
+154 -18
View File
@@ -15,9 +15,17 @@ use crate::streaming::parser_sdk_bridge::{
use crate::streaming::shred::TransactionWithSlot;
use sol_parser_sdk::grpc::parse_subscribe_update_transaction_low_latency;
use solana_sdk::pubkey::Pubkey;
use solana_sdk::signature::Signature;
use solana_sdk::transaction::VersionedTransaction;
use std::cell::RefCell;
use std::sync::Arc;
use yellowstone_grpc_proto::geyser::SubscribeUpdateTransaction;
thread_local! {
static SHRED_SDK_EVENTS: RefCell<Vec<sol_parser_sdk::DexEvent>> =
RefCell::new(Vec::with_capacity(4));
}
/// Wrap the user callback and update transaction metrics after delivery.
#[inline]
fn create_metrics_callback(
@@ -182,33 +190,68 @@ pub async fn process_shred_transaction(
let recv_us = transaction_with_slot.recv_us;
let adapter_callback = create_metrics_callback(callback);
let sdk_parse_filter = build_sdk_shred_parse_event_filter(protocols, event_type_filter);
let mut sdk_events = Vec::with_capacity(4);
sol_parser_sdk::shredstream::parse_transaction_dex_events_with_filter(
parse_shred_transaction_events(
&tx,
signature,
slot,
tx_index.unwrap_or(0),
tx_index,
recv_us,
sdk_parse_filter.as_ref(),
&mut sdk_events,
protocols,
event_type_filter,
bot_wallet,
|event| adapter_callback(event),
);
for sdk_event in sdk_events {
if let Some(mut event) =
adapt_parser_event(sdk_event, None, recv_us, protocols, event_type_filter)
{
event.metadata_mut().handle_us = elapsed_micros_since(recv_us);
event = crate::streaming::event_parser::core::event_parser::helpers::process_event(
event, bot_wallet,
);
adapter_callback(event);
}
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn parse_shred_transaction_events(
tx: &VersionedTransaction,
signature: Signature,
slot: u64,
tx_index: Option<u64>,
recv_us: i64,
protocols: &[Protocol],
event_type_filter: Option<&EventTypeFilter>,
bot_wallet: Option<Pubkey>,
mut on_event: impl FnMut(DexEvent),
) {
let sdk_parse_filter = build_sdk_shred_parse_event_filter(protocols, event_type_filter);
SHRED_SDK_EVENTS.with(|slot_events| {
let mut events = {
let mut slot_events = slot_events.borrow_mut();
std::mem::take(&mut *slot_events)
};
events.clear();
sol_parser_sdk::shredstream::parse_transaction_dex_events_with_filter(
tx,
signature,
slot,
tx_index.unwrap_or(0),
recv_us,
sdk_parse_filter.as_ref(),
&mut events,
);
for sdk_event in events.drain(..) {
if let Some(mut event) =
adapt_parser_event(sdk_event, None, recv_us, protocols, event_type_filter)
{
event.metadata_mut().tx_index = tx_index;
event.metadata_mut().handle_us = elapsed_micros_since(recv_us);
event = crate::streaming::event_parser::core::event_parser::helpers::process_event(
event, bot_wallet,
);
on_event(event);
}
}
*slot_events.borrow_mut() = events;
});
}
/// Update metrics for event processing (with optional latency check)
#[inline]
fn update_metrics(ty: MetricsEventType, count: u64, time_us: f64) {
@@ -232,3 +275,96 @@ fn update_metrics_with_latency(
block_time_ms,
);
}
#[cfg(test)]
mod tests {
use super::process_shred_transaction;
use crate::streaming::event_parser::{DexEvent, Protocol};
use crate::streaming::shred::TransactionWithSlot;
use sol_parser_sdk::instr::program_ids::PUMPFUN_PROGRAM_ID;
use solana_sdk::hash::Hash;
use solana_sdk::message::{
compiled_instruction::CompiledInstruction, v0, MessageHeader, VersionedMessage,
};
use solana_sdk::pubkey::Pubkey;
use solana_sdk::signature::Signature;
use solana_sdk::transaction::VersionedTransaction;
use std::sync::{Arc, Mutex};
fn push_string(data: &mut Vec<u8>, value: &str) {
data.extend_from_slice(&(value.len() as u32).to_le_bytes());
data.extend_from_slice(value.as_bytes());
}
fn pumpfun_create_data() -> Vec<u8> {
let mut data = Vec::new();
data.extend_from_slice(&[24, 30, 200, 40, 5, 28, 7, 119]);
push_string(&mut data, "Index Test");
push_string(&mut data, "IDX");
push_string(&mut data, "https://example.invalid/index.json");
data.extend_from_slice(Pubkey::new_unique().as_ref());
data
}
fn pumpfun_create_tx() -> VersionedTransaction {
let mut account_keys = (0..10).map(|_| Pubkey::new_unique()).collect::<Vec<_>>();
account_keys.push(PUMPFUN_PROGRAM_ID);
VersionedTransaction {
signatures: vec![Signature::default()],
message: VersionedMessage::V0(v0::Message {
header: MessageHeader {
num_required_signatures: 1,
num_readonly_signed_accounts: 0,
num_readonly_unsigned_accounts: 0,
},
account_keys,
recent_blockhash: Hash::default(),
instructions: vec![CompiledInstruction::new_from_raw_parts(
10,
pumpfun_create_data(),
(0..10).collect(),
)],
address_table_lookups: Vec::new(),
}),
}
}
async fn parse_single_shred_create(tx_index: Option<u64>) -> DexEvent {
let events = Arc::new(Mutex::new(Vec::new()));
let captured = events.clone();
let callback = Arc::new(move |event| {
captured.lock().unwrap().push(event);
});
process_shred_transaction(
TransactionWithSlot::new(pumpfun_create_tx(), 42, 1_000_000, tx_index),
&[Protocol::PumpFun],
None,
callback,
None,
)
.await
.expect("process shred transaction");
let mut events = events.lock().unwrap();
assert_eq!(events.len(), 1);
events.pop().unwrap()
}
#[tokio::test]
async fn shred_transaction_preserves_unknown_tx_index() {
let event = parse_single_shred_create(None).await;
assert!(matches!(event, DexEvent::PumpFunCreateTokenEvent(_)));
assert_eq!(event.metadata().tx_index, None);
}
#[tokio::test]
async fn shred_transaction_preserves_present_tx_index() {
let event = parse_single_shred_create(Some(42)).await;
assert!(matches!(event, DexEvent::PumpFunCreateTokenEvent(_)));
assert_eq!(event.metadata().tx_index, Some(42));
}
}
+12 -4
View File
@@ -2,7 +2,7 @@ use tokio::task::JoinHandle;
/// Subscription handle for managing and stopping subscriptions
pub struct SubscriptionHandle {
stream_handle: JoinHandle<()>,
stream_handle: Option<JoinHandle<()>>,
event_handle: Option<JoinHandle<()>>,
metrics_handle: Option<JoinHandle<()>>,
}
@@ -14,12 +14,18 @@ impl SubscriptionHandle {
event_handle: Option<JoinHandle<()>>,
metrics_handle: Option<JoinHandle<()>>,
) -> Self {
Self { stream_handle, event_handle, metrics_handle }
Self { stream_handle: Some(stream_handle), event_handle, metrics_handle }
}
pub fn metrics_only(metrics_handle: Option<JoinHandle<()>>) -> Self {
Self { stream_handle: None, event_handle: None, metrics_handle }
}
/// Stop subscription and abort all related tasks
pub fn stop(self) {
self.stream_handle.abort();
if let Some(handle) = self.stream_handle {
handle.abort();
}
if let Some(handle) = self.event_handle {
handle.abort();
}
@@ -30,7 +36,9 @@ impl SubscriptionHandle {
/// Asynchronously wait for all tasks to complete
pub async fn join(self) -> Result<(), tokio::task::JoinError> {
let _ = self.stream_handle.await;
if let Some(handle) = self.stream_handle {
let _ = handle.await;
}
if let Some(handle) = self.event_handle {
let _ = handle.await;
}
+30 -2
View File
@@ -173,7 +173,10 @@ fn push_streamer_event_sdk_grpc_types(t: &EventType, out: &mut Vec<SdkGrpcEventT
use SdkGrpcEventType as Sdk;
match t {
St::BlockMeta => out.push(Sdk::BlockMeta),
St::PumpFunCreateToken => out.push(Sdk::PumpFunCreate),
St::PumpFunCreateToken => {
out.push(Sdk::PumpFunCreate);
out.push(Sdk::PumpFunCreateV2);
}
St::PumpFunCreateV2Token => out.push(Sdk::PumpFunCreateV2),
St::PumpFunBuy => {
out.push(Sdk::PumpFunBuy);
@@ -396,7 +399,8 @@ fn event_type_matches(filter_type: &EventType, event_type: &EventType) -> bool {
filter_type == event_type
|| matches!(
(filter_type, event_type),
(EventType::PumpFunBuy, EventType::PumpFunBuyExactSolIn)
(EventType::PumpFunCreateToken, EventType::PumpFunCreateV2Token)
| (EventType::PumpFunBuy, EventType::PumpFunBuyExactSolIn)
| (EventType::TokenAccount, EventType::TokenInfo)
| (EventType::RaydiumClmmSwap, EventType::RaydiumClmmSwapV2)
| (EventType::RaydiumClmmSwapV2, EventType::RaydiumClmmSwap)
@@ -531,6 +535,30 @@ mod tests {
assert!(!sdk_f.should_include(SdkGrpcEventType::PumpFunSell));
}
#[test]
fn pumpfun_create_filter_matches_create_v2_for_backward_compat() {
let f =
EventTypeFilter { include: vec![EventType::PumpFunCreateToken], ..Default::default() };
assert!(f.passes_event_type(&EventType::PumpFunCreateToken));
assert!(f.passes_event_type(&EventType::PumpFunCreateV2Token));
let sdk_f = build_sdk_parse_event_filter(Some(&f)).expect("mapped");
assert!(sdk_f.should_include(SdkGrpcEventType::PumpFunCreate));
assert!(sdk_f.should_include(SdkGrpcEventType::PumpFunCreateV2));
assert!(!sdk_f.should_include(SdkGrpcEventType::PumpFunBuy));
let f = EventTypeFilter {
include: vec![EventType::PumpFunCreateV2Token],
..Default::default()
};
assert!(!f.passes_event_type(&EventType::PumpFunCreateToken));
assert!(f.passes_event_type(&EventType::PumpFunCreateV2Token));
let sdk_f = build_sdk_parse_event_filter(Some(&f)).expect("mapped");
assert!(!sdk_f.should_include(SdkGrpcEventType::PumpFunCreate));
assert!(sdk_f.should_include(SdkGrpcEventType::PumpFunCreateV2));
}
#[test]
fn pumpfun_buy_filter_matches_exact_sol_in_for_backward_compat() {
let f = EventTypeFilter { include: vec![EventType::PumpFunBuy], ..Default::default() };
@@ -78,37 +78,19 @@ impl EventParser {
tx_index: Option<u64>,
callback: std::sync::Arc<dyn Fn(crate::streaming::event_parser::DexEvent) + Send + Sync>,
) -> anyhow::Result<()> {
let sdk_parse_filter =
crate::streaming::event_parser::common::filter::build_sdk_shred_parse_event_filter(
protocols,
event_type_filter,
);
let mut sdk_events = Vec::with_capacity(4);
sol_parser_sdk::shredstream::parse_transaction_dex_events_with_filter(
crate::streaming::common::parse_shred_transaction_events(
transaction,
signature,
slot.unwrap_or(0),
tx_index.unwrap_or(0),
tx_index,
recv_us,
sdk_parse_filter.as_ref(),
&mut sdk_events,
);
for sdk_event in sdk_events {
if let Some(mut event) = crate::streaming::parser_sdk_bridge::adapt_parser_event(
sdk_event,
None,
recv_us,
protocols,
event_type_filter,
) {
event.metadata_mut().handle_us =
crate::streaming::event_parser::common::high_performance_clock::elapsed_micros_since(
recv_us,
);
event = helpers::process_event(event, bot_wallet);
protocols,
event_type_filter,
bot_wallet,
|event| {
callback(event);
}
}
},
);
Ok(())
}
}
+20 -24
View File
@@ -17,6 +17,9 @@ use super::ShredStreamGrpc;
impl ShredStreamGrpc {
/// 订阅ShredStream事件(支持批处理和即时处理)
///
/// Uses the SDK direct-callback path for minimum latency. User callbacks should avoid
/// blocking work; use the SDK queue API directly when heavy asynchronous processing is needed.
pub async fn shredstream_subscribe<F>(
&self,
protocols: Vec<Protocol>,
@@ -30,38 +33,25 @@ impl ShredStreamGrpc {
// 如果已有活跃订阅,先停止它
self.stop().await;
let mut metrics_handle = None;
// 启动自动性能监控(如果启用)
if self.config.enable_metrics {
metrics_handle = MetricsManager::global().start_auto_monitoring().await;
}
let sdk_parse_filter =
build_sdk_shred_parse_event_filter(&protocols, event_type_filter.as_ref());
let queue = self
.sdk_shredstream_client
.subscribe_with_filter(sdk_parse_filter)
.await
.map_err(|e| anyhow::anyhow!(e.to_string()))?;
// Wrap callback once before the async block
let callback = Arc::new(callback);
let callback_for_sdk = callback.clone();
let protocols_for_sdk = protocols;
let filter_for_sdk = event_type_filter;
let stream_task = tokio::spawn(async move {
loop {
let Some(sdk_event) = queue.pop() else {
tokio::task::yield_now().await;
continue;
};
self.sdk_shredstream_client
.subscribe_with_filter_callback(sdk_parse_filter, move |sdk_event| {
MetricsManager::global().add_tx_process_count();
let recv_wall_us = sdk_event.metadata().grpc_recv_us;
if let Some(mut event) = adapt_parser_event(
sdk_event,
None,
recv_wall_us,
&protocols,
event_type_filter.as_ref(),
&protocols_for_sdk,
filter_for_sdk.as_ref(),
) {
event.metadata_mut().handle_us = elapsed_micros_since(event.metadata().recv_us);
event =
@@ -73,7 +63,7 @@ impl ShredStreamGrpc {
let recv_us = metadata.recv_us;
let block_time_ms = metadata.block_time_ms;
callback(event);
callback_for_sdk(event);
MetricsManager::global().update_metrics_with_latency(
MetricsEventType::Transaction,
@@ -83,11 +73,17 @@ impl ShredStreamGrpc {
block_time_ms,
);
}
}
});
})
.await
.map_err(|e| anyhow::anyhow!(e.to_string()))?;
let mut metrics_handle = None;
if self.config.enable_metrics {
metrics_handle = MetricsManager::global().start_auto_monitoring().await;
}
// 保存订阅句柄
let subscription_handle = SubscriptionHandle::new(stream_task, None, metrics_handle);
let subscription_handle = SubscriptionHandle::metrics_only(metrics_handle);
let mut handle_guard = self.subscription_handle.lock().await;
*handle_guard = Some(subscription_handle);