use std::sync::Arc;
use std::sync::atomic::AtomicBool;
use std::time::{Duration, Instant};
use crate::scope_kernel::ScopeKernel;
use arc_swap::ArcSwap;
use nmbrs_metrics::cadence_reporter::{CadenceReporter, SubscriberId};
use nmbrs_metrics::snapshot::MetricSet;
use polydat::Kernel;
use polydat::ast::Value;
use polydat::kernel::PolydatProgram;
use super::phase_pulse::{PhaseStopEvaluator, PulseEvaluator, StopOutcomeCell};
use crate::phase_outcome::Outcome;
fn objective_to_f64(v: &Value) -> f64 {
match v {
Value::F64(f) => *f,
Value::U64(u) => *u as f64,
Value::Bool(b) => {
if *b {
1.0
} else {
0.0
}
}
_ => 0.0,
}
}
#[derive(Clone, Copy, Debug)]
pub struct SettleReading {
pub value: f64,
pub stable: bool,
pub pulses: u64,
}
impl Default for SettleReading {
fn default() -> Self {
Self {
value: 0.0,
stable: false,
pulses: 0,
}
}
}
pub struct SettleInterpreter {
kernel: Box<dyn Kernel>,
samples_input: usize,
window: std::collections::VecDeque<f64>,
horizon: usize,
value_wire: String,
stable_wire: String,
register: Arc<ArcSwap<SettleReading>>,
pulses: u64,
}
impl SettleInterpreter {
pub fn new(
kernel: Box<dyn Kernel>,
samples: &str,
value_wire: &str,
stable_wire: &str,
horizon: usize,
) -> Self {
let samples_input = kernel
.input_index(samples)
.unwrap_or_else(|| panic!("settle kernel has no `{samples}` input"));
Self {
kernel,
samples_input,
window: std::collections::VecDeque::with_capacity(horizon),
horizon,
value_wire: value_wire.to_string(),
stable_wire: stable_wire.to_string(),
register: Arc::new(ArcSwap::from_pointee(SettleReading::default())),
pulses: 0,
}
}
pub fn register(&self) -> Arc<ArcSwap<SettleReading>> {
self.register.clone()
}
pub fn pulse(&mut self, sample: f64) -> SettleReading {
self.pulses += 1;
if self.window.len() == self.horizon {
self.window.pop_front();
}
self.window.push_back(sample);
let window: Vec<f64> = self.window.iter().copied().collect();
self.kernel
.set_input_at(
self.samples_input,
Value::VecF64(polydat::ast::SliceArc::from_vec(window)),
)
.expect("settle kernel refused its vec_f64 `samples` extern");
let stable = self.kernel.pull(&self.stable_wire).as_u64() != 0;
let value = self.kernel.pull(&self.value_wire).as_f64();
let reading = SettleReading {
value,
stable,
pulses: self.pulses,
};
self.register.store(Arc::new(reading));
reading
}
}
pub struct SettleEvaluator {
objective: ScopeKernel,
objective_wire: String,
poke: Option<usize>,
interp: SettleInterpreter,
timeout: Duration,
min_viable: Duration,
started: Option<Instant>,
pulses: u64,
}
impl SettleEvaluator {
pub fn new(
objective: ScopeKernel,
objective_wire: &str,
poke_input: &str,
interp: SettleInterpreter,
timeout: Duration,
min_viable: Duration,
) -> Self {
let poke = objective.program().find_input(poke_input);
Self {
objective,
objective_wire: objective_wire.to_string(),
poke,
interp,
timeout,
min_viable,
started: None,
pulses: 0,
}
}
pub fn register(&self) -> Arc<ArcSwap<SettleReading>> {
self.interp.register()
}
}
impl PulseEvaluator for SettleEvaluator {
fn evaluate(&mut self, _window: &MetricSet) -> Option<Outcome> {
let start = *self.started.get_or_insert_with(Instant::now);
self.pulses += 1;
if let Some(idx) = self.poke {
let k = &mut self.objective;
let coords = k.coord_count();
if idx < coords {
let mut position: Vec<u64> = (0..coords)
.map(|i| k.input_value_at(i).map_or(0, |v| v.as_u64()))
.collect();
position[idx] = self.pulses;
k.set_inputs(&position);
k.invalidate_all();
} else if let Err(e) = k.set_input_at(idx, Value::U64(self.pulses)) {
crate::diag!(
crate::observer::LogLevel::Warn,
"settle: the objective's poke input refused pulse {}: {e}",
self.pulses
);
}
}
let obj = objective_to_f64(&self.objective.pull(&self.objective_wire));
if obj.is_nan() {
if start.elapsed() >= self.timeout {
return Some(Outcome::failed());
}
return None;
}
let reading = self.interp.pulse(obj);
if reading.stable && start.elapsed() >= self.min_viable {
return Some(Outcome::interrupted());
}
if start.elapsed() >= self.timeout {
return Some(Outcome::failed());
}
None
}
}
const READER_NODES: &[&str] = &[
"metric",
"metric_window",
"metricsql",
"metricsql_scalar",
"metricsql_vector",
"metricsql_window",
];
pub fn program_reads_live_metrics(program: &PolydatProgram) -> bool {
(0..program.node_count()).any(|i| READER_NODES.contains(&program.node_meta(i).name.as_str()))
}
const SESSION_CUMULATIVE_READER: &str = "metric";
fn program_reads_session_cumulative_metrics(program: &PolydatProgram) -> bool {
(0..program.node_count()).any(|i| program.node_meta(i).name == SESSION_CUMULATIVE_READER)
}
fn warn_once_session_cumulative(objective: &str) -> bool {
static WARNED: std::sync::LazyLock<std::sync::Mutex<std::collections::HashSet<String>>> =
std::sync::LazyLock::new(Default::default);
WARNED
.lock()
.unwrap_or_else(|e| e.into_inner())
.insert(objective.to_string())
}
const SETTLE_MARGIN: f64 = 0.05;
const SETTLE_MIN_SAMPLES: u64 = 4;
const SETTLE_HORIZON: u64 = 8;
const SETTLE_TIMEOUT: Duration = Duration::from_secs(60);
pub struct SettleHandle {
pub subscriber: SubscriberId,
pub register: Arc<ArcSwap<SettleReading>>,
pub outcome: StopOutcomeCell,
}
#[derive(Debug)]
pub enum SettleSkip {
NotWindowed,
CadenceDisabled,
Failed(String),
}
impl std::fmt::Display for SettleSkip {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
SettleSkip::NotWindowed => write!(f, "it reads no live windowed metric"),
SettleSkip::CadenceDisabled => write!(f, "the metrics cadence is disabled"),
SettleSkip::Failed(reason) => {
write!(f, "the settle detector failed to start: {reason}")
}
}
}
}
pub fn start_settle(
parent: &Arc<ScopeKernel>,
phase_kernel: &Arc<ScopeKernel>,
objective: &str,
reporter: &Arc<CadenceReporter>,
stop_flag: Arc<AtomicBool>,
) -> Result<SettleHandle, SettleSkip> {
let program = phase_kernel.program();
if !program_reads_live_metrics(program) {
return Err(SettleSkip::NotWindowed);
}
if program_reads_session_cumulative_metrics(program) && warn_once_session_cumulative(objective)
{
crate::diag!(
crate::observer::LogLevel::Warn,
"optimizer objective '{objective}' reads a session-cumulative metric \
(`metric(...)` → session_lifetime): it aggregates across coordinates and \
will not isolate per-coordinate. Use `metric_window(...)` or \
`metricsql_scalar(rate(...[W]))` for a per-coordinate objective."
);
}
let cadence = reporter.declared_cadences().smallest();
if cadence.is_zero() {
return Err(SettleSkip::CadenceDisabled);
}
let failed = |what: &str, e: &dyn std::fmt::Display| SettleSkip::Failed(format!("{what}: {e}"));
let obj_kernel = phase_kernel
.bind_under(parent.kernel(), &[])
.map_err(|e| failed("objective kernel", &e))?;
let is_stable_kernel = polydat::dsl::compile::compile_polydat(&format!(
"extern samples: vec_f64\n(stable_value, stable) := is_stable(samples, {SETTLE_MARGIN}, \
{SETTLE_MIN_SAMPLES})"
))
.map_err(|e| failed("is_stable kernel", &e))?;
let interp = SettleInterpreter::new(
is_stable_kernel,
"samples",
"stable_value",
"stable",
SETTLE_HORIZON as usize,
);
let min_viable = cadence.saturating_mul(SETTLE_HORIZON as u32);
let eval = SettleEvaluator::new(
obj_kernel,
objective,
"cycle",
interp,
SETTLE_TIMEOUT,
min_viable,
);
let register = eval.register();
let pse = PhaseStopEvaluator::new(Box::new(eval), stop_flag);
let outcome = pse.outcome_cell();
let mut opts = nmbrs_metrics::cadence_reporter::SubscriptionOpts::default();
if let Some(ctx) = crate::execution_context::try_current() {
opts.context_wrap = Some(std::sync::Arc::new(move |fut| {
Box::pin(crate::execution_context::scope(ctx.clone(), fut))
as std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>>
}));
}
let subscriber = reporter
.subscribe(cadence, Box::new(pse), opts)
.map_err(|e| failed("cadence subscription", &e))?;
Ok(SettleHandle {
subscriber,
register,
outcome,
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::phase_outcome::{Disposition, Validity};
use polydat::dsl::compile::{compile_polydat, compile_polydat_interpreter};
fn settle_interp() -> SettleInterpreter {
let kernel = compile_polydat(
"extern samples: vec_f64\n(stable_value, stable) := is_stable(samples, 0.05, 4)",
)
.expect("is_stable kernel compiles");
SettleInterpreter::new(kernel, "samples", "stable_value", "stable", 8)
}
fn obj_kernel(src: &str) -> ScopeKernel {
crate::bindings::compile_scope_kernel(src, &Default::default())
.expect("objective kernel compiles")
}
const STEADY_OBJ: &str = "input cycle: u64\nobj := 5.0";
const RAMP_OBJ: &str = "input cycle: u64\nobj := cycle";
fn window() -> MetricSet {
MetricSet::new(Duration::from_secs(1))
}
#[test]
fn interpreter_publishes_settled_value_into_the_register() {
let mut i = settle_interp();
let mut last = SettleReading::default();
for _ in 0..8 {
last = i.pulse(5.0);
}
assert!(last.stable, "a steady value settles");
assert!(
(i.register().load().value - 5.0).abs() < 1e-9,
"register holds the steady level"
);
}
#[test]
fn evaluator_yields_interrupted_when_settled() {
let mut ev = SettleEvaluator::new(
obj_kernel(STEADY_OBJ),
"obj",
"cycle",
settle_interp(),
Duration::from_secs(60),
Duration::ZERO,
);
let reg = ev.register();
let mut verdict = None;
for _ in 0..16 {
if let Some(o) = ev.evaluate(&window()) {
verdict = Some(o);
break;
}
}
let o = verdict.expect("a steady objective settles within the budget");
assert_eq!(o.disposition, Disposition::Interrupted);
assert_eq!(o.validity, Validity::Succeeded);
assert!(
(reg.load().value - 5.0).abs() < 1e-9,
"settled register reads 5.0"
);
}
#[test]
fn evaluator_yields_failed_on_settle_timeout() {
let mut ev = SettleEvaluator::new(
obj_kernel(RAMP_OBJ),
"obj",
"cycle",
settle_interp(),
Duration::from_millis(40),
Duration::ZERO,
);
assert!(
ev.evaluate(&window()).is_none(),
"no verdict before timeout"
);
std::thread::sleep(Duration::from_millis(55));
let o = ev.evaluate(&window()).expect("timeout fires a verdict");
assert_eq!(o.disposition, Disposition::Interrupted);
assert_eq!(
o.validity,
Validity::Failed,
"a settle timeout is the untrustworthy quadrant"
);
}
#[test]
fn viability_gate_withholds_settle_until_min_viable_elapses() {
let mut ev = SettleEvaluator::new(
obj_kernel(STEADY_OBJ),
"obj",
"cycle",
settle_interp(),
Duration::from_secs(60),
Duration::from_millis(60),
);
for _ in 0..32 {
assert!(
ev.evaluate(&window()).is_none(),
"stable-but-gated: a burst of pulses must not settle before min_viable wall-clock"
);
}
std::thread::sleep(Duration::from_millis(70));
let o = ev
.evaluate(&window())
.expect("settles once min_viable has elapsed");
assert_eq!(o.disposition, Disposition::Interrupted);
assert_eq!(o.validity, Validity::Succeeded);
}
#[test]
fn detects_session_cumulative_reader_only() {
let cum = compile_polydat_interpreter(r#"obj := metric("cycles_total, phase=p", "rate")"#)
.expect("metric node compiles");
assert!(program_reads_session_cumulative_metrics(cum.program()));
let win =
compile_polydat_interpreter(r#"obj := metric_window("cycles_total, phase=p", "rate")"#)
.expect("metric_window node compiles");
assert!(
!program_reads_session_cumulative_metrics(win.program()),
"metric_window is windowed, not session-cumulative"
);
let plain = compile_polydat_interpreter("obj := 5.0").expect("plain objective compiles");
assert!(!program_reads_session_cumulative_metrics(plain.program()));
}
}