1use 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}; use 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};
32use 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#[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 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 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 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 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 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 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 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 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 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#[inline(always)]
549fn get_timestamp_us() -> i64 {
550 now_micros()
551}
552
553#[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 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}