use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, RwLock};
use std::time::Duration;
use crate::scope_kernel::ScopeKernel;
use arc_swap::ArcSwap;
use nmbrs_metrics::cadence_reporter::CadenceReporter;
use nmbrs_metrics::component::Component;
use nmbrs_metrics::controls::ControlOrigin;
use super::settle::start_settle;
use super::{Budget, Coord, LexSource, OptimizerParams, PullSource, SearchSpace};
#[derive(Debug, Clone)]
pub struct ControlAxis {
pub axis_idx: usize,
pub control: String,
}
#[derive(Clone)]
pub struct ServoBest {
pub coord: Coord,
pub value: f64,
}
#[derive(Clone, Default)]
pub struct ServoOutcome {
pub best: Option<ServoBest>,
pub evals: usize,
}
#[derive(Clone)]
pub struct ServoSpec {
pub method: String,
pub params: Vec<(String, f64)>,
pub objective: String,
pub max_evals: usize,
pub seed: u64,
pub space: SearchSpace,
pub controls: Vec<ControlAxis>,
pub result: Arc<ArcSwap<ServoOutcome>>,
}
const SETTLE_POLL: Duration = Duration::from_millis(50);
async fn retarget(
phase_component: &Arc<RwLock<Component>>,
control: &str,
value: f64,
) -> Result<(), String> {
let erased = {
let guard = phase_component.read().unwrap_or_else(|e| e.into_inner());
guard.find_control_erased_up(control)
};
let Some(erased) = erased else {
return Err(format!(
"optimizer Control-class axis targets control '{control}', but the phase \
declares no such control"
));
};
erased
.set_f64(
value,
ControlOrigin::Api {
source: "optimizer".into(),
},
)
.await
.map(|_rev| ())
.map_err(|e| format!("optimizer retarget '{control}' = {value}: {e}"))
}
pub async fn servo(
spec: ServoSpec,
stop_flag: Arc<AtomicBool>,
reporter: Arc<CadenceReporter>,
parent: Arc<ScopeKernel>,
phase_kernel: Arc<ScopeKernel>,
phase_component: Arc<RwLock<Component>>,
phase_done: Arc<AtomicBool>,
) -> Result<(), String> {
let mut params = OptimizerParams::new();
for (k, v) in &spec.params {
params = params.with(k.clone(), *v);
}
let optimizer = super::by_name(&spec.method, ¶ms)
.ok_or_else(|| format!("unknown optimizer method '{}'", spec.method))?;
let budget = Budget::seeded(spec.max_evals, spec.seed);
let lex: Box<dyn PullSource> = Box::new(LexSource::new(&spec.space));
let mut src = optimizer.coordinate_source(&spec.space, &budget, lex);
let mut best_value = f64::NEG_INFINITY;
let mut best_coord: Option<Coord> = None;
let mut evals = 0usize;
let mut err: Option<String> = None;
let mut batch = crate::executor::source_next(&mut src, &[]);
'outer: while let Some(coords) = batch.take() {
let mut evaluated: Vec<(Coord, f64)> = Vec::new();
for coord in coords {
if evals >= spec.max_evals || phase_done.load(Ordering::Relaxed) {
break 'outer;
}
for ca in &spec.controls {
let value = coord[ca.axis_idx].as_num();
if let Err(e) = retarget(&phase_component, &ca.control, value).await {
err = Some(e);
break 'outer;
}
}
let settle_done = Arc::new(AtomicBool::new(false));
let handle = match start_settle(
&parent,
&phase_kernel,
&spec.objective,
&reporter,
settle_done,
) {
Ok(handle) => handle,
Err(super::settle::SettleSkip::NotWindowed) => {
err = Some(format!(
"optimizer objective '{}' is not a windowed metric — a Control-class \
sweep settles the live windowed objective per setting; use \
`metric_window(...)` or `metricsql_scalar(rate(...[W]))`",
spec.objective
));
break 'outer;
}
Err(e) => {
err = Some(format!(
"optimizer objective '{}' cannot be settled: {e}",
spec.objective
));
break 'outer;
}
};
let value = loop {
if phase_done.load(Ordering::Relaxed) {
reporter.unsubscribe(handle.subscriber);
break 'outer;
}
if handle.outcome.load().is_some() {
let v = handle.register.load().value;
reporter.unsubscribe(handle.subscriber);
break v;
}
tokio::time::sleep(SETTLE_POLL).await;
};
evals += 1;
if value > best_value {
best_value = value;
best_coord = Some(coord.clone());
}
evaluated.push((coord, value));
}
batch = crate::executor::source_next(&mut src, &evaluated);
}
spec.result.store(Arc::new(ServoOutcome {
best: best_coord.map(|coord| ServoBest {
coord,
value: best_value,
}),
evals,
}));
stop_flag.store(true, Ordering::Relaxed);
match err {
Some(e) => Err(e),
None => Ok(()),
}
}