Skip to main content

poi_rs/source/
grpc.rs

1// Copyright 2020-2026 IOTA Stiftung
2// SPDX-License-Identifier: Apache-2.0
3
4use async_trait::async_trait;
5use iota_grpc_client::read_mask_fields::{
6    CheckpointResponseField, EpochField, ObjectField, ServiceInfoField, TransactionField,
7};
8use iota_grpc_client::{Client as GrpcClient, ReadMask};
9use iota_grpc_types::proto::TryFromProtoError;
10use iota_sdk_types::{
11    CheckpointContents, CheckpointDigest, ObjectId, SignedCheckpointSummary, SignedTransaction, TransactionDigest,
12    Version,
13};
14use iota_types::committee::{Committee, EpochId};
15use iota_types::digests::ChainIdentifier;
16use iota_types::effects::TransactionEffectsAPI;
17use iota_types::messages_checkpoint::CertifiedCheckpointSummary;
18use iota_types::object::Object;
19use iota_types::transaction::Transaction;
20
21use super::{Source, SourceCheckpoint, SourceError, SourceTransaction};
22
23#[async_trait]
24impl Source for GrpcClient {
25    async fn chain_identifier(&self) -> Result<ChainIdentifier, SourceError> {
26        let service_info = self
27            .get_service_info(Some(ReadMask::from(ServiceInfoField::CHAIN_ID)))
28            .await
29            .map_err(SourceError::request)?;
30        let chain_identifier = service_info
31            .body()
32            .chain_identifier()
33            .map_err(SourceError::invalid_response)?;
34
35        Ok(ChainIdentifier::from(CheckpointDigest::new(
36            chain_identifier.into_inner(),
37        )))
38    }
39
40    async fn transaction(
41        &self,
42        transaction_digest: TransactionDigest,
43    ) -> Result<Option<SourceTransaction>, SourceError> {
44        let transactions = self
45            .get_transactions(
46                &[transaction_digest],
47                Some(ReadMask::from(&[
48                    TransactionField::TRANSACTION_BCS,
49                    TransactionField::SIGNATURES,
50                    TransactionField::EFFECTS_BCS,
51                    TransactionField::EVENTS_DIGEST,
52                    TransactionField::EVENTS_EVENTS_BCS,
53                    TransactionField::CHECKPOINT,
54                ])),
55            )
56            .await
57            .map_err(SourceError::request)?;
58        let Some(executed_transaction) = transactions.into_inner().into_iter().next() else {
59            return Ok(None);
60        };
61        let executed_transaction = match executed_transaction {
62            Ok(executed_transaction) => executed_transaction,
63            Err(error) if error.is_not_found() => return Ok(None),
64            Err(error) => return Err(SourceError::request(error)),
65        };
66        let effects = executed_transaction
67            .effects()
68            .map_err(SourceError::invalid_response)?
69            .effects()
70            .map_err(SourceError::invalid_response)?;
71
72        let transaction = executed_transaction
73            .transaction()
74            .map_err(SourceError::invalid_response)?
75            .transaction()
76            .map_err(SourceError::invalid_response)?;
77        let signatures = executed_transaction
78            .signatures()
79            .map_err(SourceError::invalid_response)?
80            .signatures
81            .iter()
82            .map(|signature| signature.signature().map_err(SourceError::invalid_response))
83            .collect::<Result<Vec<_>, SourceError>>()?;
84        let transaction: Transaction = SignedTransaction {
85            transaction,
86            signatures,
87        }
88        .into();
89        let events = if effects.events_digest().is_some() {
90            executed_transaction
91                .events()
92                .map_err(SourceError::missing_data)?
93                .events()
94                .map_err(SourceError::invalid_response)
95                .map(Some)?
96        } else {
97            None
98        };
99        let checkpoint_sequence_number = executed_transaction
100            .checkpoint_sequence_number()
101            .map_err(SourceError::missing_data)?;
102
103        Ok(Some(SourceTransaction {
104            transaction,
105            effects,
106            events,
107            checkpoint_sequence_number,
108        }))
109    }
110
111    async fn object(&self, object_id: ObjectId, version: Option<Version>) -> Result<Option<Object>, SourceError> {
112        let objects = self
113            .get_objects(&[(object_id, version)], Some(ReadMask::from(ObjectField::BCS)))
114            .await
115            .map_err(SourceError::request)?;
116        let Some(response) = objects.into_inner().into_iter().next() else {
117            return Ok(None);
118        };
119        let response = match response {
120            Ok(response) => response,
121            Err(error) if error.is_not_found() => return Ok(None),
122            Err(error) => return Err(SourceError::request(error)),
123        };
124        let object: Object = response.object().map_err(SourceError::invalid_response)?.into();
125
126        Ok(Some(object))
127    }
128
129    async fn checkpoint(&self, sequence_number: u64) -> Result<Option<SourceCheckpoint>, SourceError> {
130        let checkpoint = match self
131            .get_checkpoint_by_sequence_number(
132                sequence_number,
133                Some(ReadMask::from(&[
134                    CheckpointResponseField::CHECKPOINT_SUMMARY_BCS,
135                    CheckpointResponseField::CHECKPOINT_SIGNATURE,
136                    CheckpointResponseField::CHECKPOINT_CONTENTS_BCS,
137                ])),
138                None,
139                None,
140            )
141            .await
142        {
143            Ok(response) => response.into_inner(),
144            Err(error) if error.is_not_found() => return Ok(None),
145            Err(error) => return Err(SourceError::request(error)),
146        };
147        let summary: CertifiedCheckpointSummary = checkpoint
148            .signed_summary()
149            .map_err(SourceError::invalid_response)?
150            .try_into()
151            .map_err(SourceError::invalid_response)?;
152        let contents: CheckpointContents = checkpoint
153            .contents()
154            .map_err(SourceError::invalid_response)?
155            .contents()
156            .map_err(SourceError::invalid_response)?;
157
158        Ok(Some(SourceCheckpoint { summary, contents }))
159    }
160
161    async fn committee(&self, epoch: EpochId) -> Result<Committee, SourceError> {
162        let epoch_info = self
163            .get_epoch(Some(epoch), Some(ReadMask::from(EpochField::COMMITTEE)))
164            .await
165            .map_err(SourceError::request)?
166            .into_inner();
167        let committee = epoch_info.committee().map_err(SourceError::invalid_response)?;
168
169        Ok(committee.into())
170    }
171
172    async fn current_epoch(&self) -> Result<Option<EpochId>, SourceError> {
173        self.get_service_info(Some(ReadMask::from(ServiceInfoField::EPOCH)))
174            .await
175            .map(|response| response.body().epoch)
176            .map_err(SourceError::request)
177    }
178
179    async fn epoch_close_summary(&self, epoch: EpochId) -> Result<Option<CertifiedCheckpointSummary>, SourceError> {
180        let epoch_info = self
181            .get_epoch(
182                Some(epoch),
183                Some(ReadMask::from(EpochField::EPOCH_CLOSE_PROOF_CHECKPOINT)),
184            )
185            .await
186            .map_err(SourceError::request)?
187            .into_inner();
188        let Some(epoch_close_proof) = epoch_info.epoch_close_proof().map_err(SourceError::invalid_response)? else {
189            return Ok(None);
190        };
191        let checkpoint = epoch_close_proof.checkpoint().map_err(SourceError::missing_data)?;
192        let summary = checkpoint
193            .summary
194            .as_ref()
195            .ok_or_else(|| TryFromProtoError::missing("summary"))
196            .map_err(SourceError::missing_data)?;
197        let summary = summary.summary().map_err(SourceError::invalid_response)?;
198        let signature = checkpoint
199            .signature
200            .as_ref()
201            .ok_or_else(|| TryFromProtoError::missing("signature"))
202            .map_err(SourceError::missing_data)?;
203        let signature = signature.signature().map_err(SourceError::invalid_response)?;
204        let signed_summary = SignedCheckpointSummary {
205            checkpoint: summary,
206            signature,
207        };
208        let certified_summary = signed_summary.try_into().map_err(SourceError::invalid_response)?;
209
210        Ok(Some(certified_summary))
211    }
212}