mirror of
https://github.com/0xfnzero/solana-streamer.git
synced 2026-08-15 18:08:05 +00:00
refactor: modernize dependency stack and minimize footprint
- Migrate to native 'std::sync' primitives (LazyLock, OnceLock, RwLock). - Replace 'parking_lot' with std library synchronization types. - Switch from 'chrono' to 'std::time::SystemTime' for lightweight timestamps. - Bump core dependencies: tokio (1.50.0) and rustls (0.23.37). - Decouple 'solana-program' by utilizing 'solana-sdk' re-exports. - Prune redundant crates: lazy_static, once_cell, chrono, maplit, and parking_lot. - Remove unused dev-dependencies (criterion) and example-only utils (env_logger). - Improve compile times and reduce supply chain attack surface.
This commit is contained in:
@@ -91,7 +91,12 @@ pub async fn process_grpc_transaction(
|
||||
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());
|
||||
.unwrap_or_else(|| {
|
||||
std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap()
|
||||
.as_millis() as i64
|
||||
});
|
||||
|
||||
let block_meta_event = CommonEventParser::generate_block_meta_event(
|
||||
block_meta_pretty.slot,
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
use std::fmt::Debug;
|
||||
use std::time::Instant;
|
||||
use std::time::{Instant, SystemTime, UNIX_EPOCH};
|
||||
|
||||
/// 高性能时钟管理器,减少系统调用开销并最小化延迟
|
||||
#[derive(Debug)]
|
||||
@@ -25,12 +25,12 @@ impl HighPerformanceClock {
|
||||
// 通过多次采样来减少初始化误差
|
||||
let mut best_offset = i64::MAX;
|
||||
let mut best_instant = Instant::now();
|
||||
let mut best_timestamp = chrono::Utc::now().timestamp_micros();
|
||||
let mut best_timestamp = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_micros() as i64;
|
||||
|
||||
// 进行3次采样,选择延迟最小的
|
||||
for _ in 0..3 {
|
||||
let instant_before = Instant::now();
|
||||
let timestamp = chrono::Utc::now().timestamp_micros();
|
||||
let timestamp = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_micros() as i64;
|
||||
let instant_after = Instant::now();
|
||||
|
||||
let sample_latency = instant_after.duration_since(instant_before).as_nanos() as i64;
|
||||
@@ -69,7 +69,7 @@ impl HighPerformanceClock {
|
||||
/// 重新校准时钟,减少累积漂移
|
||||
fn recalibrate(&mut self) {
|
||||
let current_monotonic = Instant::now();
|
||||
let current_utc = chrono::Utc::now().timestamp_micros();
|
||||
let current_utc = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_micros() as i64;
|
||||
|
||||
// 计算预期的UTC时间戳(基于单调时钟)
|
||||
let expected_utc = self.base_timestamp_us
|
||||
@@ -113,8 +113,8 @@ impl Default for HighPerformanceClock {
|
||||
}
|
||||
|
||||
/// 全局高性能时钟实例
|
||||
static HIGH_PERF_CLOCK: once_cell::sync::OnceCell<HighPerformanceClock> =
|
||||
once_cell::sync::OnceCell::new();
|
||||
static HIGH_PERF_CLOCK: std::sync::OnceLock<HighPerformanceClock> =
|
||||
std::sync::OnceLock::new();
|
||||
|
||||
/// 获取全局高性能时钟实例(最简单的实现)
|
||||
#[inline(always)]
|
||||
|
||||
@@ -36,9 +36,8 @@ impl EventMetadataPool {
|
||||
}
|
||||
|
||||
// Global object pool instances
|
||||
lazy_static::lazy_static! {
|
||||
pub static ref EVENT_METADATA_POOL: EventMetadataPool = EventMetadataPool::new();
|
||||
}
|
||||
pub static EVENT_METADATA_POOL: std::sync::LazyLock<EventMetadataPool> =
|
||||
std::sync::LazyLock::new(EventMetadataPool::new);
|
||||
|
||||
#[derive(
|
||||
Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, BorshSerialize, BorshDeserialize,
|
||||
@@ -356,14 +355,13 @@ impl EventMetadata {
|
||||
}
|
||||
}
|
||||
|
||||
lazy_static::lazy_static! {
|
||||
static ref SOL_MINT: Pubkey = Pubkey::from_str("So11111111111111111111111111111111111111111").unwrap();
|
||||
static ref SYSTEM_PROGRAMS: [Pubkey; 3] = [
|
||||
Pubkey::from_str("TokenkegQfeZyiNwAJbNbGKPFXCWuBvf9Ss623VQ5DA").unwrap(),
|
||||
Pubkey::from_str("TokenzQdBNbLqP5VEhdkAS6EPFLC1PHnBqCXEpPxuEb").unwrap(),
|
||||
Pubkey::from_str("11111111111111111111111111111111").unwrap(),
|
||||
];
|
||||
}
|
||||
static SOL_MINT: std::sync::LazyLock<Pubkey> =
|
||||
std::sync::LazyLock::new(|| Pubkey::from_str("So11111111111111111111111111111111111111111").unwrap());
|
||||
static SYSTEM_PROGRAMS: std::sync::LazyLock<[Pubkey; 3]> = std::sync::LazyLock::new(|| [
|
||||
Pubkey::from_str("TokenkegQfeZyiNwAJbNbGKPFXCWuBvf9Ss623VQ5DA").unwrap(),
|
||||
Pubkey::from_str("TokenzQdBNbLqP5VEhdkAS6EPFLC1PHnBqCXEpPxuEb").unwrap(),
|
||||
Pubkey::from_str("11111111111111111111111111111111").unwrap(),
|
||||
]);
|
||||
|
||||
/// Parse token transfer data from next instructions
|
||||
pub fn parse_swap_data_from_next_instructions(
|
||||
|
||||
@@ -179,8 +179,8 @@ impl Default for GlobalState {
|
||||
}
|
||||
|
||||
/// Global state instance
|
||||
static GLOBAL_STATE: once_cell::sync::Lazy<GlobalState> =
|
||||
once_cell::sync::Lazy::new(GlobalState::new);
|
||||
static GLOBAL_STATE: std::sync::LazyLock<GlobalState> =
|
||||
std::sync::LazyLock::new(GlobalState::new);
|
||||
|
||||
/// Get global state instance
|
||||
pub fn get_global_state() -> &'static GlobalState {
|
||||
|
||||
@@ -54,8 +54,8 @@ impl CacheKey {
|
||||
|
||||
/// 全局程序ID缓存(使用读写锁保护)
|
||||
static GLOBAL_PROGRAM_IDS_CACHE: LazyLock<
|
||||
parking_lot::RwLock<HashMap<CacheKey, Arc<Vec<Pubkey>>>>,
|
||||
> = LazyLock::new(|| parking_lot::RwLock::new(HashMap::new()));
|
||||
std::sync::RwLock<HashMap<CacheKey, Arc<Vec<Pubkey>>>>,
|
||||
> = LazyLock::new(|| std::sync::RwLock::new(HashMap::new()));
|
||||
|
||||
/// 获取指定协议的程序ID列表
|
||||
///
|
||||
@@ -68,7 +68,7 @@ pub fn get_global_program_ids(
|
||||
|
||||
// 快速路径:尝试读取缓存
|
||||
{
|
||||
let cache = GLOBAL_PROGRAM_IDS_CACHE.read();
|
||||
let cache = GLOBAL_PROGRAM_IDS_CACHE.read().unwrap();
|
||||
if let Some(program_ids) = cache.get(&cache_key) {
|
||||
return program_ids.clone();
|
||||
}
|
||||
@@ -78,7 +78,7 @@ pub fn get_global_program_ids(
|
||||
let program_ids = Arc::new(EventDispatcher::get_program_ids(protocols));
|
||||
|
||||
// 缓存结果(写锁)
|
||||
GLOBAL_PROGRAM_IDS_CACHE.write().insert(cache_key, program_ids.clone());
|
||||
GLOBAL_PROGRAM_IDS_CACHE.write().unwrap().insert(cache_key, program_ids.clone());
|
||||
|
||||
program_ids
|
||||
}
|
||||
|
||||
@@ -406,9 +406,8 @@ impl EventPrettyPool {
|
||||
}
|
||||
|
||||
// 全局池管理器实例
|
||||
lazy_static::lazy_static! {
|
||||
pub static ref GLOBAL_POOL_MANAGER: PoolManager = PoolManager::new();
|
||||
}
|
||||
pub static GLOBAL_POOL_MANAGER: std::sync::LazyLock<PoolManager> =
|
||||
std::sync::LazyLock::new(PoolManager::new);
|
||||
|
||||
/// 便捷的全局工厂函数
|
||||
pub mod factory {
|
||||
|
||||
@@ -1,5 +1,4 @@
|
||||
use futures::{channel::mpsc, sink::Sink, Stream};
|
||||
use maplit::hashmap;
|
||||
use std::{collections::HashMap, time::Duration};
|
||||
use tonic::{transport::channel::ClientTlsConfig, Status};
|
||||
use yellowstone_grpc_client::{GeyserGrpcClient, Interceptor};
|
||||
@@ -55,11 +54,11 @@ impl SubscriptionManager {
|
||||
)> {
|
||||
let blocks_meta =
|
||||
if event_type_filter.is_some() && event_type_filter.unwrap().include_block_event() {
|
||||
hashmap! { "".to_owned() => SubscribeRequestFilterBlocksMeta {} }
|
||||
HashMap::from([("".to_owned(), SubscribeRequestFilterBlocksMeta {})])
|
||||
} else if event_type_filter.is_none() {
|
||||
hashmap! { "".to_owned() => SubscribeRequestFilterBlocksMeta {} }
|
||||
HashMap::from([("".to_owned(), SubscribeRequestFilterBlocksMeta {})])
|
||||
} else {
|
||||
hashmap! {}
|
||||
HashMap::new()
|
||||
};
|
||||
let subscribe_request = SubscribeRequest {
|
||||
accounts: accounts.unwrap_or_default(),
|
||||
|
||||
@@ -137,9 +137,8 @@ impl Default for ShredPoolManager {
|
||||
}
|
||||
|
||||
// 全局 Shred 池管理器实例
|
||||
lazy_static::lazy_static! {
|
||||
pub static ref GLOBAL_SHRED_POOL_MANAGER: ShredPoolManager = ShredPoolManager::new();
|
||||
}
|
||||
pub static GLOBAL_SHRED_POOL_MANAGER: std::sync::LazyLock<ShredPoolManager> =
|
||||
std::sync::LazyLock::new(ShredPoolManager::new);
|
||||
|
||||
/// 便捷的全局工厂函数
|
||||
pub mod factory {
|
||||
|
||||
@@ -8,7 +8,7 @@ use crate::streaming::event_parser::{Protocol, DexEvent};
|
||||
use crate::streaming::grpc::pool::factory;
|
||||
use crate::streaming::grpc::{EventPretty, SubscriptionManager};
|
||||
use anyhow::anyhow;
|
||||
use chrono::Local;
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
use futures::channel::mpsc;
|
||||
use futures::{SinkExt, StreamExt};
|
||||
use log::error;
|
||||
@@ -247,10 +247,12 @@ impl YellowstoneGrpc {
|
||||
})
|
||||
.await;
|
||||
}
|
||||
log::debug!("service is ping: {}", Local::now());
|
||||
let ts = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_secs();
|
||||
log::debug!("service is ping: {}", ts);
|
||||
}
|
||||
Some(UpdateOneof::Pong(_)) => {
|
||||
log::debug!("service is pong: {}", Local::now());
|
||||
let ts = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_secs();
|
||||
log::debug!("service is pong: {}", ts);
|
||||
}
|
||||
_ => {
|
||||
log::debug!("Received other message type");
|
||||
|
||||
@@ -7,7 +7,7 @@ use crate::{
|
||||
};
|
||||
use futures::{SinkExt, StreamExt};
|
||||
use log::error;
|
||||
use solana_program::pubkey;
|
||||
use solana_sdk::pubkey;
|
||||
use solana_sdk::pubkey::Pubkey;
|
||||
use yellowstone_grpc_proto::geyser::SubscribeUpdateTransactionInfo;
|
||||
use yellowstone_grpc_proto::geyser::{
|
||||
|
||||
Reference in New Issue
Block a user