Skip to main content

carbon_yellowstone_grpc_datasource/
lib.rs

1use {
2    async_trait::async_trait,
3    carbon_core::{
4        datasource::{
5            AccountDeletion, AccountUpdate, Datasource, DatasourceDisconnection, DatasourceId,
6            TransactionUpdate, Update, UpdateType,
7        },
8        error::CarbonResult,
9        metrics::{Counter, Histogram, MetricsRegistry},
10    },
11    chrono::{DateTime, Utc},
12    futures::{sink::SinkExt, StreamExt},
13    solana_account::Account,
14    solana_pubkey::Pubkey,
15    solana_signature::Signature,
16    std::{
17        collections::{HashMap, HashSet},
18        convert::TryFrom,
19        sync::{Arc, LazyLock},
20        time::Duration,
21    },
22    tokio::sync::{mpsc, mpsc::Sender, RwLock},
23    tokio_util::sync::CancellationToken,
24    yellowstone_grpc_client::{GeyserGrpcBuilder, GeyserGrpcBuilderResult, GeyserGrpcClient},
25    yellowstone_grpc_proto::{
26        convert_from::{create_tx_meta, create_tx_versioned},
27        geyser::{
28            subscribe_update::UpdateOneof, CommitmentLevel, SubscribeRequest,
29            SubscribeRequestFilterAccounts, SubscribeRequestFilterBlocks,
30            SubscribeRequestFilterTransactions, SubscribeRequestPing, SubscribeUpdateAccountInfo,
31            SubscribeUpdateTransactionInfo,
32        },
33        tonic::{codec::CompressionEncoding, transport::ClientTlsConfig},
34    },
35};
36
37static ACCOUNT_PROCESS_TIME_NANOS: LazyLock<Histogram> = LazyLock::new(|| {
38    Histogram::new(
39        "yellowstone_grpc_account_process_time_nanoseconds",
40        "Time taken to process account updates in nanoseconds",
41        vec![
42            1_000.0,
43            10_000.0,
44            100_000.0,
45            1_000_000.0,
46            10_000_000.0,
47            100_000_000.0,
48            1_000_000_000.0,
49        ],
50    )
51});
52static ACCOUNT_UPDATES_RECEIVED: Counter = Counter::new(
53    "yellowstone_grpc_account_updates_received_total",
54    "Total account updates received from Yellowstone gRPC",
55);
56static TRANSACTION_PROCESS_TIME_NANOS: LazyLock<Histogram> = LazyLock::new(|| {
57    Histogram::new(
58        "yellowstone_grpc_transaction_process_time_nanoseconds",
59        "Time taken to process transaction updates in nanoseconds",
60        vec![
61            1_000.0,
62            10_000.0,
63            100_000.0,
64            1_000_000.0,
65            10_000_000.0,
66            100_000_000.0,
67            1_000_000_000.0,
68        ],
69    )
70});
71static TRANSACTION_UPDATES_RECEIVED: Counter = Counter::new(
72    "yellowstone_grpc_transaction_updates_received_total",
73    "Total transaction updates received from Yellowstone gRPC",
74);
75
76fn register_yellowstone_metrics() {
77    let registry = MetricsRegistry::global();
78    registry.register_counter(&ACCOUNT_UPDATES_RECEIVED);
79    registry.register_counter(&TRANSACTION_UPDATES_RECEIVED);
80    registry.register_histogram(&ACCOUNT_PROCESS_TIME_NANOS);
81    registry.register_histogram(&TRANSACTION_PROCESS_TIME_NANOS);
82}
83
84/// Default timeout for detecting stale connections (30 seconds)
85pub const DEFAULT_STREAM_TIMEOUT_SECS: u64 = 30;
86
87#[derive(Debug)]
88pub struct YellowstoneGrpcGeyserClient {
89    pub endpoint: String,
90    pub x_token: Option<String>,
91    pub commitment: Option<CommitmentLevel>,
92    pub account_filters: HashMap<String, SubscribeRequestFilterAccounts>,
93    pub transaction_filters: HashMap<String, SubscribeRequestFilterTransactions>,
94    pub block_filters: BlockFilters,
95    pub account_deletions_tracked: Arc<RwLock<HashSet<Pubkey>>>,
96    pub geyser_config: YellowstoneGrpcClientConfig,
97    pub disconnect_notifier: Option<mpsc::Sender<DatasourceDisconnection>>,
98    /// Timeout for detecting hung/stale connections. Default: 30 seconds.
99    pub stream_timeout: Duration,
100}
101
102#[derive(Debug, Clone)]
103pub struct YellowstoneGrpcClientConfig {
104    pub compression: Option<CompressionEncoding>,
105    pub connect_timeout: Option<Duration>,
106    pub timeout: Option<Duration>,
107    pub max_decoding_message_size: Option<usize>,
108    pub tls_config: Option<ClientTlsConfig>,
109    pub tcp_nodelay: Option<bool>,
110}
111
112impl Default for YellowstoneGrpcClientConfig {
113    fn default() -> Self {
114        Self {
115            compression: None,
116            connect_timeout: Some(Duration::from_secs(15)),
117            timeout: Some(Duration::from_secs(15)),
118            max_decoding_message_size: None,
119            tls_config: None,
120            tcp_nodelay: None,
121        }
122    }
123}
124
125#[derive(Default, Debug, Clone)]
126pub struct BlockFilters {
127    pub filters: HashMap<String, SubscribeRequestFilterBlocks>,
128    pub failed_transactions: Option<bool>,
129}
130
131impl YellowstoneGrpcGeyserClient {
132    /// Creates a new YellowstoneGrpcGeyserClient with optional stream timeout.
133    /// If `stream_timeout` is None, defaults to 30 seconds.
134    #[allow(clippy::too_many_arguments)]
135    pub fn new(
136        endpoint: String,
137        x_token: Option<String>,
138        commitment: Option<CommitmentLevel>,
139        account_filters: HashMap<String, SubscribeRequestFilterAccounts>,
140        transaction_filters: HashMap<String, SubscribeRequestFilterTransactions>,
141        block_filters: BlockFilters,
142        account_deletions_tracked: Arc<RwLock<HashSet<Pubkey>>>,
143        geyser_config: YellowstoneGrpcClientConfig,
144        disconnect_notifier: Option<mpsc::Sender<DatasourceDisconnection>>,
145        stream_timeout: Option<Duration>,
146    ) -> Self {
147        YellowstoneGrpcGeyserClient {
148            endpoint,
149            x_token,
150            commitment,
151            account_filters,
152            transaction_filters,
153            block_filters,
154            account_deletions_tracked,
155            geyser_config,
156            disconnect_notifier,
157            stream_timeout: stream_timeout
158                .unwrap_or(Duration::from_secs(DEFAULT_STREAM_TIMEOUT_SECS)),
159        }
160    }
161}
162
163impl YellowstoneGrpcClientConfig {
164    pub const fn new(
165        compression: Option<CompressionEncoding>,
166        connect_timeout: Option<Duration>,
167        timeout: Option<Duration>,
168        max_decoding_message_size: Option<usize>,
169        tls_config: Option<ClientTlsConfig>,
170        tcp_nodelay: Option<bool>,
171    ) -> Self {
172        YellowstoneGrpcClientConfig {
173            compression,
174            connect_timeout,
175            timeout,
176            max_decoding_message_size,
177            tls_config,
178            tcp_nodelay,
179        }
180    }
181
182    pub fn geyser_config_builder(
183        &self,
184        mut builder: GeyserGrpcBuilder,
185    ) -> GeyserGrpcBuilderResult<GeyserGrpcBuilder> {
186        builder = builder.connect_timeout(self.connect_timeout.unwrap_or(Duration::from_secs(15)));
187
188        builder = builder.timeout(self.timeout.unwrap_or(Duration::from_secs(15)));
189        let tls = self
190            .tls_config
191            .clone()
192            .unwrap_or_else(|| ClientTlsConfig::new().with_enabled_roots());
193        builder = builder.tls_config(tls)?;
194
195        if let Some(compression) = self.compression {
196            builder = builder
197                .send_compressed(compression)
198                .accept_compressed(compression);
199        }
200        if let Some(val) = self.max_decoding_message_size {
201            builder = builder.max_decoding_message_size(val);
202        }
203
204        if let Some(val) = self.tcp_nodelay {
205            builder = builder.tcp_nodelay(val);
206        }
207        Ok(builder)
208    }
209}
210
211#[async_trait]
212impl Datasource for YellowstoneGrpcGeyserClient {
213    async fn consume(
214        &self,
215        id: DatasourceId,
216        sender: Sender<(Update, DatasourceId)>,
217        cancellation_token: CancellationToken,
218    ) -> CarbonResult<()> {
219        register_yellowstone_metrics();
220        let endpoint = self.endpoint.clone();
221        let x_token = self.x_token.clone();
222        let commitment = self.commitment;
223        let account_filters = self.account_filters.clone();
224        let transaction_filters = self.transaction_filters.clone();
225        let account_deletions_tracked = self.account_deletions_tracked.clone();
226        let BlockFilters {
227            filters,
228            failed_transactions: block_failed_transactions,
229        } = self.block_filters.clone();
230        let retain_block_failed_transactions = block_failed_transactions.unwrap_or(true);
231
232        let builder = GeyserGrpcClient::build_from_shared(endpoint)
233            .map_err(|err| carbon_core::error::Error::FailedToConsumeDatasource(err.to_string()))?
234            .x_token(x_token)
235            .map_err(|err| carbon_core::error::Error::FailedToConsumeDatasource(err.to_string()))?;
236
237        let mut geyser_client = self
238            .geyser_config
239            .geyser_config_builder(builder)
240            .map_err(|err| carbon_core::error::Error::FailedToConsumeDatasource(err.to_string()))?
241            .connect()
242            .await
243            .map_err(|err| carbon_core::error::Error::FailedToConsumeDatasource(err.to_string()))?;
244
245        let disconnect_tx_clone = self.disconnect_notifier.clone();
246        let stream_timeout = self.stream_timeout;
247
248        tokio::spawn(async move {
249            let subscribe_request = SubscribeRequest {
250                slots: HashMap::new(),
251                accounts: account_filters,
252                transactions: transaction_filters,
253                transactions_status: HashMap::new(),
254                entry: HashMap::new(),
255                blocks: filters,
256                blocks_meta: HashMap::new(),
257                commitment: commitment.map(|x| x as i32),
258                accounts_data_slice: vec![],
259                ping: None,
260                from_slot: None,
261            };
262
263            let id_for_loop = id.clone();
264
265            let mut last_disconnect_time: Option<DateTime<Utc>> = None;
266            let mut last_slot_before_disconnect: Option<u64> = None;
267            let mut last_processed_slot: u64 = 0;
268
269            loop {
270                tokio::select! {
271                    _ = cancellation_token.cancelled() => {
272                        log::info!("Cancelling Yellowstone gRPC subscription.");
273                        break;
274                    }
275                    result = geyser_client.subscribe_with_request(Some(subscribe_request.clone())) => {
276                        match result {
277                            Ok((mut subscribe_tx, mut stream)) => {
278                                let mut first_message_after_reconnect = last_disconnect_time.is_some();
279
280                                loop {
281                                    if cancellation_token.is_cancelled() {
282                                        break;
283                                    }
284
285                                    let message_result = tokio::time::timeout(
286                                        stream_timeout,
287                                        stream.next()
288                                    ).await;
289
290                                    let message = match message_result {
291                                        Ok(Some(msg)) => msg,
292                                        Ok(None) => {
293                                            log::warn!("Stream closed");
294                                            if last_disconnect_time.is_none() {
295                                                last_disconnect_time = Some(Utc::now());
296                                                last_slot_before_disconnect = Some(last_processed_slot);
297                                                log::warn!("Disconnected at slot {last_processed_slot}");
298                                            }
299                                            break;
300                                        }
301                                        Err(_) => {
302                                            log::warn!("Stream timeout - no messages for {stream_timeout:?}");
303                                            if last_disconnect_time.is_none() {
304                                                last_disconnect_time = Some(Utc::now());
305                                                last_slot_before_disconnect = Some(last_processed_slot);
306                                                log::warn!("Disconnected at slot {last_processed_slot} (timeout)");
307                                            }
308                                            break;
309                                        }
310                                    };
311
312                                    match message {
313                                        Ok(msg) => {
314                                            if first_message_after_reconnect {
315                                                first_message_after_reconnect = false;
316
317                                                let current_slot = match &msg.update_oneof {
318                                                    Some(UpdateOneof::Account(ref update)) => Some(update.slot),
319                                                    Some(UpdateOneof::Transaction(ref update)) => Some(update.slot),
320                                                    Some(UpdateOneof::Block(ref update)) => Some(update.slot),
321                                                    _ => None,
322                                                };
323
324                                                if let Some(slot) = current_slot {
325                                                    if let (Some(disconnect_time), Some(last_slot)) =
326                                                        (last_disconnect_time.take(), last_slot_before_disconnect.take())
327                                                    {
328                                                        let missed = slot.saturating_sub(last_slot);
329
330                                                        let disconnection = DatasourceDisconnection {
331                                                            source: "yellowstone-grpc".to_string(),
332                                                            disconnect_time,
333                                                            last_slot_before_disconnect: last_slot,
334                                                            first_slot_after_reconnect: slot,
335                                                            missed_slots: missed,
336                                                        };
337
338                                                        if let Some(tx) = &disconnect_tx_clone {
339                                                            let _ = tx.try_send(disconnection);
340                                                        }
341
342                                                        log::info!("Reconnected. Slots: {last_slot} -> {slot} (missed: {missed})");
343                                                    }
344                                                }
345                                            }
346
347                                            match msg.update_oneof {
348                                            Some(UpdateOneof::Account(account_update)) => {
349                                                last_processed_slot = account_update.slot;
350                                                send_subscribe_account_update_info(
351                                                    account_update.account,
352                                                    &sender,
353                                                    id_for_loop.clone(),
354                                                    account_update.slot,
355                                                    &account_deletions_tracked,
356                                                )
357                                                .await
358                                            }
359
360                                            Some(UpdateOneof::Transaction(transaction_update)) => {
361                                                last_processed_slot = transaction_update.slot;
362                                                send_subscribe_update_transaction_info(
363                                                    transaction_update.transaction,
364                                                    &sender,
365                                                    id_for_loop.clone(),
366                                                    transaction_update.slot,
367                                                    None,
368                                                )
369                                                .await
370                                            }
371                                            Some(UpdateOneof::Block(block_update)) => {
372                                                last_processed_slot = block_update.slot;
373                                                let block_time = block_update.block_time.map(|ts| ts.timestamp);
374
375                                                for transaction_update in block_update.transactions {
376                                                    if retain_block_failed_transactions || transaction_update.meta.as_ref().map(|meta| meta.err.is_none()).unwrap_or(false) {
377                                                        send_subscribe_update_transaction_info(Some(transaction_update), &sender, id_for_loop.clone(), block_update.slot, block_time).await
378                                                    }
379                                                }
380
381                                                for account_info in block_update.accounts {
382                                                    send_subscribe_account_update_info(
383                                                        Some(account_info),
384                                                        &sender,
385                                                        id_for_loop.clone(),
386                                                        block_update.slot,
387                                                        &account_deletions_tracked,
388                                                    )
389                                                    .await;
390                                                }
391                                            }
392
393                                            Some(UpdateOneof::Ping(_)) => {
394                                                match subscribe_tx
395                                                    .send(SubscribeRequest {
396                                                        ping: Some(SubscribeRequestPing { id: 1 }),
397                                                        ..Default::default()
398                                                    })
399                                                    .await {
400                                                        Ok(()) => (),
401                                                        Err(error) => {
402                                                            log::error!("Failed to send ping error: {error:?}");
403                                                            break;
404                                                        },
405                                                    }
406                                            }
407
408                                            _ => {}
409                                        }
410                                        }
411                                        Err(error) => {
412                                            log::error!("Geyser stream error: {error:?}");
413
414                                            if last_disconnect_time.is_none() {
415                                                last_disconnect_time = Some(Utc::now());
416                                                last_slot_before_disconnect = Some(last_processed_slot);
417                                                log::error!("Disconnected at slot {last_processed_slot}");
418                                            }
419
420                                            break;
421                                        }
422                                    }
423                                }
424                            }
425                            Err(e) => {
426                                log::error!("Failed to subscribe: {e:?}");
427
428                                if last_disconnect_time.is_none() {
429                                    last_disconnect_time = Some(Utc::now());
430                                    last_slot_before_disconnect = Some(last_processed_slot);
431                                }
432
433                            }
434                        }
435                    }
436                }
437            }
438        });
439
440        Ok(())
441    }
442
443    fn update_types(&self) -> Vec<UpdateType> {
444        vec![
445            UpdateType::AccountUpdate,
446            UpdateType::Transaction,
447            UpdateType::AccountDeletion,
448        ]
449    }
450}
451
452async fn send_subscribe_account_update_info(
453    account_update_info: Option<SubscribeUpdateAccountInfo>,
454    sender: &Sender<(Update, DatasourceId)>,
455    id: DatasourceId,
456    slot: u64,
457    account_deletions_tracked: &RwLock<HashSet<Pubkey>>,
458) {
459    let start_time = std::time::Instant::now();
460
461    if let Some(account_info) = account_update_info {
462        let Ok(account_pubkey) = Pubkey::try_from(account_info.pubkey) else {
463            return;
464        };
465
466        let Ok(account_owner_pubkey) = Pubkey::try_from(account_info.owner) else {
467            return;
468        };
469
470        let account = Account {
471            lamports: account_info.lamports,
472            data: account_info.data,
473            owner: account_owner_pubkey,
474            executable: account_info.executable,
475            rent_epoch: account_info.rent_epoch,
476        };
477
478        if account.lamports == 0
479            && account.data.is_empty()
480            && account_owner_pubkey == solana_system_interface::program::ID
481        {
482            let accounts = account_deletions_tracked.read().await;
483            if accounts.contains(&account_pubkey) {
484                let account_deletion = AccountDeletion {
485                    pubkey: account_pubkey,
486                    slot,
487                    transaction_signature: account_info
488                        .txn_signature
489                        .and_then(|sig| Signature::try_from(sig).ok()),
490                };
491                if let Err(e) = sender.try_send((Update::AccountDeletion(account_deletion), id)) {
492                    log::error!(
493                        "Failed to send account deletion update for pubkey {account_pubkey:?} at slot {slot}: {e:?}"
494                    );
495                }
496            }
497        } else {
498            let update = Update::Account(AccountUpdate {
499                pubkey: account_pubkey,
500                account,
501                slot,
502                transaction_signature: account_info
503                    .txn_signature
504                    .and_then(|sig| Signature::try_from(sig).ok()),
505            });
506
507            if let Err(e) = sender.try_send((update, id)) {
508                log::error!(
509                    "Failed to send account update for pubkey {account_pubkey:?} at slot {slot}: {e:?}"
510                );
511            }
512        }
513
514        ACCOUNT_PROCESS_TIME_NANOS.record(start_time.elapsed().as_nanos() as f64);
515        ACCOUNT_UPDATES_RECEIVED.inc();
516    } else {
517        log::error!("No account info in UpdateOneof::Account at slot {slot}");
518    }
519}
520
521async fn send_subscribe_update_transaction_info(
522    transaction_info: Option<SubscribeUpdateTransactionInfo>,
523    sender: &Sender<(Update, DatasourceId)>,
524    id: DatasourceId,
525    slot: u64,
526    block_time: Option<i64>,
527) {
528    let start_time = std::time::Instant::now();
529
530    if let Some(transaction_info) = transaction_info {
531        let Ok(signature) = Signature::try_from(transaction_info.signature) else {
532            return;
533        };
534        let Some(yellowstone_transaction) = transaction_info.transaction else {
535            return;
536        };
537        let Some(yellowstone_tx_meta) = transaction_info.meta else {
538            return;
539        };
540        let Ok(versioned_transaction) = create_tx_versioned(yellowstone_transaction) else {
541            return;
542        };
543        let meta_original = match create_tx_meta(yellowstone_tx_meta) {
544            Ok(meta) => meta,
545            Err(err) => {
546                log::error!("Failed to create transaction meta: {err:?}");
547                return;
548            }
549        };
550        let update = Update::Transaction(Box::new(TransactionUpdate {
551            signature,
552            transaction: versioned_transaction,
553            meta: meta_original,
554            is_vote: transaction_info.is_vote,
555            slot,
556            index: Some(transaction_info.index),
557            block_time,
558            block_hash: None,
559        }));
560        if let Err(e) = sender.try_send((update, id)) {
561            log::error!(
562                "Failed to send transaction update with signature {signature:?} at slot {slot}: {e:?}"
563            );
564            return;
565        }
566
567        TRANSACTION_PROCESS_TIME_NANOS.record(start_time.elapsed().as_nanos() as f64);
568        TRANSACTION_UPDATES_RECEIVED.inc();
569    } else {
570        log::error!("No transaction info in `UpdateOneof::Transaction` at slot {slot}");
571    }
572}