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}