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