1use std::collections::{BTreeMap, BTreeSet, HashMap};
6
7use async_graphql::{connection::Connection, dataloader::Loader, *};
8use diesel::{ExpressionMethods, OptionalExtension, QueryDsl, SelectableHelper};
9use fastcrypto::encoding::{Base58, Encoding};
10use iota_indexer::{models::epoch::QueryableEpochInfo, schema::epochs};
11use iota_sdk_types::CheckpointCommitment as EpochCommitment;
12
13use crate::{
14 config::DEFAULT_PAGE_SIZE,
15 connection::ScanConnection,
16 context_data::db_data_provider::PgManager,
17 data::{DataLoader, Db, DbConnection, QueryExecutor},
18 error::Error,
19 server::watermark_task::Watermark,
20 types::{
21 big_int::BigInt,
22 checkpoint::{self, Checkpoint},
23 cursor::Page,
24 date_time::DateTime,
25 protocol_config::ProtocolConfigs,
26 system_state_summary::{NativeStateValidatorInfo, SystemStateSummary},
27 transaction_block::{self, TransactionBlock, TransactionBlockFilter},
28 uint53::UInt53,
29 validator_set::ValidatorSet,
30 },
31};
32
33#[derive(Clone)]
34pub(crate) struct Epoch {
35 pub stored: QueryableEpochInfo,
36 pub checkpoint_viewed_at: u64,
37}
38
39#[derive(Copy, Clone, Hash, Eq, PartialEq, Debug)]
42struct EpochKey {
43 pub epoch_id: u64,
44 pub checkpoint_viewed_at: u64,
45}
46
47#[Object]
55impl Epoch {
56 async fn epoch_id(&self) -> Result<UInt53> {
59 UInt53::try_from(self.stored.epoch as u64).extend()
60 }
61
62 async fn reference_gas_price(&self) -> Option<BigInt> {
65 Some(BigInt::from(self.stored.reference_gas_price as u64))
66 }
67
68 async fn validator_set(&self, ctx: &Context<'_>) -> Result<Option<ValidatorSet>> {
73 let system_state = ctx
74 .data_unchecked::<PgManager>()
75 .fetch_iota_system_state(Some(self.stored.epoch as u64))
76 .await?;
77
78 let validator_set = NativeStateValidatorInfo::from(system_state).into_validator_set(
79 self.stored.total_stake as u64,
80 self.checkpoint_viewed_at,
81 self.stored.epoch as u64,
82 )?;
83 Ok(Some(validator_set))
84 }
85
86 async fn start_timestamp(&self) -> Result<DateTime, Error> {
88 DateTime::from_ms(self.stored.epoch_start_timestamp)
89 }
90
91 async fn end_timestamp(&self) -> Result<Option<DateTime>, Error> {
93 self.stored
94 .epoch_end_timestamp
95 .map(DateTime::from_ms)
96 .transpose()
97 }
98
99 async fn total_checkpoints(&self, ctx: &Context<'_>) -> Result<Option<UInt53>> {
101 let last = match self.stored.last_checkpoint_id {
102 Some(last) => last as u64,
103 None => {
104 let Watermark { checkpoint, .. } = *ctx.data_unchecked();
105 checkpoint
106 }
107 };
108
109 Ok(Some(
110 UInt53::try_from(last - self.stored.first_checkpoint_id as u64 + 1).extend()?,
111 ))
112 }
113
114 async fn total_transactions(&self) -> Result<Option<UInt53>> {
116 self.stored
118 .epoch_total_transactions()
119 .map(|v| UInt53::try_from(v as u64))
120 .transpose()
121 .extend()
122 }
123
124 async fn total_gas_fees(&self) -> Option<BigInt> {
126 self.stored.total_gas_fees.map(BigInt::from)
127 }
128
129 async fn total_stake_rewards(&self) -> Option<BigInt> {
131 self.stored
132 .total_stake_rewards_distributed
133 .map(BigInt::from)
134 }
135
136 async fn fund_size(&self) -> Option<BigInt> {
140 Some(BigInt::from(self.stored.storage_fund_balance))
141 }
142
143 async fn net_inflow(&self) -> Option<BigInt> {
146 if let (Some(fund_inflow), Some(fund_outflow)) =
147 (self.stored.storage_charge, self.stored.storage_rebate)
148 {
149 Some(BigInt::from(fund_inflow - fund_outflow))
150 } else {
151 None
152 }
153 }
154
155 async fn fund_inflow(&self) -> Option<BigInt> {
157 self.stored.storage_charge.map(BigInt::from)
158 }
159
160 async fn fund_outflow(&self) -> Option<BigInt> {
163 self.stored.storage_rebate.map(BigInt::from)
164 }
165
166 async fn protocol_configs(&self, ctx: &Context<'_>) -> Result<ProtocolConfigs> {
169 ProtocolConfigs::query(ctx.data_unchecked(), Some(self.protocol_version()))
170 .await
171 .extend()
172 }
173
174 #[graphql(flatten)]
175 async fn system_state_summary(&self, ctx: &Context<'_>) -> Result<SystemStateSummary> {
176 let state = ctx
177 .data_unchecked::<PgManager>()
178 .fetch_iota_system_state(Some(self.stored.epoch as u64))
179 .await?;
180 Ok(SystemStateSummary { native: state })
181 }
182
183 async fn live_object_set_digest(&self) -> Result<Option<String>> {
187 let Some(commitments) = self.stored.epoch_commitments.as_ref() else {
188 return Ok(None);
189 };
190 let commitments: Vec<EpochCommitment> = bcs::from_bytes(commitments).map_err(|e| {
191 Error::Internal(format!("Error deserializing commitments: {e}")).extend()
192 })?;
193
194 let digest = commitments.into_iter().next().map(|commitment| {
195 let EpochCommitment::EcmhLiveObjectSet { digest } = commitment else {
196 panic!("a new CheckpointCommitment variant was added and must be handled")
197 };
198 Base58::encode(digest.into_bytes())
199 });
200
201 Ok(digest)
202 }
203
204 async fn checkpoints(
206 &self,
207 ctx: &Context<'_>,
208 first: Option<u64>,
209 after: Option<checkpoint::Cursor>,
210 last: Option<u64>,
211 before: Option<checkpoint::Cursor>,
212 ) -> Result<Connection<String, Checkpoint>> {
213 let page = Page::from_params(ctx.data_unchecked(), first, after, last, before)?;
214 let epoch = self.stored.epoch as u64;
215 Checkpoint::paginate(
216 ctx.data_unchecked(),
217 page,
218 Some(epoch),
219 self.checkpoint_viewed_at,
220 )
221 .await
222 .extend()
223 }
224
225 #[graphql(
254 complexity = "first.or(last).unwrap_or(DEFAULT_PAGE_SIZE as u64) as usize * child_complexity"
255 )]
256 async fn transaction_blocks(
257 &self,
258 ctx: &Context<'_>,
259 first: Option<u64>,
260 after: Option<transaction_block::Cursor>,
261 last: Option<u64>,
262 before: Option<transaction_block::Cursor>,
263 filter: Option<TransactionBlockFilter>,
264 #[graphql(
265 deprecation = "`scanLimit` will be removed with v1.38, along with the support for combining complex filters."
266 )]
267 scan_limit: Option<u64>,
268 ) -> Result<ScanConnection<String, TransactionBlock>> {
269 let page = Page::from_params(ctx.data_unchecked(), first, after, last, before)?;
270
271 let Some(filter) = filter
272 .unwrap_or_default()
273 .intersect(TransactionBlockFilter {
274 after_checkpoint: (self.stored.first_checkpoint_id > 0)
276 .then(|| UInt53::try_from(self.stored.first_checkpoint_id as u64 - 1))
277 .transpose()
278 .extend()?,
279 before_checkpoint: self
280 .stored
281 .last_checkpoint_id
282 .map(|id| UInt53::try_from(id as u64 + 1))
283 .transpose()
284 .extend()?,
285 ..Default::default()
286 })
287 else {
288 return Ok(ScanConnection::new(false, false));
289 };
290
291 TransactionBlock::paginate(ctx, page, filter, self.checkpoint_viewed_at, scan_limit)
292 .await
293 .extend()
294 }
295}
296
297impl Epoch {
298 pub(crate) fn protocol_version(&self) -> u64 {
300 self.stored.protocol_version as u64
301 }
302
303 pub(crate) async fn query(
306 ctx: &Context<'_>,
307 filter: Option<u64>,
308 checkpoint_viewed_at: u64,
309 ) -> Result<Option<Self>, Error> {
310 if let Some(epoch_id) = filter {
311 let DataLoader(dl) = ctx.data_unchecked();
312 dl.load_one(EpochKey {
313 epoch_id,
314 checkpoint_viewed_at,
315 })
316 .await
317 } else {
318 Self::query_latest_at(ctx.data_unchecked(), checkpoint_viewed_at).await
319 }
320 }
321
322 pub(crate) async fn query_latest_at(
326 db: &Db,
327 checkpoint_viewed_at: u64,
328 ) -> Result<Option<Self>, Error> {
329 use epochs::dsl;
330
331 let stored: Option<QueryableEpochInfo> = db
332 .execute(move |conn| {
333 conn.first(move || {
334 dsl::epochs
338 .select(QueryableEpochInfo::as_select())
339 .filter(dsl::first_checkpoint_id.le(checkpoint_viewed_at as i64))
340 .order_by(dsl::first_checkpoint_id.desc())
341 })
342 .optional()
343 })
344 .await
345 .map_err(|e| Error::Internal(format!("Failed to fetch epoch: {e}")))?;
346
347 Ok(stored.map(|stored| Epoch {
348 stored,
349 checkpoint_viewed_at,
350 }))
351 }
352}
353
354impl Loader<EpochKey> for Db {
355 type Value = Epoch;
356 type Error = Error;
357
358 async fn load(&self, keys: &[EpochKey]) -> Result<HashMap<EpochKey, Epoch>, Error> {
359 use epochs::dsl;
360
361 let epoch_ids: BTreeSet<_> = keys.iter().map(|key| key.epoch_id as i64).collect();
362 let epochs: Vec<QueryableEpochInfo> = self
363 .execute_repeatable(move |conn| {
364 conn.results(move || {
365 dsl::epochs
366 .select(QueryableEpochInfo::as_select())
367 .filter(dsl::epoch.eq_any(epoch_ids.iter().cloned()))
368 })
369 })
370 .await
371 .map_err(|e| Error::Internal(format!("Failed to fetch epochs: {e}")))?;
372
373 let epoch_id_to_stored: BTreeMap<_, _> = epochs
374 .into_iter()
375 .map(|stored| (stored.epoch as u64, stored))
376 .collect();
377
378 Ok(keys
379 .iter()
380 .filter_map(|key| {
381 let stored = epoch_id_to_stored.get(&key.epoch_id).cloned()?;
382 let epoch = Epoch {
383 stored,
384 checkpoint_viewed_at: key.checkpoint_viewed_at,
385 };
386
387 let start = epoch.stored.first_checkpoint_id as u64;
393 (key.checkpoint_viewed_at >= start).then_some((*key, epoch))
394 })
395 .collect())
396 }
397}