use std::sync::Arc;
use crate::adapter::WrappingDispenser;
use crate::adapter::{AdapterError, ExecutionError, OpDispenser, OpResult};
use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
pub const NAME: WrapperName = WrapperName::new("poll");
fn triggers(s: WrapperSubject) -> bool {
let Some(template) = s.op() else {
return false;
};
template
.params
.get("poll")
.map(|v| v.is_string() || v.is_object())
.unwrap_or(false)
}
fn describe_assignment(s: WrapperSubject) -> Option<String> {
let template = s.op()?;
let poll_val = template.params.get("poll")?;
let (mode, interval, timeout): (String, u64, u64) = match poll_val {
v if v.is_string() => (v.as_str().unwrap().to_string(), 1000, 300_000),
v if v.is_object() => {
let m = v.as_object().unwrap();
let mode = m
.get("mode")
.and_then(|x| x.as_str())
.unwrap_or("await_empty")
.to_string();
let interval = m
.get("interval_ms")
.and_then(crate::wrapper_registrations::json_to_u64)
.unwrap_or(1000);
let timeout = m
.get("timeout_ms")
.and_then(crate::wrapper_registrations::json_to_u64)
.unwrap_or(300_000);
(mode, interval, timeout)
}
_ => return None,
};
Some(format!(
"poll: every {}ms, timeout {}ms, on `{mode}`",
interval, timeout
))
}
inventory::submit! {
WrapperRegistration {
name: NAME,
owned_fields: &[
"poll",
],
triggers,
requires_inner: &[super::traverse::NAME],
forbids_outer: &[],
mutually_exclusive_with: &[],
describe_assignment,
levels: &[crate::wrapper_registry::WrapperLevel::Op],
}
}
pub struct PollingDispenser {
inner: Arc<dyn OpDispenser>,
poll_interval: std::time::Duration,
timeout: std::time::Duration,
max_error_retries: u32,
stop: crate::session_signals::StopView,
metric_name: Option<String>,
min_rows: u64,
max_rows: u64,
json_path: Option<String>,
each_memo: Option<String>,
memo_state: Option<Arc<arc_swap::ArcSwap<String>>>,
progress_template: Option<String>,
each_gutter: Option<(crate::wrappers::gutter::GutterKind, String)>,
gutter_state: Option<Arc<arc_swap::ArcSwapOption<crate::wrappers::gutter::GutterSpec>>>,
iteration_gauges:
Option<Arc<arc_swap::ArcSwapOption<Vec<crate::wrappers::metrics::MetricSlot>>>>,
until: bool,
on_done: Vec<(String, String)>,
activity_metrics: Option<Arc<crate::activity::ActivityMetrics>>,
pub metrics: Arc<PollingMetrics>,
}
pub struct PollingMetrics {
pub polls_total: std::sync::atomic::AtomicU64,
pub poll_elapsed_ms: std::sync::atomic::AtomicU64,
pub condition_met: std::sync::atomic::AtomicU64,
pub poll_metric: std::sync::atomic::AtomicU64,
}
impl PollingMetrics {
fn new() -> Self {
Self {
polls_total: std::sync::atomic::AtomicU64::new(0),
poll_elapsed_ms: std::sync::atomic::AtomicU64::new(0),
condition_met: std::sync::atomic::AtomicU64::new(0),
poll_metric: std::sync::atomic::AtomicU64::new(0),
}
}
}
impl PollingDispenser {
#[allow(clippy::too_many_arguments)]
pub fn wrap(
inner: Arc<dyn OpDispenser>,
poll_interval_ms: u64,
timeout_ms: u64,
max_error_retries: u32,
metric_name: Option<String>,
min_rows: u64,
max_rows: u64,
json_path: Option<String>,
) -> (Arc<dyn OpDispenser>, Arc<PollingMetrics>) {
Self::wrap_with_status(
inner,
poll_interval_ms,
timeout_ms,
max_error_retries,
metric_name,
min_rows,
max_rows,
json_path,
None,
None,
None,
None,
None,
None,
None,
Vec::new(),
false,
crate::session_signals::StopView::default(),
)
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn wrap_with_status(
inner: Arc<dyn OpDispenser>,
poll_interval_ms: u64,
timeout_ms: u64,
max_error_retries: u32,
metric_name: Option<String>,
min_rows: u64,
max_rows: u64,
json_path: Option<String>,
each_memo: Option<String>,
memo_state: Option<Arc<arc_swap::ArcSwap<String>>>,
progress_template: Option<String>,
activity_metrics: Option<Arc<crate::activity::ActivityMetrics>>,
each_gutter: Option<(crate::wrappers::gutter::GutterKind, String)>,
gutter_state: Option<Arc<arc_swap::ArcSwapOption<crate::wrappers::gutter::GutterSpec>>>,
iteration_gauges: Option<
Arc<arc_swap::ArcSwapOption<Vec<crate::wrappers::metrics::MetricSlot>>>,
>,
on_done: Vec<(String, String)>,
until: bool,
stop: crate::session_signals::StopView,
) -> (Arc<dyn OpDispenser>, Arc<PollingMetrics>) {
let metrics = Arc::new(PollingMetrics::new());
let dispenser = Arc::new(Self {
inner,
poll_interval: std::time::Duration::from_millis(poll_interval_ms),
timeout: std::time::Duration::from_millis(timeout_ms),
max_error_retries,
metric_name,
min_rows,
max_rows,
json_path,
each_memo,
memo_state,
progress_template,
activity_metrics,
each_gutter,
gutter_state,
iteration_gauges,
on_done,
until,
stop,
metrics: metrics.clone(),
});
(dispenser, metrics)
}
fn publish_iteration_status(
&self,
wires: &dyn crate::wires::WireSource,
base_memo: &str,
row_count: u64,
polls: u64,
elapsed_secs: f64,
) {
if let Some(memo) = &self.memo_state {
let rendered = self.each_memo.as_deref().and_then(|t| {
match crate::wires::substitute_via_wires(t, wires) {
Ok(s) => Some(s),
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Debug,
"poll.memo: substitution failed for '{t}': {e}"
);
None
}
}
});
let text = rendered.unwrap_or_else(|| if self.until {
format!(
"{base_memo} — waiting on `until:` (not yet satisfied) · \
{row_count} row(s) observed · poll {polls}, {elapsed_secs:.0}s")
} else {
format!(
"{base_memo} — measured {row_count} row(s) [target {}..={}] · poll {polls}, {elapsed_secs:.0}s",
self.min_rows, self.max_rows)
});
memo.store(Arc::new(text));
}
if let Some(handle) = &self.iteration_gauges {
if let Some(slots) = handle.load_full() {
crate::wrappers::metrics::publish_gauges_lenient(&slots, wires);
}
}
if let (Some(state), Some((kind, template))) = (&self.gutter_state, &self.each_gutter) {
if let Some(spec) = crate::wrappers::gutter::render_spec(*kind, template, wires) {
state.store(Some(Arc::new(spec)));
}
}
if let (Some(metrics), Some(t)) =
(&self.activity_metrics, self.progress_template.as_deref())
{
match crate::wires::substitute_via_wires(t, wires) {
Ok(s) => match s.trim().parse::<f64>() {
Ok(f) => metrics.set_progress_override_with_elapsed(f, elapsed_secs),
Err(_) => {
crate::diag!(
crate::observer::LogLevel::Debug,
"poll.progress: '{t}' rendered to non-numeric '{s}'"
);
}
},
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Debug,
"poll.progress: substitution failed for '{t}': {e}"
);
}
}
}
}
}
struct ProgressOverrideClear<'m>(Option<&'m crate::activity::ActivityMetrics>);
impl Drop for ProgressOverrideClear<'_> {
fn drop(&mut self) {
if let Some(m) = self.0 {
m.set_progress_override(None);
}
}
}
impl WrappingDispenser for PollingDispenser {}
impl OpDispenser for PollingDispenser {
fn execute<'a>(
&'a self,
cycle: u64,
ctx: &'a crate::fixture::ExecCtx<'a>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
> {
Box::pin(async move {
let start = std::time::Instant::now();
let mut polls = 0u64;
let mut retryable_errors_consumed: u32 = 0;
let base_memo: String = self
.memo_state
.as_ref()
.map(|m| strip_measurement_suffix(m.load().as_str()).to_string())
.unwrap_or_default();
let _progress_clear = ProgressOverrideClear(self.activity_metrics.as_deref());
loop {
if self.stop.stopped() {
return Err(ExecutionError::Op(AdapterError {
error_name: "shutdown_cancelled".into(),
message: format!(
"session stop requested — poll abandoned after {} poll(s), {:.1}s",
polls,
start.elapsed().as_secs_f64()
),
retryable: false,
}));
}
let result = match self.inner.execute(cycle, ctx).await {
Ok(r) => r,
Err(e) => {
let retryable = match &e {
ExecutionError::Op(ad) => ad.retryable,
ExecutionError::Adapter(_) => false,
};
if !retryable {
return Err(e);
}
if retryable_errors_consumed >= self.max_error_retries {
return Err(e);
}
retryable_errors_consumed += 1;
let indent = crate::scene_tree::running_phase_indent();
let color = crate::observer::use_color();
let yellow = if color { "\x1b[33m" } else { "" };
let reset = if color { "\x1b[0m" } else { "" };
crate::diag!(
crate::observer::LogLevel::Warn,
"{indent}{yellow}poll retry {retryable_errors_consumed}/{}{reset} after retryable error: {}",
self.max_error_retries,
match &e {
ExecutionError::Op(ad) => &ad.message,
ExecutionError::Adapter(ad) => &ad.message,
},
);
tokio::time::sleep(self.poll_interval).await;
continue;
}
};
polls += 1;
let row_count = match (&result.body, self.json_path.as_deref()) {
(Some(body), Some(path)) => {
let json = body.to_json();
count_from_json_pointer(&json, path)
}
(Some(body), None) => body.element_count(),
(None, _) => 0,
};
let elapsed_ms = start.elapsed().as_millis() as u64;
ctx.wires
.write("poll_elapsed_ms", polydat::ast::Value::U64(elapsed_ms));
ctx.wires
.write("poll_count", polydat::ast::Value::U64(polls));
let is_done = if self.until {
match crate::wrappers::condition::holds(
ctx.wires,
crate::wrappers::condition::UNTIL_BINDING,
) {
Some(done) => done,
None => {
return Err(ExecutionError::Op(crate::adapter::AdapterError {
error_name: "poll_until_unresolved".into(),
message: format!(
"poll `until:` predicate did not resolve through \
ctx.wires (binding '{}') — scope synthesis should \
have lowered it into this node's kernel",
crate::wrappers::condition::UNTIL_BINDING
),
retryable: false,
}));
}
}
} else {
row_count >= self.min_rows && row_count <= self.max_rows
};
self.metrics
.poll_metric
.store(row_count, std::sync::atomic::Ordering::Relaxed);
if !is_done {
let indent = crate::scene_tree::running_phase_indent();
crate::diag!(
crate::observer::LogLevel::Debug,
"{indent}awaiting: {row_count} row(s), need [{}..={}] ({:.0}s elapsed)",
self.min_rows,
self.max_rows,
start.elapsed().as_secs_f64()
);
let _ = ctx
.wires
.write("poll_count", polydat::ast::Value::U64(polls));
let _ = ctx.wires.write(
"poll_elapsed_ms",
polydat::ast::Value::U64(start.elapsed().as_millis() as u64),
);
self.publish_iteration_status(
ctx.wires,
&base_memo,
row_count,
polls,
start.elapsed().as_secs_f64(),
);
}
if is_done {
let _ = ctx
.wires
.write("poll_count", polydat::ast::Value::U64(polls));
let _ = ctx.wires.write(
"poll_elapsed_ms",
polydat::ast::Value::U64(start.elapsed().as_millis() as u64),
);
for (name, expr) in &self.on_done {
match crate::wires::substitute_via_wires(expr, ctx.wires) {
Ok(rendered) => match rendered.trim().parse::<f64>() {
Ok(v) => {
let _ = ctx.wires.write(name, polydat::ast::Value::F64(v));
}
Err(_) => crate::diag!(
crate::observer::LogLevel::Debug,
"poll.on_done: '{name}: {expr}' rendered to \
non-numeric '{rendered}'"
),
},
Err(e) => crate::diag!(
crate::observer::LogLevel::Debug,
"poll.on_done: substitution failed for \
'{name}: {expr}': {e}"
),
}
}
self.publish_iteration_status(
ctx.wires,
&base_memo,
row_count,
polls,
start.elapsed().as_secs_f64(),
);
}
if is_done {
let elapsed = start.elapsed();
let elapsed_secs = elapsed.as_secs_f64();
self.metrics
.polls_total
.fetch_add(polls, std::sync::atomic::Ordering::Relaxed);
self.metrics.poll_elapsed_ms.store(
elapsed.as_millis() as u64,
std::sync::atomic::Ordering::Relaxed,
);
self.metrics
.condition_met
.store(1, std::sync::atomic::Ordering::Relaxed);
let indent = crate::scene_tree::running_phase_indent();
let color = crate::observer::use_color();
let dim = if color { "\x1b[2m" } else { "" };
let green = if color { "\x1b[32m" } else { "" };
let reset = if color { "\x1b[0m" } else { "" };
crate::observer::log(
crate::observer::LogLevel::Info,
&format!(
"{indent}{green}poll complete{reset}: {polls} polls {dim}in {elapsed_secs:.1}s{reset}"
),
);
let _ = ctx
.wires
.write("poll_count", polydat::ast::Value::U64(polls));
let _ = ctx.wires.write(
"poll_elapsed_ms",
polydat::ast::Value::U64(elapsed.as_millis() as u64),
);
if let Some(ref name) = self.metric_name {
let value = duration_value_for_metric_name(name, elapsed_secs);
let _ = ctx.wires.write(name, polydat::ast::Value::F64(value));
}
return Ok(OpResult {
body: None,
skipped: false,
});
}
if start.elapsed() > self.timeout {
return Err(ExecutionError::Op(AdapterError {
error_name: "poll_timeout".into(),
message: format!(
"polling timed out after {:.1}s ({} polls). Last result had rows.",
start.elapsed().as_secs_f64(),
polls
),
retryable: false,
}));
}
tokio::time::sleep(self.poll_interval).await;
}
})
}
fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
Some(self.inner.as_ref())
}
}
pub(crate) fn duration_value_for_metric_name(name: &str, elapsed_secs: f64) -> f64 {
if name.ends_with("_ns") {
elapsed_secs * 1e9
} else if name.ends_with("_us") {
elapsed_secs * 1e6
} else if name.ends_with("_ms") {
elapsed_secs * 1e3
} else if name.ends_with("_s") {
elapsed_secs
} else if name.ends_with("_m") {
elapsed_secs / 60.0
} else if name.ends_with("_h") {
elapsed_secs / 3600.0
} else {
elapsed_secs
}
}
pub(crate) fn count_from_json_pointer(json: &serde_json::Value, path: &str) -> u64 {
let Some(v) = json.pointer(path) else {
return 0;
};
match v {
serde_json::Value::Array(a) => a.len() as u64,
serde_json::Value::Number(n) => n
.as_u64()
.or_else(|| n.as_i64().map(|i| i.max(0) as u64))
.or_else(|| n.as_f64().map(|f| f.max(0.0) as u64))
.unwrap_or(0),
serde_json::Value::Object(_) => 1,
serde_json::Value::Bool(b) => {
if *b {
1
} else {
0
}
}
serde_json::Value::String(s) if s.is_empty() => 0,
serde_json::Value::String(_) => 1,
serde_json::Value::Null => 0,
}
}
pub(crate) fn parse_on_done(
cfg: Option<&serde_json::Map<String, serde_json::Value>>,
) -> Vec<(String, String)> {
let Some(map) = cfg
.and_then(|m| m.get("on_done"))
.and_then(|v| v.as_object())
else {
return Vec::new();
};
let mut out: Vec<(String, String)> = map
.iter()
.map(|(k, v)| {
let text = match v {
serde_json::Value::String(s) => s.clone(),
other => other.to_string(),
};
(k.clone(), text)
})
.collect();
out.sort_by(|a, b| a.0.cmp(&b.0));
out
}
fn strip_measurement_suffix(memo: &str) -> &str {
match memo.find(" — measured ") {
Some(i) if memo[i..].contains(" row(s) [target ") => &memo[..i],
_ => memo,
}
}
#[cfg(test)]
mod on_done_config_tests {
use super::parse_on_done;
fn cfg(json: &str) -> serde_json::Map<String, serde_json::Value> {
match serde_json::from_str(json).expect("test json") {
serde_json::Value::Object(m) => m,
_ => panic!("expected an object"),
}
}
#[test]
fn absent_or_empty_yields_nothing() {
assert!(parse_on_done(None).is_empty());
assert!(parse_on_done(Some(&cfg(r#"{"mode":"await_empty"}"#))).is_empty());
assert!(parse_on_done(Some(&cfg(r#"{"on_done":{}}"#))).is_empty());
}
#[test]
fn numeric_and_string_values_both_read() {
let got = parse_on_done(Some(&cfg(
r#"{"on_done":{"completion_ratio":1.0,"progress":"total"}}"#,
)));
assert_eq!(
got,
vec![
("completion_ratio".to_string(), "1.0".to_string()),
("progress".to_string(), "total".to_string()),
]
);
}
#[test]
fn entries_are_ordered_by_key() {
let got = parse_on_done(Some(&cfg(r#"{"on_done":{"z":1,"a":2,"m":3}}"#)));
let keys: Vec<&str> = got.iter().map(|(k, _)| k.as_str()).collect();
assert_eq!(keys, vec!["a", "m", "z"]);
}
}
#[cfg(test)]
mod status_publish_tests {
use super::*;
use std::sync::Arc;
struct RatioWire(f64);
impl crate::wires::WireSource for RatioWire {
fn get(&self, name: &str) -> Option<polydat::ast::Value> {
(name == "completion_ratio").then(|| polydat::ast::Value::F64(self.0))
}
fn names(&self) -> Box<dyn Iterator<Item = String> + '_> {
Box::new(std::iter::once("completion_ratio".to_string()))
}
}
fn dispenser_with_status(
each_memo: Option<&str>,
progress: Option<&str>,
memo: &Arc<arc_swap::ArcSwap<String>>,
metrics: &Arc<crate::activity::ActivityMetrics>,
) -> PollingDispenser {
struct NoopInner;
impl OpDispenser for NoopInner {
fn execute<'a>(
&'a self,
_cycle: u64,
_ctx: &'a crate::fixture::ExecCtx<'a>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
> {
Box::pin(async move {
Ok(OpResult {
body: None,
skipped: false,
})
})
}
}
PollingDispenser {
inner: Arc::new(NoopInner),
poll_interval: std::time::Duration::from_millis(1),
timeout: std::time::Duration::from_millis(10),
max_error_retries: 0,
metric_name: None,
min_rows: 0,
max_rows: 0,
json_path: None,
each_memo: each_memo.map(String::from),
memo_state: Some(memo.clone()),
progress_template: progress.map(String::from),
activity_metrics: Some(metrics.clone()),
each_gutter: None,
gutter_state: None,
iteration_gauges: None,
on_done: Vec::new(),
until: false,
stop: crate::session_signals::StopView::default(),
metrics: Arc::new(PollingMetrics::new()),
}
}
struct MapWires(std::sync::Mutex<std::collections::HashMap<String, polydat::ast::Value>>);
impl MapWires {
fn new(seed: &[(&str, f64)]) -> Self {
Self(std::sync::Mutex::new(
seed.iter()
.map(|(k, v)| (k.to_string(), polydat::ast::Value::F64(*v)))
.collect(),
))
}
}
impl crate::wires::WireSource for MapWires {
fn get(&self, name: &str) -> Option<polydat::ast::Value> {
self.0.lock().unwrap().get(name).cloned()
}
fn names(&self) -> Box<dyn Iterator<Item = String> + '_> {
let v: Vec<String> = self.0.lock().unwrap().keys().cloned().collect();
Box::new(v.into_iter())
}
fn write(&self, name: &str, value: polydat::ast::Value) -> crate::wires::WriteOutcome {
self.0.lock().unwrap().insert(name.to_string(), value);
crate::wires::WriteOutcome::Stored
}
}
fn run_to_done(d: &PollingDispenser, wires: &dyn crate::wires::WireSource) {
let _guard = crate::session_signals::STOP_GLOBAL_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
crate::session_signals::clear_session_stop_for_test();
let fields = crate::adapter::ResolvedFields::new(Vec::new(), Vec::new());
let pulls = crate::fixture::ResolvedPulls::empty();
let ctx = crate::fixture::ExecCtx::with_wires(&fields, &pulls, wires);
let rt = tokio::runtime::Builder::new_current_thread()
.enable_time()
.build()
.expect("test runtime");
rt.block_on(d.execute(0, &ctx)).expect("poll completes");
}
#[test]
fn the_terminating_poll_publishes_its_measurement() {
let memo = Arc::new(arc_swap::ArcSwap::from_pointee(String::from("base")));
let metrics = test_metrics();
let mut d = dispenser_with_status(None, None, &memo, &metrics);
let (slot, gauge) =
crate::wrappers::metrics::test_gauge_slot("completion", "completion_ratio");
d.iteration_gauges = Some(Arc::new(arc_swap::ArcSwapOption::from_pointee(vec![slot])));
let wires = MapWires::new(&[("completion_ratio", 0.42)]);
run_to_done(&d, &wires);
assert_eq!(
gauge.get(),
0.42,
"the terminating poll's measurement must reach the gauge"
);
}
#[test]
fn on_done_proxies_the_completion_a_remote_view_cannot_show() {
let memo = Arc::new(arc_swap::ArcSwap::from_pointee(String::from("base")));
let metrics = test_metrics();
let mut d = dispenser_with_status(None, None, &memo, &metrics);
let (slot, gauge) =
crate::wrappers::metrics::test_gauge_slot("completion", "completion_ratio");
d.iteration_gauges = Some(Arc::new(arc_swap::ArcSwapOption::from_pointee(vec![slot])));
d.on_done = vec![("completion_ratio".into(), "1.0".into())];
let wires = MapWires::new(&[("completion_ratio", 0.42)]);
run_to_done(&d, &wires);
assert_eq!(
gauge.get(),
1.0,
"on_done must override the stale in-flight value"
);
}
#[test]
fn republishing_the_terminating_measurement_is_idempotent() {
let memo = Arc::new(arc_swap::ArcSwap::from_pointee(String::from("base")));
let metrics = test_metrics();
let mut d = dispenser_with_status(None, None, &memo, &metrics);
let (slot, gauge) =
crate::wrappers::metrics::test_gauge_slot("completion", "completion_ratio");
d.iteration_gauges = Some(Arc::new(arc_swap::ArcSwapOption::from_pointee(vec![slot])));
d.on_done = vec![("completion_ratio".into(), "1.0".into())];
let wires = MapWires::new(&[("completion_ratio", 0.42)]);
run_to_done(&d, &wires);
let first = gauge.get();
let memo_after_first = memo.load_full();
run_to_done(&d, &wires);
let second = gauge.get();
assert_eq!(
first, second,
"a repeated terminal publish must not shift the value"
);
assert_eq!(
*memo_after_first,
*memo.load_full(),
"a repeated terminal publish must not accumulate into the memo"
);
}
fn test_metrics() -> Arc<crate::activity::ActivityMetrics> {
Arc::new(crate::activity::ActivityMetrics::new(
&nmbrs_metrics::labels::Labels::default(),
))
}
#[test]
fn progress_template_publishes_override_and_memo_renders() {
let memo = Arc::new(arc_swap::ArcSwap::from_pointee(String::from("base")));
let metrics = test_metrics();
let d = dispenser_with_status(
Some("ratio {completion_ratio}"),
Some("{completion_ratio}"),
&memo,
&metrics,
);
let wires = RatioWire(0.42);
d.publish_iteration_status(&wires, "base", 1, 3, 15.0);
assert_eq!(
metrics.progress_override(),
Some(0.42),
"progress template must publish the derived override"
);
assert_eq!(memo.load().as_str(), "ratio 0.42");
}
#[test]
fn default_memo_suffix_without_template() {
let memo = Arc::new(arc_swap::ArcSwap::from_pointee(String::from("waiting")));
let metrics = test_metrics();
let d = dispenser_with_status(None, None, &memo, &metrics);
let wires = RatioWire(0.9);
d.publish_iteration_status(&wires, "waiting", 4, 7, 33.0);
assert!(
memo.load().contains("measured 4 row(s)"),
"default memo must carry the measurement: {}",
memo.load()
);
assert_eq!(
metrics.progress_override(),
None,
"no progress template -> no override"
);
}
#[test]
fn iteration_status_republishes_gutter_during_form() {
let memo = Arc::new(arc_swap::ArcSwap::from_pointee(String::new()));
let metrics = test_metrics();
let gutter_state: Arc<arc_swap::ArcSwapOption<crate::wrappers::gutter::GutterSpec>> =
Arc::new(arc_swap::ArcSwapOption::empty());
let mut d = dispenser_with_status(None, None, &memo, &metrics);
d.each_gutter = Some((
crate::wrappers::gutter::GutterKind::Text,
"ratio {completion_ratio}".into(),
));
d.gutter_state = Some(gutter_state.clone());
d.publish_iteration_status(&RatioWire(0.25), "waiting", 4, 1, 5.0);
assert_eq!(
gutter_state.load().as_deref(),
Some(&crate::wrappers::gutter::GutterSpec::Text(
"ratio 0.25".into()
)),
"first iteration publishes the rendered during form"
);
d.publish_iteration_status(&RatioWire(0.75), "waiting", 2, 2, 10.0);
assert_eq!(
gutter_state.load().as_deref(),
Some(&crate::wrappers::gutter::GutterSpec::Text(
"ratio 0.75".into()
)),
"each iteration refreshes the cell"
);
d.each_gutter = Some((
crate::wrappers::gutter::GutterKind::Spark,
"{no_such_wire}".into(),
));
d.publish_iteration_status(&RatioWire(0.9), "waiting", 1, 3, 15.0);
assert_eq!(
gutter_state.load().as_deref(),
Some(&crate::wrappers::gutter::GutterSpec::Text(
"ratio 0.75".into()
)),
"render failure must not clobber the cell"
);
}
}
#[cfg(test)]
mod poll_wire_tests {
#[test]
fn until_memo_does_not_quote_the_unused_row_window() {
let src =
std::fs::read_to_string(concat!(env!("CARGO_MANIFEST_DIR"), "/src/wrappers/poll.rs"))
.expect("read own source");
let until_arm = src
.find("waiting on `until:` (not yet satisfied)")
.expect("the until: memo arm must exist");
let window_arm = src
.find("measured {row_count} row(s) [target")
.expect("the row-window memo arm must exist");
assert!(
until_arm < window_arm,
"the until: arm must be the FIRST branch, so a declared condition \
never falls through to the row-window text"
);
assert!(
src.contains("if self.until {"),
"the memo must branch on whether an until: was declared"
);
}
#[test]
fn poll_progress_wires_are_published_before_the_predicate() {
let src =
std::fs::read_to_string(concat!(env!("CARGO_MANIFEST_DIR"), "/src/wrappers/poll.rs"))
.expect("read own source");
let write_at = src
.find(".write(\"poll_elapsed_ms\"")
.expect("poll_elapsed_ms must be published");
let predicate_at = src
.find("let is_done = if self.until {")
.expect("predicate evaluation site");
assert!(
write_at < predicate_at,
"the poll-progress wires must be written BEFORE the until: predicate \
is evaluated, or a self-bounding condition reads a stale 0"
);
assert!(
src.contains(".write(\"poll_count\""),
"poll_count must be published alongside poll_elapsed_ms"
);
}
}