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}