use std::time::Instant;
const SIGNIFICANT_PCT: u64 = 2;
const RELEASE_PCT: u64 = 1;
const SETTLE_GROWTH_PCT: u64 = 1;
const SETTLE_STREAK: u32 = 4;
const SETTLE_DEADLINE_SAMPLES: u32 = 240;
#[derive(Clone, Copy, PartialEq, Eq, Debug, Default)]
pub enum DriftVerdict {
#[default]
Settled,
Heap,
OutsideHeap,
Both,
}
impl DriftVerdict {
pub fn label(self) -> &'static str {
match self {
DriftVerdict::Settled => "settled",
DriftVerdict::Heap => "heap",
DriftVerdict::OutsideHeap => "outside-heap",
DriftVerdict::Both => "heap and outside-heap",
}
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub struct MemoryDrift {
pub heap_growth_bytes: i64,
pub outside_heap_growth_bytes: i64,
pub window_secs: u64,
pub verdict: DriftVerdict,
}
#[derive(Clone, Copy, Debug)]
struct Baseline {
at: Instant,
rss_bytes: u64,
heap_live_bytes: u64,
}
#[derive(Debug, Default)]
pub(crate) struct DriftTracker {
baseline: Option<Baseline>,
last_rss: Option<u64>,
flat_streak: u32,
samples_before_baseline: u32,
heap_moved: bool,
outside_heap_moved: bool,
}
impl DriftTracker {
pub(crate) fn sample(
&mut self,
rss_bytes: u64,
heap_live_bytes: u64,
budget_bytes: u64,
) -> Option<MemoryDrift> {
self.sample_at(rss_bytes, heap_live_bytes, budget_bytes, Instant::now())
}
fn sample_at(
&mut self,
rss_bytes: u64,
heap_live_bytes: u64,
budget_bytes: u64,
now: Instant,
) -> Option<MemoryDrift> {
let Some(base) = self.baseline else {
self.settle(rss_bytes, heap_live_bytes, now);
return None;
};
let heap_growth = growth(heap_live_bytes, base.heap_live_bytes);
let outside_heap_growth = growth(rss_bytes, base.rss_bytes) - heap_growth;
self.heap_moved = moved(heap_growth, budget_bytes, self.heap_moved);
self.outside_heap_moved = moved(outside_heap_growth, budget_bytes, self.outside_heap_moved);
Some(MemoryDrift {
heap_growth_bytes: heap_growth,
outside_heap_growth_bytes: outside_heap_growth,
window_secs: now.saturating_duration_since(base.at).as_secs(),
verdict: verdict(self.heap_moved, self.outside_heap_moved),
})
}
fn settle(&mut self, rss_bytes: u64, heap_live_bytes: u64, now: Instant) {
self.samples_before_baseline = self.samples_before_baseline.saturating_add(1);
let flat = self.last_rss.is_some_and(|prev| {
rss_bytes.saturating_sub(prev) <= prev.saturating_mul(SETTLE_GROWTH_PCT) / 100
});
self.flat_streak = if flat { self.flat_streak + 1 } else { 0 };
self.last_rss = Some(rss_bytes);
if self.flat_streak >= SETTLE_STREAK
|| self.samples_before_baseline >= SETTLE_DEADLINE_SAMPLES
{
self.baseline = Some(Baseline {
at: now,
rss_bytes,
heap_live_bytes,
});
}
}
}
fn growth(now: u64, base: u64) -> i64 {
(now as i128 - base as i128).clamp(i64::MIN as i128, i64::MAX as i128) as i64
}
fn moved(growth: i64, budget_bytes: u64, was_moved: bool) -> bool {
let pct = if was_moved {
RELEASE_PCT
} else {
SIGNIFICANT_PCT
};
let threshold = budget_bytes.saturating_mul(pct) / 100;
threshold > 0 && growth > 0 && growth as u64 > threshold
}
fn verdict(heap_moved: bool, outside_heap_moved: bool) -> DriftVerdict {
match (heap_moved, outside_heap_moved) {
(true, true) => DriftVerdict::Both,
(true, false) => DriftVerdict::Heap,
(false, true) => DriftVerdict::OutsideHeap,
(false, false) => DriftVerdict::Settled,
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
const MIB: u64 = 1024 * 1024;
const BUDGET: u64 = 1000 * MIB;
const SIGNIFICANT: u64 = 20 * MIB;
fn settled(rss: u64, heap: u64, at: Instant) -> DriftTracker {
let mut t = DriftTracker::default();
for _ in 0..=SETTLE_STREAK {
assert_eq!(t.sample_at(rss, heap, BUDGET, at), None);
}
assert!(
t.baseline.is_some(),
"a run of steady samples settles the baseline"
);
t
}
#[test]
fn a_climbing_rss_against_a_flat_heap_reads_as_outside_the_heap() {
let start = Instant::now();
let mut t = settled(2000 * MIB, 400 * MIB, start);
let d = t
.sample_at(
2400 * MIB,
400 * MIB,
BUDGET,
start + Duration::from_secs(3600),
)
.expect("the baseline is captured");
assert_eq!(d.verdict, DriftVerdict::OutsideHeap);
assert_eq!(d.heap_growth_bytes, 0);
assert_eq!(d.outside_heap_growth_bytes, (400 * MIB) as i64);
assert_eq!(d.window_secs, 3600);
}
#[test]
fn a_climbing_heap_carrying_the_rss_reads_as_ours() {
let start = Instant::now();
let mut t = settled(2000 * MIB, 400 * MIB, start);
let d = t
.sample_at(
2400 * MIB,
800 * MIB,
BUDGET,
start + Duration::from_secs(60),
)
.expect("the baseline is captured");
assert_eq!(d.verdict, DriftVerdict::Heap);
assert_eq!(d.heap_growth_bytes, (400 * MIB) as i64);
assert_eq!(d.outside_heap_growth_bytes, 0);
}
#[test]
fn both_terms_growing_are_reported_as_both() {
let start = Instant::now();
let mut t = settled(2000 * MIB, 400 * MIB, start);
let d = t
.sample_at(2600 * MIB, 700 * MIB, BUDGET, start)
.expect("the baseline is captured");
assert_eq!(d.verdict, DriftVerdict::Both);
assert_eq!(d.heap_growth_bytes, (300 * MIB) as i64);
assert_eq!(d.outside_heap_growth_bytes, (300 * MIB) as i64);
}
#[test]
fn movement_under_the_threshold_stays_settled() {
let start = Instant::now();
let mut t = settled(2000 * MIB, 400 * MIB, start);
let d = t
.sample_at(
2000 * MIB + SIGNIFICANT - 1,
400 * MIB + SIGNIFICANT - 1,
BUDGET,
start,
)
.expect("the baseline is captured");
assert_eq!(d.verdict, DriftVerdict::Settled);
}
#[test]
fn a_growth_hovering_at_the_threshold_holds_its_reading() {
let start = Instant::now();
let mut t = settled(2000 * MIB, 400 * MIB, start);
let rss_at = |outside: u64| 2000 * MIB + outside;
let d = t
.sample_at(rss_at(SIGNIFICANT + MIB), 400 * MIB, BUDGET, start)
.expect("the baseline is captured");
assert_eq!(d.verdict, DriftVerdict::OutsideHeap);
for outside in [SIGNIFICANT - MIB, SIGNIFICANT + MIB, SIGNIFICANT - MIB] {
let d = t
.sample_at(rss_at(outside), 400 * MIB, BUDGET, start)
.expect("the baseline is captured");
assert_eq!(
d.verdict,
DriftVerdict::OutsideHeap,
"a term at {outside} bytes flipped inside the hysteresis band"
);
}
let released = BUDGET * RELEASE_PCT / 100 - MIB;
let d = t
.sample_at(rss_at(released), 400 * MIB, BUDGET, start)
.expect("the baseline is captured");
assert_eq!(d.verdict, DriftVerdict::Settled);
}
#[test]
fn shrinking_is_reported_but_never_read_as_drift() {
let start = Instant::now();
let mut t = settled(2000 * MIB, 400 * MIB, start);
let d = t
.sample_at(1500 * MIB, 300 * MIB, BUDGET, start)
.expect("the baseline is captured");
assert_eq!(d.verdict, DriftVerdict::Settled);
assert_eq!(d.heap_growth_bytes, -((100 * MIB) as i64));
assert_eq!(d.outside_heap_growth_bytes, -((400 * MIB) as i64));
}
#[test]
fn a_climbing_startup_does_not_become_the_baseline() {
let start = Instant::now();
let mut t = DriftTracker::default();
let mut rss = 100 * MIB;
for _ in 0..20 {
assert_eq!(t.sample_at(rss, 50 * MIB, BUDGET, start), None);
rss += rss / 10;
}
assert!(t.baseline.is_none(), "a climbing session has not settled");
for _ in 0..=SETTLE_STREAK {
assert_eq!(t.sample_at(rss, 50 * MIB, BUDGET, start), None);
}
assert!(t.baseline.is_some(), "a flat run settles the baseline");
assert!(t.sample_at(rss, 50 * MIB, BUDGET, start).is_some());
}
#[test]
fn one_flat_sample_inside_a_climb_does_not_settle_it() {
let start = Instant::now();
let mut t = DriftTracker::default();
let mut rss = 100 * MIB;
for _ in 0..10 {
t.sample_at(rss, 50 * MIB, BUDGET, start);
t.sample_at(rss, 50 * MIB, BUDGET, start);
rss += rss / 4;
t.sample_at(rss, 50 * MIB, BUDGET, start);
}
assert!(
t.baseline.is_none(),
"a climb interrupted by single flat samples is not settled"
);
}
#[test]
fn a_session_that_never_settles_captures_a_baseline_at_the_deadline() {
let start = Instant::now();
let mut t = DriftTracker::default();
let mut rss = 100 * MIB;
for _ in 0..SETTLE_DEADLINE_SAMPLES {
t.sample_at(rss, 50 * MIB, BUDGET, start);
rss += rss / 10;
}
assert!(
t.baseline.is_some(),
"the deadline captures a baseline even while RSS climbs"
);
}
#[test]
fn a_zero_budget_reports_movement_without_a_verdict() {
let start = Instant::now();
let mut t = DriftTracker::default();
for _ in 0..=SETTLE_STREAK {
t.sample_at(2000 * MIB, 400 * MIB, 0, start);
}
let d = t
.sample_at(9000 * MIB, 400 * MIB, 0, start)
.expect("the baseline is captured");
assert_eq!(d.verdict, DriftVerdict::Settled);
assert_eq!(d.outside_heap_growth_bytes, (7000 * MIB) as i64);
}
#[test]
fn every_verdict_has_a_label() {
for v in [
DriftVerdict::Settled,
DriftVerdict::Heap,
DriftVerdict::OutsideHeap,
DriftVerdict::Both,
] {
assert!(!v.label().is_empty());
}
}
}