use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread::{self, JoinHandle};
use std::time::Duration;
use super::core::Phase;
pub struct Heartbeat {
handle: Option<JoinHandle<()>>,
stop: Arc<AtomicBool>,
floor: Arc<std::sync::atomic::AtomicU64>,
}
impl Heartbeat {
pub fn spawn(phase: Phase, interval: Duration) -> Self {
Self::spawn_with_backoff(phase, interval, interval)
}
pub fn spawn_with_backoff(phase: Phase, initial: Duration, max: Duration) -> Self {
let initial = initial.max(Duration::from_millis(250));
let max = max.max(initial);
let stop = Arc::new(AtomicBool::new(false));
let floor = Arc::new(std::sync::atomic::AtomicU64::new(0));
let stop_clone = stop.clone();
let floor_clone = floor.clone();
let phase_clone = phase.clone();
let handle = thread::spawn(move || {
let mut count: u64 = 0;
let mut interval = initial;
while !stop_clone.load(Ordering::Relaxed) {
thread::sleep(interval);
if stop_clone.load(Ordering::Relaxed) {
break;
}
count = count.saturating_add(1);
let f = floor_clone.load(Ordering::Relaxed);
let value = count.max(f);
phase_clone.tick(value);
interval = interval.checked_mul(2).unwrap_or(max).min(max);
}
});
Self {
handle: Some(handle),
stop,
floor,
}
}
pub fn raise_floor(&self, value: u64) {
let prev = self.floor.load(Ordering::Relaxed);
if value > prev {
self.floor.store(value, Ordering::Relaxed);
}
}
pub fn stop(mut self) {
self.stop.store(true, Ordering::Relaxed);
if let Some(handle) = self.handle.take() {
let _ = handle.join();
}
}
}
impl Drop for Heartbeat {
fn drop(&mut self) {
self.stop.store(true, Ordering::Relaxed);
if let Some(handle) = self.handle.take() {
let _ = handle.join();
}
}
}