use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WtActivity {
Active,
Idle,
Deactivated,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ActivityConfig {
pub idle_after: Duration,
pub deactivate_after: Duration,
}
impl ActivityConfig {
pub const DEFAULT_IDLE_SECS: u64 = 300;
pub const DEFAULT_DEACTIVATE_SECS: u64 = 600;
pub const FLOOR_SECS: u64 = 30;
pub fn defaults() -> Self {
Self {
idle_after: Duration::from_secs(Self::DEFAULT_IDLE_SECS),
deactivate_after: Duration::from_secs(Self::DEFAULT_DEACTIVATE_SECS),
}
}
pub fn from_env(env: &dyn Fn(&str) -> Option<String>) -> Self {
Self {
idle_after: Duration::from_secs(resolve_secs(
env("TF_WT_IDLE_SECS").as_deref(),
Self::DEFAULT_IDLE_SECS,
)),
deactivate_after: Duration::from_secs(resolve_secs(
env("TF_WT_DEACTIVATE_SECS").as_deref(),
Self::DEFAULT_DEACTIVATE_SECS,
)),
}
}
}
impl Default for ActivityConfig {
fn default() -> Self {
Self::defaults()
}
}
pub fn resolve_secs(v: Option<&str>, default: u64) -> u64 {
match v {
Some(s) => s
.parse::<u64>()
.map(|n| n.max(ActivityConfig::FLOOR_SECS))
.unwrap_or(default),
None => default,
}
}
#[derive(Debug, Clone)]
pub struct WtLifecycle {
cfg: ActivityConfig,
state: WtActivity,
last_activity: Instant,
idle_since: Option<Instant>,
}
impl WtLifecycle {
pub fn new(now: Instant, cfg: ActivityConfig) -> Self {
Self {
cfg,
state: WtActivity::Active,
last_activity: now,
idle_since: None,
}
}
pub fn state(&self) -> WtActivity {
self.state
}
pub fn touch(&mut self, now: Instant) {
self.state = WtActivity::Active;
self.last_activity = now;
self.idle_since = None;
}
pub fn tick(&mut self, now: Instant) -> WtActivity {
let idle_for = now.saturating_duration_since(self.last_activity);
if idle_for < self.cfg.idle_after {
if self.state == WtActivity::Active {
self.idle_since = None;
}
return self.state;
}
if self.state == WtActivity::Active {
self.state = WtActivity::Idle;
self.idle_since = Some(self.last_activity + self.cfg.idle_after);
}
if self.state == WtActivity::Idle {
let since = self
.idle_since
.unwrap_or(self.last_activity + self.cfg.idle_after);
if now.saturating_duration_since(since) >= self.cfg.deactivate_after {
self.state = WtActivity::Deactivated;
}
}
self.state
}
}
#[derive(Debug, Default)]
pub struct ActivityCounters {
deactivations: AtomicU64,
reactivations: AtomicU64,
deactivated_ms: AtomicU64,
}
impl ActivityCounters {
pub fn new() -> Self {
Self::default()
}
pub fn record_deactivation(&self) {
self.deactivations.fetch_add(1, Ordering::Relaxed);
}
pub fn record_reactivation(&self) {
self.reactivations.fetch_add(1, Ordering::Relaxed);
}
pub fn add_deactivated(&self, d: Duration) {
self.deactivated_ms
.fetch_add(d.as_millis() as u64, Ordering::Relaxed);
}
pub fn snapshot(&self) -> (u64, u64, u64) {
(
self.deactivations.load(Ordering::Relaxed),
self.reactivations.load(Ordering::Relaxed),
self.deactivated_ms.load(Ordering::Relaxed),
)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn cfg(idle: u64, deact: u64) -> ActivityConfig {
ActivityConfig {
idle_after: Duration::from_secs(idle),
deactivate_after: Duration::from_secs(deact),
}
}
#[test]
fn starts_active_and_stays_active_under_threshold() {
let t0 = Instant::now();
let mut wt = WtLifecycle::new(t0, cfg(300, 600));
assert_eq!(wt.state(), WtActivity::Active);
assert_eq!(wt.tick(t0 + Duration::from_secs(299)), WtActivity::Active);
}
#[test]
fn active_to_idle_to_deactivated_progression() {
let t0 = Instant::now();
let mut wt = WtLifecycle::new(t0, cfg(300, 600));
assert_eq!(wt.tick(t0 + Duration::from_secs(300)), WtActivity::Idle);
assert_eq!(wt.tick(t0 + Duration::from_secs(899)), WtActivity::Idle);
assert_eq!(
wt.tick(t0 + Duration::from_secs(900)),
WtActivity::Deactivated
);
}
#[test]
fn coarse_tick_cannot_postpone_deactivation() {
let t0 = Instant::now();
let mut wt = WtLifecycle::new(t0, cfg(300, 600));
assert_eq!(
wt.tick(t0 + Duration::from_secs(5000)),
WtActivity::Deactivated
);
}
#[test]
fn touch_reactivates_from_any_state() {
let t0 = Instant::now();
let mut wt = WtLifecycle::new(t0, cfg(300, 600));
wt.tick(t0 + Duration::from_secs(2000)); assert_eq!(wt.state(), WtActivity::Deactivated);
wt.touch(t0 + Duration::from_secs(2001));
assert_eq!(wt.state(), WtActivity::Active);
assert_eq!(
wt.tick(t0 + Duration::from_secs(2001 + 299)),
WtActivity::Active
);
assert_eq!(
wt.tick(t0 + Duration::from_secs(2001 + 300)),
WtActivity::Idle,
"fresh idle window measured from the re-activation, not t0"
);
}
#[test]
fn non_monotonic_now_never_regresses_state() {
let t0 = Instant::now();
let mut wt = WtLifecycle::new(t0 + Duration::from_secs(1000), cfg(300, 600));
assert_eq!(
wt.tick(t0 + Duration::from_secs(500)),
WtActivity::Active,
"earlier now ⇒ saturating 0 elapsed ⇒ stays Active, no panic/regress"
);
}
#[test]
fn resolve_secs_rule_default_clamp_floor() {
assert_eq!(resolve_secs(None, 300), 300, "unset ⇒ default");
assert_eq!(
resolve_secs(Some("nope"), 300),
300,
"non-numeric ⇒ default"
);
assert_eq!(resolve_secs(Some("900"), 300), 900, "sane value honored");
assert_eq!(
resolve_secs(Some("5"), 300),
ActivityConfig::FLOOR_SECS,
"below floor ⇒ clamped"
);
assert_eq!(resolve_secs(Some("30"), 300), 30, "floor exact honored");
}
#[test]
fn from_env_threads_both_knobs() {
let env = |k: &str| match k {
"TF_WT_IDLE_SECS" => Some("120".to_string()),
"TF_WT_DEACTIVATE_SECS" => Some("7".to_string()), _ => None,
};
let c = ActivityConfig::from_env(&env);
assert_eq!(c.idle_after, Duration::from_secs(120));
assert_eq!(
c.deactivate_after,
Duration::from_secs(ActivityConfig::FLOOR_SECS),
"below-floor deactivate clamped, not honored"
);
let c2 = ActivityConfig::from_env(&|_| None);
assert_eq!(c2, ActivityConfig::defaults());
}
#[test]
fn counters_track_edges_and_time() {
let c = ActivityCounters::new();
assert_eq!(c.snapshot(), (0, 0, 0));
c.record_deactivation();
c.add_deactivated(Duration::from_millis(2000));
c.record_reactivation();
c.record_deactivation();
c.add_deactivated(Duration::from_millis(500));
assert_eq!(c.snapshot(), (2, 1, 2500));
}
}