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