iota_graphql_rpc/context_data/
db_data_provider.rs1use std::time::Duration;
6
7use iota_indexer::{
8 apis::GovernanceReadApi,
9 db::{ConnectionPoolConfig, DbUrl},
10 historical_fallback::HistoricalFallbackReader,
11 metrics::IndexerMetrics,
12 pruning::watermark_task::{WatermarkCache, WatermarkTask},
13 read::IndexerReader,
14 store::PgIndexerStore,
15};
16use iota_json_rpc_types::Stake as RpcStakedIota;
17use iota_types::{
18 governance::StakedIota as NativeStakedIota,
19 iota_system_state::iota_system_state_summary::IotaSystemStateSummary as NativeIotaSystemStateSummary,
20};
21use prometheus_filtered::Registry;
22use tokio_util::sync::CancellationToken;
23use tracing::info;
24
25use crate::{config::HistoricFallbackOptions, error::Error};
26
27pub(crate) struct PgManager {
28 pub inner: IndexerReader,
29}
30
31impl PgManager {
32 pub(crate) fn new(inner: IndexerReader) -> Self {
33 Self { inner }
34 }
35
36 pub(crate) fn reader_with_config(
41 db_url: &DbUrl,
42 pool_size: u32,
43 timeout_ms: u64,
44 indexer_metrics: IndexerMetrics,
45 cancellation_token: CancellationToken,
46 historic_fallback: &HistoricFallbackOptions,
47 registry: &Registry,
48 ) -> Result<IndexerReader, Error> {
49 let mut config = ConnectionPoolConfig::default();
50 config.set_pool_size(pool_size);
51 config.set_statement_timeout(Duration::from_millis(timeout_ms));
52
53 let connection_pool = iota_indexer::db::new_connection_pool(db_url, &config)
55 .map_err(|e| Error::Internal(format!("Failed to create connection pool: {e}")))?;
56
57 let store = PgIndexerStore::new(connection_pool.clone(), indexer_metrics);
59 let watermark_cache = WatermarkCache::new();
60
61 let watermark_task = WatermarkTask::new(store, watermark_cache.clone());
63 watermark_task.start(cancellation_token);
64
65 let mut reader = IndexerReader::new(connection_pool, watermark_cache);
67
68 if let HistoricFallbackOptions {
69 fallback_kv_url: Some(url),
70 fallback_kv_multi_fetch_batch_size,
71 fallback_kv_concurrent_fetches,
72 fallback_kv_cache_size,
73 } = historic_fallback
74 {
75 let fallback = HistoricalFallbackReader::new(
76 url.as_str(),
77 *fallback_kv_cache_size,
78 reader.package_resolver().clone(),
79 *fallback_kv_multi_fetch_batch_size,
80 *fallback_kv_concurrent_fetches,
81 registry,
82 )
83 .map_err(|e| {
84 Error::Internal(format!(
85 "Failed to construct historical fallback reader: {e}"
86 ))
87 })?;
88 info!("HistoricalFallbackReader initialized with URL: {url}");
89 reader.with_fallback_reader(fallback);
90 } else {
91 info!("No config for HistoricalFallbackReader provided, skipping...");
92 }
93
94 Ok(reader)
95 }
96}
97
98impl PgManager {
100 pub(crate) async fn fetch_iota_system_state(
103 &self,
104 epoch_id: Option<u64>,
105 ) -> Result<NativeIotaSystemStateSummary, Error> {
106 let latest_iota_system_state = self
107 .inner
108 .spawn_blocking(move |this| this.get_latest_iota_system_state())
109 .await?;
110
111 if epoch_id.is_none() || epoch_id.is_some_and(|id| id == latest_iota_system_state.epoch()) {
112 Ok(latest_iota_system_state)
113 } else {
114 Ok(self
115 .inner
116 .spawn_blocking(move |this| this.get_epoch_iota_system_state(epoch_id))
117 .await?)
118 }
119 }
120
121 pub(crate) async fn fetch_rpc_staked_iota(
125 &self,
126 stake: NativeStakedIota,
127 ) -> Result<RpcStakedIota, Error> {
128 let governance_api = GovernanceReadApi::new(self.inner.clone());
129
130 let mut delegated_stakes = governance_api
131 .get_delegated_stakes(vec![stake])
132 .await
133 .map_err(|e| Error::Internal(format!("Error fetching delegated stake. {e}")))?;
134
135 let Some(mut delegated_stake) = delegated_stakes.pop() else {
136 return Err(Error::Internal(
137 "Error fetching delegated stake. No pools returned.".to_string(),
138 ));
139 };
140
141 let Some(stake) = delegated_stake.stakes.pop() else {
142 return Err(Error::Internal(
143 "Error fetching delegated stake. No stake in pool.".to_string(),
144 ));
145 };
146
147 Ok(stake)
148 }
149}