nmbrs_runtime/wrappers/
delay.rs1use 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
20pub const NAME: WrapperName = WrapperName::new("delay");
22
23fn triggers(s: WrapperSubject) -> bool {
25 let Some(template) = s.op() else {
26 return false;
27 };
28 template.delay.is_some()
29}
30
31fn 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
67pub struct DelayDispenser {
75 inner: Arc<dyn OpDispenser>,
76 before_handle: Option<crate::fixture::PullHandle>,
78 after_handle: Option<crate::fixture::PullHandle>,
80}
81
82impl DelayDispenser {
83 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 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}