use std::sync::Arc;
use crate::{
chaindef::{ScriptHash, TokenID},
indexes::{outputindex::OutputIndexRow, unspentindex::UnspentIndexRow, DBRow},
query::{
queryutil::{multi_get_utxo, outpoints_are_spent, token_from_outpoint},
BUFFER_SIZE, CHUNK_SIZE,
},
store::DBStore,
};
use super::queryfilter::QueryFilter;
use anyhow::{Context, Result};
use bitcoin_hashes::hex::ToHex;
use futures::{stream, try_join, StreamExt};
use rayon::iter::{IntoParallelRefIterator, ParallelIterator};
use tokio::sync::mpsc::Sender;
pub type UtxoEntry = (
UnspentIndexRow,
OutputIndexRow,
Option<(TokenID, i64, Option<Vec<u8>>)>,
);
pub async fn get_utxos_in_db(
db: &Arc<DBStore>,
filter_db: Option<&Arc<DBStore>>,
scripthash: &ScriptHash,
filter: &QueryFilter,
filter_token: &Option<TokenID>,
sender: Sender<Option<UtxoEntry>>,
is_last_task: bool,
) -> Result<()> {
let (query, stream) = db
.scan(
UnspentIndexRow::CF,
UnspentIndexRow::scripthash_filter(scripthash),
None,
)
.await;
let utxos: Vec<UnspentIndexRow> = stream
.ready_chunks(CHUNK_SIZE)
.flat_map_unordered(BUFFER_SIZE, move |utxos| {
let utxos: Vec<UnspentIndexRow> = utxos.iter().map(UnspentIndexRow::from_row).collect();
let utxos = if filter.token_only {
utxos.into_iter().filter(|u| u.has_token()).collect()
} else if filter.exclude_tokens {
utxos.into_iter().filter(|u| !u.has_token()).collect()
} else {
utxos
};
stream::iter(utxos.into_iter())
})
.collect()
.await;
query.await??;
let utxos: Vec<UnspentIndexRow> = if let Some(filter_db) = filter_db {
stream::iter(utxos.into_iter())
.ready_chunks(CHUNK_SIZE)
.then(|utxos| {
let filter_db = filter_db.clone();
async move {
let outpoints = utxos
.par_iter()
.map(|u| u.outpointhash())
.collect::<Vec<_>>();
let spends = outpoints_are_spent(&filter_db, &outpoints).await;
let filtered = utxos
.into_iter()
.zip(spends.into_iter())
.filter_map(|(o, spent)| if spent { None } else { Some(o) })
.collect::<Vec<_>>();
stream::iter(filtered.into_iter())
}
})
.flatten_unordered(BUFFER_SIZE)
.collect()
.await
} else {
utxos
};
let utxos = if filter.exclude_tokens {
utxos.into_iter().map(|u| (u, None)).collect::<Vec<_>>()
} else {
let with_token: Vec<(UnspentIndexRow, Option<_>)> = stream::iter(utxos.into_iter())
.filter_map(|utxo| {
let store = db.clone();
async move {
let token = token_from_outpoint(&store, &utxo.outpointhash())
.await
.expect("failed to fetch token");
if let (Some(ft), Some((t, _, _))) = (filter_token, token.as_ref()) {
if t != ft {
return None;
}
};
Some((utxo, token))
}
})
.collect()
.await;
with_token
};
let mut results = stream::iter(utxos.into_iter())
.ready_chunks(CHUNK_SIZE)
.map(|unspent_with_token| {
let store = db.clone();
let sender = sender.clone();
async move {
let outpoints = unspent_with_token
.iter()
.map(|(unspent, _)| unspent.outpointhash())
.collect::<Vec<_>>();
let utxos = multi_get_utxo(&store, &outpoints).await;
for ((unspent, token), utxo) in
unspent_with_token.into_iter().zip(utxos.into_iter())
{
match utxo {
Some(u) => {
sender
.send(Some((unspent, u, token)))
.await
.context("failed to send UTXO")?;
}
None => {
warn!(
"unspent: did not find utxo {}",
unspent.outpointhash().to_hex()
);
}
}
}
anyhow::Ok(())
}
})
.buffer_unordered(BUFFER_SIZE);
while let Some(result) = results.next().await {
result?;
}
if is_last_task {
sender
.send(None)
.await
.context("failed to send done signal")?;
}
Ok(())
}
pub async fn get_utxos_at_tip(
confirmed: &Arc<DBStore>,
mempool: &Arc<DBStore>,
scripthash: &ScriptHash,
filter: &QueryFilter,
filter_token: &Option<TokenID>,
sender: Sender<Option<UtxoEntry>>,
) -> Result<()>
where
{
assert!(
!filter.has_height_filter(),
"unspent index does not have historical info"
);
try_join!(
get_utxos_in_db(
confirmed,
Some(mempool),
scripthash,
filter,
filter_token,
sender.clone(),
false
),
get_utxos_in_db(
mempool,
None,
scripthash,
filter,
filter_token,
sender.clone(),
false
)
)?;
sender
.send(None)
.await
.context("failed to send done signal")?;
Ok(())
}