nmbrs_runtime/readout_context.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! [`ActivityReadoutContext`] — concrete `ReadoutContext`
5//! impl built from the activity-side data the
6//! ✓ DONE block already gathers.
7//!
8//! Push 1 surface only. Built up by `nmbrs-runtime::activity`
9//! at end-of-activity right before invoking the `phase_outcome`
10//! readout. Each later push grows this struct as new
11//! built-ins (and new `ReadoutContext` methods) arrive.
12//!
13//! Owned, not borrowed: every field is computed at the call
14//! site (counter snapshots, the rendered chip string, the
15//! depth-indent string) and parked here for the readout's
16//! duration. This keeps the readout call free of borrow
17//! plumbing through the activity's locks.
18
19use std::sync::Arc;
20use std::sync::atomic::Ordering;
21
22use crate::lifecycle::EventType;
23use crate::readouts::{LifecycleState, ReadoutContext};
24
25/// Snapshot of everything `phase_outcome` needs to render at
26/// `Lod::Labeled / ContentMode::Value`. Constructed at
27/// end-of-activity in `nmbrs-runtime::activity`; thrown
28/// away after the render returns.
29pub struct ActivityReadoutContext {
30 pub phase_name: String,
31 pub phase_seq: Option<(usize, usize)>,
32 pub phase_labels: String,
33 pub cycles_completed: u64,
34 pub cycles_total: u64,
35 pub ops_ok: u64,
36 /// SKIPPED ops (`skips_total`) — excluded from the ok% denominator
37 /// (a skip is neither a success nor a failure).
38 pub skips: u64,
39 pub errors: u64,
40 pub retries: u64,
41 pub concurrency: usize,
42 pub elapsed_secs: f64,
43 pub consumed: u64,
44 pub status_metric_chips: String,
45 pub depth_indent: String,
46 pub use_color: bool,
47 /// Snapshot of the activity's memo at context-build time.
48 /// Empty when no `memo:` wrapper is active on any op.
49 pub memo: String,
50 /// SRD-76 / SRD-82 Part 1 — the terminal two-axis outcome. The
51 /// executor sets this when it installs the
52 /// [`crate::phase_outcome::PhaseOutcome`] on the scene tree,
53 /// before firing the on_phase_end binder.
54 pub outcome: crate::phase_outcome::Outcome,
55 /// SRD-76 — chronologically ordered error list. Empty
56 /// for `Completed`/`Skipped`; non-empty for `Failed`.
57 /// Drives the failure-flavoured rendering of the
58 /// [`crate::readouts::builtins::phase_outcome`] readout.
59 pub outcome_errors: Vec<crate::phase_outcome::PhaseErrorDetail>,
60 /// SRD-76 — cursor-resume payload, when the phase
61 /// supports it. `None` for the common case.
62 pub outcome_resume_cursor: Option<crate::phase_outcome::ResumeCursor>,
63 /// True for a daemon (open-ended) phase: there is no "done" to
64 /// meter, so completion renders no percentage (SRD-92 — the
65 /// same rule the live meter slot follows). Without this the
66 /// trait-default `false` let `progress_fraction()` fall through
67 /// to the cycles basis and a stopped daemon printed a
68 /// meaningless `N%` (cycles over its wall-clock ceiling).
69 pub open_ended: bool,
70}
71
72impl ReadoutContext for ActivityReadoutContext {
73 fn subject_name(&self) -> &str {
74 &self.phase_name
75 }
76 fn subject_seq(&self) -> Option<(usize, usize)> {
77 self.phase_seq
78 }
79 fn subject_labels(&self) -> &str {
80 &self.phase_labels
81 }
82 fn open_ended(&self) -> bool {
83 self.open_ended
84 }
85 fn cycles_completed(&self) -> u64 {
86 self.cycles_completed
87 }
88 fn cycles_total(&self) -> u64 {
89 self.cycles_total
90 }
91 fn ops_ok(&self) -> u64 {
92 self.ops_ok
93 }
94 fn skips(&self) -> u64 {
95 self.skips
96 }
97 fn errors(&self) -> u64 {
98 self.errors
99 }
100 fn retries(&self) -> u64 {
101 self.retries
102 }
103 fn concurrency(&self) -> usize {
104 self.concurrency
105 }
106 fn elapsed_secs(&self) -> f64 {
107 self.elapsed_secs
108 }
109 fn consumed(&self) -> u64 {
110 self.consumed
111 }
112 fn status_metric_chips(&self) -> String {
113 self.status_metric_chips.clone()
114 }
115 fn depth_indent(&self) -> &str {
116 &self.depth_indent
117 }
118 fn use_color(&self) -> bool {
119 self.use_color
120 }
121 fn event(&self) -> EventType {
122 EventType::PhaseEnd
123 }
124 fn subject_state(&self) -> LifecycleState {
125 // Mirror the outcome status onto the lifecycle axis
126 // so existing consumers that branch on `subject_state`
127 // see Failed when the phase failed (today they'd see
128 // Completed because the binder fired before the
129 // executor recorded the failure). SRD-76 unifies
130 // the two surfaces.
131 match self.outcome.validity {
132 crate::phase_outcome::Validity::Succeeded => LifecycleState::Completed,
133 crate::phase_outcome::Validity::Failed => LifecycleState::Failed(
134 self.outcome_errors
135 .first()
136 .map(|e| e.message.clone())
137 .unwrap_or_else(|| "phase failed".into()),
138 ),
139 }
140 }
141 fn phase_memo(&self) -> &str {
142 &self.memo
143 }
144 fn outcome(&self) -> crate::phase_outcome::Outcome {
145 self.outcome.clone()
146 }
147 fn outcome_errors(&self) -> &[crate::phase_outcome::PhaseErrorDetail] {
148 &self.outcome_errors
149 }
150 fn outcome_resume_cursor(&self) -> Option<&crate::phase_outcome::ResumeCursor> {
151 self.outcome_resume_cursor.as_ref()
152 }
153}
154
155/// Per-event context for lifecycle fires (Push 9a):
156/// `on_session_start` / `on_session_end`,
157/// `on_phase_start`, `on_each_start` / `on_each_end`,
158/// `on_scope_start` / `on_scope_end`.
159///
160/// Carries just the fields a structural readout
161/// (`scope_header`, `session_banner`, `each_close`, …)
162/// needs — subject name, root-first labels, depth indent,
163/// colour flag, plus the firing event so a wildcard-bound
164/// readout can branch.
165///
166/// Counter-shaped methods all return zero / empty since
167/// lifecycle readouts don't depend on per-cycle progress;
168/// the `Default` impl on the trait handles those.
169pub struct LifecycleContext {
170 pub event: crate::lifecycle::EventType,
171 pub subject_name: String,
172 pub subject_labels: String,
173 pub depth_indent: String,
174 pub use_color: bool,
175 /// SRD-106 — the session id the `stick_session` rung
176 /// re-attached to; empty everywhere except the SessionStart
177 /// fire of a stick-engaged run. Read by `session_notice`.
178 pub stick_reattached: String,
179}
180
181impl ReadoutContext for LifecycleContext {
182 fn subject_name(&self) -> &str {
183 &self.subject_name
184 }
185 fn subject_seq(&self) -> Option<(usize, usize)> {
186 None
187 }
188 fn subject_labels(&self) -> &str {
189 &self.subject_labels
190 }
191 fn cycles_completed(&self) -> u64 {
192 0
193 }
194 fn cycles_total(&self) -> u64 {
195 0
196 }
197 fn ops_ok(&self) -> u64 {
198 0
199 }
200 fn errors(&self) -> u64 {
201 0
202 }
203 fn retries(&self) -> u64 {
204 0
205 }
206 fn concurrency(&self) -> usize {
207 0
208 }
209 fn elapsed_secs(&self) -> f64 {
210 0.0
211 }
212 fn consumed(&self) -> u64 {
213 0
214 }
215 fn status_metric_chips(&self) -> String {
216 String::new()
217 }
218 fn depth_indent(&self) -> &str {
219 &self.depth_indent
220 }
221 fn use_color(&self) -> bool {
222 self.use_color
223 }
224 fn event(&self) -> crate::lifecycle::EventType {
225 self.event
226 }
227 fn stick_reattached_session(&self) -> &str {
228 &self.stick_reattached
229 }
230 fn subject_state(&self) -> LifecycleState {
231 // Lifecycle events fire at the boundary; the
232 // subject is in transition. `Running` is the safe
233 // default for `on_*_start` (the subject is now
234 // in flight); `_end` events technically transition
235 // to `Completed` but the readouts that fire here
236 // (scope_header, session_banner, etc.) don't
237 // branch on subject_state anyway, so a single
238 // default keeps things simple.
239 LifecycleState::Running
240 }
241}
242
243/// Per-tick context for the inline-status refresh thread
244/// (Push 2). Identifies as [`EventType::Update`]; carries a
245/// monotonic refresh tick for spinner cycling, the full
246/// activity name with leaf coord, and pre-formatted
247/// adapter / batch tails (the iteration over registered
248/// dispensers stays in the surface for now — Push 4
249/// migrates the trait to expose the typed iterator).
250pub struct InlineRefreshContext {
251 pub phase_name: String,
252 pub activity_name: String,
253 pub phase_seq: Option<(usize, usize)>,
254 pub phase_labels: String,
255 pub cycles_completed: u64,
256 pub cycles_total: u64,
257 pub ops_started: u64,
258 pub ops_finished: u64,
259 pub ops_ok: u64,
260 /// SKIPPED ops (`skips_total`) — `if:`-gated ops that ran no
261 /// adapter call. Excluded from the `ok%` denominator: a skip is
262 /// neither a success nor a failure.
263 pub skips: u64,
264 pub errors: u64,
265 pub retries: u64,
266 /// SRD-91 attempt-level tallies — successful and failed
267 /// RESOLVED attempts (both observed at attempt end).
268 /// `attempt_ok / (attempt_ok + attempt_failed)` is the
269 /// attempt success rate the status line surfaces beside the
270 /// result-level `ok%`; in-flight attempts are excluded so it
271 /// doesn't skew low the way the dispatch-time counter would.
272 pub attempt_ok: u64,
273 pub attempt_failed: u64,
274 pub concurrency: usize,
275 pub elapsed_secs: f64,
276 pub consumed: u64,
277 /// Cursor ordinals consumed / cursor extent for a data-driven
278 /// phase (polydat `global_consumed()` / `global_extent()`).
279 /// Both `0` for non-cursor phases (plain `cycles:`), where the
280 /// display keeps the op-denominated `cycles:` chip. `rows_total
281 /// > 0` selects the row-denominated `rows:{consumed}/{total}`
282 /// chip + rows/s rate.
283 pub rows_consumed: u64,
284 pub rows_total: u64,
285 pub status_metric_chips: String,
286 pub adapter_counters_text: String,
287 pub batch_info_text: String,
288 pub depth_indent: String,
289 pub refresh_tick: u64,
290 pub use_color: bool,
291 /// Snapshot of the activity's memo at tick build time.
292 /// Empty when no `memo:` wrapper has published anything.
293 pub memo: String,
294 /// Derived-progress override snapshot (see
295 /// [`crate::activity::ActivityMetrics::progress_override`]).
296 pub progress_override: Option<f64>,
297 /// Producer-elapsed seconds recorded with the override — the
298 /// measured-basis ETA's time denominator.
299 pub progress_override_elapsed: Option<f64>,
300 /// Open-ended subject (daemon / background poll): no progress
301 /// meter; latency summary renders in its place.
302 pub open_ended: bool,
303 /// Live service-time percentiles (nanos) from the activity's
304 /// timer, for the open-ended latency chip. 0 = no data yet.
305 pub lat_p50_nanos: u64,
306 pub lat_p99_nanos: u64,
307}
308
309impl ReadoutContext for InlineRefreshContext {
310 fn subject_name(&self) -> &str {
311 &self.phase_name
312 }
313 fn activity_name(&self) -> &str {
314 &self.activity_name
315 }
316 fn subject_seq(&self) -> Option<(usize, usize)> {
317 self.phase_seq
318 }
319 fn subject_labels(&self) -> &str {
320 &self.phase_labels
321 }
322 fn cycles_completed(&self) -> u64 {
323 self.cycles_completed
324 }
325 fn cycles_total(&self) -> u64 {
326 self.cycles_total
327 }
328 fn ops_started(&self) -> u64 {
329 self.ops_started
330 }
331 fn ops_finished(&self) -> u64 {
332 self.ops_finished
333 }
334 fn ops_ok(&self) -> u64 {
335 self.ops_ok
336 }
337 fn skips(&self) -> u64 {
338 self.skips
339 }
340 fn errors(&self) -> u64 {
341 self.errors
342 }
343 fn retries(&self) -> u64 {
344 self.retries
345 }
346 fn attempt_ok(&self) -> u64 {
347 self.attempt_ok
348 }
349 fn attempt_failed(&self) -> u64 {
350 self.attempt_failed
351 }
352 fn concurrency(&self) -> usize {
353 self.concurrency
354 }
355 fn elapsed_secs(&self) -> f64 {
356 self.elapsed_secs
357 }
358 fn consumed(&self) -> u64 {
359 self.consumed
360 }
361 fn rows_consumed(&self) -> u64 {
362 self.rows_consumed
363 }
364 fn rows_total(&self) -> u64 {
365 self.rows_total
366 }
367 fn status_metric_chips(&self) -> String {
368 self.status_metric_chips.clone()
369 }
370 fn adapter_counters_text(&self) -> String {
371 self.adapter_counters_text.clone()
372 }
373 fn batch_info_text(&self) -> String {
374 self.batch_info_text.clone()
375 }
376 fn depth_indent(&self) -> &str {
377 &self.depth_indent
378 }
379 fn use_color(&self) -> bool {
380 self.use_color
381 }
382 fn event(&self) -> EventType {
383 EventType::Update
384 }
385 fn refresh_tick(&self) -> u64 {
386 self.refresh_tick
387 }
388 fn phase_memo(&self) -> &str {
389 &self.memo
390 }
391 fn progress_override(&self) -> Option<f64> {
392 self.progress_override
393 }
394 fn open_ended(&self) -> bool {
395 self.open_ended
396 }
397 fn latency_p50_nanos(&self) -> u64 {
398 self.lat_p50_nanos
399 }
400 fn latency_p99_nanos(&self) -> u64 {
401 self.lat_p99_nanos
402 }
403 /// SRD-63 Push 9f: derive ETA from `cycles_total -
404 /// ops_finished` divided by the observed throughput
405 /// rate (`ops_finished / elapsed`). `None` when the
406 /// extent isn't known (sourceless phase running by
407 /// time / open-ended) or no progress has been made
408 /// yet (rate would divide-by-zero).
409 ///
410 /// A derived-progress override with a recorded producer
411 /// elapsed wins: ETA = `elapsed × (1−f)/f` — the measured
412 /// basis for a single long op (a poll-driven drain) whose
413 /// cycle accounting stands still. Guarded to `f` in
414 /// `(0, 1)`: at 0 nothing is measurable yet, at 1 the
415 /// after-state takes over momentarily.
416 fn eta_secs(&self) -> Option<f64> {
417 // Open-ended subjects have no completion, hence no ETA.
418 if self.open_ended {
419 return None;
420 }
421 if let (Some(f), Some(e)) = (self.progress_override, self.progress_override_elapsed)
422 && f > 0.0
423 && f < 1.0
424 && e > 0.0
425 {
426 return Some(e * (1.0 - f) / f);
427 }
428 // Cursor-driven phase: rows are the authoritative ordinal
429 // basis. Ops stride N rows each, so an op-denominated rate
430 // against the row-denominated extent would overstate the
431 // ETA by the stride factor (e.g. ~64/s ops vs 7.8K/s rows
432 // → 35h instead of 18m).
433 if self.rows_total > 0 && self.elapsed_secs > 0.0 {
434 if self.rows_consumed == 0 {
435 return None;
436 }
437 let rate = self.rows_consumed as f64 / self.elapsed_secs;
438 let remaining = self.rows_total.saturating_sub(self.rows_consumed) as f64;
439 return Some(remaining / rate);
440 }
441 if self.cycles_total == 0 || self.elapsed_secs <= 0.0 {
442 return None;
443 }
444 let rate = self.ops_finished as f64 / self.elapsed_secs;
445 if rate <= 0.0 {
446 return None;
447 }
448 let remaining = self.cycles_total.saturating_sub(self.ops_finished) as f64;
449 Some(remaining / rate)
450 }
451}
452
453/// One-shot lifecycle fire helper. Builds a binder for
454/// `event` against `bindings`, runs every bound body
455/// against `ctx`, writes the rendered text via
456/// `crate::diag!` (so it lands in stderr / log file
457/// uniformly), and captures to the snapshot store via
458/// `subject_kind` / `subject_id`.
459///
460/// Best-effort: errors building the binder log a warning
461/// and the fire is skipped — a malformed `readouts:`
462/// binding never blocks the run. Bindings that resolve
463/// to zero bodies (the usual case for structural slots
464/// with no built-in default and no workload binding)
465/// produce no output.
466pub fn fire_lifecycle(
467 event: crate::lifecycle::EventType,
468 bindings: &nmbrs_workload::model::ReadoutsBindings,
469 default: Option<crate::readouts::BakedBody>,
470 ctx: &dyn crate::readouts::ReadoutContext,
471 snapshot_writer: Option<&crate::readouts::snapshot::SnapshotWriter>,
472) {
473 use crate::readouts::ReadoutBinder;
474
475 // Use the built-in default when supplied (currently
476 // only PhaseEnd/Update have defaults); otherwise the
477 // slot starts empty and falls through to whatever the
478 // workload bound. `build_event_binder` always seeds
479 // the default — pass an empty body when none exists
480 // so unbound slots stay quiet.
481 let seed = default.unwrap_or_default();
482 let mut binder = match crate::readouts::build_event_binder(bindings, event, seed) {
483 Ok(b) => b,
484 Err(e) => {
485 crate::diag!(
486 crate::observer::LogLevel::Warn,
487 "readouts: failed to bind {slot} — {e}",
488 slot = event.slot_name()
489 );
490 return;
491 }
492 };
493 let mut sink = crate::readouts::StringSink::with_capacity(128);
494 binder.fire(event, ctx, &mut sink);
495 let rendered = sink.take();
496 if rendered.trim().is_empty() {
497 return; // no bound body for this slot — quiet exit
498 }
499 // The firing lifecycle slot IS the tag's attachment axis —
500 // this readout render is definitionally attached to the
501 // boundary that fired it. Sinks derive their rules from the
502 // axes (the terminal sink keeps `PhaseStart`-attached renders
503 // out of scrollback — its managed phase-history region mirrors
504 // them; scope / iteration / session boundaries have no region
505 // counterpart and flow through like in-flight lines).
506 let tag = crate::observer::EventTag::at(event, crate::observer::EventCategory::General);
507 crate::observer::log_tagged(crate::observer::LogLevel::Info, tag, &rendered);
508
509 // Snapshot capture per Push 6. Subject identity comes
510 // straight from the context: `subject_kind` from the
511 // firing event (the sole source of truth for which
512 // table dimension this row belongs to), `subject_id`
513 // from `ctx.subject_id()`'s default `name@labels`
514 // shape (overridden for session-scope contexts that
515 // collapse to a literal `"session"`). Replay reads
516 // stable tuples (slot, subject_kind, subject_id, ...).
517 let subject_id = ctx.subject_id();
518 crate::readouts::snapshot::capture(
519 snapshot_writer,
520 event.slot_name(),
521 ctx.subject_exec_id(),
522 event.subject_kind().as_str(),
523 &subject_id,
524 "binder",
525 crate::readouts::snapshot::lod_str(crate::readouts::Lod::Labeled),
526 &rendered,
527 );
528}
529
530/// Build an [`InlineRefreshContext`] from the per-tick
531/// counter snapshots the inline-status thread takes. This
532/// preserves the byte-equivalence target by constructing
533/// the same intermediate values the prior `format!()`
534/// inlined (adapter counter chips, batch info, the
535/// scene-tree-walk for `seq` + depth indent) — they're
536/// each derived once per tick, then handed to the
537/// [`crate::readouts::builtins::phase_status::PhaseStatus`]
538/// readout for actual rendering.
539// reason: cohesive per-tick context builder — each argument is a distinct
540// counter snapshot/handle taken once per refresh tick; grouping them into a
541// struct would only relocate the same fields.
542/// Resolve a phase's `(seq, total)` pre-map coordinate + depth indent
543/// from the GLOBAL scene tree by NAME, matching the first Running node.
544///
545/// Retained for the legacy inline-status / phase-end callers that only
546/// have the activity name. The executor's on-task render-handle attach
547/// (SRD-100 P2) resolves these from the dispatch-time `SceneNodeId`
548/// instead — race-safe under concurrent same-name dispatch, where this
549/// first-Running-match could pick the wrong sibling.
550pub fn resolve_phase_coord_by_name(activity_name: &str) -> (Option<(usize, usize)>, String) {
551 let bare_name = activity_name
552 .split_once(" (")
553 .map(|(n, _)| n)
554 .unwrap_or(activity_name);
555 crate::scene_tree::current()
556 .and_then(|t| {
557 let node = t
558 .dfs_phases()
559 .find(|n| {
560 n.name == bare_name
561 && matches!(n.status, crate::scene_tree::PhaseStatus::Running)
562 })?
563 .clone();
564 let seq = node.seq?;
565 let depth = node.depth.saturating_sub(1);
566 Some((Some((seq, t.total_phases())), " ".repeat(depth)))
567 })
568 .unwrap_or((None, String::new()))
569}
570
571/// Resolve a phase's `(seq, total)` pre-map coordinate + depth indent
572/// from the GLOBAL scene tree by its dispatch-time [`SceneNodeId`](crate::scene_tree::SceneNodeId).
573///
574/// SRD-100 P2 — the race-safe replacement for [`resolve_phase_coord_by_name`]:
575/// keying on the node id (allocated at dispatch, P1c) addresses the exact
576/// node, so concurrent same-name siblings (sweep cells, comprehension
577/// iterations, daemon+foreground) each resolve their OWN coordinate. The
578/// executor calls this once at render-handle attach time.
579pub fn resolve_phase_coord_by_id(
580 scene_node_id: crate::scene_tree::SceneNodeId,
581) -> (Option<(usize, usize)>, String) {
582 crate::scene_tree::current()
583 .and_then(|t| {
584 let node = t.nodes.get(scene_node_id)?;
585 let seq = node.seq?;
586 let depth = node.depth.saturating_sub(1);
587 Some((Some((seq, t.total_phases())), " ".repeat(depth)))
588 })
589 .unwrap_or((None, String::new()))
590}
591
592/// Average batch size for the ` rows/batch:` chip — rows written per
593/// successful batch op.
594///
595/// Prefers `rows_inserted / batch_writes`, where `batch_writes` is the
596/// number of ops that actually wrote ≥1 row (published by the CQL
597/// batch dispensers alongside `rows_inserted`). This is the true
598/// per-op stride: retried, failed, and non-inserting ops never touch
599/// `batch_writes`, so the denominator can't drift the way
600/// `stanzas_total` (one inc per op *attempt*, regardless of
601/// success/failure/type) does.
602///
603/// Falls back to `rows_inserted / stanzas_total` when no `batch_writes`
604/// counter is present — non-CQL or older adapter paths that never
605/// learned to publish it. Returns `None` when no batched write was
606/// observed (average would be ≤1 row/op, i.e. not a batch).
607///
608/// Kept as a free function (not inline at the call site) so the three
609/// display paths that render this chip — the inline refresh here and
610/// the two TUI progress-thread snapshots in `executor.rs` — share one
611/// formula, and so the formula is unit-testable in isolation.
612/// Whether a dispenser status counter is INTERNAL — published only to feed a
613/// derived display metric, never rendered as its own `<name>/s` throughput
614/// chip. The convention is a **leading underscore**: `_batch_writes` backs the
615/// `rows/batch` average (see [`rows_per_batch`]) but must not clutter the
616/// operator-facing chip row as `_batch_writes/s`. Every surface that turns a
617/// dispenser counter into a chip filters on this one predicate, so the
618/// convention has a single point of truth. `find_counter`-style lookups pass
619/// the underscore name explicitly, so the counter stays available to the
620/// derived metric it exists for.
621pub fn is_internal_counter(name: &str) -> bool {
622 name.starts_with('_')
623}
624
625pub(crate) fn rows_per_batch(
626 rows_inserted: Option<u64>,
627 batch_writes: Option<u64>,
628 stanzas: u64,
629) -> Option<f64> {
630 let rows = rows_inserted?;
631 match batch_writes {
632 // A batch write count is present: divide by it directly. Show
633 // only when a real batch (>1 row/op) was observed.
634 Some(batches) if batches > 0 => (rows > batches).then(|| rows as f64 / batches as f64),
635 // Legacy / non-CQL fallback: attempt-count denominator.
636 _ => (stanzas > 0 && rows > stanzas).then(|| rows as f64 / stanzas as f64),
637 }
638}
639
640#[allow(clippy::too_many_arguments)]
641pub fn build_inline_refresh_context(
642 progress_metrics: &Arc<crate::activity::ActivityMetrics>,
643 activity_name: &str,
644 concurrency: usize,
645 total_extent: u64,
646 // Row-level cursor progress for a data-driven phase
647 // (`global_consumed()` / `global_extent()`); both `0` for
648 // non-cursor phases so the readout keeps the `cycles:` chip.
649 rows_consumed: u64,
650 rows_total: u64,
651 elapsed_secs: f64,
652 refresh_tick: u64,
653 status_metrics: &[String],
654 memo: &arc_swap::ArcSwap<String>,
655 // SRD-100 P2 — pre-map `(seq, total)` coordinate + depth indent,
656 // resolved by the caller. The producer no longer walks the global
657 // scene tree by name (a first-Running-match that raced under
658 // concurrent same-name dispatch); the executor's on-task attach
659 // resolves these from the dispatch-time `SceneNodeId` instead.
660 phase_seq: Option<(usize, usize)>,
661 depth_indent: String,
662 // Open-ended (daemon) subject: suppress progress metering, carry
663 // the latency chip instead.
664 open_ended: bool,
665) -> InlineRefreshContext {
666 // Counter snapshots — must match the prior inline-status
667 // formulas so byte equivalence holds.
668 let started = progress_metrics.ops_started.load(Ordering::Relaxed);
669 let finished = progress_metrics.ops_finished.load(Ordering::Relaxed);
670 let ops_completed = progress_metrics.cycles_completed();
671 // SRD-91: terminal-success count = `result_success.count()`;
672 // `errors_total` is RESULT-level (one inc per terminal failure),
673 // so it drives the `e:` count directly.
674 let successes = progress_metrics.result_success.count();
675 let errors = progress_metrics.errors_total.get();
676 let failed_ops = ops_completed
677 .saturating_sub(successes)
678 .saturating_sub(progress_metrics.skips_total.get());
679 let consumed = finished;
680 // SRD-91 attempt-level tallies, owned by the innermost
681 // `TriesDispenser` (or the error-handler wrapper for
682 // single-attempt ops). ALL attempt instruments count when an
683 // attempt RETURNS (2026-07-10 — attempt_total moved from
684 // dispatch to resolution, same discipline as the result
685 // instruments), so `attempt_total == attempt_success +
686 // attempt_failure` holds at every read and
687 // `attempt_ok / (attempt_ok + attempt_failed)` is the exact
688 // attempt success rate — dropping below the result-level
689 // `ok%` exactly when retries burn attempts to keep results
690 // green.
691 let attempt_ok = progress_metrics.attempt_success.count();
692 let attempt_failed = progress_metrics.attempt_failure.count();
693 // Retries = failed attempts that were NOT the terminal outcome.
694 // `errors_total` went RESULT-level with the TriesDispenser
695 // refactor, so the old `errors - failed_ops` derivation
696 // collapsed to ~0; the per-attempt failure count now lives in
697 // `attempt_failure`, and `attempt_failed - failed_ops` is the
698 // true retry count (each non-terminal failed attempt spawned a
699 // retry).
700 let retries = attempt_failed.saturating_sub(failed_ops);
701
702 // Adapter-status chips: ` <name>:<rate>/s` per registered
703 // dispenser counter. `collect_status_counters` aggregates
704 // every dispenser's typed counters into a flat
705 // `(name, total)` list — same data the inline thread used
706 // to read directly from `progress_metrics.dispensers`,
707 // exposed through the public accessor.
708 let mut adapter_counters_text = String::new();
709 let counters = progress_metrics.collect_status_counters();
710 for (name, total) in &counters {
711 // Internal counters (`_batch_writes`) feed derived metrics only —
712 // never their own chip. See `is_internal_counter`.
713 if is_internal_counter(name) {
714 continue;
715 }
716 let item_rate = if elapsed_secs > 0.0 {
717 *total as f64 / elapsed_secs
718 } else {
719 0.0
720 };
721 let rate_str = if item_rate >= 1_000_000.0 {
722 format!("{:.1}M", item_rate / 1_000_000.0)
723 } else if item_rate >= 1_000.0 {
724 format!("{:.1}K", item_rate / 1_000.0)
725 } else {
726 format!("{:.0}", item_rate)
727 };
728 adapter_counters_text.push_str(&format!(" {name}:{rate_str}/s"));
729 }
730
731 // Batch info: ` rows/batch:` = true average batch size (rows per
732 // successful batch op). Prefers `rows_inserted / batch_writes`
733 // when the CQL batch dispensers publish a `batch_writes` counter;
734 // otherwise falls back to the attempt-based `rows_inserted /
735 // stanzas_total`. See [`rows_per_batch`].
736 let stanzas = progress_metrics.stanzas_total.get();
737 let find_counter = |want: &str| counters.iter().find(|(n, _)| n == want).map(|(_, t)| *t);
738 let batch_info_text = rows_per_batch(
739 find_counter("rows_inserted"),
740 find_counter("_batch_writes"),
741 stanzas,
742 )
743 .map(|avg| format!(" rows/batch:{avg:.1}"))
744 .unwrap_or_default();
745
746 // Pre-rendered status-metric chip string.
747 let status_metric_chips = progress_metrics
748 .collect_status_values(status_metrics)
749 .concat();
750
751 // Activity name carries the leaf coord; the bare phase name is the
752 // readout's subject identity. (`phase_seq` / `depth_indent` now arrive
753 // as params — resolved race-safely by the caller.)
754 let bare_name = activity_name
755 .split_once(" (")
756 .map(|(n, _)| n)
757 .unwrap_or(activity_name);
758
759 let memo_snapshot: String = memo.load().as_str().to_string();
760 // Live latency percentiles for the open-ended chip (and any future
761 // consumer): a peek at the service-time HDR — no reset, cheap at
762 // display cadence.
763 let (lat_p50, lat_p99) = {
764 let snap = progress_metrics.service_time.peek_snapshot();
765 let h = &snap.histogram;
766 if h.is_empty() {
767 (0, 0)
768 } else {
769 (h.value_at_quantile(0.50), h.value_at_quantile(0.99))
770 }
771 };
772 InlineRefreshContext {
773 phase_name: bare_name.to_string(),
774 activity_name: activity_name.to_string(),
775 phase_seq,
776 phase_labels: String::new(),
777 cycles_completed: ops_completed,
778 cycles_total: total_extent,
779 ops_started: started,
780 ops_finished: finished,
781 ops_ok: successes,
782 skips: progress_metrics.skips_total.get(),
783 errors,
784 retries,
785 attempt_ok,
786 attempt_failed,
787 concurrency,
788 elapsed_secs,
789 consumed,
790 rows_consumed,
791 rows_total,
792 status_metric_chips,
793 adapter_counters_text,
794 batch_info_text,
795 depth_indent,
796 refresh_tick,
797 use_color: crate::observer::use_color(),
798 memo: memo_snapshot,
799 progress_override: progress_metrics.progress_override(),
800 progress_override_elapsed: progress_metrics.progress_override_elapsed_secs(),
801 open_ended,
802 lat_p50_nanos: lat_p50,
803 lat_p99_nanos: lat_p99,
804 }
805}
806
807#[cfg(test)]
808mod tests {
809 #[test]
810 fn eta_prefers_measured_basis_when_override_present() {
811 // 25% done after 60s of measured work → 180s remain,
812 // regardless of the standing-still cycle accounting.
813 let ctx = super::InlineRefreshContext {
814 phase_name: String::new(),
815 activity_name: String::new(),
816 phase_seq: None,
817 phase_labels: String::new(),
818 cycles_completed: 3,
819 cycles_total: 4,
820 ops_started: 4,
821 ops_finished: 3,
822 ops_ok: 3,
823 skips: 0,
824 errors: 0,
825 retries: 0,
826 attempt_ok: 3,
827 attempt_failed: 0,
828 concurrency: 1,
829 elapsed_secs: 600.0,
830 consumed: 3,
831 rows_consumed: 0,
832 rows_total: 0,
833 status_metric_chips: String::new(),
834 adapter_counters_text: String::new(),
835 batch_info_text: String::new(),
836 depth_indent: String::new(),
837 refresh_tick: 0,
838 use_color: false,
839 memo: String::new(),
840 progress_override: Some(0.25),
841 progress_override_elapsed: Some(60.0),
842 open_ended: false,
843 lat_p50_nanos: 0,
844 lat_p99_nanos: 0,
845 };
846 use crate::readouts::ReadoutContext;
847 let eta = ctx.eta_secs().expect("measured ETA");
848 assert!((eta - 180.0).abs() < 1e-9, "eta={eta}");
849 // Without the elapsed companion, fall back to cycle basis.
850 let ctx2 = super::InlineRefreshContext {
851 progress_override_elapsed: None,
852 ..ctx
853 };
854 let eta2 = ctx2.eta_secs().expect("cycle ETA");
855 assert!((eta2 - 200.0).abs() < 1e-9, "eta2={eta2}");
856 }
857
858 #[test]
859 fn eta_uses_row_basis_for_cursor_phases() {
860 // Cursor phase, ops stride 100 rows each: 2M of 10M rows in
861 // 200s → 8M remain at 10K rows/s → 800s. The op basis
862 // (10K ops finished vs a 10M extent) would claim ~200,000s —
863 // the stride-factor overstatement this pins against.
864 let ctx = super::InlineRefreshContext {
865 phase_name: String::new(),
866 activity_name: String::new(),
867 phase_seq: None,
868 phase_labels: String::new(),
869 cycles_completed: 10_000,
870 cycles_total: 10_000_000,
871 ops_started: 10_000,
872 ops_finished: 10_000,
873 ops_ok: 10_000,
874 skips: 0,
875 errors: 0,
876 retries: 0,
877 attempt_ok: 10_000,
878 attempt_failed: 0,
879 concurrency: 1,
880 elapsed_secs: 200.0,
881 consumed: 10_000,
882 rows_consumed: 2_000_000,
883 rows_total: 10_000_000,
884 status_metric_chips: String::new(),
885 adapter_counters_text: String::new(),
886 batch_info_text: String::new(),
887 depth_indent: String::new(),
888 refresh_tick: 0,
889 use_color: false,
890 memo: String::new(),
891 progress_override: None,
892 progress_override_elapsed: None,
893 open_ended: false,
894 lat_p50_nanos: 0,
895 lat_p99_nanos: 0,
896 };
897 use crate::readouts::ReadoutContext;
898 let eta = ctx.eta_secs().expect("row-basis ETA");
899 assert!((eta - 800.0).abs() < 1e-6, "eta={eta}");
900 // No rows consumed yet → unknown, not a fabricated op-basis ETA.
901 let ctx2 = super::InlineRefreshContext {
902 rows_consumed: 0,
903 ..ctx
904 };
905 assert!(
906 ctx2.eta_secs().is_none(),
907 "zero-row cursor phase must have no ETA"
908 );
909 }
910
911 use super::{is_internal_counter, rows_per_batch};
912
913 /// The leading-underscore convention: `_batch_writes` is an internal
914 /// denominator (hidden from the chip row), while a plain counter like
915 /// `rows_inserted` is a visible throughput chip. The chip-render loops
916 /// filter on exactly this predicate.
917 #[test]
918 fn internal_counter_is_underscore_prefixed() {
919 assert!(is_internal_counter("_batch_writes"));
920 assert!(!is_internal_counter("rows_inserted"));
921 assert!(!is_internal_counter("queries"));
922 }
923
924 /// With a `batch_writes` counter present, `rows/batch` is the true
925 /// average batch size (`rows_inserted / batch_writes`), NOT the
926 /// attempt-based `rows_inserted / stanzas_total`. Here 1000 rows
927 /// across 5 successful batch ops ⇒ 200.0, even though 40 op
928 /// attempts (stanzas) were recorded — the stanzas formula would
929 /// have wrongly shown 25.0.
930 #[test]
931 fn prefers_batch_writes_over_stanzas() {
932 let avg = rows_per_batch(Some(1000), Some(5), 40);
933 assert_eq!(avg, Some(200.0));
934 }
935
936 /// Without a `batch_writes` counter (non-CQL / older paths), the
937 /// formula falls back to `rows_inserted / stanzas_total`.
938 #[test]
939 fn falls_back_to_stanzas_without_batch_writes() {
940 let avg = rows_per_batch(Some(1000), None, 40);
941 assert_eq!(avg, Some(25.0));
942 }
943
944 /// `batch_writes = 0` behaves like "not present" — the counter is
945 /// only published once it has ticked, but guard against a zero
946 /// denominator either way and use the fallback.
947 #[test]
948 fn zero_batch_writes_uses_fallback() {
949 let avg = rows_per_batch(Some(1000), Some(0), 40);
950 assert_eq!(avg, Some(25.0));
951 }
952
953 /// No batched write observed (avg would be ≤1 row/op) ⇒ no chip.
954 #[test]
955 fn no_batch_observed_is_none() {
956 // batch_writes path: rows == batches ⇒ average of 1, not a batch.
957 assert_eq!(rows_per_batch(Some(5), Some(5), 40), None);
958 // stanzas fallback: rows == stanzas ⇒ not a batch.
959 assert_eq!(rows_per_batch(Some(40), None, 40), None);
960 // no rows_inserted counter at all ⇒ nothing to show.
961 assert_eq!(rows_per_batch(None, Some(5), 40), None);
962 }
963}