Skip to main content

sol_parser_sdk/grpc/
client.rs

1//! Yellowstone gRPC 客户端 - 超低延迟 DEX 事件订阅
2//!
3//! 支持多种事件输出模式:
4//! - Unordered: 10-20μs 极低延迟
5//! - MicroBatch: 50-200μs 微批次有序
6//! - StreamingOrdered: 0.1-5ms 流式有序
7//! - Ordered: 1-50ms 完全有序
8
9use super::buffers::{MicroBatchBuffer, SlotBuffer};
10use super::subscribe_builder::{
11    build_subscribe_request, build_subscribe_request_with_event_filter,
12};
13use super::types::*;
14use crate::core::{now_micros, EventMetadata}; // 导入高性能时钟
15use crate::instr::read_pubkey_fast;
16use crate::logs::timestamp_to_microseconds;
17use crate::DexEvent;
18use crossbeam_queue::ArrayQueue;
19use futures::{SinkExt, StreamExt};
20use log::error;
21use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
22use std::sync::Arc;
23use tokio::sync::{mpsc, Mutex};
24use tokio::task::JoinHandle;
25use tokio::time::{Duration, Instant};
26// Note: ClientTlsConfig moved to yellowstone_grpc_client in newer versions
27use yellowstone_grpc_client::{ClientTlsConfig, GeyserGrpcClient};
28use yellowstone_grpc_proto::prelude::*;
29
30static GRPC_DROPPED_EVENTS: AtomicU64 = AtomicU64::new(0);
31
32#[inline]
33fn push_queue(queue: &ArrayQueue<DexEvent>, event: DexEvent) {
34    if queue.push(event).is_err() {
35        let dropped = GRPC_DROPPED_EVENTS.fetch_add(1, Ordering::Relaxed) + 1;
36        if dropped <= 10 || dropped.is_power_of_two() {
37            log::warn!(
38                target: "sol_parser_sdk::grpc",
39                "gRPC event queue is full; dropped event count={dropped}"
40            );
41        }
42    }
43}
44
45// ==================== YellowstoneGrpc 客户端 ====================
46
47#[derive(Clone)]
48pub struct YellowstoneGrpc {
49    endpoint: String,
50    token: Option<String>,
51    config: ClientConfig,
52    control_tx: Arc<Mutex<Option<mpsc::Sender<SubscribeRequest>>>>,
53    subscription_handle: Arc<Mutex<Option<JoinHandle<()>>>>,
54    subscription_lifecycle: Arc<Mutex<()>>,
55    stop_signal: Arc<Mutex<Option<Arc<AtomicBool>>>>,
56}
57
58impl YellowstoneGrpc {
59    pub fn new(
60        endpoint: String,
61        token: Option<String>,
62    ) -> Result<Self, Box<dyn std::error::Error>> {
63        crate::warmup::warmup_parser();
64        Ok(Self {
65            endpoint,
66            token,
67            config: ClientConfig::default(),
68            control_tx: Arc::new(Mutex::new(None)),
69            subscription_handle: Arc::new(Mutex::new(None)),
70            subscription_lifecycle: Arc::new(Mutex::new(())),
71            stop_signal: Arc::new(Mutex::new(None)),
72        })
73    }
74
75    pub fn new_with_config(
76        endpoint: String,
77        token: Option<String>,
78        config: ClientConfig,
79    ) -> Result<Self, Box<dyn std::error::Error>> {
80        crate::warmup::warmup_parser();
81        Ok(Self {
82            endpoint,
83            token,
84            config,
85            control_tx: Arc::new(Mutex::new(None)),
86            subscription_handle: Arc::new(Mutex::new(None)),
87            subscription_lifecycle: Arc::new(Mutex::new(())),
88            stop_signal: Arc::new(Mutex::new(None)),
89        })
90    }
91
92    /// 订阅 DEX 事件(自动重连)
93    pub async fn subscribe_dex_events(
94        &self,
95        transaction_filters: Vec<TransactionFilter>,
96        account_filters: Vec<AccountFilter>,
97        event_type_filter: Option<EventTypeFilter>,
98    ) -> Result<Arc<ArrayQueue<DexEvent>>, Box<dyn std::error::Error>> {
99        let _lifecycle = self.subscription_lifecycle.lock().await;
100        self.stop_without_lifecycle_lock().await;
101
102        let queue = Arc::new(ArrayQueue::new(self.config.buffer_size.max(1)));
103        let queue_clone = Arc::clone(&queue);
104        let self_clone = self.clone();
105        let stop_signal = Arc::new(AtomicBool::new(false));
106        *self.stop_signal.lock().await = Some(Arc::clone(&stop_signal));
107
108        let handle = tokio::spawn(async move {
109            let mut delay = 1u64;
110            loop {
111                if stop_signal.load(Ordering::SeqCst) {
112                    break;
113                }
114
115                match self_clone
116                    .stream_events(
117                        &transaction_filters,
118                        &account_filters,
119                        &event_type_filter,
120                        &queue_clone,
121                    )
122                    .await
123                {
124                    Ok(_) => delay = 1,
125                    Err(e) => {
126                        if stop_signal.load(Ordering::SeqCst) {
127                            break;
128                        }
129                        error!("Grpc error: {} - retry in {}s", e, delay);
130                    }
131                }
132
133                if stop_signal.load(Ordering::SeqCst) {
134                    break;
135                }
136                tokio::time::sleep(Duration::from_secs(delay)).await;
137                delay = (delay * 2).min(60);
138            }
139        });
140
141        *self.subscription_handle.lock().await = Some(handle);
142        Ok(queue)
143    }
144
145    /// 动态更新订阅过滤器
146    pub async fn update_subscription(
147        &self,
148        transaction_filters: Vec<TransactionFilter>,
149        account_filters: Vec<AccountFilter>,
150    ) -> Result<(), Box<dyn std::error::Error>> {
151        let sender = self.control_tx.lock().await.as_ref().ok_or("No active subscription")?.clone();
152
153        let request = build_subscribe_request(&transaction_filters, &account_filters);
154        sender.send(request).await.map_err(|e| e.to_string())?;
155        Ok(())
156    }
157
158    pub async fn stop(&self) {
159        let _lifecycle = self.subscription_lifecycle.lock().await;
160        self.stop_without_lifecycle_lock().await;
161    }
162
163    async fn stop_without_lifecycle_lock(&self) {
164        if let Some(stop_signal) = self.stop_signal.lock().await.take() {
165            stop_signal.store(true, Ordering::SeqCst);
166        }
167        self.control_tx.lock().await.take();
168        let handle = self.subscription_handle.lock().await.take();
169        if let Some(handle) = handle {
170            handle.abort();
171            let _ = handle.await;
172        }
173    }
174
175    // ==================== 核心事件流处理 ====================
176
177    async fn stream_events(
178        &self,
179        tx_filters: &[TransactionFilter],
180        acc_filters: &[AccountFilter],
181        event_filter: &Option<EventTypeFilter>,
182        queue: &Arc<ArrayQueue<DexEvent>>,
183    ) -> Result<(), String> {
184        let _ = rustls::crypto::ring::default_provider().install_default();
185
186        // 构建客户端
187        let mut builder = GeyserGrpcClient::build_from_shared(self.endpoint.clone())
188            .map_err(|e| e.to_string())?
189            .x_token(self.token.clone())
190            .map_err(|e| e.to_string())?
191            .max_decoding_message_size(1024 * 1024 * 1024);
192
193        if self.config.connection_timeout_ms > 0 {
194            builder =
195                builder.connect_timeout(Duration::from_millis(self.config.connection_timeout_ms));
196        }
197        if self.config.enable_tls {
198            builder = builder
199                .tls_config(ClientTlsConfig::new().with_native_roots())
200                .map_err(|e| e.to_string())?;
201        }
202
203        let mut client = builder.connect().await.map_err(|e| e.to_string())?;
204        let request = build_subscribe_request_with_event_filter(
205            tx_filters,
206            acc_filters,
207            event_filter.as_ref(),
208            CommitmentLevel::Processed,
209        );
210
211        let (subscribe_tx, mut stream) =
212            client.subscribe_with_request(Some(request)).await.map_err(|e| e.to_string())?;
213
214        self.print_mode_info();
215
216        // 设置控制通道
217        let (control_tx, mut control_rx) = mpsc::channel::<SubscribeRequest>(100);
218        *self.control_tx.lock().await = Some(control_tx);
219        let subscribe_tx = Arc::new(Mutex::new(subscribe_tx));
220
221        // 初始化缓冲区
222        let mut slot_buffer = SlotBuffer::new();
223        let mut micro_batch = MicroBatchBuffer::new();
224        let mut last_slot = 0u64;
225
226        let order_mode = self.config.order_mode;
227        let timeout_ms = self.config.order_timeout_ms;
228        let batch_us = self.config.micro_batch_us;
229        let check_interval = match order_mode {
230            OrderMode::MicroBatch => Duration::from_micros(batch_us.max(1)),
231            _ => Duration::from_millis((timeout_ms / 2).max(1)),
232        };
233        let mut next_check = Instant::now() + check_interval;
234
235        loop {
236            tokio::select! {
237                msg = stream.next() => {
238                    match msg {
239                        Some(Ok(update)) => {
240                            // Geyser 会周期性下发 ping;必须在同一 subscribe 流上回写 SubscribeRequest.ping,否则公共节点 / LB 可能 RST_STREAM。
241                            if matches!(
242                                update.update_oneof.as_ref(),
243                                Some(subscribe_update::UpdateOneof::Ping(_))
244                            ) {
245                                if let Err(e) = subscribe_tx
246                                    .lock()
247                                    .await
248                                    .send(SubscribeRequest {
249                                        ping: Some(SubscribeRequestPing { id: 1 }),
250                                        ..Default::default()
251                                    })
252                                    .await
253                                {
254                                    self.control_tx.lock().await.take();
255                                    return Err(e.to_string());
256                                }
257                                continue;
258                            }
259                            self.handle_update(
260                                update, order_mode, event_filter, queue,
261                                &mut slot_buffer, &mut micro_batch, &mut last_slot, batch_us
262                            );
263                        }
264                        Some(Err(e)) => {
265                            error!("Grpc Stream error: {:?}", e);
266                            self.flush_on_disconnect(
267                                order_mode,
268                                &mut slot_buffer,
269                                &mut micro_batch,
270                                queue,
271                            );
272                            self.control_tx.lock().await.take();
273                            return Err(e.to_string());
274                        }
275                        None => {
276                            self.flush_on_disconnect(
277                                order_mode,
278                                &mut slot_buffer,
279                                &mut micro_batch,
280                                queue,
281                            );
282                            self.control_tx.lock().await.take();
283                            return Ok(());
284                        }
285                    }
286                }
287                Some(req) = control_rx.recv() => {
288                    if let Err(e) = subscribe_tx.lock().await.send(req).await {
289                        self.control_tx.lock().await.take();
290                        return Err(e.to_string());
291                    }
292                }
293                _ = tokio::time::sleep_until(next_check) => {
294                    self.check_timeout(
295                        order_mode,
296                        &mut slot_buffer,
297                        &mut micro_batch,
298                        queue,
299                        timeout_ms,
300                        batch_us,
301                        &mut next_check,
302                        check_interval,
303                    );
304                }
305            }
306        }
307    }
308
309    fn print_mode_info(&self) {
310        match self.config.order_mode {
311            OrderMode::Unordered => println!("✅ Unordered Mode (10-20μs)"),
312            OrderMode::Ordered => {
313                println!("✅ Ordered Mode (timeout={}ms)", self.config.order_timeout_ms)
314            }
315            OrderMode::StreamingOrdered => {
316                println!("✅ StreamingOrdered Mode (timeout={}ms)", self.config.order_timeout_ms)
317            }
318            OrderMode::MicroBatch => {
319                println!("✅ MicroBatch Mode (window={}μs)", self.config.micro_batch_us)
320            }
321        }
322    }
323
324    #[inline]
325    fn check_timeout(
326        &self,
327        mode: OrderMode,
328        slot_buf: &mut SlotBuffer,
329        micro_buf: &mut MicroBatchBuffer,
330        queue: &Arc<ArrayQueue<DexEvent>>,
331        timeout_ms: u64,
332        batch_us: u64,
333        next_check: &mut Instant,
334        interval: Duration,
335    ) {
336        if Instant::now() < *next_check {
337            return;
338        }
339        *next_check = Instant::now() + interval;
340
341        match mode {
342            OrderMode::Ordered => {
343                if slot_buf.should_timeout(timeout_ms) {
344                    for e in slot_buf.flush_all() {
345                        push_queue(queue, e);
346                    }
347                }
348            }
349            OrderMode::StreamingOrdered => {
350                if slot_buf.should_timeout(timeout_ms) {
351                    for e in slot_buf.flush_streaming_timeout() {
352                        push_queue(queue, e);
353                    }
354                }
355            }
356            OrderMode::MicroBatch => {
357                // Periodic flush for MicroBatch mode
358                let now_us = get_timestamp_us();
359                if micro_buf.should_flush(now_us, batch_us) {
360                    for e in micro_buf.flush() {
361                        push_queue(queue, e);
362                    }
363                }
364            }
365            OrderMode::Unordered => {}
366        }
367    }
368
369    fn flush_on_disconnect(
370        &self,
371        mode: OrderMode,
372        buffer: &mut SlotBuffer,
373        micro_batch: &mut MicroBatchBuffer,
374        queue: &Arc<ArrayQueue<DexEvent>>,
375    ) {
376        let events = match mode {
377            OrderMode::Ordered => buffer.flush_all(),
378            OrderMode::StreamingOrdered => buffer.flush_streaming_timeout(),
379            OrderMode::MicroBatch => micro_batch.flush(),
380            OrderMode::Unordered => Vec::new(),
381        };
382        for event in events {
383            push_queue(queue, event);
384        }
385    }
386
387    #[inline]
388    fn handle_update(
389        &self,
390        update_msg: SubscribeUpdate,
391        mode: OrderMode,
392        filter: &Option<EventTypeFilter>,
393        queue: &Arc<ArrayQueue<DexEvent>>,
394        slot_buf: &mut SlotBuffer,
395        micro_buf: &mut MicroBatchBuffer,
396        last_slot: &mut u64,
397        batch_us: u64,
398    ) {
399        let created_at = update_msg.created_at.unwrap_or_default();
400        let block_time_us = timestamp_to_microseconds(created_at.seconds, created_at.nanos) as i64;
401        let grpc_recv_us = get_timestamp_us();
402
403        let Some(update) = update_msg.update_oneof else { return };
404
405        match update {
406            subscribe_update::UpdateOneof::Transaction(tx) => {
407                self.handle_transaction(
408                    tx,
409                    mode,
410                    filter,
411                    queue,
412                    slot_buf,
413                    micro_buf,
414                    last_slot,
415                    batch_us,
416                    grpc_recv_us,
417                    block_time_us,
418                );
419            }
420            subscribe_update::UpdateOneof::Account(acc) => {
421                Self::handle_account(acc, filter, queue, grpc_recv_us, block_time_us);
422            }
423            subscribe_update::UpdateOneof::BlockMeta(block_meta) => {
424                Self::handle_block_meta(block_meta, filter, queue, grpc_recv_us, block_time_us);
425            }
426            _ => {}
427        }
428    }
429
430    #[inline]
431    fn handle_transaction(
432        &self,
433        tx: SubscribeUpdateTransaction,
434        mode: OrderMode,
435        filter: &Option<EventTypeFilter>,
436        queue: &Arc<ArrayQueue<DexEvent>>,
437        slot_buf: &mut SlotBuffer,
438        micro_buf: &mut MicroBatchBuffer,
439        last_slot: &mut u64,
440        batch_us: u64,
441        grpc_us: i64,
442        block_us: i64,
443    ) {
444        let slot = tx.slot;
445
446        match mode {
447            OrderMode::Unordered => {
448                for e in crate::grpc::parse_subscribe_update_transaction_low_latency(
449                    &tx,
450                    grpc_us,
451                    Some(block_us),
452                    filter.as_ref(),
453                ) {
454                    push_queue(queue, e);
455                }
456            }
457            OrderMode::Ordered => {
458                if slot > *last_slot && *last_slot > 0 {
459                    for e in slot_buf.flush_before(slot) {
460                        push_queue(queue, e);
461                    }
462                }
463                *last_slot = slot;
464                for (idx, e) in
465                    parse_transaction_to_vec(&tx, grpc_us, Some(block_us), filter.as_ref())
466                {
467                    slot_buf.push(slot, idx, e);
468                }
469            }
470            OrderMode::StreamingOrdered => {
471                for (idx, e) in
472                    parse_transaction_to_vec(&tx, grpc_us, Some(block_us), filter.as_ref())
473                {
474                    for evt in slot_buf.push_streaming(slot, idx, e) {
475                        push_queue(queue, evt);
476                    }
477                }
478            }
479            OrderMode::MicroBatch => {
480                for (idx, e) in
481                    parse_transaction_to_vec(&tx, grpc_us, Some(block_us), filter.as_ref())
482                {
483                    if micro_buf.push(slot, idx, e, grpc_us, batch_us) {
484                        for evt in micro_buf.flush() {
485                            push_queue(queue, evt);
486                        }
487                    }
488                }
489            }
490        }
491    }
492
493    #[inline]
494    fn handle_account(
495        acc: SubscribeUpdateAccount,
496        filter: &Option<EventTypeFilter>,
497        queue: &Arc<ArrayQueue<DexEvent>>,
498        grpc_us: i64,
499        block_us: i64,
500    ) {
501        let Some(info) = acc.account else { return };
502        let data = crate::accounts::AccountData {
503            pubkey: read_pubkey_fast(&info.pubkey),
504            executable: info.executable,
505            lamports: info.lamports,
506            owner: read_pubkey_fast(&info.owner),
507            rent_epoch: info.rent_epoch,
508            data: info.data,
509        };
510        let meta = EventMetadata {
511            signature: Default::default(),
512            slot: acc.slot,
513            tx_index: 0,
514            block_time_us: block_us,
515            grpc_recv_us: grpc_us,
516            recent_blockhash: None,
517        };
518        if let Some(e) = crate::accounts::parse_account_unified(&data, meta, filter.as_ref()) {
519            push_queue(queue, e);
520        }
521    }
522
523    #[inline]
524    fn handle_block_meta(
525        block_meta: SubscribeUpdateBlockMeta,
526        filter: &Option<EventTypeFilter>,
527        queue: &Arc<ArrayQueue<DexEvent>>,
528        grpc_us: i64,
529        fallback_block_us: i64,
530    ) {
531        let block_time_us = block_meta
532            .block_time
533            .as_ref()
534            .map(|t| t.timestamp.saturating_mul(1_000_000))
535            .unwrap_or(fallback_block_us);
536        let event = DexEvent::BlockMeta(crate::core::events::BlockMetaEvent {
537            metadata: EventMetadata {
538                signature: Default::default(),
539                slot: block_meta.slot,
540                tx_index: 0,
541                block_time_us,
542                grpc_recv_us: grpc_us,
543                recent_blockhash: (!block_meta.blockhash.is_empty())
544                    .then_some(block_meta.blockhash),
545            },
546        });
547        if filter.as_ref().map(|f| f.should_include_dex_event(&event)).unwrap_or(true) {
548            push_queue(queue, event);
549        }
550    }
551}
552
553// ==================== 辅助函数 ====================
554
555/// 获取当前时间戳(微秒)
556///
557/// 使用高性能时钟,避免系统调用开销
558///
559/// # 性能优势
560/// - 旧实现:使用 libc::clock_gettime,每次调用约 1-2μs
561/// - 新实现:使用高性能时钟,每次调用约 10-50ns
562/// - 性能提升:20-100 倍
563#[inline(always)]
564fn get_timestamp_us() -> i64 {
565    now_micros()
566}
567
568// ==================== 交易解析 ====================
569
570#[inline]
571fn parse_transaction_to_vec(
572    tx: &SubscribeUpdateTransaction,
573    grpc_us: i64,
574    block_us: Option<i64>,
575    filter: Option<&EventTypeFilter>,
576) -> Vec<(u64, DexEvent)> {
577    let idx = tx.transaction.as_ref().map(|t| t.index).unwrap_or(0);
578    crate::grpc::parse_subscribe_update_transaction_low_latency(tx, grpc_us, block_us, filter)
579        .into_iter()
580        .map(|event| (idx, event))
581        .collect()
582}
583
584#[cfg(test)]
585mod tests {
586    use super::*;
587
588    fn test_event(slot: u64) -> DexEvent {
589        DexEvent::BlockMeta(crate::core::events::BlockMetaEvent {
590            metadata: EventMetadata { slot, ..Default::default() },
591        })
592    }
593
594    #[tokio::test]
595    async fn stop_clears_subscription_state_and_aborts_handle() {
596        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_string(), None).unwrap();
597        let (tx, _rx) = mpsc::channel::<SubscribeRequest>(1);
598        let handle = tokio::spawn(async {
599            std::future::pending::<()>().await;
600        });
601
602        let stop_signal = Arc::new(AtomicBool::new(false));
603        *grpc.control_tx.lock().await = Some(tx);
604        *grpc.subscription_handle.lock().await = Some(handle);
605        *grpc.stop_signal.lock().await = Some(Arc::clone(&stop_signal));
606
607        grpc.stop().await;
608
609        assert!(stop_signal.load(Ordering::SeqCst));
610        assert!(grpc.stop_signal.lock().await.is_none());
611        assert!(grpc.control_tx.lock().await.is_none());
612        assert!(grpc.subscription_handle.lock().await.is_none());
613    }
614
615    #[test]
616    fn micro_batch_is_flushed_on_disconnect() {
617        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_string(), None).unwrap();
618        let queue = Arc::new(ArrayQueue::new(2));
619        let mut slot_buffer = SlotBuffer::new();
620        let mut micro_batch = MicroBatchBuffer::new();
621        assert!(!micro_batch.push(7, 0, test_event(7), 10, 100));
622
623        grpc.flush_on_disconnect(OrderMode::MicroBatch, &mut slot_buffer, &mut micro_batch, &queue);
624
625        assert!(micro_batch.is_empty());
626        assert!(matches!(queue.pop(), Some(DexEvent::BlockMeta(_))));
627    }
628
629    #[test]
630    fn full_queue_increments_drop_counter() {
631        let queue = ArrayQueue::new(1);
632        push_queue(&queue, test_event(1));
633        let before = GRPC_DROPPED_EVENTS.load(Ordering::Relaxed);
634        push_queue(&queue, test_event(2));
635        assert_eq!(GRPC_DROPPED_EVENTS.load(Ordering::Relaxed), before + 1);
636    }
637}