Skip to main content

dash_sdk/core/
transaction.rs

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/// A Core transaction fetched by id, plus the finality metadata needed to
26/// reconstruct an asset-lock proof from it (an InstantSend proof when the
27/// InstantLock is known, otherwise a ChainLock proof once chain-locked).
28#[derive(Clone, Debug)]
29pub struct FetchedCoreTransaction {
30    /// The decoded transaction.
31    pub transaction: Transaction,
32    /// Height of the block the transaction was mined in (0 if unconfirmed).
33    pub height: u32,
34    /// Whether the transaction's block is ChainLocked.
35    pub is_chain_locked: bool,
36    /// Whether the transaction is InstantSend-locked. Deliberately surfaced but
37    /// not required by the invitation claim: the proof carries the islock from
38    /// the link, and consensus re-verifies it — this flag is informational.
39    pub is_instant_locked: bool,
40}
41
42/// Where a Core transaction was mined, as the queried DAPI node reports it.
43///
44/// Everything here is self-reported by one node. In particular `block_hash`
45/// is not verified: a caller that builds anything on the placement must check
46/// the hash against a header chain it verified itself.
47#[derive(Clone, Debug)]
48#[non_exhaustive]
49pub struct CoreTransactionPlacement {
50    /// The decoded transaction.
51    pub transaction: Transaction,
52    /// Height of the block the transaction was mined in (0 if unconfirmed).
53    pub height: u32,
54    /// Hash of the block the transaction was mined in; `None` when unconfirmed
55    /// or when the reported bytes are not a 32-byte hash.
56    pub block_hash: Option<BlockHash>,
57    /// Whether the transaction's block is ChainLocked.
58    pub is_chain_locked: bool,
59    /// Whether the transaction is InstantSend-locked.
60    pub is_instant_locked: bool,
61}
62
63/// Whether an SDK error is a gRPC `NOT_FOUND` (the requested tx is unknown to
64/// the node), as opposed to a transient/transport failure. Used to distinguish
65/// "retry with a reversed txid" from "surface the error".
66fn 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
76/// DAPI fills `GetTransactionResponse.block_hash` by hex-decoding Core's
77/// display string, so the bytes arrive reversed relative to the hash's
78/// internal order.
79fn 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
85/// Decode the consensus-encoded transaction bytes of a `getTransaction` reply.
86fn 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    /// Fetch a Core transaction by its id via DAPI `getTransaction`.
92    ///
93    /// `txid` is the transaction id as a hex string (big-endian display form).
94    /// Returns `Ok(Some(..))` with the decoded transaction plus its
95    /// confirmation/lock metadata; `Ok(None)` when the node does not know the tx
96    /// (empty response or gRPC `NOT_FOUND`) so the caller can retry with the id
97    /// byte-reversed; and `Err` for a transient/transport failure that must not
98    /// be masked by a doomed reversed-id retry.
99    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    /// Fetch where a Core transaction was mined via DAPI `getTransaction`,
119    /// including the block hash the node reports.
120    ///
121    /// Same `txid` form and `Ok(None)` / `Err` contract as
122    /// [`Sdk::get_transaction`]; `settings` override the SDK's request
123    /// settings for this call, so a caller on a deadline can bound it. The
124    /// placement is unverified — see [`CoreTransactionPlacement`].
125    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    /// Fetch a Core block by hash via DAPI `getBlock`.
144    ///
145    /// Returns `Ok(None)` when the node does not serve the block (gRPC
146    /// `NOT_FOUND` or an empty reply) and `Err` for any other failure,
147    /// including bytes that do not decode as a block. The bytes are
148    /// self-reported by the queried node: check `block_hash()` against a
149    /// header chain you trust before relying on anything in the block.
150    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                    // Core's `getblock` takes the display-order hex `BlockHash` formats to.
159                    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    /// Run `getTransaction`, mapping an unknown transaction (gRPC `NOT_FOUND`
186    /// or an empty reply) to `Ok(None)` and every other failure to `Err`.
187    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    /// Starts the stream to listen for instant send lock messages
220    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        // create the bloom filter
228        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, // Subscribing to new transactions as well
251            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    /// Waits for a response for the asset lock proof
263    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        // Define an inner async block to handle the stream processing.
280        let stream_processing = async {
281            loop {
282                // TODO: We should retry if Err is returned
283                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                        // Continue receiving messages until we find the transaction
352                        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                        // TODO: This a temporary implementation until we have headers stream running in background
365                        //  so we can always get actual height and chain locks
366
367                        // Wait until the block is chainlocked
368                        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 // TODO: We need better way to handle execution errors
382                                .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                        // Wait until platform chain is on the block's chain locked height
401                        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        // Apply the timeout if `time_out_ms` is Some, otherwise just await the processing.
444        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    /// DAPI's display-order bytes come back as the hash's internal order.
458    #[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    /// Anything but a 32-byte hash, including an unconfirmed tx's empty field,
470    /// yields no hash.
471    #[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}