1use super::buffers::{MicroBatchBuffer, SlotBuffer};
10use super::subscribe_builder::build_subscribe_request_with_event_filter;
11use super::types::*;
12use crate::core::{now_micros, EventMetadata}; use 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};
24use 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#[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#[derive(Clone)]
98struct SubscriptionFilters {
99 transactions: Vec<TransactionFilter>,
100 accounts: Vec<AccountFilter>,
101 events: Option<EventTypeFilter>,
102}
103
104struct 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 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 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 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 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 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 self.health.invalidate();
276 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 self.health.disconnected();
304 self.subscription_filters.lock().await.take();
305 }
306
307 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 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 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 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 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 _ = 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 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 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 let normalized = if data.lamports == 0 {
661 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#[inline(always)]
731fn get_timestamp_us() -> i64 {
732 now_micros()
733}
734
735#[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}