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
18pub type QuorumListExtendedInfo = HashMap<QuorumHash, ExtendedQuorumDetails>;
20
21const COINBASE_TRANSACTION_TYPE: u16 = 5;
23
24pub(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
50pub type CoreHeight = u32;
52#[cfg_attr(any(feature = "mocks", test), mockall::automock)]
54pub trait CoreRPCLike {
55 fn get_block_hash(&self, height: CoreHeight) -> Result<BlockHash, Error>;
57
58 fn get_block_header(&self, block_hash: &BlockHash) -> Result<Header, Error>;
60
61 fn get_block_time_from_height(&self, height: CoreHeight) -> Result<TimestampMillis, Error>;
63
64 fn get_best_chain_lock(&self) -> Result<ChainLock, Error>;
66
67 fn submit_chain_lock(&self, chain_lock: &ChainLock) -> Result<u32, Error>;
69
70 fn get_transaction(&self, tx_id: &Txid) -> Result<Transaction, Error>;
72
73 fn get_asset_unlock_statuses(
75 &self,
76 indices: &[u64],
77 core_chain_locked_height: u32,
78 ) -> Result<Vec<AssetUnlockStatusResult>, Error>;
79
80 fn get_transaction_extended_info(&self, tx_id: &Txid)
82 -> Result<GetRawTransactionResult, Error>;
83
84 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 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 fn get_fork_info(&self, name: &str) -> Result<Option<SoftforkInfo>, Error>;
105
106 fn get_block(&self, block_hash: &BlockHash) -> Result<Block, Error>;
108
109 fn get_block_json(&self, block_hash: &BlockHash) -> Result<Value, Error>;
111
112 fn get_chain_tips(&self) -> Result<GetChainTipsResult, Error>;
114
115 fn get_quorum_listextended(
119 &self,
120 height: Option<CoreHeight>,
121 ) -> Result<ExtendedQuorumListResult, Error>;
122
123 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 fn get_protx_diff_with_masternodes(
135 &self,
136 base_block: Option<u32>,
137 block: u32,
138 ) -> Result<MasternodeListDiff, Error>;
139
140 fn verify_instant_lock(
147 &self,
148 instant_lock: &InstantLock,
149 max_height: Option<u32>,
150 ) -> Result<bool, Error>;
151
152 fn verify_chain_lock(&self, chain_lock: &ChainLock) -> Result<bool, Error>;
154
155 fn masternode_sync_status(&self) -> Result<MnSyncStatus, Error>;
157
158 fn send_raw_transaction(&self, transaction: &[u8]) -> Result<Txid, Error>;
160
161 fn get_credit_pool_balance(&self, height: CoreHeight) -> Result<u64, Error>;
165}
166
167#[derive(Debug)]
168pub struct DefaultCoreRPC {
170 inner: Client,
171 prefetcher: Option<CorePrefetcher>,
174}
175
176pub const CORE_RPC_TX_CONSENSUS_ERROR: i32 = -26;
180pub const CORE_RPC_TX_ALREADY_IN_CHAIN: i32 = -27;
182pub const CORE_RPC_ERROR_IN_WARMUP: i32 = -28;
184pub const CORE_RPC_CLIENT_NOT_CONNECTED: i32 = -9;
186pub const CORE_RPC_CLIENT_IN_INITIAL_DOWNLOAD: i32 = -10;
188pub const CORE_RPC_PARSE_ERROR: i32 = -32700;
190pub const CORE_RPC_INVALID_ADDRESS_OR_KEY: i32 = -5;
192pub const CORE_RPC_INVALID_PARAMETER: i32 = -8;
194
195macro_rules! retry {
196 ($action:expr) => {{
197 const MAX_RETRIES: u32 = 4;
199 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 dpp::dashcore_rpc::jsonrpc::error::Error::Transport(_)
219 | dpp::dashcore_rpc::jsonrpc::error::Error::Rpc(
220 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 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 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 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 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 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 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 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 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 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 #[test]
504 fn should_read_the_balance_of_a_version_4_coinbase() {
505 let version_3 = coinbase(Some(version_3_payload()));
506 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 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}