Skip to main content

drive_abci/rpc/
core.rs

1use crate::rpc::prefetch::CorePrefetcher;
2use dpp::dashcore::consensus::encode::deserialize_partial;
3use dpp::dashcore::ephemerealdata::chain_lock::ChainLock;
4use dpp::dashcore::transaction::special_transaction::TransactionPayload;
5use dpp::dashcore::{Block, BlockHash, QuorumHash, Transaction, Txid};
6use dpp::dashcore::{Header, InstantLock};
7use dpp::dashcore_rpc::dashcore_rpc_json::{
8    AssetUnlockStatusResult, ExtendedQuorumDetails, ExtendedQuorumListResult, GetChainTipsResult,
9    MasternodeListDiff, MnSyncStatus, QuorumInfoResult, QuorumType, SoftforkInfo,
10};
11use dpp::dashcore_rpc::json::GetRawTransactionResult;
12use dpp::dashcore_rpc::{Auth, Client, Error, RpcApi};
13use dpp::prelude::TimestampMillis;
14use serde_json::Value;
15use std::collections::HashMap;
16use std::time::Duration;
17
18/// Information returned by QuorumListExtended
19pub type QuorumListExtendedInfo = HashMap<QuorumHash, ExtendedQuorumDetails>;
20
21/// The special transaction type of a coinbase (`TRANSACTION_COINBASE` in Dash Core).
22const COINBASE_TRANSACTION_TYPE: u16 = 5;
23
24/// Reads Core's credit pool balance after a block, in duffs, from the block's serialized
25/// coinbase transaction; `0` when it carries no payload or its payload predates version 3,
26/// before the credit pool existed, which is how Core's own unlock limit reads such a block.
27///
28/// Decoded with `deserialize_partial`: the pinned payload decoder reads the fields of version 3
29/// for every later version and stops after the balance, while the version 4 payload Core v24
30/// requires appends `merkleRootAssetUnlocks` after it. A strict `deserialize`, and so a whole
31/// `Block` decode, refuses those unread bytes. Core has only ever appended fields to the
32/// payload, and its consensus rules (`CheckCbTx`) refuse versions it does not know, so the
33/// balance stays where version 3 put it unless a Core release moves it, which Platform would
34/// have to follow anyway.
35pub(crate) fn credit_pool_balance_from_coinbase(coinbase: &[u8]) -> Result<u64, String> {
36    let (transaction, _) = deserialize_partial::<Transaction>(coinbase)
37        .map_err(|e| format!("coinbase cannot be decoded: {e}"))?;
38    match transaction.special_transaction_payload {
39        None => Ok(0),
40        Some(TransactionPayload::CoinbasePayloadType(payload)) => {
41            Ok(payload.asset_locked_amount.unwrap_or_default())
42        }
43        Some(payload) => Err(format!(
44            "coinbase carries a {:?} payload",
45            payload.get_type()
46        )),
47    }
48}
49
50/// Core height must be of type u32 (Platform heights are u64)
51pub type CoreHeight = u32;
52/// Core RPC interface
53#[cfg_attr(any(feature = "mocks", test), mockall::automock)]
54pub trait CoreRPCLike {
55    /// Get block hash by height
56    fn get_block_hash(&self, height: CoreHeight) -> Result<BlockHash, Error>;
57
58    /// Get block hash by height
59    fn get_block_header(&self, block_hash: &BlockHash) -> Result<Header, Error>;
60
61    /// Get block time of a chain locked core height
62    fn get_block_time_from_height(&self, height: CoreHeight) -> Result<TimestampMillis, Error>;
63
64    /// Get the best chain lock
65    fn get_best_chain_lock(&self) -> Result<ChainLock, Error>;
66
67    /// Submit a chain lock
68    fn submit_chain_lock(&self, chain_lock: &ChainLock) -> Result<u32, Error>;
69
70    /// Get transaction
71    fn get_transaction(&self, tx_id: &Txid) -> Result<Transaction, Error>;
72
73    /// Get asset unlock statuses
74    fn get_asset_unlock_statuses(
75        &self,
76        indices: &[u64],
77        core_chain_locked_height: u32,
78    ) -> Result<Vec<AssetUnlockStatusResult>, Error>;
79
80    /// Get transaction
81    fn get_transaction_extended_info(&self, tx_id: &Txid)
82        -> Result<GetRawTransactionResult, Error>;
83
84    /// Get optional transaction extended info
85    /// Returns None if transaction doesn't exists
86    fn get_optional_transaction_extended_info(
87        &self,
88        transaction_id: &Txid,
89    ) -> Result<Option<GetRawTransactionResult>, Error> {
90        match self.get_transaction_extended_info(transaction_id) {
91            Ok(transaction_info) => Ok(Some(transaction_info)),
92            // Return None if transaction with specified tx id is not present
93            Err(Error::JsonRpc(dpp::dashcore_rpc::jsonrpc::error::Error::Rpc(
94                dpp::dashcore_rpc::jsonrpc::error::RpcError {
95                    code: CORE_RPC_INVALID_ADDRESS_OR_KEY,
96                    ..
97                },
98            ))) => Ok(None),
99            Err(e) => Err(e),
100        }
101    }
102
103    /// Get block by hash
104    fn get_fork_info(&self, name: &str) -> Result<Option<SoftforkInfo>, Error>;
105
106    /// Get block by hash
107    fn get_block(&self, block_hash: &BlockHash) -> Result<Block, Error>;
108
109    /// Get block by hash in JSON format
110    fn get_block_json(&self, block_hash: &BlockHash) -> Result<Value, Error>;
111
112    /// Get chain tips
113    fn get_chain_tips(&self) -> Result<GetChainTipsResult, Error>;
114
115    /// Get list of quorums by type at a given height.
116    ///
117    /// See <https://dashcore.readme.io/v19.0.0/docs/core-api-ref-remote-procedure-calls-evo#quorum-listextended>
118    fn get_quorum_listextended(
119        &self,
120        height: Option<CoreHeight>,
121    ) -> Result<ExtendedQuorumListResult, Error>;
122
123    /// Get quorum information.
124    ///
125    /// See <https://dashcore.readme.io/v19.0.0/docs/core-api-ref-remote-procedure-calls-evo#quorum-info>
126    fn get_quorum_info(
127        &self,
128        quorum_type: QuorumType,
129        hash: &QuorumHash,
130        include_secret_key_share: Option<bool>,
131    ) -> Result<QuorumInfoResult, Error>;
132
133    /// Get the difference in masternode list, return masternodes as diff elements
134    fn get_protx_diff_with_masternodes(
135        &self,
136        base_block: Option<u32>,
137        block: u32,
138    ) -> Result<MasternodeListDiff, Error>;
139
140    // /// Get the detailed information about a deterministic masternode
141    // fn get_protx_info(&self, pro_tx_hash: &ProTxHash) -> Result<ProTxInfo, Error>;
142
143    /// Verify Instant Lock signature
144    /// If `max_height` is provided the chain lock will be verified
145    /// against quorums available at this height
146    fn verify_instant_lock(
147        &self,
148        instant_lock: &InstantLock,
149        max_height: Option<u32>,
150    ) -> Result<bool, Error>;
151
152    /// Verify a chain lock signature
153    fn verify_chain_lock(&self, chain_lock: &ChainLock) -> Result<bool, Error>;
154
155    /// Returns masternode sync status
156    fn masternode_sync_status(&self) -> Result<MnSyncStatus, Error>;
157
158    /// Sends raw transaction to the network
159    fn send_raw_transaction(&self, transaction: &[u8]) -> Result<Txid, Error>;
160
161    /// Get Core's credit pool balance after the block at `height`, in duffs, read from the
162    /// block's coinbase (only the coinbase is transferred). Only ask for a chain locked height:
163    /// the answer is then the same on every node.
164    fn get_credit_pool_balance(&self, height: CoreHeight) -> Result<u64, Error>;
165}
166
167#[derive(Debug)]
168/// Default implementation of Dash Core RPC using DashCoreRPC client
169pub struct DefaultCoreRPC {
170    inner: Client,
171    /// Speculative fetcher for the next core height, on its own connection.
172    /// `None` when a second connection could not be opened.
173    prefetcher: Option<CorePrefetcher>,
174}
175
176// TODO: Create errors for these error codes in dashcore_rpc
177
178/// TX is invalid due to consensus rules
179pub const CORE_RPC_TX_CONSENSUS_ERROR: i32 = -26;
180/// Tx already broadcasted and included in the chain
181pub const CORE_RPC_TX_ALREADY_IN_CHAIN: i32 = -27;
182/// Client still warming up
183pub const CORE_RPC_ERROR_IN_WARMUP: i32 = -28;
184/// Dash is not connected
185pub const CORE_RPC_CLIENT_NOT_CONNECTED: i32 = -9;
186/// Still downloading initial blocks
187pub const CORE_RPC_CLIENT_IN_INITIAL_DOWNLOAD: i32 = -10;
188/// Parse error
189pub const CORE_RPC_PARSE_ERROR: i32 = -32700;
190/// Invalid address or key
191pub const CORE_RPC_INVALID_ADDRESS_OR_KEY: i32 = -5;
192/// Invalid, missing or duplicate parameter
193pub const CORE_RPC_INVALID_PARAMETER: i32 = -8;
194
195macro_rules! retry {
196    ($action:expr) => {{
197        /// Maximum number of retry attempts
198        const MAX_RETRIES: u32 = 4;
199        /// // Multiplier for Fibonacci sequence
200        const FIB_MULTIPLIER: u64 = 1;
201
202        fn fibonacci(n: u32) -> u64 {
203            match n {
204                0 => 0,
205                1 => 1,
206                _ => fibonacci(n - 1) + fibonacci(n - 2),
207            }
208        }
209
210        let mut last_err = None;
211        let result = (0..MAX_RETRIES).find_map(|i| {
212            match $action {
213                Ok(result) => Some(Ok(result)),
214                Err(e) => {
215                    match e {
216                        dpp::dashcore_rpc::Error::JsonRpc(
217                            // Retry on transport connection error
218                            dpp::dashcore_rpc::jsonrpc::error::Error::Transport(_)
219                            | dpp::dashcore_rpc::jsonrpc::error::Error::Rpc(
220                                // Retry on Core RPC "not ready" errors
221                                dpp::dashcore_rpc::jsonrpc::error::RpcError {
222                                    code:
223                                        CORE_RPC_ERROR_IN_WARMUP
224                                        | CORE_RPC_CLIENT_NOT_CONNECTED
225                                        | CORE_RPC_CLIENT_IN_INITIAL_DOWNLOAD,
226                                    ..
227                                },
228                            ),
229                        ) => {
230                            // Delay before next try
231                            last_err = Some(e);
232                            let delay = fibonacci(i + 2) * FIB_MULTIPLIER;
233                            std::thread::sleep(Duration::from_secs(delay));
234                            None
235                        }
236                        _ => Some(Err(e)),
237                    }
238                }
239            }
240        });
241
242        result.unwrap_or_else(|| Err(last_err.unwrap()))
243    }};
244}
245
246impl DefaultCoreRPC {
247    /// Create new instance
248    pub fn open(url: &str, username: String, password: String) -> Result<Self, Error> {
249        let prefetcher = CorePrefetcher::new(url, username.clone(), password.clone());
250        if prefetcher.is_none() {
251            tracing::warn!(
252                "could not open a second Core RPC connection; masternode and quorum updates will be fetched on the critical path"
253            );
254        }
255        Ok(DefaultCoreRPC {
256            inner: Client::new(url, Auth::UserPass(username, password))?,
257            prefetcher,
258        })
259    }
260}
261
262impl CoreRPCLike for DefaultCoreRPC {
263    fn get_block_hash(&self, height: u32) -> Result<BlockHash, Error> {
264        retry!(self.inner.get_block_hash(height))
265    }
266
267    fn get_block_header(&self, block_hash: &BlockHash) -> Result<Header, Error> {
268        retry!(self.inner.get_block_header(block_hash))
269    }
270
271    fn get_block_time_from_height(&self, height: CoreHeight) -> Result<TimestampMillis, Error> {
272        let block_hash = self.get_block_hash(height)?;
273        let block_header = self.get_block_header(&block_hash)?;
274        let block_time = block_header.time as u64 * 1000;
275        Ok(block_time)
276    }
277
278    fn get_best_chain_lock(&self) -> Result<ChainLock, Error> {
279        retry!(self.inner.get_best_chain_lock())
280    }
281
282    fn submit_chain_lock(&self, chain_lock: &ChainLock) -> Result<u32, Error> {
283        retry!(self.inner.submit_chain_lock(chain_lock))
284    }
285
286    fn get_transaction(&self, tx_id: &Txid) -> Result<Transaction, Error> {
287        retry!(self.inner.get_raw_transaction(tx_id, None))
288    }
289
290    fn get_transaction_extended_info(
291        &self,
292        tx_id: &Txid,
293    ) -> Result<GetRawTransactionResult, Error> {
294        retry!(self.inner.get_raw_transaction_info(tx_id, None))
295    }
296
297    fn get_fork_info(&self, name: &str) -> Result<Option<SoftforkInfo>, Error> {
298        retry!(self
299            .inner
300            .get_blockchain_info()
301            .map(|blockchain_info| blockchain_info.softforks.get(name).cloned()))
302    }
303
304    fn get_block(&self, block_hash: &BlockHash) -> Result<Block, Error> {
305        retry!(self.inner.get_block(block_hash))
306    }
307
308    fn get_block_json(&self, block_hash: &BlockHash) -> Result<Value, Error> {
309        retry!(self.inner.get_block_json(block_hash))
310    }
311
312    fn get_chain_tips(&self) -> Result<GetChainTipsResult, Error> {
313        retry!(self.inner.get_chain_tips())
314    }
315
316    fn get_quorum_listextended(
317        &self,
318        height: Option<CoreHeight>,
319    ) -> Result<ExtendedQuorumListResult, Error> {
320        // Block sync walks core heights in order, so the next call is almost
321        // always for height + 1. Take the speculative answer when it is for the
322        // height we were asked about, and start the next guess either way. The
323        // prefetcher declines a guess past the chain lock, so at the tip this is
324        // a no-op until Core locks the next block.
325        let prefetched = height
326            .zip(self.prefetcher.as_ref())
327            .and_then(|(height, prefetcher)| prefetcher.take_quorum_list(height));
328
329        let result = match prefetched {
330            Some(list) => Ok(list),
331            None => retry!(self.inner.get_quorum_listextended_reversed(height)),
332        };
333
334        if let (Ok(_), Some(height), Some(prefetcher)) = (&result, height, self.prefetcher.as_ref())
335        {
336            prefetcher.start_quorum_list(height + 1);
337        }
338
339        result
340    }
341
342    fn get_quorum_info(
343        &self,
344        quorum_type: QuorumType,
345        hash: &QuorumHash,
346        include_secret_key_share: Option<bool>,
347    ) -> Result<QuorumInfoResult, Error> {
348        retry!(self
349            .inner
350            .get_quorum_info_reversed(quorum_type, hash, include_secret_key_share))
351    }
352
353    fn get_protx_diff_with_masternodes(
354        &self,
355        base_block: Option<u32>,
356        block: u32,
357    ) -> Result<MasternodeListDiff, Error> {
358        let base = base_block.unwrap_or(1);
359
360        // Same reasoning as get_quorum_listextended: the next diff a syncing
361        // node asks for is from this block to the one after it.
362        let prefetched = self
363            .prefetcher
364            .as_ref()
365            .and_then(|prefetcher| prefetcher.take_protx_diff(base, block));
366
367        let result = match prefetched {
368            Some(diff) => Ok(diff),
369            None => retry!(self.inner.get_protx_listdiff(base, block)),
370        };
371
372        if let (Ok(_), Some(prefetcher)) = (&result, self.prefetcher.as_ref()) {
373            prefetcher.start_protx_diff(block, block + 1);
374        }
375
376        result
377    }
378
379    /// Verify Instant Lock signature
380    /// If `max_height` is provided the chain lock will be verified
381    /// against quorums available at this height
382    fn verify_instant_lock(
383        &self,
384        instant_lock: &InstantLock,
385        max_height: Option<u32>,
386    ) -> Result<bool, Error> {
387        let request_id = instant_lock.request_id()?.to_string();
388        let transaction_id = instant_lock.txid.to_string();
389        let signature = hex::encode(instant_lock.signature);
390
391        retry!(self
392            .inner
393            .get_verifyislock(&request_id, &transaction_id, &signature, max_height))
394    }
395
396    /// Verify a chain lock signature
397    fn verify_chain_lock(&self, chain_lock: &ChainLock) -> Result<bool, Error> {
398        let block_hash = chain_lock.block_hash.to_string();
399        let signature = hex::encode(chain_lock.signature);
400
401        retry!(self.inner.get_verifychainlock(
402            block_hash.as_str(),
403            &signature,
404            Some(chain_lock.block_height)
405        ))
406    }
407
408    /// Returns masternode sync status
409    fn masternode_sync_status(&self) -> Result<MnSyncStatus, Error> {
410        retry!(self.inner.mnsync_status())
411    }
412
413    fn send_raw_transaction(&self, transaction: &[u8]) -> Result<Txid, Error> {
414        retry!(self.inner.send_raw_transaction(transaction))
415    }
416
417    fn get_asset_unlock_statuses(
418        &self,
419        indices: &[u64],
420        core_chain_locked_height: u32,
421    ) -> Result<Vec<AssetUnlockStatusResult>, Error> {
422        retry!(self
423            .inner
424            .get_asset_unlock_statuses(indices, Some(core_chain_locked_height)))
425    }
426
427    fn get_credit_pool_balance(&self, height: CoreHeight) -> Result<u64, Error> {
428        let block_hash = self.get_block_hash(height)?;
429        // Only the coinbase, as raw hex: the special transactions of type 5, the first one,
430        // verbosity 1. An empty answer is a coinbase without a payload, before DIP3.
431        let args = [
432            Value::String(block_hash.to_string()),
433            Value::from(COINBASE_TRANSACTION_TYPE),
434            Value::from(1),
435            Value::from(0),
436            Value::from(1),
437        ];
438        let coinbases = retry!(self.inner.call::<Vec<String>>("getspecialtxes", &args))?;
439        let Some(coinbase) = coinbases.first() else {
440            return Ok(0);
441        };
442        let coinbase = hex::decode(coinbase).map_err(|e| {
443            Error::UnexpectedStructure(format!("getspecialtxes answered invalid hex: {e}"))
444        })?;
445        credit_pool_balance_from_coinbase(&coinbase).map_err(Error::UnexpectedStructure)
446    }
447}
448
449#[cfg(test)]
450mod tests {
451    use super::credit_pool_balance_from_coinbase;
452    use dpp::dashcore::bls_sig_utils::BLSSignature;
453    use dpp::dashcore::consensus::serialize;
454    use dpp::dashcore::hash_types::{MerkleRootMasternodeList, MerkleRootQuorums};
455    use dpp::dashcore::hashes::Hash;
456    use dpp::dashcore::transaction::special_transaction::coinbase::CoinbasePayload;
457    use dpp::dashcore::transaction::special_transaction::TransactionPayload;
458    use dpp::dashcore::{OutPoint, ScriptBuf, Transaction, TxIn, TxOut};
459
460    const BALANCE_DUFFS: u64 = 3_700_000_000_000;
461
462    /// A coinbase as the pinned rust-dashcore encodes it.
463    fn coinbase(payload: Option<TransactionPayload>) -> Vec<u8> {
464        serialize(&Transaction {
465            version: 3,
466            lock_time: 0,
467            input: vec![TxIn {
468                previous_output: OutPoint::null(),
469                script_sig: ScriptBuf::from(vec![0x51, 0x51]),
470                sequence: u32::MAX,
471                witness: Default::default(),
472            }],
473            output: vec![TxOut {
474                value: 5_000_000,
475                script_pubkey: ScriptBuf::from(vec![0x76; 25]),
476            }],
477            special_transaction_payload: payload,
478        })
479    }
480
481    fn version_3_payload() -> TransactionPayload {
482        TransactionPayload::CoinbasePayloadType(CoinbasePayload {
483            version: 3,
484            height: 1_000,
485            merkle_root_masternode_list: MerkleRootMasternodeList::from_byte_array([1; 32]),
486            merkle_root_quorums: MerkleRootQuorums::from_byte_array([2; 32]),
487            best_cl_height: Some(30),
488            best_cl_signature: Some(BLSSignature::from([3; 96])),
489            asset_locked_amount: Some(BALANCE_DUFFS),
490        })
491    }
492
493    #[test]
494    fn should_read_the_balance_of_a_version_3_coinbase() {
495        assert_eq!(
496            credit_pool_balance_from_coinbase(&coinbase(Some(version_3_payload()))),
497            Ok(BALANCE_DUFFS)
498        );
499    }
500
501    /// Core v24 blocks carry a version 4 payload, which appends `merkleRootAssetUnlocks`
502    /// after the balance; the pinned transaction decoder cannot read it.
503    #[test]
504    fn should_read_the_balance_of_a_version_4_coinbase() {
505        let version_3 = coinbase(Some(version_3_payload()));
506        // The payload is last: its 1-byte length (175), then the payload itself.
507        let payload_start = version_3.len() - 175;
508        assert_eq!(version_3[payload_start - 1], 175);
509        let mut version_4 = version_3[..payload_start - 1].to_vec();
510        version_4.push(175 + 32);
511        version_4.extend_from_slice(&4u16.to_le_bytes());
512        version_4.extend_from_slice(&version_3[payload_start + 2..]);
513        version_4.extend_from_slice(&[0xaa; 32]);
514
515        assert_eq!(
516            credit_pool_balance_from_coinbase(&version_4),
517            Ok(BALANCE_DUFFS)
518        );
519    }
520
521    #[test]
522    fn should_read_no_balance_from_a_coinbase_before_the_credit_pool() {
523        let version_2 = TransactionPayload::CoinbasePayloadType(CoinbasePayload {
524            version: 2,
525            height: 1_000,
526            merkle_root_masternode_list: MerkleRootMasternodeList::from_byte_array([1; 32]),
527            merkle_root_quorums: MerkleRootQuorums::from_byte_array([2; 32]),
528            best_cl_height: None,
529            best_cl_signature: None,
530            asset_locked_amount: None,
531        });
532        assert_eq!(
533            credit_pool_balance_from_coinbase(&coinbase(Some(version_2))),
534            Ok(0)
535        );
536        assert_eq!(credit_pool_balance_from_coinbase(&coinbase(None)), Ok(0));
537    }
538
539    #[test]
540    fn should_fail_on_a_truncated_coinbase() {
541        let full = coinbase(Some(version_3_payload()));
542        // Cut inside the payload, before the balance.
543        assert!(credit_pool_balance_from_coinbase(&full[..full.len() - 60]).is_err());
544        assert!(credit_pool_balance_from_coinbase(&full[..40]).is_err());
545    }
546}