Skip to main content

nmbrs_runtime/wrappers/
poll.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Polling / await wrapper. Re-executes the inner op until its
5//! row count (optionally projected through a JSON-Pointer path)
6//! falls into the configured `[min_rows, max_rows]` window, or
7//! the timeout fires. Used for waiting on backend state to
8//! settle: SAI index build, compactions, etc.
9
10use std::sync::Arc;
11
12use crate::adapter::WrappingDispenser;
13use crate::adapter::{AdapterError, ExecutionError, OpDispenser, OpResult};
14use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
15
16/// SRD-32a wrapper name.
17pub const NAME: WrapperName = WrapperName::new("poll");
18
19/// Trigger: `poll:` may be a bare string (mode only, defaults
20/// for everything else) or a map carrying the full config —
21/// either form turns the wrapper on.
22fn triggers(s: WrapperSubject) -> bool {
23    let Some(template) = s.op() else {
24        return false;
25    };
26    template
27        .params
28        .get("poll")
29        .map(|v| v.is_string() || v.is_object())
30        .unwrap_or(false)
31}
32
33fn describe_assignment(s: WrapperSubject) -> Option<String> {
34    let template = s.op()?;
35    let poll_val = template.params.get("poll")?;
36    let (mode, interval, timeout): (String, u64, u64) = match poll_val {
37        v if v.is_string() => (v.as_str().unwrap().to_string(), 1000, 300_000),
38        v if v.is_object() => {
39            let m = v.as_object().unwrap();
40            let mode = m
41                .get("mode")
42                .and_then(|x| x.as_str())
43                .unwrap_or("await_empty")
44                .to_string();
45            let interval = m
46                .get("interval_ms")
47                .and_then(crate::wrapper_registrations::json_to_u64)
48                .unwrap_or(1000);
49            let timeout = m
50                .get("timeout_ms")
51                .and_then(crate::wrapper_registrations::json_to_u64)
52                .unwrap_or(300_000);
53            (mode, interval, timeout)
54        }
55        _ => return None,
56    };
57    Some(format!(
58        "poll: every {}ms, timeout {}ms, on `{mode}`",
59        interval, timeout
60    ))
61}
62
63inventory::submit! {
64    WrapperRegistration {
65        name: NAME,
66        owned_fields: &[
67            // `poll:` is the single discriminant for the poll
68            // wrapper; every knob (interval_ms, timeout_ms,
69            // min_rows, max_rows, json_path, metric_name,
70            // max_error_retries) lives under it as a map. The
71            // flat `poll_*`-prefix surface was retired.
72            "poll",
73        ],
74        triggers,
75        requires_inner: &[super::traverse::NAME],
76        forbids_outer: &[],
77        mutually_exclusive_with: &[],
78        describe_assignment,
79        levels: &[crate::wrapper_registry::WrapperLevel::Op],
80    }
81}
82
83/// Wraps an inner dispenser and re-executes it until the result
84/// body is empty (zero rows). Used for awaiting conditions like
85/// SAI index compaction completing.
86///
87/// Configured via op params:
88/// - `poll_interval_ms`: delay between polls (default: 1000)
89/// - `timeout_ms`: maximum total wait (default: 300000 = 5 min)
90/// - `poll_condition`: when to stop: "empty" (default) = stop when 0 rows
91/// - `poll_max_error_retries`: how many retryable errors to swallow
92///   before propagating (default: 0 — strict: any inner error fails
93///   the poll immediately, per SRD-03 §"Status-Determination
94///   Invariant")
95///
96/// Per SRD-03 §"Status-Determination Invariant", this wrapper
97/// short-circuits on every non-positive case:
98///
99/// - **Positive case**: inner op returns `OpResult` with an empty
100///   body → poll succeeds, this dispenser returns success.
101/// - **Any other case**: inner op returns a non-retryable
102///   `ExecutionError`, OR a retryable error past the retry limit,
103///   OR the timeout fires while non-empty bodies are still coming
104///   back → this dispenser returns the error, the activity error
105///   router sees it, and (under default `errors:` policy) the
106///   phase + the run stop. Errors are never swallowed behind the
107///   poll.
108pub struct PollingDispenser {
109    inner: Arc<dyn OpDispenser>,
110    poll_interval: std::time::Duration,
111    timeout: std::time::Duration,
112    /// Cap on consecutive retryable inner-op errors before the
113    /// wrapper propagates upstream. `0` means strict: any
114    /// inner-op error fails the poll immediately.
115    max_error_retries: u32,
116    /// SRD-92 cooperative-stop view (same aspect as the `while:` and
117    /// `tries` wrappers). Checked at the top of every poll iteration
118    /// so a session/walk/daemon stop abandons the poll instead of
119    /// waiting out the cadence or timeout. Injected at wrap time.
120    stop: crate::session_signals::StopView,
121    /// Named metric for the poll elapsed time (e.g., "index_build_time").
122    metric_name: Option<String>,
123    /// Threshold for "done": the poll is considered satisfied
124    /// when the inner op's row-count is in `[min_rows, max_rows]`.
125    /// Default `max_rows=0, min_rows=0` reproduces the historical
126    /// `await_empty` semantics (zero rows = done). Use
127    /// `min_rows=1, max_rows=1` for "settled to a single row"
128    /// cases such as SAI's `sai_sstable_count == 1` after
129    /// memtable flush + compaction (without the lower bound the
130    /// poll would exit too early at count=0, before the
131    /// memtable has flushed).
132    min_rows: u64,
133    max_rows: u64,
134    /// Optional JSON-Pointer path (RFC 6901, e.g. `/value`) that
135    /// drills into the result body before computing the count
136    /// for the `[min_rows, max_rows]` check. Use this when the
137    /// op's body wraps the meaningful payload in an envelope —
138    /// notably Jolokia, whose every response is
139    /// `{request, value, status, timestamp}` and the actual
140    /// answer lives under `.value`. When the addressed sub-tree
141    /// is an array, count is its length; a number maps directly
142    /// to count; an object or null maps to 1 / 0. Default `None`
143    /// uses `body.element_count()` as-is.
144    json_path: Option<String>,
145    /// `poll.memo` — optional template re-rendered after EVERY poll
146    /// iteration (against the wires, so this iteration's captures are
147    /// visible) and published to the activity memo. Without it, a
148    /// long await shows only the memo wrapper's static `before:` text
149    /// while the measured values sit unrendered in wires — the memo
150    /// is where the operator is looking, so the measurement belongs
151    /// there. When absent but `memo_state` is wired, a generic
152    /// `<base memo> — measured N row(s) …` suffix is published
153    /// instead, so every poll surfaces its live measurement without
154    /// workload changes.
155    each_memo: Option<String>,
156    /// The activity's memo slot (same ArcSwap the memo wrapper
157    /// writes). `None` in tests / callers that don't surface memos.
158    memo_state: Option<Arc<arc_swap::ArcSwap<String>>>,
159    /// `poll.progress` — optional template whose rendered value is
160    /// parsed as an `f64` completion fraction in `[0.0, 1.0]` and
161    /// published to the activity's derived-progress override each
162    /// iteration (e.g. `"{completion_ratio}"`). Drives the phase
163    /// completion bar for phases whose one long op measures its own
164    /// progress; cleared when the poll completes so the cycle-based
165    /// accounting takes back over.
166    progress_template: Option<String>,
167    /// Metrics of the owning activity — target of the
168    /// derived-progress override. `None` in bare tests.
169    /// The op's `gutter:` DURING form, re-published per poll
170    /// iteration against that poll's wires — a single long drain op
171    /// (`await_empty` over an hours-long compaction) keeps a live
172    /// cell instead of one publish at op end. The final-form/`final:`
173    /// semantics are untouched (the activity epilogue owns those).
174    each_gutter: Option<(crate::wrappers::gutter::GutterKind, String)>,
175    /// The activity's shared gutter slot; `None` in bare tests.
176    gutter_state: Option<Arc<arc_swap::ArcSwapOption<crate::wrappers::gutter::GutterSpec>>>,
177    /// The op's compiled `metrics:` GAUGE slots, re-published per
178    /// poll iteration (see `metrics::publish_gauges_lenient`).
179    /// Filled by the metrics wrapper's cascade arm AFTER this
180    /// dispenser is built (metrics wraps outside poll), hence the
181    /// late-bound swap; empty/`None` when the op has no metrics.
182    iteration_gauges:
183        Option<Arc<arc_swap::ArcSwapOption<Vec<crate::wrappers::metrics::MetricSlot>>>>,
184    /// When set, the poll is DONE the first time this predicate reads truthy,
185    /// and the row-count window is not consulted.
186    ///
187    /// The predicate itself lives in the executing node's kernel under
188    /// [`crate::wrappers::condition::UNTIL_BINDING`], put there by scope
189    /// synthesis. This wrapper never learns whether that kernel belongs to an
190    /// op template or a phase — it reads one wire through `ctx.wires`, and
191    /// scoping decides which binding answers.
192    until: bool,
193    /// Values written to the wires on the TERMINATING poll only,
194    /// as `(wire_name, expression)` pairs from `poll.on_done:`.
195    ///
196    /// A poll that watches remote work usually cannot observe that
197    /// work in a completed state: `system_views.compactions` (and
198    /// every view like it) lists what is RUNNING, so the observation
199    /// that means "done" is an empty result — the finished tier's
200    /// attributes are gone with it. The terminal sample is therefore
201    /// real ("nothing is running") but says nothing about the thing
202    /// that finished, and a `max(completion_ratio)` over the series
203    /// keeps reporting the last in-flight fraction for work that has
204    /// been done for an hour.
205    ///
206    /// `on_done` proxies the measurement the remote view cannot give
207    /// us: at termination these expressions are written to the wires
208    /// before the final publish, so the gauges fed by them record the
209    /// completed state as if it had been observed.
210    on_done: Vec<(String, String)>,
211    activity_metrics: Option<Arc<crate::activity::ActivityMetrics>>,
212    /// Externally visible metrics for the polling operation.
213    pub metrics: Arc<PollingMetrics>,
214}
215
216/// Metrics surfaced by the polling wrapper.
217pub struct PollingMetrics {
218    /// Total polls executed across all invocations.
219    pub polls_total: std::sync::atomic::AtomicU64,
220    /// Total time spent polling (milliseconds).
221    pub poll_elapsed_ms: std::sync::atomic::AtomicU64,
222    /// Whether the condition has been met (0 = waiting, 1 = done).
223    pub condition_met: std::sync::atomic::AtomicU64,
224    /// The last observed value from the poll condition (e.g., number of
225    /// remaining tasks). This is the metric that determines completion.
226    pub poll_metric: std::sync::atomic::AtomicU64,
227}
228
229impl PollingMetrics {
230    fn new() -> Self {
231        Self {
232            polls_total: std::sync::atomic::AtomicU64::new(0),
233            poll_elapsed_ms: std::sync::atomic::AtomicU64::new(0),
234            condition_met: std::sync::atomic::AtomicU64::new(0),
235            poll_metric: std::sync::atomic::AtomicU64::new(0),
236        }
237    }
238}
239
240impl PollingDispenser {
241    /// Wrap an inner dispenser with polling behavior.
242    /// Returns the wrapped dispenser and a handle to the metrics.
243    ///
244    /// `metric_name`: if set, the elapsed poll time is captured as a named
245    /// gauge (in seconds) for the summary report.
246    /// `max_error_retries`: cap on consecutive retryable inner errors
247    /// (default 0 = strict).
248    // reason: cohesive wrapper constructor — each argument is a distinct poll
249    // policy knob (interval, timeout, retry cap, metric name, row bounds);
250    // bundling them into a struct would only relocate the same fields.
251    #[allow(clippy::too_many_arguments)]
252    pub fn wrap(
253        inner: Arc<dyn OpDispenser>,
254        poll_interval_ms: u64,
255        timeout_ms: u64,
256        max_error_retries: u32,
257        metric_name: Option<String>,
258        min_rows: u64,
259        max_rows: u64,
260        json_path: Option<String>,
261    ) -> (Arc<dyn OpDispenser>, Arc<PollingMetrics>) {
262        Self::wrap_with_status(
263            inner,
264            poll_interval_ms,
265            timeout_ms,
266            max_error_retries,
267            metric_name,
268            min_rows,
269            max_rows,
270            json_path,
271            None,
272            None,
273            None,
274            None,
275            None,
276            None,
277            None,
278            Vec::new(),
279            false,
280            crate::session_signals::StopView::default(),
281        )
282    }
283
284    /// As [`Self::wrap`], plus the live-status handles: the
285    /// per-iteration memo template + activity memo slot, and the
286    /// derived-progress template + activity metrics. See the field
287    /// docs (`each_memo`, `progress_template`) for semantics.
288    // pub(crate), not pub: the returned handles include `MetricSlot`, which is crate-private, and
289    // every caller (activity.rs wiring, `wrap` above) is in-crate. Declaring it `pub` leaked a
290    // private type through a public signature.
291    #[allow(clippy::too_many_arguments)]
292    pub(crate) fn wrap_with_status(
293        inner: Arc<dyn OpDispenser>,
294        poll_interval_ms: u64,
295        timeout_ms: u64,
296        max_error_retries: u32,
297        metric_name: Option<String>,
298        min_rows: u64,
299        max_rows: u64,
300        json_path: Option<String>,
301        each_memo: Option<String>,
302        memo_state: Option<Arc<arc_swap::ArcSwap<String>>>,
303        progress_template: Option<String>,
304        activity_metrics: Option<Arc<crate::activity::ActivityMetrics>>,
305        each_gutter: Option<(crate::wrappers::gutter::GutterKind, String)>,
306        gutter_state: Option<Arc<arc_swap::ArcSwapOption<crate::wrappers::gutter::GutterSpec>>>,
307        iteration_gauges: Option<
308            Arc<arc_swap::ArcSwapOption<Vec<crate::wrappers::metrics::MetricSlot>>>,
309        >,
310        on_done: Vec<(String, String)>,
311        until: bool,
312        stop: crate::session_signals::StopView,
313    ) -> (Arc<dyn OpDispenser>, Arc<PollingMetrics>) {
314        let metrics = Arc::new(PollingMetrics::new());
315        let dispenser = Arc::new(Self {
316            inner,
317            poll_interval: std::time::Duration::from_millis(poll_interval_ms),
318            timeout: std::time::Duration::from_millis(timeout_ms),
319            max_error_retries,
320            metric_name,
321            min_rows,
322            max_rows,
323            json_path,
324            each_memo,
325            memo_state,
326            progress_template,
327            activity_metrics,
328            each_gutter,
329            gutter_state,
330            iteration_gauges,
331            on_done,
332            until,
333            stop,
334            metrics: metrics.clone(),
335        });
336        (dispenser, metrics)
337    }
338
339    /// Per-iteration status publish: memo + derived progress.
340    ///
341    /// Runs after EVERY poll — including the terminating one, whose
342    /// observation is the one that detects completion and is no less
343    /// a measurement than the ones before it. Skipping it left every
344    /// series ending on the last incomplete sample, so a finished
345    /// unit of work read as 93% done forever.
346    ///
347    /// Idempotent by construction: every publish here is a set, not
348    /// an accumulate — `publish_gauges_lenient` handles gauges only,
349    /// and the memo / gutter / progress-override slots are
350    /// last-write-wins. Publishing the same iteration twice therefore
351    /// lands the same state; only the sample timestamp moves.
352    ///
353    /// Captures for this iteration are already on the wires.
354    /// Substitution failures degrade to a debug log — status must
355    /// never fail the poll.
356    fn publish_iteration_status(
357        &self,
358        wires: &dyn crate::wires::WireSource,
359        base_memo: &str,
360        row_count: u64,
361        polls: u64,
362        elapsed_secs: f64,
363    ) {
364        if let Some(memo) = &self.memo_state {
365            let rendered = self.each_memo.as_deref().and_then(|t| {
366                match crate::wires::substitute_via_wires(t, wires) {
367                    Ok(s) => Some(s),
368                    Err(e) => {
369                        crate::diag!(
370                            crate::observer::LogLevel::Debug,
371                            "poll.memo: substitution failed for '{t}': {e}"
372                        );
373                        None
374                    }
375                }
376            });
377            // Default (no template / render failure): keep the memo the
378            // operator already sees and append the measurement to it.
379            let text = rendered.unwrap_or_else(|| if self.until {
380                // A declared `until:` REPLACED the row-count window, so quoting
381                // that window here would describe a gate that is not in effect —
382                // and worse, "measured 0 row(s) [target 0..=0]" reads as ALREADY
383                // SATISFIED while the poll keeps waiting, which is the most
384                // confusing thing a status line can say. Report the condition
385                // that is actually holding it, and the observation feeding it.
386                format!(
387                    "{base_memo} — waiting on `until:` (not yet satisfied) · \
388                     {row_count} row(s) observed · poll {polls}, {elapsed_secs:.0}s")
389            } else {
390                format!(
391                    "{base_memo} — measured {row_count} row(s) [target {}..={}] · poll {polls}, {elapsed_secs:.0}s",
392                    self.min_rows, self.max_rows)
393            });
394            memo.store(Arc::new(text));
395        }
396        // Gauges first: the gutter/memo templates may read metricsql
397        // over the very series these samples feed.
398        if let Some(handle) = &self.iteration_gauges {
399            if let Some(slots) = handle.load_full() {
400                crate::wrappers::metrics::publish_gauges_lenient(&slots, wires);
401            }
402        }
403        if let (Some(state), Some((kind, template))) = (&self.gutter_state, &self.each_gutter) {
404            if let Some(spec) = crate::wrappers::gutter::render_spec(*kind, template, wires) {
405                state.store(Some(Arc::new(spec)));
406            }
407        }
408        if let (Some(metrics), Some(t)) =
409            (&self.activity_metrics, self.progress_template.as_deref())
410        {
411            match crate::wires::substitute_via_wires(t, wires) {
412                Ok(s) => match s.trim().parse::<f64>() {
413                    // Elapsed rides along so the display can derive the
414                    // measured-basis ETA (`elapsed × (1−f)/f`) — the
415                    // cycle-based ETA stands still for one long measured op.
416                    Ok(f) => metrics.set_progress_override_with_elapsed(f, elapsed_secs),
417                    Err(_) => {
418                        crate::diag!(
419                            crate::observer::LogLevel::Debug,
420                            "poll.progress: '{t}' rendered to non-numeric '{s}'"
421                        );
422                    }
423                },
424                Err(e) => {
425                    crate::diag!(
426                        crate::observer::LogLevel::Debug,
427                        "poll.progress: substitution failed for '{t}': {e}"
428                    );
429                }
430            }
431        }
432    }
433}
434
435/// Drop guard clearing the derived-progress override when the poll
436/// future ends — by completion, timeout, error, OR cancellation (a
437/// daemon drop mid-await). Without it a finished/cancelled poll's
438/// stale fraction would keep driving the phase bar.
439struct ProgressOverrideClear<'m>(Option<&'m crate::activity::ActivityMetrics>);
440
441impl Drop for ProgressOverrideClear<'_> {
442    fn drop(&mut self) {
443        if let Some(m) = self.0 {
444            m.set_progress_override(None);
445        }
446    }
447}
448
449impl WrappingDispenser for PollingDispenser {}
450
451impl OpDispenser for PollingDispenser {
452    fn execute<'a>(
453        &'a self,
454        cycle: u64,
455        ctx: &'a crate::fixture::ExecCtx<'a>,
456    ) -> std::pin::Pin<
457        Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
458    > {
459        Box::pin(async move {
460            let start = std::time::Instant::now();
461            let mut polls = 0u64;
462            let mut retryable_errors_consumed: u32 = 0;
463            // Memo text as of poll start — the memo wrapper is OUTER,
464            // so its `before:` template is already published; the
465            // default per-iteration publish appends the measurement
466            // to this base rather than compounding onto itself.
467            //
468            // Stripping a previously-appended suffix is what makes that
469            // hold ACROSS executions too: within one poll the base is
470            // captured once, but a second execution of the same op would
471            // otherwise read back its own decorated text as the base and
472            // grow the memo without bound.
473            let base_memo: String = self
474                .memo_state
475                .as_ref()
476                .map(|m| strip_measurement_suffix(m.load().as_str()).to_string())
477                .unwrap_or_default();
478            let _progress_clear = ProgressOverrideClear(self.activity_metrics.as_deref());
479
480            loop {
481                // Session shutdown: abandon the poll. The cooperative
482                // drain waits for in-flight WORK, not for a poll
483                // cadence that may have minutes of timeout left —
484                // without this check a Ctrl-C leaves the drain stuck
485                // behind the next interval/retry sleep.
486                if self.stop.stopped() {
487                    return Err(ExecutionError::Op(AdapterError {
488                        error_name: "shutdown_cancelled".into(),
489                        message: format!(
490                            "session stop requested — poll abandoned after {} poll(s), {:.1}s",
491                            polls,
492                            start.elapsed().as_secs_f64()
493                        ),
494                        retryable: false,
495                    }));
496                }
497                // SRD-03 §"Status-Determination Invariant":
498                // every inner-op outcome other than the
499                // specific positive case (empty body, signalling
500                // "no remaining build tasks") short-circuits
501                // upstream as an error. Retryable errors get a
502                // bounded retry budget (`max_error_retries`);
503                // non-retryable errors never get retried — they
504                // propagate on the first occurrence.
505                let result = match self.inner.execute(cycle, ctx).await {
506                    Ok(r) => r,
507                    Err(e) => {
508                        let retryable = match &e {
509                            ExecutionError::Op(ad) => ad.retryable,
510                            ExecutionError::Adapter(_) => false,
511                        };
512                        if !retryable {
513                            return Err(e);
514                        }
515                        if retryable_errors_consumed >= self.max_error_retries {
516                            return Err(e);
517                        }
518                        retryable_errors_consumed += 1;
519                        let indent = crate::scene_tree::running_phase_indent();
520                        let color = crate::observer::use_color();
521                        let yellow = if color { "\x1b[33m" } else { "" };
522                        let reset = if color { "\x1b[0m" } else { "" };
523                        crate::diag!(
524                            crate::observer::LogLevel::Warn,
525                            "{indent}{yellow}poll retry {retryable_errors_consumed}/{}{reset} after retryable error: {}",
526                            self.max_error_retries,
527                            match &e {
528                                ExecutionError::Op(ad) => &ad.message,
529                                ExecutionError::Adapter(ad) => &ad.message,
530                            },
531                        );
532                        // Backoff before retry — same as the
533                        // between-polls cadence so a flapping
534                        // backend doesn't burn the retry budget
535                        // in a tight loop.
536                        tokio::time::sleep(self.poll_interval).await;
537                        continue;
538                    }
539                };
540                polls += 1;
541
542                // Check condition: row count in [min_rows, max_rows] = done.
543                // Default `min_rows=0, max_rows=0` reproduces the legacy
544                // `await_empty` (exactly 0 = done) semantics.
545                //
546                // When `json_path` is set, the count comes from the
547                // addressed sub-tree of the body's JSON projection
548                // — array length, raw number, or 1 (object) / 0
549                // (null/missing). This is the Jolokia-poll path:
550                // `getCompactions` returns `{value: [...], status: 200, ...}`
551                // and we want the length of `.value`, not 1 (the
552                // envelope object).
553                let row_count = match (&result.body, self.json_path.as_deref()) {
554                    (Some(body), Some(path)) => {
555                        let json = body.to_json();
556                        count_from_json_pointer(&json, path)
557                    }
558                    (Some(body), None) => body.element_count(),
559                    (None, _) => 0,
560                };
561                // A declared `until:` REPLACES the row-count window: the
562                // workload stated the completion condition, so counting rows
563                // would be second-guessing it. An unresolved predicate is a
564                // wiring fault, not a false one — failing loudly beats
565                // spinning to the timeout with no explanation.
566                // Publish the poll's own progress as wires BEFORE evaluating the
567                // predicate, so an `until:` can bound its own patience. Without
568                // these a condition can only describe the observed world, never
569                // "and give up waiting for it" — which turns any never-satisfied
570                // predicate into a wait to `timeout_ms`. `poll_elapsed_ms` is
571                // the time spent polling THIS invocation; `poll_count` is the
572                // number of iterations completed.
573                //
574                // Written through the same `WireSource::write` path captures
575                // use, so the names resolve exactly like any other wire and an
576                // `extern poll_elapsed_ms: u64 = 0` declaration picks them up.
577                // A workload that never mentions them has no slot, the write
578                // is a no-op, and nothing changes.
579                let elapsed_ms = start.elapsed().as_millis() as u64;
580                ctx.wires
581                    .write("poll_elapsed_ms", polydat::ast::Value::U64(elapsed_ms));
582                ctx.wires
583                    .write("poll_count", polydat::ast::Value::U64(polls));
584                let is_done = if self.until {
585                    match crate::wrappers::condition::holds(
586                        ctx.wires,
587                        crate::wrappers::condition::UNTIL_BINDING,
588                    ) {
589                        Some(done) => done,
590                        None => {
591                            return Err(ExecutionError::Op(crate::adapter::AdapterError {
592                                error_name: "poll_until_unresolved".into(),
593                                message: format!(
594                                    "poll `until:` predicate did not resolve through \
595                                     ctx.wires (binding '{}') — scope synthesis should \
596                                     have lowered it into this node's kernel",
597                                    crate::wrappers::condition::UNTIL_BINDING
598                                ),
599                                retryable: false,
600                            }));
601                        }
602                    }
603                } else {
604                    row_count >= self.min_rows && row_count <= self.max_rows
605                };
606
607                self.metrics
608                    .poll_metric
609                    .store(row_count, std::sync::atomic::Ordering::Relaxed);
610
611                if !is_done {
612                    // Per-poll progress goes to the durable
613                    // session log at Debug — direct `eprint!`
614                    // here would clobber the TUI's render
615                    // surface. The TUI surfaces poll progress
616                    // via the `poll_metric` gauge (live row
617                    // count) which is already updated above.
618                    let indent = crate::scene_tree::running_phase_indent();
619                    crate::diag!(
620                        crate::observer::LogLevel::Debug,
621                        "{indent}awaiting: {row_count} row(s), need [{}..={}] ({:.0}s elapsed)",
622                        self.min_rows,
623                        self.max_rows,
624                        start.elapsed().as_secs_f64()
625                    );
626                    // Live status: this iteration's captures are on the
627                    // wires; expose the poll's own counters alongside
628                    // them (slot-absent writes no-op) and publish the
629                    // measured values to the memo + the derived
630                    // phase-progress override.
631                    let _ = ctx
632                        .wires
633                        .write("poll_count", polydat::ast::Value::U64(polls));
634                    let _ = ctx.wires.write(
635                        "poll_elapsed_ms",
636                        polydat::ast::Value::U64(start.elapsed().as_millis() as u64),
637                    );
638                    self.publish_iteration_status(
639                        ctx.wires,
640                        &base_memo,
641                        row_count,
642                        polls,
643                        start.elapsed().as_secs_f64(),
644                    );
645                }
646                if is_done {
647                    // The terminating observation is a measurement too, and it
648                    // is the only one that carries the completion time. It used
649                    // to be dropped: the loop published on `!is_done` only, so
650                    // every series ended on the last INCOMPLETE sample and
651                    // finished work read as partially done forever.
652                    //
653                    // `on_done` first, then the publish: the wires it writes are
654                    // what the gauges read. This is the proxy for a completed
655                    // state a remote view cannot show (see the field docs) — the
656                    // terminal captures describe an empty result set, not the
657                    // work that just finished.
658                    let _ = ctx
659                        .wires
660                        .write("poll_count", polydat::ast::Value::U64(polls));
661                    let _ = ctx.wires.write(
662                        "poll_elapsed_ms",
663                        polydat::ast::Value::U64(start.elapsed().as_millis() as u64),
664                    );
665                    for (name, expr) in &self.on_done {
666                        match crate::wires::substitute_via_wires(expr, ctx.wires) {
667                            Ok(rendered) => match rendered.trim().parse::<f64>() {
668                                Ok(v) => {
669                                    let _ = ctx.wires.write(name, polydat::ast::Value::F64(v));
670                                }
671                                Err(_) => crate::diag!(
672                                    crate::observer::LogLevel::Debug,
673                                    "poll.on_done: '{name}: {expr}' rendered to \
674                                     non-numeric '{rendered}'"
675                                ),
676                            },
677                            Err(e) => crate::diag!(
678                                crate::observer::LogLevel::Debug,
679                                "poll.on_done: substitution failed for \
680                                 '{name}: {expr}': {e}"
681                            ),
682                        }
683                    }
684                    self.publish_iteration_status(
685                        ctx.wires,
686                        &base_memo,
687                        row_count,
688                        polls,
689                        start.elapsed().as_secs_f64(),
690                    );
691                }
692
693                if is_done {
694                    let elapsed = start.elapsed();
695                    let elapsed_secs = elapsed.as_secs_f64();
696                    self.metrics
697                        .polls_total
698                        .fetch_add(polls, std::sync::atomic::Ordering::Relaxed);
699                    self.metrics.poll_elapsed_ms.store(
700                        elapsed.as_millis() as u64,
701                        std::sync::atomic::Ordering::Relaxed,
702                    );
703                    self.metrics
704                        .condition_met
705                        .store(1, std::sync::atomic::Ordering::Relaxed);
706                    let indent = crate::scene_tree::running_phase_indent();
707                    let color = crate::observer::use_color();
708                    let dim = if color { "\x1b[2m" } else { "" };
709                    let green = if color { "\x1b[32m" } else { "" };
710                    let reset = if color { "\x1b[0m" } else { "" };
711                    crate::observer::log(
712                        crate::observer::LogLevel::Info,
713                        &format!(
714                            "{indent}{green}poll complete{reset}: {polls} polls {dim}in {elapsed_secs:.1}s{reset}"
715                        ),
716                    );
717                    // Captures land on the per-fiber kernel directly
718                    // via ctx.wires.write — wrappers above this layer
719                    // see the values through wires.get on the same
720                    // cycle. Slot-absent writes silently no-op
721                    // (closure-binding economy).
722                    let _ = ctx
723                        .wires
724                        .write("poll_count", polydat::ast::Value::U64(polls));
725                    let _ = ctx.wires.write(
726                        "poll_elapsed_ms",
727                        polydat::ast::Value::U64(elapsed.as_millis() as u64),
728                    );
729                    // Emit named metric. The recorded value is the
730                    // elapsed wait duration; if `metric_name` carries
731                    // a recognized unit suffix (`_ns` / `_us` / `_ms`
732                    // / `_s` / `_m` / `_h`), the seconds are
733                    // converted so the metric reads in the unit its
734                    // name advertises. Names without a recognized
735                    // suffix fall through as seconds (legacy
736                    // behaviour, used by e.g. `index_build_time`).
737                    if let Some(ref name) = self.metric_name {
738                        let value = duration_value_for_metric_name(name, elapsed_secs);
739                        let _ = ctx.wires.write(name, polydat::ast::Value::F64(value));
740                    }
741                    return Ok(OpResult {
742                        body: None,
743                        skipped: false,
744                    });
745                }
746
747                // Check timeout
748                if start.elapsed() > self.timeout {
749                    return Err(ExecutionError::Op(AdapterError {
750                        error_name: "poll_timeout".into(),
751                        message: format!(
752                            "polling timed out after {:.1}s ({} polls). Last result had rows.",
753                            start.elapsed().as_secs_f64(),
754                            polls
755                        ),
756                        retryable: false,
757                    }));
758                }
759
760                // Wait before next poll
761                tokio::time::sleep(self.poll_interval).await;
762            }
763        })
764    }
765    fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
766        Some(self.inner.as_ref())
767    }
768}
769
770/// Map a named metric's elapsed-time value to the numeric value the
771/// metric name advertises. The metric name suffix selects the unit
772/// (`_ns`, `_us`, `_ms`, `_s`, `_m`, `_h`); the elapsed seconds are
773/// scaled accordingly so a metric called `index_build_time_ms`
774/// reads in milliseconds.
775///
776/// Names without a recognised suffix fall through as seconds
777/// (preserves the historical contract — e.g. `index_build_time`
778/// used to be emitted raw in seconds and still is).
779///
780/// Longest suffixes are tested first so `_ms` doesn't tail-bind
781/// to a more-permissive `_s` rule by accident.
782pub(crate) fn duration_value_for_metric_name(name: &str, elapsed_secs: f64) -> f64 {
783    if name.ends_with("_ns") {
784        elapsed_secs * 1e9
785    } else if name.ends_with("_us") {
786        elapsed_secs * 1e6
787    } else if name.ends_with("_ms") {
788        elapsed_secs * 1e3
789    } else if name.ends_with("_s") {
790        elapsed_secs
791    } else if name.ends_with("_m") {
792        elapsed_secs / 60.0
793    } else if name.ends_with("_h") {
794        elapsed_secs / 3600.0
795    } else {
796        elapsed_secs
797    }
798}
799
800/// Drill into a JSON tree via JSON-Pointer path (RFC 6901, e.g.
801/// `/value`, `/value/results/0`) and reduce the addressed
802/// sub-tree to a u64 count for the polling threshold:
803///
804/// - Array → `len()` (use case: "list of running jobs is empty").
805/// - Number → the integer value (use case: a numeric counter
806///   like `Compaction.PendingTasks.Value` reaches zero).
807/// - Object → 1 (the addressed payload exists; for "wait until
808///   *something* is present" patterns).
809/// - Null / missing path → 0 (treat as "nothing there").
810///
811/// An empty path string addresses the root, matching
812/// `serde_json::Value::pointer("")`'s contract.
813pub(crate) fn count_from_json_pointer(json: &serde_json::Value, path: &str) -> u64 {
814    let Some(v) = json.pointer(path) else {
815        return 0;
816    };
817    match v {
818        serde_json::Value::Array(a) => a.len() as u64,
819        serde_json::Value::Number(n) => n
820            .as_u64()
821            .or_else(|| n.as_i64().map(|i| i.max(0) as u64))
822            .or_else(|| n.as_f64().map(|f| f.max(0.0) as u64))
823            .unwrap_or(0),
824        serde_json::Value::Object(_) => 1,
825        serde_json::Value::Bool(b) => {
826            if *b {
827                1
828            } else {
829                0
830            }
831        }
832        serde_json::Value::String(s) if s.is_empty() => 0,
833        serde_json::Value::String(_) => 1,
834        serde_json::Value::Null => 0,
835    }
836}
837
838/// Read `poll.on_done:` — the wire values a terminating poll writes before
839/// its final publish. Values may be numbers or strings; strings are kept
840/// verbatim so they can be wire expressions (`"total"`) rather than only
841/// literals.
842///
843/// Ordered by key so the writes are deterministic across runs — a map
844/// iteration order that varied would make one `on_done` entry able to
845/// shadow another differently from run to run.
846pub(crate) fn parse_on_done(
847    cfg: Option<&serde_json::Map<String, serde_json::Value>>,
848) -> Vec<(String, String)> {
849    let Some(map) = cfg
850        .and_then(|m| m.get("on_done"))
851        .and_then(|v| v.as_object())
852    else {
853        return Vec::new();
854    };
855    let mut out: Vec<(String, String)> = map
856        .iter()
857        .map(|(k, v)| {
858            let text = match v {
859                serde_json::Value::String(s) => s.clone(),
860                other => other.to_string(),
861            };
862            (k.clone(), text)
863        })
864        .collect();
865    out.sort_by(|a, b| a.0.cmp(&b.0));
866    out
867}
868
869/// Drop the default per-iteration measurement suffix this wrapper appends,
870/// so re-deriving the base from a published memo is a fixed point.
871/// Matched on both halves of the default shape — a memo that merely
872/// contains an em-dash keeps it.
873fn strip_measurement_suffix(memo: &str) -> &str {
874    match memo.find(" — measured ") {
875        Some(i) if memo[i..].contains(" row(s) [target ") => &memo[..i],
876        _ => memo,
877    }
878}
879
880#[cfg(test)]
881mod on_done_config_tests {
882    use super::parse_on_done;
883
884    fn cfg(json: &str) -> serde_json::Map<String, serde_json::Value> {
885        match serde_json::from_str(json).expect("test json") {
886            serde_json::Value::Object(m) => m,
887            _ => panic!("expected an object"),
888        }
889    }
890
891    #[test]
892    fn absent_or_empty_yields_nothing() {
893        assert!(parse_on_done(None).is_empty());
894        assert!(parse_on_done(Some(&cfg(r#"{"mode":"await_empty"}"#))).is_empty());
895        assert!(parse_on_done(Some(&cfg(r#"{"on_done":{}}"#))).is_empty());
896    }
897
898    /// A YAML `completion_ratio: 1.0` arrives as a NUMBER, not a string — the
899    /// spelling an operator actually writes must not be silently dropped.
900    #[test]
901    fn numeric_and_string_values_both_read() {
902        let got = parse_on_done(Some(&cfg(
903            r#"{"on_done":{"completion_ratio":1.0,"progress":"total"}}"#,
904        )));
905        assert_eq!(
906            got,
907            vec![
908                ("completion_ratio".to_string(), "1.0".to_string()),
909                ("progress".to_string(), "total".to_string()),
910            ]
911        );
912    }
913
914    /// Deterministic order: two entries must be applied the same way on every
915    /// run, so a later write shadowing an earlier one is reproducible.
916    #[test]
917    fn entries_are_ordered_by_key() {
918        let got = parse_on_done(Some(&cfg(r#"{"on_done":{"z":1,"a":2,"m":3}}"#)));
919        let keys: Vec<&str> = got.iter().map(|(k, _)| k.as_str()).collect();
920        assert_eq!(keys, vec!["a", "m", "z"]);
921    }
922}
923
924#[cfg(test)]
925mod status_publish_tests {
926    use super::*;
927    use std::sync::Arc;
928
929    /// Minimal WireSource: one f64 wire named `completion_ratio`.
930    struct RatioWire(f64);
931    impl crate::wires::WireSource for RatioWire {
932        fn get(&self, name: &str) -> Option<polydat::ast::Value> {
933            (name == "completion_ratio").then(|| polydat::ast::Value::F64(self.0))
934        }
935        fn names(&self) -> Box<dyn Iterator<Item = String> + '_> {
936            Box::new(std::iter::once("completion_ratio".to_string()))
937        }
938    }
939
940    fn dispenser_with_status(
941        each_memo: Option<&str>,
942        progress: Option<&str>,
943        memo: &Arc<arc_swap::ArcSwap<String>>,
944        metrics: &Arc<crate::activity::ActivityMetrics>,
945    ) -> PollingDispenser {
946        struct NoopInner;
947        impl OpDispenser for NoopInner {
948            fn execute<'a>(
949                &'a self,
950                _cycle: u64,
951                _ctx: &'a crate::fixture::ExecCtx<'a>,
952            ) -> std::pin::Pin<
953                Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
954            > {
955                Box::pin(async move {
956                    Ok(OpResult {
957                        body: None,
958                        skipped: false,
959                    })
960                })
961            }
962        }
963        PollingDispenser {
964            inner: Arc::new(NoopInner),
965            poll_interval: std::time::Duration::from_millis(1),
966            timeout: std::time::Duration::from_millis(10),
967            max_error_retries: 0,
968            metric_name: None,
969            min_rows: 0,
970            max_rows: 0,
971            json_path: None,
972            each_memo: each_memo.map(String::from),
973            memo_state: Some(memo.clone()),
974            progress_template: progress.map(String::from),
975            activity_metrics: Some(metrics.clone()),
976            each_gutter: None,
977            gutter_state: None,
978            iteration_gauges: None,
979            on_done: Vec::new(),
980            until: false,
981            stop: crate::session_signals::StopView::default(),
982            metrics: Arc::new(PollingMetrics::new()),
983        }
984    }
985
986    /// A writable wire bag — the terminating-publish tests need
987    /// `write` to actually land, which `NullWireSource` refuses.
988    struct MapWires(std::sync::Mutex<std::collections::HashMap<String, polydat::ast::Value>>);
989    impl MapWires {
990        fn new(seed: &[(&str, f64)]) -> Self {
991            Self(std::sync::Mutex::new(
992                seed.iter()
993                    .map(|(k, v)| (k.to_string(), polydat::ast::Value::F64(*v)))
994                    .collect(),
995            ))
996        }
997    }
998    impl crate::wires::WireSource for MapWires {
999        fn get(&self, name: &str) -> Option<polydat::ast::Value> {
1000            self.0.lock().unwrap().get(name).cloned()
1001        }
1002        fn names(&self) -> Box<dyn Iterator<Item = String> + '_> {
1003            let v: Vec<String> = self.0.lock().unwrap().keys().cloned().collect();
1004            Box::new(v.into_iter())
1005        }
1006        fn write(&self, name: &str, value: polydat::ast::Value) -> crate::wires::WriteOutcome {
1007            self.0.lock().unwrap().insert(name.to_string(), value);
1008            crate::wires::WriteOutcome::Stored
1009        }
1010    }
1011
1012    /// Run a poll to completion against `wires`. `NoopInner` returns an
1013    /// empty body, so the FIRST poll is the terminating one — exactly the
1014    /// iteration whose measurement used to be dropped.
1015    fn run_to_done(d: &PollingDispenser, wires: &dyn crate::wires::WireSource) {
1016        // The poll abandons on a session stop, and the stop flag is
1017        // process-global: hold it clear against sibling tests that set it.
1018        let _guard = crate::session_signals::STOP_GLOBAL_TEST_LOCK
1019            .lock()
1020            .unwrap_or_else(|e| e.into_inner());
1021        crate::session_signals::clear_session_stop_for_test();
1022        let fields = crate::adapter::ResolvedFields::new(Vec::new(), Vec::new());
1023        let pulls = crate::fixture::ResolvedPulls::empty();
1024        let ctx = crate::fixture::ExecCtx::with_wires(&fields, &pulls, wires);
1025        let rt = tokio::runtime::Builder::new_current_thread()
1026            .enable_time()
1027            .build()
1028            .expect("test runtime");
1029        rt.block_on(d.execute(0, &ctx)).expect("poll completes");
1030    }
1031
1032    /// The measurement that DETECTS completion must reach metrics like every
1033    /// incomplete one before it. Without this the last sample published is the
1034    /// final not-yet-done poll, so a finished unit of work reads as partially
1035    /// done for as long as the series is kept.
1036    #[test]
1037    fn the_terminating_poll_publishes_its_measurement() {
1038        let memo = Arc::new(arc_swap::ArcSwap::from_pointee(String::from("base")));
1039        let metrics = test_metrics();
1040        let mut d = dispenser_with_status(None, None, &memo, &metrics);
1041        let (slot, gauge) =
1042            crate::wrappers::metrics::test_gauge_slot("completion", "completion_ratio");
1043        d.iteration_gauges = Some(Arc::new(arc_swap::ArcSwapOption::from_pointee(vec![slot])));
1044        let wires = MapWires::new(&[("completion_ratio", 0.42)]);
1045
1046        run_to_done(&d, &wires);
1047
1048        assert_eq!(
1049            gauge.get(),
1050            0.42,
1051            "the terminating poll's measurement must reach the gauge"
1052        );
1053    }
1054
1055    /// A remote view of in-flight work cannot show the finished item — it is
1056    /// gone from the view, which is *why* the poll ended. `on_done` proxies
1057    /// that unobservable completed state onto the final sample.
1058    #[test]
1059    fn on_done_proxies_the_completion_a_remote_view_cannot_show() {
1060        let memo = Arc::new(arc_swap::ArcSwap::from_pointee(String::from("base")));
1061        let metrics = test_metrics();
1062        let mut d = dispenser_with_status(None, None, &memo, &metrics);
1063        let (slot, gauge) =
1064            crate::wrappers::metrics::test_gauge_slot("completion", "completion_ratio");
1065        d.iteration_gauges = Some(Arc::new(arc_swap::ArcSwapOption::from_pointee(vec![slot])));
1066        d.on_done = vec![("completion_ratio".into(), "1.0".into())];
1067        // The last value the view ever showed: 42% done, then it vanished.
1068        let wires = MapWires::new(&[("completion_ratio", 0.42)]);
1069
1070        run_to_done(&d, &wires);
1071
1072        assert_eq!(
1073            gauge.get(),
1074            1.0,
1075            "on_done must override the stale in-flight value"
1076        );
1077    }
1078
1079    /// Idempotence: these publishes are sets, not accumulates, so the same
1080    /// terminating observation applied twice lands the same state. That is what
1081    /// makes it safe to publish the final measurement in addition to whatever
1082    /// the cadence already emitted for the same moment.
1083    #[test]
1084    fn republishing_the_terminating_measurement_is_idempotent() {
1085        let memo = Arc::new(arc_swap::ArcSwap::from_pointee(String::from("base")));
1086        let metrics = test_metrics();
1087        let mut d = dispenser_with_status(None, None, &memo, &metrics);
1088        let (slot, gauge) =
1089            crate::wrappers::metrics::test_gauge_slot("completion", "completion_ratio");
1090        d.iteration_gauges = Some(Arc::new(arc_swap::ArcSwapOption::from_pointee(vec![slot])));
1091        d.on_done = vec![("completion_ratio".into(), "1.0".into())];
1092        let wires = MapWires::new(&[("completion_ratio", 0.42)]);
1093
1094        run_to_done(&d, &wires);
1095        let first = gauge.get();
1096        let memo_after_first = memo.load_full();
1097        run_to_done(&d, &wires);
1098        let second = gauge.get();
1099
1100        assert_eq!(
1101            first, second,
1102            "a repeated terminal publish must not shift the value"
1103        );
1104        assert_eq!(
1105            *memo_after_first,
1106            *memo.load_full(),
1107            "a repeated terminal publish must not accumulate into the memo"
1108        );
1109    }
1110
1111    fn test_metrics() -> Arc<crate::activity::ActivityMetrics> {
1112        Arc::new(crate::activity::ActivityMetrics::new(
1113            &nmbrs_metrics::labels::Labels::default(),
1114        ))
1115    }
1116
1117    #[test]
1118    fn progress_template_publishes_override_and_memo_renders() {
1119        let memo = Arc::new(arc_swap::ArcSwap::from_pointee(String::from("base")));
1120        let metrics = test_metrics();
1121        let d = dispenser_with_status(
1122            Some("ratio {completion_ratio}"),
1123            Some("{completion_ratio}"),
1124            &memo,
1125            &metrics,
1126        );
1127        let wires = RatioWire(0.42);
1128        d.publish_iteration_status(&wires, "base", 1, 3, 15.0);
1129        assert_eq!(
1130            metrics.progress_override(),
1131            Some(0.42),
1132            "progress template must publish the derived override"
1133        );
1134        assert_eq!(memo.load().as_str(), "ratio 0.42");
1135    }
1136
1137    #[test]
1138    fn default_memo_suffix_without_template() {
1139        let memo = Arc::new(arc_swap::ArcSwap::from_pointee(String::from("waiting")));
1140        let metrics = test_metrics();
1141        let d = dispenser_with_status(None, None, &memo, &metrics);
1142        let wires = RatioWire(0.9);
1143        d.publish_iteration_status(&wires, "waiting", 4, 7, 33.0);
1144        assert!(
1145            memo.load().contains("measured 4 row(s)"),
1146            "default memo must carry the measurement: {}",
1147            memo.load()
1148        );
1149        assert_eq!(
1150            metrics.progress_override(),
1151            None,
1152            "no progress template -> no override"
1153        );
1154    }
1155
1156    #[test]
1157    fn iteration_status_republishes_gutter_during_form() {
1158        // SRD-92: an op with a `gutter:` DURING form inside a poll
1159        // drain refreshes its cell per poll iteration — the
1160        // GutterDispenser alone publishes only once, at the end of
1161        // the (potentially hours-long) drain op. Substitution runs
1162        // against the iteration's wires, so the cell tracks the
1163        // measured state exactly like the poll memo.
1164        let memo = Arc::new(arc_swap::ArcSwap::from_pointee(String::new()));
1165        let metrics = test_metrics();
1166        let gutter_state: Arc<arc_swap::ArcSwapOption<crate::wrappers::gutter::GutterSpec>> =
1167            Arc::new(arc_swap::ArcSwapOption::empty());
1168        let mut d = dispenser_with_status(None, None, &memo, &metrics);
1169        d.each_gutter = Some((
1170            crate::wrappers::gutter::GutterKind::Text,
1171            "ratio {completion_ratio}".into(),
1172        ));
1173        d.gutter_state = Some(gutter_state.clone());
1174
1175        d.publish_iteration_status(&RatioWire(0.25), "waiting", 4, 1, 5.0);
1176        assert_eq!(
1177            gutter_state.load().as_deref(),
1178            Some(&crate::wrappers::gutter::GutterSpec::Text(
1179                "ratio 0.25".into()
1180            )),
1181            "first iteration publishes the rendered during form"
1182        );
1183
1184        // Next iteration's wires supersede the cell.
1185        d.publish_iteration_status(&RatioWire(0.75), "waiting", 2, 2, 10.0);
1186        assert_eq!(
1187            gutter_state.load().as_deref(),
1188            Some(&crate::wrappers::gutter::GutterSpec::Text(
1189                "ratio 0.75".into()
1190            )),
1191            "each iteration refreshes the cell"
1192        );
1193
1194        // A failed substitution must leave the last good value.
1195        d.each_gutter = Some((
1196            crate::wrappers::gutter::GutterKind::Spark,
1197            "{no_such_wire}".into(),
1198        ));
1199        d.publish_iteration_status(&RatioWire(0.9), "waiting", 1, 3, 15.0);
1200        assert_eq!(
1201            gutter_state.load().as_deref(),
1202            Some(&crate::wrappers::gutter::GutterSpec::Text(
1203                "ratio 0.75".into()
1204            )),
1205            "render failure must not clobber the cell"
1206        );
1207    }
1208}
1209
1210#[cfg(test)]
1211mod poll_wire_tests {
1212    /// With a declared `until:`, the memo must NOT quote the row-count window.
1213    /// "measured 0 row(s) [target 0..=0]" reads as satisfied while the poll is
1214    /// still waiting — the single most misleading thing a status line can say,
1215    /// and exactly what an operator saw during a 1m28s drain.
1216    #[test]
1217    fn until_memo_does_not_quote_the_unused_row_window() {
1218        let src =
1219            std::fs::read_to_string(concat!(env!("CARGO_MANIFEST_DIR"), "/src/wrappers/poll.rs"))
1220                .expect("read own source");
1221        let until_arm = src
1222            .find("waiting on `until:` (not yet satisfied)")
1223            .expect("the until: memo arm must exist");
1224        let window_arm = src
1225            .find("measured {row_count} row(s) [target")
1226            .expect("the row-window memo arm must exist");
1227        assert!(
1228            until_arm < window_arm,
1229            "the until: arm must be the FIRST branch, so a declared condition \
1230             never falls through to the row-window text"
1231        );
1232        // And the two must be distinct branches of one `if self.until`.
1233        assert!(
1234            src.contains("if self.until {"),
1235            "the memo must branch on whether an until: was declared"
1236        );
1237    }
1238
1239    /// `poll_elapsed_ms` and `poll_count` must be WRITTEN before the predicate
1240    /// is evaluated, not after. A condition that bounds its own patience
1241    /// (`... || poll_elapsed_ms >= 60000`) is the only way to assert something
1242    /// that may never be observed without risking a wait to `timeout_ms` —
1243    /// which in the compaction drain is 48 hours. Publishing them a moment too
1244    /// late would leave the first evaluation reading 0 forever.
1245    #[test]
1246    fn poll_progress_wires_are_published_before_the_predicate() {
1247        let src =
1248            std::fs::read_to_string(concat!(env!("CARGO_MANIFEST_DIR"), "/src/wrappers/poll.rs"))
1249                .expect("read own source");
1250        // Needle tolerates rustfmt splitting the receiver from the call —
1251        // `ctx.wires\n    .write(...)` — which is exactly what broke the
1252        // original `ctx.wires.write("poll_elapsed_ms"` form.
1253        let write_at = src
1254            .find(".write(\"poll_elapsed_ms\"")
1255            .expect("poll_elapsed_ms must be published");
1256        let predicate_at = src
1257            .find("let is_done = if self.until {")
1258            .expect("predicate evaluation site");
1259        assert!(
1260            write_at < predicate_at,
1261            "the poll-progress wires must be written BEFORE the until: predicate \
1262             is evaluated, or a self-bounding condition reads a stale 0"
1263        );
1264        assert!(
1265            src.contains(".write(\"poll_count\""),
1266            "poll_count must be published alongside poll_elapsed_ms"
1267        );
1268    }
1269}