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 = 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                for (idx, e) in
606                    parse_transaction_to_vec(&tx, grpc_us, Some(block_us), filter.as_ref())
607                {
608                    for evt in slot_buf.push_streaming(slot, idx, e) {
609                        self.push_queue(queue, evt);
610                    }
611                }
612            }
613            OrderMode::MicroBatch => {
614                for (idx, e) in
615                    parse_transaction_to_vec(&tx, grpc_us, Some(block_us), filter.as_ref())
616                {
617                    if micro_buf.push(slot, idx, e, grpc_us, batch_us) {
618                        for evt in micro_buf.flush() {
619                            self.push_queue(queue, evt);
620                        }
621                    }
622                }
623            }
624        }
625    }
626
627    #[inline]
628    fn handle_account(
629        &self,
630        acc: SubscribeUpdateAccount,
631        filter: &Option<EventTypeFilter>,
632        queue: &Arc<ArrayQueue<DexEvent>>,
633        grpc_us: i64,
634        block_us: i64,
635    ) {
636        let Some(info) = acc.account else { return };
637        // Malformed identities must not become default keys in subscription caches.
638        if info.pubkey.len() != 32 || info.owner.len() != 32 {
639            self.health.dropped();
640            return;
641        }
642        let data = crate::accounts::AccountData {
643            pubkey: read_pubkey_fast(&info.pubkey),
644            executable: info.executable,
645            lamports: info.lamports,
646            owner: read_pubkey_fast(&info.owner),
647            rent_epoch: info.rent_epoch,
648            data: info.data,
649        };
650        let meta = EventMetadata {
651            signature: Default::default(),
652            slot: acc.slot,
653            tx_index: 0,
654            block_time_us: block_us,
655            grpc_recv_us: grpc_us,
656            recent_blockhash: None,
657        };
658        // Parse while borrowing bytes, then move them into the opt-in raw event.
659        // This also avoids a full account-data clone when both outputs are requested.
660        let normalized = if data.lamports == 0 {
661            // A closed account can retain its former bytes in a Geyser update.
662            // Preserve the tombstone in raw output, never emit it as live state.
663            None
664        } else {
665            crate::accounts::parse_account_unified(&data, meta.clone(), filter.as_ref())
666        };
667        if filter.as_ref().is_some_and(|f| {
668            f.include_only
669                .as_ref()
670                .is_some_and(|types| types.contains(&crate::grpc::EventType::AccountRawSnapshot))
671        }) {
672            self.push_queue(
673                queue,
674                DexEvent::RawAccountSnapshot(Box::new(
675                    crate::accounts::liquidity_snapshot::RawAccountSnapshotEvent {
676                        metadata: meta,
677                        account: data,
678                        write_version: info.write_version,
679                        is_startup: acc.is_startup,
680                    },
681                )),
682            );
683        }
684        if let Some(e) = normalized {
685            self.push_queue(queue, e);
686        }
687    }
688
689    #[inline]
690    fn handle_block_meta(
691        &self,
692        block_meta: SubscribeUpdateBlockMeta,
693        filter: &Option<EventTypeFilter>,
694        queue: &Arc<ArrayQueue<DexEvent>>,
695        grpc_us: i64,
696        fallback_block_us: i64,
697    ) {
698        let block_time_us = block_meta
699            .block_time
700            .as_ref()
701            .map(|t| t.timestamp.saturating_mul(1_000_000))
702            .unwrap_or(fallback_block_us);
703        let event = DexEvent::BlockMeta(crate::core::events::BlockMetaEvent {
704            metadata: EventMetadata {
705                signature: Default::default(),
706                slot: block_meta.slot,
707                tx_index: 0,
708                block_time_us,
709                grpc_recv_us: grpc_us,
710                recent_blockhash: (!block_meta.blockhash.is_empty())
711                    .then_some(block_meta.blockhash),
712            },
713        });
714        if filter.as_ref().map(|f| f.should_include_dex_event(&event)).unwrap_or(true) {
715            self.push_queue(queue, event);
716        }
717    }
718}
719
720// ==================== 辅助函数 ====================
721
722/// 获取当前时间戳(微秒)
723///
724/// 使用高性能时钟,避免系统调用开销
725///
726/// # 性能优势
727/// - 旧实现:使用 libc::clock_gettime,每次调用约 1-2μs
728/// - 新实现:使用高性能时钟,每次调用约 10-50ns
729/// - 性能提升:20-100 倍
730#[inline(always)]
731fn get_timestamp_us() -> i64 {
732    now_micros()
733}
734
735// ==================== 交易解析 ====================
736
737#[inline]
738fn parse_transaction_to_vec(
739    tx: &SubscribeUpdateTransaction,
740    grpc_us: i64,
741    block_us: Option<i64>,
742    filter: Option<&EventTypeFilter>,
743) -> Vec<(u64, DexEvent)> {
744    let idx = tx.transaction.as_ref().map(|t| t.index).unwrap_or(0);
745    crate::grpc::parse_subscribe_update_transaction_low_latency(tx, grpc_us, block_us, filter)
746        .into_iter()
747        .map(|event| (idx, event))
748        .collect()
749}
750
751#[cfg(test)]
752mod tests {
753    use super::*;
754    #[test]
755    fn closed_share_retained_bytes_emit_only_raw_tombstone() {
756        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
757        let mut bytes = vec![0; 145];
758        bytes[..8]
759            .copy_from_slice(crate::accounts::raydium_cpmm::discriminators::CREATOR_FEE_SHARE);
760        bytes[73..81].copy_from_slice(&300_000u64.to_le_bytes());
761        for raw in [false, true] {
762            let queue = Arc::new(ArrayQueue::new(4));
763            let mut types = vec![EventType::AccountRaydiumCpmmCreatorFeeShare];
764            if raw {
765                types.push(EventType::AccountRawSnapshot);
766            }
767            let filter = Some(EventTypeFilter::include_only(types));
768            let mut update = SubscribeUpdateAccount {
769                account: Some(SubscribeUpdateAccountInfo {
770                    pubkey: solana_sdk::pubkey::Pubkey::new_unique().to_bytes().to_vec(),
771                    owner: crate::instr::program_ids::RAYDIUM_CPMM_PROGRAM_ID.to_bytes().to_vec(),
772                    data: bytes.clone(),
773                    lamports: 1,
774                    write_version: 1,
775                    ..Default::default()
776                }),
777                slot: 100,
778                ..Default::default()
779            };
780            grpc.handle_account(update.clone(), &filter, &queue, 1, 2);
781            if raw {
782                assert!(matches!(queue.pop(), Some(DexEvent::RawAccountSnapshot(_))));
783            }
784            assert!(matches!(queue.pop(), Some(DexEvent::RaydiumCpmmCreatorFeeShareAccount(_))));
785            update.slot = 101;
786            let info = update.account.as_mut().unwrap();
787            info.lamports = 0;
788            info.write_version = 2;
789            grpc.handle_account(update, &filter, &queue, 3, 4);
790            if raw {
791                let DexEvent::RawAccountSnapshot(closed) = queue.pop().unwrap() else {
792                    panic!("tombstone")
793                };
794                assert_eq!(closed.account.lamports, 0);
795                assert_eq!(closed.account.data, bytes);
796                assert_eq!((closed.metadata.slot, closed.write_version), (101, 2));
797            }
798            assert!(queue.is_empty(), "closed account was decoded as live state");
799            assert_eq!(grpc.subscription_status().dropped_events, 0);
800        }
801    }
802    #[test]
803    fn reconnect_backoff_honors_low_latency_config_and_resets_after_recovery() {
804        let base = ClientConfig::low_latency().retry_delay_ms;
805        let mut backoff = ReconnectBackoff::new(base);
806        assert_eq!(backoff.next_delay(false), Duration::from_millis(base));
807        assert_eq!(backoff.next_delay(false), Duration::from_millis(base * 2));
808        for _ in 0..20 {
809            backoff.next_delay(false);
810        }
811        assert_eq!(backoff.next_delay(false), Duration::from_secs(60));
812        assert_eq!(backoff.next_delay(true), Duration::from_millis(base));
813        assert_eq!(backoff.next_delay(true), Duration::from_millis(base));
814        assert_eq!(ReconnectBackoff::new(0).next_delay(false), Duration::from_millis(1));
815        assert_eq!(ReconnectBackoff::new(u64::MAX).next_delay(false), Duration::from_secs(60));
816    }
817    #[test]
818    fn raw_snapshot_moves_original_buffer_and_can_coexist_with_normalized_output() {
819        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
820        for both in [false, true] {
821            let mut data = vec![0; 236];
822            data[..8].copy_from_slice(crate::accounts::raydium_cpmm::discriminators::AMM_CONFIG);
823            let original = data.as_ptr();
824            let mut types = vec![EventType::AccountRawSnapshot];
825            if both {
826                types.push(EventType::AccountRaydiumCpmmAmmConfig);
827            }
828            let queue = Arc::new(ArrayQueue::new(4));
829            grpc.handle_account(
830                SubscribeUpdateAccount {
831                    account: Some(SubscribeUpdateAccountInfo {
832                        pubkey: solana_sdk::pubkey::Pubkey::new_unique().to_bytes().to_vec(),
833                        owner: crate::instr::program_ids::RAYDIUM_CPMM_PROGRAM_ID
834                            .to_bytes()
835                            .to_vec(),
836                        data,
837                        lamports: 1,
838                        ..Default::default()
839                    }),
840                    ..Default::default()
841                },
842                &Some(EventTypeFilter::include_only(types)),
843                &queue,
844                1,
845                2,
846            );
847            let DexEvent::RawAccountSnapshot(raw) = queue.pop().unwrap() else { panic!("raw") };
848            assert_eq!(raw.account.data.as_ptr(), original);
849            if both {
850                assert!(matches!(queue.pop(), Some(DexEvent::RaydiumCpmmAmmConfigAccount(_))));
851            }
852            assert!(queue.is_empty());
853        }
854    }
855    #[test]
856    fn malformed_account_identities_invalidate_continuity_without_default_key_events() {
857        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
858        let queue = Arc::new(ArrayQueue::new(4));
859        for (key_len, owner_len) in [(31, 32), (32, 33), (0, 32)] {
860            grpc.handle_account(
861                SubscribeUpdateAccount {
862                    account: Some(SubscribeUpdateAccountInfo {
863                        pubkey: vec![1; key_len],
864                        owner: vec![1; owner_len],
865                        ..Default::default()
866                    }),
867                    ..Default::default()
868                },
869                &Some(EventTypeFilter::include_only(vec![EventType::AccountRawSnapshot])),
870                &queue,
871                1,
872                2,
873            );
874        }
875        assert!(queue.is_empty());
876        assert_eq!(grpc.subscription_status().dropped_events, 3);
877        assert_eq!(grpc.subscription_status().continuity_revision, 3);
878    }
879    #[test]
880    fn raw_snapshots_are_opt_in_and_preserve_version_and_closures() {
881        let key = solana_sdk::pubkey::Pubkey::new_unique();
882        let account = SubscribeUpdateAccount {
883            account: Some(SubscribeUpdateAccountInfo {
884                pubkey: key.to_bytes().to_vec(),
885                owner: solana_sdk::pubkey::Pubkey::default().to_bytes().to_vec(),
886                lamports: 0,
887                data: Vec::new(),
888                write_version: 42,
889                ..Default::default()
890            }),
891            slot: 123,
892            is_startup: true,
893        };
894        let queue = Arc::new(ArrayQueue::new(4));
895        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
896        grpc.handle_account(account.clone(), &None, &queue, 1, 2);
897        assert!(queue.pop().is_none());
898        let filter =
899            Some(EventTypeFilter::include_only(vec![crate::grpc::EventType::AccountRawSnapshot]));
900        grpc.handle_account(account, &filter, &queue, 1, 2);
901        let DexEvent::RawAccountSnapshot(event) = queue.pop().unwrap() else {
902            panic!("raw snapshot")
903        };
904        assert_eq!(event.write_version, 42);
905        assert_eq!(event.metadata.slot, 123);
906        assert_eq!(event.account.pubkey, key);
907        assert_eq!(event.account.lamports, 0);
908        assert!(event.account.data.is_empty() && event.is_startup);
909        assert!(queue.pop().is_none());
910    }
911
912    fn test_event(slot: u64) -> DexEvent {
913        DexEvent::BlockMeta(crate::core::events::BlockMetaEvent {
914            metadata: EventMetadata { slot, ..Default::default() },
915        })
916    }
917
918    #[tokio::test]
919    async fn stop_clears_subscription_state_and_aborts_handle() {
920        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_string(), None).unwrap();
921        let (tx, _rx) = mpsc::channel::<SubscribeRequest>(1);
922        let handle = tokio::spawn(async {
923            std::future::pending::<()>().await;
924        });
925
926        grpc.health.connected();
927        let stop_signal = Arc::new(AtomicBool::new(false));
928        *grpc.control_tx.lock().await = Some(tx);
929        *grpc.subscription_handle.lock().await = Some(handle);
930        *grpc.stop_signal.lock().await = Some(Arc::clone(&stop_signal));
931
932        grpc.stop().await;
933
934        assert!(stop_signal.load(Ordering::SeqCst));
935        assert!(grpc.stop_signal.lock().await.is_none());
936        assert!(grpc.control_tx.lock().await.is_none());
937        assert!(grpc.subscription_handle.lock().await.is_none());
938        assert!(!grpc.subscription_status().connected);
939        assert_eq!(grpc.subscription_status().disconnects, 1);
940    }
941
942    #[test]
943    fn micro_batch_is_flushed_on_disconnect() {
944        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_string(), None).unwrap();
945        let queue = Arc::new(ArrayQueue::new(2));
946        let mut slot_buffer = SlotBuffer::new();
947        let mut micro_batch = MicroBatchBuffer::new();
948        assert!(!micro_batch.push(7, 0, test_event(7), 10, 100));
949
950        grpc.flush_on_disconnect(OrderMode::MicroBatch, &mut slot_buffer, &mut micro_batch, &queue);
951
952        assert!(micro_batch.is_empty());
953        assert!(matches!(queue.pop(), Some(DexEvent::BlockMeta(_))));
954    }
955
956    #[test]
957    fn full_queue_increments_drop_counter() {
958        let queue = ArrayQueue::new(1);
959        push_queue(&queue, test_event(1));
960        let before = GRPC_DROPPED_EVENTS.load(Ordering::Relaxed);
961        push_queue(&queue, test_event(2));
962        assert!(GRPC_DROPPED_EVENTS.load(Ordering::Relaxed) >= before + 1);
963    }
964    #[test]
965    fn continuity_status_is_shared_by_clones_and_loss_is_client_local() {
966        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
967        let clone = grpc.clone();
968        let other = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
969        grpc.health.connected();
970        let initial = grpc.subscription_status();
971        let queue = ArrayQueue::new(1);
972        grpc.push_queue(&queue, test_event(1));
973        assert_eq!(clone.subscription_status(), initial);
974        grpc.push_queue(&queue, test_event(2));
975        let loss = clone.subscription_status();
976        assert_eq!(loss.dropped_events, 1);
977        assert!(loss.continuity_revision > initial.continuity_revision);
978        assert_eq!(other.subscription_status(), GrpcSubscriptionStatus::default());
979        grpc.health.disconnected();
980        let disconnected = clone.subscription_status();
981        assert!(!disconnected.connected);
982        assert_eq!(disconnected.disconnects, 1);
983        grpc.health.disconnected();
984        assert_eq!(clone.subscription_status(), disconnected);
985        grpc.health.connected();
986        let reconnected = clone.subscription_status();
987        assert!(reconnected.connected);
988        assert_eq!(reconnected.generation, 2);
989        assert!(reconnected.continuity_revision > disconnected.continuity_revision);
990    }
991
992    #[test]
993    fn raw_account_overflow_marks_subscription_state_incomplete() {
994        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
995        let queue = Arc::new(ArrayQueue::new(1));
996        queue.push(test_event(1)).unwrap();
997        let account = SubscribeUpdateAccount {
998            account: Some(SubscribeUpdateAccountInfo {
999                pubkey: solana_sdk::pubkey::Pubkey::new_unique().to_bytes().to_vec(),
1000                owner: solana_sdk::pubkey::Pubkey::default().to_bytes().to_vec(),
1001                ..Default::default()
1002            }),
1003            ..Default::default()
1004        };
1005        grpc.handle_account(
1006            account,
1007            &Some(EventTypeFilter::include_only(vec![EventType::AccountRawSnapshot])),
1008            &queue,
1009            0,
1010            0,
1011        );
1012        assert_eq!(grpc.subscription_status().dropped_events, 1);
1013        assert_eq!(queue.len(), 1);
1014    }
1015    #[tokio::test]
1016    async fn full_update_queue_does_not_block_stop_or_replace_reconnect_filters() {
1017        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_owned(), None).unwrap();
1018        *grpc.subscription_filters.lock().await = Some(SubscriptionFilters {
1019            transactions: vec![],
1020            accounts: vec![AccountFilter::new().add_account("old")],
1021            events: Some(EventTypeFilter::include_only(vec![EventType::BlockMeta])),
1022        });
1023        let (tx, _rx) = mpsc::channel(1);
1024        tx.try_send(SubscribeRequest::default()).unwrap();
1025        *grpc.control_tx.lock().await = Some(tx);
1026        let error = tokio::time::timeout(
1027            Duration::from_secs(1),
1028            grpc.update_subscription(vec![], vec![AccountFilter::new().add_account("new")]),
1029        )
1030        .await
1031        .unwrap()
1032        .unwrap_err();
1033        assert!(error.to_string().contains("queue is full"));
1034        assert_eq!(
1035            grpc.subscription_filters.lock().await.as_ref().unwrap().accounts[0].account,
1036            vec!["old"]
1037        );
1038        tokio::time::timeout(Duration::from_secs(1), grpc.stop()).await.unwrap();
1039        assert!(grpc.subscription_filters.lock().await.is_none());
1040    }
1041}