nmbrs_runtime/optimize/servo.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! SRD-86 §4–§6 — the **Control-class actuation daemon**.
5//!
6//! Where a `Coordinate`-class axis is actuated by *re-running* the phase once
7//! per coordinate (`executor::dispatch_optimization`'s default loop), a
8//! `Control`-class axis is actuated by **live-retargeting** the phase's dynamic
9//! control (SRD-23 `concurrency` / `rate`) on **one continuous phase** — no
10//! restart. The optimizer becomes a servoing daemon: pull the next setting →
11//! retarget the live control (confirmed-apply, `await`ed) → wait for the
12//! windowed objective to *settle* at that setting → read it → `step`. Repeat
13//! until the budget is spent, then stop the phase.
14//!
15//! ## Why an async task, not a cadence callback
16//!
17//! The settle detector ([`super::settle`]) rides the **sync** metrics-cadence
18//! callback, but [`Control::set`](nmbrs_metrics::controls::Control::set) is
19//! **async** (confirmed-apply over async appliers). So the servo is a
20//! concurrent **async future**, [`tokio::join!`]'d with the activity loop inside
21//! `run_phase`. It reuses [`start_settle`] (which legitimately rides the sync
22//! callback) only to *read* the settled objective; the async retarget lives
23//! here, `await`ed directly so setting *N* is confirmed in effect before it is
24//! settled and read.
25//!
26//! The live control is resolved off the **phase component** — where the fiber
27//! pool / rate limiter declared it in `Activity::attach_component`, which runs
28//! *before* the activity loop, so the handle always exists by the time the
29//! servo retargets.
30
31use std::sync::atomic::{AtomicBool, Ordering};
32use std::sync::{Arc, RwLock};
33use std::time::Duration;
34
35use crate::scope_kernel::ScopeKernel;
36use arc_swap::ArcSwap;
37use nmbrs_metrics::cadence_reporter::CadenceReporter;
38use nmbrs_metrics::component::Component;
39use nmbrs_metrics::controls::ControlOrigin;
40
41use super::settle::start_settle;
42use super::{Budget, Coord, LexSource, OptimizerParams, PullSource, SearchSpace};
43
44/// One `Control`-class axis: its index in a coordinate tuple + the name of the
45/// live control (`"concurrency"` / `"rate"`) it servos.
46#[derive(Debug, Clone)]
47pub struct ControlAxis {
48 pub axis_idx: usize,
49 pub control: String,
50}
51
52/// The best coordinate + objective value a servoing run found.
53#[derive(Clone)]
54pub struct ServoBest {
55 pub coord: Coord,
56 pub value: f64,
57}
58
59/// What a servoing run produced — read by the dispatch after the continuous
60/// phase returns. `best` is `None` if no setting ever yielded a value (e.g. the
61/// phase ended before the first settle).
62#[derive(Clone, Default)]
63pub struct ServoOutcome {
64 pub best: Option<ServoBest>,
65 pub evals: usize,
66}
67
68/// The servoing job, handed from `dispatch_optimization` to `run_phase` through
69/// `ctx.optimize_servo`. The `result` cell is written by [`servo`] and read by
70/// the dispatch after the continuous phase returns. (`Clone` only so the
71/// enclosing `ExecCtx` stays `Clone`; the spec is moved, not cloned, in use.)
72#[derive(Clone)]
73pub struct ServoSpec {
74 pub method: String,
75 pub params: Vec<(String, f64)>,
76 pub objective: String,
77 pub max_evals: usize,
78 pub seed: u64,
79 pub space: SearchSpace,
80 pub controls: Vec<ControlAxis>,
81 pub result: Arc<ArcSwap<ServoOutcome>>,
82}
83
84/// Poll interval for a per-setting settle verdict (the settle detector runs on
85/// the cadence worker; the servo awaits its `outcome` cell).
86const SETTLE_POLL: Duration = Duration::from_millis(50);
87
88/// Retarget one live control to `value` (SRD-23 confirmed-apply). Resolves the
89/// erased handle off the phase component, where the applier declared it.
90async fn retarget(
91 phase_component: &Arc<RwLock<Component>>,
92 control: &str,
93 value: f64,
94) -> Result<(), String> {
95 let erased = {
96 let guard = phase_component.read().unwrap_or_else(|e| e.into_inner());
97 guard.find_control_erased_up(control)
98 };
99 let Some(erased) = erased else {
100 return Err(format!(
101 "optimizer Control-class axis targets control '{control}', but the phase \
102 declares no such control"
103 ));
104 };
105 erased
106 .set_f64(
107 value,
108 ControlOrigin::Api {
109 source: "optimizer".into(),
110 },
111 )
112 .await
113 .map(|_rev| ())
114 .map_err(|e| format!("optimizer retarget '{control}' = {value}: {e}"))
115}
116
117/// Drive the optimizer over **one continuous phase** by live-retargeting its
118/// controls (SRD-86 Control-class actuation). Runs concurrent with the activity
119/// loop (`tokio::join!` in `run_phase`):
120///
121/// 1. pull the next coordinate, **retarget** each control axis to it (`await`ed,
122/// confirmed-apply);
123/// 2. **settle** the windowed objective at that setting via [`start_settle`]
124/// (its verdict raises a *throwaway* per-setting flag, never the phase stop
125/// flag), reading the stabilized value from the settle register;
126/// 3. `step` the optimizer and track the best.
127///
128/// On budget exhaustion (or when the phase ends first — `phase_done`), it
129/// publishes the best into `spec.best` and raises `stop_flag` to end the phase.
130pub async fn servo(
131 spec: ServoSpec,
132 stop_flag: Arc<AtomicBool>,
133 reporter: Arc<CadenceReporter>,
134 parent: Arc<ScopeKernel>,
135 phase_kernel: Arc<ScopeKernel>,
136 phase_component: Arc<RwLock<Component>>,
137 phase_done: Arc<AtomicBool>,
138) -> Result<(), String> {
139 let mut params = OptimizerParams::new();
140 for (k, v) in &spec.params {
141 params = params.with(k.clone(), *v);
142 }
143 let optimizer = super::by_name(&spec.method, ¶ms)
144 .ok_or_else(|| format!("unknown optimizer method '{}'", spec.method))?;
145 let budget = Budget::seeded(spec.max_evals, spec.seed);
146 let lex: Box<dyn PullSource> = Box::new(LexSource::new(&spec.space));
147 let mut src = optimizer.coordinate_source(&spec.space, &budget, lex);
148
149 let mut best_value = f64::NEG_INFINITY;
150 let mut best_coord: Option<Coord> = None;
151 let mut evals = 0usize;
152
153 // Accumulate any fatal servoing error here; ALL exit paths still publish the
154 // best-so-far and stop the phase, so the join never hangs and the dispatch
155 // always reads a result.
156 let mut err: Option<String> = None;
157 let mut batch = crate::executor::source_next(&mut src, &[]);
158 'outer: while let Some(coords) = batch.take() {
159 let mut evaluated: Vec<(Coord, f64)> = Vec::new();
160 for coord in coords {
161 if evals >= spec.max_evals || phase_done.load(Ordering::Relaxed) {
162 break 'outer;
163 }
164
165 // (1) Retarget every control axis to this coordinate.
166 for ca in &spec.controls {
167 let value = coord[ca.axis_idx].as_num();
168 if let Err(e) = retarget(&phase_component, &ca.control, value).await {
169 err = Some(e);
170 break 'outer;
171 }
172 }
173
174 // (2) Settle the windowed objective at this setting. A throwaway
175 // flag absorbs the settle's own stop verdict so it does NOT end the
176 // phase (the servo owns the phase stop, on budget exhaustion).
177 let settle_done = Arc::new(AtomicBool::new(false));
178 let handle = match start_settle(
179 &parent,
180 &phase_kernel,
181 &spec.objective,
182 &reporter,
183 settle_done,
184 ) {
185 Ok(handle) => handle,
186 Err(super::settle::SettleSkip::NotWindowed) => {
187 err = Some(format!(
188 "optimizer objective '{}' is not a windowed metric — a Control-class \
189 sweep settles the live windowed objective per setting; use \
190 `metric_window(...)` or `metricsql_scalar(rate(...[W]))`",
191 spec.objective
192 ));
193 break 'outer;
194 }
195 Err(e) => {
196 err = Some(format!(
197 "optimizer objective '{}' cannot be settled: {e}",
198 spec.objective
199 ));
200 break 'outer;
201 }
202 };
203 let value = loop {
204 if phase_done.load(Ordering::Relaxed) {
205 reporter.unsubscribe(handle.subscriber);
206 break 'outer;
207 }
208 if handle.outcome.load().is_some() {
209 let v = handle.register.load().value;
210 reporter.unsubscribe(handle.subscriber);
211 break v;
212 }
213 tokio::time::sleep(SETTLE_POLL).await;
214 };
215
216 // (3) Step the optimizer and track the best.
217 evals += 1;
218 if value > best_value {
219 best_value = value;
220 best_coord = Some(coord.clone());
221 }
222 evaluated.push((coord, value));
223 }
224 batch = crate::executor::source_next(&mut src, &evaluated);
225 }
226
227 spec.result.store(Arc::new(ServoOutcome {
228 best: best_coord.map(|coord| ServoBest {
229 coord,
230 value: best_value,
231 }),
232 evals,
233 }));
234 // End the continuous phase (read at its next cycle boundary).
235 stop_flag.store(true, Ordering::Relaxed);
236 match err {
237 Some(e) => Err(e),
238 None => Ok(()),
239 }
240}