Skip to main content

nmbrs_runtime/wrappers/
delay.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Per-cycle delay wrapper. Reads delay values through the
5//! cycle's pull plan and sleeps before and/or after delegating
6//! to the inner op. u64 → nanoseconds; f64 → milliseconds.
7//!
8//! Two surface forms (see [`nmbrs_workload::model::DelaySpec`]):
9//! - `delay: <name>` — single pre-op delay
10//! - `delay: { before: <name>, after: <name> }` — independent
11//!   pre-op and/or post-op delays
12
13use std::sync::Arc;
14
15use crate::adapter::WrappingDispenser;
16use crate::adapter::{ExecutionError, OpDispenser, OpResult};
17use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
18use nmbrs_workload::model::DelaySpec;
19
20/// SRD-32a wrapper name.
21pub const NAME: WrapperName = WrapperName::new("delay");
22
23/// Trigger: an op declares any `delay:` spec.
24fn triggers(s: WrapperSubject) -> bool {
25    let Some(template) = s.op() else {
26        return false;
27    };
28    template.delay.is_some()
29}
30
31/// One-line assignment summary for init-time diagnostics.
32fn describe_assignment(s: WrapperSubject) -> Option<String> {
33    let template = s.op()?;
34    template.delay.as_ref().map(|spec| match spec {
35        DelaySpec::Before(name) => {
36            let trimmed = crate::wrapper_registrations::trim_braces(name);
37            format!("delay: delay binding `{trimmed}`")
38        }
39        DelaySpec::BeforeAfter { before, after } => {
40            let b = before.as_deref().map(|n| {
41                let t = crate::wrapper_registrations::trim_braces(n);
42                format!("before=`{t}`")
43            });
44            let a = after.as_deref().map(|n| {
45                let t = crate::wrapper_registrations::trim_braces(n);
46                format!("after=`{t}`")
47            });
48            let parts: Vec<String> = [b, a].into_iter().flatten().collect();
49            format!("delay: {}", parts.join(", "))
50        }
51    })
52}
53
54inventory::submit! {
55    WrapperRegistration {
56        name: NAME,
57        owned_fields: &["delay"],
58        triggers,
59        requires_inner: &[super::traverse::NAME],
60        forbids_outer: &[],
61        mutually_exclusive_with: &[],
62        describe_assignment,
63        levels: &[crate::wrapper_registry::WrapperLevel::Op],
64    }
65}
66
67/// Wraps an inner OpDispenser with per-cycle delays.
68///
69/// Reads delay values via `PullHandle`s from the cycle's
70/// `ResolvedPulls`. u64 values are interpreted as nanoseconds;
71/// f64 values are interpreted as milliseconds. Delays are
72/// invisible to the inner adapter — they're never in
73/// `ResolvedFields`.
74pub struct DelayDispenser {
75    inner: Arc<dyn OpDispenser>,
76    /// Pre-op delay handle, when configured.
77    before_handle: Option<crate::fixture::PullHandle>,
78    /// Post-op delay handle, when configured.
79    after_handle: Option<crate::fixture::PullHandle>,
80}
81
82impl DelayDispenser {
83    /// Wrap an inner dispenser with a single pre-op delay
84    /// binding. Backwards-compatible entry point for the bare-
85    /// string `delay: <name>` form.
86    pub fn wrap(
87        inner: Arc<dyn OpDispenser>,
88        delay_field: &str,
89        fx: &mut crate::fixture::ScopeFixture,
90    ) -> Result<Arc<dyn OpDispenser>, String> {
91        let before_handle = Some(
92            fx.register_pull(delay_field)
93                .map_err(|e| format!("delay: {e}"))?,
94        );
95        Ok(Arc::new(Self {
96            inner,
97            before_handle,
98            after_handle: None,
99        }))
100    }
101
102    /// Wrap an inner dispenser with optional pre-op and/or
103    /// post-op delays. At least one must be set; the caller is
104    /// expected to have validated that.
105    pub fn wrap_before_after(
106        inner: Arc<dyn OpDispenser>,
107        before_name: Option<&str>,
108        after_name: Option<&str>,
109        fx: &mut crate::fixture::ScopeFixture,
110    ) -> Result<Arc<dyn OpDispenser>, String> {
111        let before_handle = match before_name {
112            Some(name) => Some(
113                fx.register_pull(name)
114                    .map_err(|e| format!("delay.before: {e}"))?,
115            ),
116            None => None,
117        };
118        let after_handle = match after_name {
119            Some(name) => Some(
120                fx.register_pull(name)
121                    .map_err(|e| format!("delay.after: {e}"))?,
122            ),
123            None => None,
124        };
125        if before_handle.is_none() && after_handle.is_none() {
126            return Err("delay: empty before/after — at least one must be set".into());
127        }
128        Ok(Arc::new(Self {
129            inner,
130            before_handle,
131            after_handle,
132        }))
133    }
134}
135
136fn value_to_nanos(value: &polydat::ast::Value) -> u64 {
137    match value {
138        polydat::ast::Value::U64(ns) => *ns,
139        polydat::ast::Value::F64(ms) => (*ms * 1_000_000.0) as u64,
140        _ => 0,
141    }
142}
143
144impl WrappingDispenser for DelayDispenser {}
145
146impl OpDispenser for DelayDispenser {
147    fn execute<'a>(
148        &'a self,
149        cycle: u64,
150        ctx: &'a crate::fixture::ExecCtx<'a>,
151    ) -> std::pin::Pin<
152        Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
153    > {
154        Box::pin(async move {
155            if let Some(h) = self.before_handle {
156                let nanos = value_to_nanos(ctx.pulls.get(h));
157                if nanos > 0 {
158                    tokio::time::sleep(std::time::Duration::from_nanos(nanos)).await;
159                }
160            }
161            let result = self.inner.execute(cycle, ctx).await?;
162            if let Some(h) = self.after_handle {
163                let nanos = value_to_nanos(ctx.pulls.get(h));
164                if nanos > 0 {
165                    tokio::time::sleep(std::time::Duration::from_nanos(nanos)).await;
166                }
167            }
168            Ok(result)
169        })
170    }
171    fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
172        Some(self.inner.as_ref())
173    }
174}