use super::disk_index::DiskIndex;
use super::memory_age::{MemoryAge, Pin};
use super::memory_run::MemoryRun;
use super::types::{OutputId, OutputKV};
use std::cell::Cell;
use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::{AtomicI32, AtomicU64, AtomicUsize, Ordering};
use std::time::Instant;
thread_local! {
static LAST_QUERY_AGES_MS: Cell<u64> = const { Cell::new(0) };
static LAST_QUERY_DISK_MS: Cell<u64> = const { Cell::new(0) };
static LAST_DISK_PREADS: Cell<u64> = const { Cell::new(0) };
static LAST_DISK_PREAD_KB: Cell<u64> = const { Cell::new(0) };
static LAST_DISK_MAX_PREAD_KB: Cell<u64> = const { Cell::new(0) };
static LAST_DISK_CANDS: Cell<u64> = const { Cell::new(0) };
static LAST_DISK_SEGS: Cell<u64> = const { Cell::new(0) };
}
pub fn take_last_query_split_ms() -> (u64, u64) {
(
LAST_QUERY_AGES_MS.with(Cell::get),
LAST_QUERY_DISK_MS.with(Cell::get),
)
}
pub fn take_last_disk_io_stats() -> (u64, u64, u64, u64, u64) {
(
LAST_DISK_PREADS.with(Cell::get),
LAST_DISK_PREAD_KB.with(Cell::get),
LAST_DISK_MAX_PREAD_KB.with(Cell::get),
LAST_DISK_CANDS.with(Cell::get),
LAST_DISK_SEGS.with(Cell::get),
)
}
pub(crate) static COMPACTER_INFLIGHT_BYTES: AtomicU64 = AtomicU64::new(0);
#[cfg(test)]
extern crate tempfile;
const K_AGES: usize = 7;
const K_MUTABLE_AGES: usize = 3;
const K_FAN_IN: usize = 8;
const K_COMPACTER_THREADS: usize = 5;
fn spill_merge_take_from_env() -> usize {
static CACHED: std::sync::OnceLock<usize> = std::sync::OnceLock::new();
*CACHED.get_or_init(|| {
std::env::var("BLVM_IBD_SPILL_MERGE_TAKE")
.ok()
.and_then(|s| s.parse().ok())
.filter(|&n| n <= K_FAN_IN)
.unwrap_or(0)
})
}
fn elevated_no_demote_from_env() -> bool {
matches!(
std::env::var("BLVM_IBD_ELEVATED_NO_DEMOTE")
.ok()
.as_deref()
.map(str::trim),
Some("1") | Some("true") | Some("yes") | Some("on")
)
}
fn critical_no_demote_from_env() -> bool {
matches!(
std::env::var("BLVM_IBD_CRITICAL_NO_DEMOTE")
.ok()
.as_deref()
.map(str::trim),
Some("1") | Some("true") | Some("yes") | Some("on")
)
}
fn critical_soft_demote_from_env() -> bool {
matches!(
std::env::var("BLVM_IBD_CRITICAL_SOFT_DEMOTE")
.ok()
.as_deref()
.map(str::trim),
Some("1") | Some("true") | Some("yes") | Some("on")
)
}
fn hot_pin_keep_on_critical_from_env() -> bool {
matches!(
std::env::var("BLVM_IBD_HOT_PIN_KEEP_ON_CRITICAL")
.ok()
.as_deref()
.map(str::trim),
Some("1") | Some("true") | Some("yes") | Some("on")
)
}
fn hot_pin_emergency_hold_ms_from_env() -> u64 {
std::env::var("BLVM_IBD_HOT_PIN_EMERGENCY_HOLD_MS")
.ok()
.and_then(|s| s.trim().parse().ok())
.unwrap_or(0)
}
fn tip_resident_from_env() -> bool {
matches!(
std::env::var("BLVM_IBD_TIP_RESIDENT")
.ok()
.as_deref()
.map(str::trim),
Some("1") | Some("true") | Some("yes") | Some("on")
)
}
fn oldest_accumulate_from_env() -> bool {
matches!(
std::env::var("BLVM_IBD_OLDEST_ACCUMULATE")
.ok()
.as_deref()
.map(str::trim),
Some("1") | Some("true") | Some("yes") | Some("on")
)
}
fn critical_demote_hold_ms_from_env() -> u64 {
std::env::var("BLVM_IBD_CRITICAL_DEMOTE_HOLD_MS")
.ok()
.and_then(|s| s.trim().parse().ok())
.unwrap_or(0)
}
fn unix_now_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
fn force_merge_spill_to_disk(source_bytes: u64, merged_len: usize) -> bool {
let mb_thr: u64 = std::env::var("BLVM_IBD_FORCE_DISK_MERGE_MB")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(1536);
let entry_thr: usize = std::env::var("BLVM_IBD_FORCE_DISK_MERGE_ENTRIES")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(6_000_000);
(mb_thr > 0 && source_bytes >= mb_thr.saturating_mul(1024 * 1024))
|| (entry_thr > 0 && merged_len >= entry_thr)
}
fn choose_eviction_age(avail_mb: u64) -> usize {
if let Ok(s) = std::env::var("BLVM_IBD_ENGINE_EVICTION_AGE") {
if let Ok(n) = s.trim().parse::<usize>() {
let clamped = n.clamp(K_MUTABLE_AGES, K_AGES);
tracing::info!(
"UTXO engine: eviction age = {} (from BLVM_IBD_ENGINE_EVICTION_AGE)",
clamped
);
return clamped;
}
}
let total_mb = proc_mem_total_mb().unwrap_or(avail_mb);
let age = if avail_mb >= 65 * 1024 {
5
} else if avail_mb >= 20 * 1024 {
4
} else {
3
};
tracing::info!(
"UTXO engine: eviction age = {} (auto-detected: {:.1} GiB physical RAM, {:.1} GiB available)",
age,
total_mb as f64 / 1024.0,
avail_mb as f64 / 1024.0,
);
age
}
fn proc_mem_total_mb() -> Option<u64> {
let content = std::fs::read_to_string("/proc/meminfo").ok()?;
for line in content.lines() {
if line.starts_with("MemTotal:") {
let kb: u64 = line.split_whitespace().nth(1)?.parse().ok()?;
return Some(kb / 1024);
}
}
None
}
struct Compacter {
tx: crossbeam_channel::Sender<usize>,
eviction_age_live: Arc<AtomicUsize>,
_threads: Vec<std::thread::JoinHandle<()>>,
}
impl Compacter {
fn start(
ages: Arc<[MemoryAge; K_AGES]>,
disk_index: Arc<DiskIndex>,
boot_eviction_age: usize,
) -> Self {
let eviction_age_live = Arc::new(AtomicUsize::new(boot_eviction_age));
let (tx, rx) = crossbeam_channel::unbounded::<usize>();
let mut threads = Vec::with_capacity(K_COMPACTER_THREADS);
for _ in 0..K_COMPACTER_THREADS {
let rx = rx.clone();
let tx = tx.clone();
let ages = Arc::clone(&ages);
let disk_index = Arc::clone(&disk_index);
let eviction_age_live = Arc::clone(&eviction_age_live);
let handle = std::thread::Builder::new()
.name("utxo-compacter".to_string())
.spawn(move || {
while let Ok(age_idx) = rx.recv() {
if age_idx == usize::MAX {
break; }
let eviction_age = eviction_age_live.load(Ordering::Acquire);
run_merge_for_age(&ages, age_idx, &disk_index, &tx, eviction_age);
}
})
.expect("spawn compacter thread");
threads.push(handle);
}
Self {
tx,
eviction_age_live,
_threads: threads,
}
}
fn enqueue(&self, age_idx: usize) {
let _ = self.tx.try_send(age_idx);
}
fn shutdown(&self) {
for _ in 0..K_COMPACTER_THREADS {
let _ = self.tx.send(usize::MAX);
}
}
}
fn run_merge_for_age(
ages: &[MemoryAge; K_AGES],
age_idx: usize,
disk_index: &Arc<DiskIndex>,
tx: &crossbeam_channel::Sender<usize>,
eviction_age: usize,
) {
let age = &ages[age_idx];
let runs_before = age.run_count();
let Some(runs) = age.take_for_merge() else {
if runs_before >= age.merge_fan_in * 2 {
age.merge_ready_logged(); }
return;
};
let t_merge = std::time::Instant::now();
let source_bytes: u64 = runs.iter().map(|r| r.mem_bytes() as u64).sum();
COMPACTER_INFLIGHT_BYTES.fetch_add(source_bytes, Ordering::Relaxed);
let merged = MemoryRun::merge(&runs);
COMPACTER_INFLIGHT_BYTES.fetch_sub(source_bytes, Ordering::Relaxed);
let merge_ms = t_merge.elapsed().as_millis() as u64;
let max_h = runs
.iter()
.map(|r| r.height_range().1)
.max()
.unwrap_or(i32::MIN);
let merged_len = merged.len();
let emergency_floor = eviction_age <= K_MUTABLE_AGES;
let tip_hold = tip_resident_from_env() && !emergency_floor;
let oldest_acc = oldest_accumulate_from_env() && !emergency_floor;
let force_disk = !tip_hold && force_merge_spill_to_disk(source_bytes, merged_len);
let effective_eviction = if tip_hold && !oldest_acc {
K_AGES
} else {
eviction_age
};
let would_leave_ram = (age_idx + 1) >= effective_eviction || force_disk;
let fold_oldest = oldest_acc && would_leave_ram && !force_disk && !merged.is_empty();
let to_disk = !merged.is_empty() && would_leave_ram && !fold_oldest;
if !merged.is_empty() {
let next_idx = age_idx + 1;
if fold_oldest {
let dest = effective_eviction.saturating_sub(1).min(K_AGES - 1);
ages[dest].push_frozen_run(Arc::new(merged));
if ages[dest].merge_ready() {
let _ = tx.send(dest);
}
} else if next_idx < effective_eviction && !force_disk {
ages[next_idx].push_frozen_run(Arc::new(merged));
if ages[next_idx].merge_ready() {
let _ = tx.send(next_idx);
}
} else {
if force_disk && (age_idx + 1) < eviction_age {
tracing::warn!(
"UTXO compacter: age[{}] force-spill to disk (source_mb={} entries={}) — avoid merge RAM spike",
age_idx,
source_bytes / (1024 * 1024),
merged_len
);
}
if let Err(e) = disk_index.push_run_no_compact(merged) {
tracing::error!("UTXO engine: disk eviction failed — data may be lost: {e}");
}
}
}
let runs_before_complete = age.run_count();
age.complete_merge(max_h, &runs);
let runs_after = age.run_count();
drop(runs);
if to_disk || merge_ms >= 500 {
tracing::info!(
"UTXO compacter: age[{}] merged {}→{} runs, out_entries={}, merge_ms={}, to_disk={}",
age_idx,
runs_before,
runs_after,
merged_len,
merge_ms,
to_disk,
);
} else {
tracing::debug!(
"UTXO compacter: age[{}] merged {}→{} runs, out_entries={}, merge_ms={}",
age_idx,
runs_before,
runs_after,
merged_len,
merge_ms,
);
}
if to_disk {
disk_index.compact_oldest_async();
}
if age.merge_ready() {
let _ = tx.send(age_idx);
}
}
pub struct UtxoIndex {
ages: Arc<[MemoryAge; K_AGES]>,
compacter: Compacter,
disk_index: Arc<DiskIndex>,
contiguous_length: AtomicI32,
boot_eviction_age: usize,
critical_entered_ms: AtomicU64,
emergency_pin_entered_ms: AtomicU64,
}
impl UtxoIndex {
pub fn open(seg_dir: &Path, avail_mb: u64) -> anyhow::Result<Self> {
let (disk_index, restored_cl) = DiskIndex::new(seg_dir)?;
Self::open_with_disk(Arc::new(disk_index), avail_mb, restored_cl, None)
}
pub(super) fn open_with_disk(
disk_index: Arc<DiskIndex>,
avail_mb: u64,
restored_cl: i32,
table_path: Option<&Path>,
) -> anyhow::Result<Self> {
let eviction_age = choose_eviction_age(avail_mb);
let index_epoch = Arc::new(AtomicU64::new(0));
let ages_raw: [MemoryAge; K_AGES] = std::array::from_fn(|i| {
let is_mutable = i < K_MUTABLE_AGES;
let enqueue = None; MemoryAge::new_with_hooks(
is_mutable,
K_FAN_IN,
enqueue,
Some(Arc::clone(&index_epoch)),
)
});
let ages = Arc::new(ages_raw);
let compacter = Compacter::start(Arc::clone(&ages), Arc::clone(&disk_index), eviction_age);
{
let take = spill_merge_take_from_env();
if take > 0 {
tracing::info!(
"UTXO engine: spill-tier early merge take={take} \
(BLVM_IBD_SPILL_MERGE_TAKE; 0 disables)"
);
}
for i in 0..K_AGES {
let is_spill = i + 1 >= eviction_age;
ages[i].set_spill_early_take(if is_spill { take } else { 0 });
}
}
let mut restored_cl = restored_cl;
if let Some(tp) = table_path {
if let Some(file_cl) = super::meta::read_contiguous_length_sidecar(tp) {
if restored_cl < 0 && file_cl >= 0 {
restored_cl = file_cl;
} else if file_cl > restored_cl {
tracing::warn!(
"UTXO engine: contiguous_length sidecar={file_cl} > segment max={restored_cl} \
— clamping to segment max (in-memory tail not persisted)"
);
}
}
}
let contiguous_length = AtomicI32::new(restored_cl);
if restored_cl >= 0 {
super::set_gc_fence(restored_cl);
tracing::info!(
"UTXO engine: restored contiguous_length={} (segments + sidecar)",
restored_cl,
);
}
Ok(Self {
ages,
compacter,
disk_index,
contiguous_length,
boot_eviction_age: eviction_age,
critical_entered_ms: AtomicU64::new(0),
emergency_pin_entered_ms: AtomicU64::new(0),
})
}
pub fn mem_bytes(&self) -> usize {
let age_bytes: usize = self.ages.iter().map(|a| a.mem_bytes()).sum();
let disk_bytes = self.disk_index.bloom_bytes_total();
age_bytes + disk_bytes
}
pub fn age_run_count(&self, age_idx: usize) -> usize {
if age_idx < K_AGES {
self.ages[age_idx].run_count()
} else {
0
}
}
pub fn age_is_merging(&self, age_idx: usize) -> bool {
if age_idx < K_AGES {
self.ages[age_idx].is_merging.load(Ordering::Relaxed)
} else {
false
}
}
pub fn disk_is_compacting(&self) -> bool {
self.disk_index.is_compacting()
}
pub fn spill_io_busy(&self) -> bool {
self.disk_index.spill_io_busy()
}
pub fn disk_segment_count(&self) -> usize {
self.disk_index.segment_count()
}
pub fn age_detail(&self) -> (Vec<(usize, u64)>, (usize, u64)) {
let ages: Vec<(usize, u64)> = self
.ages
.iter()
.map(|a| {
let guard = a.runs.read();
let bytes: usize = guard.iter().map(|r| r.mem_bytes()).sum();
(guard.len(), bytes as u64 / (1024 * 1024))
})
.collect();
let disk_segs = self.disk_index.segment_count();
let disk_bloom_mb = self.disk_index.bloom_bytes_total() as u64 / (1024 * 1024);
(ages, (disk_segs, disk_bloom_mb))
}
pub fn memory_pressure_tick(&self, level_u8: u8) {
let boot = self.boot_eviction_age;
let tip_res = tip_resident_from_env();
let oldest_acc = oldest_accumulate_from_env();
let target = match level_u8 {
3 if crate::node::parallel_ibd::tip_stage::tip_crawl_supply_healthy_now() => {
self.critical_entered_ms.store(0, Ordering::Relaxed);
boot
}
3 => {
self.critical_entered_ms.store(0, Ordering::Relaxed);
K_MUTABLE_AGES }
2 if tip_res
|| oldest_acc
|| critical_no_demote_from_env()
|| crate::node::parallel_ibd::tip_stage::tip_crawl_supply_healthy_now() =>
{
self.critical_entered_ms.store(0, Ordering::Relaxed);
boot
}
2 if critical_soft_demote_from_env() => {
self.critical_entered_ms.store(0, Ordering::Relaxed);
boot.saturating_sub(1).max(K_MUTABLE_AGES)
}
2 => {
let hold = critical_demote_hold_ms_from_env();
if hold == 0 {
self.critical_entered_ms.store(0, Ordering::Relaxed);
K_MUTABLE_AGES
} else {
let now = unix_now_ms();
let since = self.critical_entered_ms.load(Ordering::Relaxed);
let since = if since == 0 {
self.critical_entered_ms.store(now, Ordering::Relaxed);
now
} else {
since
};
if now.saturating_sub(since) < hold {
boot } else {
K_MUTABLE_AGES
}
}
}
1 if tip_res || oldest_acc || elevated_no_demote_from_env() => {
self.critical_entered_ms.store(0, Ordering::Relaxed);
boot
}
1 => {
self.critical_entered_ms.store(0, Ordering::Relaxed);
boot.saturating_sub(1).max(K_MUTABLE_AGES)
}
_ => {
self.critical_entered_ms.store(0, Ordering::Relaxed);
boot }
};
let live = &self.compacter.eviction_age_live;
let prev = live.load(Ordering::Relaxed);
if target != prev {
live.store(target, Ordering::Release);
let take = spill_merge_take_from_env();
for i in 0..K_AGES {
let is_spill = i + 1 >= target;
self.ages[i].set_spill_early_take(if is_spill { take } else { 0 });
}
if target < prev {
tracing::warn!(
"UTXO engine: memory pressure level {} — eviction age {} → {} (spilling index to disk)",
level_u8,
prev,
target
);
} else if level_u8 == 0 {
tracing::info!(
"UTXO engine: memory pressure cleared — eviction age restored to {}",
target
);
}
}
if level_u8 >= 2 {
let drop_pin = if level_u8 >= 3 {
let hold = hot_pin_emergency_hold_ms_from_env();
if hold == 0 {
self.emergency_pin_entered_ms.store(0, Ordering::Relaxed);
true
} else {
let now = unix_now_ms();
let since = self.emergency_pin_entered_ms.load(Ordering::Relaxed);
let since = if since == 0 {
self.emergency_pin_entered_ms.store(now, Ordering::Relaxed);
now
} else {
since
};
let sustained = now.saturating_sub(since) >= hold;
if !sustained {
tracing::info!(
"DiskSegment: hot-pin clear deferred (emergency hold {}ms, elapsed {}ms)",
hold,
now.saturating_sub(since)
);
}
sustained
}
} else {
self.emergency_pin_entered_ms.store(0, Ordering::Relaxed);
!hot_pin_keep_on_critical_from_env()
};
if drop_pin {
self.disk_index.clear_hot_pins_keep_seed();
}
for i in 0..K_AGES {
if self.ages[i].merge_ready() {
self.compacter.enqueue(i);
}
}
} else {
self.emergency_pin_entered_ms.store(0, Ordering::Relaxed);
}
}
#[cfg(test)]
pub fn new_for_test() -> Self {
let tmp = tempfile::tempdir().expect("tempdir");
let idx = Self::open(tmp.path(), 8 * 1024).expect("UtxoIndex::open"); std::mem::forget(tmp);
idx
}
pub fn append(&self, entries: Vec<OutputKV>, height: i32) -> Pin {
let pin = self.ages[0].pin_height(height);
self.ages[0].append(entries, height);
self.contiguous_length.fetch_max(height, Ordering::Relaxed);
let spill_hi = self
.compacter
.eviction_age_live
.load(Ordering::Relaxed)
.saturating_sub(1)
.min(K_AGES - 1);
for i in 0..=spill_hi {
if self.ages[i].merge_ready() {
self.compacter.enqueue(i);
}
}
pin
}
pub fn eviction_age_live(&self) -> usize {
self.compacter.eviction_age_live.load(Ordering::Relaxed)
}
pub fn lookup_key(&self, key: &[u8; 36]) -> Option<OutputId> {
for age in self.ages.iter() {
if let Some(id) = age.lookup_key(key, 0, i32::MAX) {
if id == super::types::OUTPUT_ID_DELETED {
return None;
}
return Some(id);
}
}
None
}
pub fn batch_query(&self, keys: &[[u8; 36]], ids: &mut [OutputId], before: i32) {
debug_assert_eq!(keys.len(), ids.len());
let t_ages = Instant::now();
for i in 0..K_AGES {
if !ids.contains(&OutputId::MAX) {
break;
}
self.ages[i].batch_query(keys, ids, 0, before);
}
let ages_ms = t_ages.elapsed().as_millis() as u64;
let t_disk = Instant::now();
if ids.contains(&OutputId::MAX) {
self.disk_index.batch_query(keys, ids, before);
}
let disk_ms = t_disk.elapsed().as_millis() as u64;
LAST_QUERY_AGES_MS.with(|c| c.set(ages_ms));
LAST_QUERY_DISK_MS.with(|c| c.set(disk_ms));
let (preads, pread_kb, max_kb, cands, segs) = super::disk_segment::take_disk_io_stats();
LAST_DISK_PREADS.with(|c| c.set(preads));
LAST_DISK_PREAD_KB.with(|c| c.set(pread_kb));
LAST_DISK_MAX_PREAD_KB.with(|c| c.set(max_kb));
LAST_DISK_CANDS.with(|c| c.set(cands));
LAST_DISK_SEGS.with(|c| c.set(segs));
}
pub fn wait_for_height(&self, height: i32) {
while self.contiguous_length.load(Ordering::Relaxed) < height {
std::thread::sleep(std::time::Duration::from_millis(1));
}
}
pub fn contiguous_length(&self) -> i32 {
self.contiguous_length.load(Ordering::Relaxed)
}
pub fn seed_checkpoint(&self, mut entries: Vec<OutputKV>, checkpoint_height: i32) {
if !entries.is_empty() {
entries.sort_unstable();
match self.disk_index.push_sorted_segment_owned(entries) {
Ok(()) => {
}
Err(e) => {
tracing::error!(
"seed_checkpoint: disk write failed ({e:#}) — UTXO seed incomplete"
);
}
}
}
self.contiguous_length
.store(checkpoint_height, Ordering::Release);
super::set_gc_fence(checkpoint_height);
}
pub fn alloc_seed_seg(&self) -> (usize, std::path::PathBuf) {
self.disk_index.alloc_seg()
}
pub fn finalize_seed(&self, seg: super::disk_segment::DiskSegment, checkpoint_height: i32) {
self.disk_index.register_seg(seg);
self.contiguous_length
.store(checkpoint_height, Ordering::Release);
super::set_gc_fence(checkpoint_height);
}
pub fn erase_since(&self, since: i32) {
for i in 0..K_MUTABLE_AGES {
self.ages[i].erase_since(since);
}
self.contiguous_length
.fetch_min(since - 1, Ordering::Relaxed);
}
pub fn scan_all_live(&self) -> Vec<OutputKV> {
let (disk_segs, mem_entries_pre) = {
let guard = self.disk_index.segments.read();
let disk = Arc::clone(&*guard);
let mut mem: Vec<OutputKV> = Vec::new();
for a in self.ages.iter().rev() {
a.collect_entries_into(&mut mem);
}
for run in self.disk_index.pending_spills.read().iter() {
mem.extend_from_slice(&run.entries);
}
(disk, mem)
};
if disk_segs.len() > 2 || mem_entries_pre.len() > 1_000_000 {
tracing::warn!(
"scan_all_live: materializing full index (segs={}, mem_entries={}) — \
tip-scale callers must use streaming checkpoint export (run_watermark_export)",
disk_segs.len(),
mem_entries_pre.len()
);
}
let mut all_entries: Vec<OutputKV> = mem_entries_pre;
for seg in disk_segs.iter() {
let entries = match seg.read_all_entries() {
Ok(e) => e,
Err(err) => {
tracing::warn!("scan_all_live: skipping segment {:?}: {err}", seg.path);
continue;
}
};
all_entries.extend_from_slice(&entries);
}
all_entries.sort_unstable();
let mut result: Vec<OutputKV> = Vec::new();
let mut i = 0;
while i < all_entries.len() {
let first = all_entries[i];
if first.is_add() {
result.push(first);
}
let key = first.key;
i += 1;
while i < all_entries.len() && all_entries[i].key == key {
i += 1;
}
}
result
}
pub fn scan_live_at_height(&self, max_height: i32) -> Vec<OutputKV> {
self.disk_index.compact_for_checkpoint_sync();
let (disk_segs, mut all_entries) = {
let guard = self.disk_index.segments.read();
let disk = Arc::clone(&*guard);
let mut mem: Vec<OutputKV> = Vec::new();
for a in self.ages.iter().rev() {
a.collect_entries_at_or_below_into(max_height, &mut mem);
}
for run in self.disk_index.pending_spills.read().iter() {
for e in &run.entries {
if e.height <= max_height {
mem.push(*e);
}
}
}
(disk, mem)
};
for seg in disk_segs.iter() {
let mut reader = seg.stream();
loop {
match reader.advance() {
Ok(Some(entry)) => {
if entry.height <= max_height {
all_entries.push(entry);
}
}
Ok(None) => break,
Err(err) => {
tracing::warn!(
"scan_live_at_height: read error on segment {:?}: {err}",
seg.path
);
break;
}
}
}
}
all_entries.sort_unstable();
let mut result: Vec<OutputKV> = Vec::new();
let mut i = 0;
while i < all_entries.len() {
let first = all_entries[i];
if first.is_add() {
result.push(first);
}
let key = first.key;
i += 1;
while i < all_entries.len() && all_entries[i].key == key {
i += 1;
}
}
result
}
pub fn iter_live_at_height(
&self,
max_height: i32,
) -> anyhow::Result<(CheckpointStream, u64, u64)> {
let t_compact = std::time::Instant::now();
self.disk_index.compact_for_checkpoint_sync();
let compact_ms = t_compact.elapsed().as_millis() as u64;
let t_scan_prep = std::time::Instant::now();
let (disk_segs, mut mem_entries) = {
let guard = self.disk_index.segments.read();
let disk = Arc::clone(&*guard);
let mut mem: Vec<OutputKV> = Vec::new();
for a in self.ages.iter().rev() {
a.collect_entries_at_or_below_into(max_height, &mut mem);
}
(disk, mem)
};
mem_entries.sort_unstable();
let mut readers: Vec<super::disk_segment::SegmentReader> =
disk_segs.iter().map(|seg| seg.stream()).collect();
let mut heads: Vec<Option<OutputKV>> = Vec::with_capacity(readers.len());
for reader in &mut readers {
heads.push(reader.advance()?);
}
let scan_prep_ms = t_scan_prep.elapsed().as_millis() as u64;
Ok((
CheckpointStream {
mem: mem_entries,
mem_pos: 0,
readers,
heads,
max_height,
last_key: None,
},
compact_ms,
scan_prep_ms,
))
}
pub fn compact_for_checkpoint_sync_with_sink<F>(
&self,
checkpoint_height: i32,
on_live: Option<F>,
) -> anyhow::Result<u64>
where
F: FnMut(OutputKV) -> anyhow::Result<()>,
{
let t_compact = std::time::Instant::now();
self.disk_index
.compact_for_checkpoint_sync_with_sink(checkpoint_height, on_live)?;
Ok(t_compact.elapsed().as_millis() as u64)
}
pub fn collect_memory_entries_at_or_below(&self, max_height: i32) -> Vec<OutputKV> {
let mut mem: Vec<OutputKV> = Vec::new();
for a in self.ages.iter().rev() {
a.collect_entries_at_or_below_into(max_height, &mut mem);
}
mem.sort_unstable();
mem
}
}
pub struct CheckpointStream {
mem: Vec<OutputKV>,
mem_pos: usize,
readers: Vec<super::disk_segment::SegmentReader>,
heads: Vec<Option<OutputKV>>,
max_height: i32,
last_key: Option<[u8; 36]>,
}
impl CheckpointStream {
pub fn next_live(&mut self) -> anyhow::Result<Option<OutputKV>> {
loop {
let entry = match self.pick_min()? {
Some(e) => e,
None => return Ok(None),
};
if entry.height > self.max_height {
continue;
}
if Some(entry.key) == self.last_key {
continue;
}
self.last_key = Some(entry.key);
if entry.is_add() {
return Ok(Some(entry));
}
}
}
fn pick_min(&mut self) -> anyhow::Result<Option<OutputKV>> {
let mut best: Option<OutputKV> = None;
let mut best_is_disk = false;
let mut best_disk_idx: usize = 0;
if let Some(&me) = self.mem.get(self.mem_pos) {
best = Some(me);
}
for i in 0..self.heads.len() {
if let Some(de) = self.heads[i] {
let take = match best {
None => true,
Some(cur) => de < cur,
};
if take {
best = Some(de);
best_is_disk = true;
best_disk_idx = i;
}
}
}
match best {
None => Ok(None),
Some(e) => {
if best_is_disk {
self.heads[best_disk_idx] = self.readers[best_disk_idx].advance()?;
} else {
self.mem_pos += 1;
}
Ok(Some(e))
}
}
}
}
impl Drop for UtxoIndex {
fn drop(&mut self) {
self.compacter.shutdown();
}
}
#[cfg(test)]
mod tests {
use super::super::types::OutputKV;
use super::*;
fn make_key(n: u8) -> [u8; 36] {
let mut k = [0u8; 36];
k[0] = n;
k
}
#[serial_test::serial(ibd)]
#[test]
fn w52_force_merge_spill_thresholds() {
assert!(!force_merge_spill_to_disk(1024 * 1024 * 1024, 1_000_000));
assert!(force_merge_spill_to_disk(1536 * 1024 * 1024, 1));
assert!(force_merge_spill_to_disk(1, 6_000_000));
}
#[serial_test::serial(ibd)]
#[test]
fn test_append_and_query() {
let idx = UtxoIndex::new_for_test();
let k = make_key(1);
let _pin = idx.append(vec![OutputKV::new_add(k, 100, 42)], 100);
assert_eq!(idx.lookup_key(&k), Some(42));
assert_eq!(idx.lookup_key(&make_key(2)), None);
}
#[serial_test::serial(ibd)]
#[test]
fn test_batch_query() {
let idx = UtxoIndex::new_for_test();
let k1 = make_key(1);
let k2 = make_key(2);
let _p1 = idx.append(vec![OutputKV::new_add(k1, 100, 10)], 100);
let _p2 = idx.append(vec![OutputKV::new_add(k2, 101, 20)], 101);
let mut ids = [OutputId::MAX; 2];
idx.batch_query(&[k1, k2], &mut ids, i32::MAX);
assert_eq!(ids[0], 10);
assert_eq!(ids[1], 20);
}
#[serial_test::serial(ibd)]
#[test]
fn test_contiguous_length() {
let idx = UtxoIndex::new_for_test();
assert_eq!(idx.contiguous_length(), -1);
let k = make_key(1);
let _pin = idx.append(vec![OutputKV::new_add(k, 50, 1)], 50);
assert_eq!(idx.contiguous_length(), 50);
}
#[serial_test::serial(ibd)]
#[test]
fn test_erase_since() {
let idx = UtxoIndex::new_for_test();
let k1 = make_key(1);
let k2 = make_key(2);
let _p1 = idx.append(vec![OutputKV::new_add(k1, 50, 1)], 50);
let _p2 = idx.append(vec![OutputKV::new_add(k2, 100, 2)], 100);
idx.erase_since(75);
assert_eq!(idx.lookup_key(&k1), Some(1));
assert_eq!(idx.lookup_key(&k2), None);
}
fn demote_env_lock() -> std::sync::MutexGuard<'static, ()> {
static LOCK: std::sync::OnceLock<std::sync::Mutex<()>> = std::sync::OnceLock::new();
LOCK.get_or_init(|| std::sync::Mutex::new(()))
.lock()
.unwrap_or_else(|e| e.into_inner())
}
#[serial_test::serial(ibd)]
#[test]
fn critical_pressure_lowers_eviction_age_to_mutable_floor() {
use std::sync::atomic::Ordering;
let _guard = demote_env_lock();
unsafe {
std::env::remove_var("BLVM_IBD_ELEVATED_NO_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_NO_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_SOFT_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_DEMOTE_HOLD_MS");
}
let tmp = tempfile::tempdir().expect("tempdir");
let (disk_index, restored_cl) = DiskIndex::new(tmp.path()).expect("DiskIndex::new");
let idx = UtxoIndex::open_with_disk(Arc::new(disk_index), 24 * 1024, restored_cl, None)
.expect("open_with_disk");
assert_eq!(idx.boot_eviction_age, 4);
idx.memory_pressure_tick(2);
assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
K_MUTABLE_AGES,
);
std::mem::forget(tmp);
}
#[serial_test::serial(ibd)]
#[test]
fn tip_crawl_supply_healthy_holds_boot_under_critical() {
use std::sync::atomic::Ordering;
let _guard = demote_env_lock();
let _tip = crate::node::parallel_ibd::tip_stage::test_tip_atomics_lock();
crate::node::parallel_ibd::tip_stage::test_reset_tip_stage();
unsafe {
std::env::remove_var("BLVM_IBD_ELEVATED_NO_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_NO_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_SOFT_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_DEMOTE_HOLD_MS");
}
crate::node::parallel_ibd::tip_stage::publish_wan_body_tip(100);
crate::node::parallel_ibd::tip_stage::mark_needed(200);
crate::node::parallel_ibd::tip_stage::test_seed_getdata_body_ewma(40, 32);
let tmp = tempfile::tempdir().expect("tempdir");
let (disk_index, restored_cl) = DiskIndex::new(tmp.path()).expect("DiskIndex::new");
let idx = UtxoIndex::open_with_disk(Arc::new(disk_index), 24 * 1024, restored_cl, None)
.expect("open_with_disk");
let boot = idx.boot_eviction_age;
idx.memory_pressure_tick(2); assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
boot,
"tip-crawl healthy supply must not floor eviction age on Critical"
);
idx.memory_pressure_tick(3); assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
boot,
"KEEP C0: Emergency + healthy supply must not floor ages (view-double)"
);
crate::node::parallel_ibd::tip_stage::test_reset_getdata_body_ewma();
crate::node::parallel_ibd::tip_stage::publish_wan_body_tip(100);
crate::node::parallel_ibd::tip_stage::mark_needed(200);
idx.memory_pressure_tick(3); assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
K_MUTABLE_AGES,
"Emergency + unhealthy supply must still floor"
);
crate::node::parallel_ibd::tip_stage::test_reset_tip_stage();
crate::node::parallel_ibd::tip_stage::test_reset_getdata_body_ewma();
std::mem::forget(tmp);
}
#[serial_test::serial(ibd)]
#[test]
fn elevated_no_demote_keeps_boot_age_under_elevated() {
use std::sync::atomic::Ordering;
let _guard = demote_env_lock();
unsafe {
std::env::set_var("BLVM_IBD_ELEVATED_NO_DEMOTE", "1");
std::env::remove_var("BLVM_IBD_CRITICAL_NO_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_SOFT_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_DEMOTE_HOLD_MS");
}
let tmp = tempfile::tempdir().expect("tempdir");
let (disk_index, restored_cl) = DiskIndex::new(tmp.path()).expect("DiskIndex::new");
let idx = UtxoIndex::open_with_disk(Arc::new(disk_index), 24 * 1024, restored_cl, None)
.expect("open_with_disk");
let boot = idx.boot_eviction_age;
assert!(boot >= 4);
idx.memory_pressure_tick(1); assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
boot,
);
idx.memory_pressure_tick(2); assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
K_MUTABLE_AGES,
);
unsafe {
std::env::remove_var("BLVM_IBD_ELEVATED_NO_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_NO_DEMOTE");
}
std::mem::forget(tmp);
}
#[serial_test::serial(ibd)]
#[test]
fn critical_no_demote_keeps_boot_age_under_critical() {
use std::sync::atomic::Ordering;
let _guard = demote_env_lock();
unsafe {
std::env::set_var("BLVM_IBD_ELEVATED_NO_DEMOTE", "1");
std::env::set_var("BLVM_IBD_CRITICAL_NO_DEMOTE", "1");
std::env::remove_var("BLVM_IBD_CRITICAL_SOFT_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_DEMOTE_HOLD_MS");
}
let tmp = tempfile::tempdir().expect("tempdir");
let (disk_index, restored_cl) = DiskIndex::new(tmp.path()).expect("DiskIndex::new");
let idx = UtxoIndex::open_with_disk(Arc::new(disk_index), 24 * 1024, restored_cl, None)
.expect("open_with_disk");
let boot = idx.boot_eviction_age;
assert!(boot >= 4);
idx.memory_pressure_tick(2); assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
boot,
);
idx.memory_pressure_tick(3); assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
K_MUTABLE_AGES,
);
unsafe {
std::env::remove_var("BLVM_IBD_ELEVATED_NO_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_NO_DEMOTE");
}
std::mem::forget(tmp);
}
#[serial_test::serial(ibd)]
#[test]
fn critical_soft_demote_steps_one_from_boot() {
use std::sync::atomic::Ordering;
let _guard = demote_env_lock();
unsafe {
std::env::remove_var("BLVM_IBD_CRITICAL_NO_DEMOTE");
std::env::set_var("BLVM_IBD_CRITICAL_SOFT_DEMOTE", "1");
std::env::set_var("BLVM_IBD_ENGINE_EVICTION_AGE", "5");
}
let tmp = tempfile::tempdir().expect("tempdir");
let (disk_index, restored_cl) = DiskIndex::new(tmp.path()).expect("DiskIndex::new");
let idx = UtxoIndex::open_with_disk(Arc::new(disk_index), 24 * 1024, restored_cl, None)
.expect("open_with_disk");
assert_eq!(idx.boot_eviction_age, 5);
idx.memory_pressure_tick(2); assert_eq!(idx.compacter.eviction_age_live.load(Ordering::Relaxed), 4,);
idx.memory_pressure_tick(3); assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
K_MUTABLE_AGES,
);
unsafe {
std::env::remove_var("BLVM_IBD_CRITICAL_SOFT_DEMOTE");
std::env::remove_var("BLVM_IBD_ENGINE_EVICTION_AGE");
}
std::mem::forget(tmp);
}
#[serial_test::serial(ibd)]
#[test]
fn critical_demote_hold_defers_floor_then_demotes() {
use std::sync::atomic::Ordering;
let _guard = demote_env_lock();
unsafe {
std::env::remove_var("BLVM_IBD_CRITICAL_NO_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_SOFT_DEMOTE");
std::env::set_var("BLVM_IBD_CRITICAL_DEMOTE_HOLD_MS", "200");
std::env::set_var("BLVM_IBD_ENGINE_EVICTION_AGE", "4");
}
let tmp = tempfile::tempdir().expect("tempdir");
let (disk_index, restored_cl) = DiskIndex::new(tmp.path()).expect("DiskIndex::new");
let idx = UtxoIndex::open_with_disk(Arc::new(disk_index), 24 * 1024, restored_cl, None)
.expect("open_with_disk");
let boot = idx.boot_eviction_age;
assert_eq!(boot, 4);
idx.memory_pressure_tick(2); assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
boot,
);
std::thread::sleep(std::time::Duration::from_millis(250));
idx.memory_pressure_tick(2); assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
K_MUTABLE_AGES,
);
idx.memory_pressure_tick(0); assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
boot,
);
unsafe {
std::env::remove_var("BLVM_IBD_CRITICAL_DEMOTE_HOLD_MS");
std::env::remove_var("BLVM_IBD_ENGINE_EVICTION_AGE");
}
std::mem::forget(tmp);
}
#[serial_test::serial(ibd)]
#[test]
fn hot_pin_emergency_hold_defers_clear() {
use std::sync::atomic::Ordering;
let _guard = demote_env_lock();
unsafe {
std::env::set_var("BLVM_IBD_HOT_PIN", "1");
std::env::set_var("BLVM_IBD_HOT_PIN_KEEP_ON_CRITICAL", "1");
std::env::set_var("BLVM_IBD_HOT_PIN_EMERGENCY_HOLD_MS", "200");
std::env::set_var("BLVM_IBD_ENGINE_EVICTION_AGE", "4");
}
let tmp = tempfile::tempdir().expect("tempdir");
let (disk_index, restored_cl) = DiskIndex::new(tmp.path()).expect("DiskIndex::new");
let idx = UtxoIndex::open_with_disk(Arc::new(disk_index), 24 * 1024, restored_cl, None)
.expect("open_with_disk");
idx.memory_pressure_tick(3); assert!(
idx.emergency_pin_entered_ms.load(Ordering::Relaxed) > 0,
"emergency hold streak should start"
);
std::thread::sleep(std::time::Duration::from_millis(250));
idx.memory_pressure_tick(3); idx.memory_pressure_tick(0); assert_eq!(idx.emergency_pin_entered_ms.load(Ordering::Relaxed), 0);
unsafe {
std::env::remove_var("BLVM_IBD_HOT_PIN");
std::env::remove_var("BLVM_IBD_HOT_PIN_KEEP_ON_CRITICAL");
std::env::remove_var("BLVM_IBD_HOT_PIN_EMERGENCY_HOLD_MS");
std::env::remove_var("BLVM_IBD_ENGINE_EVICTION_AGE");
}
std::mem::forget(tmp);
}
#[serial_test::serial(ibd)]
#[test]
fn tip_resident_keeps_boot_under_critical_emergency_floors() {
unsafe {
std::env::set_var("BLVM_IBD_TIP_RESIDENT", "1");
std::env::set_var("BLVM_IBD_ENGINE_EVICTION_AGE", "5");
std::env::remove_var("BLVM_IBD_CRITICAL_NO_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_SOFT_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_DEMOTE_HOLD_MS");
std::env::remove_var("BLVM_IBD_ELEVATED_NO_DEMOTE");
std::env::remove_var("BLVM_IBD_OLDEST_ACCUMULATE");
}
let tmp = tempfile::tempdir().expect("tempdir");
let (disk_index, restored_cl) = DiskIndex::new(tmp.path()).expect("DiskIndex::new");
let idx = UtxoIndex::open_with_disk(Arc::new(disk_index), 24 * 1024, restored_cl, None)
.expect("open_with_disk");
assert_eq!(idx.boot_eviction_age, 5);
idx.memory_pressure_tick(1); assert_eq!(idx.compacter.eviction_age_live.load(Ordering::Relaxed), 5,);
idx.memory_pressure_tick(2); assert_eq!(idx.compacter.eviction_age_live.load(Ordering::Relaxed), 5,);
idx.memory_pressure_tick(3); assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
K_MUTABLE_AGES,
);
unsafe {
std::env::remove_var("BLVM_IBD_TIP_RESIDENT");
std::env::remove_var("BLVM_IBD_ENGINE_EVICTION_AGE");
}
std::mem::forget(tmp);
}
#[serial_test::serial(ibd)]
#[test]
fn oldest_accumulate_keeps_boot_under_critical_emergency_floors() {
unsafe {
std::env::set_var("BLVM_IBD_OLDEST_ACCUMULATE", "1");
std::env::set_var("BLVM_IBD_ENGINE_EVICTION_AGE", "4");
std::env::remove_var("BLVM_IBD_TIP_RESIDENT");
std::env::remove_var("BLVM_IBD_CRITICAL_NO_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_SOFT_DEMOTE");
std::env::remove_var("BLVM_IBD_CRITICAL_DEMOTE_HOLD_MS");
std::env::remove_var("BLVM_IBD_ELEVATED_NO_DEMOTE");
}
let tmp = tempfile::tempdir().expect("tempdir");
let (disk_index, restored_cl) = DiskIndex::new(tmp.path()).expect("DiskIndex::new");
let idx = UtxoIndex::open_with_disk(Arc::new(disk_index), 24 * 1024, restored_cl, None)
.expect("open_with_disk");
assert_eq!(idx.boot_eviction_age, 4);
idx.memory_pressure_tick(2);
assert_eq!(idx.compacter.eviction_age_live.load(Ordering::Relaxed), 4,);
idx.memory_pressure_tick(3);
assert_eq!(
idx.compacter.eviction_age_live.load(Ordering::Relaxed),
K_MUTABLE_AGES,
);
unsafe {
std::env::remove_var("BLVM_IBD_OLDEST_ACCUMULATE");
std::env::remove_var("BLVM_IBD_ENGINE_EVICTION_AGE");
}
std::mem::forget(tmp);
}
}