use super::*;
fn fresh_cap(initial: usize) -> Arc<AtomicUsize> {
Arc::new(AtomicUsize::new(initial))
}
fn fresh_last_adapt_zero() -> Arc<AtomicU64> {
Arc::new(AtomicU64::new(0))
}
#[serial_test::serial(ibd)]
#[test]
fn classify_binder_supply_vs_engine() {
assert_eq!(
classify_ibd_binder(0, 12, 0, 0, Some(90), PressureLevel::None, false),
"SUPPLY_TIP_HOLE_STAGED",
"R-302: feeder=0 holes>0 contig=0 await_ms=0 is staged (in our bridge), not absent"
);
assert_eq!(
classify_ibd_binder(0, 0, 0, 250, None, PressureLevel::None, false),
"SUPPLY_TIP_HOLE_ABSENT",
"R-302: await_ms>=200 is H has not arrived"
);
assert_eq!(
classify_ibd_binder(0, 0, 0, 0, None, PressureLevel::None, false),
"SUPPLY_EMPTY_TIP"
);
assert_eq!(
classify_ibd_binder(0, 0, 8, 0, None, PressureLevel::None, false),
"SUPPLY_FEEDER_STARVE",
"no gd sample → still classic starve"
);
assert_eq!(
classify_ibd_binder(0, 0, 66, 0, Some(29), PressureLevel::None, false),
"PIPE_DRAINED",
"H3 C3: feeder=0 + healthy gd + contig ≠ supply starve"
);
assert_eq!(
classify_ibd_binder(64, 0, 16, 0, Some(80), PressureLevel::None, false),
"ENGINE_OR_SCRIPTS"
);
assert_eq!(
classify_ibd_binder(64, 0, 16, 0, None, PressureLevel::Emergency, false),
"ENGINE_PRESSURE"
);
assert_eq!(
classify_ibd_binder(4, 0, 0, 0, Some(400), PressureLevel::None, false),
"SUPPLY_GD_SLOW"
);
}
#[serial_test::serial(ibd)]
#[test]
fn adapt_always_runs_with_positive_nominal() {
let cap = fresh_cap(5_000_000);
let last = fresh_last_adapt_zero();
adapt_max_pending_ops_tick(&cap, 5_000_000, PressureLevel::Emergency, 5_000_000, &last);
assert!(cap.load(Ordering::Relaxed) < 5_000_000);
}
#[serial_test::serial(ibd)]
#[test]
fn adapt_emergency_halves_cap() {
let cap = fresh_cap(8_000_000);
let last = fresh_last_adapt_zero();
adapt_max_pending_ops_tick(&cap, 8_000_000, PressureLevel::Emergency, 8_000_000, &last);
let new = cap.load(Ordering::Relaxed);
assert!(new < 8_000_000, "Emergency must shrink");
assert!(new >= 100_000, "Emergency must respect 100k floor");
}
#[serial_test::serial(ibd)]
#[test]
fn pipeline_depth_pressure_and_engine_append_throttle() {
use crate::storage::ibd_engine::memory_age::{
bump_append_stats_detailed_for_test, reset_append_diagnostics_for_test,
set_append_window_baseline_for_test,
};
reset_append_diagnostics_for_test();
let _tip = super::super::tip_stage::test_tip_atomics_lock();
super::super::tip_stage::test_reset_tip_stage();
super::super::tip_stage::test_reset_getdata_body_ewma();
super::memory::test_seed_ibd_rss_anon_mb(0);
assert_eq!(pipeline_depth_for_pressure(PressureLevel::None, 32), 32);
assert_eq!(pipeline_depth_for_pressure(PressureLevel::Emergency, 32), 8);
assert_eq!(pipeline_depth_for_pressure(PressureLevel::Critical, 32), 16);
assert_eq!(pipeline_depth_for_pressure(PressureLevel::Elevated, 32), 24);
super::super::tip_stage::publish_wan_body_tip(100);
super::super::tip_stage::mark_needed(200);
super::super::tip_stage::test_seed_getdata_body_ewma(40, 32);
assert_eq!(
pipeline_depth_for_pressure(PressureLevel::Elevated, 32),
24,
"Elevated must not soft-cap (C2 tip30 regress)"
);
assert_eq!(
pipeline_depth_for_pressure(PressureLevel::Critical, 32),
24,
"peak: raw Critical + tip-crawl healthy holds at Elevated depth"
);
assert_eq!(
pipeline_depth_for_pressure(PressureLevel::Emergency, 32),
8,
"Emergency depth unchanged (real reclaim)"
);
assert_eq!(
engine_pressure_poll_interval(PressureLevel::Critical),
16,
"tip-crawl healthy supply must poll Critical like Elevated"
);
super::memory::test_seed_ibd_rss_anon_mb(8377);
assert_eq!(
pipeline_depth_for_pressure(PressureLevel::Elevated, 32),
24,
"peak: Elevated stays 24 at anon 8G (r28 hold off)"
);
assert_eq!(
engine_pressure_poll_interval(PressureLevel::Elevated),
16,
"must not remap poll — that is r26"
);
assert_eq!(
pipeline_depth_for_pressure(PressureLevel::Critical, 32),
24,
"must not lift Critical off Land E 24"
);
assert_eq!(
pipeline_depth_for_pressure(PressureLevel::Emergency, 32),
8,
"must not hold Emergency"
);
super::memory::test_seed_ibd_rss_anon_mb(17174);
assert_eq!(
pipeline_depth_for_pressure(PressureLevel::Elevated, 32),
24,
"real Elevated anon 17G stays depth 24"
);
super::memory::test_seed_ibd_rss_anon_mb(8377);
super::super::tip_stage::test_reset_getdata_body_ewma();
super::super::tip_stage::publish_wan_body_tip(100);
super::super::tip_stage::mark_needed(200);
assert_eq!(
pipeline_depth_for_pressure(PressureLevel::Elevated, 32),
24,
"unhealthy supply still Elevated 24"
);
super::memory::test_seed_ibd_rss_anon_mb(0);
super::super::tip_stage::test_reset_tip_stage();
super::super::tip_stage::test_reset_getdata_body_ewma();
reset_append_diagnostics_for_test();
bump_append_stats_detailed_for_test(70_000, 35_000, 0);
assert_eq!(pipeline_depth_for_engine_append(16), 16);
reset_append_diagnostics_for_test();
bump_append_stats_detailed_for_test(90_000, 140_600, 140_600);
set_append_window_baseline_for_test(90_000, 140_600);
bump_append_stats_detailed_for_test(0, 256, 256);
assert_eq!(pipeline_depth_for_engine_append(32), 1);
}
#[serial_test::serial(ibd)]
#[test]
fn engine_pressure_poll_interval_tightens_with_pressure() {
assert_eq!(engine_pressure_poll_interval(PressureLevel::None), 32);
assert_eq!(engine_pressure_poll_interval(PressureLevel::Emergency), 1);
assert_eq!(engine_pressure_poll_interval(PressureLevel::Critical), 4);
}
#[serial_test::serial(ibd)]
#[test]
fn adapt_critical_multiplies_by_three_quarters() {
let cap = fresh_cap(8_000_000);
let last = fresh_last_adapt_zero();
adapt_max_pending_ops_tick(&cap, 8_000_000, PressureLevel::Critical, 8_000_000, &last);
let new = cap.load(Ordering::Relaxed);
assert!(new < 8_000_000);
assert!(
new >= 8_000_000 / 8,
"Critical must respect nominal/8 floor"
);
}
#[serial_test::serial(ibd)]
#[test]
fn adapt_elevated_is_hold() {
let cap = fresh_cap(8_000_000);
let last = fresh_last_adapt_zero();
adapt_max_pending_ops_tick(&cap, 8_000_000, PressureLevel::Elevated, 8_000_000, &last);
assert_eq!(cap.load(Ordering::Relaxed), 8_000_000);
}
#[serial_test::serial(ibd)]
#[test]
fn adapt_none_grows_when_drain_keeps_up() {
let nominal = 8_000_000;
let cap = fresh_cap(nominal);
let last = fresh_last_adapt_zero();
adapt_max_pending_ops_tick(&cap, nominal, PressureLevel::None, 100_000, &last);
let new = cap.load(Ordering::Relaxed);
assert!(new > nominal, "None + drain-ahead must grow cap");
let ceiling = nominal.saturating_mul(11).saturating_div(10);
assert!(
new <= ceiling,
"Must respect 1.1× nominal ceiling (got {new}, ceiling {ceiling})"
);
}
#[serial_test::serial(ibd)]
#[test]
fn adapt_none_holds_when_pending_full() {
let cap = fresh_cap(8_000_000);
let last = fresh_last_adapt_zero();
adapt_max_pending_ops_tick(&cap, 8_000_000, PressureLevel::None, 7_000_000, &last);
assert_eq!(cap.load(Ordering::Relaxed), 8_000_000);
}
#[serial_test::serial(ibd)]
#[test]
fn adapt_throttle_skips_recent_calls() {
let cap = fresh_cap(8_000_000);
let now_ms = crate::utils::time::current_timestamp_millis();
let last = Arc::new(AtomicU64::new(now_ms));
adapt_max_pending_ops_tick(&cap, 8_000_000, PressureLevel::Emergency, 8_000_000, &last);
assert_eq!(cap.load(Ordering::Relaxed), 8_000_000);
}
#[serial_test::serial(ibd)]
#[test]
fn adapt_emergency_respects_floor_under_repeat() {
let nominal = 8_000_000;
let cap = fresh_cap(nominal);
for _ in 0..50 {
let last = fresh_last_adapt_zero();
adapt_max_pending_ops_tick(&cap, nominal, PressureLevel::Emergency, nominal, &last);
}
let final_cap = cap.load(Ordering::Relaxed);
let expected_floor = (nominal / 2).max(1_000_000);
assert!(
final_cap >= expected_floor,
"must respect floor {expected_floor} (got {final_cap})",
);
assert!(
final_cap <= nominal,
"must not exceed nominal (got {final_cap} for nominal {nominal})",
);
}
#[serial_test::serial(ibd)]
#[test]
fn join_all_utxo_flush_handles_releases_mutex_before_join() {
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread;
use std::time::{Duration, Instant};
let started = Arc::new(AtomicBool::new(false));
let release = Arc::new(AtomicBool::new(false));
let started_join = Arc::clone(&started);
let release_join = Arc::clone(&release);
let slow = thread::spawn(move || {
started_join.store(true, Ordering::Release);
while !release_join.load(Ordering::Acquire) {
thread::sleep(Duration::from_millis(1));
}
Ok(blvm_muhash::MuHash3072::new())
});
let utxo_flush_handles = Arc::new(Mutex::new(VecDeque::new()));
utxo_flush_handles.lock().push_back(slow);
let handles_for_join = Arc::clone(&utxo_flush_handles);
let joiner = thread::spawn(move || join_all_utxo_flush_handles(&handles_for_join, "test"));
let wait_start = Instant::now();
while !started.load(Ordering::Acquire) {
assert!(
wait_start.elapsed() < Duration::from_secs(2),
"slow flush worker did not start"
);
thread::sleep(Duration::from_millis(1));
}
let lock_start = Instant::now();
{
let _guard = utxo_flush_handles.lock();
}
assert!(
lock_start.elapsed() < Duration::from_millis(500),
"utxo_flush_handles mutex still held during join"
);
release.store(true, Ordering::Release);
joiner.join().expect("join thread").expect("join flushes");
assert!(utxo_flush_handles.lock().is_empty());
}
#[test]
fn r361_deferred_dropper_frees_last_refs() {
use blvm_consensus::{Block, BlockHeader};
let mk = || {
Arc::new(Block {
header: BlockHeader {
version: 4,
..Default::default()
},
transactions: Vec::new().into(),
})
};
let block = mk();
let weak = Arc::downgrade(&block);
let mut dropper = DeferredDropper::spawn();
assert!(dropper.send(DeferredDrop {
block,
witnesses: Arc::new(Vec::new()),
undo_log: Some(blvm_consensus::reorganization::BlockUndoLog::new()),
}));
dropper.close_and_join();
assert!(
weak.upgrade().is_none(),
"dropper must have freed the block"
);
dropper.close_and_join();
let block2 = mk();
let weak2 = Arc::downgrade(&block2);
assert!(!dropper.send(DeferredDrop {
block: block2,
witnesses: Arc::new(Vec::new()),
undo_log: None,
}));
assert!(weak2.upgrade().is_none());
let off = DeferredDropper::disabled();
let block3 = mk();
let weak3 = Arc::downgrade(&block3);
assert!(!off.send(DeferredDrop {
block: block3,
witnesses: Arc::new(Vec::new()),
undo_log: None,
}));
assert!(weak3.upgrade().is_none(), "disabled dropper frees inline");
}
#[test]
fn r362_tip_syncer_coalesces_and_never_drops_the_last_tip() {
use std::sync::atomic::{AtomicU64, Ordering};
let calls = Arc::new(AtomicU64::new(0));
let max_seen = Arc::new(AtomicU64::new(0));
let (c, m) = (calls.clone(), max_seen.clone());
let mut syncer = TipSyncer::spawn_with(move |h| {
std::thread::sleep(std::time::Duration::from_millis(5));
c.fetch_add(1, Ordering::SeqCst);
m.fetch_max(h, Ordering::SeqCst);
Ok(())
});
for h in 1..=200u64 {
assert!(syncer.request(h * 1000), "request {} must be covered", h);
}
syncer.close_and_join();
let n = calls.load(Ordering::SeqCst);
assert!(n >= 1, "at least one sync ran");
assert!(n < 200, "requests must coalesce (ran {n})");
assert!(
max_seen.load(Ordering::SeqCst) >= 1000,
"a real tip was synced"
);
assert!(!syncer.request(201_000));
syncer.close_and_join();
let inline = TipSyncer::inline();
assert!(!inline.request(1000));
let mut failing = TipSyncer::spawn_with(|_h| Err(anyhow::anyhow!("disk gone")));
assert!(failing.request(5000));
failing.close_and_join();
}
#[test]
fn r365_mtp_window_tracks_dispatch_not_drain() {
use blvm_protocol::bip113::get_median_time_past;
use std::collections::VecDeque;
let ts = |h: u64| 1_400_000_000u64 + h * 600;
let hdr = |h: u64| {
Arc::new(BlockHeader {
version: 4,
timestamp: ts(h),
..Default::default()
})
};
let start: u64 = 100;
let mut window: VecDeque<Arc<BlockHeader>> = VecDeque::with_capacity(12);
for h in start - 11..start {
push_dispatched_header(&mut window, hdr(h));
}
assert_eq!(window.len(), MTP_WINDOW_HEADERS);
let depth = 64usize;
let mut drained_push_window: VecDeque<Arc<BlockHeader>> = window.clone();
let mut pending: VecDeque<Arc<BlockHeader>> = VecDeque::new();
for h in start..start + 500 {
let snap: Vec<Arc<BlockHeader>> = window.iter().cloned().collect();
assert_eq!(snap.len(), MTP_WINDOW_HEADERS, "h={h}");
assert_eq!(
snap.last().unwrap().timestamp,
ts(h - 1),
"last header must be h-1 at h={h}"
);
assert_eq!(
snap[0].timestamp,
ts(h - 11),
"first header must be h-11 at h={h}"
);
assert_eq!(
get_median_time_past(&snap),
ts(h - 6),
"MTP(h-1) is the 6th of 11 at h={h}"
);
push_dispatched_header(&mut window, hdr(h));
let stale_snap: Vec<Arc<BlockHeader>> = drained_push_window.iter().cloned().collect();
pending.push_back(hdr(h));
if pending.len() > depth {
push_dispatched_header(&mut drained_push_window, pending.pop_front().unwrap());
}
if h >= start + depth as u64 {
let stale = get_median_time_past(&stale_snap);
assert!(
stale < ts(h - 6),
"drain-time window must be strictly stale at depth {depth} (h={h}): {stale} vs {}",
ts(h - 6)
);
let lock_time = ts(h - 6) - 1;
assert!(lock_time >= stale && lock_time < ts(h - 6));
}
}
}
#[test]
fn r360_in_order_jobs_reassembles_prep_pool_output() {
let mut q: InOrderJobs<&'static str> = InOrderJobs::new(100);
assert!(q.pop_ready().is_none());
q.push(102, "c").unwrap();
q.push(101, "b").unwrap();
assert!(q.pop_ready().is_none(), "100 not yet arrived");
assert_eq!(q.pending_len(), 2);
q.push(100, "a").unwrap();
let mut out = Vec::new();
while let Some(j) = q.pop_ready() {
out.push(j);
}
assert_eq!(out, vec!["a", "b", "c"]);
assert_eq!(q.next_height(), 103);
assert_eq!(q.pending_len(), 0);
assert_eq!(q.push(102, "dup"), Err((103, "dup")));
q.push(104, "e").unwrap();
assert!(q.pop_ready().is_none(), "103 missing blocks 104");
q.push(103, "d").unwrap();
assert_eq!(q.pop_ready(), Some("d"));
assert_eq!(q.pop_ready(), Some("e"));
}
#[test]
fn r360_engine_prep_deferred_matches_dispatch_prep() {
use blvm_consensus::{Block, BlockHeader, Transaction, TransactionOutput};
let tx = |value: i64| Transaction {
version: 1,
inputs: blvm_protocol::tx_inputs![],
outputs: blvm_protocol::tx_outputs![
TransactionOutput {
value,
script_pubkey: vec![0x51],
},
TransactionOutput {
value: value / 2,
script_pubkey: vec![0x52],
}
],
lock_time: 0,
};
let block = Block {
header: BlockHeader {
version: 4,
timestamp: 1_600_000_000,
..Default::default()
},
transactions: vec![tx(50_0000_0000), tx(25_0000_0000), tx(7)].into(),
};
let mut expect_ids = Vec::new();
crate::storage::disk_utxo::compute_tx_ids_only(&block, &mut expect_ids);
assert_eq!(expect_ids.len(), 3);
let expect_cache = blvm_consensus::utxo_overlay::build_block_output_utxo_cache(
&block,
expect_ids.as_slice(),
1_000,
);
let mut ids = Vec::new();
let cache = engine_prep_deferred(&block, 1_000, 912_683, &mut ids).expect("cache below AV");
assert_eq!(ids, expect_ids);
assert_eq!(cache.len(), expect_cache.len());
assert_eq!(cache.len(), 6);
for (op, u) in expect_cache.iter() {
let got = cache.get(op).expect("outpoint present");
assert_eq!(got.value, u.value);
assert_eq!(got.script_pubkey, u.script_pubkey);
}
let mut ids_hi = Vec::new();
assert!(engine_prep_deferred(&block, 912_683, 912_683, &mut ids_hi).is_none());
assert_eq!(ids_hi, expect_ids);
let mut pre = expect_ids.clone();
pre.reverse();
let _ = engine_prep_deferred(&block, 1_000, 912_683, &mut pre);
assert_ne!(pre, expect_ids);
}
#[serial_test::serial(ibd)]
#[test]
fn join_all_utxo_flush_handles_empty_queue_is_noop() {
let utxo_flush_handles = Arc::new(Mutex::new(VecDeque::new()));
let combined = join_all_utxo_flush_handles(&utxo_flush_handles, "test").expect("empty join");
assert_eq!(
combined.finalize(),
blvm_muhash::MuHash3072::new().finalize()
);
}