nmbrs-runtime 0.4.0

Workload execution runtime for nmbrs
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
// Copyright 2024-2026 Jonathan Shook
// SPDX-License-Identifier: Apache-2.0

//! SRD-86 §"Settling via the cadence pulse" — the per-pulse settle
//! interpreter and its [`PulseEvaluator`] adapter.
//!
//! For a *volatile* optimizer objective (one defined over a windowed
//! run-produced metric such as `metric_window("...","errors")` or a
//! `metricsql_*` reader), the objective value chases the live cadence
//! window and at phase completion the trailing window is empty. The
//! objective therefore cannot be read by a one-shot post-execution
//! pull; it must be *settled* across the run and held in a register the
//! executor reads at completion.
//!
//! [`SettleInterpreter`] is the settle-signal engine. It owns one
//! persistent settle-kernel — typically the phase's objective bindings
//! plus `(stable_value, stable) := is_stable(<objective>, …)` — and is
//! driven once per cadence pulse:
//!
//! 1. `set_input` on the kernel's poke input advances the generation,
//!    which dirties **every** non-deterministic node (engine rule, see
//!    `kernel::engines::set_input`): the embedded volatile objective
//!    reader re-reads the latest published window and
//!    [`is_stable`](polydat::library::stability) re-evaluates — exactly
//!    one new sample on its ring.
//! 2. Pull `stable` then `stable_value`; both land on that single
//!    evaluation (generation cache), so the multi-output node
//!    contributes one sample per pulse, not two.
//! 3. Publish `stable_value` into the shared register (an [`ArcSwap`],
//!    lock-free for the executor's completion read).
//!
//! The interpreter does **not** decide phase stop — that is the
//! [`SettleEvaluator`]'s job (settled ⇒ `interrupted`; timed out ⇒
//! `failed`), driven through the general
//! [`super::phase_pulse::PhaseStopEvaluator`] callback registered on the
//! metrics cadence feed.

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;

/// Convert a pulled objective wire to f64 with the same numeric
/// coercion as `read_objective_at_completion` (F64 as-is, U64 widened,
/// Bool 0/1); a non-numeric objective reads as 0.0 (the
/// volatile-objective gate upstream only admits numeric readers).
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,
    }
}

/// The latest reading a [`SettleInterpreter`] publishes.
#[derive(Clone, Copy, Debug)]
pub struct SettleReading {
    /// The stabilized objective value — `is_stable`'s `stable_value`
    /// output (the median of the recent window). This is what the
    /// phase executor reads as the objective at completion.
    pub value: f64,
    /// Whether the signal reached steady state on the latest pulse.
    pub stable: bool,
    /// Count of pulses delivered so far (viability / diagnostics).
    pub pulses: u64,
}

impl Default for SettleReading {
    fn default() -> Self {
        Self {
            value: 0.0,
            stable: false,
            pulses: 0,
        }
    }
}

/// Per-pulse interpreter of objective settling. See the module docs.
pub struct SettleInterpreter {
    /// A standalone kernel on whatever engine the compile chose — it
    /// takes no part in the scope tree, so it is driven through the
    /// engine-neutral [`Kernel`] trait.
    kernel: Box<dyn Kernel>,
    /// Index of the kernel's `samples: vec_f64` input, the window
    /// `is_stable` judges.
    samples_input: usize,
    /// The most recent samples, oldest first, at most `horizon` of them.
    /// `is_stable` is a pure function of its window, so the interpreter
    /// owns the history.
    window: std::collections::VecDeque<f64>,
    horizon: usize,
    value_wire: String,
    stable_wire: String,
    register: Arc<ArcSwap<SettleReading>>,
    pulses: u64,
}

impl SettleInterpreter {
    /// Build an interpreter over a compiled settle kernel. `samples` is
    /// its `vec_f64` window input; `value_wire` / `stable_wire` are the
    /// `is_stable` multi-output wire names; `horizon` is how many recent
    /// samples the window holds.
    ///
    /// # Panics
    /// When the kernel has no `samples` input: a malformed settle kernel.
    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,
        }
    }

    /// The shared register handle. The executor reads the settled
    /// objective from this at phase completion; the interpreter
    /// publishes into it on every pulse.
    pub fn register(&self) -> Arc<ArcSwap<SettleReading>> {
        self.register.clone()
    }

    /// Deliver one cadence pulse: append `sample` to the window (dropping
    /// the oldest past `horizon`), write the window to the kernel, read
    /// `is_stable`'s verdict on it, publish the stabilized value, and
    /// return the reading.
    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();
        // `samples` is a `vec_f64` extern, so a `VecF64` write cannot be
        // refused; a refusal is a malformed settle kernel.
        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
    }
}

