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"))
}
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(¤t)? else {
break;
};
if blockstore.get_hash_by_height(entry.height)?.as_ref() == Some(¤t) {
break;
}
blockstore.store_height(entry.height, ¤t)?;
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)
}
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");
}
}
#[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,
¤t_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)
}