Skip to main content

nmbrs_runtime/optimize/
phase_pulse.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! SRD-86 §"Settling" — a cadence-pulse phase evaluator.
5//!
6//! A [`PhaseStopEvaluator`] is a callback registered on the **metrics
7//! cadence feed** ([`nmbrs_metrics::cadence_reporter::CadenceReporter::subscribe`],
8//! SRD-42). It runs once per **cadence pulse** against the running
9//! phase's current state + metrics, and at any pulse may set a
10//! **terminal disposition** on the phase — stopping it (cooperatively,
11//! by raising the phase stop flag, which the activity loop reads at its
12//! next cycle boundary) and recording *which* [`Outcome`] disposition it
13//! stopped with. Once it has set a terminal disposition it
14//! **unregisters itself**: it reports [`finished`](Reporter::finished),
15//! so the cadence-feed dispatch worker stops delivering further pulses
16//! (no self-join deadlock — the evaluator runs on that worker thread).
17//!
18//! The per-pulse decision is a [`PulseEvaluator`]: it returns
19//! `Some(outcome)` to stop the phase with that disposition, or `None` to
20//! let it keep running. The settle detector
21//! ([`super::settle::SettleEvaluator`]) is the first implementation;
22//! the mechanism is general (any cadence-driven phase-stop policy).
23
24use std::sync::Arc;
25use std::sync::atomic::{AtomicBool, Ordering};
26
27use arc_swap::ArcSwap;
28use nmbrs_metrics::scheduler::Reporter;
29use nmbrs_metrics::snapshot::MetricSet;
30
31use crate::phase_outcome::Outcome;
32
33/// Per-pulse evaluation of a running phase against the just-published
34/// cadence window. Returns `Some(outcome)` to terminally stop the phase
35/// with that disposition, or `None` to let it keep running.
36///
37/// Implementors read the phase's live state however they need — the
38/// `window` is the freshly-closed [`MetricSet`] that triggered this
39/// pulse (a metrics reader may instead re-read the live feed, which
40/// resolves to the same window the cadence feed just published).
41pub trait PulseEvaluator: Send {
42    fn evaluate(&mut self, window: &MetricSet) -> Option<Outcome>;
43}
44
45/// The terminal disposition a [`PhaseStopEvaluator`] publishes. The
46/// executor reads it at phase completion: `None` ⇒ the evaluator never
47/// fired (the phase ran to its natural end); `Some(outcome)` ⇒ the
48/// evaluator stopped the phase with that disposition.
49pub type StopOutcomeCell = Arc<ArcSwap<Option<Outcome>>>;
50
51/// A cadence-feed callback that drives a [`PulseEvaluator`] against a
52/// running phase and terminally stops it on the first verdict. See the
53/// module docs.
54pub struct PhaseStopEvaluator {
55    eval: Box<dyn PulseEvaluator>,
56    stop_flag: Arc<AtomicBool>,
57    outcome: StopOutcomeCell,
58    done: bool,
59}
60
61impl PhaseStopEvaluator {
62    /// Build an evaluator over the phase's `stop_flag` (raised
63    /// cooperatively when the evaluator yields a terminal disposition).
64    pub fn new(eval: Box<dyn PulseEvaluator>, stop_flag: Arc<AtomicBool>) -> Self {
65        Self {
66            eval,
67            stop_flag,
68            outcome: Arc::new(ArcSwap::from_pointee(None)),
69            done: false,
70        }
71    }
72
73    /// The cell the executor reads at phase completion for the terminal
74    /// disposition. Shared (cloned `Arc`) so the executor holds a handle
75    /// while the evaluator runs on the cadence-feed worker.
76    pub fn outcome_cell(&self) -> StopOutcomeCell {
77        self.outcome.clone()
78    }
79}
80
81impl Reporter for PhaseStopEvaluator {
82    fn report(&mut self, window: &MetricSet) {
83        if self.done {
84            return;
85        }
86        // The delivery fiber already runs inside this subscription's
87        // execution context (SRD-88 — set as the cadence subscription's
88        // `context_wrap` at subscribe time), so the objective read here
89        // scopes to the owning execution without any per-call rebinding.
90        if let Some(outcome) = self.eval.evaluate(window) {
91            self.outcome.store(Arc::new(Some(outcome)));
92            // Cooperative stop: the activity loop reads the flag at its
93            // next cycle boundary.
94            self.stop_flag.store(true, Ordering::Relaxed);
95            self.done = true;
96        }
97    }
98
99    fn finished(&self) -> bool {
100        self.done
101    }
102}
103
104#[cfg(test)]
105mod tests {
106    use super::*;
107    use crate::phase_outcome::{Disposition, Validity};
108    use std::time::Duration;
109
110    /// An evaluator that fires a fixed outcome on the Nth pulse.
111    struct FireOnNth {
112        n: u64,
113        seen: u64,
114        outcome: Outcome,
115    }
116    impl PulseEvaluator for FireOnNth {
117        fn evaluate(&mut self, _w: &MetricSet) -> Option<Outcome> {
118            self.seen += 1;
119            (self.seen >= self.n).then(|| self.outcome.clone())
120        }
121    }
122
123    fn empty_window() -> MetricSet {
124        MetricSet::new(Duration::from_secs(1))
125    }
126
127    #[test]
128    fn fires_once_then_reports_finished_and_no_ops() {
129        let stop = Arc::new(AtomicBool::new(false));
130        let mut ev = PhaseStopEvaluator::new(
131            Box::new(FireOnNth {
132                n: 3,
133                seen: 0,
134                outcome: Outcome::interrupted(),
135            }),
136            stop.clone(),
137        );
138        let cell = ev.outcome_cell();
139
140        ev.report(&empty_window());
141        ev.report(&empty_window());
142        assert!(!stop.load(Ordering::Relaxed), "no stop before the verdict");
143        assert!(!ev.finished(), "not finished before the verdict");
144        assert!(cell.load().is_none(), "no outcome before the verdict");
145
146        ev.report(&empty_window()); // third pulse fires
147        assert!(
148            stop.load(Ordering::Relaxed),
149            "stop flag raised on the verdict"
150        );
151        assert!(ev.finished(), "self-unregisters after the verdict");
152        let got = (**cell.load()).clone().expect("outcome published");
153        assert_eq!(got.disposition, Disposition::Interrupted);
154        assert_eq!(got.validity, Validity::Succeeded);
155
156        // Further pulses are no-ops — the disposition is not overwritten.
157        ev.report(&empty_window());
158        assert_eq!(
159            (**cell.load()).clone().expect("still set").disposition,
160            Disposition::Interrupted
161        );
162    }
163
164    #[test]
165    fn timeout_disposition_differs_from_settle() {
166        // A failed() verdict (timeout) publishes Interrupted+Failed.
167        let stop = Arc::new(AtomicBool::new(false));
168        let mut ev = PhaseStopEvaluator::new(
169            Box::new(FireOnNth {
170                n: 1,
171                seen: 0,
172                outcome: Outcome::failed(),
173            }),
174            stop.clone(),
175        );
176        let cell = ev.outcome_cell();
177        ev.report(&empty_window());
178        let got = (**cell.load()).clone().expect("outcome published");
179        assert_eq!(got.disposition, Disposition::Interrupted);
180        assert_eq!(
181            got.validity,
182            Validity::Failed,
183            "timeout is the untrustworthy quadrant"
184        );
185    }
186}