use crate::{
chaindef::{Block, BlockHash, BlockHeader, OutPointHash, ScriptHash, Transaction},
daemon::Daemon,
encode::compute_outpoint_hash,
index::index_block,
index::last_indexed_block,
indexes::{
headerindex::HeaderRow, outputindex::OutputIndexRow,
outtoscriptindex::OutToScripthashIndex, unspentindex::UnspentIndexRow, DBRow,
},
store::{DBStore, Row, META_CF},
undo::blockundoer::BlockUndoer,
writebatch::{PrefilledSpentData, SpentUpdateBatch, WriteBatch},
};
use anyhow::{Context, Result};
use async_trait::async_trait;
use bitcoin_hashes::hex::ToHex;
use bitcoincash::Txid;
use rayon::iter::{IndexedParallelIterator, IntoParallelRefIterator, ParallelIterator};
use std::{
collections::{HashMap, HashSet},
sync::{Arc, Mutex},
};
pub struct StoreBlockUndoer {
store: Arc<DBStore>,
daemon: Arc<Daemon>,
header_chain_cache: Arc<Mutex<Option<Vec<BlockHeader>>>>,
}
impl StoreBlockUndoer {
pub fn new(store: Arc<DBStore>, daemon: Arc<Daemon>) -> Self {
Self {
store,
daemon,
header_chain_cache: Arc::new(Mutex::new(None)),
}
}
async fn get_header_chain(
&self,
tip: &BlockHash,
blockhash: &BlockHash,
) -> Result<Vec<BlockHeader>> {
{
let mut cache = self.header_chain_cache.lock().unwrap();
if let Some(ref cached_chain) = *cache {
if cached_chain
.iter()
.any(|header| header.block_hash() == *blockhash)
{
return Ok(cached_chain.clone());
}
*cache = None;
}
}
let header_chain = HeaderRow::headers_up_to_tip(self.store.as_ref(), tip)
.await
.context("failed to build header chain for undo")?;
{
let mut cache = self.header_chain_cache.lock().unwrap();
*cache = Some(header_chain.clone());
}
Ok(header_chain)
}
fn collect_created_outputs(block: &Block) -> HashSet<OutPointHash> {
block
.txdata
.par_iter()
.map(|tx| {
#[cfg(bch)]
let txid_or_idem = tx.txid();
#[cfg(nexa)]
let txid_or_idem = tx.txidem();
tx.output
.par_iter()
.enumerate()
.map(|(n, _out)| compute_outpoint_hash(&txid_or_idem, n as u32))
.collect::<Vec<OutPointHash>>()
})
.flatten()
.collect()
}
fn collect_spent_outputs(block: &Block) -> HashSet<OutPointHash> {
block
.txdata
.par_iter()
.map(|tx| {
#[cfg(nexa)]
let ophs = tx
.input
.iter()
.map(|input| input.previous_output.hash)
.collect::<Vec<OutPointHash>>();
#[cfg(bch)]
let ophs = tx
.input
.par_iter()
.map(|input| {
compute_outpoint_hash(
&input.previous_output.txid,
input.previous_output.vout,
)
})
.collect::<Vec<OutPointHash>>();
ophs
})
.flatten()
.collect()
}
async fn get_scripthash_from_outpoints(
&self,
outpoints: &HashSet<OutPointHash>,
blockhash: &BlockHash,
) -> Result<HashMap<OutPointHash, PrefilledSpentData>> {
let header_chain = self.get_header_chain(blockhash, blockhash).await?;
let mut output_rows = HashMap::new();
for outpoint in outpoints {
let output_key = OutputIndexRow::filter_by_outpointhash(outpoint);
let (_, output_value) = self.store.get(OutputIndexRow::CF, output_key).await;
if let Some(value) = output_value {
let row = OutputIndexRow::from_row(&Row {
key: OutputIndexRow::filter_by_outpointhash(outpoint).into_boxed_slice(),
value: value.into_boxed_slice(),
});
output_rows.insert(*outpoint, row);
}
}
let mut blockhash_to_outpoints: HashMap<BlockHash, Vec<(OutPointHash, Txid, u32)>> =
HashMap::new();
for (outpoint, output_row) in &output_rows {
let txid = output_row.txid();
let vout = output_row.index();
let height = output_row.height() as usize;
if let Some(header) = header_chain.get(height) {
let blockhash = header.block_hash();
blockhash_to_outpoints
.entry(blockhash)
.or_default()
.push((*outpoint, txid, vout));
}
}
let mut result = HashMap::new();
info!(
"undo: Fetching undo data from {} blocks for {} outpoints",
blockhash_to_outpoints.len(),
outpoints.len()
);
for (blockhash, outpoints_for_block) in blockhash_to_outpoints {
let block = match self.daemon.getblock(&blockhash).await {
Ok(block) => block,
Err(e) => {
warn!(
"could not get block {} from daemon: {}, skipping {} outpoints",
blockhash.to_hex(),
e,
outpoints_for_block.len()
);
continue;
}
};
let txcache: HashMap<Txid, &Transaction> =
block.txdata.iter().map(|tx| (tx.txid(), tx)).collect();
for (outpoint, txid, vout) in outpoints_for_block {
if let Some(tx) = txcache.get(&txid) {
if let Some(txout) = tx.output.get(vout as usize) {
result.insert(
outpoint,
PrefilledSpentData {
scripthash: ScriptHash::normalized_from_txout(txout),
has_token: txout.has_token(),
},
);
} else {
warn!(
"skipping output {}: transaction {} does not have output index {}",
outpoint.to_hex(),
txid.to_hex(),
vout
);
}
} else {
warn!(
"skipping output {}: could not find transaction {} in block {}",
outpoint.to_hex(),
txid.to_hex(),
blockhash.to_hex()
);
}
}
}
Ok(result)
}
fn recreate_unspent_data(
&self,
outpoints: &HashSet<OutPointHash>,
prefilled_spent_data: &HashMap<OutPointHash, PrefilledSpentData>,
) -> Vec<Row> {
outpoints
.iter()
.filter_map(|outpoint| {
prefilled_spent_data.get(outpoint).map(|data| {
UnspentIndexRow::new(&data.scripthash, outpoint, data.has_token).to_row()
})
})
.collect()
}
fn recreate_out_to_scripthash_data(
outpoints: &HashSet<OutPointHash>,
prefilled_spent_data: &HashMap<OutPointHash, PrefilledSpentData>,
) -> Vec<Row> {
outpoints
.iter()
.filter_map(|outpoint| {
prefilled_spent_data.get(outpoint).map(|data| {
OutToScripthashIndex::new(data.scripthash, *outpoint, data.has_token).to_row()
})
})
.collect()
}
}
#[async_trait]
impl BlockUndoer for StoreBlockUndoer {
async fn undo_block(
&self,
blockhash: &BlockHash,
height: u64,
new_tip: &BlockHash,
) -> Result<()> {
info!(
"undo: Undoing block {} (height {})",
blockhash.to_hex(),
height
);
let block = self.daemon.getblock(blockhash).await?;
let write = WriteBatch::new();
let mut spentupdatebatch = SpentUpdateBatch::new().set_erase_only();
let recreate_unspent = WriteBatch::new();
let spent_outputs = Self::collect_spent_outputs(&block);
let prefilled_spent_data = self
.get_scripthash_from_outpoints(&spent_outputs, blockhash)
.await?;
spentupdatebatch = spentupdatebatch.add_prefilled_spent_data(prefilled_spent_data.clone());
index_block(&write, &spentupdatebatch, &block, height as usize);
if height != 0 {
let created = Self::collect_created_outputs(&block);
let spent = spent_outputs;
let old_spent: HashSet<OutPointHash> = spent.difference(&created).copied().collect();
let recreated_unspent_rows =
self.recreate_unspent_data(&old_spent, &prefilled_spent_data);
recreate_unspent.insert(UnspentIndexRow::CF, recreated_unspent_rows);
let recreated_scripthash_rows =
Self::recreate_out_to_scripthash_data(&old_spent, &prefilled_spent_data);
recreate_unspent.insert(OutToScripthashIndex::CF, recreated_scripthash_rows);
};
self.store.update_spends(&spentupdatebatch);
self.store.erase_batch(&write);
self.store.write_batch(&recreate_unspent);
let batch = WriteBatch::new();
batch.insert(META_CF, rayon::iter::once(last_indexed_block(new_tip)));
self.store.write_batch(&batch);
Ok(())
}
}