use crate::chain::Chain;
use crate::chaindef::Block;
use crate::chaindef::BlockHeader;
use crate::chaindef::OutPointHash;
use crate::discovery::announcer::PeerAnnouncer;
use crate::discovery::discoverer::PeerDiscoverer;
use crate::index::get_last_indexed_block;
use crate::indexes::outputindex::OutputIndexRow;
use crate::indexes::scripthashindex::ScriptHashIndexRow;
use crate::indexes::DBRow;
use crate::signal::Waiter;
use crate::store::DBStore;
use futures_util::stream;
use queryfilter::QueryFilter;
use rayon::prelude::*;
use crate::chaindef::ScriptHash;
use crate::chaindef::Transaction;
use crate::metrics;
use bitcoincash::hash_types::{TxMerkleNode, Txid};
use bitcoincash::hashes::sha256d::Hash as Sha256dHash;
use bitcoincash::hashes::Hash;
use bitcoincash::network::constants::Network;
use futures::stream::StreamExt;
use serde_json::Value;
use std::collections::HashSet;
use std::sync::Arc;
use std::time::Duration;
use std::time::Instant;
use crate::app::App;
use crate::cache::TransactionCache;
use crate::chaindef::BlockHash;
use crate::mempool::Tracker;
use crate::metrics::Metrics;
#[cfg(nexa)]
use crate::nexa::hash_types::TxIdem;
use crate::query::header::HeaderQuery;
use crate::query::tx::TxQuery;
use anyhow::{Context, Result};
use self::queryutil::get_tx_funding_prevout;
use self::queryutil::get_tx_spending_prevout;
pub mod balance;
pub mod header;
pub mod queryfilter;
pub mod queryutil;
pub mod scripthash_subscriptions;
pub mod status;
pub mod tip;
pub mod token;
pub mod tx;
pub mod unspent;
pub mod volume;
pub const BUFFER_SIZE: usize = 10;
pub const CHUNK_SIZE: usize = 1000;
fn merklize<T: Hash>(left: T, right: T) -> T {
let data = [&left[..], &right[..]].concat();
<T as Hash>::hash(&data)
}
fn create_merkle_branch_and_root<T: Hash>(mut hashes: Vec<T>, mut index: usize) -> (Vec<T>, T) {
let mut merkle = vec![];
while hashes.len() > 1 {
if !hashes.len().is_multiple_of(2) {
let last = *hashes.last().unwrap();
hashes.push(last);
}
index = if index.is_multiple_of(2) {
index + 1
} else {
index - 1
};
merkle.push(hashes[index]);
index /= 2;
hashes = hashes
.chunks(2)
.map(|pair| merklize(pair[0], pair[1]))
.collect()
}
(merkle, hashes[0])
}
pub struct Query {
app: Arc<App>,
tracker: Arc<Tracker>,
duration: Arc<prometheus::HistogramVec>,
tx: TxQuery,
header: Arc<HeaderQuery>,
}
impl Query {
pub fn new(
app: Arc<App>,
metrics: &Metrics,
tx_cache: TransactionCache,
tracker: Arc<Tracker>,
network: Network,
) -> Result<Arc<Query>> {
let daemon = app.daemon().reconnect()?;
let duration = Arc::new(metrics.histogram_vec(
"rostrum_query_duration",
"Request duration (in seconds)",
&["type"],
metrics::default_duration_buckets(),
));
let header = Arc::new(HeaderQuery::new(app.clone()));
let tx = TxQuery::new(
tx_cache,
daemon,
tracker.clone(),
header.clone(),
duration.clone(),
network,
Arc::clone(app.index().read_store()),
);
Ok(Arc::new(Query {
app,
tracker,
duration,
tx,
header,
}))
}
pub(crate) fn announcer(&self) -> &PeerAnnouncer {
self.app.announcer()
}
pub(crate) fn chain(&self) -> &Chain {
self.app.index().chain()
}
pub(crate) fn discoverer(&self) -> &PeerDiscoverer {
self.app.discoverer()
}
#[cfg(nexa)]
async fn get_block_height(&self, header: &BlockHeader) -> Option<u64> {
Some(header.height())
}
#[cfg(bch)]
async fn get_block_height(&self, header: &BlockHeader) -> Option<u64> {
self.chain().get_block_height(&header.block_hash()).await
}
pub async fn get_confirmed_blockhash(&self, tx_hash: &Txid) -> Result<Value> {
let header = self.header.get_by_txid(tx_hash, None).await?;
if header.is_none() {
bail!("tx {} is unconfirmed or does not exist", tx_hash);
}
let header = header.unwrap();
let height = self
.get_block_height(&header)
.await
.context(anyhow!("No height for block {}", header.block_hash()))?;
Ok(json!({
"block_hash": header.block_hash(),
"block_height": height,
}))
}
pub async fn get_headers(&self, heights: &[usize]) -> Vec<BlockHeader> {
let _timer = self
.duration
.with_label_values(&["get_headers"])
.start_timer();
let index = self.app.index();
stream::iter(heights)
.filter_map(|height| async move { index.chain().get_block_header(*height).await })
.collect()
.await
}
pub(crate) async fn get_best_header(&self) -> Result<(u64, BlockHeader)> {
let now = Instant::now();
loop {
let last_indexed = get_last_indexed_block(self.confirmed_index())
.context("failed to fetch last indexed block".to_string())?;
match self
.app
.index()
.chain()
.get_block_height_and_header(&last_indexed)
.await
{
Some((height, header)) => {
return Ok((height, header));
}
None => {
if now.elapsed() > Duration::from_secs(10) {
bail!("Failed to get best header (busy with reorg?)");
}
Waiter::wait_or_shutdown(Duration::from_millis(10)).await?;
}
};
}
}
pub async fn getblocktxids(&self, blockhash: &BlockHash) -> Result<Vec<Txid>> {
self.app.daemon().getblocktxids(blockhash).await
}
pub async fn get_block(&self, blockhash: &BlockHash) -> Result<Block> {
self.app.daemon().getblock(blockhash).await
}
pub async fn get_merkle_proof(
&self,
tx_hash: &Txid,
height: usize,
) -> Result<(Vec<TxMerkleNode>, usize)> {
let header_entry = self
.app
.index()
.chain()
.get_block_header(height)
.await
.context(format!("missing block #{}", height))?;
let txids = self
.app
.daemon()
.getblocktxids(&header_entry.block_hash())
.await?;
let pos = txids
.iter()
.position(|txid| txid == tx_hash)
.context(format!("missing txid {}", tx_hash))?;
let tx_nodes: Vec<TxMerkleNode> = txids
.into_iter()
.map(|txid| TxMerkleNode::from_inner(txid.into_inner()))
.collect();
let (branch, _root) = create_merkle_branch_and_root(tx_nodes, pos);
Ok((branch, pos))
}
pub async fn get_header_merkle_proof(
&self,
height: usize,
cp_height: usize,
) -> Result<(Vec<Sha256dHash>, Sha256dHash)> {
if cp_height < height {
bail!("cp_height #{} < height #{}", cp_height, height);
}
let best_height = self.app.index().chain().height().await;
if best_height < (cp_height as u64) {
bail!(
"cp_height #{} above best block height #{}",
cp_height,
best_height
);
}
let heights: Vec<usize> = (0..=cp_height).collect();
let header_hashes: Vec<BlockHash> = self
.get_headers(&heights)
.await
.into_par_iter()
.map(|h| h.block_hash())
.collect();
let merkle_nodes: Vec<Sha256dHash> = header_hashes
.par_iter()
.map(|block_hash| Sha256dHash::from_inner(block_hash.into_inner()))
.collect();
assert_eq!(header_hashes.len(), heights.len());
Ok(create_merkle_branch_and_root(merkle_nodes, height))
}
pub async fn get_id_from_pos(
&self,
height: usize,
tx_pos: usize,
want_merkle: bool,
) -> Result<(Txid, Vec<TxMerkleNode>)> {
let header_entry = self
.app
.index()
.chain()
.get_block_header(height)
.await
.context(format!("missing block #{}", height))?;
let txids = self
.app
.daemon()
.getblocktxids(&header_entry.block_hash())
.await?;
let txid = *txids.get(tx_pos).context(format!(
"No tx in position #{} in block #{}",
tx_pos, height
))?;
let tx_nodes = txids
.into_iter()
.map(|txid| TxMerkleNode::from_inner(txid.into_inner()))
.collect();
let branch = if want_merkle {
create_merkle_branch_and_root(tx_nodes, tx_pos).0
} else {
vec![]
};
Ok((txid, branch))
}
#[cfg(bch)]
pub async fn broadcast(&self, txn: &Transaction) -> Result<Txid> {
self.app.daemon().broadcast(txn).await
}
#[cfg(nexa)]
pub async fn broadcast(&self, txn: &Transaction) -> Result<TxIdem> {
self.app.daemon().broadcast(txn).await
}
pub async fn broadcast_p2p(&self, txn: Transaction) -> Result<()> {
self.app.daemon().broadcast_p2p(txn).await
}
pub async fn update_mempool(&self) -> Result<HashSet<Txid>> {
let _timer = self
.duration
.with_label_values(&["update_mempool"])
.start_timer();
self.tracker
.update_from_daemon(self.app.daemon(), self.tx())
.await
}
pub async fn estimate_fee(&self, blocks: usize) -> f64 {
let mut total_vsize = 0u32;
let mut last_fee_rate = 0.0;
let blocks_in_vbytes = (blocks * 1_000_000) as u32; for (fee_rate, vsize) in self.tracker.fee_histogram().await {
last_fee_rate = fee_rate;
total_vsize += vsize;
if total_vsize >= blocks_in_vbytes {
break; }
}
(last_fee_rate as f64) * 1e-5 }
pub async fn get_banner(&self) -> Result<String> {
self.app.get_banner().await
}
pub async fn scripthash_first_use(
&self,
scripthash: ScriptHash,
filter: &QueryFilter,
) -> Result<Option<(u32, Txid)>> {
let get_tx = |store: Arc<DBStore>| async move {
let (scan_prefix, scan_bitmask) =
ScriptHashIndexRow::filter_by_outputs(scripthash.into_inner(), filter);
let (query, stream) = store
.scan(ScriptHashIndexRow::CF, scan_prefix, scan_bitmask)
.await;
let mut best_height = u32::MAX;
let mut best_candidates: Vec<OutPointHash> = vec![];
stream
.for_each(|o| {
let row = ScriptHashIndexRow::from_row(&o);
let h = row.get_height();
match h.cmp(&best_height) {
std::cmp::Ordering::Less => {
best_height = h;
best_candidates = vec![row.outpointhash()];
}
std::cmp::Ordering::Equal => {
best_candidates.push(row.outpointhash());
}
std::cmp::Ordering::Greater => {
}
}
futures::future::ready(())
})
.await;
query.await??;
if best_candidates.is_empty() {
return Ok(None);
}
let outpoint_keys: Vec<Vec<u8>> = best_candidates
.iter()
.map(OutputIndexRow::filter_by_outpointhash)
.collect();
let outpoint_values = store.multi_get(OutputIndexRow::CF, outpoint_keys).await;
let txids: Vec<Txid> = outpoint_values
.into_iter()
.flatten()
.map(|v| OutputIndexRow::txid_from_value(&v).unwrap())
.collect();
let txid = txids.into_iter().min_by(|a, b| b.cmp(a)).unwrap();
Ok(Some((best_height, txid)))
};
if let Some(tx) = get_tx(self.app.index().read_store().clone()).await? {
return Ok(Some(tx));
}
let mempool = self.tracker.index().clone();
get_tx(mempool).await
}
pub async fn scripthash_first_spend(
&self,
scripthash: ScriptHash,
filter: &QueryFilter,
) -> Result<Option<(u32, Txid)>> {
let get_tx = |store: Arc<DBStore>| async move {
let (scan_prefix, scan_bitmask) =
ScriptHashIndexRow::filter_by_spends(scripthash.into_inner(), filter);
let (query, mut stream) = store
.scan(ScriptHashIndexRow::CF, scan_prefix, scan_bitmask)
.await;
let first_row = stream.next().await;
match first_row {
Some(row) => {
query.await??;
let row = ScriptHashIndexRow::from_row(&row);
let txid = row
.spending_txid()
.ok_or_else(|| anyhow::anyhow!("spending row has no txid"))?;
Ok(Some((row.get_height(), Txid::from_inner(txid))))
}
None => {
query.await??;
Ok(None)
}
}
};
if let Some(tx) = get_tx(self.app.index().read_store().clone()).await? {
return Ok(Some(tx));
}
let mempool = self.tracker.index().clone();
get_tx(mempool).await
}
pub async fn get_relayfee(&self) -> Result<f64> {
self.app.daemon().get_relayfee().await
}
pub fn tx(&self) -> &TxQuery {
&self.tx
}
pub fn header(&self) -> &HeaderQuery {
&self.header
}
pub async fn get_tx_spending_prevout(
&self,
outpoint: &OutPointHash,
) -> Result<
Option<(
Transaction,
u32, /* input index */
u32, /* confirmation height */
)>,
> {
{
let mempool = self.tracker.index();
let spent = get_tx_spending_prevout(mempool, &self.tx, outpoint).await?;
if spent.is_some() {
return Ok(spent);
}
}
let store = self.app.index().read_store();
get_tx_spending_prevout(store, self.tx(), outpoint).await
}
pub async fn get_tx_funding_prevout(
&self,
prevout: &OutPointHash,
) -> Result<
Option<(
Transaction,
u32, /* output index */
u32, /* confirmation height */
)>,
> {
{
let mempool = self.tracker.index();
let spent = get_tx_funding_prevout(mempool, &self.tx, prevout).await?;
if spent.is_some() {
return Ok(spent);
}
}
get_tx_funding_prevout(self.app.index().read_store(), &self.tx, prevout).await
}
pub fn confirmed_index(&self) -> &Arc<DBStore> {
self.app.index().read_store()
}
pub fn unconfirmed_index(&self) -> &Tracker {
&self.tracker
}
pub fn unconfirmed_index_cloned(&self) -> Arc<Tracker> {
self.tracker.clone()
}
}