blvm-node 0.1.64

Bitcoin Commons BLVM: Minimal Bitcoin node implementation using blvm-protocol and blvm-consensus
//! Chain reorganization: disconnect/connect via `blvm-consensus` and persist results.

use crate::node::chain_selector::should_activate_over_active_tip;
use crate::node::event_publisher::EventPublisher;
use crate::storage::Storage;
use crate::storage::blockstore::BlockStore;
use crate::utils::current_timestamp;
use anyhow::{Context, Result};
use blvm_consensus::reorganization::{BlockUndoLog, reorganize_chain_with_witnesses};
use blvm_consensus::types::Network;
use blvm_protocol::segwit::Witness;
use blvm_protocol::{BitcoinProtocolEngine, Block, Hash, UtxoSet};
use std::sync::Arc;
use tracing::{info, warn};

fn protocol_network(protocol: &BitcoinProtocolEngine) -> Network {
    protocol.get_protocol_version().consensus_network()
}

fn load_reorg_block(blockstore: &BlockStore, hash: &Hash) -> Result<Block> {
    blockstore
        .get_block(hash)?
        .with_context(|| format!("missing block {hash:?} for reorg"))
}

/// Blocks from the common ancestor through each tip, oldest first.
///
/// The two slices start on the same block. Walking a fixed number of parents
/// from each tip separately shifts the windows when the tips differ in height,
/// and the ancestor search then reports that the chains do not meet.
pub fn collect_fork_chains(
    storage: &Storage,
    blockstore: &BlockStore,
    active_tip: &Hash,
    candidate_tip: &Hash,
) -> Result<(Vec<Block>, Vec<Block>)> {
    let index = storage.chain().block_index();
    let active_entry = index
        .get(active_tip)?
        .context("active tip missing from block index")?;
    let candidate_entry = index
        .get(candidate_tip)?
        .context("candidate tip missing from block index")?;

    let mut active_hash = *active_tip;
    let mut candidate_hash = *candidate_tip;
    let mut active_height = active_entry.height;
    let mut candidate_height = candidate_entry.height;
    let mut active_above = Vec::new();
    let mut candidate_above = Vec::new();

    while active_height > candidate_height {
        active_above.push(load_reorg_block(blockstore, &active_hash)?);
        let entry = index
            .get(&active_hash)?
            .with_context(|| format!("missing index entry {active_hash:?}"))?;
        active_hash = entry.prev_hash;
        active_height -= 1;
    }
    while candidate_height > active_height {
        candidate_above.push(load_reorg_block(blockstore, &candidate_hash)?);
        let entry = index
            .get(&candidate_hash)?
            .with_context(|| format!("missing index entry {candidate_hash:?}"))?;
        candidate_hash = entry.prev_hash;
        candidate_height -= 1;
    }

    let mut active_fork = Vec::new();
    let mut candidate_fork = Vec::new();
    loop {
        if active_hash == candidate_hash {
            let ancestor = load_reorg_block(blockstore, &active_hash)?;
            active_fork.push(ancestor.clone());
            candidate_fork.push(ancestor);
            break;
        }
        active_fork.push(load_reorg_block(blockstore, &active_hash)?);
        candidate_fork.push(load_reorg_block(blockstore, &candidate_hash)?);
        if active_height == 0 {
            anyhow::bail!("chains do not share a common ancestor");
        }
        let active_entry = index
            .get(&active_hash)?
            .with_context(|| format!("missing index entry {active_hash:?}"))?;
        let candidate_entry = index
            .get(&candidate_hash)?
            .with_context(|| format!("missing index entry {candidate_hash:?}"))?;
        active_hash = active_entry.prev_hash;
        candidate_hash = candidate_entry.prev_hash;
        active_height -= 1;
        candidate_height -= 1;
    }

    active_fork.reverse();
    candidate_fork.reverse();
    active_above.reverse();
    candidate_above.reverse();
    active_fork.extend(active_above);
    candidate_fork.extend(candidate_above);
    Ok((active_fork, candidate_fork))
}

fn witnesses_for_chain(
    blockstore: &BlockStore,
    storage: &Storage,
    chain: &[Block],
    protocol: &BitcoinProtocolEngine,
) -> Result<Vec<Vec<Vec<Witness>>>> {
    let protocol_version = protocol.get_protocol_version();
    let mut out = Vec::with_capacity(chain.len());
    for block in chain {
        let hash = blockstore.get_block_hash(block);
        let height = storage
            .chain()
            .block_index()
            .get(&hash)?
            .map(|e| e.height)
            .unwrap_or(0);
        let witnesses = crate::node::block_processor::load_witnesses_for_block(
            blockstore,
            block,
            height,
            protocol_version,
        )?;
        out.push(witnesses);
    }
    Ok(out)
}

fn refresh_active_height_index(
    storage: &Storage,
    blockstore: &BlockStore,
    tip: &Hash,
) -> Result<()> {
    let index = storage.chain().block_index();
    let mut current = *tip;
    let mut seen = std::collections::HashSet::new();
    loop {
        if !seen.insert(current) {
            break;
        }
        let Some(entry) = index.get(&current)? else {
            break;
        };
        if blockstore.get_hash_by_height(entry.height)?.as_ref() == Some(&current) {
            break;
        }
        blockstore.store_height(entry.height, &current)?;
        if entry.height == 0 {
            break;
        }
        current = entry.prev_hash;
    }
    Ok(())
}

fn block_height_in_index(storage: &Storage, blockstore: &BlockStore, block: &Block) -> u64 {
    let hash = blockstore.get_block_hash(block);
    storage
        .chain()
        .block_index()
        .get(&hash)
        .ok()
        .flatten()
        .map(|e| e.height)
        .unwrap_or(0)
}

