use super::latch_env;
use super::types::{SharedBlock, SharedWitnesses};
use crate::storage::blockstore::BlockStore;
use anyhow::Result;
use blvm_protocol::features::FeatureRegistry;
use blvm_protocol::types::ARC_BLOCK_CREATED;
use blvm_protocol::{Block, Hash, ProtocolVersion, segwit::Witness};
use std::collections::BTreeMap;
use std::fmt;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, OnceLock};
use tracing::{debug, info, warn};
static FLUSH_PATH_BODY_TIP: AtomicU64 = AtomicU64::new(0);
static IBD_PRUNE_HORIZON: AtomicU64 = AtomicU64::new(0);
static GAP_PERSIST_PRUNE_SKIP_LOGGED: AtomicBool = AtomicBool::new(false);
pub fn cached_feature_registry(protocol_version: ProtocolVersion) -> &'static FeatureRegistry {
static MAINNET: OnceLock<FeatureRegistry> = OnceLock::new();
static TESTNET3: OnceLock<FeatureRegistry> = OnceLock::new();
static REGTEST: OnceLock<FeatureRegistry> = OnceLock::new();
static SIGNET: OnceLock<FeatureRegistry> = OnceLock::new();
match protocol_version {
ProtocolVersion::BitcoinV1 => {
MAINNET.get_or_init(|| FeatureRegistry::for_protocol(ProtocolVersion::BitcoinV1))
}
ProtocolVersion::Testnet3 => {
TESTNET3.get_or_init(|| FeatureRegistry::for_protocol(ProtocolVersion::Testnet3))
}
ProtocolVersion::Regtest => {
REGTEST.get_or_init(|| FeatureRegistry::for_protocol(ProtocolVersion::Regtest))
}
ProtocolVersion::Signet => {
SIGNET.get_or_init(|| FeatureRegistry::for_protocol(ProtocolVersion::Signet))
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum LocalBlockMiss {
NotInStore,
HashMismatch { computed: Hash },
WitnessMissing,
WitnessEmptyStale,
HeightHashUnavailable,
}
impl fmt::Display for LocalBlockMiss {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::NotInStore => write!(f, "block body not in store"),
Self::HashMismatch { computed } => {
write!(
f,
"block hash mismatch (computed {})",
hex::encode(computed)
)
}
Self::WitnessMissing => write!(f, "witness missing for segwit height"),
Self::WitnessEmptyStale => write!(f, "witness blob all-empty (stale MSG_BLOCK fetch)"),
Self::HeightHashUnavailable => write!(f, "height→hash index missing"),
}
}
}
pub fn probe_confirmed_body_height(blockstore: &BlockStore) -> Result<u64> {
let header_max = blockstore.highest_stored_height()?.unwrap_or(0);
if header_max == 0 {
return Ok(0);
}
let has_body = |h: u64| -> Result<bool> {
if h == 0 {
return Ok(false);
}
match blockstore.get_hash_by_height(h)? {
Some(hash) => Ok(blockstore.get_block(&hash)?.is_some()),
None => Ok(false),
}
};
if !has_body(1)? {
return Ok(0);
}
let mut lo = 1u64;
let mut hi = header_max;
while lo < hi {
let mid = lo + (hi - lo).div_ceil(2);
if has_body(mid)? {
lo = mid;
} else {
hi = mid - 1;
}
}
Ok(lo)
}
pub fn body_warehouse_enabled() -> bool {
latch_env!(bool, {
!matches!(
std::env::var("BLVM_IBD_BODY_WAREHOUSE").as_deref(),
Ok("0") | Ok("false") | Ok("FALSE")
)
})
}
pub fn flush_path_body_tip() -> u64 {
FLUSH_PATH_BODY_TIP.load(Ordering::Relaxed)
}
fn flush_path_height_has_body(blockstore: &BlockStore, height: u64) -> Result<bool> {
let Some(hash) = blockstore.get_hash_by_height(height)? else {
return Ok(false);
};
blockstore.has_block_body(&hash)
}
fn bump_flush_path_if_next(height: u64) {
let mut tip = FLUSH_PATH_BODY_TIP.load(Ordering::Relaxed);
if height <= tip || height != tip.saturating_add(1) {
return;
}
loop {
match FLUSH_PATH_BODY_TIP.compare_exchange_weak(
tip,
height,
Ordering::Release,
Ordering::Relaxed,
) {
Ok(_) => return,
Err(actual) => {
tip = actual;
if height <= tip || height != tip.saturating_add(1) {
return;
}
}
}
}
}
fn catch_up_flush_path_on_disk(blockstore: &BlockStore) {
let mut tip = FLUSH_PATH_BODY_TIP.load(Ordering::Relaxed);
loop {
let next = tip.saturating_add(1);
match flush_path_height_has_body(blockstore, next) {
Ok(true) => {}
_ => break,
}
match FLUSH_PATH_BODY_TIP.compare_exchange(tip, next, Ordering::Release, Ordering::Relaxed)
{
Ok(_) => tip = next,
Err(actual) => {
if actual > tip {
tip = actual;
} else {
break;
}
}
}
}
}
fn publish_flush_path_tip() {
let tip = FLUSH_PATH_BODY_TIP.load(Ordering::Relaxed);
if tip > 0 {
super::tip_stage::publish_wan_body_tip(tip);
}
}
pub fn note_available_body(height: u64) {
if !body_warehouse_enabled() || height == 0 {
return;
}
bump_flush_path_if_next(height);
publish_flush_path_tip();
}
pub fn note_flush_path_bodies(blockstore: &BlockStore, heights: &[u64]) {
if !body_warehouse_enabled() {
return;
}
for &height in heights {
if height == 0 {
continue;
}
bump_flush_path_if_next(height);
}
catch_up_flush_path_on_disk(blockstore);
publish_flush_path_tip();
}
#[cfg(test)]
pub(crate) fn reset_flush_path_body_tip_for_test() {
FLUSH_PATH_BODY_TIP.store(0, Ordering::Relaxed);
}
pub fn extend_contiguous_body_tip(
blockstore: &BlockStore,
from: u64,
max_steps: u64,
) -> Result<u64> {
if from == 0 || max_steps == 0 {
return Ok(from);
}
let mut tip = from;
let limit = max_steps.min(4096);
for _ in 0..limit {
let next = tip.saturating_add(1);
let Some(hash) = blockstore.get_hash_by_height(next)? else {
break;
};
if !blockstore.has_block_body(&hash)? {
break;
}
tip = next;
}
Ok(tip)
}
pub fn first_stored_body_height(blockstore: &BlockStore) -> Result<u64> {
let tree = blockstore.height_tree()?;
for item in tree.iter() {
let (key, _) = item?;
if key.len() != 8 {
continue;
}
let h = u64::from_be_bytes(key[..].try_into().unwrap_or([0u8; 8]));
if h == 0 {
continue;
}
let Some(hash) = blockstore.get_hash_by_height(h)? else {
continue;
};
if blockstore.has_block_body(&hash)? {
return Ok(h);
}
}
Ok(0)
}
pub fn contiguous_body_range_tip(blockstore: &BlockStore) -> Result<u64> {
let start = first_stored_body_height(blockstore)?;
if start == 0 {
return Ok(0);
}
contiguous_body_tip_from(blockstore, start)
}
pub fn contiguous_body_tip_from(blockstore: &BlockStore, from: u64) -> Result<u64> {
if from == 0 {
return Ok(0);
}
let Some(hash) = blockstore.get_hash_by_height(from)? else {
return Ok(from.saturating_sub(1));
};
if !blockstore.has_block_body(&hash)? {
return Ok(from.saturating_sub(1));
}
let mut tip = from;
loop {
let next = tip.saturating_add(1);
let Some(hash) = blockstore.get_hash_by_height(next)? else {
break;
};
if !blockstore.has_block_body(&hash)? {
break;
}
tip = next;
}
Ok(tip)
}
pub fn probe_highest_stored_body_height(blockstore: &BlockStore) -> Result<u64> {
let header_max = blockstore.highest_stored_height()?.unwrap_or(0);
if header_max == 0 {
return Ok(0);
}
for h in (1..=header_max).rev() {
let Some(hash) = blockstore.get_hash_by_height(h)? else {
continue;
};
if blockstore.has_block_body(&hash)? {
return Ok(h);
}
}
Ok(0)
}
pub fn pruned_midchain_skip_horizon(header_tip: u64, prune_window: u64) -> u64 {
if prune_window == 0 || header_tip <= prune_window {
return 0;
}
header_tip.saturating_sub(prune_window)
}
pub fn incremental_ibd_prune_window(
incremental_during_ibd: bool,
pruning_enabled: bool,
prune_window_size: u64,
min_blocks_to_keep: u64,
) -> u64 {
if !incremental_during_ibd || !pruning_enabled {
return 0;
}
let window = prune_window_size.max(min_blocks_to_keep);
if window == 0 { 0 } else { window }
}
pub fn ibd_prune_horizon(header_tip: u64, end_height: u64, window: u64) -> u64 {
if window == 0 {
return 0;
}
pruned_midchain_skip_horizon(header_tip.max(end_height), window)
}
pub fn ibd_prune_gc_height(validated: u64, window: u64, min_height: u64) -> u64 {
if window == 0 || validated < min_height {
return 0;
}
validated.saturating_sub(window)
}
pub fn should_skip_block_store_write(
blockstore: &BlockStore,
height: u64,
block_hash: &Hash,
local_replay_max_height: u64,
prune_horizon: u64,
) -> Result<bool> {
if height > 0 && height <= local_replay_max_height {
return Ok(true);
}
if prune_horizon > 0 && height < prune_horizon {
return Ok(true);
}
blockstore.has_block_body(block_hash)
}
pub fn publish_ibd_prune_horizon(horizon: u64) {
IBD_PRUNE_HORIZON.store(horizon, Ordering::Relaxed);
}
pub fn published_prune_horizon() -> u64 {
IBD_PRUNE_HORIZON.load(Ordering::Relaxed)
}
pub(crate) fn height_below_prune_horizon(height: u64, horizon: u64) -> bool {
horizon > 0 && height < horizon
}
pub fn gc_pruned_window_gap_persist(blockstore: &BlockStore, gc_height: u64) -> Result<bool> {
if gc_height == 0 {
return Ok(false);
}
let Some(hash) = blockstore.get_hash_by_height(gc_height)? else {
return Ok(false);
};
if !blockstore.has_block_body(&hash)? {
return Ok(false);
}
blockstore.remove_block_body(&hash)?;
let _ = blockstore.remove_witness(&hash);
Ok(true)
}
pub fn has_real_witnesses(w: &[Vec<Witness>]) -> bool {
w.iter().any(|tx_w| tx_w.iter().any(|s| !s.is_empty()))
}
pub fn coinbase_has_witness_commitment(block: &Block) -> bool {
const MAGIC: [u8; 4] = [0xaa, 0x21, 0xa9, 0xed];
let Some(coinbase) = block.transactions.first() else {
return false;
};
coinbase.outputs.iter().any(|o| {
let s = o.script_pubkey.as_slice();
s.len() >= 38 && s[0] == 0x6a && s[1] == 0x24 && s[2..6] == MAGIC
})
}
pub fn empty_witness_unacceptable(
block: &Block,
witnesses: &[Vec<Witness>],
segwit_on: bool,
) -> bool {
segwit_on && !has_real_witnesses(witnesses) && coinbase_has_witness_commitment(block)
}
pub fn empty_witness_stacks_for_block(block: &Block) -> Vec<Vec<Witness>> {
vec![Vec::new(); block.transactions.len()]
}
pub fn try_repair_missing_witness(
blockstore: &BlockStore,
height: u64,
block_hash: Hash,
witnesses: &[Vec<Witness>],
protocol_version: ProtocolVersion,
block_in_memory: Option<&Block>,
) -> Result<bool> {
if !has_real_witnesses(witnesses) {
return Ok(false);
}
if !blockstore.has_witness_blob(&block_hash)? {
if block_in_memory.is_none() && blockstore.get_block(&block_hash)?.is_none() {
return Ok(false);
}
} else if let Some(w) = blockstore.get_witness(&block_hash)? {
if has_real_witnesses(&w) {
return Ok(false);
}
}
let header_ts = match block_in_memory {
Some(b) => b.header.timestamp,
None => {
blockstore
.get_block(&block_hash)?
.ok_or_else(|| {
anyhow::anyhow!("block disappeared during witness repair at {height}")
})?
.header
.timestamp
}
};
let registry = cached_feature_registry(protocol_version);
if !registry.is_feature_active("segwit", height, header_ts) {
return Ok(false);
}
blockstore.store_witness_at_height(&block_hash, height, witnesses)?;
debug!(
"[IBD_WITNESS_REPAIR] stored missing witness for height {}",
height
);
Ok(true)
}
fn gap_persist_lookahead() -> u64 {
latch_env!(u64, {
std::env::var("BLVM_IBD_GAP_PERSIST_LOOKAHEAD")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(128)
.clamp(1, 512)
})
}
fn should_persist_gap_height(blockstore: &BlockStore, val_h: u64, height: u64) -> bool {
if height == 0 {
return false;
}
if height > val_h {
return true;
}
if height == 1 {
return true;
}
flush_path_height_has_body(blockstore, height.saturating_sub(1)).unwrap_or(false)
}
pub fn try_persist_gap_block_for_local_inject(
blockstore: &BlockStore,
validation_height: Option<&std::sync::Arc<std::sync::atomic::AtomicU64>>,
height: u64,
block_hash: Hash,
block: &Block,
witnesses: &[Vec<Witness>],
protocol_version: ProtocolVersion,
) -> Result<bool> {
try_persist_gap_block_for_local_inject_with_wire(
blockstore,
validation_height,
height,
block_hash,
block,
witnesses,
protocol_version,
None,
)
}
pub fn try_persist_gap_block_for_local_inject_with_wire(
blockstore: &BlockStore,
validation_height: Option<&std::sync::Arc<std::sync::atomic::AtomicU64>>,
height: u64,
block_hash: Hash,
block: &Block,
witnesses: &[Vec<Witness>],
protocol_version: ProtocolVersion,
wire_payload: Option<&[u8]>,
) -> Result<bool> {
match gap_persist_gate(
blockstore,
validation_height,
height,
block_hash,
block,
witnesses,
protocol_version,
)? {
GapPersistGate::Skip => Ok(false),
GapPersistGate::OnDisk(repaired) => Ok(repaired),
GapPersistGate::Write => gap_persist_write_one(
blockstore,
height,
block_hash,
block,
witnesses,
wire_payload,
),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum GapPersistGate {
Skip,
OnDisk(bool),
Write,
}
pub(crate) fn gap_persist_gate(
blockstore: &BlockStore,
validation_height: Option<&std::sync::Arc<std::sync::atomic::AtomicU64>>,
height: u64,
block_hash: Hash,
block: &Block,
witnesses: &[Vec<Witness>],
protocol_version: ProtocolVersion,
) -> Result<GapPersistGate> {
let Some(vh) = validation_height else {
return Ok(GapPersistGate::Skip);
};
let val_h = vh.load(std::sync::atomic::Ordering::Relaxed);
let horizon = published_prune_horizon();
if height_below_prune_horizon(height, horizon) {
if !GAP_PERSIST_PRUNE_SKIP_LOGGED.swap(true, Ordering::Relaxed) || height % 10_000 == 0 {
info!("[IBD_GAP_PERSIST_PRUNE_SKIP] h={height} horizon={horizon} val={val_h}");
}
return Ok(GapPersistGate::Skip);
}
let lookahead = gap_persist_lookahead();
if !should_persist_gap_height(blockstore, val_h, height) {
return Ok(GapPersistGate::Skip);
}
if height > val_h.saturating_add(lookahead) {
debug!(
"[IBD_GAP_PERSIST] height {} outside lookahead {} (val={}) — persist-ahead",
height, lookahead, val_h
);
}
let registry = cached_feature_registry(protocol_version);
let segwit_on = registry.is_feature_active("segwit", height, block.header.timestamp);
if empty_witness_unacceptable(block, witnesses, segwit_on) {
warn!(
"[IBD_GAP_PERSIST_SKIP] height {} (hash {}): empty witness with commitment — not persisting",
height,
hex::encode(block_hash)
);
return Ok(GapPersistGate::Skip);
}
if blockstore.has_block_body(&block_hash)? {
let repaired = try_repair_missing_witness(
blockstore,
height,
block_hash,
witnesses,
protocol_version,
Some(block),
)?;
return Ok(GapPersistGate::OnDisk(repaired));
}
Ok(GapPersistGate::Write)
}
pub(crate) fn gap_persist_write_one(
blockstore: &BlockStore,
height: u64,
block_hash: Hash,
block: &Block,
witnesses: &[Vec<Witness>],
wire_payload: Option<&[u8]>,
) -> Result<bool> {
if wire_bytes_store_enabled() {
if let Some(payload) = wire_payload.filter(|p| !p.is_empty()) {
blockstore.store_block_wire_bytes(block, height, payload)?;
if has_real_witnesses(witnesses) {
blockstore.store_witness_at_height(&block_hash, height, witnesses)?;
}
debug!(
"[IBD_GAP_PERSIST] stored wire-bytes gap block height {} ({} payload bytes)",
height,
payload.len()
);
return Ok(true);
}
warn!(
"[IBD_GAP_PERSIST] WIRE_BYTES_STORE=1 but no payload at height {} — falling back to bincode",
height
);
}
blockstore.store_block_with_witness(block, witnesses, height)?;
debug!(
"[IBD_GAP_PERSIST] stored gap block height {} for coordinator local inject",
height
);
Ok(true)
}
pub fn is_local_witness_hole(
blockstore: &BlockStore,
height: u64,
block_hash: Hash,
protocol_version: ProtocolVersion,
) -> Result<bool> {
if !blockstore.has_block_body(&block_hash)? {
return Ok(false);
}
match try_load_local_ibd_block_with_reason(blockstore, height, block_hash, protocol_version)? {
Ok(_) => Ok(false),
Err(LocalBlockMiss::WitnessMissing | LocalBlockMiss::WitnessEmptyStale) => Ok(true),
Err(_) => Ok(false),
}
}
pub fn wire_bytes_store_enabled() -> bool {
latch_env!(bool, {
matches!(
std::env::var("BLVM_IBD_WIRE_BYTES_STORE").as_deref(),
Ok("1") | Ok("true") | Ok("TRUE") | Ok("yes") | Ok("YES")
)
})
}
pub fn try_load_local_ibd_block_with_reason(
blockstore: &BlockStore,
height: u64,
expected_hash: Hash,
protocol_version: ProtocolVersion,
) -> Result<Result<(Block, Vec<Vec<Witness>>), LocalBlockMiss>> {
let Some((block, witnesses)) = blockstore.get_block_and_witnesses(&expected_hash)? else {
return Ok(Err(LocalBlockMiss::NotInStore));
};
let computed = blockstore.get_block_hash(&block);
if computed != expected_hash {
return Ok(Err(LocalBlockMiss::HashMismatch { computed }));
}
let registry = cached_feature_registry(protocol_version);
let segwit_on = registry.is_feature_active("segwit", height, block.header.timestamp);
match (witnesses.is_empty(), has_real_witnesses(&witnesses)) {
(_, true) => Ok(Ok((block, witnesses))),
(false, false) if segwit_on && coinbase_has_witness_commitment(&block) => {
Ok(Err(LocalBlockMiss::WitnessEmptyStale))
}
(false, false) => Ok(Ok((block, witnesses))),
(true, _) if !segwit_on => Ok(Ok((block, Vec::new()))),
(true, _) if coinbase_has_witness_commitment(&block) => {
Ok(Err(LocalBlockMiss::WitnessMissing))
}
(true, _) => {
let empty = empty_witness_stacks_for_block(&block);
Ok(Ok((block, empty)))
}
}
}
pub fn store_apply_enabled() -> bool {
latch_env!(bool, {
matches!(
std::env::var("BLVM_IBD_STORE_APPLY").as_deref(),
Ok("1") | Ok("true") | Ok("TRUE") | Ok("on") | Ok("yes")
)
})
}
static STORE_APPLY_N: AtomicU64 = AtomicU64::new(0);
pub fn store_apply_count() -> u64 {
STORE_APPLY_N.load(Ordering::Relaxed)
}
pub fn note_store_apply(height: u64) {
let n = STORE_APPLY_N
.fetch_add(1, Ordering::Relaxed)
.saturating_add(1);
if n <= 3 || n % 256 == 0 {
info!("[IBD_STORE_APPLY] h={} n={}", height, n);
}
}
pub fn leftover_stall_try_load(
blockstore: &BlockStore,
height: u64,
protocol_version: ProtocolVersion,
) -> Result<Option<(SharedBlock, SharedWitnesses)>> {
let Some(expected_hash) = blockstore.get_hash_by_height(height)? else {
return Ok(None);
};
if !blockstore.has_block_body(&expected_hash)? {
return Ok(None);
}
Ok(
try_load_local_ibd_block(blockstore, height, expected_hash, protocol_version)?
.map(|(block, witnesses)| (Arc::new(block), Arc::new(witnesses))),
)
}
pub fn try_load_local_ibd_block(
blockstore: &BlockStore,
height: u64,
expected_hash: Hash,
protocol_version: ProtocolVersion,
) -> Result<Option<(Block, Vec<Vec<Witness>>)>> {
match try_load_local_ibd_block_with_reason(blockstore, height, expected_hash, protocol_version)?
{
Ok(pair) => Ok(Some(pair)),
Err(_) => Ok(None),
}
}
pub fn ibd_stall_abort_gap_fetch_on_confirmed_bodies() -> bool {
matches!(
std::env::var("BLVM_IBD_STALL_ABORT_GAP_FETCH")
.ok()
.as_deref(),
Some("1") | Some("true") | Some("TRUE")
)
}
pub fn ibd_stall_aborts_inflight_gap_fetch(
wan_multi_peer: bool,
confirmed_body_height: u64,
gap_height: u64,
) -> bool {
match std::env::var("BLVM_IBD_STALL_ABORT_GAP_FETCH")
.ok()
.as_deref()
{
Some("1") | Some("true") | Some("TRUE") => true,
Some("0") | Some("false") | Some("FALSE") => false,
_ => {
if super::memory::ibd_pressure_is_critical_or_worse() {
static LAST: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let prev = LAST.load(std::sync::atomic::Ordering::Relaxed);
if now.saturating_sub(prev) >= 15 {
LAST.store(now, std::sync::atomic::Ordering::Relaxed);
super::memory::log_pressure_behavior(
"stall_abort_gap_fetch",
"abort",
"reason=pressure_critical_or_worse",
);
}
return true;
}
if wan_multi_peer {
let _ = (confirmed_body_height, gap_height);
false
} else {
!(confirmed_body_height > 0 && gap_height <= confirmed_body_height)
}
}
}
}
pub fn ibd_local_gap_fill_enabled() -> bool {
!matches!(
std::env::var("BLVM_IBD_LOCAL_GAP_FILL").ok().as_deref(),
Some("0") | Some("false") | Some("FALSE")
)
}
pub fn ibd_local_gap_fill_max_height(_confirmed_body_height_at_start: u64) -> u64 {
std::env::var("BLVM_IBD_LOCAL_GAP_FILL_MAX_HEIGHT")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(u64::MAX)
}
fn height_eligible_for_local_gap_fill(height: u64, confirmed_body_height_at_start: u64) -> bool {
if !ibd_local_gap_fill_enabled() || height == 0 {
return false;
}
height <= ibd_local_gap_fill_max_height(confirmed_body_height_at_start)
}
pub fn gap_inject_lookahead_pub() -> u64 {
gap_inject_lookahead()
}
fn gap_inject_lookahead() -> u64 {
latch_env!(u64, {
std::env::var("BLVM_IBD_GAP_INJECT_LOOKAHEAD")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(128)
.clamp(1, 256)
})
}
fn coordinator_inject_one(
blockstore: &BlockStore,
protocol_version: ProtocolVersion,
height: u64,
confirmed_body_height_at_start: u64,
reorder_buffer: &mut BTreeMap<u64, (SharedBlock, SharedWitnesses)>,
already_dispatched: &rustc_hash::FxHashSet<u64>,
log_miss: &mut rustc_hash::FxHashSet<u64>,
reload_if_dispatched: bool,
) -> Result<bool> {
if reorder_buffer.contains_key(&height) {
return Ok(true); }
if already_dispatched.contains(&height) && !reload_if_dispatched {
return Ok(true);
}
if !height_eligible_for_local_gap_fill(height, confirmed_body_height_at_start) {
return Ok(false);
}
let Some(expected_hash) = blockstore.get_hash_by_height(height)? else {
if log_miss.insert(height) {
warn!(
"[IBD_LOCAL_GAP] height {}: {}",
height,
LocalBlockMiss::HeightHashUnavailable
);
}
return Ok(false);
};
match try_load_local_ibd_block_with_reason(blockstore, height, expected_hash, protocol_version)?
{
Ok((block, witnesses)) => {
ARC_BLOCK_CREATED.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
reorder_buffer.insert(height, (Arc::new(block), Arc::new(witnesses)));
crate::node::parallel_ibd::tip_stage::mark_reorder(height);
log_miss.remove(&height);
debug!("[IBD_LOCAL_GAP] injected local block height {}", height);
Ok(true)
}
Err(miss) => {
if log_miss.insert(height) {
warn!("[IBD_LOCAL_GAP] height {}: {miss}", height);
}
Ok(false)
}
}
}
pub fn coordinator_inject_local_gap(
blockstore: &BlockStore,
protocol_version: ProtocolVersion,
height: u64,
confirmed_body_height_at_start: u64,
validation_height: u64,
reorder_buffer: &mut BTreeMap<u64, (SharedBlock, SharedWitnesses)>,
already_dispatched: &rustc_hash::FxHashSet<u64>,
log_miss: &mut rustc_hash::FxHashSet<u64>,
tip_in_pipeline: bool,
) -> Result<bool> {
if height != validation_height.saturating_add(1) {
return Ok(false);
}
let lookahead = gap_inject_lookahead();
let mut any_success = false;
let mut newly_injected = 0u64;
let mut last_new_h = height;
for i in 0..lookahead {
let h = height + i;
let had_in_reorder = reorder_buffer.contains_key(&h);
let reload_if_dispatched = i == 0 && !tip_in_pipeline;
match coordinator_inject_one(
blockstore,
protocol_version,
h,
confirmed_body_height_at_start,
reorder_buffer,
already_dispatched,
log_miss,
reload_if_dispatched,
)? {
true => {
any_success = true;
if !had_in_reorder && reorder_buffer.contains_key(&h) {
newly_injected += 1;
last_new_h = h;
}
}
false => break, }
}
super::feeder_miss::note_inject(newly_injected, height, tip_in_pipeline);
if newly_injected > 1 {
info!(
"[IBD_INJECT_CHAIN] from {} chained {} height(s) (lookahead={}, stopped_at={})",
height,
newly_injected,
lookahead,
last_new_h.saturating_add(1)
);
}
Ok(any_success)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::storage::blockstore::BlockStore;
use crate::storage::database::{create_database, default_backend};
use blvm_protocol::{
Block, BlockHeader, OutPoint, Transaction, TransactionInput, TransactionOutput,
};
use std::sync::Arc;
use std::sync::atomic::AtomicU64;
use tempfile::TempDir;
fn temp_blockstore() -> BlockStore {
let dir = TempDir::new().unwrap();
let db: Arc<dyn crate::storage::database::Database> =
Arc::from(create_database(dir.path(), default_backend(), None).unwrap());
std::mem::forget(dir);
BlockStore::new(db).unwrap()
}
#[test]
fn note_store_apply_counts() {
let before = store_apply_count();
note_store_apply(7);
assert!(store_apply_count() > before);
}
#[test]
fn leftover_stall_try_load_hits_stored_body() {
let blockstore = temp_blockstore();
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(1, &hash).unwrap();
blockstore.store_block_with_witness(&block, &[], 1).unwrap();
let loaded = leftover_stall_try_load(&blockstore, 1, ProtocolVersion::BitcoinV1)
.unwrap()
.expect("body on disk");
assert_eq!(blockstore.get_block_hash(loaded.0.as_ref()), hash);
assert!(
leftover_stall_try_load(&blockstore, 2, ProtocolVersion::BitcoinV1)
.unwrap()
.is_none()
);
}
#[test]
fn local_gap_fill_allowed_without_contiguous_body_probe() {
assert!(height_eligible_for_local_gap_fill(338_305, 0));
assert!(!height_eligible_for_local_gap_fill(0, 0));
}
#[test]
fn incremental_ibd_prune_window_is_zero_unless_enabled() {
assert_eq!(incremental_ibd_prune_window(false, true, 144, 144), 0);
assert_eq!(incremental_ibd_prune_window(true, false, 144, 144), 0);
assert_eq!(incremental_ibd_prune_window(true, true, 0, 0), 0);
assert_eq!(
incremental_ibd_prune_window(true, true, 144, 50_000),
50_000
);
assert_eq!(incremental_ibd_prune_window(true, true, 144, 144), 144);
}
#[test]
fn ibd_prune_horizon_and_gc_follow_the_keep_window() {
assert_eq!(ibd_prune_horizon(961_637, 250_000, 144), 961_493);
assert_eq!(ibd_prune_horizon(100, 250_000, 144), 249_856);
assert_eq!(ibd_prune_horizon(961_637, 250_000, 0), 0);
assert_eq!(ibd_prune_gc_height(287, 144, 288), 0);
assert_eq!(ibd_prune_gc_height(288, 144, 288), 144);
assert_eq!(ibd_prune_gc_height(50_000, 50_000, 288), 0);
assert_eq!(ibd_prune_gc_height(50_144, 144, 288), 50_000);
}
#[test]
fn pruned_midchain_horizon_skips_below_window() {
assert_eq!(pruned_midchain_skip_horizon(961_637, 144), 961_493);
assert_eq!(pruned_midchain_skip_horizon(223_000, 144), 222_856);
assert_eq!(pruned_midchain_skip_horizon(100, 144), 0);
assert_eq!(pruned_midchain_skip_horizon(0, 144), 0);
assert_eq!(pruned_midchain_skip_horizon(961_637, 0), 0);
}
#[test]
fn gap_persist_prune_skip_matches_flush_horizon() {
assert!(height_below_prune_horizon(223_000, 961_493));
assert!(!height_below_prune_horizon(961_493, 961_493));
assert!(!height_below_prune_horizon(961_494, 961_493));
assert!(!height_below_prune_horizon(223_000, 0));
}
fn should_skip_pruned_midchain_without_body_on_disk() {
let blockstore = temp_blockstore();
let missing = [0xCCu8; 32];
assert!(should_skip_block_store_write(&blockstore, 223_000, &missing, 0, 961_493).unwrap());
assert!(
!should_skip_block_store_write(&blockstore, 961_493, &missing, 0, 961_493).unwrap()
);
assert!(!should_skip_block_store_write(&blockstore, 223_000, &missing, 0, 0).unwrap());
}
#[test]
fn gc_pruned_window_gap_persist_drops_body_keeps_header() {
let blockstore = temp_blockstore();
let height = 200u64;
let block = Block {
header: BlockHeader {
version: 1,
timestamp: 1_234_567,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(height, &hash).unwrap();
blockstore
.store_block_with_witness(&block, &[vec![]], height)
.unwrap();
assert!(blockstore.has_block_body(&hash).unwrap());
assert!(gc_pruned_window_gap_persist(&blockstore, height).unwrap());
assert!(!blockstore.has_block_body(&hash).unwrap());
assert_eq!(blockstore.get_hash_by_height(height).unwrap(), Some(hash));
assert!(!gc_pruned_window_gap_persist(&blockstore, height).unwrap());
assert!(!gc_pruned_window_gap_persist(&blockstore, 0).unwrap());
}
#[test]
fn w5_wire_bytes_persist_and_inject_coexist_with_bincode() {
use crate::storage::blockstore::{
decode_wire_body_blob, encode_wire_body_blob, is_wire_body_blob,
};
use blvm_protocol::serialization::{
deserialize_block_with_witnesses, serialize_block_with_witnesses,
};
let blockstore = temp_blockstore();
let height = 500u64;
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let witnesses: Vec<Vec<Witness>> = vec![vec![]];
let payload = serialize_block_with_witnesses(&block, &witnesses, true);
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(height, &hash).unwrap();
blockstore
.store_block_wire_bytes(&block, height, &payload)
.unwrap();
let blob = blockstore
.get_block_and_witnesses(&hash)
.unwrap()
.expect("wire body present");
assert_eq!(blockstore.get_block_hash(&blob.0), hash);
assert!(is_wire_body_blob(&encode_wire_body_blob(&payload)));
assert_eq!(
decode_wire_body_blob(&encode_wire_body_blob(&payload)).unwrap(),
payload.as_slice()
);
let loaded =
try_load_local_ibd_block(&blockstore, height, hash, ProtocolVersion::BitcoinV1)
.unwrap()
.expect("inject load");
assert_eq!(blockstore.get_block_hash(&loaded.0), hash);
let h2 = 501u64;
let hash2 = {
let mut b2 = block.clone();
b2.header.nonce = 99;
let h = blockstore.get_block_hash(&b2);
blockstore.store_height(h2, &h).unwrap();
blockstore
.store_block_with_witness(&b2, &witnesses, h2)
.unwrap();
h
};
let legacy = try_load_local_ibd_block(&blockstore, h2, hash2, ProtocolVersion::BitcoinV1)
.unwrap()
.expect("bincode inject");
assert_eq!(blockstore.get_block_hash(&legacy.0), hash2);
let t0 = std::time::Instant::now();
for _ in 0..200 {
let body = bincode::serialize(&block).unwrap();
let wit = bincode::serialize(&witnesses).unwrap();
let _: Block = bincode::deserialize(&body).unwrap();
let _: Vec<Vec<Witness>> = bincode::deserialize(&wit).unwrap();
}
let bincode_ns = t0.elapsed().as_nanos() / 200;
let t1 = std::time::Instant::now();
for _ in 0..200 {
let tagged = encode_wire_body_blob(&payload);
let p = decode_wire_body_blob(&tagged).unwrap();
let _ = deserialize_block_with_witnesses(p).unwrap();
}
let wire_ns = t1.elapsed().as_nanos() / 200;
eprintln!(
"[W5 micro] bincode ser+de×2 ≈{bincode_ns} ns/op; wire memcpy+deser ≈{wire_ns} ns/op"
);
let _ = (bincode_ns, wire_ns);
}
#[test]
fn sparse_body_probe_finds_max_without_height_one() {
let blockstore = temp_blockstore();
let height = 500u64;
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
let placeholder = [0xAAu8; 32];
for h in 0..height {
blockstore.store_height(h, &placeholder).unwrap();
}
blockstore
.store_block_with_witness(&block, &[], height)
.unwrap();
blockstore.store_height(height, &hash).unwrap();
assert_eq!(probe_confirmed_body_height(&blockstore).unwrap(), 0);
assert_eq!(
probe_highest_stored_body_height(&blockstore).unwrap(),
height
);
assert!(should_skip_block_store_write(&blockstore, height, &hash, 0, 0).unwrap());
assert!(
!should_skip_block_store_write(&blockstore, height + 1, &[0xBBu8; 32], 0, 0).unwrap()
);
}
#[test]
fn contiguous_body_probe_unchanged_when_height_one_present() {
let blockstore = temp_blockstore();
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
blockstore.store_block_with_witness(&block, &[], 1).unwrap();
blockstore.store_height(0, &[0u8; 32]).unwrap();
blockstore
.store_height(1, &blockstore.get_block_hash(&block))
.unwrap();
assert_eq!(probe_confirmed_body_height(&blockstore).unwrap(), 1);
assert_eq!(probe_highest_stored_body_height(&blockstore).unwrap(), 1);
}
#[test]
fn cached_feature_registry_is_process_latched() {
let a = cached_feature_registry(ProtocolVersion::BitcoinV1);
let b = cached_feature_registry(ProtocolVersion::BitcoinV1);
assert!(std::ptr::eq(a, b), "same protocol must reuse one registry");
assert!(a.is_feature_active("segwit", 500_000, 1_600_000_000));
let r = cached_feature_registry(ProtocolVersion::Regtest);
assert!(
!std::ptr::eq(a, r),
"distinct protocols get distinct caches"
);
}
#[test]
fn stall_abort_policy() {
use super::super::memory::{self, PressureLevel};
let prev = std::env::var("BLVM_IBD_STALL_ABORT_GAP_FETCH").ok();
unsafe { std::env::remove_var("BLVM_IBD_STALL_ABORT_GAP_FETCH") };
memory::publish_ibd_pressure(PressureLevel::None);
assert!(!ibd_stall_aborts_inflight_gap_fetch(true, 0, 265_553));
assert!(!ibd_stall_aborts_inflight_gap_fetch(true, 640_000, 640_001));
assert!(!ibd_stall_aborts_inflight_gap_fetch(true, 640_000, 640_000));
assert!(ibd_stall_aborts_inflight_gap_fetch(false, 0, 265_553));
memory::publish_ibd_pressure(PressureLevel::Critical);
assert!(ibd_stall_aborts_inflight_gap_fetch(true, 0, 265_553));
memory::publish_ibd_pressure(PressureLevel::None);
unsafe { std::env::set_var("BLVM_IBD_STALL_ABORT_GAP_FETCH", "1") };
assert!(ibd_stall_aborts_inflight_gap_fetch(true, 0, 265_553));
if let Some(v) = prev {
unsafe { std::env::set_var("BLVM_IBD_STALL_ABORT_GAP_FETCH", v) };
} else {
unsafe { std::env::remove_var("BLVM_IBD_STALL_ABORT_GAP_FETCH") };
}
memory::publish_ibd_pressure(PressureLevel::None);
}
#[test]
fn local_block_miss_display() {
let m = LocalBlockMiss::WitnessEmptyStale;
assert!(m.to_string().contains("empty"));
}
#[test]
fn contiguous_body_tip_from_stops_at_first_hole() {
let blockstore = temp_blockstore();
for h in [1u64, 2, 3, 5] {
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000 + h,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(h, &hash).unwrap();
blockstore.store_block_with_witness(&block, &[], h).unwrap();
}
assert_eq!(
contiguous_body_tip_from(&blockstore, 1).unwrap(),
3,
"must stop before leftover cheese at 5"
);
assert_eq!(
contiguous_body_tip_from(&blockstore, 5).unwrap(),
5,
"walk from a later island stays on that island"
);
}
#[test]
fn contiguous_body_range_tip_does_not_require_genesis() {
let blockstore = temp_blockstore();
for h in [180_001u64, 180_002, 180_003, 180_005] {
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000 + h,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(h, &hash).unwrap();
blockstore.store_block_with_witness(&block, &[], h).unwrap();
}
assert_eq!(first_stored_body_height(&blockstore).unwrap(), 180_001);
assert_eq!(
contiguous_body_range_tip(&blockstore).unwrap(),
180_003,
"range tip is the first island, not 0 and not leftover cheese"
);
assert_eq!(
contiguous_body_tip_from(&blockstore, 1).unwrap(),
0,
"walk-from-1 still misses a slice that does not start at genesis"
);
}
#[test]
fn extend_contiguous_body_tip_stops_at_hole() {
let blockstore = temp_blockstore();
for h in [100u64, 101, 102, 104] {
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000 + h,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(h, &hash).unwrap();
blockstore.store_block_with_witness(&block, &[], h).unwrap();
}
assert_eq!(
extend_contiguous_body_tip(&blockstore, 100, 256).unwrap(),
102,
"must stop before hole at 103"
);
assert_eq!(
extend_contiguous_body_tip(&blockstore, 102, 256).unwrap(),
102,
"no further contiguous from hole edge"
);
}
fn store_body_at(blockstore: &BlockStore, h: u64) {
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000 + h,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(h, &hash).unwrap();
blockstore.store_block_with_witness(&block, &[], h).unwrap();
}
#[test]
fn available_body_bumps_without_disk_or_flush_commit() {
reset_flush_path_body_tip_for_test();
note_available_body(1);
assert_eq!(flush_path_body_tip(), 1, "feeder-take / insert of 1 bumps");
note_available_body(2);
assert_eq!(flush_path_body_tip(), 2, "sequential availability bumps");
note_available_body(4);
assert_eq!(flush_path_body_tip(), 2, "hole at 3 must not jump");
note_available_body(3);
assert_eq!(
flush_path_body_tip(),
3,
"fill 3; 4 is not auto-caught without disk"
);
reset_flush_path_body_tip_for_test();
}
#[test]
fn flush_path_tip_advances_on_write_not_sparse_jump() {
reset_flush_path_body_tip_for_test();
let blockstore = temp_blockstore();
for h in [3u64, 5] {
store_body_at(&blockstore, h);
note_flush_path_bodies(&blockstore, &[h]);
}
assert_eq!(
flush_path_body_tip(),
0,
"out-of-order persist must not jump a hole at 1"
);
store_body_at(&blockstore, 1);
note_flush_path_bodies(&blockstore, &[1]);
assert_eq!(flush_path_body_tip(), 1, "tip+1 persist bumps once");
store_body_at(&blockstore, 2);
note_flush_path_bodies(&blockstore, &[2]);
assert_eq!(
flush_path_body_tip(),
3,
"tip+1 persist catch-up walks already-on-disk 3"
);
store_body_at(&blockstore, 4);
note_flush_path_bodies(&blockstore, &[4]);
assert_eq!(flush_path_body_tip(), 5, "catch-up reaches 5 after 4 lands");
reset_flush_path_body_tip_for_test();
}
#[test]
fn flush_path_tip_batch_1_to_50_holds() {
reset_flush_path_body_tip_for_test();
let blockstore = temp_blockstore();
let heights: Vec<u64> = (1..=50).collect();
for &h in &heights {
store_body_at(&blockstore, h);
}
note_flush_path_bodies(&blockstore, &heights);
assert_eq!(
flush_path_body_tip(),
50,
"sorted flush batch 1..=50 must publish tip=50"
);
reset_flush_path_body_tip_for_test();
}
#[test]
fn persist_gap_block_enables_local_inject() {
let blockstore = temp_blockstore();
let vh = Arc::new(AtomicU64::new(499));
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(500, &hash).unwrap();
assert!(!blockstore.has_block_body(&hash).unwrap());
let persisted = try_persist_gap_block_for_local_inject(
&blockstore,
Some(&vh),
500,
hash,
&block,
&[],
ProtocolVersion::BitcoinV1,
)
.unwrap();
assert!(persisted);
assert!(blockstore.has_block_body(&hash).unwrap());
let mut reorder_buffer = std::collections::BTreeMap::new();
let already_dispatched = rustc_hash::FxHashSet::default();
let mut log_miss = rustc_hash::FxHashSet::default();
assert!(
coordinator_inject_local_gap(
&blockstore,
ProtocolVersion::BitcoinV1,
500,
0,
499,
&mut reorder_buffer,
&already_dispatched,
&mut log_miss,
false,
)
.unwrap()
);
assert!(reorder_buffer.contains_key(&500));
}
fn persist_toy_at(blockstore: &BlockStore, vh: &Arc<AtomicU64>, h: u64) -> bool {
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000 + h,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(h, &hash).unwrap();
try_persist_gap_block_for_local_inject(
blockstore,
Some(vh),
h,
hash,
&block,
&[],
ProtocolVersion::BitcoinV1,
)
.unwrap()
}
fn toy_block_at(h: u64) -> Block {
Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000 + h,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000 + h as i64,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
}
}
#[test]
fn r348_batch_store_round_trips_and_gate_sees_on_disk() {
let bs = temp_blockstore();
let blocks: Vec<Block> = (100u64..132).map(toy_block_at).collect();
for (i, b) in blocks.iter().enumerate() {
bs.store_height(100 + i as u64, &bs.get_block_hash(b))
.unwrap();
}
let empty: Vec<Vec<blvm_protocol::segwit::Witness>> = Vec::new();
let items: Vec<(&Block, &[Vec<blvm_protocol::segwit::Witness>], u64)> = blocks
.iter()
.enumerate()
.map(|(i, b)| (b, empty.as_slice(), 100 + i as u64))
.collect();
bs.store_blocks_with_witness_batch(&items).unwrap();
let vh = Arc::new(AtomicU64::new(99));
for (i, b) in blocks.iter().enumerate() {
let h = 100 + i as u64;
let hash = bs.get_block_hash(b);
assert!(bs.has_block_body(&hash).unwrap(), "body {h} on disk");
let (got, _w) =
try_load_local_ibd_block_with_reason(&bs, h, hash, ProtocolVersion::BitcoinV1)
.unwrap()
.unwrap_or_else(|m| panic!("load {h}: {m:?}"));
assert_eq!(got.header.timestamp, b.header.timestamp);
let gate =
gap_persist_gate(&bs, Some(&vh), h, hash, b, &[], ProtocolVersion::BitcoinV1)
.unwrap();
assert!(
matches!(gate, GapPersistGate::OnDisk(_)),
"gate after batch: {gate:?}"
);
}
let nb = toy_block_at(200);
let nh = bs.get_block_hash(&nb);
bs.store_height(200, &nh).unwrap();
assert_eq!(
gap_persist_gate(
&bs,
Some(&vh),
200,
nh,
&nb,
&[],
ProtocolVersion::BitcoinV1
)
.unwrap(),
GapPersistGate::Write
);
}
fn persist_wrote_body(blockstore: &BlockStore, h: u64) -> bool {
blockstore
.get_hash_by_height(h)
.ok()
.flatten()
.and_then(|hash| blockstore.has_block_body(&hash).ok())
.unwrap_or(false)
}
#[test]
fn r196_persist_ahead_prefix() {
reset_flush_path_body_tip_for_test();
let ahead = temp_blockstore();
let vh0 = Arc::new(AtomicU64::new(0));
for h in 1u64..=5 {
assert!(
persist_toy_at(&ahead, &vh0, h),
"persist {h} while val=0 must write"
);
assert!(persist_wrote_body(&ahead, h), "body {h} on disk");
}
assert_eq!(
flush_path_body_tip(),
0,
"GAP_PERSIST must not publish leftover warehouse tip"
);
reset_flush_path_body_tip_for_test();
let late = temp_blockstore();
let vh600 = Arc::new(AtomicU64::new(600));
assert!(
persist_toy_at(&late, &vh600, 1),
"late tip+1 must persist after apply passed"
);
assert!(persist_wrote_body(&late, 1));
assert_eq!(flush_path_body_tip(), 0);
assert!(
persist_toy_at(&late, &vh600, 2),
"sequential glue continues after apply passed"
);
assert!(persist_wrote_body(&late, 2));
assert_eq!(flush_path_body_tip(), 0);
reset_flush_path_body_tip_for_test();
let farm = temp_blockstore();
let vh10 = Arc::new(AtomicU64::new(10));
assert!(
persist_toy_at(&farm, &vh10, 2000),
"farm tile outside lookahead must persist"
);
assert!(persist_wrote_body(&farm, 2000));
assert_eq!(flush_path_body_tip(), 0);
reset_flush_path_body_tip_for_test();
}
#[test]
fn w1_persist_skips_empty_witness_when_commitment_present() {
let blockstore = temp_blockstore();
let vh = Arc::new(AtomicU64::new(500_000));
let height = 500_001u64;
let mut commitment_script = vec![0x6a, 0x24, 0xaa, 0x21, 0xa9, 0xed];
commitment_script.extend_from_slice(&[0u8; 32]);
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![TransactionInput {
prevout: OutPoint {
hash: [0u8; 32],
index: 0xffff_ffff,
},
script_sig: vec![0x01].into(),
sequence: 0xffff_ffff,
}],
outputs: blvm_protocol::tx_outputs![
TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
},
TransactionOutput {
value: 0,
script_pubkey: commitment_script.into(),
}
],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(height, &hash).unwrap();
let persisted = try_persist_gap_block_for_local_inject(
&blockstore,
Some(&vh),
height,
hash,
&block,
&[], ProtocolVersion::BitcoinV1,
)
.unwrap();
assert!(
!persisted,
"W1 must refuse empty-witness persist when commitment present"
);
assert!(
!blockstore.has_block_body(&hash).unwrap(),
"body must not be written"
);
}
#[test]
fn w1_persist_allows_empty_witness_without_commitment() {
let blockstore = temp_blockstore();
let vh = Arc::new(AtomicU64::new(500_000));
let height = 500_001u64;
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(height, &hash).unwrap();
let persisted = try_persist_gap_block_for_local_inject(
&blockstore,
Some(&vh),
height,
hash,
&block,
&[], ProtocolVersion::BitcoinV1,
)
.unwrap();
assert!(persisted, "empty witness without commitment must persist");
}
#[test]
fn wire_persist_stores_witness_blob_for_commitment_block() {
use blvm_protocol::serialization::serialize_block_with_witnesses;
let prev_wire = std::env::var("BLVM_IBD_WIRE_BYTES_STORE").ok();
unsafe { std::env::set_var("BLVM_IBD_WIRE_BYTES_STORE", "1") };
let blockstore = temp_blockstore();
let vh = Arc::new(AtomicU64::new(500_000));
let height = 500_001u64;
let mut commitment_script = vec![0x6a, 0x24, 0xaa, 0x21, 0xa9, 0xed];
commitment_script.extend_from_slice(&[0u8; 32]);
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![TransactionInput {
prevout: OutPoint {
hash: [0u8; 32],
index: 0xffff_ffff,
},
script_sig: vec![0x01].into(),
sequence: 0xffff_ffff,
}],
outputs: blvm_protocol::tx_outputs![
TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
},
TransactionOutput {
value: 0,
script_pubkey: commitment_script.into(),
}
],
lock_time: 0,
}]
.into(),
};
let witnesses: Vec<Vec<Witness>> = vec![vec![vec![vec![0x01u8, 0x02, 0x03]]]];
let payload = serialize_block_with_witnesses(&block, &witnesses, true);
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(height, &hash).unwrap();
let persisted = try_persist_gap_block_for_local_inject_with_wire(
&blockstore,
Some(&vh),
height,
hash,
&block,
&witnesses,
ProtocolVersion::BitcoinV1,
Some(payload.as_slice()),
)
.unwrap();
assert!(persisted, "wire persist with real witnesses must succeed");
assert!(
blockstore.has_witness_blob(&hash).unwrap(),
"wire path must write witness row"
);
unsafe {
match prev_wire {
Some(v) => std::env::set_var("BLVM_IBD_WIRE_BYTES_STORE", v),
None => std::env::remove_var("BLVM_IBD_WIRE_BYTES_STORE"),
}
}
}
#[test]
fn local_load_synthesizes_empty_witness_without_commitment() {
let blockstore = temp_blockstore();
let height = 500_001u64;
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(height, &hash).unwrap();
blockstore
.store_block_with_witness(&block, &[], height)
.unwrap();
let loaded =
try_load_local_ibd_block(&blockstore, height, hash, ProtocolVersion::BitcoinV1)
.unwrap();
assert!(
loaded.is_some(),
"body without commitment must load with empty witnesses"
);
assert!(
!is_local_witness_hole(&blockstore, height, hash, ProtocolVersion::BitcoinV1).unwrap(),
"no commitment → not a witness hole"
);
}
#[test]
fn is_local_witness_hole_detects_body_without_witness() {
let blockstore = temp_blockstore();
let height = 500_001u64;
let mut commitment_script = vec![0x6a, 0x24, 0xaa, 0x21, 0xa9, 0xed];
commitment_script.extend_from_slice(&[0u8; 32]);
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![TransactionInput {
prevout: OutPoint {
hash: [0u8; 32],
index: 0xffff_ffff,
},
script_sig: vec![0x01].into(),
sequence: 0xffff_ffff,
}],
outputs: blvm_protocol::tx_outputs![
TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
},
TransactionOutput {
value: 0,
script_pubkey: commitment_script.into(),
}
],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(height, &hash).unwrap();
blockstore
.store_block_with_witness(&block, &[], height)
.unwrap();
assert!(
is_local_witness_hole(&blockstore, height, hash, ProtocolVersion::BitcoinV1).unwrap()
);
}
#[test]
fn inject_past_start_confirmed_body_cap_when_uncapped() {
let blockstore = temp_blockstore();
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_400_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(500, &hash).unwrap();
blockstore
.store_block_with_witness(&block, &[], 500)
.unwrap();
let mut reorder_buffer = std::collections::BTreeMap::new();
let already_dispatched = rustc_hash::FxHashSet::default();
let mut log_miss = rustc_hash::FxHashSet::default();
assert!(
coordinator_inject_local_gap(
&blockstore,
ProtocolVersion::BitcoinV1,
500,
400, 499,
&mut reorder_buffer,
&already_dispatched,
&mut log_miss,
false,
)
.unwrap()
);
assert!(reorder_buffer.contains_key(&500));
}
#[test]
fn inject_skips_already_dispatched_heights() {
let blockstore = temp_blockstore();
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(500, &hash).unwrap();
blockstore
.store_block_with_witness(&block, &[], 500)
.unwrap();
let mut reorder_buffer = std::collections::BTreeMap::new();
let mut already_dispatched = rustc_hash::FxHashSet::default();
already_dispatched.insert(500);
let mut log_miss = rustc_hash::FxHashSet::default();
assert!(
coordinator_inject_local_gap(
&blockstore,
ProtocolVersion::BitcoinV1,
500,
0,
499,
&mut reorder_buffer,
&already_dispatched,
&mut log_miss,
true,
)
.unwrap()
);
assert!(
reorder_buffer.is_empty(),
"already-dispatched height must not be re-loaded into reorder_buffer"
);
}
#[test]
fn inject_chain_skips_counting_already_present_heights() {
let blockstore = temp_blockstore();
let mk = |ts| Block {
header: BlockHeader {
version: 4,
timestamp: ts,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let b500 = mk(1_600_000_000);
let h500 = blockstore.get_block_hash(&b500);
blockstore.store_height(500, &h500).unwrap();
blockstore
.store_block_with_witness(&b500, &[], 500)
.unwrap();
let b501 = mk(1_600_000_001);
let h501 = blockstore.get_block_hash(&b501);
blockstore.store_height(501, &h501).unwrap();
blockstore
.store_block_with_witness(&b501, &[], 501)
.unwrap();
let mut reorder_buffer = std::collections::BTreeMap::new();
reorder_buffer.insert(500, (Arc::new(b500), Arc::new(Vec::new())));
let mut already_dispatched = rustc_hash::FxHashSet::default();
already_dispatched.insert(500);
already_dispatched.insert(501);
let mut log_miss = rustc_hash::FxHashSet::default();
assert!(
coordinator_inject_local_gap(
&blockstore,
ProtocolVersion::BitcoinV1,
500,
0,
499,
&mut reorder_buffer,
&already_dispatched,
&mut log_miss,
true,
)
.unwrap()
);
assert_eq!(
reorder_buffer.len(),
1,
"already-present chain must not grow reorder_buffer"
);
assert!(reorder_buffer.contains_key(&500));
assert!(!reorder_buffer.contains_key(&501));
}
#[test]
fn inject_tip_not_in_pipeline_reloads_from_disk() {
let blockstore = temp_blockstore();
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(500, &hash).unwrap();
blockstore
.store_block_with_witness(&block, &[], 500)
.unwrap();
let mut reorder_buffer = std::collections::BTreeMap::new();
let mut already_dispatched = rustc_hash::FxHashSet::default();
already_dispatched.insert(500);
let mut log_miss = rustc_hash::FxHashSet::default();
assert!(
coordinator_inject_local_gap(
&blockstore,
ProtocolVersion::BitcoinV1,
500,
0,
499,
&mut reorder_buffer,
&already_dispatched,
&mut log_miss,
false, )
.unwrap()
);
assert!(
reorder_buffer.contains_key(&500),
"unconfirmed dispatched tip must reload from disk"
);
}
#[test]
fn witness_repair_stores_row_key_when_body_exists() {
let blockstore = temp_blockstore();
let height = 500_000u64;
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![TransactionOutput {
value: 50_0000_0000,
script_pubkey: vec![0x51],
}],
lock_time: 0,
}]
.into(),
};
let hash = blockstore.get_block_hash(&block);
blockstore.store_height(height, &hash).unwrap();
blockstore
.store_block_with_witness(&block, &[], height)
.unwrap();
assert!(!blockstore.has_witness_blob(&hash).unwrap());
let witnesses: Vec<Vec<Witness>> = vec![vec![vec![vec![0x51u8]]]];
let repaired = try_repair_missing_witness(
&blockstore,
height,
hash,
witnesses.as_slice(),
ProtocolVersion::BitcoinV1,
Some(&block),
)
.unwrap();
assert!(repaired);
assert!(blockstore.has_witness_blob(&hash).unwrap());
let loaded = blockstore.get_witness(&hash).unwrap().unwrap();
assert!(has_real_witnesses(&loaded));
}
}