use crate::budget::current_rss_mb;
use std::sync::atomic::{AtomicI64, AtomicU64, Ordering};
use std::time::{Duration, Instant};
static LAST_ACTIVITY_NANOS: AtomicU64 = AtomicU64::new(0);
static LAST_RSS_MB: AtomicI64 = AtomicI64::new(0);
fn gc_max_rss_mb() -> u64 {
std::env::var("LEANKG_GC_MAX_RSS_MB")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(4_096)
}
fn gc_idle_after_secs() -> u64 {
std::env::var("LEANKG_GC_IDLE_AFTER_SECS")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(60)
}
fn gc_poll_secs() -> u64 {
std::env::var("LEANKG_GC_POLL_SECS")
.ok()
.and_then(|v| v.parse().ok())
.filter(|n| *n > 0)
.unwrap_or(10)
}
const FORCE_TRIM_COOLDOWN_SECS: u64 = 30;
pub type ReleaseFn = Box<dyn Fn() -> bool + Send + Sync + 'static>;
pub struct MemoryGuard {
started: Instant,
last_check: Instant,
last_force_trim: Instant,
last_idle_trim_activity: u64,
release_fn: Option<ReleaseFn>,
}
impl MemoryGuard {
pub fn new(release_fn: Option<ReleaseFn>) -> Self {
let now = Instant::now();
Self::record_activity();
Self {
started: now,
last_check: now,
last_force_trim: now,
last_idle_trim_activity: LAST_ACTIVITY_NANOS.load(Ordering::Relaxed),
release_fn,
}
}
pub fn poll_interval() -> Duration {
Duration::from_secs(gc_poll_secs())
}
pub fn touch() {
Self::record_activity();
}
fn record_activity() {
let now_nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as u64)
.unwrap_or(0);
LAST_ACTIVITY_NANOS.store(now_nanos, Ordering::Relaxed);
}
fn activity_nanos() -> u64 {
LAST_ACTIVITY_NANOS.load(Ordering::Relaxed)
}
fn idle_secs() -> u64 {
let last = Self::activity_nanos();
if last == 0 {
return 0;
}
let now_nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as u64)
.unwrap_or(0);
(now_nanos.saturating_sub(last)) / 1_000_000_000
}
pub fn rss_mb() -> u64 {
LAST_RSS_MB.load(Ordering::Relaxed).max(0) as u64
}
pub fn tick(&mut self) -> GcAction {
let now = Instant::now();
if now.duration_since(self.last_check) < Duration::from_secs(gc_poll_secs()) {
return GcAction::Skipped;
}
self.last_check = now;
let rss = match current_rss_mb() {
Ok(rss) => {
LAST_RSS_MB.store(rss as i64, Ordering::Relaxed);
rss
}
Err(_) => Self::rss_mb(),
};
let idle = Self::idle_secs();
let activity = Self::activity_nanos();
if rss >= gc_max_rss_mb() {
if now.duration_since(self.last_force_trim)
< Duration::from_secs(FORCE_TRIM_COOLDOWN_SECS)
{
return GcAction::Skipped;
}
self.last_force_trim = now;
if self.run_release() {
eprintln!(
"leankg::gc: RSS {} MB >= max {} MB; force-trimmed caches",
rss,
gc_max_rss_mb()
);
return GcAction::ForceTrim { rss_mb: rss };
}
return GcAction::NoOp { rss_mb: rss };
}
if idle >= gc_idle_after_secs() && activity != self.last_idle_trim_activity {
self.last_idle_trim_activity = activity;
if self.run_release() {
eprintln!(
"leankg::gc: idle for {}s; trimmed caches (RSS {} MB)",
idle, rss
);
return GcAction::IdleTrim {
idle_secs: idle,
rss_mb: rss,
};
}
return GcAction::NoOp { rss_mb: rss };
}
GcAction::NoOp { rss_mb: rss }
}
fn run_release(&self) -> bool {
match &self.release_fn {
Some(f) => f(),
None => false,
}
}
pub fn uptime(&self) -> Duration {
self.started.elapsed()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum GcAction {
Skipped,
NoOp { rss_mb: u64 },
IdleTrim { idle_secs: u64, rss_mb: u64 },
ForceTrim { rss_mb: u64 },
}
pub fn trim_heap() -> bool {
#[cfg(target_os = "linux")]
unsafe {
extern "C" {
fn malloc_trim(pad: usize) -> i32;
}
malloc_trim(0) == 1
}
#[cfg(not(target_os = "linux"))]
{
false
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
#[test]
fn guard_touch_resets_idle() {
MemoryGuard::touch();
assert_eq!(MemoryGuard::idle_secs(), 0);
}
#[test]
fn rss_mb_returns_nonzero_or_zero() {
let _ = MemoryGuard::rss_mb();
}
#[test]
fn trim_heap_does_not_panic() {
let _ = trim_heap();
}
#[test]
fn tick_returns_a_variant() {
let mut g = MemoryGuard::new(None);
let _ = g.tick();
}
#[test]
fn tick_skips_when_poll_window_not_elapsed() {
let mut g = MemoryGuard::new(None);
g.last_check = Instant::now();
let action = g.tick();
assert_eq!(action, GcAction::Skipped);
}
#[test]
fn poll_interval_is_positive() {
assert!(MemoryGuard::poll_interval() >= Duration::from_secs(1));
}
#[test]
fn idle_trim_runs_at_most_once_per_activity_period() {
let called = std::sync::Arc::new(AtomicUsize::new(0));
let called2 = called.clone();
let mut g = MemoryGuard::new(Some(Box::new(move || {
called2.fetch_add(1, Ordering::Relaxed);
true
})));
g.last_check = Instant::now() - Duration::from_secs(gc_poll_secs() + 1);
let past = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as u64)
.unwrap_or(0)
.saturating_sub((gc_idle_after_secs() + 5) * 1_000_000_000);
LAST_ACTIVITY_NANOS.store(past, Ordering::Relaxed);
g.last_idle_trim_activity = 0;
let first = g.tick();
assert!(
matches!(first, GcAction::IdleTrim { .. }),
"expected IdleTrim, got {:?}",
first
);
assert_eq!(called.load(Ordering::Relaxed), 1);
g.last_check = Instant::now() - Duration::from_secs(gc_poll_secs() + 1);
let second = g.tick();
assert!(
matches!(second, GcAction::NoOp { .. } | GcAction::Skipped),
"expected NoOp/Skipped after once-per-idle, got {:?}",
second
);
assert_eq!(called.load(Ordering::Relaxed), 1);
MemoryGuard::touch();
let past2 = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as u64)
.unwrap_or(0)
.saturating_sub((gc_idle_after_secs() + 5) * 1_000_000_000);
LAST_ACTIVITY_NANOS.store(past2, Ordering::Relaxed);
g.last_check = Instant::now() - Duration::from_secs(gc_poll_secs() + 1);
let third = g.tick();
assert!(
matches!(third, GcAction::IdleTrim { .. }),
"expected IdleTrim after new activity, got {:?}",
third
);
assert_eq!(called.load(Ordering::Relaxed), 2);
}
#[test]
fn release_returning_false_is_noop_not_idle_trim() {
let mut g = MemoryGuard::new(Some(Box::new(|| false)));
g.last_check = Instant::now() - Duration::from_secs(gc_poll_secs() + 1);
let past = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as u64)
.unwrap_or(0)
.saturating_sub((gc_idle_after_secs() + 5) * 1_000_000_000);
LAST_ACTIVITY_NANOS.store(past, Ordering::Relaxed);
g.last_idle_trim_activity = 0;
let action = g.tick();
assert!(
matches!(action, GcAction::NoOp { .. }),
"expected NoOp when release returns false, got {:?}",
action
);
assert_eq!(g.last_idle_trim_activity, past);
}
}