Skip to main content

iota_graphql_rpc/types/
available_range.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5use async_graphql::*;
6use diesel::{ExpressionMethods, QueryDsl, QueryResult};
7use iota_indexer::schema::{checkpoints, watermarks};
8
9use crate::{
10    data::{Conn, Db, DbConnection, QueryExecutor},
11    error::Error,
12    types::checkpoint::{Checkpoint, CheckpointId},
13};
14#[derive(Clone, Debug, PartialEq, Eq, Copy)]
15pub(crate) struct AvailableRange {
16    pub first: u64,
17    pub last: u64,
18}
19
20/// Range of checkpoints that the RPC is guaranteed to produce a consistent
21/// response for.
22#[Object]
23impl AvailableRange {
24    async fn first(&self, ctx: &Context<'_>) -> Result<Option<Checkpoint>> {
25        let id = CheckpointId::by_seq_num(self.first).extend()?;
26        Checkpoint::query(ctx, id, self.last).await.extend()
27    }
28
29    async fn last(&self, ctx: &Context<'_>) -> Result<Option<Checkpoint>> {
30        let id = CheckpointId::by_seq_num(self.last).extend()?;
31        Checkpoint::query(ctx, id, self.last).await.extend()
32    }
33}
34
35impl AvailableRange {
36    /// Look up the available range when viewing the data consistently at
37    /// `checkpoint_viewed_at`. Uses the executor's configured
38    /// `max_available_range`.
39    pub(crate) async fn query(db: &Db, checkpoint_viewed_at: u64) -> Result<Self, Error> {
40        let max_available_range = db.max_available_range;
41        let Some(range): Option<Self> = db
42            .execute(move |conn| Self::result(conn, checkpoint_viewed_at, max_available_range))
43            .await
44            .map_err(|e| Error::Internal(format!("Failed to fetch available range: {e}")))?
45        else {
46            return Err(Error::Client(format!(
47                "Requesting data at checkpoint {checkpoint_viewed_at}, outside the available \
48                 range.",
49            )));
50        };
51
52        Ok(range)
53    }
54
55    /// Computes the backward-diff retention window. Returns `Some(range)`
56    /// when `checkpoint_viewed_at` falls inside it, `None` otherwise. The
57    /// lower bound is the greater of the backward history watermark
58    /// (`min_available_cp`) and `latest_checkpoint - max_available_range`; the
59    /// upper bound is `checkpoint_viewed_at` itself, so the returned range
60    /// never extends past the request's captured watermark.
61    ///
62    /// `max_available_range` caps how far back a consistent view can be
63    /// requested, for performance reasons (the further back, the more
64    /// backward history entries to traverse).
65    pub(crate) fn result(
66        conn: &mut Conn,
67        checkpoint_viewed_at: u64,
68        max_available_range: u64,
69    ) -> QueryResult<Option<Self>> {
70        use checkpoints::dsl as cp;
71        use watermarks::dsl as wm;
72
73        let last: Option<i64> = conn
74            .results(|| {
75                cp::checkpoints
76                    .select(cp::sequence_number)
77                    .order(cp::sequence_number.desc())
78                    .limit(1)
79                    .into_boxed()
80            })?
81            .into_iter()
82            .next();
83
84        let watermark_first: Option<i64> = conn
85            .results(|| {
86                wm::watermarks
87                    .filter(wm::entity.eq(crate::backward_view::BACKWARD_HISTORY_WATERMARK_ENTITY))
88                    .select(wm::min_available_cp)
89                    .limit(1)
90                    .into_boxed()
91            })?
92            .into_iter()
93            .next();
94
95        let last = last.unwrap_or(0) as u64;
96        let watermark_first = watermark_first.unwrap_or(0) as u64;
97        let lag_first = last.saturating_sub(max_available_range);
98        let first = watermark_first.max(lag_first);
99
100        if checkpoint_viewed_at < first || checkpoint_viewed_at > last {
101            return Ok(None);
102        }
103        // Cap `last` to `checkpoint_viewed_at` so a single request never
104        // reports a range extending past its captured watermark.
105        Ok(Some(Self {
106            first,
107            last: checkpoint_viewed_at,
108        }))
109    }
110
111    pub(crate) fn is_checkpoint_in_backward_history_range(
112        conn: &mut Conn,
113        checkpoint_viewed_at: u64,
114        max_available_range: u64,
115    ) -> QueryResult<bool> {
116        Ok(Self::result(conn, checkpoint_viewed_at, max_available_range)?.is_some())
117    }
118}