1use 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}