Skip to main content

nmbrs_runtime/wrappers/
interval.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Interval wrapper — the first PHASE-level wrapper (SRD-82/92 cross-level;
5//! see `docs/cross-level-wrapper-cascade-scope.md`). `interval:` is its
6//! sigil: re-run this phase, dwelling `interval` between runs, bounded by
7//! `repeat:`.
8//!
9//! Why a PHASE wrapper and not a rate: a rate paces the units *inside* a
10//! phase (for a recall phase, individual queries), so `rate: 1/300` would
11//! trickle one query every 5 minutes — a single glacial pass, not a
12//! measurement every 5 minutes. Pacing the *phase* is one layer out: each
13//! iteration is a whole phase run (a complete recall measurement), and the
14//! repeats live in ONE session — one metrics timeline — which an external
15//! `while … sleep` loop cannot give you (each `nmbrs run` is its own session).
16//!
17//! - **No `interval:`** → the phase runs exactly once (unchanged).
18//! - **`interval: <dur>` + `repeat: N`** → up to `N` runs, dwelling
19//!   `interval` between them.
20//! - **`interval: <dur>` with no `repeat:`** → runs until the session stops.
21//!
22//! The dwell is cooperative: it wakes on a short tick to observe the session
23//! stop, so Ctrl-C / a `stop_when` `action: abort` cuts the wait immediately
24//! instead of waiting out the remaining interval. A failing run ends the
25//! repeat (the phase's own outcome propagates unchanged).
26//!
27//! Registration is declared here at `WrapperLevel::Phase` so the field
28//! validation + telemetry stay consistent with the op wrappers; the
29//! construction is hooked at the phase seam (`PhaseShell::run`) because a
30//! phase layer wraps an `ExecShell`, not an `OpDispenser`.
31
32use nmbrs_workload::model::WorkloadPhase;
33
34use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
35
36pub const NAME: WrapperName = WrapperName::new("interval");
37
38/// Registry trigger: the phase's own `interval:` field. `repeat:` alone is
39/// inert (it bounds an interval that isn't there) — the misplaced-field
40/// guard reports it against THIS registration.
41fn triggers(s: WrapperSubject) -> bool {
42    s.phase().is_some_and(|p| p.interval.is_some())
43}
44
45fn describe_assignment(s: WrapperSubject) -> Option<String> {
46    let p = s.phase()?;
47    let every = p.interval.as_deref()?;
48    Some(match p.repeat {
49        Some(n) => format!("interval: every {every} × {n}"),
50        None => format!("interval: every {every}, until session stop"),
51    })
52}
53
54inventory::submit! {
55    WrapperRegistration {
56        name: NAME,
57        owned_fields: &["interval", "repeat"],
58        triggers,
59        requires_inner: &[],
60        forbids_outer: &[],
61        mutually_exclusive_with: &[],
62        describe_assignment,
63        // The layer ([`IntervalShell`]) is implemented GENERICALLY over
64        // `ExecShell`, so `levels` is a pure TYPE FILTER over the subjects it
65        // accepts — not an implementation limit. Any shell level can carry
66        // it: re-running a scenario every N is the same operation as
67        // re-running a phase, and the layer is identical. `Op` is excluded
68        // only because the op leaf is deliberately NOT an `ExecShell` (it
69        // sits below the Outcome projection boundary — SRD-82 Decision B),
70        // so there is no op shell to decorate; `Stanza` likewise has no shell.
71        levels: &[
72            crate::wrapper_registry::WrapperLevel::Phase,
73            crate::wrapper_registry::WrapperLevel::Scenario,
74            crate::wrapper_registry::WrapperLevel::Session,
75        ],
76    }
77}
78
79/// The `interval:` LAYER — a generic [`ExecShell`] decorator.
80///
81/// It knows nothing about phases: it wraps ANY inner shell and re-runs it,
82/// dwelling `interval` between runs, bounded by `repeat`. Placing it around a
83/// `ScenarioShell` (or a future `SessionShell`) requires no change here —
84/// only a subject at that level to resolve the schedule from. `shell_kind`
85/// delegates to the inner shell because a layer decorates a level, it does
86/// not change what level the thing IS.
87///
88/// Semantics: a failing run ends the schedule (its outcome propagates
89/// unchanged, exactly as an unwrapped run's would); the dwell is cooperative
90/// so a stop cuts the wait instead of stranding the run for the remainder.
91pub(crate) struct IntervalShell<'i> {
92    inner: &'i dyn crate::executor::ExecShell,
93    spec: IntervalSpec,
94    /// The wrapped unit's name, for the schedule's summary line.
95    label: &'i str,
96}
97
98impl<'i> IntervalShell<'i> {
99    pub(crate) fn new(
100        inner: &'i dyn crate::executor::ExecShell,
101        spec: IntervalSpec,
102        label: &'i str,
103    ) -> Self {
104        Self { inner, spec, label }
105    }
106}
107
108impl<'i> crate::executor::ExecShell for IntervalShell<'i> {
109    fn run<'a>(
110        &'a self,
111        ctx: &'a mut crate::executor::ExecCtx,
112    ) -> std::pin::Pin<
113        Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>,
114    > {
115        Box::pin(async move {
116            let mut runs: u64 = 0;
117            let mut last;
118            loop {
119                // Reborrow per iteration — the inner shell takes `&mut ExecCtx`
120                // and we drive it repeatedly.
121                last = self.inner.run(&mut *ctx).await;
122                runs += 1;
123                // A failing run ends the schedule.
124                if last.is_failure() {
125                    break;
126                }
127                // Bound reached; `repeat: None` = until the session stops.
128                if self.spec.repeat.is_some_and(|r| runs >= r) {
129                    break;
130                }
131                // Cooperative dwell — a stop ends the schedule immediately.
132                if !dwell(self.spec.interval_ms).await {
133                    break;
134                }
135            }
136            crate::diag!(
137                crate::observer::LogLevel::Info,
138                "interval: '{}' schedule ended after {runs} run(s)",
139                self.label
140            );
141            last
142        })
143    }
144
145    fn shell_kind(&self) -> crate::executor::ShellKind {
146        self.inner.shell_kind()
147    }
148}
149
150/// A phase's resolved repeat schedule.
151#[derive(Debug, Clone, Copy, PartialEq, Eq)]
152pub(crate) struct IntervalSpec {
153    /// Dwell between runs, milliseconds.
154    pub interval_ms: u64,
155    /// Total runs. `None` = until the session stops.
156    pub repeat: Option<u64>,
157}
158
159/// Resolve the schedule for `phase_name`, or `None` when the phase is not
160/// scheduled (the overwhelmingly common case — it runs once).
161///
162/// The raw `interval:` goes through `{param}` interpolation first, so a
163/// phase SHARED between an unscheduled use and a scheduled one needs no
164/// duplicate: it declares `interval: "{some_param}"`, the param defaults to
165/// the disable value, and a run opts in from the CLI (`some_param=5m`).
166///
167/// `0` / empty is the DISABLE knob — run once, silently (the same idiom as
168/// `retry_backoff: 0` disabling retry pacing). A zero dwell would otherwise
169/// spin the phase with no pacing at all. An unparseable duration is a loud
170/// error that still degrades to running once, never to a spin.
171pub(crate) fn for_phase(
172    phases: &std::collections::HashMap<String, WorkloadPhase>,
173    workload_params: &std::collections::HashMap<String, String>,
174    phase_name: &str,
175) -> Option<IntervalSpec> {
176    let p = phases.get(phase_name)?;
177    let declared = p.interval.as_deref()?;
178    let expanded = crate::runner::expand_workload_params(declared, workload_params);
179    let raw = expanded.trim();
180    if raw.is_empty() || raw == "0" {
181        return None;
182    }
183    match crate::timeval::parse_time_ms(raw) {
184        Ok(0) => None,
185        Ok(ms) => Some(IntervalSpec {
186            interval_ms: ms,
187            repeat: p.repeat,
188        }),
189        Err(e) => {
190            crate::diag!(
191                crate::observer::LogLevel::Error,
192                "phase '{phase_name}': `interval: {raw}` is not a duration ({e}) \
193                 — running once"
194            );
195            None
196        }
197    }
198}
199
200/// Cooperative dwell. Sleeps `total_ms`, waking every [`TICK_MS`] to observe
201/// the session stop. Returns `false` if a stop was seen (the caller must not
202/// start another run) — so an abort or Ctrl-C during a long interval cuts the
203/// wait instead of stranding the run for the remainder.
204pub(crate) async fn dwell(total_ms: u64) -> bool {
205    const TICK_MS: u64 = 250;
206    let mut remaining = total_ms;
207    while remaining > 0 {
208        if crate::session_signals::stop_requested() {
209            return false;
210        }
211        let step = remaining.min(TICK_MS);
212        tokio::time::sleep(std::time::Duration::from_millis(step)).await;
213        remaining -= step;
214    }
215    !crate::session_signals::stop_requested()
216}
217
218#[cfg(test)]
219mod tests {
220    use super::*;
221    use std::collections::HashMap;
222
223    fn phase_with(interval: Option<&str>, repeat: Option<u64>) -> HashMap<String, WorkloadPhase> {
224        let mut m = HashMap::new();
225        m.insert(
226            "p".to_string(),
227            WorkloadPhase {
228                interval: interval.map(str::to_string),
229                repeat,
230                ..Default::default()
231            },
232        );
233        m
234    }
235
236    fn no_params() -> HashMap<String, String> {
237        HashMap::new()
238    }
239
240    /// No `interval:` → no schedule (the phase runs once).
241    #[test]
242    fn absent_interval_yields_no_schedule() {
243        assert_eq!(for_phase(&phase_with(None, None), &no_params(), "p"), None);
244        // `repeat:` alone does not conjure a schedule.
245        assert_eq!(
246            for_phase(&phase_with(None, Some(5)), &no_params(), "p"),
247            None
248        );
249        // An unknown phase name is simply not scheduled.
250        assert_eq!(
251            for_phase(&phase_with(Some("5m"), None), &no_params(), "nope"),
252            None
253        );
254    }
255
256    /// A duration + bound resolves to milliseconds.
257    #[test]
258    fn interval_resolves_to_millis_with_bound() {
259        assert_eq!(
260            for_phase(&phase_with(Some("5m"), Some(288)), &no_params(), "p"),
261            Some(IntervalSpec {
262                interval_ms: 300_000,
263                repeat: Some(288)
264            })
265        );
266        // No `repeat:` = until session stop.
267        assert_eq!(
268            for_phase(&phase_with(Some("250ms"), None), &no_params(), "p"),
269            Some(IntervalSpec {
270                interval_ms: 250,
271                repeat: None
272            })
273        );
274    }
275
276    /// `0` / empty is the DISABLE knob — run once, so a shared phase stays
277    /// unscheduled by default. A malformed duration also degrades to once,
278    /// never to a spin.
279    #[test]
280    fn zero_or_bad_interval_degrades_to_run_once() {
281        assert_eq!(
282            for_phase(&phase_with(Some("0"), Some(10)), &no_params(), "p"),
283            None
284        );
285        assert_eq!(
286            for_phase(&phase_with(Some(""), None), &no_params(), "p"),
287            None
288        );
289        assert_eq!(
290            for_phase(&phase_with(Some("banana"), Some(10)), &no_params(), "p"),
291            None
292        );
293    }
294
295    /// The `{param}` interpolation that lets ONE shared phase be scheduled
296    /// per-run: the declared `interval: "{recall_interval}"` is disabled by
297    /// the param's default and opted into from the CLI.
298    #[test]
299    fn interval_interpolates_a_workload_param() {
300        let phases = phase_with(Some("{recall_interval}"), None);
301        // Default (disabled) → unscheduled: the shared phase runs once.
302        let off = HashMap::from([("recall_interval".to_string(), "0".to_string())]);
303        assert_eq!(for_phase(&phases, &off, "p"), None);
304        // Opted in → scheduled.
305        let on = HashMap::from([("recall_interval".to_string(), "5m".to_string())]);
306        assert_eq!(
307            for_phase(&phases, &on, "p"),
308            Some(IntervalSpec {
309                interval_ms: 300_000,
310                repeat: None
311            })
312        );
313    }
314
315    /// The trigger fires on a phase with `interval:`, and never on an op
316    /// subject (this wrapper is `WrapperLevel::Phase`).
317    #[test]
318    fn triggers_only_on_a_phase_declaring_interval() {
319        let with = WorkloadPhase {
320            interval: Some("5m".into()),
321            ..Default::default()
322        };
323        let without = WorkloadPhase::default();
324        assert!(triggers(WrapperSubject::Phase(&with)));
325        assert!(!triggers(WrapperSubject::Phase(&without)));
326        assert!(describe_assignment(WrapperSubject::Phase(&with)).is_some());
327    }
328
329    /// A stand-in shell at an arbitrary level — exists to prove the layer is
330    /// a generic `ExecShell` decorator, not a phase-specific loop.
331    struct FakeShell(crate::executor::ShellKind);
332
333    impl crate::executor::ExecShell for FakeShell {
334        fn run<'a>(
335            &'a self,
336            _ctx: &'a mut crate::executor::ExecCtx,
337        ) -> std::pin::Pin<
338            Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>,
339        > {
340            Box::pin(async { crate::phase_outcome::Outcome::skipped() })
341        }
342        fn shell_kind(&self) -> crate::executor::ShellKind {
343            self.0
344        }
345    }
346
347    /// The layer is generic over ANY shell — it wraps a SCENARIO shell just
348    /// as readily as a phase one (this only compiles because `IntervalShell`
349    /// decorates `dyn ExecShell`, with no phase in its type). `levels:` is a
350    /// filter over which subjects may trigger it, NOT a limit on what the
351    /// layer can wrap. And a layer DECORATES a level — it never changes what
352    /// level the wrapped thing is, so `shell_kind` delegates.
353    #[test]
354    fn layer_is_generic_over_any_shell() {
355        use crate::executor::{ExecShell, ShellKind};
356        let spec = IntervalSpec {
357            interval_ms: 1,
358            repeat: Some(1),
359        };
360
361        let scenario = FakeShell(ShellKind::Scenario);
362        assert_eq!(
363            IntervalShell::new(&scenario, spec, "s").shell_kind(),
364            ShellKind::Scenario,
365            "wrapping a scenario keeps it a scenario"
366        );
367
368        let phase = FakeShell(ShellKind::Phase);
369        assert_eq!(
370            IntervalShell::new(&phase, spec, "p").shell_kind(),
371            ShellKind::Phase,
372            "wrapping a phase keeps it a phase"
373        );
374
375        let session = FakeShell(ShellKind::Session);
376        assert_eq!(
377            IntervalShell::new(&session, spec, "sess").shell_kind(),
378            ShellKind::Session,
379            "wrapping a session keeps it a session"
380        );
381    }
382
383    /// The registration's level filter is permissive by default: a layer that
384    /// works over a phase is allowed over a scenario and a session too.
385    #[test]
386    fn levels_filter_admits_every_shell_level() {
387        let reg = inventory::iter::<WrapperRegistration>
388            .into_iter()
389            .find(|r| r.name == NAME)
390            .expect("interval wrapper is registered");
391        use crate::wrapper_registry::WrapperLevel;
392        assert!(reg.applies_at(WrapperLevel::Phase));
393        assert!(reg.applies_at(WrapperLevel::Scenario));
394        assert!(reg.applies_at(WrapperLevel::Session));
395        // The op leaf is NOT an ExecShell (below the Outcome projection
396        // boundary), so there is no op shell to decorate.
397        assert!(!reg.applies_at(WrapperLevel::Op));
398    }
399
400    /// The dwell returns promptly (and reports stopped) when the session
401    /// stop is already latched — it must not wait out the interval.
402    #[tokio::test]
403    async fn dwell_short_circuits_on_session_stop() {
404        let _g = crate::session_signals::STOP_GLOBAL_TEST_LOCK
405            .lock()
406            .unwrap_or_else(|e| e.into_inner());
407        crate::session_signals::request_stop();
408        let t = std::time::Instant::now();
409        // A 10s dwell must return immediately, not after 10s.
410        assert!(!dwell(10_000).await, "a latched stop must end the dwell");
411        assert!(
412            t.elapsed() < std::time::Duration::from_secs(1),
413            "dwell must not wait out the interval after a stop"
414        );
415        crate::session_signals::clear_session_stop_for_test();
416    }
417}