use std::time::Duration;
use tokio::sync::watch;
use zaino_primitives::types::{
rpc, AddressBalance, AddressDelta, Block, BlockHash, BlockVerbose, BlockchainInfo, Difficulty,
Height, OutputIndex, PreIndexCompactBlock, ShieldedPool, SubtreeRoot, TransactionId, TreeRoots,
Treestate, Utxo,
};
use zaino_source::*;
use zaino_source_zebra_readstate::ZebraReadStateAdapter;
use zaino_source_zebra_rpc::ZebraRpcAdapter;
use crate::fallback::retry_on_slow_path;
pub struct ZebraValidator {
rpc: ZebraRpcAdapter,
readstate: Option<ZebraReadStateAdapter>,
tip: Option<PolledChainTip>,
}
impl ZebraValidator {
pub fn rpc_only(rpc: ZebraRpcAdapter) -> Self {
Self {
rpc,
readstate: None,
tip: None,
}
}
pub fn with_read_state(rpc: ZebraRpcAdapter, readstate: ZebraReadStateAdapter) -> Self {
Self {
rpc,
readstate: Some(readstate),
tip: None,
}
}
pub async fn with_tip_polling<S>(
mut self,
source: S,
interval: Duration,
) -> Result<Self, QueryError<GetChainTipError>>
where
S: OneShotGetChainTip + Send + 'static,
{
self.tip = Some(PolledChainTip::spawn(source, interval).await?);
Ok(self)
}
fn fast(&self) -> Option<&ZebraReadStateAdapter> {
self.readstate.as_ref()
}
#[cfg(feature = "test_dependencies")]
pub fn read_state(&self) -> Option<&ZebraReadStateAdapter> {
self.readstate.as_ref()
}
}
macro_rules! fast_then_slow {
($self:ident, $method:ident $(, $arg:expr)*) => {{
if let Some(fast) = $self.fast() {
let result = fast.$method($($arg),*).await;
if !retry_on_slow_path(&result) {
return result;
}
}
$self.rpc.$method($($arg),*).await
}};
}
macro_rules! fast_or_slow {
($self:ident, $method:ident $(, $arg:expr)*) => {{
match $self.fast() {
Some(fast) => fast.$method($($arg),*).await,
None => $self.rpc.$method($($arg),*).await,
}
}};
}
impl OneShotGetBlock for ZebraValidator {
async fn get_block(&self, height: Height) -> Result<Block, QueryError<GetBlockError>> {
fast_or_slow!(self, get_block, height)
}
}
impl OneShotGetBlockByHash for ZebraValidator {
async fn get_block_by_hash(
&self,
hash: BlockHash,
) -> Result<Block, QueryError<GetBlockByHashError>> {
fast_then_slow!(self, get_block_by_hash, hash)
}
}
impl OneShotGetRawBlock for ZebraValidator {
async fn get_raw_block(&self, height: Height) -> Result<Vec<u8>, QueryError<GetBlockError>> {
fast_or_slow!(self, get_raw_block, height)
}
}
impl OneShotGetRawBlockByHash for ZebraValidator {
async fn get_raw_block_by_hash(
&self,
hash: BlockHash,
) -> Result<Vec<u8>, QueryError<GetBlockByHashError>> {
fast_then_slow!(self, get_raw_block_by_hash, hash)
}
}
impl OneShotGetChainTip for ZebraValidator {
async fn get_chain_tip(&self) -> Result<(BlockHash, Height), QueryError<GetChainTipError>> {
fast_or_slow!(self, get_chain_tip)
}
}
impl OneShotGetBestBlockHeight for ZebraValidator {
async fn get_best_block_height(&self) -> Result<Height, QueryError<GetBestBlockHeightError>> {
fast_or_slow!(self, get_best_block_height)
}
}
impl OneShotGetPreIndexCompactBlock for ZebraValidator {
async fn get_pre_index_compact_block(
&self,
height: Height,
) -> Result<PreIndexCompactBlock, QueryError<GetBlockError>> {
fast_or_slow!(self, get_pre_index_compact_block, height)
}
}
impl OneShotGetTransaction for ZebraValidator {
async fn get_transaction(
&self,
txid: TransactionId,
) -> Result<TransactionResponse, QueryError<GetTransactionError>> {
fast_then_slow!(self, get_transaction, txid)
}
}
impl OneShotGetTreestate for ZebraValidator {
async fn get_treestate(
&self,
height: Height,
) -> Result<Treestate, QueryError<GetTreestateError>> {
fast_or_slow!(self, get_treestate, height)
}
}
impl OneShotGetTreestateByHash for ZebraValidator {
async fn get_treestate_by_hash(
&self,
hash: BlockHash,
) -> Result<Treestate, QueryError<GetTreestateByHashError>> {
fast_or_slow!(self, get_treestate_by_hash, hash)
}
}
impl OneShotGetCommitmentTreeRoots for ZebraValidator {
async fn get_commitment_tree_roots(
&self,
block: BlockHash,
) -> Result<TreeRoots, QueryError<GetCommitmentTreeRootsError>> {
fast_or_slow!(self, get_commitment_tree_roots, block)
}
}
impl OneShotGetSubtreeRoots for ZebraValidator {
async fn get_subtree_roots(
&self,
pool: ShieldedPool,
start_index: u16,
limit: Option<u16>,
) -> Result<Vec<SubtreeRoot>, QueryError<GetSubtreeRootsError>> {
fast_or_slow!(self, get_subtree_roots, pool, start_index, limit)
}
}
impl OneShotGetAddressBalance for ZebraValidator {
async fn get_address_balance(
&self,
addresses: Vec<String>,
) -> Result<AddressBalance, QueryError<GetAddressBalanceError>> {
match self.fast() {
Some(fast) => fast.get_address_balance(addresses).await,
None => self.rpc.get_address_balance(addresses).await,
}
}
}
impl OneShotGetAddressTxids for ZebraValidator {
async fn get_address_txids(
&self,
addresses: Vec<String>,
start: Height,
end: Height,
) -> Result<Vec<TransactionId>, QueryError<GetAddressTxidsError>> {
match self.fast() {
Some(fast) => fast.get_address_txids(addresses, start, end).await,
None => self.rpc.get_address_txids(addresses, start, end).await,
}
}
}
impl OneShotGetAddressUtxos for ZebraValidator {
async fn get_address_utxos(
&self,
addresses: Vec<String>,
) -> Result<Vec<Utxo>, QueryError<GetAddressUtxosError>> {
match self.fast() {
Some(fast) => fast.get_address_utxos(addresses).await,
None => self.rpc.get_address_utxos(addresses).await,
}
}
}
impl OneShotGetAddressDeltas for ZebraValidator {
async fn get_address_deltas(
&self,
addresses: Vec<String>,
start: Height,
end: Height,
) -> Result<Vec<AddressDelta>, QueryError<GetAddressDeltasError>> {
match self.fast() {
Some(fast) => fast.get_address_deltas(addresses, start, end).await,
None => self.rpc.get_address_deltas(addresses, start, end).await,
}
}
}
impl OneShotGetMempoolTxids for ZebraValidator {
async fn get_mempool_txids(
&self,
) -> Result<Vec<TransactionId>, QueryError<GetMempoolTxidsError>> {
self.rpc.get_mempool_txids().await
}
}
impl OneShotGetMempoolMetadata for ZebraValidator {
async fn get_mempool_metadata(
&self,
) -> Result<Vec<MempoolTxMeta>, QueryError<GetMempoolMetadataError>> {
self.rpc.get_mempool_metadata().await
}
}
impl OneShotGetRawMempoolTransaction for ZebraValidator {
async fn get_raw_mempool_transaction(
&self,
txid: TransactionId,
) -> Result<Vec<u8>, QueryError<GetRawMempoolTransactionError>> {
self.rpc.get_raw_mempool_transaction(txid).await
}
}
impl OneShotGetMempoolSourceTip for ZebraValidator {
async fn get_mempool_source_tip(
&self,
) -> Result<(BlockHash, Height), QueryError<std::convert::Infallible>> {
self.rpc.get_mempool_source_tip().await
}
}
impl OneShotGetChainTips for ZebraValidator {
async fn get_chain_tips(&self) -> Result<Vec<rpc::ChainTip>, QueryError<GetChainTipsError>> {
self.rpc.get_chain_tips().await
}
}
impl OneShotGetBlockVerbose for ZebraValidator {
async fn get_block_verbose(
&self,
height: Height,
) -> Result<BlockVerbose, QueryError<GetBlockVerboseError>> {
self.rpc.get_block_verbose(height).await
}
}
impl OneShotGetBlockVerboseByHash for ZebraValidator {
async fn get_block_verbose_by_hash(
&self,
hash: BlockHash,
) -> Result<BlockVerbose, QueryError<GetBlockVerboseError>> {
self.rpc.get_block_verbose_by_hash(hash).await
}
}
impl OneShotGetBlockHeader for ZebraValidator {
async fn get_block_header(
&self,
hash: BlockHash,
) -> Result<rpc::BlockHeaderVerbose, QueryError<GetBlockHeaderError>> {
self.rpc.get_block_header(hash).await
}
}
impl OneShotGetRawBlockHeader for ZebraValidator {
async fn get_raw_block_header(
&self,
hash: BlockHash,
) -> Result<Vec<u8>, QueryError<GetBlockHeaderError>> {
self.rpc.get_raw_block_header(hash).await
}
}
impl OneShotGetBlockDeltas for ZebraValidator {
async fn get_block_deltas(
&self,
hash: BlockHash,
) -> Result<rpc::BlockDeltas, QueryError<GetBlockDeltasError>> {
fast_then_slow!(self, get_block_deltas, hash)
}
}
impl OneShotGetBlockSubsidy for ZebraValidator {
async fn get_block_subsidy(
&self,
height: Height,
) -> Result<rpc::BlockSubsidy, QueryError<GetBlockSubsidyError>> {
self.rpc.get_block_subsidy(height).await
}
}
impl OneShotGetNodeInfo for ZebraValidator {
async fn get_node_info(&self) -> Result<rpc::NodeInfo, QueryError<GetNodeInfoError>> {
self.rpc.get_node_info().await
}
}
impl OneShotGetPeerInfo for ZebraValidator {
async fn get_peer_info(&self) -> Result<Vec<rpc::PeerInfo>, QueryError<GetPeerInfoError>> {
self.rpc.get_peer_info().await
}
}
impl OneShotGetMiningInfo for ZebraValidator {
async fn get_mining_info(&self) -> Result<rpc::MiningInfo, QueryError<GetMiningInfoError>> {
self.rpc.get_mining_info().await
}
}
impl OneShotGetNetworkSolPs for ZebraValidator {
async fn get_network_sol_ps(
&self,
blocks: Option<u32>,
height: Option<Height>,
) -> Result<u64, QueryError<GetNetworkSolPsError>> {
self.rpc.get_network_sol_ps(blocks, height).await
}
}
impl OneShotGetTxOut for ZebraValidator {
async fn get_tx_out(
&self,
txid: TransactionId,
index: OutputIndex,
include_mempool: bool,
) -> Result<Option<rpc::TxOut>, QueryError<GetTxOutError>> {
self.rpc.get_tx_out(txid, index, include_mempool).await
}
}
impl OneShotGetSpentInfo for ZebraValidator {
async fn get_spent_info(
&self,
outpoint: rpc::SpentOutpoint,
) -> Result<rpc::SpentInfo, QueryError<GetSpentInfoError>> {
self.rpc.get_spent_info(outpoint).await
}
}
impl OneShotSendRawTransaction for ZebraValidator {
async fn send_raw_transaction(
&self,
transaction: Vec<u8>,
) -> Result<TransactionId, QueryError<SendRawTransactionError>> {
self.rpc.send_raw_transaction(transaction).await
}
}
impl OneShotGetDifficulty for ZebraValidator {
async fn get_difficulty(&self) -> Result<Difficulty, QueryError<GetDifficultyError>> {
fast_or_slow!(self, get_difficulty)
}
}
impl OneShotGetBlockchainInfo for ZebraValidator {
async fn get_blockchain_info(
&self,
) -> Result<BlockchainInfo, QueryError<GetBlockchainInfoError>> {
fast_or_slow!(self, get_blockchain_info)
}
}
impl SubscribeChainTip for ZebraValidator {
fn subscribe_to_chain_tip(&self) -> Option<watch::Receiver<TipObservation>> {
self.readstate
.as_ref()
.and_then(|readstate| readstate.subscribe_to_chain_tip())
.or_else(|| {
self.tip
.as_ref()
.and_then(|tip| tip.subscribe_to_chain_tip())
})
}
}
impl SubscribeBlocks for ZebraValidator {
fn subscribe_to_blocks_received(&self) -> Option<watch::Receiver<()>> {
None
}
}
impl SourceLifecycle for ZebraValidator {
fn shutdown(&self) {
self.rpc.shutdown();
if let Some(readstate) = &self.readstate {
readstate.shutdown();
}
}
}