use nmbrs_workload::model::WorkloadPhase;
use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
pub const NAME: WrapperName = WrapperName::new("interval");
fn triggers(s: WrapperSubject) -> bool {
s.phase().is_some_and(|p| p.interval.is_some())
}
fn describe_assignment(s: WrapperSubject) -> Option<String> {
let p = s.phase()?;
let every = p.interval.as_deref()?;
Some(match p.repeat {
Some(n) => format!("interval: every {every} × {n}"),
None => format!("interval: every {every}, until session stop"),
})
}
inventory::submit! {
WrapperRegistration {
name: NAME,
owned_fields: &["interval", "repeat"],
triggers,
requires_inner: &[],
forbids_outer: &[],
mutually_exclusive_with: &[],
describe_assignment,
levels: &[
crate::wrapper_registry::WrapperLevel::Phase,
crate::wrapper_registry::WrapperLevel::Scenario,
crate::wrapper_registry::WrapperLevel::Session,
],
}
}
pub(crate) struct IntervalShell<'i> {
inner: &'i dyn crate::executor::ExecShell,
spec: IntervalSpec,
label: &'i str,
}
impl<'i> IntervalShell<'i> {
pub(crate) fn new(
inner: &'i dyn crate::executor::ExecShell,
spec: IntervalSpec,
label: &'i str,
) -> Self {
Self { inner, spec, label }
}
}
impl<'i> crate::executor::ExecShell for IntervalShell<'i> {
fn run<'a>(
&'a self,
ctx: &'a mut crate::executor::ExecCtx,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>,
> {
Box::pin(async move {
let mut runs: u64 = 0;
let mut last;
loop {
last = self.inner.run(&mut *ctx).await;
runs += 1;
if last.is_failure() {
break;
}
if self.spec.repeat.is_some_and(|r| runs >= r) {
break;
}
if !dwell(self.spec.interval_ms).await {
break;
}
}
crate::diag!(
crate::observer::LogLevel::Info,
"interval: '{}' schedule ended after {runs} run(s)",
self.label
);
last
})
}
fn shell_kind(&self) -> crate::executor::ShellKind {
self.inner.shell_kind()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct IntervalSpec {
pub interval_ms: u64,
pub repeat: Option<u64>,
}
pub(crate) fn for_phase(
phases: &std::collections::HashMap<String, WorkloadPhase>,
workload_params: &std::collections::HashMap<String, String>,
phase_name: &str,
) -> Option<IntervalSpec> {
let p = phases.get(phase_name)?;
let declared = p.interval.as_deref()?;
let expanded = crate::runner::expand_workload_params(declared, workload_params);
let raw = expanded.trim();
if raw.is_empty() || raw == "0" {
return None;
}
match crate::timeval::parse_time_ms(raw) {
Ok(0) => None,
Ok(ms) => Some(IntervalSpec {
interval_ms: ms,
repeat: p.repeat,
}),
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Error,
"phase '{phase_name}': `interval: {raw}` is not a duration ({e}) \
— running once"
);
None
}
}
}
pub(crate) async fn dwell(total_ms: u64) -> bool {
const TICK_MS: u64 = 250;
let mut remaining = total_ms;
while remaining > 0 {
if crate::session_signals::stop_requested() {
return false;
}
let step = remaining.min(TICK_MS);
tokio::time::sleep(std::time::Duration::from_millis(step)).await;
remaining -= step;
}
!crate::session_signals::stop_requested()
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap;
fn phase_with(interval: Option<&str>, repeat: Option<u64>) -> HashMap<String, WorkloadPhase> {
let mut m = HashMap::new();
m.insert(
"p".to_string(),
WorkloadPhase {
interval: interval.map(str::to_string),
repeat,
..Default::default()
},
);
m
}
fn no_params() -> HashMap<String, String> {
HashMap::new()
}
#[test]
fn absent_interval_yields_no_schedule() {
assert_eq!(for_phase(&phase_with(None, None), &no_params(), "p"), None);
assert_eq!(
for_phase(&phase_with(None, Some(5)), &no_params(), "p"),
None
);
assert_eq!(
for_phase(&phase_with(Some("5m"), None), &no_params(), "nope"),
None
);
}
#[test]
fn interval_resolves_to_millis_with_bound() {
assert_eq!(
for_phase(&phase_with(Some("5m"), Some(288)), &no_params(), "p"),
Some(IntervalSpec {
interval_ms: 300_000,
repeat: Some(288)
})
);
assert_eq!(
for_phase(&phase_with(Some("250ms"), None), &no_params(), "p"),
Some(IntervalSpec {
interval_ms: 250,
repeat: None
})
);
}
#[test]
fn zero_or_bad_interval_degrades_to_run_once() {
assert_eq!(
for_phase(&phase_with(Some("0"), Some(10)), &no_params(), "p"),
None
);
assert_eq!(
for_phase(&phase_with(Some(""), None), &no_params(), "p"),
None
);
assert_eq!(
for_phase(&phase_with(Some("banana"), Some(10)), &no_params(), "p"),
None
);
}
#[test]
fn interval_interpolates_a_workload_param() {
let phases = phase_with(Some("{recall_interval}"), None);
let off = HashMap::from([("recall_interval".to_string(), "0".to_string())]);
assert_eq!(for_phase(&phases, &off, "p"), None);
let on = HashMap::from([("recall_interval".to_string(), "5m".to_string())]);
assert_eq!(
for_phase(&phases, &on, "p"),
Some(IntervalSpec {
interval_ms: 300_000,
repeat: None
})
);
}
#[test]
fn triggers_only_on_a_phase_declaring_interval() {
let with = WorkloadPhase {
interval: Some("5m".into()),
..Default::default()
};
let without = WorkloadPhase::default();
assert!(triggers(WrapperSubject::Phase(&with)));
assert!(!triggers(WrapperSubject::Phase(&without)));
assert!(describe_assignment(WrapperSubject::Phase(&with)).is_some());
}
struct FakeShell(crate::executor::ShellKind);
impl crate::executor::ExecShell for FakeShell {
fn run<'a>(
&'a self,
_ctx: &'a mut crate::executor::ExecCtx,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>,
> {
Box::pin(async { crate::phase_outcome::Outcome::skipped() })
}
fn shell_kind(&self) -> crate::executor::ShellKind {
self.0
}
}
#[test]
fn layer_is_generic_over_any_shell() {
use crate::executor::{ExecShell, ShellKind};
let spec = IntervalSpec {
interval_ms: 1,
repeat: Some(1),
};
let scenario = FakeShell(ShellKind::Scenario);
assert_eq!(
IntervalShell::new(&scenario, spec, "s").shell_kind(),
ShellKind::Scenario,
"wrapping a scenario keeps it a scenario"
);
let phase = FakeShell(ShellKind::Phase);
assert_eq!(
IntervalShell::new(&phase, spec, "p").shell_kind(),
ShellKind::Phase,
"wrapping a phase keeps it a phase"
);
let session = FakeShell(ShellKind::Session);
assert_eq!(
IntervalShell::new(&session, spec, "sess").shell_kind(),
ShellKind::Session,
"wrapping a session keeps it a session"
);
}
#[test]
fn levels_filter_admits_every_shell_level() {
let reg = inventory::iter::<WrapperRegistration>
.into_iter()
.find(|r| r.name == NAME)
.expect("interval wrapper is registered");
use crate::wrapper_registry::WrapperLevel;
assert!(reg.applies_at(WrapperLevel::Phase));
assert!(reg.applies_at(WrapperLevel::Scenario));
assert!(reg.applies_at(WrapperLevel::Session));
assert!(!reg.applies_at(WrapperLevel::Op));
}
#[tokio::test]
async fn dwell_short_circuits_on_session_stop() {
let _g = crate::session_signals::STOP_GLOBAL_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
crate::session_signals::request_stop();
let t = std::time::Instant::now();
assert!(!dwell(10_000).await, "a latched stop must end the dwell");
assert!(
t.elapsed() < std::time::Duration::from_secs(1),
"dwell must not wait out the interval after a stop"
);
crate::session_signals::clear_session_stop_for_test();
}
}