Skip to main content

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, &params)
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}