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())
.unwrap_or(10)
}
pub type ReleaseFn = Box<dyn Fn() + Send + Sync + 'static>;
pub struct MemoryGuard {
started: Instant,
last_check: Instant,
last_idle_trim: Instant,
last_force_trim: Instant,
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_idle_trim: now,
last_force_trim: now,
release_fn,
}
}
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 idle_secs() -> u64 {
let last = LAST_ACTIVITY_NANOS.load(Ordering::Relaxed);
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;
if let Ok(rss) = current_rss_mb() {
LAST_RSS_MB.store(rss as i64, Ordering::Relaxed);
}
let rss = current_rss_mb().unwrap_or(0);
let idle = Self::idle_secs();
if rss >= gc_max_rss_mb() {
if now.duration_since(self.last_force_trim) >= Duration::from_secs(30) {
self.last_force_trim = now;
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::Skipped;
}
if idle >= gc_idle_after_secs()
&& now.duration_since(self.last_idle_trim) >= Duration::from_secs(30)
{
self.last_idle_trim = now;
self.run_release();
eprintln!(
"leankg::gc: idle for {}s; trimmed caches (RSS {} MB)",
idle, rss
);
return GcAction::IdleTrim {
idle_secs: idle,
rss_mb: rss,
};
}
GcAction::NoOp { rss_mb: rss }
}
fn run_release(&self) {
if let Some(f) = &self.release_fn {
f();
}
}
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_runs_release_on_idle() {
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);
})));
g.last_idle_trim = Instant::now() - Duration::from_secs(60);
assert_eq!(called.load(Ordering::Relaxed), 0);
}
#[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);
}
}