/// The settle detector as a [`PulseEvaluator`]. It holds the phase's
/// **objective kernel** (a clone of node X's kernel — already bound to
/// this evaluation's coordinate, the same kernel
/// `read_objective_at_completion` pulls) and a fixed `is_stable`
/// engine. Each cadence pulse:
///
/// 1. positions the objective kernel at the pulse's ordinal on its
///    coordinate input (typically `cycle`); a volatile objective reader
///    re-reads the latest published window on every pull regardless;
/// 2. pulls the objective wire — the fresh windowed objective value;
/// 3. feeds it to [`SettleInterpreter`] (`is_stable`).
///
/// It yields a terminal [`Outcome`] when the objective settles
/// (`interrupted` — stopped early, register trustworthy) or when the
/// settle `timeout` elapses without settling (`failed` — SRD-86 §6
/// step 5). `None` while the loop should hold.
pub struct SettleEvaluator {
    objective: ScopeKernel,
    objective_wire: String,
    poke: Option<usize>,
    interp: SettleInterpreter,
    timeout: Duration,
    /// SRD-86 viability gate — minimum WALL-CLOCK a coordinate must run before a
    /// stable verdict is trusted, so the windowed objective's rollup has cleared
    /// the prior coordinate (and the leading transient). Wall-clock, not pulse
    /// count: under concurrent scheduling cadence pulses are delivered in
    /// bursts (many per cadence interval), so a pulse gate collapses to far less
    /// than the window — the gate must measure real time.
    min_viable: Duration,
    started: Option<Instant>,
    pulses: u64,
}

impl SettleEvaluator {
    /// `objective` is the phase's objective kernel (node X clone);
    /// `objective_wire` the objective output; `poke_input` the coordinate
    /// input positioned at each pulse's ordinal (typically `cycle`);
    /// `interp` the `is_stable` engine fed the objective value.
    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,
        }
    }

    /// The settled-value register (grab it before boxing the evaluator).
    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;
        // Position the objective kernel at this pulse, then read it; a
        // volatile reader in its cone re-reads the latest published
        // window on every pull.
        if let Some(idx) = self.poke {
            let k = &mut self.objective;
            let coords = k.coord_count();
            if idx < coords {
                // A coordinate is positioned with `set_inputs`; the pulse
                // then invalidates every output so a volatile reader
                // re-reads (native_scope_trees.md §3, the settle pulse).
                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));
        // SRD-89 — a NaN objective is a windowed metric reading **no data** (an
        // empty `rate(...[W])` lookback — see `nodes::no_data_value`), distinct
        // from a real 0. HOLD on it: do not feed the stability detector (a
        // fabricated 0 would let `is_stable` settle on the empty leading reads,
        // mis-converging the optimizer under concurrency where early windows are
        // routinely empty) and do not advance the register. The timeout still
        // bounds the wait, so a window that never produces data fails rather
        // than hanging.
        if obj.is_nan() {
            if start.elapsed() >= self.timeout {
                return Some(Outcome::failed());
            }
            return None;
        }
        let reading = self.interp.pulse(obj);
        // SRD-86 viability gate — do NOT trust a stable verdict until the
        // coordinate has run for `min_viable` of WALL-CLOCK, so the windowed
        // objective's rollup has cleared the prior coordinate (and its own
        // leading transient). At a coordinate's START the windowed objective is
        // a stable run of stale data — `rate(errors_total[W])` reads ~0 before
        // the first error registers at warmup, and at a transition it reads the
        // PRIOR coordinate's value drifting out of the window. `is_stable`
        // (which fires on `SETTLE_MIN_SAMPLES`, and whose relative margin admits
        // a slow drift as "stable") would latch that stale value, mis-converging
        // the optimizer (it keeps a saturating `concurrency` because it "saw no
        // errors", or accepts a half-cleared transition value). The gate is
        // wall-clock, not pulse count: under concurrent scheduling cadence
        // pulses arrive in bursts, so a pulse gate collapses to far less than
        // the window — only real elapsed time guarantees the window has cleared.
        if reading.stable && start.elapsed() >= self.min_viable {
            return Some(Outcome::interrupted());
        }
        if start.elapsed() >= self.timeout {
            return Some(Outcome::failed());
        }
        None
    }
}

/// Metrics-reader node names whose presence makes a phase's objective
/// potentially volatile (read from the live cadence feed). The
/// `metric` / `metric_window` stat-readers and the four `metricsql_*`
/// readers. (First-push heuristic: program-wide presence, not
/// objective-cone-precise — a non-reader objective in a phase that also
/// reads metrics elsewhere would settle trivially on its constant.)
const READER_NODES: &[&str] = &[
    "metric",
    "metric_window",
    "metricsql",
    "metricsql_scalar",
    "metricsql_vector",
    "metricsql_window",
];

/// True if `program` contains a metrics-reader node — the signal that
/// the phase's objective may be a volatile windowed metric that the
/// one-shot post-completion read cannot capture.
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()))
}

/// Node name of the session-cumulative reader (`metric(...)`), which
/// reads `MetricsQuery::session_lifetime` — a running total over ALL
/// coordinates. Unlike the windowed readers it has no bounded lookback,
/// so the viability gate cannot scope it to one coordinate.
const SESSION_CUMULATIVE_READER: &str = "metric";