/// Publish disconnect notifications (tip-first) and a single `ChainReorg` when a runtime is available.
pub fn publish_reorg_events(
    publisher: Option<&Arc<EventPublisher>>,
    storage: &Storage,
    blockstore: &BlockStore,
    old_tip: &Hash,
    new_tip: &Hash,
    disconnected_blocks: &[Block],
) {
    let Some(publisher) = publisher else {
        return;
    };

    let mut disconnects: Vec<(Hash, u64)> = Vec::with_capacity(disconnected_blocks.len());
    for block in disconnected_blocks.iter().rev() {
        let hash = blockstore.get_block_hash(block);
        let height = block_height_in_index(storage, blockstore, block);
        disconnects.push((hash, height));
    }
    let old_tip = *old_tip;
    let new_tip = *new_tip;
    let publisher = Arc::clone(publisher);

    let publish = async move {
        for (hash, height) in disconnects {
            publisher.publish_block_disconnected(&hash, height).await;
        }
        publisher.publish_chain_reorg(&old_tip, &new_tip).await;
    };

    if let Ok(handle) = tokio::runtime::Handle::try_current() {
        handle.spawn(publish);
    } else if let Ok(rt) = tokio::runtime::Runtime::new() {
        rt.block_on(publish);
    } else {
        warn!("No tokio runtime; skipping reorg event publish");
    }
}

/// Switch to `candidate_tip` when it carries more work than the active tip.
#[cfg(feature = "production")]
pub fn try_activate_heavier_fork(
    storage: &Storage,
    blockstore: &BlockStore,
    protocol: &BitcoinProtocolEngine,
    candidate_tip: &Hash,
    utxo_set: &mut UtxoSet,
    event_publisher: Option<&Arc<EventPublisher>>,
) -> Result<bool> {
    if !should_activate_over_active_tip(storage, candidate_tip)? {
        return Ok(false);
    }

    let (active_tip, active_height) = storage.chain().get_tip_hash_and_height()?;
    let (current_chain, new_chain) =
        collect_fork_chains(storage, blockstore, &active_tip, candidate_tip)?;
    if current_chain.is_empty() || new_chain.is_empty() {
        return Ok(false);
    }

    let mut new_chain = new_chain;
    let mut new_witnesses = witnesses_for_chain(blockstore, storage, &new_chain, protocol)?;
    for (block, witnesses) in new_chain.iter_mut().zip(new_witnesses.iter_mut()) {
        let hash = blockstore.get_block_hash(block);
        let height = storage
            .chain()
            .block_index()
            .get(&hash)
            .ok()
            .flatten()
            .map(|e| e.height)
            .unwrap_or(0);
        let (restored, w) = crate::module::pipeline::try_rehydrate_block_for_consensus(
            height,
            hash,
            block.clone(),
            std::mem::take(witnesses),
        );
        *block = restored;
        *witnesses = w;
    }
    let network = protocol_network(protocol);
    let utxo_backup = utxo_set.clone();
    let owned_utxo = std::mem::take(utxo_set);

    let blockstore_ref = blockstore;
    let get_undo = |hash: &Hash| blockstore_ref.get_undo_log(hash).ok().flatten();
    let put_undo = |hash: &Hash, log: &BlockUndoLog| {
        blockstore_ref.store_undo_log(hash, log).map_err(|e| {
            blvm_consensus::error::ConsensusError::BlockValidation(e.to_string().into())
        })
    };

    let mtp_store = blockstore.clone();
    let difficulty_store = blockstore.clone();
    let mut connect_context =
        move |_height: u64,
              recent_headers: Option<&[blvm_consensus::types::BlockHeader]>,
              network_time: u64,
              net: Network|
              -> blvm_consensus::block::BlockValidationContext {
            let mut ctx = blvm_consensus::block::block_validation_context_for_connect_ibd(
                recent_headers,
                network_time,
                net,
            );
            ctx.sequence_prev_mtp = Some(mtp_store.sequence_prev_mtp_lookup());
            ctx.difficulty_ancestor = Some(difficulty_store.difficulty_ancestor_lookup());
            ctx
        };

    let result = reorganize_chain_with_witnesses(
        &new_chain,
        &new_witnesses,
        None,
        &current_chain,
        owned_utxo,
        active_height,
        None::<fn(&Block) -> Option<Vec<Witness>>>,
        None::<fn(u64) -> Option<Vec<blvm_protocol::BlockHeader>>>,
        Some(get_undo),
        Some(put_undo),
        current_timestamp(),
        network,
        Some(&mut connect_context),
    )
    .map_err(|e| {
        *utxo_set = utxo_backup;
        anyhow::anyhow!("reorganize_chain_with_witnesses: {e}")
    })?;

    *utxo_set = result.new_utxo_set;

    let new_tip_block = result
        .connected_blocks
        .last()
        .or(new_chain.last())
        .context("reorg produced no connected blocks")?;
    let new_tip_hash = blockstore.get_block_hash(new_tip_block);
    let new_height = result.new_height;

    storage
        .chain()
        .update_tip(&new_tip_hash, &new_tip_block.header, new_height)?;
    refresh_active_height_index(storage, blockstore, &new_tip_hash)?;
    storage.utxos().store_utxo_set(utxo_set)?;

    publish_reorg_events(
        event_publisher,
        storage,
        blockstore,
        &active_tip,
        &new_tip_hash,
        &result.disconnected_blocks,
    );

    info!(
        "Chain reorg: depth {}, new tip height {}, hash {:?}",
        result.reorganization_depth, new_height, new_tip_hash
    );
    Ok(true)
}