use std::sync::{
Arc,
atomic::{AtomicU32, AtomicU64, Ordering},
};
use tracing::{debug, info, warn};
use crate::ActError;
pub(crate) const DEGRADE_AFTER: u64 = 3;
pub(crate) const MAX_BACKOFF_TICKS: u32 = 8;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum HealthState {
Healthy,
Degraded,
}
impl HealthState {
pub(crate) fn as_str(self) -> &'static str {
match self {
HealthState::Healthy => "healthy",
HealthState::Degraded => "degraded",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct LoopHealth {
pub state: HealthState,
pub consecutive_failures: u64,
pub backoff_ticks: u32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct SchedulerHealth {
pub trigger: LoopHealth,
pub retry: LoopHealth,
}
pub(crate) struct LoopGuard {
name: &'static str,
failures: AtomicU64,
skips: AtomicU32,
}
impl LoopGuard {
pub(crate) fn new(name: &'static str) -> Arc<Self> {
Arc::new(Self {
name,
failures: AtomicU64::new(0),
skips: AtomicU32::new(0),
})
}
pub(crate) fn attempt(&self) -> bool {
loop {
let skips = self.skips.load(Ordering::Acquire);
if skips == 0 {
return true;
}
if self
.skips
.compare_exchange_weak(skips, skips - 1, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
{
return false;
}
}
}
pub(crate) fn failed(&self, err: &ActError) {
let failures = self.failures.fetch_add(1, Ordering::AcqRel) + 1;
if failures < DEGRADE_AFTER {
return;
}
let backoff = Self::backoff_ticks(failures);
self.skips.store(backoff, Ordering::Release);
let widened = backoff > Self::backoff_ticks(failures - 1);
if failures == DEGRADE_AFTER || widened {
warn!(
loop_name = self.name,
state = HealthState::Degraded.as_str(),
consecutive_failures = failures,
backoff_ticks = backoff,
error = %err,
"store loop degraded"
);
}
}
pub(crate) fn recovered(&self) {
let failures = self.failures.swap(0, Ordering::AcqRel);
self.skips.store(0, Ordering::Release);
if failures >= DEGRADE_AFTER {
info!(
loop_name = self.name,
state = HealthState::Healthy.as_str(),
consecutive_failures = failures,
"store loop recovered"
);
} else if failures > 0 {
debug!(
loop_name = self.name,
state = HealthState::Healthy.as_str(),
consecutive_failures = failures,
"store loop recovered"
);
}
}
pub(crate) fn snapshot(&self) -> LoopHealth {
let consecutive_failures = self.failures.load(Ordering::Acquire);
LoopHealth {
state: if consecutive_failures >= DEGRADE_AFTER {
HealthState::Degraded
} else {
HealthState::Healthy
},
consecutive_failures,
backoff_ticks: self.skips.load(Ordering::Acquire),
}
}
fn backoff_ticks(failures: u64) -> u32 {
let steps = failures.saturating_sub(DEGRADE_AFTER).min(31) as u32;
(1u32 << steps).min(MAX_BACKOFF_TICKS)
}
}