Skip to main content

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}