Release solana-streamer-sdk v1.5.3

This commit is contained in:
0xfnzero
2026-05-28 03:44:42 +08:00
parent 02c8a6d168
commit 4cb22592e8
7 changed files with 201 additions and 76 deletions
+40 -9
View File
@@ -110,7 +110,7 @@ pub(crate) fn build_sdk_parse_event_filter(
if !f.exclude.is_empty() && f.include.is_empty() {
let mut raw: Vec<SdkGrpcEventType> = Vec::with_capacity(f.exclude.len());
for et in &f.exclude {
push_streamer_event_sdk_grpc_types(et, &mut raw);
push_streamer_event_sdk_grpc_types(et, &mut raw, FilterMapMode::Exclude);
}
dedup_sdk_grpc_event_types(&mut raw);
return (!raw.is_empty()).then(|| SdkGrpcEventTypeFilter::exclude_types(raw));
@@ -121,7 +121,7 @@ pub(crate) fn build_sdk_parse_event_filter(
}
let mut raw: Vec<SdkGrpcEventType> = Vec::with_capacity(f.include.len());
for et in &f.include {
if !push_streamer_event_sdk_grpc_types(et, &mut raw) {
if !push_streamer_event_sdk_grpc_types(et, &mut raw, FilterMapMode::Include) {
return None;
}
}
@@ -168,7 +168,17 @@ fn dedup_sdk_grpc_event_types(v: &mut Vec<SdkGrpcEventType>) {
}
}
fn push_streamer_event_sdk_grpc_types(t: &EventType, out: &mut Vec<SdkGrpcEventType>) -> bool {
#[derive(Clone, Copy)]
enum FilterMapMode {
Include,
Exclude,
}
fn push_streamer_event_sdk_grpc_types(
t: &EventType,
out: &mut Vec<SdkGrpcEventType>,
mode: FilterMapMode,
) -> bool {
use EventType as St;
use SdkGrpcEventType as Sdk;
match t {
@@ -179,11 +189,24 @@ fn push_streamer_event_sdk_grpc_types(t: &EventType, out: &mut Vec<SdkGrpcEventT
}
St::PumpFunCreateV2Token => out.push(Sdk::PumpFunCreateV2),
St::PumpFunBuy => {
if matches!(mode, FilterMapMode::Include) {
out.push(Sdk::PumpFunTrade);
}
out.push(Sdk::PumpFunBuy);
out.push(Sdk::PumpFunBuyExactSolIn);
}
St::PumpFunBuyExactSolIn => out.push(Sdk::PumpFunBuyExactSolIn),
St::PumpFunSell => out.push(Sdk::PumpFunSell),
St::PumpFunBuyExactSolIn => {
if matches!(mode, FilterMapMode::Include) {
out.push(Sdk::PumpFunTrade);
}
out.push(Sdk::PumpFunBuyExactSolIn);
}
St::PumpFunSell => {
if matches!(mode, FilterMapMode::Include) {
out.push(Sdk::PumpFunTrade);
}
out.push(Sdk::PumpFunSell);
}
St::PumpFunMigrate => out.push(Sdk::PumpFunMigrate),
St::PumpFeesCreateFeeSharingConfig => out.push(Sdk::PumpFeesCreateFeeSharingConfig),
St::PumpFeesInitializeFeeConfig => out.push(Sdk::PumpFeesInitializeFeeConfig),
@@ -401,6 +424,7 @@ fn event_type_matches(filter_type: &EventType, event_type: &EventType) -> bool {
(filter_type, event_type),
(EventType::PumpFunCreateToken, EventType::PumpFunCreateV2Token)
| (EventType::PumpFunBuy, EventType::PumpFunBuyExactSolIn)
| (EventType::PumpFunBuyExactSolIn, EventType::PumpFunBuy)
| (EventType::TokenAccount, EventType::TokenInfo)
| (EventType::RaydiumClmmSwap, EventType::RaydiumClmmSwapV2)
| (EventType::RaydiumClmmSwapV2, EventType::RaydiumClmmSwap)
@@ -530,9 +554,12 @@ mod tests {
fn build_sdk_filter_pumpfun_buy_only() {
let f = EventTypeFilter { include: vec![EventType::PumpFunBuy], ..Default::default() };
let sdk_f = build_sdk_parse_event_filter(Some(&f)).expect("mapped");
assert!(sdk_f.should_include(SdkGrpcEventType::PumpFunTrade));
assert!(sdk_f.should_include(SdkGrpcEventType::PumpFunBuy));
assert!(sdk_f.should_include(SdkGrpcEventType::PumpFunBuyExactSolIn));
assert!(!sdk_f.should_include(SdkGrpcEventType::PumpFunSell));
assert!(sdk_f.should_include(SdkGrpcEventType::PumpFunSell));
assert!(f.passes_event_type(&EventType::PumpFunBuy));
assert!(!f.passes_event_type(&EventType::PumpFunSell));
}
#[test]
@@ -569,7 +596,7 @@ mod tests {
include: vec![EventType::PumpFunBuyExactSolIn],
..Default::default()
};
assert!(!f.passes_event_type(&EventType::PumpFunBuy));
assert!(f.passes_event_type(&EventType::PumpFunBuy));
assert!(f.passes_event_type(&EventType::PumpFunBuyExactSolIn));
}
@@ -602,9 +629,13 @@ mod tests {
..Default::default()
};
let sdk_f = build_sdk_parse_event_filter(Some(&f)).expect("mapped");
assert!(!sdk_f.should_include(SdkGrpcEventType::PumpFunBuy));
assert!(sdk_f.should_include(SdkGrpcEventType::PumpFunBuy));
assert!(sdk_f.should_include(SdkGrpcEventType::PumpFunTrade));
assert!(sdk_f.should_include(SdkGrpcEventType::PumpFunBuyExactSolIn));
assert!(!sdk_f.should_include(SdkGrpcEventType::PumpFunSell));
assert!(sdk_f.should_include(SdkGrpcEventType::PumpFunSell));
assert!(f.passes_event_type(&EventType::PumpFunBuy));
assert!(!f.passes_event_type(&EventType::PumpFunSell));
assert!(f.passes_event_type(&EventType::PumpFunBuyExactSolIn));
}
#[test]
+3 -3
View File
@@ -142,17 +142,17 @@ mod tests {
}
#[test]
fn converts_pumpfun_buy_exact_quote_in_as_buy_event() {
fn converts_pumpfun_buy_exact_quote_in_v2_as_buy_event() {
let mut t = PbPumpTrade::default();
t.metadata = EventMetadata::default();
t.is_buy = true;
t.ix_name = "buy_exact_quote_in".to_string();
t.ix_name = "buy_exact_quote_in_v2".to_string();
let ev = convert_parser_event(PbDexEvent::PumpFunBuy(t), None, 999).expect("convert");
match ev {
DexEvent::PumpFunTradeEvent(st) => {
assert_eq!(st.metadata.event_type, EventType::PumpFunBuy);
assert_eq!(st.ix_name, "buy_exact_quote_in");
assert_eq!(st.ix_name, "buy_exact_quote_in_v2");
}
_ => panic!("expected PumpFunTradeEvent"),
}
+3 -1
View File
@@ -25,7 +25,9 @@ impl ShredStreamGrpc {
/// 创建客户端,使用自定义配置
pub async fn new_with_config(endpoint: String, config: StreamClientConfig) -> AnyResult<Self> {
let shredstream_client = ShredstreamProxyClient::connect(endpoint.clone()).await?;
let shredstream_client = ShredstreamProxyClient::connect(endpoint.clone())
.await?
.max_decoding_message_size(config.connection.max_decoding_message_size);
let sdk_config = sol_parser_sdk::shredstream::ShredStreamConfig {
connection_timeout_ms: config.connection.connect_timeout.saturating_mul(1000),
request_timeout_ms: config.connection.request_timeout.saturating_mul(1000),
+132 -48
View File
@@ -1,25 +1,30 @@
//! ShredStream 订阅入口:底层订阅与热路径解析直接复用 `sol-parser-sdk::shredstream`
//! 本模块负责把 SDK `DexEvent` 适配回 streamer 原有 callback API。
use std::sync::Arc;
//! 本模块负责把 SDK `DexEvent` 适配回 streamer 原有 callback API。
#![allow(deprecated)]
use std::sync::Arc;
use std::time::Duration;
use futures::StreamExt;
use solana_entry::entry::Entry as SolanaEntry;
use solana_sdk::pubkey::Pubkey;
use tokio::time::timeout;
use crate::common::AnyResult;
use crate::streaming::common::parse_shred_transaction_events;
use crate::streaming::common::{MetricsEventType, MetricsManager, SubscriptionHandle};
use crate::streaming::event_parser::common::filter::{
build_sdk_shred_parse_event_filter, EventTypeFilter,
};
use crate::streaming::event_parser::common::high_performance_clock::elapsed_micros_since;
use crate::streaming::event_parser::common::filter::EventTypeFilter;
use crate::streaming::event_parser::common::high_performance_clock::get_high_perf_clock;
use crate::streaming::event_parser::{DexEvent, Protocol};
use crate::streaming::parser_sdk_bridge::adapt_parser_event;
use super::ShredStreamGrpc;
impl ShredStreamGrpc {
/// 订阅ShredStream事件(支持批处理和即时处理)
/// Subscribe to ShredStream events on the direct low-latency path.
///
/// 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.
/// The callback runs on the stream read task; keep it non-blocking. `tx_index` is an
/// entry-local best-effort index because ShredStream entries do not carry the slot-level
/// Yellowstone transaction index.
pub async fn shredstream_subscribe<F>(
&self,
protocols: Vec<Protocol>,
@@ -33,37 +38,128 @@ impl ShredStreamGrpc {
// 如果已有活跃订阅,先停止它
self.stop().await;
let sdk_parse_filter =
build_sdk_shred_parse_event_filter(&protocols, event_type_filter.as_ref());
// Wrap callback once before the async block
let request_timeout = Duration::from_secs(self.config.connection.request_timeout);
let client = self.shredstream_client.as_ref().clone();
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 {
let mut delay_ms = 100u64;
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_for_sdk,
filter_for_sdk.as_ref(),
) {
event.metadata_mut().handle_us = elapsed_micros_since(event.metadata().recv_us);
event =
crate::streaming::event_parser::core::event_parser::helpers::process_event(
event, bot_wallet,
loop {
let mut client = client.clone();
let request =
tonic::Request::new(crate::protos::shredstream::SubscribeEntriesRequest {});
let response = if request_timeout.is_zero() {
client.subscribe_entries(request).await
} else {
match timeout(request_timeout, client.subscribe_entries(request)).await {
Ok(response) => response,
Err(_) => {
log::error!(
"ShredStream subscribe request timed out after {}s - retry in {}ms",
request_timeout.as_secs(),
delay_ms
);
tokio::time::sleep(Duration::from_millis(delay_ms)).await;
delay_ms = (delay_ms * 2).min(60_000);
continue;
}
}
};
let mut stream = match response {
Ok(response) => response.into_inner(),
Err(error) => {
log::error!(
"ShredStream subscribe failed: {error} - retry in {delay_ms}ms"
);
tokio::time::sleep(Duration::from_millis(delay_ms)).await;
delay_ms = (delay_ms * 2).min(60_000);
continue;
}
};
while let Some(message) = stream.next().await {
let entry = match message {
Ok(entry) => entry,
Err(error) => {
log::error!("ShredStream stream error: {error} - reconnecting");
break;
}
};
delay_ms = 100;
process_shred_entry_direct(
entry,
&protocols,
event_type_filter.as_ref(),
bot_wallet,
callback.as_ref(),
);
}
log::warn!("ShredStream stream ended - reconnecting in {delay_ms}ms");
tokio::time::sleep(Duration::from_millis(delay_ms)).await;
delay_ms = (delay_ms * 2).min(60_000);
}
});
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 mut handle_guard = self.subscription_handle.lock().await;
*handle_guard = Some(subscription_handle);
Ok(())
}
}
#[inline]
fn process_shred_entry_direct(
entry: crate::protos::shredstream::Entry,
protocols: &[Protocol],
event_type_filter: Option<&EventTypeFilter>,
bot_wallet: Option<Pubkey>,
callback: &(dyn Fn(DexEvent) + Send + Sync),
) {
let slot = entry.slot;
let recv_us = get_high_perf_clock();
let entries = match bincode::deserialize::<Vec<SolanaEntry>>(&entry.entries) {
Ok(entries) => entries,
Err(error) => {
log::debug!("Failed to deserialize ShredStream entries: {error}");
return;
}
};
let mut tx_index = 0u64;
for entry in entries {
for transaction in &entry.transactions {
if transaction.signatures.is_empty() {
tx_index += 1;
continue;
}
MetricsManager::global().add_tx_process_count();
let signature = transaction.signatures[0];
parse_shred_transaction_events(
transaction,
signature,
slot,
Some(tx_index),
recv_us,
protocols,
event_type_filter,
bot_wallet,
|event| {
let metadata = event.metadata();
let processing_time_us = metadata.handle_us as f64;
let recv_us = metadata.recv_us;
let block_time_ms = metadata.block_time_ms;
callback_for_sdk(event);
callback(event);
MetricsManager::global().update_metrics_with_latency(
MetricsEventType::Transaction,
@@ -72,21 +168,9 @@ impl ShredStreamGrpc {
recv_us,
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;
},
);
tx_index += 1;
}
// 保存订阅句柄
let subscription_handle = SubscriptionHandle::metrics_only(metrics_handle);
let mut handle_guard = self.subscription_handle.lock().await;
*handle_guard = Some(subscription_handle);
Ok(())
}
}