fn shared_empty_witness_stacks(n_tx: usize) -> Arc<Vec<Vec<Witness>>> {
EMPTY_WITNESS_STACKS.with(|cell| {
let mut g = cell.borrow_mut();
if let Some(a) = g.get(&n_tx) {
return Arc::clone(a);
}
let arc = Arc::new(vec![Vec::new(); n_tx]);
if g.len() > 512 {
g.clear();
}
g.insert(n_tx, Arc::clone(&arc));
arc
})
}
static LAST_IBD_HEAP_TRIM_WALL_MS: AtomicU64 = AtomicU64::new(0);
const IBD_HEAP_TRIM_MIN_INTERVAL_MS: u64 = 2_000;
fn async_engine_append_enabled() -> bool {
match std::env::var("BLVM_IBD_ASYNC_ENGINE_APPEND") {
Ok(v) => {
let t = v.trim();
!(t == "0" || t.eq_ignore_ascii_case("false") || t.eq_ignore_ascii_case("off"))
}
Err(_) => true,
}
}
fn binder_log_enabled() -> bool {
latch_env!(bool, {
matches!(
std::env::var("BLVM_IBD_BINDER_LOG")
.ok()
.as_deref()
.map(str::trim),
Some("1") | Some("true") | Some("on") | Some("yes")
)
})
}
fn classify_ibd_binder(
feeder: usize,
holes: u64,
contig: u64,
await_ms: u64,
gd_ewma_ms: Option<u64>,
pressure: PressureLevel,
tip_failover: bool,
) -> &'static str {
if matches!(pressure, PressureLevel::Critical | PressureLevel::Emergency) {
return "ENGINE_PRESSURE";
}
if await_ms >= 200 || (feeder == 0 && holes > 0 && contig == 0) {
return "SUPPLY_TIP_HOLE";
}
if tip_failover {
return "SUPPLY_FAILOVER";
}
if feeder == 0 && contig == 0 {
return "SUPPLY_EMPTY_TIP";
}
if feeder == 0 {
let gd_ok = gd_ewma_ms.map(|ms| ms < 200).unwrap_or(false);
if holes == 0 && contig > 0 && gd_ok {
return "PIPE_DRAINED";
}
return "SUPPLY_FEEDER_STARVE";
}
if let Some(ms) = gd_ewma_ms {
if ms >= 200 && feeder < 16 {
return "SUPPLY_GD_SLOW";
}
}
if contig == 0 && feeder < 8 {
return "SUPPLY_THIN_RUNWAY";
}
"ENGINE_OR_SCRIPTS"
}
static IBD_EMERGENCY_EVICT_BLOCKS_SEEN: AtomicU64 = AtomicU64::new(0);
const IBD_EMERGENCY_EVICT_EVERY_N_BLOCKS: u64 = 8;
const IBD_EMERGENCY_EVICT_MIN_UNPROTECTED: usize = 32_768;
fn ibd_maybe_heap_trim() {
let now_ms = crate::utils::time::current_timestamp_millis();
loop {
let prev = LAST_IBD_HEAP_TRIM_WALL_MS.load(Ordering::Relaxed);
if now_ms.saturating_sub(prev) < IBD_HEAP_TRIM_MIN_INTERVAL_MS {
return;
}
if LAST_IBD_HEAP_TRIM_WALL_MS
.compare_exchange_weak(prev, now_ms, Ordering::Relaxed, Ordering::Relaxed)
.is_ok()
{
break;
}
}
}
use super::IbdBlockFlushOpts;
use super::ParallelIBD;
use super::ibd_staging::empty_utxo_delta;
use blvm_protocol::block::UtxoDelta;
#[allow(clippy::too_many_arguments)]
pub(crate) fn ibd_v2_retire_pre_lock(
next_height: u64,
store: &IbdUtxoStore,
blocks_buf: &[Arc<Block>],
keys_buf: &mut Vec<OutPointKey>,
keys_seen: &mut rustc_hash::FxHashSet<OutPointKey>,
evict_scratch: &mut Vec<(OutPointKey, u64)>,
) {
let cache_len = store.len() as u64;
let evict_interval: u64 = if cache_len > 5_000_000 {
64
} else if cache_len > 2_000_000 {
32
} else {
16
};
if next_height % evict_interval == 0 {
store.maybe_evict(evict_scratch);
}
if store.is_dynamic_eviction() {
block_input_keys_batch_into_arc(blocks_buf, keys_buf, keys_seen);
store.protect_keys_for_next_blocks(keys_buf);
store.evict_if_needed(next_height);
}
}
pub(crate) fn ibd_v2_retire_post_lock(
store: &IbdUtxoStore,
new_cap: usize,
pre_tune_len: usize,
current_height: u64,
) {
store.tune_max_entries_for_pressure(new_cap, current_height);
let evicted = pre_tune_len.saturating_sub(store.len());
if evicted > 32_768 {
let _ = evicted; }
}
fn ibd_empty_checkpoint_package(boundary_height: u64) -> PendingFlushPackage {
PendingFlushPackage {
ops: Arc::new(Vec::new()),
max_block_height: boundary_height,
heights: Arc::new(FxHashSet::default()),
}
}
fn ibd_formal_checkpoint_flush_batch(
store: &IbdUtxoStore,
next_height: u64,
) -> Option<PendingFlushPackage> {
let pending_before = store.pending_len();
let batch = store.take_flush_batch_adds_only().or_else(|| {
if store.pending_len() > 0 {
Some(ibd_empty_checkpoint_package(next_height))
} else {
None
}
});
warn!(
"[CAPPED_DRAIN] path=at_checkpoint h={next_height} pending_before={pending_before} drained={}",
batch.as_ref().map(|p| p.ops.len()).unwrap_or(0)
);
batch
}
pub(crate) fn ibd_v2_retire_apply_utxo_delta(
next_height: u64,
store: &IbdUtxoStore,
mem_guard: &mut MemoryGuard,
max_ahead_live: &Arc<AtomicU64>,
nominal_max_ahead: u64,
ibd_defer_flush: bool,
ibd_defer_checkpoint: u64,
) -> (
u64,
u64,
Option<PendingFlushPackage>,
bool,
Option<(usize, usize)>,
) {
let pressure_level = mem_guard.should_flush(Some((max_ahead_live, nominal_max_ahead)));
memory::publish_ibd_pressure(pressure_level);
let cap_change: Option<(usize, usize)> = mem_guard
.compute_adaptive_cache_cap()
.map(|new_cap| (new_cap, store.len()));
if ibd_defer_flush && next_height > 0 && next_height % ibd_defer_checkpoint == 0 {
let batch = ibd_formal_checkpoint_flush_batch(store, next_height);
return (0u64, 0u64, batch, true, cap_change);
}
let rss_pressure = pressure_level >= PressureLevel::Critical;
if rss_pressure {
let pending_now = store.pending_len();
info!(
"[IBD_V2] height={} RSS pressure ({:?}, cache={}, pending={}), forcing flush",
next_height,
pressure_level,
store.len(),
pending_now
);
if pressure_level == PressureLevel::Emergency {
let n = IBD_EMERGENCY_EVICT_BLOCKS_SEEN.fetch_add(1, Ordering::Relaxed);
if n % IBD_EMERGENCY_EVICT_EVERY_N_BLOCKS == 0 {
let cache_now = store.len();
let protected_now = store.protected_len();
let evictable = cache_now.saturating_sub(protected_now);
if evictable >= IBD_EMERGENCY_EVICT_MIN_UNPROTECTED {
store.evict_aggressive_for_rss();
}
}
}
let pending_before = store.pending_len();
let batch = store.maybe_take_flush_batch_adds_only();
if batch.is_some() {
warn!(
"[CAPPED_DRAIN] path=rss_pressure level={:?} h={next_height} pending_before={pending_before} drained={} force_ckpt=false",
pressure_level,
batch.as_ref().map(|p| p.ops.len()).unwrap_or(0),
);
}
ibd_maybe_heap_trim();
(0u64, 0u64, batch, false, cap_change)
} else {
let (batch, force_durability) = ibd_retire_pick_flush_batch(store);
(0u64, 0u64, batch, force_durability, cap_change)
}
}
#[cfg(feature = "profile")]
#[inline]
fn ibd_profile_height_matches_sample(sample: u64, height: u64) -> bool {
sample == 1 || (sample > 0 && height % sample == 0)
}
#[inline]
fn dynamic_utxo_cap(level: PressureLevel, nominal: usize) -> usize {
if nominal == usize::MAX {
return usize::MAX;
}
match level {
PressureLevel::Emergency => (nominal / 4).max(8_192),
PressureLevel::Critical => (nominal * 2 / 3).max(nominal / 2),
PressureLevel::Elevated => (nominal * 9 / 10).max(nominal * 4 / 5),
PressureLevel::None => nominal,
}
}
#[inline]
fn dynamic_prefetch_lookahead(level: PressureLevel, nominal: usize) -> usize {
let n = nominal.clamp(1, 128);
match level {
PressureLevel::Emergency => 8,
PressureLevel::Critical => (n / 2).clamp(12, 48),
PressureLevel::Elevated => ((n * 2 / 3).max(24)).min(n),
PressureLevel::None => n,
}
}
#[inline]
fn pipeline_depth_for_engine_append(nominal: usize) -> usize {
let slow_pct = crate::storage::ibd_engine::memory_age::memory_age_throttle_slow_pct();
if slow_pct == 0 {
return nominal;
}
if slow_pct >= 95 {
1
} else if slow_pct >= 50 {
(nominal / 4).max(4)
} else if slow_pct >= 25 {
(nominal / 2).max(8)
} else {
nominal
}
}
fn tip_concurrency_adapt_from_env() -> bool {
latch_env!(bool, {
matches!(
std::env::var("BLVM_IBD_TIP_CONCURRENCY_ADAPT")
.ok()
.as_deref()
.map(str::trim),
Some("1") | Some("true") | Some("yes") | Some("on")
)
})
}
#[inline]
fn effective_pressure_for_tip_crawl(level: PressureLevel) -> PressureLevel {
if matches!(level, PressureLevel::Critical) && super::tip_stage::tip_crawl_supply_healthy_now()
{
PressureLevel::Elevated
} else {
level
}
}
#[inline]
fn pipeline_depth_for_pressure(level: PressureLevel, nominal: usize) -> usize {
let mut nominal = pipeline_depth_for_engine_append(nominal);
if tip_concurrency_adapt_from_env() {
let segs = crate::storage::ibd_engine::tip_disk_segs_hint();
let disk_ms = crate::storage::ibd_engine::tip_disk_ms_hint();
if segs >= 2 || disk_ms >= 8 {
nominal = (nominal / 2).max(8);
}
}
match effective_pressure_for_tip_crawl(level) {
PressureLevel::Emergency => (nominal / 4).max(4),
PressureLevel::Critical => (nominal / 2).max(8),
PressureLevel::Elevated => (nominal * 3 / 4).max(12),
PressureLevel::None => nominal,
}
}
#[inline]
fn engine_pressure_poll_interval(level: PressureLevel) -> u64 {
match effective_pressure_for_tip_crawl(level) {
PressureLevel::Emergency => 1,
PressureLevel::Critical => 4,
PressureLevel::Elevated => 16,
PressureLevel::None => 32,
}
}
pub(super) struct DurabilityRequest {
pub(super) pkg: PendingFlushPackage,
pub(super) trigger_height: u64,
pub(super) is_checkpoint: bool,
}
fn ibd_durability_may_merge_request(batch: &[DurabilityRequest], next: &DurabilityRequest) -> bool {
if batch.is_empty() {
return true;
}
let batch_is_ckpt = batch.iter().any(|r| r.is_checkpoint);
next.is_checkpoint == batch_is_ckpt
}
#[allow(clippy::too_many_arguments)]
fn run_ibd_durability_loop(
store: Arc<IbdUtxoStore>,
storage_wm: Arc<Storage>,
utxo_flush_handles: Arc<Mutex<VecDeque<JoinHandle<Result<blvm_muhash::MuHash3072>>>>>,
ibd_muhash: Arc<Mutex<blvm_muhash::MuHash3072>>,
retire_err: Arc<Mutex<Option<anyhow::Error>>>,
in_flight_counter: Arc<AtomicUsize>,
ibd_defer_checkpoint: u64,
rx: std::sync::mpsc::Receiver<DurabilityRequest>,
) {
const BATCH_MAX: usize = 8;
const BATCH_MAX_OPS: usize = 200_000;
loop {
let first = match rx.recv() {
Ok(r) => r,
Err(_) => return, };
let mut total_ops = first.pkg.ops.len();
let mut batch: Vec<DurabilityRequest> = Vec::with_capacity(BATCH_MAX);
batch.push(first);
while batch.len() < BATCH_MAX && total_ops < BATCH_MAX_OPS {
match rx.try_recv() {
Ok(r) => {
if !ibd_durability_may_merge_request(&batch, &r) {
break;
}
total_ops += r.pkg.ops.len();
batch.push(r);
}
Err(_) => break,
}
}
let n_reqs = batch.len();
let has_checkpoint = batch.iter().any(|r| r.is_checkpoint);
let total_ops: usize = batch.iter().map(|r| r.pkg.ops.len()).sum();
let first_h = batch[0].trigger_height;
let t0 = std::time::Instant::now();
in_flight_counter.fetch_add(1, Ordering::Relaxed);
if has_checkpoint {
let cache_entries = store.cache_len();
let est_cache_mb = (cache_entries as u64 * 216) / (1_024 * 1_024);
let formal = first_h > 0 && first_h % ibd_defer_checkpoint == 0;
info!(
"[IBD_DURABILITY] h={first_h}: checkpoint batch started \
(requests={n_reqs}, total_ops={total_ops}, formal={formal}, \
cache_entries={cache_entries} ~{est_cache_mb}MB)"
);
}
{
let mut combined_sub_mh = blvm_muhash::MuHash3072::new();
loop {
let handle = { utxo_flush_handles.lock().pop_front() };
let Some(handle) = handle else { break };
match join_utxo_flush_handle_collect_sub_mh(handle) {
Ok(sub) => combined_sub_mh = combined_sub_mh.multiply(&sub),
Err(e) => {
*retire_err.lock() = Some(e);
return;
}
}
}
{
let mut mh_guard = ibd_muhash.lock();
*mh_guard = std::mem::take(&mut *mh_guard).multiply(&combined_sub_mh);
}
}
let t_after_handle_join = t0.elapsed().as_millis();
let codec = store.value_codec();
let mut prepared_pkgs: Vec<(
Arc<FxHashSet<u32>>,
crate::storage::ibd_utxo_store::PreparedFlushPackage,
)> = Vec::with_capacity(n_reqs);
for req in &batch {
let heights = Arc::clone(&req.pkg.heights);
match req.pkg.prepare_for_disk(codec) {
Ok(p) => prepared_pkgs.push((heights, p)),
Err(e) => {
*retire_err.lock() = Some(e);
return;
}
}
}
let t_after_prepare = t0.elapsed().as_millis();
let mut local_mh = blvm_muhash::MuHash3072::new();
for (_, prepared) in &prepared_pkgs {
if let Err(e) = store.compute_package_muhash(prepared, &mut local_mh) {
*retire_err.lock() = Some(e);
return;
}
}
let max_height = prepared_pkgs
.iter()
.map(|(_, p)| p.max_block_height)
.max()
.or_else(|| batch.iter().map(|r| r.trigger_height).max())
.unwrap_or(0);
let total_rows: usize = prepared_pkgs.iter().map(|(_, p)| p.rows.len()).sum();
let checkpoint_wm = has_checkpoint
.then(|| ibd_checkpoint_watermark_for_batch(&batch, max_height, ibd_defer_checkpoint));
let muhash_running_opt: Option<[u8; blvm_muhash::MUHASH_RUNNING_STATE_BYTES]> = {
let mut mh_guard = ibd_muhash.lock();
*mh_guard = std::mem::take(&mut *mh_guard).multiply(&local_mh);
if has_checkpoint {
Some(mh_guard.serialize_running_state())
} else {
None
}
};
let t_after_muhash = t0.elapsed().as_millis();
for (_, prepared) in &prepared_pkgs {
if let Err(e) = store.flush_prepared_package_adds_only(prepared) {
*retire_err.lock() = Some(e);
return;
}
}
let t_after_adds = t0.elapsed().as_millis();
for (heights, prepared) in &prepared_pkgs {
store.release_protected_heights(heights);
store.note_utxo_flush_completed(prepared.max_block_height);
}
let t_after_release = t0.elapsed().as_millis();
let skip_del_lmdb: bool = std::env::var("BLVM_IBD_SKIP_DEL_LMDB")
.map(|v| v != "0")
.unwrap_or(false);
let (t_after_dels, t_after_sync) = if has_checkpoint {
if !skip_del_lmdb {
for (_, prepared) in &prepared_pkgs {
if let Err(e) = store.flush_prepared_package_dels_only(prepared) {
*retire_err.lock() = Some(e);
return;
}
}
}
let t_dels = t0.elapsed().as_millis();
if let Err(e) = store.flush_disk() {
*retire_err.lock() = Some(e);
return;
}
let t_sync = t0.elapsed().as_millis();
(t_dels, t_sync)
} else {
(t_after_adds, t_after_adds)
};
let prep_ms = t_after_prepare.saturating_sub(t_after_handle_join);
let release_early_ms = t_after_release.saturating_sub(t_after_adds);
let adds_ms = t_after_adds.saturating_sub(t_after_muhash);
let dels_ms = t_after_dels.saturating_sub(t_after_release);
let sync_ms = t_after_sync.saturating_sub(t_after_dels);
if has_checkpoint {
info!(
"[IBD_DURABILITY_TIMING] h={first_h} \
prep={prep_ms}ms release_early={release_early_ms}ms \
adds={adds_ms}ms dels={dels_ms}ms sync={sync_ms}ms"
);
}
let (t_after_del_backlog, del_batch_n) = if let Some(checkpoint_wm) = checkpoint_wm {
let mut del_batch_n = 0usize;
if !skip_del_lmdb {
match ibd_flush_del_backlog_drain(&store, &ibd_muhash, checkpoint_wm, true) {
Ok(n) => del_batch_n = n,
Err(e) => {
*retire_err.lock() = Some(e);
return;
}
}
} else {
if let Err(e) = ibd_flush_leftover_adds_through_watermark(&store, checkpoint_wm) {
*retire_err.lock() = Some(e);
return;
}
let discarded = store.discard_del_backlog_through_watermark(checkpoint_wm);
if discarded > 0 {
tracing::debug!(
"[SKIP_DEL_PURGE] h={first_h} wm={checkpoint_wm} discarded={discarded} DEL tombstones"
);
}
}
if store.has_pending_adds_at_or_below(checkpoint_wm) {
*retire_err.lock() = Some(anyhow::anyhow!(
"checkpoint incomplete: ADDs remain at/below formal wm={checkpoint_wm} (h={first_h})"
));
return;
}
if let Err(e) =
ibd_persist_checkpoint_watermark(&storage_wm, &store, &ibd_muhash, checkpoint_wm)
{
*retire_err.lock() = Some(e);
return;
}
(t0.elapsed().as_millis(), del_batch_n)
} else {
(t_after_release, 0)
};
in_flight_counter.fetch_sub(1, Ordering::Relaxed);
let elapsed_ms = t0.elapsed().as_millis();
if has_checkpoint {
let del_backlog_ms = t_after_del_backlog.saturating_sub(t_after_sync);
let rss_kb = std::fs::read_to_string("/proc/self/status")
.ok()
.and_then(|s| {
s.lines()
.find(|l| l.starts_with("VmRSS:"))
.and_then(|l| l.split_whitespace().nth(1))
.and_then(|v| v.parse::<u64>().ok())
})
.unwrap_or(0);
let wm_logged = checkpoint_wm.unwrap_or(max_height);
let del_path = if ibd_del_backlog_use_collapse_path(wm_logged) {
"collapse"
} else {
"fast"
};
info!(
"[IBD_DURABILITY] h={first_h}: checkpoint batch complete \
(requests={n_reqs}, rows={total_rows}, wm={wm_logged}, elapsed={elapsed_ms}ms \
del_backlog={del_backlog_ms}ms del_batches={del_batch_n} del_path={del_path} \
pending_after={} rss={}MB)",
store.pending_len(),
rss_kb / 1024
);
}
}
}
pub(crate) struct IbdRetireWork {
pub(crate) height: u64,
pub(crate) blocks_buf: Vec<Arc<Block>>,
pub(crate) block: Option<Arc<Block>>,
}
const IBD_MIN_ADDS_ONLY_BATCH: usize = 1_000;
fn ibd_retire_pick_flush_batch(store: &IbdUtxoStore) -> (Option<PendingFlushPackage>, bool) {
let pending = store.pending_len();
let threshold = store.flush_threshold();
let Some(pkg) = store.maybe_take_flush_batch_adds_only() else {
return (None, false);
};
let remaining = store.pending_len();
if pkg.ops.len() < IBD_MIN_ADDS_ONLY_BATCH && remaining > threshold {
return (Some(pkg), false);
}
(Some(pkg), false)
}
fn ibd_formal_checkpoint_boundary(trigger_height: u64, interval: u64) -> Option<u64> {
if trigger_height > 0 && interval > 0 && trigger_height % interval == 0 {
Some(trigger_height)
} else {
None
}
}
fn ibd_checkpoint_watermark_for_batch(
batch: &[DurabilityRequest],
batch_max_prepared: u64,
ibd_defer_checkpoint: u64,
) -> u64 {
let ckpt_triggers: Vec<u64> = batch
.iter()
.filter(|r| r.is_checkpoint)
.map(|r| r.trigger_height)
.collect();
if ckpt_triggers.is_empty() {
return batch_max_prepared;
}
if let Some(formal) = ckpt_triggers
.iter()
.filter_map(|&h| ibd_formal_checkpoint_boundary(h, ibd_defer_checkpoint))
.min()
{
return formal;
}
let max_trigger = *ckpt_triggers.iter().max().unwrap_or(&0);
max_trigger.min(batch_max_prepared)
}
fn ibd_persist_checkpoint_watermark(
storage_wm: &Arc<Storage>,
store: &Arc<IbdUtxoStore>,
ibd_muhash: &Arc<Mutex<blvm_muhash::MuHash3072>>,
watermark: u64,
) -> Result<()> {
let muhash_running = ibd_muhash.lock().serialize_running_state();
storage_wm
.chain()
.persist_ibd_utxo_flush_checkpoint(watermark, &muhash_running)?;
store.note_utxo_flush_completed(watermark);
Ok(())
}
fn ibd_flush_leftover_adds_through_watermark(
store: &Arc<IbdUtxoStore>,
watermark: u64,
) -> Result<usize> {
let cap = store.adaptive_drain_cap();
let mut batch_n = 0usize;
while store.has_pending_adds_at_or_below(watermark) {
let Some(follow) = store.take_flush_batch_adds_only_through_capped(watermark, cap) else {
break;
};
batch_n += 1;
let heights = Arc::clone(&follow.heights);
let prepared = follow.prepare_for_disk(store.value_codec())?;
store.flush_prepared_package_adds_only(&prepared)?;
store.flush_disk()?;
store.release_protected_heights(&heights);
store.note_utxo_flush_completed(prepared.max_block_height);
}
Ok(batch_n)
}
fn ibd_del_backlog_use_collapse_path(watermark: u64) -> bool {
let min_h = std::env::var("BLVM_IBD_DEL_COLLAPSE_MIN_HEIGHT")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(350_000);
min_h == 0 || watermark >= min_h
}
fn ibd_flush_del_backlog_collapse_drain(
store: &Arc<IbdUtxoStore>,
ibd_muhash: &Arc<Mutex<blvm_muhash::MuHash3072>>,
watermark: u64,
) -> Result<usize> {
ibd_flush_leftover_adds_through_watermark(store, watermark)?;
let cap = store.adaptive_drain_cap();
let mut batch_n = 0usize;
loop {
let Some(follow) = store.take_flush_batch_dels_only_through_capped(watermark, cap) else {
break;
};
batch_n += 1;
let ops_len = follow.ops.len();
let heights = Arc::clone(&follow.heights);
let prepared = follow.prepare_for_disk(store.value_codec())?;
if crate::storage::ibd_utxo_store::ibd_per_op_muhash_enabled() {
let mut local_mh = blvm_muhash::MuHash3072::new();
store.compute_package_muhash(&prepared, &mut local_mh)?;
let mut mh_guard = ibd_muhash.lock();
*mh_guard = std::mem::take(&mut *mh_guard).multiply(&local_mh);
}
store.flush_prepared_package_dels_only(&prepared)?;
store.release_protected_heights(&heights);
debug!("[DEL_BACKLOG] wm={watermark} batch={batch_n} ops={ops_len} cap={cap}");
}
if batch_n > 0 {
store.flush_disk_sync_only()?;
}
Ok(batch_n)
}
fn ibd_flush_del_backlog_fast_drain(
store: &Arc<IbdUtxoStore>,
ibd_muhash: &Arc<Mutex<blvm_muhash::MuHash3072>>,
watermark: u64,
) -> Result<usize> {
let Some(follow) = store.take_flush_batch_force_through(watermark) else {
return Ok(0);
};
if follow.ops.is_empty() {
return Ok(0);
}
let heights = Arc::clone(&follow.heights);
let prepared = follow.prepare_for_disk(store.value_codec())?;
let mut local_mh = blvm_muhash::MuHash3072::new();
store.compute_package_muhash(&prepared, &mut local_mh)?;
{
let mut mh_guard = ibd_muhash.lock();
*mh_guard = std::mem::take(&mut *mh_guard).multiply(&local_mh);
}
store.flush_prepared_package_adds_only(&prepared)?;
store.flush_prepared_package_dels_only(&prepared)?;
store.flush_disk()?;
store.release_protected_heights(&heights);
store.note_utxo_flush_completed(watermark);
Ok(1)
}
fn ibd_flush_del_backlog_drain(
store: &Arc<IbdUtxoStore>,
ibd_muhash: &Arc<Mutex<blvm_muhash::MuHash3072>>,
watermark: u64,
force_collapse: bool,
) -> Result<usize> {
if force_collapse || ibd_del_backlog_use_collapse_path(watermark) {
ibd_flush_del_backlog_collapse_drain(store, ibd_muhash, watermark)
} else {
ibd_flush_del_backlog_fast_drain(store, ibd_muhash, watermark)
}
}
fn ibd_flush_del_backlog_through_watermark(
store: &Arc<IbdUtxoStore>,
storage_wm: &Arc<Storage>,
ibd_muhash: &Arc<Mutex<blvm_muhash::MuHash3072>>,
watermark: u64,
) -> Result<usize> {
let batch_n = ibd_flush_del_backlog_drain(store, ibd_muhash, watermark, false)?;
if batch_n > 0 || !store.has_pending_adds_at_or_below(watermark) {
ibd_persist_checkpoint_watermark(storage_wm, store, ibd_muhash, watermark)?;
}
Ok(batch_n)
}
fn retire_flush_batch_size() -> usize {
latch_env!(usize, {
std::env::var("BLVM_IBD_RETIRE_FLUSH_BATCH")
.ok()
.and_then(|s| s.parse().ok())
.filter(|n: &usize| *n >= 1)
.unwrap_or(16)
})
}
fn retire_do_durability(
force_durability: bool,
batch_count: usize,
n: usize,
durability_on_channel: bool,
) -> bool {
if durability_on_channel {
force_durability
} else {
force_durability || batch_count <= 1 || n % batch_count == 0
}
}
fn read_mem_available_kb() -> u64 {
std::fs::read_to_string("/proc/meminfo")
.ok()
.and_then(|s| {
s.lines()
.find(|l| l.starts_with("MemAvailable:"))
.and_then(|l| l.split_whitespace().nth(1))
.and_then(|v| v.parse::<u64>().ok())
})
.unwrap_or(0)
}
fn ibd_durability_channel_cap(retire_height: u64) -> usize {
let env_cap = std::env::var("BLVM_IBD_DURABILITY_CHANNEL_CAP")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(4)
.clamp(2, 16);
if retire_height >= 481_000 {
let spare_mb = read_mem_available_kb() / 1024;
let ram_cap = (spare_mb / 1500).max(2) as usize;
env_cap.min(ram_cap)
} else {
env_cap
}
}
fn join_utxo_flush_handle_mul_sub_mh(
handle: JoinHandle<Result<blvm_muhash::MuHash3072>>,
ibd_muhash: &Arc<Mutex<blvm_muhash::MuHash3072>>,
) -> Result<()> {
match handle.join() {
Ok(Ok(sub_mh)) => {
let mut mh_guard = ibd_muhash.lock();
*mh_guard = std::mem::take(&mut *mh_guard).multiply(&sub_mh);
Ok(())
}
Ok(Err(e)) => Err(e),
Err(e) => Err(anyhow::anyhow!("UTXO flush panicked: {:?}", e)),
}
}
fn join_utxo_flush_handle_collect_sub_mh(
handle: JoinHandle<Result<blvm_muhash::MuHash3072>>,
) -> Result<blvm_muhash::MuHash3072> {
match handle.join() {
Ok(Ok(sub_mh)) => Ok(sub_mh),
Ok(Err(e)) => Err(e),
Err(e) => Err(anyhow::anyhow!("UTXO flush panicked: {:?}", e)),
}
}
struct UtxoFlushGuard(Arc<Mutex<VecDeque<JoinHandle<Result<blvm_muhash::MuHash3072>>>>>);
impl Drop for UtxoFlushGuard {
fn drop(&mut self) {
let handles: Vec<_> = self.0.lock().drain(..).collect();
if handles.is_empty() {
return;
}
warn!(
"[IBD_FLUSH_GUARD] joining {} leaked UTXO flush thread(s) on validation-loop exit \
(prevents ibd_utxo_store LMDB handle leak into next IBD session)",
handles.len()
);
for h in handles {
if let Err(e) = h.join() {
warn!(
"[IBD_FLUSH_GUARD] UTXO flush thread panicked during cleanup: {:?}",
e
);
}
}
}
}
fn join_all_utxo_flush_handles(
utxo_flush_handles: &Arc<Mutex<VecDeque<JoinHandle<Result<blvm_muhash::MuHash3072>>>>>,
log_label: &str,
) -> Result<blvm_muhash::MuHash3072> {
let handles: Vec<_> = utxo_flush_handles.lock().drain(..).collect();
let n = handles.len();
if n > 0 {
info!("IBD shutdown: joining {n} in-flight UTXO flush thread(s) ({log_label})");
}
let mut combined = blvm_muhash::MuHash3072::new();
for (i, handle) in handles.into_iter().enumerate() {
debug!(
"IBD shutdown: UTXO flush join {}/{} ({log_label})",
i + 1,
n
);
let sub = join_utxo_flush_handle_collect_sub_mh(handle)?;
combined = combined.multiply(&sub);
}
Ok(combined)
}
fn push_utxo_flush_from_retire(
store: &Arc<IbdUtxoStore>,
storage_wm: &Arc<Storage>,
utxo_flush_handles: &Arc<Mutex<VecDeque<JoinHandle<Result<blvm_muhash::MuHash3072>>>>>,
retire_flush_counter: &Arc<AtomicUsize>,
next_height: u64,
max_utxo_flushes_in_flight: usize,
pkg: PendingFlushPackage,
ibd_muhash: &Arc<Mutex<blvm_muhash::MuHash3072>>,
force_durability: bool,
durability_tx: Option<&std::sync::mpsc::SyncSender<DurabilityRequest>>,
) -> Result<()> {
let flush_limit = memory::utxo_flush_concurrency_cap(max_utxo_flushes_in_flight).max(1);
let batch_count = retire_flush_batch_size();
let n = retire_flush_counter.fetch_add(1, Ordering::Relaxed);
let do_durability =
retire_do_durability(force_durability, batch_count, n, durability_tx.is_some());
if let Some(tx) = durability_tx {
let ops_len = pkg.ops.len();
let est_mb = (ops_len * 200) / (1024 * 1024);
debug!(
"[IBD_DURABILITY] h={next_height}: queuing flush \
(ops={ops_len}, ~{est_mb}MB, checkpoint={do_durability})"
);
let t_send = std::time::Instant::now();
if tx
.send(DurabilityRequest {
pkg,
trigger_height: next_height,
is_checkpoint: do_durability,
})
.is_err()
{
return Err(anyhow::anyhow!(
"IBD durability thread disconnected at h={next_height}"
));
}
let send_ms = t_send.elapsed().as_millis();
if send_ms > 500 {
let rss_kb = std::fs::read_to_string("/proc/self/status")
.ok()
.and_then(|s| {
s.lines()
.find(|l| l.starts_with("VmRSS:"))
.and_then(|l| l.split_whitespace().nth(1))
.and_then(|v| v.parse::<u64>().ok())
})
.unwrap_or(0);
warn!(
"[IBD_DURABILITY] h={next_height}: channel FULL — send blocked {send_ms}ms \
(ops={ops_len} ~{est_mb}MB, rss={}MB). \
Durability thread is behind; back-pressure is expected but sustained \
blocking indicates durability is the IBD bottleneck.",
rss_kb / 1024
);
}
return Ok(());
}
loop {
let handle = {
let mut q = utxo_flush_handles.lock();
if q.len() < flush_limit {
None
} else {
q.pop_front()
}
};
let Some(handle) = handle else {
break;
};
join_utxo_flush_handle_mul_sub_mh(handle, ibd_muhash)?;
}
let batch_size = pkg.ops.len();
let heights = Arc::clone(&pkg.heights);
if do_durability {
let mut combined_sub_mh = blvm_muhash::MuHash3072::new();
loop {
let handle = {
let mut q = utxo_flush_handles.lock();
q.pop_front()
};
let Some(handle) = handle else {
break;
};
let sub = join_utxo_flush_handle_collect_sub_mh(handle)?;
combined_sub_mh = combined_sub_mh.multiply(&sub);
}
let prepared = pkg.prepare_for_disk(store.value_codec())?;
drop(pkg);
let mut shutdown_local_mh = blvm_muhash::MuHash3072::new();
store.compute_package_muhash(&prepared, &mut shutdown_local_mh)?;
let muhash_running = {
let mut mh_guard = ibd_muhash.lock();
*mh_guard = std::mem::take(&mut *mh_guard)
.multiply(&combined_sub_mh)
.multiply(&shutdown_local_mh);
mh_guard.serialize_running_state()
};
store.flush_prepared_package_adds_only(&prepared)?;
store.flush_disk()?;
storage_wm
.chain()
.persist_ibd_utxo_flush_checkpoint(prepared.max_block_height, &muhash_running)?;
store.flush_prepared_package_dels_only(&prepared)?;
store.flush_disk()?;
store.release_protected_heights(&heights);
store.note_utxo_flush_completed(prepared.max_block_height);
let watermark = prepared.max_block_height;
ibd_flush_del_backlog_through_watermark(store, storage_wm, ibd_muhash, watermark)?;
debug!(
"[IBD_DEBUG] Block {}: durability flush boundary (batch_size={}, n={})",
next_height, batch_size, n,
);
} else {
let prepared = pkg.prepare_for_disk(store.value_codec())?;
drop(pkg);
let store_clone = Arc::clone(store);
utxo_flush_handles
.lock()
.push_back(std::thread::spawn(move || {
let mut local_mh = blvm_muhash::MuHash3072::new();
store_clone.compute_package_muhash(&prepared, &mut local_mh)?;
store_clone.flush_prepared_package_adds_only(&prepared)?;
store_clone.release_protected_heights(&heights);
store_clone.note_utxo_flush_completed(prepared.max_block_height);
Ok(local_mh)
}));
debug!(
"[IBD_DEBUG] Block {}: async commit (batch_size={}, in_flight={}, n={})",
next_height,
batch_size,
utxo_flush_handles.lock().len(),
n,
);
}
Ok(())
}
fn retire_thread_shutdown(
retire_dispatcher: &mut super::retire_dispatcher::RetireDispatcher,
retire_err: &Arc<Mutex<Option<anyhow::Error>>>,
) -> Result<()> {
retire_dispatcher.shutdown_and_join()?;
if let Some(e) = retire_err.lock().take() {
return Err(e);
}
Ok(())
}
struct EngineValidateJob {
height: u64,
block_arc: Arc<Block>,
witnesses_storage: Arc<Vec<Vec<Witness>>>,
bip30_index: Bip30Index,
recent_headers: Arc<Vec<Arc<BlockHeader>>>,
tx_ids: Vec<Hash>,
best_header_chainwork: blvm_consensus::pow::U256,
cached_network_time: u64,
partial_session: PartialSpendSession,
engine_append_ms: u64,
ibd_block_outputs:
Option<Arc<rustc_hash::FxHashMap<blvm_consensus::OutPoint, Arc<blvm_consensus::UTXO>>>>,
}
struct EngineAppendJob {
height: u64,
db: Arc<UtxoDatabase>,
block_arc: Arc<Block>,
witnesses_storage: Arc<Vec<Vec<Witness>>>,
bip30_index: Bip30Index,
recent_headers: Arc<Vec<Arc<BlockHeader>>>,
tx_ids: Vec<Hash>,
best_header_chainwork: blvm_consensus::pow::U256,
cached_network_time: u64,
ibd_block_outputs:
Option<Arc<rustc_hash::FxHashMap<blvm_consensus::OutPoint, Arc<blvm_consensus::UTXO>>>>,
}
struct LegacyValidateJob {
height: u64,
block_arc: Arc<Block>,
witnesses_storage: Arc<Vec<Vec<Witness>>>,
bip30_index: Bip30Index,
recent_headers: Arc<Vec<Arc<BlockHeader>>>,
tx_ids: Vec<Hash>,
best_header_chainwork: blvm_consensus::pow::U256,
cached_network_time: u64,
keys: Vec<OutPointKey>,
spec_adds_snapshot: Vec<(u64, Arc<UtxoSet>)>,
prefetched: PrefetchedUtxoMap,
ibd_block_outputs:
Option<Arc<rustc_hash::FxHashMap<blvm_consensus::OutPoint, Arc<blvm_consensus::UTXO>>>>,
}
enum ValidateJob {
Engine(EngineValidateJob),
Legacy(LegacyValidateJob),
}
struct ValidateResult {
height: u64,
result: Result<Option<UtxoDelta>>,
undo_log: blvm_consensus::reorganization::BlockUndoLog,
bip30_post: Bip30Index,
elapsed: std::time::Duration,
view_build_ms: u64,
engine_append_ms: u64,
engine_complete_ms: u64,
block_muhash: Option<blvm_muhash::MuHash3072>,
}
struct InFlightEntry {
height: u64,
block_arc: Arc<Block>,
witnesses_storage: Arc<Vec<Vec<Witness>>>,
feeder_est_bytes: usize,
utxo_base_ms: u64,
utxo_base_tune_ms: u64,
prefetch_ms: u64,
apply_pending_ms: u64,
input_keys: Option<Vec<OutPointKey>>,
}
#[allow(clippy::too_many_arguments)]
fn run_validation_worker_shared(
rx: crossbeam_channel::Receiver<ValidateJob>,
tx: crossbeam_channel::Sender<ValidateResult>,
parallel_ibd: Arc<super::ParallelIBD>,
blockstore: Arc<crate::storage::blockstore::BlockStore>,
protocol: Arc<blvm_protocol::BitcoinProtocolEngine>,
store: Arc<IbdUtxoStore>,
last_retired: Arc<AtomicU64>,
max_pending_ops: Arc<AtomicUsize>,
) {
let mut utxo_base: UtxoSet = UtxoSet::default();
let mut supplement_cache_buf: Vec<OutPointKey> = Vec::new();
let mut keys_missing_buf: Vec<OutPointKey> = Vec::new();
let mut del_scratch: Vec<OutPointKey> = Vec::new();
let mut add_scratch: Vec<(OutPointKey, Arc<UTXO>)> = Vec::new();
struct WavePending {
vr: ValidateResult,
rx: std::sync::mpsc::Receiver<std::result::Result<(), blvm_consensus::ConsensusError>>,
}
let mut wave_pending: Vec<WavePending> = Vec::new();
let wave_depth: usize = std::env::var("BLVM_ECDSA_WAVE_DEPTH")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(8)
.clamp(1, 32);
let flush_wave_ready =
|pending: &mut Vec<WavePending>, out: &crossbeam_channel::Sender<ValidateResult>| {
let mut i = 0;
while i < pending.len() {
match pending[i].rx.try_recv() {
Ok(Ok(())) => {
let WavePending { vr, .. } = pending.swap_remove(i);
let _ = out.send(vr);
}
Ok(Err(e)) => {
let WavePending { mut vr, .. } = pending.swap_remove(i);
vr.result = Err(anyhow::anyhow!("ECDSA wave: {e:?}"));
let _ = out.send(vr);
}
Err(std::sync::mpsc::TryRecvError::Empty) => i += 1,
Err(std::sync::mpsc::TryRecvError::Disconnected) => {
let WavePending { mut vr, .. } = pending.swap_remove(i);
vr.result = Err(anyhow::anyhow!("ECDSA wave disconnected"));
let _ = out.send(vr);
}
}
}
};
let block_wave_oldest =
|pending: &mut Vec<WavePending>, out: &crossbeam_channel::Sender<ValidateResult>| {
if pending.is_empty() {
return;
}
let WavePending { mut vr, rx } = pending.remove(0);
match rx.recv() {
Ok(Ok(())) => {
let _ = out.send(vr);
}
Ok(Err(e)) => {
vr.result = Err(anyhow::anyhow!("ECDSA wave: {e:?}"));
let _ = out.send(vr);
}
Err(_) => {
vr.result = Err(anyhow::anyhow!("ECDSA wave disconnected"));
let _ = out.send(vr);
}
}
};
loop {
flush_wave_ready(&mut wave_pending, &tx);
while wave_pending.len() >= wave_depth {
block_wave_oldest(&mut wave_pending, &tx);
flush_wave_ready(&mut wave_pending, &tx);
}
let job = match rx.recv() {
Ok(j) => j,
Err(_) => {
while !wave_pending.is_empty() {
block_wave_oldest(&mut wave_pending, &tx);
}
break;
}
};
match job {
ValidateJob::Legacy(mut lj) => {
let height = lj.height;
let t_view = std::time::Instant::now();
utxo_base.clear();
let n_keys = lj.keys.len();
let n_prefetched = lj.prefetched.len();
utxo_base.reserve(n_prefetched);
let still_missing = &mut keys_missing_buf;
still_missing.clear();
if n_prefetched > 0 {
match Arc::try_unwrap(lj.prefetched) {
Ok(mut map) => {
for (k, arc) in map.drain() {
utxo_base.insert(key_to_outpoint(&k), arc);
}
}
Err(shared) => {
for (k, arc) in shared.iter() {
utxo_base.insert(key_to_outpoint(k), Arc::clone(arc));
}
}
}
}
if n_prefetched < n_keys {
for k in lj.keys.iter() {
if !utxo_base.contains_key(&key_to_outpoint(k)) {
still_missing.push(*k);
}
}
if !still_missing.is_empty() {
still_missing.retain(|k| {
if let Some(ref r) = store.cache_get(k) {
let op = key_to_outpoint(k);
utxo_base.insert(op, Arc::clone(&r.utxo));
return false;
}
true
});
}
if !still_missing.is_empty() && !lj.spec_adds_snapshot.is_empty() {
still_missing.retain(|k| {
let op = key_to_outpoint(k);
for (_sh, set) in lj.spec_adds_snapshot.iter().rev() {
if let Some(u) = set.get(&op) {
utxo_base.insert(op, Arc::clone(u));
return false;
}
}
true
});
}
if !still_missing.is_empty() {
store.supplement_utxo_map_with_buf(
&mut utxo_base,
still_missing,
&mut supplement_cache_buf,
);
}
}
let view_build_ms = t_view.elapsed().as_millis() as u64;
let recent_opt: Option<&[Arc<BlockHeader>]> = if lj.recent_headers.is_empty() {
None
} else {
Some(lj.recent_headers.as_slice())
};
let t_val = std::time::Instant::now();
let raw = parallel_ibd.validate_block_only(
&blockstore,
protocol.as_ref(),
&mut utxo_base,
Some(&mut lj.bip30_index),
lj.block_arc.as_ref(),
Some(Arc::clone(&lj.block_arc)),
lj.witnesses_storage.as_slice(),
Some(&lj.witnesses_storage),
lj.height,
recent_opt,
lj.cached_network_time,
Some(&lj.tx_ids),
Some(lj.best_header_chainwork),
None,
lj.ibd_block_outputs.clone(),
);
let elapsed = t_val.elapsed();
let (result, undo_log) = match raw {
Ok((_ids, delta, undo)) => (Ok(delta), undo),
Err(e) => (Err(e), blvm_consensus::reorganization::BlockUndoLog::new()),
};
if let Ok(Some(delta)) = &result {
store.worker_cache_put_protected(&delta.additions, height);
store.apply_utxo_delta(delta, height, &mut del_scratch, &mut add_scratch, true);
}
let _ = tx.send(ValidateResult {
height,
result,
undo_log,
bip30_post: lj.bip30_index,
elapsed,
view_build_ms,
engine_append_ms: 0,
engine_complete_ms: 0,
block_muhash: None,
});
let cap = max_pending_ops.load(Ordering::Relaxed);
if cap > 0 {
let mut spins = 0u32;
let spin_start = std::time::Instant::now();
const WORKER_SPIN_MAX: std::time::Duration = std::time::Duration::from_secs(60);
const WORKER_SPIN_LOG_EVERY: std::time::Duration =
std::time::Duration::from_secs(5);
let mut last_log = spin_start;
loop {
let pending_now = store.pending_len();
let cap_now = max_pending_ops.load(Ordering::Relaxed);
if pending_now <= cap_now {
break;
}
if spin_start.elapsed() >= WORKER_SPIN_MAX {
warn!(
"IBD worker: pending cap spin timeout at h={height} pending={pending_now} cap={cap_now}"
);
break;
}
if last_log.elapsed() >= WORKER_SPIN_LOG_EVERY {
warn!(
"IBD worker: waiting for pending drain at h={height} pending={pending_now} cap={cap_now} spins={spins}"
);
last_log = std::time::Instant::now();
}
spins = spins.saturating_add(1);
std::thread::yield_now();
}
}
}
ValidateJob::Engine(mut ej) => {
let height = ej.height;
let t_view = std::time::Instant::now();
utxo_base.clear();
let t_complete = std::time::Instant::now();
let session = match ej.partial_session.complete() {
Ok(s) => s,
Err(e) => {
let _ = tx.send(ValidateResult {
height,
result: Err(e.context("IBD engine Phase 2 failed")),
undo_log: blvm_consensus::reorganization::BlockUndoLog::new(),
bip30_post: ej.bip30_index,
elapsed: std::time::Duration::ZERO,
view_build_ms: 0,
engine_append_ms: ej.engine_append_ms,
engine_complete_ms: 0,
block_muhash: None,
});
continue;
}
};
let use_lookup = crate::storage::ibd_engine::spend_session_lookup_enabled();
let t_fill = std::time::Instant::now();
let fill_ms = if use_lookup {
0u64
} else {
session_fill_utxo_set(&session, &mut utxo_base);
t_fill.elapsed().as_millis() as u64
};
let engine_complete_ms = t_complete.elapsed().as_millis() as u64;
let view_build_ms = t_view.elapsed().as_millis() as u64;
crate::storage::ibd_engine::note_tip_disk_hints(session.disk_segs, session.disk_ms);
let hotpath_n = hotpath_timer_sample();
if hotpath_n > 0 && height % hotpath_n == 0 {
info!(
"[IBD_HOTPATH] height={} query_ms={} ages_ms={} disk_ms={} preads={} pread_kb={} max_pread_kb={} cands={} segs={} fetch_ms={} map_ms={} fill_ms={} complete_ms={} inputs={}",
height,
session.query_ms,
session.ages_ms,
session.disk_ms,
session.disk_preads,
session.disk_pread_kb,
session.disk_max_pread_kb,
session.disk_cands,
session.disk_segs,
session.fetch_ms,
session.map_ms,
fill_ms,
engine_complete_ms,
session.details.len() + session.local_spends.len()
);
}
let recent_opt: Option<&[Arc<BlockHeader>]> = if ej.recent_headers.is_empty() {
None
} else {
Some(ej.recent_headers.as_slice())
};
let t_val = std::time::Instant::now();
let lookup = crate::storage::ibd_engine::SpendSessionLookup(&session);
let raw = if use_lookup {
utxo_base.clear();
parallel_ibd.validate_block_only(
&blockstore,
protocol.as_ref(),
&mut utxo_base,
Some(&mut ej.bip30_index),
ej.block_arc.as_ref(),
Some(Arc::clone(&ej.block_arc)),
ej.witnesses_storage.as_slice(),
Some(&ej.witnesses_storage),
ej.height,
recent_opt,
ej.cached_network_time,
Some(&ej.tx_ids),
Some(ej.best_header_chainwork),
Some(&lookup),
ej.ibd_block_outputs.clone(),
)
} else {
parallel_ibd.validate_block_only(
&blockstore,
protocol.as_ref(),
&mut utxo_base,
Some(&mut ej.bip30_index),
ej.block_arc.as_ref(),
Some(Arc::clone(&ej.block_arc)),
ej.witnesses_storage.as_slice(),
Some(&ej.witnesses_storage),
ej.height,
recent_opt,
ej.cached_network_time,
Some(&ej.tx_ids),
Some(ej.best_header_chainwork),
None,
ej.ibd_block_outputs.clone(),
)
};
let elapsed = t_val.elapsed();
let (result, undo_log) = match raw {
Ok((_ids, delta, undo)) => (Ok(delta), undo),
Err(e) => (Err(e), blvm_consensus::reorganization::BlockUndoLog::new()),
};
let block_muhash =
if crate::config::ibd::ibd_engine_muhash_enabled() && result.is_ok() {
let mut sub = blvm_muhash::MuHash3072::new();
crate::storage::ibd_utxo_muhash::fold_block_engine_muhash(
ej.block_arc.as_ref(),
&ej.tx_ids,
height,
&session,
&mut sub,
);
Some(sub)
} else {
None
};
let vr = ValidateResult {
height,
result,
undo_log,
bip30_post: ej.bip30_index,
elapsed,
view_build_ms,
engine_append_ms: ej.engine_append_ms,
engine_complete_ms,
block_muhash,
};
#[cfg(all(feature = "production"))]
{
if vr.result.is_ok() {
if let Some(soa) = blvm_consensus::ecdsa_wave::take_parked() {
let wrx = blvm_consensus::ecdsa_wave::submit(height, soa);
wave_pending.push(WavePending { vr, rx: wrx });
flush_wave_ready(&mut wave_pending, &tx);
if wave_pending.len() >= wave_depth.max(1) {
block_wave_oldest(&mut wave_pending, &tx);
} else if wave_pending.len() == 1 {
block_wave_oldest(&mut wave_pending, &tx);
}
continue;
}
}
}
let _ = tx.send(vr);
continue;
}
}
}
}
fn adapt_max_pending_ops_tick(
cap: &AtomicUsize,
nominal: usize,
pressure: PressureLevel,
pending_len: usize,
last_adapt_ms: &AtomicU64,
) {
const TICK_INTERVAL_MS: u64 = 500;
let now_ms = crate::utils::time::current_timestamp_millis();
let last = last_adapt_ms.load(Ordering::Relaxed);
if now_ms.saturating_sub(last) < TICK_INTERVAL_MS {
return;
}
if last_adapt_ms
.compare_exchange(last, now_ms, Ordering::Relaxed, Ordering::Relaxed)
.is_err()
{
return;
}
let current = cap.load(Ordering::Relaxed);
let pressure = effective_pressure_for_tip_crawl(pressure);
let new = match pressure {
PressureLevel::Emergency => {
(current * 9 / 10).max(nominal / 2).max(1_000_000)
}
PressureLevel::Critical => (current * 3 / 4).max(nominal / 4).max(500_000),
PressureLevel::Elevated => current,
PressureLevel::None => {
if pending_len < current / 4 {
let grown = (current as u128).saturating_mul(11) / 10;
let max = (nominal as u128).saturating_mul(11) / 10;
grown.min(max).max(nominal as u128 / 4) as usize
} else {
current
}
}
};
if new != current {
cap.store(new, Ordering::Relaxed);
if matches!(pressure, PressureLevel::Critical | PressureLevel::Emergency)
|| (pressure == PressureLevel::None && new > current)
{
tracing::debug!(
"[IBD_ADAPT] max_pending_ops {} → {} (pressure={:?}, pending={}, nominal={})",
current,
new,
pressure,
pending_len,
nominal
);
}
}
}
#[inline]
pub(crate) fn mtp_tip_window_fallback_ok(start_height: u64, tip_height: u64) -> bool {
tip_height > 0 && start_height.saturating_add(64) > tip_height
}
pub struct ValidationParams {
pub feeder_state: FeederState,
pub ibd_store: Arc<IbdUtxoStore>,
pub blockstore: Arc<BlockStore>,
pub storage: Arc<Storage>,
pub parallel_ibd: Arc<ParallelIBD>,
pub protocol: Arc<BitcoinProtocolEngine>,
pub utxo_mutex: Arc<std::sync::Mutex<UtxoSet>>,
pub effective_end_live: Arc<std::sync::atomic::AtomicU64>,
pub start_height: u64,
pub validation_height: Arc<std::sync::atomic::AtomicU64>,
pub mem_guard: MemoryGuard,
pub max_ahead_live: Arc<std::sync::atomic::AtomicU64>,
pub nominal_max_ahead: u64,
pub utxo_nominal_max_entries: usize,
pub utxo_prefetch_lookahead: usize,
pub stall_tx: tokio::sync::broadcast::Sender<u64>,
pub utxo_engine: Option<Arc<UtxoDatabase>>,
pub checkpoint_tx: Option<std::sync::mpsc::SyncSender<u64>>,
pub local_replay_max_height: u64,
pub engine_gap_export_defer_until: Option<u64>,
}
fn drain_ibd_pending_blocks_before_shutdown(
skip_storage: bool,
parallel_ibd: &ParallelIBD,
blockstore: &Arc<BlockStore>,
storage: &Arc<Storage>,
pending_blocks: &mut Vec<(
Arc<Block>,
Arc<Vec<Vec<Witness>>>,
u64,
blvm_consensus::reorganization::BlockUndoLog,
)>,
pending_storage_bytes: &mut u64,
flush_handles: &mut VecDeque<std::thread::JoinHandle<Result<()>>>,
) -> Result<()> {
while let Some(handle) = flush_handles.pop_front() {
match handle.join() {
Ok(Ok(())) => {}
Ok(Err(e)) => return Err(e),
Err(e) => {
return Err(anyhow::anyhow!(
"Block storage flush thread panicked at IBD completion: {:?}",
e
));
}
}
}
if skip_storage {
pending_blocks.clear();
*pending_storage_bytes = 0;
return Ok(());
}
if !pending_blocks.is_empty() {
info!(
"Flushing {} deferred blocks before IBD shutdown",
pending_blocks.len()
);
parallel_ibd.flush_pending_blocks_with_opts(
blockstore,
Some(storage),
pending_blocks,
IbdBlockFlushOpts::shutdown_sync(),
)?;
}
Ok(())
}