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::build_subscribe_request_with_event_filter;
11use super::types::*;
12use crate::core::{now_micros, EventMetadata}; // 导入高性能时钟
13use crate::instr::read_pubkey_fast;
14use crate::logs::timestamp_to_microseconds;
15use crate::DexEvent;
16use crossbeam_queue::ArrayQueue;
17use futures::{SinkExt, StreamExt};
18use log::error;
19use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
20use std::sync::Arc;
21use tokio::sync::{mpsc, Mutex};
22use tokio::task::JoinHandle;
23use tokio::time::{Duration, Instant};
24// Note: ClientTlsConfig moved to yellowstone_grpc_client in newer versions
25use yellowstone_grpc_client::{ClientTlsConfig, GeyserGrpcClient};
26use yellowstone_grpc_proto::prelude::*;
27
28static GRPC_DROPPED_EVENTS: AtomicU64 = AtomicU64::new(0);
29
30#[inline]
31fn push_queue(queue: &ArrayQueue<DexEvent>, event: DexEvent) -> bool {
32    if queue.push(event).is_err() {
33        let dropped = GRPC_DROPPED_EVENTS.fetch_add(1, Ordering::Relaxed) + 1;
34        if dropped <= 10 || dropped.is_power_of_two() {
35            log::warn!(
36                target: "sol_parser_sdk::grpc",
37                "gRPC event queue is full; dropped event count={dropped}"
38            );
39        }
40        false
41    } else {
42        true
43    }
44}
45
46/// Counters are local to this client and shared by its clones. A changed revision
47/// means cached state may have gaps; discard queued updates and rebuild state.
48/// This reports transport continuity, not fork finality or snapshot completeness.
49#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
50pub struct GrpcSubscriptionStatus {
51    pub connected: bool,
52    pub generation: u64,
53    pub continuity_revision: u64,
54    pub disconnects: u64,
55    pub dropped_events: u64,
56}
57
58#[derive(Default)]
59struct SubscriptionHealth(std::sync::Mutex<GrpcSubscriptionStatus>);
60impl SubscriptionHealth {
61    fn update(&self, f: impl FnOnce(&mut GrpcSubscriptionStatus)) {
62        let mut status = self.0.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
63        f(&mut status);
64    }
65    fn connected(&self) {
66        self.update(|s| {
67            s.generation += 1;
68            s.continuity_revision += 1;
69            s.connected = true;
70        });
71    }
72    fn disconnected(&self) {
73        self.update(|s| {
74            if s.connected {
75                s.connected = false;
76                s.disconnects += 1;
77                s.continuity_revision += 1;
78            }
79        });
80    }
81    fn invalidate(&self) {
82        self.update(|s| s.continuity_revision += 1);
83    }
84    fn dropped(&self) {
85        self.update(|s| {
86            s.dropped_events += 1;
87            s.continuity_revision += 1;
88        });
89    }
90    fn snapshot(&self) -> GrpcSubscriptionStatus {
91        *self.0.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
92    }
93}
94
95// ==================== YellowstoneGrpc 客户端 ====================
96
97#[derive(Clone)]
98struct SubscriptionFilters {
99    transactions: Vec<TransactionFilter>,
100    accounts: Vec<AccountFilter>,
101    events: Option<EventTypeFilter>,
102}
103
104/// Back off failed connection attempts, but reset after any established stream.
105struct ReconnectBackoff {
106    base_ms: u64,
107    next_ms: u64,
108}
109impl ReconnectBackoff {
110    fn new(retry_delay_ms: u64) -> Self {
111        let base_ms = retry_delay_ms.clamp(1, 60_000);
112        Self { base_ms, next_ms: base_ms }
113    }
114    fn next_delay(&mut self, established_stream: bool) -> Duration {
115        if established_stream {
116            self.next_ms = self.base_ms;
117        }
118        let delay = self.next_ms;
119        self.next_ms = self.next_ms.saturating_mul(2).min(60_000);
120        Duration::from_millis(delay)
121    }
122}
123
124#[derive(Clone)]
125pub struct YellowstoneGrpc {
126    endpoint: String,
127    token: Option<String>,
128    config: ClientConfig,
129    control_tx: Arc<Mutex<Option<mpsc::Sender<SubscribeRequest>>>>,
130    subscription_handle: Arc<Mutex<Option<JoinHandle<()>>>>,
131    subscription_lifecycle: Arc<Mutex<()>>,
132    stop_signal: Arc<Mutex<Option<Arc<AtomicBool>>>>,
133    health: Arc<SubscriptionHealth>,
134    subscription_filters: Arc<Mutex<Option<SubscriptionFilters>>>,
135}
136
137impl YellowstoneGrpc {
138    /// Sample before and after consuming a batch. A gap/reconnect/overflow must
139    /// invalidate dependent account caches; reconnect does not backfill updates.
140    pub fn subscription_status(&self) -> GrpcSubscriptionStatus {
141        self.health.snapshot()
142    }
143
144    #[inline]
145    fn push_queue(&self, queue: &ArrayQueue<DexEvent>, event: DexEvent) {
146        if !push_queue(queue, event) {
147            self.health.dropped();
148        }
149    }
150
151    pub fn new(
152        endpoint: String,
153        token: Option<String>,
154    ) -> Result<Self, Box<dyn std::error::Error>> {
155        crate::warmup::warmup_parser();
156        Ok(Self {
157            endpoint,
158            token,
159            config: ClientConfig::default(),
160            control_tx: Arc::new(Mutex::new(None)),
161            subscription_handle: Arc::new(Mutex::new(None)),
162            subscription_lifecycle: Arc::new(Mutex::new(())),
163            stop_signal: Arc::new(Mutex::new(None)),
164            health: Arc::new(SubscriptionHealth::default()),
165            subscription_filters: Arc::new(Mutex::new(None)),
166        })
167    }
168
169    pub fn new_with_config(
170        endpoint: String,
171        token: Option<String>,
172        config: ClientConfig,
173    ) -> Result<Self, Box<dyn std::error::Error>> {
174        crate::warmup::warmup_parser();
175        Ok(Self {
176            endpoint,
177            token,
178            config,
179            control_tx: Arc::new(Mutex::new(None)),
180            subscription_handle: Arc::new(Mutex::new(None)),
181            subscription_lifecycle: Arc::new(Mutex::new(())),
182            stop_signal: Arc::new(Mutex::new(None)),
183            health: Arc::new(SubscriptionHealth::default()),
184            subscription_filters: Arc::new(Mutex::new(None)),
185        })
186    }
187
188    /// 订阅 DEX 事件(自动重连)
189    pub async fn subscribe_dex_events(
190        &self,
191        transaction_filters: Vec<TransactionFilter>,
192        account_filters: Vec<AccountFilter>,
193        event_type_filter: Option<EventTypeFilter>,
194    ) -> Result<Arc<ArrayQueue<DexEvent>>, Box<dyn std::error::Error>> {
195        let _lifecycle = self.subscription_lifecycle.lock().await;
196        self.stop_without_lifecycle_lock().await;
197
198        *self.subscription_filters.lock().await = Some(SubscriptionFilters {
199            transactions: transaction_filters,
200            accounts: account_filters,
201            events: event_type_filter,
202        });
203        let queue = Arc::new(ArrayQueue::new(self.config.buffer_size.max(1)));
204        let queue_clone = Arc::clone(&queue);
205        let self_clone = self.clone();
206        let stop_signal = Arc::new(AtomicBool::new(false));
207        *self.stop_signal.lock().await = Some(Arc::clone(&stop_signal));
208
209        let handle = tokio::spawn(async move {
210            let mut backoff = ReconnectBackoff::new(self_clone.config.retry_delay_ms);
211            loop {
212                if stop_signal.load(Ordering::SeqCst) {
213                    break;
214                }
215
216                // Reconnect uses the latest accepted filters, not the initial ones.
217                let filters = self_clone.subscription_filters.lock().await.clone();
218                let Some(filters) = filters else { break };
219                let generation = self_clone.subscription_status().generation;
220                let result = self_clone
221                    .stream_events(
222                        &filters.transactions,
223                        &filters.accounts,
224                        &filters.events,
225                        &queue_clone,
226                    )
227                    .await;
228                self_clone.health.disconnected();
229                let delay =
230                    backoff.next_delay(self_clone.subscription_status().generation != generation);
231                match result {
232                    Ok(_) => {}
233                    Err(e) => {
234                        if stop_signal.load(Ordering::SeqCst) {
235                            break;
236                        }
237                        error!("Grpc error: {} - retry in {}ms", e, delay.as_millis());
238                    }
239                }
240
241                if stop_signal.load(Ordering::SeqCst) {
242                    break;
243                }
244                tokio::time::sleep(delay).await;
245            }
246        });
247
248        *self.subscription_handle.lock().await = Some(handle);
249        Ok(queue)
250    }
251
252    /// 动态更新订阅过滤器
253    pub async fn update_subscription(
254        &self,
255        transaction_filters: Vec<TransactionFilter>,
256        account_filters: Vec<AccountFilter>,
257    ) -> Result<(), Box<dyn std::error::Error>> {
258        // Serialize updates with stop/resubscribe and retain the original event policy.
259        let _lifecycle = self.subscription_lifecycle.lock().await;
260        let sender = self.control_tx.lock().await.as_ref().ok_or("No active subscription")?.clone();
261        let mut desired = self.subscription_filters.lock().await;
262        let current = desired.as_ref().ok_or("No active subscription filters")?;
263        let next = SubscriptionFilters {
264            transactions: transaction_filters,
265            accounts: account_filters,
266            events: current.events.clone(),
267        };
268        let request = build_subscribe_request_with_event_filter(
269            &next.transactions,
270            &next.accounts,
271            next.events.as_ref(),
272            CommitmentLevel::Processed,
273        );
274        // Invalidate before the provider can start emitting the changed account set.
275        self.health.invalidate();
276        // Do not hold the lifecycle lock while waiting for a full control queue:
277        // otherwise a stalled provider could prevent stop() from acquiring it.
278        sender.try_send(request).map_err(|error| match error {
279            mpsc::error::TrySendError::Full(_) => "Subscription update queue is full; retry later",
280            mpsc::error::TrySendError::Closed(_) => "No active subscription",
281        })?;
282        *desired = Some(next);
283        Ok(())
284    }
285
286    pub async fn stop(&self) {
287        let _lifecycle = self.subscription_lifecycle.lock().await;
288        self.stop_without_lifecycle_lock().await;
289    }
290
291    async fn stop_without_lifecycle_lock(&self) {
292        self.health.disconnected();
293        if let Some(stop_signal) = self.stop_signal.lock().await.take() {
294            stop_signal.store(true, Ordering::SeqCst);
295        }
296        self.control_tx.lock().await.take();
297        let handle = self.subscription_handle.lock().await.take();
298        if let Some(handle) = handle {
299            handle.abort();
300            let _ = handle.await;
301        }
302        // A connection may have completed while cancellation was in flight.
303        self.health.disconnected();
304        self.subscription_filters.lock().await.take();
305    }
306
307    // ==================== 核心事件流处理 ====================
308
309    async fn stream_events(
310        &self,
311        tx_filters: &[TransactionFilter],
312        acc_filters: &[AccountFilter],
313        event_filter: &Option<EventTypeFilter>,
314        queue: &Arc<ArrayQueue<DexEvent>>,
315    ) -> Result<(), String> {
316        let _ = rustls::crypto::ring::default_provider().install_default();
317
318        // 构建客户端
319        let mut builder = GeyserGrpcClient::build_from_shared(self.endpoint.clone())
320            .map_err(|e| e.to_string())?
321            .x_token(self.token.clone())
322            .map_err(|e| e.to_string())?
323            .max_decoding_message_size(1024 * 1024 * 1024);
324
325        if self.config.connection_timeout_ms > 0 {
326            builder =
327                builder.connect_timeout(Duration::from_millis(self.config.connection_timeout_ms));
328        }
329        if self.config.enable_tls {
330            builder = builder
331                .tls_config(ClientTlsConfig::new().with_native_roots())
332                .map_err(|e| e.to_string())?;
333        }
334
335        let mut client = builder.connect().await.map_err(|e| e.to_string())?;
336        let request = build_subscribe_request_with_event_filter(
337            tx_filters,
338            acc_filters,
339            event_filter.as_ref(),
340            CommitmentLevel::Processed,
341        );
342
343        let (subscribe_tx, mut stream) =
344            client.subscribe_with_request(Some(request)).await.map_err(|e| e.to_string())?;
345
346        self.health.connected();
347        self.print_mode_info();
348
349        // 设置控制通道
350        let (control_tx, mut control_rx) = mpsc::channel::<SubscribeRequest>(100);
351        *self.control_tx.lock().await = Some(control_tx);
352        let subscribe_tx = Arc::new(Mutex::new(subscribe_tx));
353
354        // 初始化缓冲区
355        let mut slot_buffer = SlotBuffer::new();
356        let mut micro_batch = MicroBatchBuffer::new();
357        let mut last_slot = 0u64;
358
359        let order_mode = self.config.order_mode;
360        let timeout_ms = self.config.order_timeout_ms;
361        let batch_us = self.config.micro_batch_us;
362        let check_interval = match order_mode {
363            OrderMode::MicroBatch => Duration::from_micros(batch_us.max(1)),
364            _ => Duration::from_millis((timeout_ms / 2).max(1)),
365        };
366        let mut next_check = Instant::now() + check_interval;
367
368        loop {
369            tokio::select! {
370                msg = stream.next() => {
371                    match msg {
372                        Some(Ok(update)) => {
373                            // Geyser 会周期性下发 ping;必须在同一 subscribe 流上回写 SubscribeRequest.ping,否则公共节点 / LB 可能 RST_STREAM。
374                            if matches!(
375                                update.update_oneof.as_ref(),
376                                Some(subscribe_update::UpdateOneof::Ping(_))
377                            ) {
378                                if let Err(e) = subscribe_tx
379                                    .lock()
380                                    .await
381                                    .send(SubscribeRequest {
382                                        ping: Some(SubscribeRequestPing { id: 1 }),
383                                        ..Default::default()
384                                    })
385                                    .await
386                                {
387                                    self.control_tx.lock().await.take();
388                                    return Err(e.to_string());
389                                }
390                                continue;
391                            }
392                            self.handle_update(
393                                update, order_mode, event_filter, queue,
394                                &mut slot_buffer, &mut micro_batch, &mut last_slot, batch_us
395                            );
396                        }
397                        Some(Err(e)) => {
398                            error!("Grpc Stream error: {:?}", e);
399                            self.flush_on_disconnect(
400                                order_mode,
401                                &mut slot_buffer,
402                                &mut micro_batch,
403                                queue,
404                            );
405                            self.control_tx.lock().await.take();
406                            return Err(e.to_string());
407                        }
408                        None => {
409                            self.flush_on_disconnect(
410                                order_mode,
411                                &mut slot_buffer,
412                                &mut micro_batch,
413                                queue,
414                            );
415                            self.control_tx.lock().await.take();
416                            return Ok(());
417                        }
418                    }
419                }
420                Some(req) = control_rx.recv() => {
421                    if let Err(e) = subscribe_tx.lock().await.send(req).await {
422                        self.control_tx.lock().await.take();
423                        return Err(e.to_string());
424                    }
425                }
426                // Unordered emits immediately and has no ordering buffer to flush.
427                _ = tokio::time::sleep_until(next_check), if order_mode != OrderMode::Unordered => {
428                    self.check_timeout(
429                        order_mode,
430                        &mut slot_buffer,
431                        &mut micro_batch,
432                        queue,
433                        timeout_ms,
434                        batch_us,
435                        &mut next_check,
436                        check_interval,
437                    );
438                }
439            }
440        }
441    }
442
443    fn print_mode_info(&self) {
444        match self.config.order_mode {
445            OrderMode::Unordered => println!("✅ Unordered Mode (10-20μs)"),
446            OrderMode::Ordered => {
447                println!("✅ Ordered Mode (timeout={}ms)", self.config.order_timeout_ms)
448            }
449            OrderMode::StreamingOrdered => {
450                println!("✅ StreamingOrdered Mode (timeout={}ms)", self.config.order_timeout_ms)
451            }
452            OrderMode::MicroBatch => {
453                println!("✅ MicroBatch Mode (window={}μs)", self.config.micro_batch_us)
454            }
455        }
456    }
457
458    #[inline]
459    fn check_timeout(
460        &self,
461        mode: OrderMode,
462        slot_buf: &mut SlotBuffer,
463        micro_buf: &mut MicroBatchBuffer,
464        queue: &Arc<ArrayQueue<DexEvent>>,
465        timeout_ms: u64,
466        batch_us: u64,
467        next_check: &mut Instant,
468        interval: Duration,
469    ) {
470        if Instant::now() < *next_check {
471            return;
472        }
473        *next_check = Instant::now() + interval;
474
475        match mode {
476            OrderMode::Ordered => {
477                if slot_buf.should_timeout(timeout_ms) {
478                    for e in slot_buf.flush_all() {
479                        self.push_queue(queue, e);
480                    }
481                }
482            }
483            OrderMode::StreamingOrdered => {
484                if slot_buf.should_timeout(timeout_ms) {
485                    for e in slot_buf.flush_streaming_timeout() {
486                        self.push_queue(queue, e);
487                    }
488                }
489            }
490            OrderMode::MicroBatch => {
491                // Periodic flush for MicroBatch mode
492                let now_us = get_timestamp_us();
493                if micro_buf.should_flush(now_us, batch_us) {
494                    for e in micro_buf.flush() {
495                        self.push_queue(queue, e);
496                    }
497                }
498            }
499            OrderMode::Unordered => {}
500        }
501    }
502
503    fn flush_on_disconnect(
504        &self,
505        mode: OrderMode,
506        buffer: &mut SlotBuffer,
507        micro_batch: &mut MicroBatchBuffer,
508        queue: &Arc<ArrayQueue<DexEvent>>,
509    ) {
510        let events = match mode {
511            OrderMode::Ordered => buffer.flush_all(),
512            OrderMode::StreamingOrdered => buffer.flush_streaming_timeout(),
513            OrderMode::MicroBatch => micro_batch.flush(),
514            OrderMode::Unordered => Vec::new(),
515        };
516        for event in events {
517            self.push_queue(queue, event);
518        }
519    }
520
521    #[inline]
522    fn handle_update(
523        &self,
524        update_msg: SubscribeUpdate,
525        mode: OrderMode,
526        filter: &Option<EventTypeFilter>,
527        queue: &Arc<ArrayQueue<DexEvent>>,
528        slot_buf: &mut SlotBuffer,
529        micro_buf: &mut MicroBatchBuffer,
530        last_slot: &mut u64,
531        batch_us: u64,
532    ) {
533        let created_at = update_msg.created_at.unwrap_or_default();
534        let block_time_us = timestamp_to_microseconds(created_at.seconds, created_at.nanos) as i64;
535        let grpc_recv_us = get_timestamp_us();
536
537        let Some(update) = update_msg.update_oneof else { return };
538
539        match update {
540            subscribe_update::UpdateOneof::Transaction(tx) => {
541                self.handle_transaction(
542                    tx,
543                    mode,
544                    filter,
545                    queue,
546                    slot_buf,
547                    micro_buf,
548                    last_slot,
549                    batch_us,
550                    grpc_recv_us,
551                    block_time_us,
552                );
553            }
554            subscribe_update::UpdateOneof::Account(acc) => {
555                self.handle_account(acc, filter, queue, grpc_recv_us, block_time_us);
556            }
557            subscribe_update::UpdateOneof::BlockMeta(block_meta) => {
558                self.handle_block_meta(block_meta, filter, queue, grpc_recv_us, block_time_us);
559            }
560            _ => {}
561        }
562    }
563
564    #[inline]
565    fn handle_transaction(
566        &self,
567        tx: SubscribeUpdateTransaction,
568        mode: OrderMode,
569        filter: &Option<EventTypeFilter>,
570        queue: &Arc<ArrayQueue<DexEvent>>,
571        slot_buf: &mut SlotBuffer,
572        micro_buf: &mut MicroBatchBuffer,
573        last_slot: &mut u64,
574        batch_us: u64,
575        grpc_us: i64,
576        block_us: i64,
577    ) {
578        let slot = tx.slot;
579
580        match mode {
581            OrderMode::Unordered => {
582                for e in crate::grpc::parse_subscribe_update_transaction_low_latency(
583                    &tx,
584                    grpc_us,
585                    Some(block_us),
586                    filter.as_ref(),
587                ) {
588                    self.push_queue(queue, e);
589                }
590            }
591            OrderMode::Ordered => {
592                if slot > *last_slot && *last_slot > 0 {
593                    for e in slot_buf.flush_before(slot) {
594                        self.push_queue(queue, e);
595                    }
596                }
597                *last_slot = (*last_slot).max(slot);
598                for (idx, e) in
599                    parse_transaction_to_vec(&tx, grpc_us, Some(block_us), filter.as_ref())
600                {
601                    slot_buf.push(slot, idx, e);
602                }
603            }
604            OrderMode::StreamingOrdered => {
605                let idx = tx.transaction.as_ref().map_or(0, |info| info.index);
606                let events = crate::grpc::parse_subscribe_update_transaction_low_latency(
607                    &tx,
608                    grpc_us,
609                    Some(block_us),
610                    filter.as_ref(),
611                );
612                for event in slot_buf.push_streaming_batch(slot, idx, events) {
613                    self.push_queue(queue, event);
614                }
615            }
616            OrderMode::MicroBatch => {
617                for (idx, e) in
618                    parse_transaction_to_vec(&tx, grpc_us, Some(block_us), filter.as_ref())
619                {
620                    if micro_buf.push(slot, idx, e, grpc_us, batch_us) {
621                        for evt in micro_buf.flush() {
622                            self.push_queue(queue, evt);
623                        }
624                    }
625                }
626            }
627        }
628    }
629
630    #[inline]
631    fn handle_account(
632        &self,
633        acc: SubscribeUpdateAccount,
634        filter: &Option<EventTypeFilter>,
635        queue: &Arc<ArrayQueue<DexEvent>>,
636        grpc_us: i64,
637        block_us: i64,
638    ) {
639        let Some(info) = acc.account else { return };
640        // Malformed identities must not become default keys in subscription caches.
641        if info.pubkey.len() != 32 || info.owner.len() != 32 {
642            self.health.dropped();
643            return;
644        }
645        let data = crate::accounts::AccountData {
646            pubkey: read_pubkey_fast(&info.pubkey),
647            executable: info.executable,
648            lamports: info.lamports,
649            owner: read_pubkey_fast(&info.owner),
650            rent_epoch: info.rent_epoch,
651            data: info.data,
652        };
653        let meta = EventMetadata {
654            signature: Default::default(),
655            slot: acc.slot,
656            tx_index: 0,
657            block_time_us: block_us,
658            grpc_recv_us: grpc_us,
659            recent_blockhash: None,
660        };
661        // Parse while borrowing bytes, then move them into the opt-in raw event.
662        // This also avoids a full account-data clone when both outputs are requested.
663        let normalized = if data.lamports == 0 {
664            // A closed account can retain its former bytes in a Geyser update.
665            // Preserve the tombstone in raw output, never emit it as live state.
666            None
667        } else {
668            crate::accounts::parse_account_unified(&data, meta.clone(), filter.as_ref())
669        };
670        if filter.as_ref().is_some_and(|f| {
671            f.include_only
672                .as_ref()
673                .is_some_and(|types| types.contains(&crate::grpc::EventType::AccountRawSnapshot))
674        }) {
675            self.push_queue(
676                queue,
677                DexEvent::RawAccountSnapshot(Box::new(
678                    crate::accounts::liquidity_snapshot::RawAccountSnapshotEvent {
679                        metadata: meta,
680                        account: data,
681                        write_version: info.write_version,
682                        is_startup: acc.is_startup,
683                    },
684                )),
685            );
686        }
687        if let Some(e) = normalized {
688            self.push_queue(queue, e);
689        }
690    }
691
692    #[inline]
693    fn handle_block_meta(
694        &self,
695        block_meta: SubscribeUpdateBlockMeta,
696        filter: &Option<EventTypeFilter>,
697        queue: &Arc<ArrayQueue<DexEvent>>,
698        grpc_us: i64,
699        fallback_block_us: i64,
700    ) {
701        let block_time_us = block_meta
702            .block_time
703            .as_ref()
704            .map(|t| t.timestamp.saturating_mul(1_000_000))
705            .unwrap_or(fallback_block_us);
706        let event = DexEvent::BlockMeta(crate::core::events::BlockMetaEvent {
707            metadata: EventMetadata {
708                signature: Default::default(),
709                slot: block_meta.slot,
710                tx_index: 0,
711                block_time_us,
712                grpc_recv_us: grpc_us,
713                recent_blockhash: (!block_meta.blockhash.is_empty())
714                    .then_some(block_meta.blockhash),
715            },
716        });
717        if filter.as_ref().map(|f| f.should_include_dex_event(&event)).unwrap_or(true) {
718            self.push_queue(queue, event);
719        }
720    }
721}
722
723// ==================== 辅助函数 ====================
724
725/// 获取当前时间戳(微秒)
726///
727/// 使用高性能时钟,避免系统调用开销
728///
729/// # 性能优势
730/// - 旧实现:使用 libc::clock_gettime,每次调用约 1-2μs
731/// - 新实现:使用高性能时钟,每次调用约 10-50ns
732/// - 性能提升:20-100 倍
733#[inline(always)]
734fn get_timestamp_us() -> i64 {
735    now_micros()
736}
737
738// ==================== 交易解析 ====================
739
740#[inline]
741fn parse_transaction_to_vec(
742    tx: &SubscribeUpdateTransaction,
743    grpc_us: i64,
744    block_us: Option<i64>,
745    filter: Option<&EventTypeFilter>,
746) -> Vec<(u64, DexEvent)> {
747    let idx = tx.transaction.as_ref().map(|t| t.index).unwrap_or(0);
748    crate::grpc::parse_subscribe_update_transaction_low_latency(tx, grpc_us, block_us, filter)
749        .into_iter()
750        .map(|event| (idx, event))
751        .collect()
752}
753
754#[cfg(test)]
755mod tests {
756    use super::*;
757    #[test]
758    fn closed_share_retained_bytes_emit_only_raw_tombstone() {
759        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
760        let mut bytes = vec![0; 145];
761        bytes[..8]
762            .copy_from_slice(crate::accounts::raydium_cpmm::discriminators::CREATOR_FEE_SHARE);
763        bytes[73..81].copy_from_slice(&300_000u64.to_le_bytes());
764        for raw in [false, true] {
765            let queue = Arc::new(ArrayQueue::new(4));
766            let mut types = vec![EventType::AccountRaydiumCpmmCreatorFeeShare];
767            if raw {
768                types.push(EventType::AccountRawSnapshot);
769            }
770            let filter = Some(EventTypeFilter::include_only(types));
771            let mut update = SubscribeUpdateAccount {
772                account: Some(SubscribeUpdateAccountInfo {
773                    pubkey: solana_sdk::pubkey::Pubkey::new_unique().to_bytes().to_vec(),
774                    owner: crate::instr::program_ids::RAYDIUM_CPMM_PROGRAM_ID.to_bytes().to_vec(),
775                    data: bytes.clone(),
776                    lamports: 1,
777                    write_version: 1,
778                    ..Default::default()
779                }),
780                slot: 100,
781                ..Default::default()
782            };
783            grpc.handle_account(update.clone(), &filter, &queue, 1, 2);
784            if raw {
785                assert!(matches!(queue.pop(), Some(DexEvent::RawAccountSnapshot(_))));
786            }
787            assert!(matches!(queue.pop(), Some(DexEvent::RaydiumCpmmCreatorFeeShareAccount(_))));
788            update.slot = 101;
789            let info = update.account.as_mut().unwrap();
790            info.lamports = 0;
791            info.write_version = 2;
792            grpc.handle_account(update, &filter, &queue, 3, 4);
793            if raw {
794                let DexEvent::RawAccountSnapshot(closed) = queue.pop().unwrap() else {
795                    panic!("tombstone")
796                };
797                assert_eq!(closed.account.lamports, 0);
798                assert_eq!(closed.account.data, bytes);
799                assert_eq!((closed.metadata.slot, closed.write_version), (101, 2));
800            }
801            assert!(queue.is_empty(), "closed account was decoded as live state");
802            assert_eq!(grpc.subscription_status().dropped_events, 0);
803        }
804    }
805    #[test]
806    fn reconnect_backoff_honors_low_latency_config_and_resets_after_recovery() {
807        let base = ClientConfig::low_latency().retry_delay_ms;
808        let mut backoff = ReconnectBackoff::new(base);
809        assert_eq!(backoff.next_delay(false), Duration::from_millis(base));
810        assert_eq!(backoff.next_delay(false), Duration::from_millis(base * 2));
811        for _ in 0..20 {
812            backoff.next_delay(false);
813        }
814        assert_eq!(backoff.next_delay(false), Duration::from_secs(60));
815        assert_eq!(backoff.next_delay(true), Duration::from_millis(base));
816        assert_eq!(backoff.next_delay(true), Duration::from_millis(base));
817        assert_eq!(ReconnectBackoff::new(0).next_delay(false), Duration::from_millis(1));
818        assert_eq!(ReconnectBackoff::new(u64::MAX).next_delay(false), Duration::from_secs(60));
819    }
820    #[test]
821    fn raw_snapshot_moves_original_buffer_and_can_coexist_with_normalized_output() {
822        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
823        for both in [false, true] {
824            let mut data = vec![0; 236];
825            data[..8].copy_from_slice(crate::accounts::raydium_cpmm::discriminators::AMM_CONFIG);
826            let original = data.as_ptr();
827            let mut types = vec![EventType::AccountRawSnapshot];
828            if both {
829                types.push(EventType::AccountRaydiumCpmmAmmConfig);
830            }
831            let queue = Arc::new(ArrayQueue::new(4));
832            grpc.handle_account(
833                SubscribeUpdateAccount {
834                    account: Some(SubscribeUpdateAccountInfo {
835                        pubkey: solana_sdk::pubkey::Pubkey::new_unique().to_bytes().to_vec(),
836                        owner: crate::instr::program_ids::RAYDIUM_CPMM_PROGRAM_ID
837                            .to_bytes()
838                            .to_vec(),
839                        data,
840                        lamports: 1,
841                        ..Default::default()
842                    }),
843                    ..Default::default()
844                },
845                &Some(EventTypeFilter::include_only(types)),
846                &queue,
847                1,
848                2,
849            );
850            let DexEvent::RawAccountSnapshot(raw) = queue.pop().unwrap() else { panic!("raw") };
851            assert_eq!(raw.account.data.as_ptr(), original);
852            if both {
853                assert!(matches!(queue.pop(), Some(DexEvent::RaydiumCpmmAmmConfigAccount(_))));
854            }
855            assert!(queue.is_empty());
856        }
857    }
858    #[test]
859    fn malformed_account_identities_invalidate_continuity_without_default_key_events() {
860        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
861        let queue = Arc::new(ArrayQueue::new(4));
862        for (key_len, owner_len) in [(31, 32), (32, 33), (0, 32)] {
863            grpc.handle_account(
864                SubscribeUpdateAccount {
865                    account: Some(SubscribeUpdateAccountInfo {
866                        pubkey: vec![1; key_len],
867                        owner: vec![1; owner_len],
868                        ..Default::default()
869                    }),
870                    ..Default::default()
871                },
872                &Some(EventTypeFilter::include_only(vec![EventType::AccountRawSnapshot])),
873                &queue,
874                1,
875                2,
876            );
877        }
878        assert!(queue.is_empty());
879        assert_eq!(grpc.subscription_status().dropped_events, 3);
880        assert_eq!(grpc.subscription_status().continuity_revision, 3);
881    }
882    #[test]
883    fn raw_snapshots_are_opt_in_and_preserve_version_and_closures() {
884        let key = solana_sdk::pubkey::Pubkey::new_unique();
885        let account = SubscribeUpdateAccount {
886            account: Some(SubscribeUpdateAccountInfo {
887                pubkey: key.to_bytes().to_vec(),
888                owner: solana_sdk::pubkey::Pubkey::default().to_bytes().to_vec(),
889                lamports: 0,
890                data: Vec::new(),
891                write_version: 42,
892                ..Default::default()
893            }),
894            slot: 123,
895            is_startup: true,
896        };
897        let queue = Arc::new(ArrayQueue::new(4));
898        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
899        grpc.handle_account(account.clone(), &None, &queue, 1, 2);
900        assert!(queue.pop().is_none());
901        let filter =
902            Some(EventTypeFilter::include_only(vec![crate::grpc::EventType::AccountRawSnapshot]));
903        grpc.handle_account(account, &filter, &queue, 1, 2);
904        let DexEvent::RawAccountSnapshot(event) = queue.pop().unwrap() else {
905            panic!("raw snapshot")
906        };
907        assert_eq!(event.write_version, 42);
908        assert_eq!(event.metadata.slot, 123);
909        assert_eq!(event.account.pubkey, key);
910        assert_eq!(event.account.lamports, 0);
911        assert!(event.account.data.is_empty() && event.is_startup);
912        assert!(queue.pop().is_none());
913    }
914
915    fn test_event(slot: u64) -> DexEvent {
916        DexEvent::BlockMeta(crate::core::events::BlockMetaEvent {
917            metadata: EventMetadata { slot, ..Default::default() },
918        })
919    }
920
921    #[tokio::test]
922    async fn stop_clears_subscription_state_and_aborts_handle() {
923        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_string(), None).unwrap();
924        let (tx, _rx) = mpsc::channel::<SubscribeRequest>(1);
925        let handle = tokio::spawn(async {
926            std::future::pending::<()>().await;
927        });
928
929        grpc.health.connected();
930        let stop_signal = Arc::new(AtomicBool::new(false));
931        *grpc.control_tx.lock().await = Some(tx);
932        *grpc.subscription_handle.lock().await = Some(handle);
933        *grpc.stop_signal.lock().await = Some(Arc::clone(&stop_signal));
934
935        grpc.stop().await;
936
937        assert!(stop_signal.load(Ordering::SeqCst));
938        assert!(grpc.stop_signal.lock().await.is_none());
939        assert!(grpc.control_tx.lock().await.is_none());
940        assert!(grpc.subscription_handle.lock().await.is_none());
941        assert!(!grpc.subscription_status().connected);
942        assert_eq!(grpc.subscription_status().disconnects, 1);
943    }
944
945    #[test]
946    fn micro_batch_is_flushed_on_disconnect() {
947        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_string(), None).unwrap();
948        let queue = Arc::new(ArrayQueue::new(2));
949        let mut slot_buffer = SlotBuffer::new();
950        let mut micro_batch = MicroBatchBuffer::new();
951        assert!(!micro_batch.push(7, 0, test_event(7), 10, 100));
952
953        grpc.flush_on_disconnect(OrderMode::MicroBatch, &mut slot_buffer, &mut micro_batch, &queue);
954
955        assert!(micro_batch.is_empty());
956        assert!(matches!(queue.pop(), Some(DexEvent::BlockMeta(_))));
957    }
958
959    #[test]
960    fn full_queue_increments_drop_counter() {
961        let queue = ArrayQueue::new(1);
962        push_queue(&queue, test_event(1));
963        let before = GRPC_DROPPED_EVENTS.load(Ordering::Relaxed);
964        push_queue(&queue, test_event(2));
965        assert!(GRPC_DROPPED_EVENTS.load(Ordering::Relaxed) >= before + 1);
966    }
967    #[test]
968    fn continuity_status_is_shared_by_clones_and_loss_is_client_local() {
969        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
970        let clone = grpc.clone();
971        let other = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
972        grpc.health.connected();
973        let initial = grpc.subscription_status();
974        let queue = ArrayQueue::new(1);
975        grpc.push_queue(&queue, test_event(1));
976        assert_eq!(clone.subscription_status(), initial);
977        grpc.push_queue(&queue, test_event(2));
978        let loss = clone.subscription_status();
979        assert_eq!(loss.dropped_events, 1);
980        assert!(loss.continuity_revision > initial.continuity_revision);
981        assert_eq!(other.subscription_status(), GrpcSubscriptionStatus::default());
982        grpc.health.disconnected();
983        let disconnected = clone.subscription_status();
984        assert!(!disconnected.connected);
985        assert_eq!(disconnected.disconnects, 1);
986        grpc.health.disconnected();
987        assert_eq!(clone.subscription_status(), disconnected);
988        grpc.health.connected();
989        let reconnected = clone.subscription_status();
990        assert!(reconnected.connected);
991        assert_eq!(reconnected.generation, 2);
992        assert!(reconnected.continuity_revision > disconnected.continuity_revision);
993    }
994
995    #[test]
996    fn raw_account_overflow_marks_subscription_state_incomplete() {
997        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
998        let queue = Arc::new(ArrayQueue::new(1));
999        queue.push(test_event(1)).unwrap();
1000        let account = SubscribeUpdateAccount {
1001            account: Some(SubscribeUpdateAccountInfo {
1002                pubkey: solana_sdk::pubkey::Pubkey::new_unique().to_bytes().to_vec(),
1003                owner: solana_sdk::pubkey::Pubkey::default().to_bytes().to_vec(),
1004                ..Default::default()
1005            }),
1006            ..Default::default()
1007        };
1008        grpc.handle_account(
1009            account,
1010            &Some(EventTypeFilter::include_only(vec![EventType::AccountRawSnapshot])),
1011            &queue,
1012            0,
1013            0,
1014        );
1015        assert_eq!(grpc.subscription_status().dropped_events, 1);
1016        assert_eq!(queue.len(), 1);
1017    }
1018    #[tokio::test]
1019    async fn full_update_queue_does_not_block_stop_or_replace_reconnect_filters() {
1020        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
1021        *grpc.subscription_filters.lock().await = Some(SubscriptionFilters {
1022            transactions: vec![],
1023            accounts: vec![AccountFilter::new().add_account("old")],
1024            events: Some(EventTypeFilter::include_only(vec![EventType::BlockMeta])),
1025        });
1026        let (tx, _rx) = mpsc::channel(1);
1027        tx.try_send(SubscribeRequest::default()).unwrap();
1028        *grpc.control_tx.lock().await = Some(tx);
1029        let error = tokio::time::timeout(
1030            Duration::from_secs(1),
1031            grpc.update_subscription(vec![], vec![AccountFilter::new().add_account("new")]),
1032        )
1033        .await
1034        .unwrap()
1035        .unwrap_err();
1036        assert!(error.to_string().contains("queue is full"));
1037        assert_eq!(
1038            grpc.subscription_filters.lock().await.as_ref().unwrap().accounts[0].account,
1039            vec!["old"]
1040        );
1041        tokio::time::timeout(Duration::from_secs(1), grpc.stop()).await.unwrap();
1042        assert!(grpc.subscription_filters.lock().await.is_none());
1043    }
1044}