Skip to main content

iota_graphql_rpc/types/
epoch.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5use 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/// `DataLoader` key for fetching an `Epoch` by its ID, optionally constrained
40/// by a consistency cursor.
41#[derive(Copy, Clone, Hash, Eq, PartialEq, Debug)]
42struct EpochKey {
43    pub epoch_id: u64,
44    pub checkpoint_viewed_at: u64,
45}
46
47/// Operation of the IOTA network is temporally partitioned into non-overlapping
48/// epochs, and the network aims to keep epochs roughly the same duration as
49/// each other. During a particular epoch the following data is fixed:
50///
51/// - the protocol version
52/// - the reference gas price
53/// - the set of participating validators
54#[Object]
55impl Epoch {
56    /// The epoch's id as a sequence number that starts at 0 and is incremented
57    /// by one at every epoch change.
58    async fn epoch_id(&self) -> Result<UInt53> {
59        UInt53::try_from(self.stored.epoch as u64).extend()
60    }
61
62    /// The minimum gas price that a quorum of validators are guaranteed to sign
63    /// a transaction for.
64    async fn reference_gas_price(&self) -> Option<BigInt> {
65        Some(BigInt::from(self.stored.reference_gas_price as u64))
66    }
67
68    /// Validator related properties, including the active validators.
69    ///
70    /// For epochs other than the current the data provided refer to the start
71    /// of the epoch.
72    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    /// The epoch's starting timestamp.
87    async fn start_timestamp(&self) -> Result<DateTime, Error> {
88        DateTime::from_ms(self.stored.epoch_start_timestamp)
89    }
90
91    /// The epoch's ending timestamp.
92    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    /// The total number of checkpoints in this epoch.
100    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    /// The total number of transaction blocks in this epoch.
115    async fn total_transactions(&self) -> Result<Option<UInt53>> {
116        // TODO: this currently returns None for the current epoch. Fix this.
117        self.stored
118            .epoch_total_transactions()
119            .map(|v| UInt53::try_from(v as u64))
120            .transpose()
121            .extend()
122    }
123
124    /// The total amount of gas fees (in NANOS) that were paid in this epoch.
125    async fn total_gas_fees(&self) -> Option<BigInt> {
126        self.stored.total_gas_fees.map(BigInt::from)
127    }
128
129    /// The total NANOS rewarded as stake.
130    async fn total_stake_rewards(&self) -> Option<BigInt> {
131        self.stored
132            .total_stake_rewards_distributed
133            .map(BigInt::from)
134    }
135
136    /// The storage fund available in this epoch.
137    /// This fund is used to redistribute storage fees from past transactions
138    /// to future validators.
139    async fn fund_size(&self) -> Option<BigInt> {
140        Some(BigInt::from(self.stored.storage_fund_balance))
141    }
142
143    /// The difference between the fund inflow and outflow, representing
144    /// the net amount of storage fees accumulated in this epoch.
145    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    /// The storage fees paid for transactions executed during the epoch.
156    async fn fund_inflow(&self) -> Option<BigInt> {
157        self.stored.storage_charge.map(BigInt::from)
158    }
159
160    /// The storage fee rebates paid to users who deleted the data associated
161    /// with past transactions.
162    async fn fund_outflow(&self) -> Option<BigInt> {
163        self.stored.storage_rebate.map(BigInt::from)
164    }
165
166    /// The epoch's corresponding protocol configuration, including the feature
167    /// flags and the configuration options.
168    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    /// A commitment by the committee at the end of epoch on the contents of the
184    /// live object set at that time. This can be used to verify state
185    /// snapshots.
186    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    /// The epoch's corresponding checkpoints.
205    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    /// The epoch's corresponding transaction blocks.
226    ///
227    /// `scanLimit` restricts the number of candidate transactions scanned when
228    /// gathering a page of results. It is required for queries that apply two
229    /// or more complex filters (on function, affected address, recipient, input
230    /// object, changed object, or wrapped or deleted object), and can be at
231    /// most `serviceConfig.maxScanLimit`. A `kind` filter cannot be
232    /// combined with any of them.
233    ///
234    /// When the scan limit is reached the page will be returned even if it has
235    /// fewer than `first` results when paginating forward (`last` when
236    /// paginating backwards). If there are more transactions to scan,
237    /// `pageInfo.hasNextPage` (or `pageInfo.hasPreviousPage`) will be set to
238    /// `true`, and `PageInfo.endCursor` (or `PageInfo.startCursor`) will be set
239    /// to the last transaction that was scanned as opposed to the last (or
240    /// first) transaction in the page.
241    ///
242    /// Requesting the next (or previous) page after this cursor will resume the
243    /// search, scanning the next `scanLimit` many transactions in the
244    /// direction of pagination, and so on until all transactions in the
245    /// scanning range have been visited.
246    ///
247    /// By default, the scanning range consists of all transactions in this
248    /// epoch.
249    ///
250    /// DEPRECATION NOTICE: Support for the combination of two or more complex
251    /// filters as discussed above will stop with the v1.38 release. `scanLimit`
252    /// will thus become obsolete and will be removed as well.
253    #[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                // If `first_checkpoint_id` is 0, we include the 0th checkpoint by leaving it None
275                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    /// The epoch's protocol version.
299    pub(crate) fn protocol_version(&self) -> u64 {
300        self.stored.protocol_version as u64
301    }
302
303    /// Look up an `Epoch` in the database, optionally filtered by its Epoch ID.
304    /// If no ID is supplied, defaults to fetching the latest epoch.
305    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    /// Look up the latest `Epoch` from the database, optionally filtered by a
323    /// consistency cursor (querying for a consistency cursor in the past
324    /// looks for the latest epoch as of that cursor).
325    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                    // Bound the query on `checkpoint_viewed_at` by filtering for the epoch
335                    // whose `first_checkpoint_id <= checkpoint_viewed_at`, selecting the epoch
336                    // with the largest `first_checkpoint_id` among the filtered set.
337                    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                // We filter by checkpoint viewed at in memory because it should be quite rare
388                // that this query actually filters something (only in edge
389                // cases), and not trying to encode it in the SQL query makes
390                // the query much simpler and therefore easier for
391                // the DB to plan.
392                let start = epoch.stored.first_checkpoint_id as u64;
393                (key.checkpoint_viewed_at >= start).then_some((*key, epoch))
394            })
395            .collect())
396    }
397}