#[cfg(all(not(target_os = "windows"), feature = "mimalloc"))]
use libmimalloc_sys;
#[cfg(target_os = "linux")]
use std::io::Read;
use std::sync::atomic::{AtomicU8, AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
#[repr(u8)]
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub(crate) enum PressureLevel {
None = 0,
Elevated = 1,
Critical = 2,
Emergency = 3,
}
impl PressureLevel {
#[inline]
pub(crate) fn from_u8(v: u8) -> Self {
match v {
1 => Self::Elevated,
2 => Self::Critical,
3 => Self::Emergency,
_ => Self::None,
}
}
}
static IBD_PRESSURE_LEVEL: AtomicU8 = AtomicU8::new(0);
static IBD_RSS_ANON_MB: AtomicU64 = AtomicU64::new(0);
static WORKLOAD_CLASS_LATCH: AtomicU8 = AtomicU8::new(0xff);
#[inline]
pub(crate) fn publish_ibd_pressure(level: PressureLevel) {
IBD_PRESSURE_LEVEL.store(level as u8, Ordering::Relaxed);
}
#[inline]
pub(crate) fn publish_ibd_rss_anon_mb(mb: u64) {
IBD_RSS_ANON_MB.store(mb, Ordering::Relaxed);
}
#[inline]
pub(crate) fn ibd_rss_anon_mb_snapshot() -> u64 {
IBD_RSS_ANON_MB.load(Ordering::Relaxed)
}
#[cfg(test)]
pub(crate) fn test_seed_ibd_rss_anon_mb(mb: u64) {
publish_ibd_rss_anon_mb(mb);
}
pub(crate) fn ibd_memory_pressure_maintenance(
mem_mtx: &parking_lot::Mutex<MemoryGuard>,
max_ahead_live: &AtomicU64,
nominal_max_ahead: u64,
storage: &crate::storage::Storage,
utxo_engine: Option<&crate::storage::ibd_engine::UtxoDatabase>,
) -> PressureLevel {
let level = {
let mut mem = mem_mtx.lock();
let level = mem.should_flush(Some((max_ahead_live, nominal_max_ahead)));
publish_ibd_pressure(level);
level
};
storage.ibd_memory_pressure_tick(level as u8);
if let Some(db) = utxo_engine {
db.memory_pressure_tick(level as u8);
}
level
}
#[inline]
pub(crate) fn ibd_pressure_is_emergency() -> bool {
IBD_PRESSURE_LEVEL.load(Ordering::Relaxed) >= PressureLevel::Emergency as u8
}
#[inline]
pub(crate) fn ibd_pressure_is_critical_or_worse() -> bool {
IBD_PRESSURE_LEVEL.load(Ordering::Relaxed) >= PressureLevel::Critical as u8
}
pub(crate) fn reset_ibd_pressure_on_session_end() {
IBD_PRESSURE_LEVEL.store(PressureLevel::None as u8, Ordering::Release);
IBD_RSS_ANON_MB.store(0, Ordering::Release);
}
fn emergency_entry_anon_mb(rss_budget_mb: u64, no_swap: bool) -> u64 {
if no_swap {
rss_budget_mb * 85 / 100
} else {
rss_budget_mb
}
}
pub(crate) fn stale_emergency_step_down_level(
anon_mb: u64,
rss_budget_mb: u64,
no_swap: bool,
) -> Option<PressureLevel> {
if rss_budget_mb == 0 {
return None;
}
let emerg_line = emergency_entry_anon_mb(rss_budget_mb, no_swap);
if anon_mb >= emerg_line {
return None;
}
let elev_line = if no_swap {
rss_budget_mb * 65 / 100
} else {
rss_budget_mb * 82 / 100
};
Some(if anon_mb < elev_line {
PressureLevel::None
} else {
PressureLevel::Critical
})
}
#[cfg(target_os = "linux")]
pub(crate) fn refresh_stale_emergency_pressure(rss_budget_mb: u64) {
if !ibd_pressure_is_emergency() || rss_budget_mb == 0 {
return;
}
let mut snap = MemorySnapshot::default();
let mut status = String::new();
let mut meminfo = String::new();
if proc_read_file("/proc/self/status", &mut status) {
proc_parse_status_into(&status, &mut snap);
}
if proc_read_file("/proc/meminfo", &mut meminfo) {
proc_parse_meminfo_into(&meminfo, &mut snap);
}
if snap.rss_mb == 0 {
return;
}
publish_ibd_rss_anon_mb(snap.rss_anon_mb);
let r = if snap.rss_anon_mb > 0 {
snap.rss_anon_mb
} else {
snap.rss_mb
};
let no_swap = snap.swap_total_mb == 0;
let Some(level) = stale_emergency_step_down_level(r, rss_budget_mb, no_swap) else {
return;
};
publish_ibd_pressure(level);
let emerg_line = emergency_entry_anon_mb(rss_budget_mb, no_swap);
tracing::info!(
"MemoryGuard: stale Emergency stepped down to {:?} — anon RSS {}MB < emerg {}MB ({})",
level,
r,
emerg_line,
snap
);
}
#[cfg(not(target_os = "linux"))]
pub(crate) fn refresh_stale_emergency_pressure(_rss_budget_mb: u64) {}
#[inline]
pub(crate) fn ibd_pressure_level_snapshot() -> PressureLevel {
PressureLevel::from_u8(IBD_PRESSURE_LEVEL.load(Ordering::Relaxed))
}
#[inline]
pub(crate) fn utxo_flush_concurrency_cap(base_max_flushes: usize) -> usize {
let base = base_max_flushes.max(1);
match ibd_pressure_level_snapshot() {
PressureLevel::None => {
let bonus = (base / 2).max(1);
(base + bonus).min(64)
}
PressureLevel::Elevated => {
let bonus = (base / 4).max(1);
(base + bonus).min(48)
}
PressureLevel::Critical | PressureLevel::Emergency => base,
}
}
#[inline]
pub(crate) fn last_reported_pressure_level(mg: &MemoryGuard) -> PressureLevel {
PressureLevel::from_u8(mg.last_reported_pressure.load(Ordering::Relaxed))
}
pub(crate) const TIDESDB_MAX_TXN_OPS: usize = 200_000;
pub(crate) static BLOCK_BUFFER_BYTES: AtomicU64 = AtomicU64::new(0);
pub(crate) static BLOCK_BUFFER_COUNT: AtomicU64 = AtomicU64::new(0);
pub(crate) static BRIDGE_PENDING_COUNT: AtomicU64 = AtomicU64::new(0);
pub(crate) static GAP_FLUSH_ON_ABORT_BLOCKS: AtomicU64 = AtomicU64::new(0);
pub(crate) static DOWNLOAD_RECEIVED_BLOCKS: AtomicU64 = AtomicU64::new(0);
pub(crate) static DOWNLOAD_RECEIVED_TRIM_BLOCKS: AtomicU64 = AtomicU64::new(0);
pub(crate) static GAP_STREAM_DEDUP_HEIGHT: AtomicU64 = AtomicU64::new(0);
#[inline]
pub(crate) fn bump_gap_stream_dedup(h: u64) {
let mut cur = GAP_STREAM_DEDUP_HEIGHT.load(Ordering::Relaxed);
while h > cur {
match GAP_STREAM_DEDUP_HEIGHT.compare_exchange(cur, h, Ordering::Relaxed, Ordering::Relaxed)
{
Ok(_) => break,
Err(actual) => cur = actual,
}
}
}
pub(crate) static GAP_STREAM_LAST_RESEND_HEIGHT: AtomicU64 = AtomicU64::new(0);
pub(crate) static GAP_STREAM_LAST_RESEND_MS: AtomicU64 = AtomicU64::new(0);
pub(crate) static BRIDGE_NEXT_EXPECTED: AtomicU64 = AtomicU64::new(u64::MAX);
pub(crate) static GAP_ADMIT_DROP_BLOCKS: AtomicU64 = AtomicU64::new(0);
pub(crate) static REORDER_EVICT_BLOCKS: AtomicU64 = AtomicU64::new(0);
pub(crate) static BRIDGE_EVICT_BLOCKS: AtomicU64 = AtomicU64::new(0);
static LAST_JEMALLOC_RETAINED_PURGE_MS: AtomicU64 = AtomicU64::new(0);
const JEMALLOC_RETAINED_PURGE_MIN_INTERVAL_MS: u64 = 60_000;
#[cfg(feature = "production")]
pub(crate) fn sync_reorder_buffer_stats(
reorder_buffer: &std::collections::BTreeMap<
u64,
(super::types::SharedBlock, super::types::SharedWitnesses),
>,
) {
use super::types::estimate_block_bytes;
let mut bytes = 0u64;
for (block, witnesses) in reorder_buffer.values() {
bytes += estimate_block_bytes(block.as_ref(), witnesses.as_ref()) as u64;
}
BLOCK_BUFFER_COUNT.store(reorder_buffer.len() as u64, Ordering::Relaxed);
BLOCK_BUFFER_BYTES.store(bytes, Ordering::Relaxed);
}
#[cfg(not(feature = "production"))]
pub(crate) fn sync_reorder_buffer_stats(
_reorder_buffer: &std::collections::BTreeMap<
u64,
(super::types::SharedBlock, super::types::SharedWitnesses),
>,
) {
}
#[cfg(all(feature = "jemalloc", not(test)))]
pub(crate) fn jemalloc_retained_excess_gb() -> u64 {
jemalloc_stats_snapshot()
.map(|s| s.retained_gb)
.unwrap_or(0)
}
#[cfg(all(feature = "jemalloc", not(test)))]
#[derive(Clone, Copy, Debug)]
struct JemallocStatsSnap {
retained_gb: u64,
retained_mb: u64,
allocated_mb: u64,
mapped_mb: u64,
resident_mb: u64,
opt_retain: bool,
background_thread: bool,
narenas: u32,
}
#[cfg(all(feature = "jemalloc", not(test)))]
fn jemalloc_stats_snapshot() -> Option<JemallocStatsSnap> {
use std::os::raw::c_void;
unsafe extern "C" {
fn _rjem_mallctl(
name: *const i8,
oldp: *mut c_void,
oldlenp: *mut usize,
newp: *mut c_void,
newlen: usize,
) -> i32;
}
unsafe {
let mut sz = std::mem::size_of::<usize>();
let epoch: usize = 1;
let _ = _rjem_mallctl(
c"epoch".as_ptr(),
std::ptr::null_mut(),
std::ptr::null_mut(),
&epoch as *const usize as *mut c_void,
sz,
);
let mut allocated: usize = 0;
let _ = _rjem_mallctl(
c"stats.allocated".as_ptr(),
&mut allocated as *mut usize as *mut c_void,
&mut sz,
std::ptr::null_mut(),
0,
);
let mut retained: usize = 0;
let _ = _rjem_mallctl(
c"stats.retained".as_ptr(),
&mut retained as *mut usize as *mut c_void,
&mut sz,
std::ptr::null_mut(),
0,
);
let mut mapped: usize = 0;
let _ = _rjem_mallctl(
c"stats.mapped".as_ptr(),
&mut mapped as *mut usize as *mut c_void,
&mut sz,
std::ptr::null_mut(),
0,
);
let mut resident: usize = 0;
let _ = _rjem_mallctl(
c"stats.resident".as_ptr(),
&mut resident as *mut usize as *mut c_void,
&mut sz,
std::ptr::null_mut(),
0,
);
let mut opt_retain: bool = false;
let mut bsz = std::mem::size_of::<bool>();
let _ = _rjem_mallctl(
c"opt.retain".as_ptr(),
&mut opt_retain as *mut bool as *mut c_void,
&mut bsz,
std::ptr::null_mut(),
0,
);
let mut background_thread: bool = false;
bsz = std::mem::size_of::<bool>();
let _ = _rjem_mallctl(
c"background_thread".as_ptr(),
&mut background_thread as *mut bool as *mut c_void,
&mut bsz,
std::ptr::null_mut(),
0,
);
let mut narenas: u32 = 0;
let mut nsz = std::mem::size_of::<u32>();
let _ = _rjem_mallctl(
c"arenas.narenas".as_ptr(),
&mut narenas as *mut u32 as *mut c_void,
&mut nsz,
std::ptr::null_mut(),
0,
);
Some(JemallocStatsSnap {
retained_gb: retained as u64 / (1024 * 1024 * 1024),
retained_mb: retained as u64 / (1024 * 1024),
allocated_mb: allocated as u64 / (1024 * 1024),
mapped_mb: mapped as u64 / (1024 * 1024),
resident_mb: resident as u64 / (1024 * 1024),
opt_retain,
background_thread,
narenas,
})
}
}
#[cfg(any(not(feature = "jemalloc"), test))]
pub(crate) fn jemalloc_retained_excess_gb() -> u64 {
0
}
#[cfg(all(feature = "jemalloc", not(test)))]
pub(crate) fn maybe_purge_jemalloc_retained(reason: &str) -> bool {
let threshold_gb: u64 = std::env::var("BLVM_IBD_JEMALLOC_RETAINED_PURGE_GB")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(16);
let before = match jemalloc_stats_snapshot() {
Some(s) if s.retained_gb >= threshold_gb => s,
_ => return false,
};
let now_ms = crate::utils::time::current_timestamp_millis();
loop {
let prev = LAST_JEMALLOC_RETAINED_PURGE_MS.load(Ordering::Relaxed);
if now_ms.saturating_sub(prev) < JEMALLOC_RETAINED_PURGE_MIN_INTERVAL_MS {
return false;
}
if LAST_JEMALLOC_RETAINED_PURGE_MS
.compare_exchange_weak(prev, now_ms, Ordering::Relaxed, Ordering::Relaxed)
.is_ok()
{
break;
}
}
use std::os::raw::c_void;
unsafe extern "C" {
fn _rjem_mallctl(
name: *const i8,
oldp: *mut c_void,
oldlenp: *mut usize,
newp: *mut c_void,
newlen: usize,
) -> i32;
}
let mut rcs: Vec<(String, i32)> = Vec::new();
unsafe {
let mut decay_ms: isize = 0;
let sz = std::mem::size_of::<isize>();
rcs.push((
"dirty_decay".into(),
_rjem_mallctl(
c"arenas.dirty_decay_ms".as_ptr(),
std::ptr::null_mut(),
std::ptr::null_mut(),
&mut decay_ms as *mut isize as *mut c_void,
sz,
),
));
rcs.push((
"muzzy_decay".into(),
_rjem_mallctl(
c"arenas.muzzy_decay_ms".as_ptr(),
std::ptr::null_mut(),
std::ptr::null_mut(),
&mut decay_ms as *mut isize as *mut c_void,
sz,
),
));
let all = u32::MAX;
let decay_all = format!("arena.{all}.decay\0");
let purge_all = format!("arena.{all}.purge\0");
rcs.push((
"decay_all".into(),
_rjem_mallctl(
decay_all.as_ptr() as *const i8,
std::ptr::null_mut(),
std::ptr::null_mut(),
std::ptr::null_mut(),
0,
),
));
rcs.push((
"purge_all".into(),
_rjem_mallctl(
purge_all.as_ptr() as *const i8,
std::ptr::null_mut(),
std::ptr::null_mut(),
std::ptr::null_mut(),
0,
),
));
let mut purged_ok = 0u32;
let mut purged_err = 0u32;
for i in 0..before.narenas.min(256) {
let decay = format!("arena.{i}.decay\0");
let _ = _rjem_mallctl(
decay.as_ptr() as *const i8,
std::ptr::null_mut(),
std::ptr::null_mut(),
std::ptr::null_mut(),
0,
);
let purge = format!("arena.{i}.purge\0");
let rc = _rjem_mallctl(
purge.as_ptr() as *const i8,
std::ptr::null_mut(),
std::ptr::null_mut(),
std::ptr::null_mut(),
0,
);
if rc == 0 {
purged_ok = purged_ok.saturating_add(1);
} else {
purged_err = purged_err.saturating_add(1);
}
}
rcs.push((format!("per_arena_ok={purged_ok}"), 0));
rcs.push((format!("per_arena_err={purged_err}"), 0));
}
#[cfg(target_os = "linux")]
unsafe {
libc::malloc_trim(0);
}
let after = jemalloc_stats_snapshot().unwrap_or(before);
if after.retained_mb + 64 < before.retained_mb {
tracing::info!(
"[JEMALLOC_RETAINED_PURGE] reason={reason} retained_before_mb={} \
retained_after_mb={} threshold_gb={threshold_gb} opt_retain={} bg_thread={} \
narenas={} mapped_mb={}->{} resident_mb={}->{} alloc_mb={}->{} rcs={rcs:?}",
before.retained_mb,
after.retained_mb,
before.opt_retain,
before.background_thread,
before.narenas,
before.mapped_mb,
after.mapped_mb,
before.resident_mb,
after.resident_mb,
before.allocated_mb,
after.allocated_mb,
);
} else {
tracing::warn!(
"[JEMALLOC_RETAINED_PURGE] NO_RECLAIM reason={reason} retained_before_mb={} \
retained_after_mb={} threshold_gb={threshold_gb} opt_retain={} bg_thread={} \
narenas={} mapped_mb={} resident_mb={} alloc_mb={} rcs={rcs:?} \
(if opt_retain=true, MALLOC_CONF retain:false did not apply)",
before.retained_mb,
after.retained_mb,
before.opt_retain,
before.background_thread,
before.narenas,
before.mapped_mb,
before.resident_mb,
before.allocated_mb,
);
}
true
}
#[cfg(any(not(feature = "jemalloc"), test))]
pub(crate) fn maybe_purge_jemalloc_retained(_reason: &str) -> bool {
false
}
static LAST_MDB_KEEP_TAIL_MADVISE_MS: AtomicU64 = AtomicU64::new(0);
const MDB_KEEP_TAIL_MADVISE_MIN_INTERVAL_MS: u64 = 60_000;
#[cfg(all(target_os = "linux", feature = "libc"))]
pub(crate) fn maybe_madvise_data_mdb_keep_tail(file_backed_mb: u64, reason: &str) -> bool {
let threshold_mb: u64 = std::env::var("BLVM_IBD_MDB_FILE_RSS_MADVISE_MB")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(4 * 1024);
if file_backed_mb < threshold_mb {
return false;
}
let now_ms = crate::utils::time::current_timestamp_millis();
loop {
let prev = LAST_MDB_KEEP_TAIL_MADVISE_MS.load(Ordering::Relaxed);
if now_ms.saturating_sub(prev) < MDB_KEEP_TAIL_MADVISE_MIN_INTERVAL_MS {
return false;
}
if LAST_MDB_KEEP_TAIL_MADVISE_MS
.compare_exchange_weak(prev, now_ms, Ordering::Relaxed, Ordering::Relaxed)
.is_ok()
{
break;
}
}
let keep_gb: u64 = std::env::var("BLVM_IBD_MDB_KEEP_TAIL_GB")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(4);
let keep_bytes = (keep_gb.max(1) as usize).saturating_mul(1024 * 1024 * 1024);
let keep_mb = keep_gb.saturating_mul(1024);
let cycle_cap_gb: usize = std::env::var("BLVM_IBD_MDB_MADVISE_MAX_GB")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(8)
.clamp(1, 64);
let max_advise_bytes = (file_backed_mb.saturating_sub(keep_mb) as usize)
.saturating_mul(1024 * 1024)
.min(cycle_cap_gb.saturating_mul(1024 * 1024 * 1024));
if max_advise_bytes == 0 {
return false;
}
let Ok(maps) = std::fs::File::open("/proc/self/maps") else {
return false;
};
use std::io::BufRead;
let mut ranges: Vec<(usize, usize)> = Vec::new();
for line in std::io::BufReader::new(maps).lines().map_while(Result::ok) {
if !line.contains("data.mdb") {
continue;
}
let Some(range) = line.split_whitespace().next() else {
continue;
};
let mut parts = range.splitn(2, '-');
let (Some(s), Some(e)) = (parts.next(), parts.next()) else {
continue;
};
let (Ok(start), Ok(end)) = (usize::from_str_radix(s, 16), usize::from_str_radix(e, 16))
else {
continue;
};
if end > start {
ranges.push((start, end));
}
}
if ranges.is_empty() {
return false;
}
let rss_before = {
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)
/ 1024
};
let mut advised_bytes: u64 = 0;
let mut ranges_touched = 0usize;
let mut budget = max_advise_bytes;
for (start, end) in ranges {
if budget == 0 {
break;
}
let len = end.saturating_sub(start);
if len <= keep_bytes {
continue;
}
let drop_len = (len - keep_bytes).min(budget);
unsafe {
libc::madvise(start as *mut libc::c_void, drop_len, libc::MADV_DONTNEED);
}
advised_bytes += drop_len as u64;
budget = budget.saturating_sub(drop_len);
ranges_touched += 1;
}
let rss_after = {
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)
/ 1024
};
tracing::info!(
"[MADVISE_KEEP_TAIL] reason={reason} file_backed_mb={file_backed_mb} \
threshold_mb={threshold_mb} keep_gb={keep_gb} cycle_cap_gb={cycle_cap_gb} \
ranges={ranges_touched} advised_gb={:.1} rss {}MB → {}MB",
advised_bytes as f64 / (1024.0 * 1024.0 * 1024.0),
rss_before,
rss_after
);
ranges_touched > 0
}
#[cfg(not(all(target_os = "linux", feature = "libc")))]
pub(crate) fn maybe_madvise_data_mdb_keep_tail(_file_backed_mb: u64, _reason: &str) -> bool {
false
}
#[derive(Default, Clone, Copy)]
pub(crate) struct MemorySnapshot {
pub rss_mb: u64,
pub rss_anon_mb: u64,
pub rss_file_mb: u64,
pub rss_shmem_mb: u64,
pub vm_size_mb: u64,
pub mem_total_mb: u64,
pub sys_avail_mb: u64,
pub swap_total_mb: u64,
pub swap_free_mb: u64,
pub vm_swap_mb: u64,
}
impl MemorySnapshot {
#[inline]
pub fn swap_used_mb(&self) -> u64 {
self.swap_total_mb.saturating_sub(self.swap_free_mb)
}
}
impl std::fmt::Display for MemorySnapshot {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"rss={}MB(anon={}MB file={}MB shm={}MB) vm={}MB mem_total={}MB sys_avail={}MB swap_used={}MB proc_swap={}MB",
self.rss_mb,
self.rss_anon_mb,
self.rss_file_mb,
self.rss_shmem_mb,
self.vm_size_mb,
self.mem_total_mb,
self.sys_avail_mb,
self.swap_used_mb(),
self.vm_swap_mb,
)
}
}
#[cfg(target_os = "linux")]
#[inline]
fn proc_field_kb_to_mb(line: &str) -> u64 {
line.split_whitespace()
.nth(1)
.and_then(|v| v.parse::<u64>().ok())
.unwrap_or(0)
/ 1024
}
#[cfg(target_os = "linux")]
fn proc_read_file(path: &str, buf: &mut String) -> bool {
buf.clear();
match std::fs::File::open(path) {
Ok(mut f) => f.read_to_string(buf).is_ok(),
Err(_) => false,
}
}
#[cfg(target_os = "linux")]
fn proc_parse_status_into(s: &str, snap: &mut MemorySnapshot) {
for line in s.lines() {
if line.starts_with("VmRSS:") {
snap.rss_mb = proc_field_kb_to_mb(line);
} else if line.starts_with("RssAnon:") {
snap.rss_anon_mb = proc_field_kb_to_mb(line);
} else if line.starts_with("RssFile:") {
snap.rss_file_mb = proc_field_kb_to_mb(line);
} else if line.starts_with("RssShmem:") {
snap.rss_shmem_mb = proc_field_kb_to_mb(line);
} else if line.starts_with("VmSize:") {
snap.vm_size_mb = proc_field_kb_to_mb(line);
} else if line.starts_with("VmSwap:") {
snap.vm_swap_mb = proc_field_kb_to_mb(line);
}
}
}
#[cfg(target_os = "linux")]
fn proc_parse_meminfo_into(s: &str, snap: &mut MemorySnapshot) {
for line in s.lines() {
if line.starts_with("MemTotal:") {
snap.mem_total_mb = proc_field_kb_to_mb(line);
} else if line.starts_with("MemAvailable:") {
snap.sys_avail_mb = proc_field_kb_to_mb(line);
} else if line.starts_with("SwapTotal:") {
snap.swap_total_mb = proc_field_kb_to_mb(line);
} else if line.starts_with("SwapFree:") {
snap.swap_free_mb = proc_field_kb_to_mb(line);
}
}
}
#[cfg(target_os = "linux")]
pub(crate) fn read_proc_anon_rss_mb() -> u64 {
if let Ok(s) = std::fs::read_to_string("/proc/self/status") {
proc_anon_rss_mb_from_status(&s)
} else {
0
}
}
#[cfg(target_os = "linux")]
pub(crate) fn read_proc_anon_and_swap_mb() -> (u64, u64) {
if let Ok(s) = std::fs::read_to_string("/proc/self/status") {
let anon = proc_anon_rss_mb_from_status(&s);
let swap = proc_vm_swap_mb_from_status(&s);
(anon, swap)
} else {
(0, 0)
}
}
#[cfg(not(target_os = "linux"))]
pub(crate) fn read_proc_anon_rss_mb() -> u64 {
0
}
#[cfg(not(target_os = "linux"))]
pub(crate) fn read_proc_anon_and_swap_mb() -> (u64, u64) {
(0, 0)
}
#[cfg(target_os = "linux")]
fn proc_anon_rss_mb_from_status(s: &str) -> u64 {
let mut rss_mb: u64 = 0;
let mut anon_mb: u64 = 0;
for line in s.lines() {
if line.starts_with("VmRSS:") {
rss_mb = proc_field_kb_to_mb(line);
} else if line.starts_with("RssAnon:") {
anon_mb = proc_field_kb_to_mb(line);
}
}
if anon_mb > 0 { anon_mb } else { rss_mb }
}
#[cfg(target_os = "linux")]
fn proc_vm_swap_mb_from_status(s: &str) -> u64 {
for line in s.lines() {
if line.starts_with("VmSwap:") {
return proc_field_kb_to_mb(line);
}
}
0
}
#[cfg(target_os = "linux")]
fn proc_rss_mb_from_status(s: &str) -> u64 {
for line in s.lines() {
if line.starts_with("VmRSS:") {
return proc_field_kb_to_mb(line);
}
}
0
}
pub(crate) struct MemoryGuard {
pub(crate) total_mb: u64,
pub(crate) avail_mb: u64,
budget_mb: u64,
utxo_cache_mb: usize,
pub(crate) utxo_max_entries: usize,
pub(crate) rss_budget_mb: u64,
last_adaptive_cap_entries: AtomicUsize,
last_adaptive_cap_check: Mutex<Instant>,
last_adaptive_cap_shrink: Mutex<Instant>,
above_threshold_consecutive: AtomicU8,
pub(crate) utxo_flush_threshold: usize,
block_buffer_base: usize,
pub(crate) storage_flush_interval: usize,
prefetch_limit: usize,
pub(crate) prefetch_queue_size: usize,
pub(crate) max_ahead_blocks: u64,
pub defer_flush: bool,
pub defer_checkpoint_interval: u64,
pub feeder_buffer_bytes_limit: usize,
pub max_utxo_flushes: usize,
pub max_block_flushes: usize,
storage_backend: crate::storage::database::DatabaseBackend,
no_swap_at_boot: bool,
pub spec_adds_bytes: Arc<AtomicU64>,
#[cfg(feature = "sysinfo")]
sys: sysinfo::System,
last_rss_check: Instant,
last_ahead_adjust: Instant,
last_reported_pressure: AtomicU8,
crit_rss_threshold_mb: u64,
#[cfg(target_os = "linux")]
proc_status_buf: String,
#[cfg(target_os = "linux")]
proc_meminfo_buf: String,
}
#[derive(Clone, Copy)]
pub(crate) struct FeederScaleSnapshot {
pub block_buffer_base: usize,
pub total_mb: u64,
pub feeder_buffer_bytes_limit: usize,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum WorkloadClass {
Shared,
Dedicated,
}
pub(crate) const ROCKSDB_PIPELINE_RESERVE_MB: u64 = 2048;
#[derive(Debug, Clone, Copy)]
pub(crate) struct IbdTuningContext {
pub storage_backend: crate::storage::database::DatabaseBackend,
pub ibd_dedicated: bool,
pub db_file_size_mb: u64,
}
impl Default for IbdTuningContext {
fn default() -> Self {
Self {
storage_backend: crate::storage::database::default_backend(),
ibd_dedicated: false,
db_file_size_mb: 0,
}
}
}
impl MemoryGuard {
pub(crate) const EXTENDED_SIXTEEN_CLASS_MB: u64 = 18 * 1024;
#[inline]
pub(crate) fn total_gb_rounded(total_mb: u64) -> u64 {
crate::utils::ram_tier::total_gb_rounded(total_mb)
}
pub(crate) fn detect_workload_class(
total_mb: u64,
avail_mb: u64,
ibd_dedicated: bool,
) -> WorkloadClass {
let exclusive = std::env::var("BLVM_IBD_EXCLUSIVE")
.map(|v| v == "1" || v.eq_ignore_ascii_case("true"))
.unwrap_or(false);
if exclusive {
return WorkloadClass::Dedicated;
}
let avail_pct = if total_mb > 0 {
avail_mb.saturating_mul(100) / total_mb
} else {
100
};
let wants_dedicated = ibd_dedicated
|| std::env::var("BLVM_DEDICATED_NODE")
.map(|v| v == "1" || v.eq_ignore_ascii_case("true"))
.unwrap_or(false);
if wants_dedicated {
const DEDICATED_MIN_AVAIL_PCT: u64 = 70;
if avail_pct < DEDICATED_MIN_AVAIL_PCT {
tracing::warn!(
"MemoryGuard: dedicated mode requested but MemAvailable is only {avail_pct}% \
of MemTotal — using Shared memory envelope (set BLVM_IBD_EXCLUSIVE=1 to force Dedicated)"
);
return WorkloadClass::Shared;
}
return WorkloadClass::Dedicated;
}
let dedicated_threshold_pct: u64 = if total_mb >= 32 * 1024 { 60 } else { 70 };
if total_mb == 0 || avail_pct >= dedicated_threshold_pct {
WorkloadClass::Dedicated
} else {
WorkloadClass::Shared
}
}
pub(crate) fn pipeline_reserve_mb(backend: crate::storage::database::DatabaseBackend) -> u64 {
use crate::storage::database::DatabaseBackend;
match backend {
DatabaseBackend::RocksDB => ROCKSDB_PIPELINE_RESERVE_MB,
DatabaseBackend::Heed3 => 512,
DatabaseBackend::TidesDB => 1024,
DatabaseBackend::Redb | DatabaseBackend::Sled => 768,
}
}
pub(crate) fn storage_flush_interval_base(
total_gb: u64,
backend: crate::storage::database::DatabaseBackend,
) -> usize {
use crate::storage::database::DatabaseBackend;
match backend {
DatabaseBackend::Heed3 => {
if total_gb >= 16 { 200 } else { 100 }
}
DatabaseBackend::RocksDB | DatabaseBackend::TidesDB => {
if total_gb >= 32 {
2000
} else {
300
}
}
DatabaseBackend::Redb | DatabaseBackend::Sled => {
if total_gb >= 32 {
512
} else {
300
}
}
}
}
pub(crate) fn compute_rss_budget_mb(
total_mb: u64,
avail_mb: u64,
workload: WorkloadClass,
) -> u64 {
if let Some(v) = std::env::var("BLVM_RSS_BUDGET_MB")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.filter(|&n| n >= 1024)
{
return v;
}
Self::compute_rss_budget_mb_auto(total_mb, avail_mb, workload)
}
pub(crate) fn compute_rss_budget_mb_auto(
total_mb: u64,
avail_mb: u64,
workload: WorkloadClass,
) -> u64 {
if total_mb <= 8 * 1024 {
let from_total = total_mb * 65 / 100;
let from_avail = avail_mb * 75 / 100;
return from_total
.max(from_avail)
.min(total_mb * 80 / 100)
.max(2048);
}
if total_mb <= Self::EXTENDED_SIXTEEN_CLASS_MB {
let from_avail = (avail_mb * 60 / 100).clamp(3000, 7000);
return from_avail.max(2048);
}
if workload == WorkloadClass::Dedicated {
let cap_pct: u64 = if total_mb >= 64 * 1024 {
50
} else if total_mb >= 32 * 1024 {
45
} else {
40
};
let from_total = total_mb * cap_pct / 100;
let os_reserve = if total_mb >= 64 * 1024 {
8192
} else {
(total_mb * 10 / 100).max(4096)
};
let from_avail = avail_mb.saturating_sub(os_reserve);
return from_total.min(from_avail.max(2048)).max(2048);
}
let os_reserve_pct: u64 = if total_mb >= 32 * 1024 { 25 } else { 22 };
let os_reserve = (total_mb * os_reserve_pct / 100).max(2816);
let from_spare = avail_mb.saturating_sub(os_reserve);
let cap_pct = match workload {
WorkloadClass::Shared => {
if total_mb >= 32 * 1024 {
25
} else {
35
}
}
WorkloadClass::Dedicated => unreachable!("handled above"),
};
from_spare.min(total_mb * cap_pct / 100).max(2048)
}
pub(crate) fn nominal_max_pending_ops(
total_mb: u64,
rss_budget_mb: u64,
utxo_cache_mb: usize,
utxo_flush_threshold: usize,
storage_backend: crate::storage::database::DatabaseBackend,
) -> usize {
if let Some(v) = std::env::var("BLVM_IBD_MAX_PENDING_OPS")
.ok()
.and_then(|s| s.parse::<usize>().ok())
{
return v.max(100_000);
}
const BYTES_PER_OP: usize = 160;
const PIPELINE_FRAC_PCT: usize = 6;
let total_gb = Self::total_gb_rounded(total_mb);
let tier_ceiling = if total_gb >= 64 {
8_000_000
} else if total_gb >= 32 {
5_000_000
} else if total_gb >= 24 {
3_000_000
} else if total_gb >= 16 {
1_500_000
} else {
1_000_000
};
let pipeline_mb = rss_budget_mb
.saturating_sub(utxo_cache_mb as u64)
.saturating_sub(Self::pipeline_reserve_mb(storage_backend));
let from_envelope =
pipeline_mb as usize * 1024 * 1024 * PIPELINE_FRAC_PCT / 100 / BYTES_PER_OP;
let floor = 400_000_usize.max(utxo_flush_threshold.saturating_mul(4));
from_envelope.clamp(floor, tier_ceiling)
}
pub(crate) fn nominal_max_pending_ops_for_guard(&self) -> usize {
Self::nominal_max_pending_ops(
self.total_mb,
self.rss_budget_mb,
self.utxo_cache_mb,
self.utxo_flush_threshold,
self.storage_backend,
)
}
pub(crate) fn new() -> Self {
Self::new_for_ibd(IbdTuningContext::default())
}
pub(crate) fn new_for_ibd(ctx: IbdTuningContext) -> Self {
#[cfg(target_os = "linux")]
let (mut total_mb, mut available_mb, startup_swap_total_mb, startup_swap_free_mb) = {
let mut total = 0u64;
let mut avail = 0u64;
let mut swap_total = 0u64;
let mut swap_free = 0u64;
if let Ok(s) = std::fs::read_to_string("/proc/meminfo") {
for line in s.lines() {
if line.starts_with("MemTotal:") {
total = proc_field_kb_to_mb(line);
} else if line.starts_with("MemAvailable:") {
avail = proc_field_kb_to_mb(line);
} else if line.starts_with("SwapTotal:") {
swap_total = proc_field_kb_to_mb(line);
} else if line.starts_with("SwapFree:") {
swap_free = proc_field_kb_to_mb(line);
}
}
}
(total, avail, swap_total, swap_free)
};
#[cfg(not(target_os = "linux"))]
let (mut total_mb, mut available_mb, startup_swap_total_mb, startup_swap_free_mb) =
(0u64, 0u64, 0u64, 0u64);
#[cfg(feature = "sysinfo")]
let mut sys = {
use sysinfo::System;
let mut s = System::new_all();
s.refresh_memory();
if total_mb == 0 {
total_mb = s.total_memory() / (1024 * 1024);
}
if available_mb == 0 {
available_mb = s.available_memory() / (1024 * 1024);
}
s
};
if let Some(mb) = std::env::var("BLVM_TOTAL_RAM_MB")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.filter(|&v| v > 0)
{
total_mb = mb;
}
if let Some(mb) = std::env::var("BLVM_SYS_AVAIL_MB")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.filter(|&v| v > 0)
{
available_mb = mb;
}
if total_mb == 0 {
total_mb = 8192;
}
if available_mb == 0 {
available_mb = (total_mb * 60 / 100).max(2048);
}
let total_gb = Self::total_gb_rounded(total_mb);
let mut budget_mb = if total_mb <= Self::EXTENDED_SIXTEEN_CLASS_MB {
(total_mb * 15 / 100).clamp(512, 2500)
} else {
(total_mb * 28 / 100).min(available_mb * 45 / 100).max(512)
};
let effective_avail = if total_mb <= Self::EXTENDED_SIXTEEN_CLASS_MB {
available_mb.min(total_mb * 40 / 100)
} else {
available_mb
};
let os_reserve_mb = (total_mb * 22 / 100).max(2816);
let spare_mb = effective_avail.saturating_sub(os_reserve_mb).max(256);
let swap_nearly_full = startup_swap_total_mb > 0
&& startup_swap_free_mb * 100 / startup_swap_total_mb.max(1) < 15;
let latched = WORKLOAD_CLASS_LATCH.load(Ordering::Relaxed);
let workload = if latched != 0xff {
let wc = if latched == 1 {
WorkloadClass::Dedicated
} else {
WorkloadClass::Shared
};
tracing::debug!(
"MemoryGuard: reusing latched workload={:?} (available={}MB; latch prevents oscillation)",
wc,
available_mb
);
wc
} else {
let wc = if swap_nearly_full && !ctx.ibd_dedicated {
tracing::warn!(
"MemoryGuard: swap is >85% full at startup ({}/{} MB free) — \
overriding workload to Dedicated to avoid cache surge on restart",
startup_swap_free_mb,
startup_swap_total_mb
);
WorkloadClass::Dedicated
} else {
Self::detect_workload_class(total_mb, available_mb, ctx.ibd_dedicated)
};
WORKLOAD_CLASS_LATCH.store(
if wc == WorkloadClass::Dedicated { 1 } else { 0 },
Ordering::Relaxed,
);
tracing::info!(
"MemoryGuard: workload class latched as {:?} for this process (available={}MB, swap={}/{}MB free)",
wc,
available_mb,
startup_swap_free_mb,
startup_swap_total_mb
);
wc
};
let auto_rss_budget_mb = Self::compute_rss_budget_mb_auto(total_mb, available_mb, workload);
let mut rss_budget_mb_raw = Self::compute_rss_budget_mb(total_mb, available_mb, workload);
if startup_swap_total_mb == 0 && total_mb >= 32 * 1024 {
let os_reserve = if total_mb >= 64 * 1024 {
8192
} else {
(total_mb * 10 / 100).max(4096)
};
let no_swap_cap = available_mb
.saturating_sub(os_reserve)
.min(total_mb * 40 / 100);
if rss_budget_mb_raw > no_swap_cap && no_swap_cap >= 2048 {
tracing::info!(
"MemoryGuard: no swap configured — capping rss_budget {} -> {} MB \
(MemAvailable={}MB)",
rss_budget_mb_raw,
no_swap_cap,
available_mb
);
rss_budget_mb_raw = no_swap_cap;
}
}
let rss_budget_mb = rss_budget_mb_raw.max(2048);
let envelope_cache_cap_mb = (rss_budget_mb * 45 / 100) as usize;
let mut utxo_cache_mb = if total_gb >= 32 {
let tier_max = match workload {
WorkloadClass::Shared => 4096,
WorkloadClass::Dedicated => 16384,
};
envelope_cache_cap_mb.min(tier_max)
} else if total_gb >= 17 && total_mb > Self::EXTENDED_SIXTEEN_CLASS_MB {
((available_mb * 50 / 100) as usize).clamp(2048, 4096)
} else if total_gb >= 16 {
((available_mb * 30 / 100) as usize).clamp(1024, 1400)
} else if total_gb >= 12 {
((available_mb * 25 / 100) as usize).clamp(768, 1536)
} else if total_gb >= 8 {
((available_mb * 20 / 100) as usize).clamp(256, 512)
} else if total_gb >= 4 {
((available_mb * 15 / 100) as usize).clamp(128, 256)
} else {
((budget_mb * 30 / 100) as usize).clamp(64, 128)
};
if total_gb < 12 || available_mb < 6144 {
let tight_cap_mb = (total_mb.saturating_mul(7) / 100).clamp(128, 384) as usize;
utxo_cache_mb = utxo_cache_mb.min(tight_cap_mb);
}
if let Some(mb) = std::env::var("BLVM_UTXO_CACHE_MAX_MB")
.ok()
.and_then(|s| s.parse::<usize>().ok())
{
if mb > 0 {
utxo_cache_mb = utxo_cache_mb.min(mb);
}
}
utxo_cache_mb = utxo_cache_mb.min(envelope_cache_cap_mb.max(128));
let utxo_max_entries = utxo_cache_mb * 1024 * 1024 / 1600;
let utxo_flush_threshold = {
const BYTES_PER_OP: usize = 160;
let target = (spare_mb as usize).saturating_mul(1024 * 1024) * 6 / 100 / BYTES_PER_OP;
let max_ops: usize = if total_gb >= 48 {
2_000_000
} else if total_gb >= 32 {
1_200_000
} else if total_gb >= 24 {
800_000
} else if total_gb >= 17 && total_mb > Self::EXTENDED_SIXTEEN_CLASS_MB {
480_000
} else if total_gb >= 16 {
320_000
} else {
120_000
};
target.clamp(40_000, max_ops)
};
let crit_rss_threshold_mb = std::env::var("BLVM_IBD_PRESSURE_CRIT_RSS_MB")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.filter(|&n| (800..=8000).contains(&n))
.unwrap_or_else(|| {
(total_mb * 22 / 100).clamp(1200, 6000)
});
let defer_flush = if std::env::var("BLVM_IBD_DEFER_FLUSH")
.map(|v| v == "0" || v.eq_ignore_ascii_case("false"))
.unwrap_or(false)
{
false
} else if std::env::var("BLVM_IBD_DEFER_FLUSH")
.map(|v| v == "1" || v.eq_ignore_ascii_case("true"))
.unwrap_or(false)
{
true
} else if matches!(workload, WorkloadClass::Shared) {
false
} else {
total_gb >= 32
};
let defer_checkpoint_interval_base = if matches!(
ctx.storage_backend,
crate::storage::database::DatabaseBackend::Heed3
) {
200u64
} else if total_gb >= 64 {
10_000
} else if total_gb >= 32 {
2_000
} else {
25_000
};
let defer_checkpoint_interval = std::env::var("BLVM_IBD_DEFER_CHECKPOINT_INTERVAL")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.filter(|&n| (20..=500_000).contains(&n))
.unwrap_or(defer_checkpoint_interval_base);
let block_buffer_base = {
let buffer_mb = budget_mb * 10 / 100;
let blocks = buffer_mb * 1024 / 500;
(blocks as usize).clamp(100, 800)
};
let mut storage_flush_interval =
Self::storage_flush_interval_base(total_gb, ctx.storage_backend);
if let Ok(s) = std::env::var("BLVM_IBD_STORAGE_FLUSH_INTERVAL") {
if let Ok(n) = s.parse::<usize>() {
storage_flush_interval = n.clamp(16, 4000);
}
}
let prefetch_queue_size = {
let hi: u64 = if total_gb <= 3 {
16
} else if total_gb <= 7 {
32
} else if total_mb <= Self::EXTENDED_SIXTEEN_CLASS_MB {
160
} else if total_gb <= 24 {
1024
} else {
2048
};
(spare_mb / 10).clamp(16, hi) as usize
};
let max_ahead_blocks = {
let mut v = (spare_mb / 8).clamp(64, 8192);
if total_gb < 32 {
v = v.min(4096);
}
let tier_cap = Self::tier_max_download_ahead_blocks(total_mb);
v.min(tier_cap)
};
let prefetch_limit = {
let cache_mb = budget_mb * 3 / 100;
let hi = if total_mb <= Self::EXTENDED_SIXTEEN_CLASS_MB {
8000
} else if total_gb <= 24 {
35_000
} else {
50_000
};
let spare_boost = ((spare_mb / 1024) as usize).saturating_mul(800);
(((cache_mb * 1024 * 1024 / 400) as usize).saturating_add(spare_boost)).clamp(5_000, hi)
};
let feeder_pct = if total_mb <= Self::EXTENDED_SIXTEEN_CLASS_MB {
2
} else {
5
};
let feeder_buffer_bytes_limit = (budget_mb * feeder_pct / 100 * 1024 * 1024) as usize;
let max_utxo_flushes_auto: usize = {
use crate::storage::database::DatabaseBackend;
match ctx.storage_backend {
DatabaseBackend::Heed3 | DatabaseBackend::Redb | DatabaseBackend::Sled => 4,
_ => {
if total_mb <= Self::EXTENDED_SIXTEEN_CLASS_MB {
8
} else if total_gb <= 24 {
12
} else if total_gb <= 32 {
16
} else {
32
}
}
}
};
let max_utxo_flushes: usize = std::env::var("BLVM_IBD_MAX_UTXO_FLUSHES")
.ok()
.and_then(|s| s.parse::<usize>().ok())
.filter(|&n| n > 0)
.map(|n| n.clamp(1, 64))
.unwrap_or(max_utxo_flushes_auto);
let max_block_flushes_auto: usize = {
use crate::storage::database::DatabaseBackend;
match ctx.storage_backend {
DatabaseBackend::Heed3 | DatabaseBackend::Redb | DatabaseBackend::Sled => 1,
DatabaseBackend::RocksDB | DatabaseBackend::TidesDB => {
if total_gb <= 24 {
max_utxo_flushes
} else if total_gb <= 32 {
max_utxo_flushes + max_utxo_flushes / 2
} else {
(max_utxo_flushes + max_utxo_flushes / 2).min(48)
}
}
}
};
let max_block_flushes: usize = std::env::var("BLVM_IBD_MAX_BLOCK_FLUSHES")
.ok()
.and_then(|s| s.parse::<usize>().ok())
.filter(|&n| n > 0)
.map(|n| n.clamp(1, 64))
.unwrap_or(max_block_flushes_auto);
tracing::info!(
"MemoryGuard: total={}MB available={}MB workload={:?} backend={:?} spare≈{}MB budget={}MB \
rss_budget={}MB (live /proc pressure) utxo_cache={}MB ({}entries) flush_threshold={} \
defer_flush={} defer_checkpoint={} buffer={} prefetch={} prefetch_queue={} \
max_ahead={} storage_flush={} pipeline_reserve={}MB feeder_bytes={}MB max_utxo_flush={} max_block_flush={}",
total_mb,
available_mb,
workload,
ctx.storage_backend,
spare_mb,
budget_mb,
rss_budget_mb,
utxo_cache_mb,
utxo_max_entries,
utxo_flush_threshold,
defer_flush,
defer_checkpoint_interval,
block_buffer_base,
prefetch_limit,
prefetch_queue_size,
max_ahead_blocks,
storage_flush_interval,
Self::pipeline_reserve_mb(ctx.storage_backend),
feeder_buffer_bytes_limit / (1024 * 1024),
max_utxo_flushes,
max_block_flushes,
);
if std::env::var("BLVM_RSS_BUDGET_MB").is_ok() {
tracing::warn!(
"MemoryGuard: BLVM_RSS_BUDGET_MB overrides auto envelope (auto would be {}MB); \
remove for backend-aware auto-tuning unless debugging OOM",
auto_rss_budget_mb
);
}
if std::env::var("BLVM_IBD_MAX_PARALLEL").is_ok() {
tracing::warn!(
"MemoryGuard: BLVM_IBD_MAX_PARALLEL set — validation worker count is pinned; \
auto scales to logical CPUs on 32+ GiB hosts"
);
}
if std::env::var("BLVM_IBD_MAX_BLOCK_FLUSHES").is_ok() {
tracing::warn!(
"MemoryGuard: BLVM_IBD_MAX_BLOCK_FLUSHES set — blockstore flush pool capped; \
auto allows up to {} on this host",
max_block_flushes_auto
);
}
Self {
total_mb,
avail_mb: available_mb,
budget_mb,
utxo_cache_mb,
utxo_max_entries,
rss_budget_mb,
last_adaptive_cap_entries: AtomicUsize::new(0),
last_adaptive_cap_check: Mutex::new(Instant::now() - Duration::from_secs(60)),
last_adaptive_cap_shrink: Mutex::new(Instant::now() - Duration::from_secs(120)),
above_threshold_consecutive: AtomicU8::new(0),
utxo_flush_threshold,
block_buffer_base,
storage_flush_interval,
prefetch_limit,
prefetch_queue_size,
max_ahead_blocks,
defer_flush,
defer_checkpoint_interval,
feeder_buffer_bytes_limit,
max_utxo_flushes,
max_block_flushes,
storage_backend: ctx.storage_backend,
no_swap_at_boot: startup_swap_total_mb == 0,
#[cfg(feature = "sysinfo")]
sys,
last_rss_check: Instant::now(),
last_ahead_adjust: Instant::now() - Duration::from_secs(1),
last_reported_pressure: AtomicU8::new(PressureLevel::None as u8),
crit_rss_threshold_mb,
#[cfg(target_os = "linux")]
proc_status_buf: String::with_capacity(4096),
#[cfg(target_os = "linux")]
proc_meminfo_buf: String::with_capacity(8192),
spec_adds_bytes: Arc::new(AtomicU64::new(0)),
}
}
pub(crate) fn feeder_scale_snapshot(&self) -> FeederScaleSnapshot {
FeederScaleSnapshot {
block_buffer_base: self.block_buffer_base,
total_mb: self.total_mb,
feeder_buffer_bytes_limit: self.feeder_buffer_bytes_limit,
}
}
#[inline]
pub(crate) fn storage_flush_interval_live(&self, pressure: PressureLevel) -> usize {
Self::storage_flush_interval_live_for(self.storage_flush_interval, pressure)
}
#[inline]
pub(crate) fn storage_flush_interval_live_for(base: usize, pressure: PressureLevel) -> usize {
match pressure {
PressureLevel::None => base,
PressureLevel::Elevated => (base * 3 / 4).max(200),
PressureLevel::Critical => (base / 2).max(128),
PressureLevel::Emergency => (base / 4).max(64),
}
}
#[inline]
pub(crate) fn storage_flush_pending_bytes_pressure_cap(
&self,
pressure: PressureLevel,
) -> Option<u64> {
Self::storage_flush_pending_bytes_pressure_cap_for(self.budget_mb, pressure)
}
#[inline]
pub(crate) fn storage_flush_pending_bytes_pressure_cap_for(
budget_mb: u64,
pressure: PressureLevel,
) -> Option<u64> {
let pct: u64 = match pressure {
PressureLevel::None => 20,
PressureLevel::Elevated => 12,
PressureLevel::Critical => 6,
PressureLevel::Emergency => 4,
};
let raw = budget_mb.saturating_mul(1024 * 1024).saturating_mul(pct) / 100;
Some(raw.max(64 * 1024 * 1024))
}
#[inline]
pub(crate) fn storage_flush_pressure_min_blocks(flush_interval_live: usize) -> usize {
flush_interval_live
.saturating_mul(2)
.saturating_div(5)
.max(96)
}
#[inline]
pub(crate) fn system_total_ram_mb(&self) -> u64 {
self.total_mb
}
#[inline]
pub(crate) fn budget_mb(&self) -> u64 {
self.budget_mb
}
#[inline]
pub(crate) fn tier_max_download_ahead_blocks(total_mb: u64) -> u64 {
let total_gb = Self::total_gb_rounded(total_mb);
if total_gb <= 3 {
64
} else if total_gb <= 7 {
128
} else if total_gb < 16 {
256
} else if total_gb <= 16 || total_mb <= Self::EXTENDED_SIXTEEN_CLASS_MB {
320
} else if total_gb < 32 {
512
} else if total_gb < 64 {
1024
} else {
2048
}
}
#[inline]
pub(crate) fn ibd_utxo_flush_queue_depth_default(&self) -> usize {
let total_gb = Self::total_gb_rounded(self.total_mb);
if self.total_mb <= Self::EXTENDED_SIXTEEN_CLASS_MB {
128
} else if total_gb <= 24 {
160
} else if total_gb <= 32 {
224
} else {
288
}
}
#[inline]
fn pressure_level_name(v: u8) -> &'static str {
match v {
x if x == PressureLevel::None as u8 => "None",
x if x == PressureLevel::Elevated as u8 => "Elevated",
x if x == PressureLevel::Critical as u8 => "Critical",
x if x == PressureLevel::Emergency as u8 => "Emergency",
_ => "?",
}
}
pub(crate) fn pressure_level_reported(&self, snap: &MemorySnapshot) -> PressureLevel {
let level = self.pressure_level(snap);
self.log_pressure_transition_if_changed(level, snap);
level
}
fn log_pressure_transition_if_changed(&self, level: PressureLevel, snap: &MemorySnapshot) {
let new = level as u8;
let prev = self.last_reported_pressure.swap(new, Ordering::Relaxed);
if prev == new {
return;
}
tracing::info!(
"MemoryGuard: pressure transition {} -> {} ({})",
Self::pressure_level_name(prev),
Self::pressure_level_name(new),
snap
);
if new >= (PressureLevel::Critical as u8) && prev < (PressureLevel::Critical as u8) {
#[cfg(all(not(target_os = "windows"), feature = "mimalloc"))]
unsafe {
libmimalloc_sys::mi_stats_print(std::ptr::null_mut());
}
}
}
pub(crate) fn pressure_level(&self, snap: &MemorySnapshot) -> PressureLevel {
let current = PressureLevel::from_u8(self.last_reported_pressure.load(Ordering::Relaxed));
let level = self.clamp_pressure_to_process_budget(
self.pressure_level_for(snap, current),
snap,
current,
);
Self::clamp_pressure_to_swap_state(level, snap)
}
fn clamp_pressure_to_swap_state(level: PressureLevel, snap: &MemorySnapshot) -> PressureLevel {
if snap.swap_total_mb == 0 {
return level; }
let our_swap_mb = snap.vm_swap_mb;
if our_swap_mb > 1024 {
return PressureLevel::Emergency;
}
if our_swap_mb > 256 {
return level.max(PressureLevel::Critical);
}
if our_swap_mb > 64 {
return level.max(PressureLevel::Elevated);
}
let ram_is_tight = snap.sys_avail_mb > 0 && snap.sys_avail_mb < 8192;
if !ram_is_tight {
return level;
}
let swap_free_pct = snap.swap_free_mb * 100 / snap.swap_total_mb.max(1);
if swap_free_pct < 5 {
PressureLevel::Emergency
} else if swap_free_pct < 15 {
level.max(PressureLevel::Critical)
} else if swap_free_pct < 35 {
level.max(PressureLevel::Elevated)
} else {
level
}
}
#[inline]
pub(crate) fn large_host_our_swap_counts(sys_avail_mb: u64) -> bool {
sys_avail_mb > 0 && sys_avail_mb < 32 * 1024
}
fn clamp_pressure_to_process_budget(
&self,
mut level: PressureLevel,
snap: &MemorySnapshot,
current: PressureLevel,
) -> PressureLevel {
let b = self.rss_budget_mb;
let r = if snap.rss_anon_mb > 0 {
snap.rss_anon_mb
} else {
snap.rss_mb
};
if b == 0 || r == 0 {
return level;
}
let emerg_line = if self.no_swap_at_boot {
b * 85 / 100
} else {
b
};
let emerg_exit = if self.no_swap_at_boot {
b * 80 / 100
} else {
b * 95 / 100
};
let crit_line = if self.no_swap_at_boot {
b * 75 / 100
} else {
b * 92 / 100
};
let elev_line = if self.no_swap_at_boot {
b * 65 / 100
} else {
b * 82 / 100
};
if r >= emerg_line {
return PressureLevel::Emergency;
}
if current == PressureLevel::Emergency && r >= emerg_exit {
return PressureLevel::Emergency;
}
if r >= crit_line {
level = level.max(PressureLevel::Critical);
} else if r >= elev_line {
level = level.max(PressureLevel::Elevated);
}
level
}
pub(crate) fn engine_avail_mb(&self) -> u64 {
const PIPELINE_RESERVE_MB: u64 = 6144;
const LEGACY_CACHE_CAP_MB: u64 = 4096;
let legacy = (self.utxo_cache_mb as u64).min(LEGACY_CACHE_CAP_MB);
let engine_mb = self
.rss_budget_mb
.saturating_sub(PIPELINE_RESERVE_MB)
.saturating_sub(legacy);
engine_mb.max(2048)
}
fn pressure_level_for(&self, snap: &MemorySnapshot, current: PressureLevel) -> PressureLevel {
let t = if snap.mem_total_mb > 0 {
snap.mem_total_mb
} else {
self.total_mb
};
let r = snap.rss_mb;
let a = snap.sys_avail_mb;
if r == 0 {
return PressureLevel::None;
}
if t <= 16 * 1024 {
let swap_used = snap.swap_used_mb();
let our_swap = snap.vm_swap_mb > 256; let crit_rss = self.crit_rss_threshold_mb; let rss_elev = (t * 30 / 100).max(2000); let rss_emerg = (t * 50 / 100).max(4000); let swap_elev_up = swap_used >= t * 5 / 100 && a > 0 && a < 4096 && our_swap;
let swap_crit_up = swap_used >= t * 12 / 100 && a > 0 && a < 3072 && our_swap;
let swap_emerg_up = swap_used >= t * 20 / 100 && a > 0 && a < 2048 && our_swap;
let swap_elev_dn = swap_used >= t * 5 / 100 && a > 0 && a < 4608 && our_swap;
let swap_crit_dn = swap_used >= t * 12 / 100 && a > 0 && a < 3584 && our_swap;
let swap_emerg_dn = swap_used >= t * 20 / 100 && a > 0 && a < 2560 && our_swap;
let emerg_up =
(r >= rss_emerg && a > 0 && a < 1024) || (a > 0 && a < 512) || swap_emerg_up;
let crit_up =
(r >= crit_rss && a > 0 && a < 1536) || (a > 0 && a < 768) || swap_crit_up;
let elev_up =
(r >= rss_elev && a > 0 && a < 2048) || (a > 0 && a < 1024) || swap_elev_up;
let emerg_dn = (a == 0 || a >= 768) && !swap_emerg_dn;
let crit_dn = (a == 0 || a >= 1024) && !swap_crit_dn;
let elev_dn = (a == 0 || a >= 1280) && !swap_elev_dn;
return match current {
PressureLevel::Emergency => {
if emerg_dn {
if crit_up {
PressureLevel::Critical
} else if elev_up {
PressureLevel::Elevated
} else {
PressureLevel::None
}
} else {
PressureLevel::Emergency
}
}
PressureLevel::Critical => {
if emerg_up {
PressureLevel::Emergency
} else if crit_dn {
if elev_up {
PressureLevel::Elevated
} else {
PressureLevel::None
}
} else {
PressureLevel::Critical
}
}
PressureLevel::Elevated => {
if emerg_up {
PressureLevel::Emergency
} else if crit_up {
PressureLevel::Critical
} else if elev_dn {
PressureLevel::None
} else {
PressureLevel::Elevated
}
}
PressureLevel::None => {
if emerg_up {
PressureLevel::Emergency
} else if crit_up {
PressureLevel::Critical
} else if elev_up {
PressureLevel::Elevated
} else {
PressureLevel::None
}
}
};
}
let avail_emerg_up: u64 = if t <= 24 * 1024 { 1536 } else { 768 };
let rss_emerg_pct_up: u64 = if t <= 24 * 1024 { 60 } else { 72 };
let avail_crit_up: u64 = if t <= 24 * 1024 { 1792 } else { 1024 };
let rss_crit_pct_up: u64 = if t <= 24 * 1024 { 55 } else { 65 };
let avail_elev_up: u64 = if t <= 24 * 1024 { 2048 } else { 1536 };
let rss_elev_pct_up: u64 = if t <= 24 * 1024 { 45 } else { 55 };
let avail_emerg_dn: u64 = avail_emerg_up + avail_emerg_up / 4;
let avail_crit_dn: u64 = avail_crit_up + avail_crit_up / 4;
let avail_elev_dn: u64 = avail_elev_up + avail_elev_up / 4;
let rss_emerg_pct_dn: u64 = rss_emerg_pct_up.saturating_sub(5);
let rss_crit_pct_dn: u64 = rss_crit_pct_up.saturating_sub(5);
let rss_elev_pct_dn: u64 = rss_elev_pct_up.saturating_sub(5);
let r_for_pct = if snap.rss_anon_mb > 0 {
snap.rss_anon_mb
} else {
r
};
let swap_pct = if snap.swap_total_mb > 0 {
snap.swap_used_mb() * 100 / snap.swap_total_mb
} else {
0
};
let our_swap_active = snap.vm_swap_mb > 256;
let our_swap_counts = our_swap_active && Self::large_host_our_swap_counts(a);
let swap_crit_up = swap_pct >= 90 && our_swap_counts;
let swap_emerg_up = swap_pct >= 98 && our_swap_counts;
let swap_crit_dn = swap_pct >= 85 && our_swap_counts;
let swap_emerg_dn = swap_pct >= 95 && our_swap_counts;
let sys_swap_full = snap.swap_total_mb >= 1024 && swap_pct >= 95;
let sys_swap_full_dn = snap.swap_total_mb >= 1024 && swap_pct >= 85;
let emerg_up = (a > 0 && a < avail_emerg_up)
|| r_for_pct > t * rss_emerg_pct_up / 100
|| swap_emerg_up;
let crit_up =
(a > 0 && a < avail_crit_up) || r_for_pct > t * rss_crit_pct_up / 100 || swap_crit_up;
let sys_swap_elev = sys_swap_full && a < 32 * 1024;
let sys_swap_elev_dn = sys_swap_full_dn && a < 32 * 1024;
let sys_swap_crit = sys_swap_full && a < 32 * 1024;
let sys_swap_crit_dn = sys_swap_full_dn && a < 32 * 1024;
let elev_up =
(a > 0 && a < avail_elev_up) || r_for_pct > t * rss_elev_pct_up / 100 || sys_swap_elev;
let crit_up = crit_up || sys_swap_crit;
let emerg_dn = (a == 0 || a >= avail_emerg_dn)
&& r_for_pct <= t * rss_emerg_pct_dn / 100
&& !swap_emerg_dn;
let crit_dn = (a == 0 || a >= avail_crit_dn)
&& r_for_pct <= t * rss_crit_pct_dn / 100
&& !swap_crit_dn
&& !sys_swap_crit_dn;
let elev_dn = (a == 0 || a >= avail_elev_dn)
&& r_for_pct <= t * rss_elev_pct_dn / 100
&& !sys_swap_elev_dn;
match current {
PressureLevel::Emergency => {
if emerg_dn {
if crit_up {
PressureLevel::Critical
} else if elev_up {
PressureLevel::Elevated
} else {
PressureLevel::None
}
} else {
PressureLevel::Emergency
}
}
PressureLevel::Critical => {
if emerg_up {
PressureLevel::Emergency
} else if crit_dn {
if elev_up {
PressureLevel::Elevated
} else {
PressureLevel::None
}
} else {
PressureLevel::Critical
}
}
PressureLevel::Elevated => {
if emerg_up {
PressureLevel::Emergency
} else if crit_up {
PressureLevel::Critical
} else if elev_dn {
PressureLevel::None
} else {
PressureLevel::Elevated
}
}
PressureLevel::None => {
if emerg_up {
PressureLevel::Emergency
} else if crit_up {
PressureLevel::Critical
} else if elev_up {
PressureLevel::Elevated
} else {
PressureLevel::None
}
}
}
}
fn adjust_max_ahead_live(&self, snap: &MemorySnapshot, live: &AtomicU64, nominal: u64) {
let cur = live.load(Ordering::Relaxed);
let nominal = nominal.max(64);
let level = self.pressure_level(snap);
let tight_ahead = self.total_mb <= 16 * 1024;
match level {
PressureLevel::Emergency => {
let target = if tight_ahead {
(nominal / 6).max(48)
} else {
(nominal / 4).max(64)
};
if cur > target {
tracing::warn!(
"MemoryGuard: EMERGENCY — download ahead {} → {} ({})",
cur,
target,
snap
);
live.store(target, Ordering::Relaxed);
}
}
PressureLevel::Critical => {
let target = if tight_ahead {
(nominal / 4).max(64)
} else {
(nominal / 3).max(96)
};
if cur > target {
tracing::warn!(
"MemoryGuard: CRITICAL — download ahead {} → {} ({})",
cur,
target,
snap
);
live.store(target, Ordering::Relaxed);
}
}
PressureLevel::Elevated => {
let floor = (nominal / 2).max(128);
let target = (cur * 3 / 4).max(floor);
if cur > target {
tracing::info!(
"MemoryGuard: elevated — download ahead {} → {} ({})",
cur,
target,
snap
);
live.store(target, Ordering::Relaxed);
}
}
PressureLevel::None => {
let tier_cap = Self::tier_max_download_ahead_blocks(self.total_mb);
let ceil = if self.total_mb <= Self::EXTENDED_SIXTEEN_CLASS_MB {
if snap.sys_avail_mb > 7_000 {
nominal.saturating_mul(2).min(tier_cap)
} else if snap.sys_avail_mb > 5_000 {
(nominal * 3 / 2).min(tier_cap.saturating_mul(3) / 4)
} else {
nominal
}
} else {
nominal.saturating_mul(2)
};
if cur < ceil {
let nxt = cur.saturating_add(16).min(ceil);
live.store(nxt, Ordering::Relaxed);
}
}
}
}
pub(crate) fn should_flush(
&mut self,
max_ahead_live: Option<(&AtomicU64, u64)>,
) -> PressureLevel {
let now = Instant::now();
let elapsed = now.duration_since(self.last_rss_check);
let cached = PressureLevel::from_u8(self.last_reported_pressure.load(Ordering::Relaxed));
if elapsed < Duration::from_millis(150) && cached < PressureLevel::Emergency {
return cached;
}
self.last_rss_check = now;
let snap = self.memory_snapshot();
publish_ibd_rss_anon_mb(snap.rss_anon_mb);
if let Some((live, nominal)) = max_ahead_live {
self.adjust_max_ahead_live(&snap, live, nominal);
}
if snap.rss_mb == 0 {
return PressureLevel::None;
}
let level = self.pressure_level(&snap);
self.log_pressure_transition_if_changed(level, &snap);
level
}
#[cfg(test)]
pub(crate) fn test_seed_pressure_level(&mut self, level: PressureLevel) {
self.last_reported_pressure
.store(level as u8, Ordering::Relaxed);
self.last_rss_check = Instant::now();
}
pub(crate) fn compute_adaptive_cache_cap(&mut self) -> Option<usize> {
let nominal = self.utxo_max_entries;
if nominal == usize::MAX {
return None;
}
let rss_mb = self.current_rss_mb();
if rss_mb == 0 || self.rss_budget_mb == 0 {
return None;
}
let budget = self.rss_budget_mb;
let stored = self.last_adaptive_cap_entries.load(Ordering::Relaxed);
let current = if stored == 0 { nominal } else { stored };
let ratio_x1000 = (rss_mb as u128 * 1000 / budget.max(1) as u128) as u64;
let is_emergency = ratio_x1000 >= 1000;
if !is_emergency {
let mut last = self
.last_adaptive_cap_check
.lock()
.expect("adaptive_cap_check");
if last.elapsed() < Duration::from_secs(2) {
return None;
}
*last = Instant::now();
}
let hard_floor = (nominal / 4).max(256 * 1024);
let is_above_shrink_threshold = ratio_x1000 >= 800;
if is_above_shrink_threshold && !is_emergency {
let prev = self
.above_threshold_consecutive
.fetch_add(1, Ordering::Relaxed);
if prev < 1 {
return None;
}
} else if !is_emergency {
self.above_threshold_consecutive.store(0, Ordering::Relaxed);
}
let target = if is_emergency {
let scaled =
(current as u128 * (budget as u128 * 650 / 1000) / rss_mb.max(1) as u128) as usize;
scaled.max(hard_floor)
} else if ratio_x1000 >= 900 {
(current * 7 / 10).max(hard_floor)
} else if ratio_x1000 >= 800 {
const SHRINK_COOLDOWN_SECS: u64 = 20;
{
let last_shrink = self
.last_adaptive_cap_shrink
.lock()
.expect("last_shrink lock");
if last_shrink.elapsed().as_secs() < SHRINK_COOLDOWN_SECS {
return None;
}
}
(current * 90 / 100).max(hard_floor)
} else if ratio_x1000 < 600 && current < nominal {
((current * 125 / 100).min(nominal)).max(hard_floor)
} else if ratio_x1000 < 700 && current < nominal {
((current * 115 / 100).min(nominal)).max(hard_floor)
} else if ratio_x1000 < 800 && current < nominal {
((current * 108 / 100).min(nominal)).max(hard_floor)
} else {
return None;
};
if target == current {
return None;
}
let delta = target.abs_diff(current);
if delta < (current / 33).max(8 * 1024) {
return None;
}
if target < current {
let mut last_shrink = self
.last_adaptive_cap_shrink
.lock()
.expect("last_shrink lock");
*last_shrink = Instant::now();
self.above_threshold_consecutive.store(0, Ordering::Relaxed);
}
self.last_adaptive_cap_entries
.store(target, Ordering::Relaxed);
tracing::info!(
"MemoryGuard: adaptive cache cap {} -> {} entries (rss={}MB / budget={}MB = {}.{}%, nominal={})",
current,
target,
rss_mb,
budget,
ratio_x1000 / 10,
ratio_x1000 % 10,
nominal,
);
Some(target)
}
pub(crate) fn current_rss_mb(&mut self) -> u64 {
#[cfg(target_os = "linux")]
{
if proc_read_file("/proc/self/status", &mut self.proc_status_buf) {
return proc_anon_rss_mb_from_status(&self.proc_status_buf);
}
0
}
#[cfg(all(not(target_os = "linux"), feature = "sysinfo"))]
{
use sysinfo::Pid;
let pid = Pid::from(std::process::id() as usize);
self.sys.refresh_process(pid);
self.sys
.process(pid)
.map(|p| p.memory() / (1024 * 1024))
.unwrap_or(0)
}
#[cfg(all(not(target_os = "linux"), not(feature = "sysinfo")))]
0u64
}
#[cfg(target_os = "linux")]
pub(crate) fn memory_snapshot(&mut self) -> MemorySnapshot {
let mut snap = MemorySnapshot::default();
if proc_read_file("/proc/self/status", &mut self.proc_status_buf) {
proc_parse_status_into(&self.proc_status_buf, &mut snap);
}
if proc_read_file("/proc/meminfo", &mut self.proc_meminfo_buf) {
proc_parse_meminfo_into(&self.proc_meminfo_buf, &mut snap);
}
snap
}
#[cfg(not(target_os = "linux"))]
pub(crate) fn memory_snapshot(&self) -> MemorySnapshot {
MemorySnapshot::default()
}
pub(crate) fn buffer_limit(&self, current_height: u64) -> usize {
Self::buffer_limit_for(self.block_buffer_base, self.total_mb, current_height)
}
pub(crate) fn buffer_limit_for(
block_buffer_base: usize,
total_mb: u64,
current_height: u64,
) -> usize {
let scale = match current_height {
0..=100_000 => 100,
100_001..=300_000 => 50,
300_001..=480_000 => 33,
480_001..=700_000 => 20,
_ => 12,
};
let min_buf = if total_mb <= 16 * 1024 { 50 } else { 200 };
(block_buffer_base * scale / 100).clamp(min_buf, 2_000)
}
pub(crate) fn feeder_bytes_limit_for_height(&self, current_height: u64) -> usize {
Self::feeder_bytes_for(
self.feeder_buffer_bytes_limit,
self.block_buffer_base,
self.total_mb,
current_height,
)
}
pub(crate) fn feeder_bytes_for(
feeder_buffer_bytes_limit: usize,
block_buffer_base: usize,
total_mb: u64,
current_height: u64,
) -> usize {
let tier = match current_height {
0..=100_000 => 100u64,
100_001..=300_000 => 72,
300_001..=480_000 => 58,
480_001..=700_000 => 48,
_ => 40,
};
let scaled = (feeder_buffer_bytes_limit as u64 * tier / 100) as usize;
let buf = Self::buffer_limit_for(block_buffer_base, total_mb, current_height);
let cap_by_est_blocks = buf.saturating_mul(900_000);
scaled.min(cap_by_est_blocks).max(32 * 1024 * 1024)
}
pub(crate) fn memory_diag(&mut self) -> Option<(u64, u64)> {
#[cfg(feature = "sysinfo")]
{
use sysinfo::Pid;
let pid = Pid::from(std::process::id() as usize);
self.sys.refresh_memory();
self.sys.refresh_process(pid);
let rss_mb = self
.sys
.process(pid)
.map(|p| p.memory() / (1024 * 1024))
.unwrap_or(0);
let avail_mb = self.sys.available_memory() / (1024 * 1024);
Some((rss_mb, avail_mb))
}
#[cfg(not(feature = "sysinfo"))]
None
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum MemReportAccountedVerdict {
AccountableEngine,
ResidualMystery,
Inconclusive,
}
#[inline]
pub(crate) fn classify_mem_report_accounted(
anon_mb: u64,
engine_total_mb: u64,
accounted_total_mb: u64,
post_rayon_unexplained_mb: u64,
) -> MemReportAccountedVerdict {
if anon_mb == 0 || accounted_total_mb == 0 {
return MemReportAccountedVerdict::Inconclusive;
}
let mystery_floor = 1024u64.max(anon_mb / 10);
if post_rayon_unexplained_mb > mystery_floor {
return MemReportAccountedVerdict::ResidualMystery;
}
let residual_ok = post_rayon_unexplained_mb <= 256u64.max(anon_mb / 20);
let engine_dominant = engine_total_mb.saturating_mul(2) >= accounted_total_mb;
if residual_ok && engine_dominant {
return MemReportAccountedVerdict::AccountableEngine;
}
MemReportAccountedVerdict::Inconclusive
}
#[inline]
pub(crate) fn c2_working_set_track_ok(verdict: MemReportAccountedVerdict) -> bool {
matches!(verdict, MemReportAccountedVerdict::AccountableEngine)
}
#[cfg(test)]
mod memory_tier_tests {
use super::{
MemReportAccountedVerdict, MemoryGuard, PressureLevel, ROCKSDB_PIPELINE_RESERVE_MB,
WorkloadClass, c2_working_set_track_ok, classify_mem_report_accounted,
emergency_entry_anon_mb, ibd_pressure_is_emergency, ibd_pressure_level_snapshot,
publish_ibd_pressure, reset_ibd_pressure_on_session_end, stale_emergency_step_down_level,
};
use crate::storage::database::DatabaseBackend;
#[test]
fn mem_report_accounted_vs_residual_c2_gate() {
let h3 = classify_mem_report_accounted(
15_500, 15_100, 15_400, 0, );
assert_eq!(h3, MemReportAccountedVerdict::AccountableEngine);
assert!(c2_working_set_track_ok(h3));
let mystery = classify_mem_report_accounted(12_000, 4_000, 5_000, 3_500);
assert_eq!(mystery, MemReportAccountedVerdict::ResidualMystery);
assert!(!c2_working_set_track_ok(mystery));
let mid = classify_mem_report_accounted(8_000, 1_000, 4_000, 200);
assert_eq!(mid, MemReportAccountedVerdict::Inconclusive);
assert!(!c2_working_set_track_ok(mid));
}
#[test]
fn shared_ninety_two_gb_envelope_and_pending_cap() {
let total = 94_162_u64;
let avail = 52_449_u64;
let workload = MemoryGuard::detect_workload_class(total, avail, false);
assert_eq!(workload, WorkloadClass::Shared);
let rss = MemoryGuard::compute_rss_budget_mb(total, avail, workload);
assert!(rss >= 18_000 && rss <= 26_000, "rss_budget={rss}");
let utxo_cache = ((rss * 45 / 100) as usize).min(4096);
assert_eq!(utxo_cache, 4096);
let pending = MemoryGuard::nominal_max_pending_ops(
total,
rss,
utxo_cache,
2_000_000,
DatabaseBackend::RocksDB,
);
assert!(
pending >= 4_000_000 && pending <= 8_000_000,
"pending={pending}"
);
}
#[test]
fn heed3_pipeline_reserve_and_flush_interval() {
assert_eq!(
MemoryGuard::pipeline_reserve_mb(DatabaseBackend::Heed3),
512
);
assert_eq!(
MemoryGuard::pipeline_reserve_mb(DatabaseBackend::RocksDB),
ROCKSDB_PIPELINE_RESERVE_MB
);
assert_eq!(
MemoryGuard::storage_flush_interval_base(64, DatabaseBackend::Heed3),
200
);
assert_eq!(
MemoryGuard::storage_flush_interval_base(64, DatabaseBackend::RocksDB),
2000
);
assert!(
MemoryGuard::pipeline_reserve_mb(DatabaseBackend::Heed3)
< MemoryGuard::pipeline_reserve_mb(DatabaseBackend::RocksDB)
);
}
#[test]
fn dedicated_workload_gets_higher_envelope_than_shared() {
let total = 94_162_u64;
let rss_shared = MemoryGuard::compute_rss_budget_mb(total, 52_449, WorkloadClass::Shared);
let rss_dedicated =
MemoryGuard::compute_rss_budget_mb(total, 90_000, WorkloadClass::Dedicated);
assert!(rss_dedicated > rss_shared);
assert_eq!(
MemoryGuard::detect_workload_class(total, 90_000, false),
WorkloadClass::Dedicated
);
assert_eq!(
MemoryGuard::detect_workload_class(total, 52_449, true),
WorkloadClass::Shared
);
}
#[test]
fn engine_avail_mb_formula_not_clamped_to_boot_avail() {
let rss_budget_mb = 37_664_u64;
let boot_avail_mb = 10_300_u64;
let engine_mb = rss_budget_mb.saturating_sub(6144).saturating_sub(4096);
assert!(engine_mb > 20_000, "engine_mb={engine_mb}");
let old_capped = (rss_budget_mb * 28 / 100)
.max(2048)
.min(engine_mb)
.min(boot_avail_mb);
assert_eq!(old_capped, boot_avail_mb, "old formula forced age-3 sizing");
let new_hint = engine_mb.max(2048);
assert_eq!(new_hint, 27_424, "new formula uses RSS budget headroom");
}
#[test]
fn reset_ibd_pressure_on_session_end_clears_emergency_latch() {
publish_ibd_pressure(PressureLevel::Emergency);
assert!(ibd_pressure_is_emergency());
reset_ibd_pressure_on_session_end();
assert!(!ibd_pressure_is_emergency());
assert_eq!(ibd_pressure_level_snapshot(), PressureLevel::None);
}
#[test]
fn stale_emergency_not_cleared_at_live_entry_threshold_no_swap() {
let budget = 37_664_u64;
let emerg = emergency_entry_anon_mb(budget, true);
assert_eq!(emerg, 32_014);
assert!(stale_emergency_step_down_level(32_074, budget, true).is_none());
assert!(stale_emergency_step_down_level(32_014, budget, true).is_none());
assert_eq!(
stale_emergency_step_down_level(31_000, budget, true),
Some(PressureLevel::Critical)
);
assert_eq!(
stale_emergency_step_down_level(20_000, budget, true),
Some(PressureLevel::None)
);
}
#[test]
fn zeus_boot_with_dedicated_config_uses_shared_envelope() {
let total = 94_162_u64;
let avail = 64_865_u64; assert_eq!(
MemoryGuard::detect_workload_class(total, avail, true),
WorkloadClass::Shared
);
let rss = MemoryGuard::compute_rss_budget_mb(total, avail, WorkloadClass::Shared);
assert!(
rss >= 18_000 && rss <= 28_000,
"rss_budget={rss} should be ~23 GiB Shared envelope"
);
}
#[test]
fn no_swap_caps_dedicated_budget() {
let total = 94_162_u64;
let avail = 64_865_u64;
let dedicated = MemoryGuard::compute_rss_budget_mb(total, avail, WorkloadClass::Dedicated);
assert!(dedicated <= 47_081, "dedicated={dedicated}");
let no_swap_cap = avail.saturating_sub(8192).min(total * 40 / 100);
assert_eq!(no_swap_cap, 37_664);
}
#[test]
fn extended_sixteen_class_gets_tight_download_ahead_cap() {
assert_eq!(MemoryGuard::tier_max_download_ahead_blocks(15921), 320);
assert_eq!(MemoryGuard::tier_max_download_ahead_blocks(18 * 1024), 320);
assert_eq!(
MemoryGuard::tier_max_download_ahead_blocks(18 * 1024 + 1),
512
);
}
#[test]
fn r24b_ample_avail_does_not_count_process_swap_as_emergency() {
assert!(
!MemoryGuard::large_host_our_swap_counts(63_509),
"r24b sys_avail=63G — vLLM-filled zram is not OOM"
);
assert!(
!MemoryGuard::large_host_our_swap_counts(32 * 1024),
"32 GiB avail is the gate, not under it"
);
assert!(
MemoryGuard::large_host_our_swap_counts(32 * 1024 - 1),
"just under 32 GiB still counts"
);
assert!(
MemoryGuard::large_host_our_swap_counts(8_192),
"tight RAM + our swap still Emergency-eligible"
);
assert!(
!MemoryGuard::large_host_our_swap_counts(0),
"unknown avail must not trip"
);
}
}