use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use arc_swap::ArcSwap;
use nmbrs_metrics::scheduler::Reporter;
use nmbrs_metrics::snapshot::MetricSet;
use crate::phase_outcome::Outcome;
pub trait PulseEvaluator: Send {
fn evaluate(&mut self, window: &MetricSet) -> Option<Outcome>;
}
pub type StopOutcomeCell = Arc<ArcSwap<Option<Outcome>>>;
pub struct PhaseStopEvaluator {
eval: Box<dyn PulseEvaluator>,
stop_flag: Arc<AtomicBool>,
outcome: StopOutcomeCell,
done: bool,
}
impl PhaseStopEvaluator {
pub fn new(eval: Box<dyn PulseEvaluator>, stop_flag: Arc<AtomicBool>) -> Self {
Self {
eval,
stop_flag,
outcome: Arc::new(ArcSwap::from_pointee(None)),
done: false,
}
}
pub fn outcome_cell(&self) -> StopOutcomeCell {
self.outcome.clone()
}
}
impl Reporter for PhaseStopEvaluator {
fn report(&mut self, window: &MetricSet) {
if self.done {
return;
}
if let Some(outcome) = self.eval.evaluate(window) {
self.outcome.store(Arc::new(Some(outcome)));
self.stop_flag.store(true, Ordering::Relaxed);
self.done = true;
}
}
fn finished(&self) -> bool {
self.done
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::phase_outcome::{Disposition, Validity};
use std::time::Duration;
struct FireOnNth {
n: u64,
seen: u64,
outcome: Outcome,
}
impl PulseEvaluator for FireOnNth {
fn evaluate(&mut self, _w: &MetricSet) -> Option<Outcome> {
self.seen += 1;
(self.seen >= self.n).then(|| self.outcome.clone())
}
}
fn empty_window() -> MetricSet {
MetricSet::new(Duration::from_secs(1))
}
#[test]
fn fires_once_then_reports_finished_and_no_ops() {
let stop = Arc::new(AtomicBool::new(false));
let mut ev = PhaseStopEvaluator::new(
Box::new(FireOnNth {
n: 3,
seen: 0,
outcome: Outcome::interrupted(),
}),
stop.clone(),
);
let cell = ev.outcome_cell();
ev.report(&empty_window());
ev.report(&empty_window());
assert!(!stop.load(Ordering::Relaxed), "no stop before the verdict");
assert!(!ev.finished(), "not finished before the verdict");
assert!(cell.load().is_none(), "no outcome before the verdict");
ev.report(&empty_window()); assert!(
stop.load(Ordering::Relaxed),
"stop flag raised on the verdict"
);
assert!(ev.finished(), "self-unregisters after the verdict");
let got = (**cell.load()).clone().expect("outcome published");
assert_eq!(got.disposition, Disposition::Interrupted);
assert_eq!(got.validity, Validity::Succeeded);
ev.report(&empty_window());
assert_eq!(
(**cell.load()).clone().expect("still set").disposition,
Disposition::Interrupted
);
}
#[test]
fn timeout_disposition_differs_from_settle() {
let stop = Arc::new(AtomicBool::new(false));
let mut ev = PhaseStopEvaluator::new(
Box::new(FireOnNth {
n: 1,
seen: 0,
outcome: Outcome::failed(),
}),
stop.clone(),
);
let cell = ev.outcome_cell();
ev.report(&empty_window());
let got = (**cell.load()).clone().expect("outcome published");
assert_eq!(got.disposition, Disposition::Interrupted);
assert_eq!(
got.validity,
Validity::Failed,
"timeout is the untrustworthy quadrant"
);
}
}