1use crate::platform::fetch_current_no_parameters::FetchCurrent;
2use crate::platform::types::epoch::Epoch;
3use crate::{Error, Sdk};
4use bip37_bloom_filter::{BloomFilter, BloomFilterData};
5use dapi_grpc::core::v0::{
6 get_block_request, transactions_with_proofs_request, transactions_with_proofs_response,
7 GetBlockRequest, GetTransactionRequest, GetTransactionResponse, TransactionsWithProofsRequest,
8 TransactionsWithProofsResponse,
9};
10use dpp::dashcore::consensus::Decodable;
11use dpp::dashcore::hashes::Hash;
12use dpp::dashcore::{
13 Address, Block, BlockHash, InstantLock, MerkleBlock, OutPoint, Transaction, Txid,
14};
15use dpp::identity::state_transition::asset_lock_proof::chain::ChainAssetLockProof;
16use dpp::identity::state_transition::asset_lock_proof::InstantAssetLockProof;
17use dpp::prelude::AssetLockProof;
18
19use dapi_grpc::tonic::Code;
20use rs_dapi_client::transport::TransportError;
21use rs_dapi_client::{DapiClientError, DapiRequestExecutor, IntoInner, RequestSettings};
22use std::time::Duration;
23use tokio::time::{sleep, timeout};
24
25#[derive(Clone, Debug)]
29pub struct FetchedCoreTransaction {
30 pub transaction: Transaction,
32 pub height: u32,
34 pub is_chain_locked: bool,
36 pub is_instant_locked: bool,
40}
41
42#[derive(Clone, Debug)]
48#[non_exhaustive]
49pub struct CoreTransactionPlacement {
50 pub transaction: Transaction,
52 pub height: u32,
54 pub block_hash: Option<BlockHash>,
57 pub is_chain_locked: bool,
59 pub is_instant_locked: bool,
61}
62
63fn error_is_not_found(err: &Error) -> bool {
67 match err {
68 Error::DapiClientError(DapiClientError::Transport(TransportError::Grpc(status))) => {
69 status.code() == Code::NotFound
70 }
71 Error::NoAvailableAddressesToRetry(inner) => error_is_not_found(inner),
72 _ => false,
73 }
74}
75
76fn block_hash_from_display_bytes(bytes: &[u8]) -> Option<BlockHash> {
80 let mut hash: [u8; 32] = bytes.try_into().ok()?;
81 hash.reverse();
82 Some(BlockHash::from_byte_array(hash))
83}
84
85fn decode_transaction(bytes: &[u8]) -> Result<Transaction, Error> {
87 Transaction::consensus_decode(&mut &bytes[..]).map_err(|e| Error::CoreError(e.into()))
88}
89
90impl Sdk {
91 pub async fn get_transaction(
100 &self,
101 txid: &str,
102 ) -> Result<Option<FetchedCoreTransaction>, Error> {
103 let Some(response) = self
104 .fetch_core_transaction(txid, RequestSettings::default())
105 .await?
106 else {
107 return Ok(None);
108 };
109
110 Ok(Some(FetchedCoreTransaction {
111 transaction: decode_transaction(&response.transaction)?,
112 height: response.height,
113 is_chain_locked: response.is_chain_locked,
114 is_instant_locked: response.is_instant_locked,
115 }))
116 }
117
118 pub async fn get_transaction_placement(
126 &self,
127 txid: &str,
128 settings: RequestSettings,
129 ) -> Result<Option<CoreTransactionPlacement>, Error> {
130 let Some(response) = self.fetch_core_transaction(txid, settings).await? else {
131 return Ok(None);
132 };
133
134 Ok(Some(CoreTransactionPlacement {
135 transaction: decode_transaction(&response.transaction)?,
136 height: response.height,
137 block_hash: block_hash_from_display_bytes(&response.block_hash),
138 is_chain_locked: response.is_chain_locked,
139 is_instant_locked: response.is_instant_locked,
140 }))
141 }
142
143 pub async fn get_block_by_hash(
151 &self,
152 hash: &BlockHash,
153 settings: RequestSettings,
154 ) -> Result<Option<Block>, Error> {
155 let response = match self
156 .execute(
157 GetBlockRequest {
158 block: Some(get_block_request::Block::Hash(hash.to_string())),
160 },
161 settings,
162 )
163 .await
164 .into_inner()
165 {
166 Ok(response) => response,
167 Err(e) => {
168 let err: Error = e.into();
169 return if error_is_not_found(&err) {
170 Ok(None)
171 } else {
172 Err(err)
173 };
174 }
175 };
176
177 if response.block.is_empty() {
178 return Ok(None);
179 }
180 Block::consensus_decode(&mut response.block.as_slice())
181 .map(Some)
182 .map_err(|e| Error::CoreError(e.into()))
183 }
184
185 async fn fetch_core_transaction(
188 &self,
189 txid: &str,
190 settings: RequestSettings,
191 ) -> Result<Option<GetTransactionResponse>, Error> {
192 let response = match self
193 .execute(
194 GetTransactionRequest {
195 id: txid.to_string(),
196 },
197 settings,
198 )
199 .await
200 .into_inner()
201 {
202 Ok(response) => response,
203 Err(e) => {
204 let err: Error = e.into();
205 return if error_is_not_found(&err) {
206 Ok(None)
207 } else {
208 Err(err)
209 };
210 }
211 };
212
213 if response.transaction.is_empty() {
214 return Ok(None);
215 }
216 Ok(Some(response))
217 }
218
219 pub async fn start_instant_send_lock_stream(
221 &self,
222 from_block_hash: Vec<u8>,
223 address: &Address,
224 ) -> Result<dapi_grpc::tonic::Streaming<TransactionsWithProofsResponse>, Error> {
225 let address_bytes = address.as_unchecked().payload_to_vec();
226
227 let bloom_filter = BloomFilter::builder(1, 0.001)
229 .expect("this FP rate allows up to 10000 items")
230 .add_element(&address_bytes)
231 .build();
232
233 let bloom_filter_proto = {
234 let BloomFilterData {
235 v_data,
236 n_hash_funcs,
237 n_tweak,
238 n_flags,
239 } = bloom_filter.into();
240 dapi_grpc::core::v0::BloomFilter {
241 v_data,
242 n_hash_funcs,
243 n_tweak,
244 n_flags,
245 }
246 };
247
248 let core_transactions_stream = TransactionsWithProofsRequest {
249 bloom_filter: Some(bloom_filter_proto),
250 count: 0, send_transaction_hashes: true,
252 from_block: Some(transactions_with_proofs_request::FromBlock::FromBlockHash(
253 from_block_hash,
254 )),
255 };
256 self.execute(core_transactions_stream, RequestSettings::default())
257 .await
258 .into_inner()
259 .map_err(|e| e.into())
260 }
261
262 pub async fn wait_for_asset_lock_proof_for_transaction(
264 &self,
265 mut stream: dapi_grpc::tonic::Streaming<TransactionsWithProofsResponse>,
266 transaction: &Transaction,
267 time_out: Option<Duration>,
268 ) -> Result<AssetLockProof, Error> {
269 let transaction_id = transaction.txid();
270
271 let _span = tracing::debug_span!(
272 "wait_for_asset_lock_proof_for_transaction",
273 transaction_id = transaction_id.to_string(),
274 )
275 .entered();
276
277 tracing::debug!("waiting for messages from stream");
278
279 let stream_processing = async {
281 loop {
282 let message = stream
284 .message()
285 .await
286 .map_err(|e| Error::Generic(format!("can't receive message: {e}")))?;
287
288 let Some(TransactionsWithProofsResponse { responses }) = message else {
289 return Err(Error::Generic("stream closed unexpectedly".to_string()));
290 };
291
292 match responses {
293 Some(
294 transactions_with_proofs_response::Responses::InstantSendLockMessages(
295 instant_send_lock_messages,
296 ),
297 ) => {
298 tracing::debug!(
299 "received {} instant lock message(s)",
300 instant_send_lock_messages.messages.len()
301 );
302
303 for instant_lock_bytes in instant_send_lock_messages.messages {
304 let instant_lock =
305 InstantLock::consensus_decode(&mut instant_lock_bytes.as_slice())
306 .map_err(|e| {
307 tracing::error!("invalid asset lock: {}", e);
308
309 Error::CoreError(e.into())
310 })?;
311
312 if instant_lock.txid == transaction_id {
313 let asset_lock_proof =
314 AssetLockProof::Instant(InstantAssetLockProof {
315 instant_lock,
316 transaction: transaction.clone(),
317 output_index: 0,
318 });
319
320 tracing::debug!(
321 ?asset_lock_proof,
322 "instant lock is matching to the broadcasted transaction, returning instant asset lock proof"
323 );
324
325 return Ok(asset_lock_proof);
326 } else {
327 tracing::debug!(
328 "instant lock is not matching, waiting for the next message"
329 );
330 }
331 }
332 }
333 Some(transactions_with_proofs_response::Responses::RawMerkleBlock(
334 raw_merkle_block,
335 )) => {
336 tracing::debug!("received merkle block");
337
338 let merkle_block =
339 MerkleBlock::consensus_decode(&mut raw_merkle_block.as_slice())
340 .map_err(|e| {
341 tracing::error!("can't decode merkle block: {}", e);
342
343 Error::CoreError(e.into())
344 })?;
345
346 let mut matches: Vec<Txid> = vec![];
347 let mut index: Vec<u32> = vec![];
348
349 merkle_block.extract_matches(&mut matches, &mut index)?;
350
351 if !matches.contains(&transaction_id) {
353 tracing::debug!(
354 "merkle block doesn't contain the transaction, waiting for the next message"
355 );
356
357 continue;
358 }
359
360 tracing::debug!(
361 "merkle block contains the transaction, obtaining core chain locked height"
362 );
363
364 let mut core_chain_locked_height;
369 loop {
370 let GetTransactionResponse {
371 height,
372 is_chain_locked,
373 ..
374 } = self
375 .execute(
376 GetTransactionRequest {
377 id: transaction_id.to_string(),
378 },
379 RequestSettings::default(),
380 )
381 .await .into_inner()?;
383
384 core_chain_locked_height = height;
385
386 if is_chain_locked {
387 break;
388 }
389
390 tracing::trace!("the transaction is on height {} but not chainlocked. try again in 1 sec", height);
391
392 sleep(Duration::from_secs(1)).await;
393 }
394
395 tracing::debug!(
396 "the transaction is chainlocked on height {}, waiting platform for reaching the same core height",
397 core_chain_locked_height
398 );
399
400 loop {
402 let (_epoch, metadata) =
403 Epoch::fetch_current_with_metadata(self).await?;
404
405 if metadata.core_chain_locked_height >= core_chain_locked_height {
406 break;
407 }
408
409 tracing::trace!(
410 "platform chain locked core height {} but we need {}. try again in 1 sec",
411 metadata.core_chain_locked_height,
412 core_chain_locked_height,
413 );
414
415 sleep(Duration::from_secs(1)).await;
416 }
417
418 let asset_lock_proof = AssetLockProof::Chain(ChainAssetLockProof {
419 core_chain_locked_height,
420 out_point: OutPoint {
421 txid: transaction.txid(),
422 vout: 0,
423 },
424 });
425
426 tracing::debug!(
427 ?asset_lock_proof,
428 "merkle block contains the broadcasted transaction, returning chain asset lock proof"
429 );
430
431 return Ok(asset_lock_proof);
432 }
433 Some(transactions_with_proofs_response::Responses::RawTransactions(_)) => {
434 tracing::trace!("received transaction(s), ignoring")
435 }
436 None => tracing::trace!(
437 "received empty response as a workaround for the bug in tonic, ignoring"
438 ),
439 }
440 }
441 };
442
443 match time_out {
445 Some(duration) => timeout(duration, stream_processing).await.map_err(|_| {
446 Error::TimeoutReached(duration, String::from("receiving asset lock proof"))
447 })?,
448 None => stream_processing.await,
449 }
450 }
451}
452
453#[cfg(test)]
454mod tests {
455 use super::*;
456
457 #[test]
459 fn block_hash_from_display_bytes_reverses_the_bytes() {
460 let mut display = [0u8; 32];
461 display[0] = 0xaa;
462 display[31] = 0x01;
463 let hash = block_hash_from_display_bytes(&display).expect("32 bytes");
464 let internal = hash.to_byte_array();
465 assert_eq!(internal[0], 0x01);
466 assert_eq!(internal[31], 0xaa);
467 }
468
469 #[test]
472 fn block_hash_from_display_bytes_rejects_wrong_lengths() {
473 assert!(block_hash_from_display_bytes(&[]).is_none());
474 assert!(block_hash_from_display_bytes(&[0u8; 31]).is_none());
475 assert!(block_hash_from_display_bytes(&[0u8; 33]).is_none());
476 }
477}