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
84pub 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 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 #[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}