/// True if `program` reads the session-cumulative `metric(...)` reader —
/// an objective aggregated across every coordinate, which the per-eval
/// warmup gate cannot isolate (no window to clear). The settle warns
/// rather than silently treating it as per-coordinate.
fn program_reads_session_cumulative_metrics(program: &PolydatProgram) -> bool {
    (0..program.node_count()).any(|i| program.node_meta(i).name == SESSION_CUMULATIVE_READER)
}

/// Warn at most once per distinct objective string that it reads a
/// session-cumulative metric. Returns `true` the first time a given
/// objective is seen (so the caller emits the diagnostic once, not once
/// per coordinate across a search).
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())
}

// First-push settle parameters (the SRD-86 §6 `settle:` YAML surface +
// finer-cadence reconfiguration are deferred). Sized for the default
// 1 s cadence and short phases: `is_stable`'s windowed median is
// published every pulse regardless of the `stable` flag, so the
// register holds a smoothed objective even before a full settle; the
// 8-deep horizon / 4-sample warmup let a steady objective settle (and
// stop the phase early) within a handful of pulses. The generous
// timeout means a phase that completes first simply uses its smoothed
// value — only a genuinely non-settling long phase trips `failed`.
//
// The settle is gated by a viability horizon (see `SettleEvaluator::evaluate`):
// a stable verdict is only honored once `SETTLE_HORIZON` pulses have been
// delivered, so the verdict is always taken over a full horizon of
// in-coordinate samples (≈ `SETTLE_HORIZON × cadence` of wall-clock, the rollup
// window by the usual `window = horizon × cadence` sizing). Without it, a
// windowed objective's leading transient — `rate(...[W])` reading ~0 before the
// coordinate's first data lands — is itself momentarily "stable", and
// `is_stable` (firing on `SETTLE_MIN_SAMPLES`) latches that phantom value; under
// concurrent scheduling a sub-window eval did exactly that.
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);

/// What the executor holds while a settle detector runs: the cadence
/// subscription to tear down, the settled-value register, and the
/// terminal-disposition cell.
pub struct SettleHandle {
    pub subscriber: SubscriberId,
    pub register: Arc<ArcSwap<SettleReading>>,
    pub outcome: StopOutcomeCell,
}

/// Why [`start_settle`] started no detector.
#[derive(Debug)]
pub enum SettleSkip {
    /// The objective reads no live metric: its one-shot read at phase
    /// completion is the value, and there is nothing to settle.
    NotWindowed,
    /// The metrics cadence is disabled, so no pulse would ever arrive.
    CadenceDisabled,
    /// The detector could not be built.
    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}")
            }
        }
    }
}

/// Start a cadence-fed settle detector for `objective` on the running
/// phase **iff** the objective reads live metrics. Builds the objective
/// kernel as node X's program rebound to `parent` (carrying the
/// coordinate — mirrors `read_objective_at_completion`), wraps it in a
/// [`PhaseStopEvaluator`], and subscribes it to the smallest cadence.
/// Starts nothing for a non-volatile objective (the one-shot read path
/// is correct there) or when the metrics cadence is disabled, and says
/// which; a detector that cannot be built is [`SettleSkip::Failed`].
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);
    }
    // A session-cumulative `metric(...)` objective has no bounded window
    // for the gate to scope, so it cannot isolate per-coordinate — warn
    // once and servo the author to a windowed reader.
    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,
    );
    // Viability gate = the stability horizon's worth of cadence intervals, in
    // WALL-CLOCK. With the usual `window = SETTLE_HORIZON × cadence` sizing this
    // is the rollup window — long enough for the objective's window to clear the
    // prior coordinate before a stable verdict is honored.
    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();
    // SRD-88 — bind this subscription to THIS execution's context so its
    // delivery fiber pulls the objective as the owning execution (the
    // metric read scopes to its own `exec_id`). Captured in the
    // execution's scope here; `None` in single-run (A1).
    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};

    /// The fixed `is_stable` engine fed the per-pulse objective value.
    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")
    }

    // A constant objective: settles regardless of the poke.
    const STEADY_OBJ: &str = "input cycle: u64\nobj := 5.0";
    // A ramping objective = the poke value (cycle): never settles.
    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,
        );
        // First pulse starts the clock; the ramp never settles.
        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() {
        // A steady objective is "stable" almost immediately, but the gate
        // withholds the settle until `min_viable` of WALL-CLOCK has elapsed —
        // so a windowed objective's rollup has cleared the prior coordinate /
        // warmup transient before its value is trusted (the bug that let a
        // sub-window concurrent eval latch a phantom score).
        let mut ev = SettleEvaluator::new(
            obj_kernel(STEADY_OBJ),
            "obj",
            "cycle",
            settle_interp(),
            Duration::from_secs(60),
            Duration::from_millis(60),
        );
        // Many pulses arrive in a burst (as under concurrent scheduling): the
        // objective is stable, but the gate holds because no real time passed.
        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() {
        // `metric(...)` reads session_lifetime (cumulative across all
        // coordinates) → flagged as un-gateable. `metric_window(...)` is
        // a bounded windowed reader, and a plain objective reads nothing
        // → neither is flagged.
        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()));
    }
}