rostrum 14.0.1

An efficient implementation of Electrum Server with token support
Documentation
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>>)>,
);

/// Retrive utxos in db, that are not spent in (optional) filter_db.
/// NOTE: If is_last_task is false, caller must call sender.send(None) to signal when the task is done.
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??;

    // filter out any spends in mempool
    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
    };

    // collect token data and filter out tokens if `filter_token` is set
    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
    };

    // finally; with all filtering done; fetch remaining utxo information
    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(())
}

/// Use the unspent index to fetch utxos matching filter at chaintip.
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(())
}