use std::sync::atomic::{AtomicBool, AtomicI64, Ordering};
use std::sync::Mutex;
#[derive(Debug, Default)]
pub struct RuntimeHealth {
last_poll_tick: AtomicI64,
watermark_paused: AtomicBool,
schedulers_enabled: AtomicBool,
started_at: AtomicI64,
db_probe: Mutex<Option<DbProbe>>,
db_probe_running: AtomicBool,
}
impl RuntimeHealth {
pub fn new() -> Self {
Self::default()
}
pub fn set_started_at(&self, now_unix: i64) {
self.started_at.store(now_unix, Ordering::Relaxed);
}
pub fn uptime_secs(&self, now_unix: i64) -> Option<i64> {
match self.started_at.load(Ordering::Relaxed) {
0 => None,
t => Some((now_unix - t).max(0)),
}
}
pub fn begin_db_probe(self: &std::sync::Arc<Self>) -> Result<DbProbeGuard, DbProbe> {
if self.db_probe_running.swap(true, Ordering::AcqRel) {
return Err(self.last_db_probe().unwrap_or(DbProbe::Unknown));
}
Ok(DbProbeGuard {
health: std::sync::Arc::clone(self),
})
}
#[cfg(test)]
pub fn record_for_test(&self, verdict: DbProbe) {
*self
.db_probe
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(verdict);
}
fn last_db_probe(&self) -> Option<DbProbe> {
self.db_probe
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub fn set_schedulers_enabled(&self, enabled: bool) {
self.schedulers_enabled.store(enabled, Ordering::Relaxed);
}
pub fn schedulers_enabled(&self) -> bool {
self.schedulers_enabled.load(Ordering::Relaxed)
}
pub fn poll_tick_completed(&self, now_unix: i64) {
self.last_poll_tick.store(now_unix, Ordering::Relaxed);
}
pub fn secs_since_poll_tick(&self, now_unix: i64) -> Option<i64> {
match self.last_poll_tick.load(Ordering::Relaxed) {
0 => None,
t => Some((now_unix - t).max(0)),
}
}
pub fn set_watermark(&self, paused: bool) {
self.watermark_paused.store(paused, Ordering::Relaxed);
}
pub fn watermark_paused(&self) -> bool {
self.watermark_paused.load(Ordering::Relaxed)
}
}
pub struct DbProbeGuard {
health: std::sync::Arc<RuntimeHealth>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum DbProbe {
Ok,
Failed(String),
Unknown,
}
impl DbProbeGuard {
pub fn record(self, verdict: DbProbe) {
*self
.health
.db_probe
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(verdict);
}
}
impl Drop for DbProbeGuard {
fn drop(&mut self) {
self.health.db_probe_running.store(false, Ordering::Release);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_fresh_record_reports_nothing_rather_than_zero() {
let h = RuntimeHealth::new();
assert_eq!(h.secs_since_poll_tick(1_000), None);
assert!(!h.watermark_paused());
assert!(!h.schedulers_enabled());
}
#[test]
fn a_tick_ages_and_a_backwards_clock_does_not_go_negative() {
let h = RuntimeHealth::new();
h.poll_tick_completed(1_000);
assert_eq!(h.secs_since_poll_tick(1_090), Some(90));
assert_eq!(
h.secs_since_poll_tick(900),
Some(0),
"a clock step backwards must read as 'just now', not as a negative age"
);
}
#[test]
fn the_watermark_verdict_round_trips_both_ways() {
let h = RuntimeHealth::new();
h.set_watermark(true);
assert!(h.watermark_paused());
h.set_watermark(false);
assert!(!h.watermark_paused());
}
}