Skip to main content

iota_graphql_rpc/context_data/
db_data_provider.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5use 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    /// Create a new underlying reader, which is used by this type as well as
37    /// other data providers. If `historic_fallback` carries a URL, an archival
38    /// fallback reader is attached so reads can be served from the REST KV
39    /// store if Postgres data is pruned.
40    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        // Create connection pool
54        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        // Create store and watermark cache for pruning support
58        let store = PgIndexerStore::new(connection_pool.clone(), indexer_metrics);
59        let watermark_cache = WatermarkCache::new();
60
61        // Start watermark task with cancellation token
62        let watermark_task = WatermarkTask::new(store, watermark_cache.clone());
63        watermark_task.start(cancellation_token);
64
65        // Create reader with watermark cache
66        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
98/// Implement methods to be used by graphql resolvers
99impl PgManager {
100    /// If no epoch was requested or if the epoch requested is in progress,
101    /// returns the latest iota system state.
102    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    /// Make a request to the RPC for its representations of the staked iota we
122    /// parsed out of the object.  Used to implement fields that are
123    /// implemented in JSON-RPC but not GraphQL (yet).
124    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}