iota_graphql_rpc/types/
available_range.rs1use 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#[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 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 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 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}