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::transaction_meta::try_yellowstone_signature;
14use super::types::*;
15use crate::core::{now_micros, EventMetadata}; // 导入高性能时钟
16use crate::instr::read_pubkey_fast;
17use crate::logs::timestamp_to_microseconds;
18use crate::DexEvent;
19use crossbeam_queue::ArrayQueue;
20use futures::{SinkExt, StreamExt};
21use log::error;
22use memchr::memmem;
23use once_cell::sync::Lazy;
24use solana_sdk::pubkey::Pubkey;
25use std::collections::HashMap;
26use std::str::FromStr;
27use std::sync::atomic::{AtomicBool, Ordering};
28use std::sync::Arc;
29use tokio::sync::{mpsc, Mutex};
30use tokio::task::JoinHandle;
31use tokio::time::{Duration, Instant};
32// Note: ClientTlsConfig moved to yellowstone_grpc_client in newer versions
33use yellowstone_grpc_client::{ClientTlsConfig, GeyserGrpcClient};
34use yellowstone_grpc_proto::prelude::*;
35
36static PROGRAM_DATA_FINDER: Lazy<memmem::Finder> =
37    Lazy::new(|| memmem::Finder::new(b"Program data: "));
38
39struct ActiveProgram<'a> {
40    encoded: &'a str,
41    pubkey: Pubkey,
42}
43
44// ==================== YellowstoneGrpc 客户端 ====================
45
46#[derive(Clone)]
47pub struct YellowstoneGrpc {
48    endpoint: String,
49    token: Option<String>,
50    config: ClientConfig,
51    control_tx: Arc<Mutex<Option<mpsc::Sender<SubscribeRequest>>>>,
52    subscription_handle: Arc<Mutex<Option<JoinHandle<()>>>>,
53    subscription_lifecycle: Arc<Mutex<()>>,
54    stop_signal: Arc<Mutex<Option<Arc<AtomicBool>>>>,
55}
56
57impl YellowstoneGrpc {
58    pub fn new(
59        endpoint: String,
60        token: Option<String>,
61    ) -> Result<Self, Box<dyn std::error::Error>> {
62        crate::warmup::warmup_parser();
63        Ok(Self {
64            endpoint,
65            token,
66            config: ClientConfig::default(),
67            control_tx: Arc::new(Mutex::new(None)),
68            subscription_handle: Arc::new(Mutex::new(None)),
69            subscription_lifecycle: Arc::new(Mutex::new(())),
70            stop_signal: Arc::new(Mutex::new(None)),
71        })
72    }
73
74    pub fn new_with_config(
75        endpoint: String,
76        token: Option<String>,
77        config: ClientConfig,
78    ) -> Result<Self, Box<dyn std::error::Error>> {
79        crate::warmup::warmup_parser();
80        Ok(Self {
81            endpoint,
82            token,
83            config,
84            control_tx: Arc::new(Mutex::new(None)),
85            subscription_handle: Arc::new(Mutex::new(None)),
86            subscription_lifecycle: Arc::new(Mutex::new(())),
87            stop_signal: Arc::new(Mutex::new(None)),
88        })
89    }
90
91    /// 订阅 DEX 事件(自动重连)
92    pub async fn subscribe_dex_events(
93        &self,
94        transaction_filters: Vec<TransactionFilter>,
95        account_filters: Vec<AccountFilter>,
96        event_type_filter: Option<EventTypeFilter>,
97    ) -> Result<Arc<ArrayQueue<DexEvent>>, Box<dyn std::error::Error>> {
98        let _lifecycle = self.subscription_lifecycle.lock().await;
99        self.stop_without_lifecycle_lock().await;
100
101        let queue = Arc::new(ArrayQueue::new(self.config.buffer_size.max(1)));
102        let queue_clone = Arc::clone(&queue);
103        let self_clone = self.clone();
104        let stop_signal = Arc::new(AtomicBool::new(false));
105        *self.stop_signal.lock().await = Some(Arc::clone(&stop_signal));
106
107        let handle = tokio::spawn(async move {
108            let mut delay = 1u64;
109            loop {
110                if stop_signal.load(Ordering::SeqCst) {
111                    break;
112                }
113
114                match self_clone
115                    .stream_events(
116                        &transaction_filters,
117                        &account_filters,
118                        &event_type_filter,
119                        &queue_clone,
120                    )
121                    .await
122                {
123                    Ok(_) => delay = 1,
124                    Err(e) => {
125                        if stop_signal.load(Ordering::SeqCst) {
126                            break;
127                        }
128                        error!("Grpc error: {} - retry in {}s", e, delay);
129                    }
130                }
131
132                if stop_signal.load(Ordering::SeqCst) {
133                    break;
134                }
135                tokio::time::sleep(Duration::from_secs(delay)).await;
136                delay = (delay * 2).min(60);
137            }
138        });
139
140        *self.subscription_handle.lock().await = Some(handle);
141        Ok(queue)
142    }
143
144    /// 动态更新订阅过滤器
145    pub async fn update_subscription(
146        &self,
147        transaction_filters: Vec<TransactionFilter>,
148        account_filters: Vec<AccountFilter>,
149    ) -> Result<(), Box<dyn std::error::Error>> {
150        let sender = self.control_tx.lock().await.as_ref().ok_or("No active subscription")?.clone();
151
152        let request = build_subscribe_request(&transaction_filters, &account_filters);
153        sender.send(request).await.map_err(|e| e.to_string())?;
154        Ok(())
155    }
156
157    pub async fn stop(&self) {
158        let _lifecycle = self.subscription_lifecycle.lock().await;
159        self.stop_without_lifecycle_lock().await;
160    }
161
162    async fn stop_without_lifecycle_lock(&self) {
163        if let Some(stop_signal) = self.stop_signal.lock().await.take() {
164            stop_signal.store(true, Ordering::SeqCst);
165        }
166        self.control_tx.lock().await.take();
167        let handle = self.subscription_handle.lock().await.take();
168        if let Some(handle) = handle {
169            handle.abort();
170            let _ = handle.await;
171        }
172    }
173
174    // ==================== 核心事件流处理 ====================
175
176    async fn stream_events(
177        &self,
178        tx_filters: &[TransactionFilter],
179        acc_filters: &[AccountFilter],
180        event_filter: &Option<EventTypeFilter>,
181        queue: &Arc<ArrayQueue<DexEvent>>,
182    ) -> Result<(), String> {
183        let _ = rustls::crypto::ring::default_provider().install_default();
184
185        // 构建客户端
186        let mut builder = GeyserGrpcClient::build_from_shared(self.endpoint.clone())
187            .map_err(|e| e.to_string())?
188            .x_token(self.token.clone())
189            .map_err(|e| e.to_string())?
190            .max_decoding_message_size(1024 * 1024 * 1024);
191
192        if self.config.connection_timeout_ms > 0 {
193            builder =
194                builder.connect_timeout(Duration::from_millis(self.config.connection_timeout_ms));
195        }
196        if self.config.enable_tls {
197            builder = builder
198                .tls_config(ClientTlsConfig::new().with_native_roots())
199                .map_err(|e| e.to_string())?;
200        }
201
202        let mut client = builder.connect().await.map_err(|e| e.to_string())?;
203        let request = build_subscribe_request_with_event_filter(
204            tx_filters,
205            acc_filters,
206            event_filter.as_ref(),
207            CommitmentLevel::Processed,
208        );
209
210        let (subscribe_tx, mut stream) =
211            client.subscribe_with_request(Some(request)).await.map_err(|e| e.to_string())?;
212
213        self.print_mode_info();
214
215        // 设置控制通道
216        let (control_tx, mut control_rx) = mpsc::channel::<SubscribeRequest>(100);
217        *self.control_tx.lock().await = Some(control_tx);
218        let subscribe_tx = Arc::new(Mutex::new(subscribe_tx));
219
220        // 初始化缓冲区
221        let mut slot_buffer = SlotBuffer::new();
222        let mut micro_batch = MicroBatchBuffer::new();
223        let mut last_slot = 0u64;
224
225        let order_mode = self.config.order_mode;
226        let timeout_ms = self.config.order_timeout_ms;
227        let batch_us = self.config.micro_batch_us;
228        let check_interval = Duration::from_millis(timeout_ms / 2);
229        let mut next_check = Instant::now() + check_interval;
230
231        loop {
232            // Periodic timeout check for ordered modes and MicroBatch
233            self.check_timeout(
234                order_mode,
235                &mut slot_buffer,
236                &mut micro_batch,
237                queue,
238                timeout_ms,
239                batch_us,
240                &mut next_check,
241                check_interval,
242            );
243
244            tokio::select! {
245                msg = stream.next() => {
246                    match msg {
247                        Some(Ok(update)) => {
248                            // Geyser 会周期性下发 ping;必须在同一 subscribe 流上回写 SubscribeRequest.ping,否则公共节点 / LB 可能 RST_STREAM。
249                            if matches!(
250                                update.update_oneof.as_ref(),
251                                Some(subscribe_update::UpdateOneof::Ping(_))
252                            ) {
253                                if let Err(e) = subscribe_tx
254                                    .lock()
255                                    .await
256                                    .send(SubscribeRequest {
257                                        ping: Some(SubscribeRequestPing { id: 1 }),
258                                        ..Default::default()
259                                    })
260                                    .await
261                                {
262                                    self.control_tx.lock().await.take();
263                                    return Err(e.to_string());
264                                }
265                                continue;
266                            }
267                            self.handle_update(
268                                update, order_mode, event_filter, queue,
269                                &mut slot_buffer, &mut micro_batch, &mut last_slot, batch_us
270                            );
271                        }
272                        Some(Err(e)) => {
273                            error!("Grpc Stream error: {:?}", e);
274                            self.flush_on_disconnect(order_mode, &mut slot_buffer, queue);
275                            self.control_tx.lock().await.take();
276                            return Err(e.to_string());
277                        }
278                        None => {
279                            self.flush_on_disconnect(order_mode, &mut slot_buffer, queue);
280                            self.control_tx.lock().await.take();
281                            return Ok(());
282                        }
283                    }
284                }
285                Some(req) = control_rx.recv() => {
286                    if let Err(e) = subscribe_tx.lock().await.send(req).await {
287                        self.control_tx.lock().await.take();
288                        return Err(e.to_string());
289                    }
290                }
291            }
292        }
293    }
294
295    fn print_mode_info(&self) {
296        match self.config.order_mode {
297            OrderMode::Unordered => println!("✅ Unordered Mode (10-20μs)"),
298            OrderMode::Ordered => {
299                println!("✅ Ordered Mode (timeout={}ms)", self.config.order_timeout_ms)
300            }
301            OrderMode::StreamingOrdered => {
302                println!("✅ StreamingOrdered Mode (timeout={}ms)", self.config.order_timeout_ms)
303            }
304            OrderMode::MicroBatch => {
305                println!("✅ MicroBatch Mode (window={}μs)", self.config.micro_batch_us)
306            }
307        }
308    }
309
310    #[inline]
311    fn check_timeout(
312        &self,
313        mode: OrderMode,
314        slot_buf: &mut SlotBuffer,
315        micro_buf: &mut MicroBatchBuffer,
316        queue: &Arc<ArrayQueue<DexEvent>>,
317        timeout_ms: u64,
318        batch_us: u64,
319        next_check: &mut Instant,
320        interval: Duration,
321    ) {
322        if Instant::now() < *next_check {
323            return;
324        }
325        *next_check = Instant::now() + interval;
326
327        match mode {
328            OrderMode::Ordered => {
329                if slot_buf.should_timeout(timeout_ms) {
330                    for e in slot_buf.flush_all() {
331                        let _ = queue.push(e);
332                    }
333                }
334            }
335            OrderMode::StreamingOrdered => {
336                if slot_buf.should_timeout(timeout_ms) {
337                    for e in slot_buf.flush_streaming_timeout() {
338                        let _ = queue.push(e);
339                    }
340                }
341            }
342            OrderMode::MicroBatch => {
343                // Periodic flush for MicroBatch mode
344                let now_us = get_timestamp_us();
345                if micro_buf.should_flush(now_us, batch_us) {
346                    for e in micro_buf.flush() {
347                        let _ = queue.push(e);
348                    }
349                }
350            }
351            OrderMode::Unordered => {}
352        }
353    }
354
355    fn flush_on_disconnect(
356        &self,
357        mode: OrderMode,
358        buffer: &mut SlotBuffer,
359        queue: &Arc<ArrayQueue<DexEvent>>,
360    ) {
361        if matches!(mode, OrderMode::Ordered | OrderMode::StreamingOrdered) {
362            let events = match mode {
363                OrderMode::StreamingOrdered => buffer.flush_streaming_timeout(),
364                _ => buffer.flush_all(),
365            };
366            for e in events {
367                let _ = queue.push(e);
368            }
369        }
370    }
371
372    #[inline]
373    fn handle_update(
374        &self,
375        update_msg: SubscribeUpdate,
376        mode: OrderMode,
377        filter: &Option<EventTypeFilter>,
378        queue: &Arc<ArrayQueue<DexEvent>>,
379        slot_buf: &mut SlotBuffer,
380        micro_buf: &mut MicroBatchBuffer,
381        last_slot: &mut u64,
382        batch_us: u64,
383    ) {
384        let created_at = update_msg.created_at.unwrap_or_default();
385        let block_time_us = timestamp_to_microseconds(created_at.seconds, created_at.nanos) as i64;
386        let grpc_recv_us = get_timestamp_us();
387
388        let Some(update) = update_msg.update_oneof else { return };
389
390        match update {
391            subscribe_update::UpdateOneof::Transaction(tx) => {
392                self.handle_transaction(
393                    tx,
394                    mode,
395                    filter,
396                    queue,
397                    slot_buf,
398                    micro_buf,
399                    last_slot,
400                    batch_us,
401                    grpc_recv_us,
402                    block_time_us,
403                );
404            }
405            subscribe_update::UpdateOneof::Account(acc) => {
406                Self::handle_account(acc, filter, queue, grpc_recv_us, block_time_us);
407            }
408            subscribe_update::UpdateOneof::BlockMeta(block_meta) => {
409                Self::handle_block_meta(block_meta, filter, queue, grpc_recv_us, block_time_us);
410            }
411            _ => {}
412        }
413    }
414
415    #[inline]
416    fn handle_transaction(
417        &self,
418        tx: SubscribeUpdateTransaction,
419        mode: OrderMode,
420        filter: &Option<EventTypeFilter>,
421        queue: &Arc<ArrayQueue<DexEvent>>,
422        slot_buf: &mut SlotBuffer,
423        micro_buf: &mut MicroBatchBuffer,
424        last_slot: &mut u64,
425        batch_us: u64,
426        grpc_us: i64,
427        block_us: i64,
428    ) {
429        let slot = tx.slot;
430
431        match mode {
432            OrderMode::Unordered => {
433                for e in crate::grpc::parse_subscribe_update_transaction_low_latency(
434                    &tx,
435                    grpc_us,
436                    Some(block_us),
437                    filter.as_ref(),
438                ) {
439                    let _ = queue.push(e);
440                }
441            }
442            OrderMode::Ordered => {
443                if slot > *last_slot && *last_slot > 0 {
444                    for e in slot_buf.flush_before(slot) {
445                        let _ = queue.push(e);
446                    }
447                }
448                *last_slot = slot;
449                for (idx, e) in
450                    parse_transaction_to_vec(&tx, grpc_us, Some(block_us), filter.as_ref())
451                {
452                    slot_buf.push(slot, idx, e);
453                }
454            }
455            OrderMode::StreamingOrdered => {
456                for (idx, e) in
457                    parse_transaction_to_vec(&tx, grpc_us, Some(block_us), filter.as_ref())
458                {
459                    for evt in slot_buf.push_streaming(slot, idx, e) {
460                        let _ = queue.push(evt);
461                    }
462                }
463            }
464            OrderMode::MicroBatch => {
465                for (idx, e) in
466                    parse_transaction_to_vec(&tx, grpc_us, Some(block_us), filter.as_ref())
467                {
468                    if micro_buf.push(slot, idx, e, grpc_us, batch_us) {
469                        for evt in micro_buf.flush() {
470                            let _ = queue.push(evt);
471                        }
472                    }
473                }
474            }
475        }
476    }
477
478    #[inline]
479    fn handle_account(
480        acc: SubscribeUpdateAccount,
481        filter: &Option<EventTypeFilter>,
482        queue: &Arc<ArrayQueue<DexEvent>>,
483        grpc_us: i64,
484        block_us: i64,
485    ) {
486        let Some(info) = acc.account else { return };
487        let data = crate::accounts::AccountData {
488            pubkey: read_pubkey_fast(&info.pubkey),
489            executable: info.executable,
490            lamports: info.lamports,
491            owner: read_pubkey_fast(&info.owner),
492            rent_epoch: info.rent_epoch,
493            data: info.data,
494        };
495        let meta = EventMetadata {
496            signature: Default::default(),
497            slot: acc.slot,
498            tx_index: 0,
499            block_time_us: block_us,
500            grpc_recv_us: grpc_us,
501            recent_blockhash: None,
502        };
503        if let Some(e) = crate::accounts::parse_account_unified(&data, meta, filter.as_ref()) {
504            let _ = queue.push(e);
505        }
506    }
507
508    #[inline]
509    fn handle_block_meta(
510        block_meta: SubscribeUpdateBlockMeta,
511        filter: &Option<EventTypeFilter>,
512        queue: &Arc<ArrayQueue<DexEvent>>,
513        grpc_us: i64,
514        fallback_block_us: i64,
515    ) {
516        let block_time_us = block_meta
517            .block_time
518            .as_ref()
519            .map(|t| t.timestamp.saturating_mul(1_000_000))
520            .unwrap_or(fallback_block_us);
521        let event = DexEvent::BlockMeta(crate::core::events::BlockMetaEvent {
522            metadata: EventMetadata {
523                signature: Default::default(),
524                slot: block_meta.slot,
525                tx_index: 0,
526                block_time_us,
527                grpc_recv_us: grpc_us,
528                recent_blockhash: (!block_meta.blockhash.is_empty())
529                    .then_some(block_meta.blockhash),
530            },
531        });
532        if filter.as_ref().map(|f| f.should_include_dex_event(&event)).unwrap_or(true) {
533            let _ = queue.push(event);
534        }
535    }
536}
537
538// ==================== 辅助函数 ====================
539
540/// 获取当前时间戳(微秒)
541///
542/// 使用高性能时钟,避免系统调用开销
543///
544/// # 性能优势
545/// - 旧实现:使用 libc::clock_gettime,每次调用约 1-2μs
546/// - 新实现:使用高性能时钟,每次调用约 10-50ns
547/// - 性能提升:20-100 倍
548#[inline(always)]
549fn get_timestamp_us() -> i64 {
550    now_micros()
551}
552
553// ==================== 交易解析 ====================
554
555#[inline]
556fn parse_transaction_to_vec(
557    tx: &SubscribeUpdateTransaction,
558    grpc_us: i64,
559    block_us: Option<i64>,
560    filter: Option<&EventTypeFilter>,
561) -> Vec<(u64, DexEvent)> {
562    let idx = tx.transaction.as_ref().map(|t| t.index).unwrap_or(0);
563    parse_transaction_core(tx, grpc_us, block_us, filter).into_iter().map(|e| (idx, e)).collect()
564}
565
566#[inline]
567fn parse_transaction_core(
568    tx: &SubscribeUpdateTransaction,
569    grpc_us: i64,
570    block_us: Option<i64>,
571    filter: Option<&EventTypeFilter>,
572) -> Vec<DexEvent> {
573    let Some(info) = &tx.transaction else { return Vec::new() };
574    let Some(meta) = &info.meta else { return Vec::new() };
575
576    let sig = extract_signature(&info.signature);
577    let slot = tx.slot;
578    let idx = info.index;
579    let needs_pumpfun = filter.map(EventTypeFilter::includes_pumpfun).unwrap_or(true);
580    let is_created_buy =
581        needs_pumpfun && crate::logs::optimized_matcher::detect_pumpfun_create(&meta.log_messages);
582
583    let log_events = parse_logs(
584        meta,
585        &info.transaction,
586        &meta.log_messages,
587        sig,
588        slot,
589        idx,
590        block_us,
591        grpc_us,
592        filter,
593        is_created_buy,
594    );
595    let instr_events = parse_instructions(
596        meta,
597        &info.transaction,
598        sig,
599        slot,
600        idx,
601        block_us,
602        grpc_us,
603        filter,
604        is_created_buy,
605    );
606
607    let mut events =
608        crate::grpc::log_instr_dedup::dedupe_log_instruction_events(log_events, instr_events);
609    crate::grpc::transaction_meta::fill_recent_blockhash(&mut events, &info.transaction);
610    for event in &mut events {
611        crate::core::common_filler::fill_token_balances(event, meta, &info.transaction);
612    }
613    if let Some(filter) = filter {
614        events.into_iter().map(|e| filter.normalize_dex_event(e)).collect()
615    } else {
616        events
617    }
618}
619
620#[inline(always)]
621fn extract_signature(bytes: &[u8]) -> solana_sdk::signature::Signature {
622    try_yellowstone_signature(bytes).expect("yellowstone signature must be 64 bytes")
623}
624
625#[inline]
626fn parse_logs(
627    meta: &TransactionStatusMeta,
628    transaction: &Option<yellowstone_grpc_proto::prelude::Transaction>,
629    logs: &[String],
630    sig: solana_sdk::signature::Signature,
631    slot: u64,
632    tx_idx: u64,
633    block_us: Option<i64>,
634    grpc_us: i64,
635    filter: Option<&EventTypeFilter>,
636    is_created_buy: bool,
637) -> Vec<DexEvent> {
638    let mut outer_idx: i32 = -1;
639    let mut inner_idx: i32 = -1;
640    let mut invokes: HashMap<Pubkey, Vec<(i32, i32)>> = HashMap::with_capacity(8);
641    let mut active_program_stack: Vec<ActiveProgram<'_>> = Vec::with_capacity(8);
642    let mut result = Vec::with_capacity(4);
643
644    for log in logs {
645        if let Some((pid, depth)) = crate::logs::optimized_matcher::parse_invoke_info(log) {
646            if depth == 1 {
647                inner_idx = -1;
648                outer_idx += 1;
649            } else {
650                inner_idx += 1;
651            }
652            let program_id = crate::grpc::program_ids::known_program_id(pid)
653                .or_else(|| Pubkey::from_str(pid).ok());
654            if let Some(pk) = program_id {
655                active_program_stack.truncate(depth.saturating_sub(1));
656                active_program_stack.push(ActiveProgram { encoded: pid, pubkey: pk });
657                if crate::grpc::program_ids::needs_invoke_context(&pk) {
658                    invokes.entry(pk).or_default().push((outer_idx, inner_idx));
659                }
660            }
661        }
662
663        if PROGRAM_DATA_FINDER.find(log.as_bytes()).is_some() {
664            let current_program = active_program_stack.last().map(|active| &active.pubkey);
665            if let Some(mut e) = crate::logs::parse_log_with_program_id(
666                log,
667                sig,
668                slot,
669                tx_idx,
670                block_us,
671                grpc_us,
672                filter,
673                is_created_buy,
674                None,
675                current_program,
676            ) {
677                crate::core::account_dispatcher::fill_accounts_with_owned_keys(
678                    &mut e,
679                    meta,
680                    transaction,
681                    &invokes,
682                );
683                crate::core::common_filler::fill_data(&mut e, meta, transaction, &invokes);
684                result.push(e);
685            }
686        }
687
688        if let Some(pid) = crate::logs::optimized_matcher::parse_program_complete_info(log) {
689            if let Some(pos) = active_program_stack.iter().rposition(|active| active.encoded == pid)
690            {
691                active_program_stack.truncate(pos);
692            }
693        }
694    }
695    result
696}
697
698#[inline]
699fn parse_instructions(
700    meta: &TransactionStatusMeta,
701    transaction: &Option<yellowstone_grpc_proto::prelude::Transaction>,
702    sig: solana_sdk::signature::Signature,
703    slot: u64,
704    tx_idx: u64,
705    block_us: Option<i64>,
706    grpc_us: i64,
707    filter: Option<&EventTypeFilter>,
708    is_created_buy: bool,
709) -> Vec<DexEvent> {
710    // 使用增强的 instruction 解析器
711    // 支持:
712    // - 主指令解析(8字节 discriminator)
713    // - Inner instruction 解析(16字节 discriminator)
714    // - 自动事件合并(instruction + inner instruction)
715    crate::grpc::instruction_parser::parse_instructions_enhanced_with_created_buy(
716        meta,
717        transaction,
718        sig,
719        slot,
720        tx_idx,
721        block_us,
722        grpc_us,
723        filter,
724        is_created_buy,
725    )
726}
727
728#[cfg(test)]
729mod tests {
730    use super::*;
731
732    #[tokio::test]
733    async fn stop_clears_subscription_state_and_aborts_handle() {
734        let grpc = YellowstoneGrpc::new("http://127.0.0.1:1".to_string(), None).unwrap();
735        let (tx, _rx) = mpsc::channel::<SubscribeRequest>(1);
736        let handle = tokio::spawn(async {
737            std::future::pending::<()>().await;
738        });
739
740        let stop_signal = Arc::new(AtomicBool::new(false));
741        *grpc.control_tx.lock().await = Some(tx);
742        *grpc.subscription_handle.lock().await = Some(handle);
743        *grpc.stop_signal.lock().await = Some(Arc::clone(&stop_signal));
744
745        grpc.stop().await;
746
747        assert!(stop_signal.load(Ordering::SeqCst));
748        assert!(grpc.stop_signal.lock().await.is_none());
749        assert!(grpc.control_tx.lock().await.is_none());
750        assert!(grpc.subscription_handle.lock().await.is_none());
751    }
752}