Skip to main content

nmbrs_runtime/
activity.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Activity: the unit of concurrent execution.
5
6use std::sync::Arc;
7use std::sync::atomic::Ordering;
8use std::time::{Duration, Instant};
9
10use nmbrs_metrics::instruments::counter::Counter;
11use nmbrs_metrics::instruments::histogram::Histogram;
12use nmbrs_metrics::instruments::outcome::{MetricDetail, MetricDetailConfig, OutcomeInstrument};
13use nmbrs_metrics::instruments::timer::Timer;
14use nmbrs_metrics::labels::Labels;
15use nmbrs_metrics::snapshot::MetricSet;
16use nmbrs_rate::RateLimiter;
17
18use crate::adapter::{DriverAdapter, OpDispenser};
19// CycleSource removed — all iteration goes through DataSourceFactory
20use crate::opseq::{OpSequence, SequencerType};
21use crate::validation;
22
23/// Configuration for an activity.
24pub struct ActivityConfig {
25    pub name: String,
26    pub cycles: u64,
27    /// Number of fibers (tokio tasks) executing stanzas concurrently.
28    pub concurrency: usize,
29    /// Target ops/sec for the single activity-level rate
30    /// limiter. `None` disables rate limiting. There is one
31    /// rate limiter per activity — no separate stanza-rate
32    /// mechanism.
33    pub rate: Option<f64>,
34    pub sequencer: SequencerType,
35    pub error_spec: String,
36    /// Error-rate circuit breaker: fail the phase early when the
37    /// fraction of errored ops exceeds this threshold (e.g. `0.1`
38    /// = 10%). Evaluated only after at least
39    /// [`ERROR_RATE_MIN_OPS`] ops so a small phase can't trip on a
40    /// single error. `None` disables it; a value `>= 1.0` also
41    /// effectively disables it (the rate never exceeds 1.0).
42    /// Resolved per phase as the workload's `error_rate_max:` field
43    /// over the session-wide default. Installed as the default SRD-83
44    /// stop condition (`error_rate > error_rate_max`).
45    pub error_rate_max: Option<f64>,
46    /// SRD-83 — the phase's declared stop-condition predicates (the
47    /// `when:` of each `stop_when:` entry). Compiled into scope-bound
48    /// `ScopedPredicate`s alongside the default error-rate condition and
49    /// evaluated per tick.
50    pub stop_when: Vec<crate::stop_conditions::StopConditionDecl>,
51    /// SRD-83 §throttle — the phase's adaptive backpressure governor
52    /// spec: keep the windowed attempt-failure fraction under a bound
53    /// by walking a dynamic control (`concurrency`/`rate`). Evaluated
54    /// on the same drain-loop tick as the stop conditions.
55    pub throttle: Option<nmbrs_workload::model::ThrottleSpec>,
56    /// Inherited total-attempts budget for this activity's ops (phase
57    /// `tries:` or the workload-root `tries` param). An op's own `tries:`
58    /// field overrides it. `None` = no budget in scope → ops WITHOUT their
59    /// own `tries:` run single-attempt with no tries wrapper (SRD-82 Part 3b
60    /// sigil). `0` = ops fail without executing; `1` = explicit
61    /// single-attempt.
62    pub tries: Option<u32>,
63    /// Retry-backoff overrides from the phase-level `tries:` map form
64    /// (`tries: {count, backoff: {ratio, min, max}}`). `None` = none
65    /// declared at the phase → the tries wrapper falls back to the op's
66    /// standalone `retry_backoff*` params, then the built-in defaults.
67    pub tries_backoff: Option<nmbrs_workload::model::BackoffSpec>,
68    /// Maximum number of ops within a stanza that execute concurrently.
69    pub stanza_concurrency: usize,
70    /// Source factory for data-driven phases. When present, fibers pull
71    /// from this source instead of the cycle counter. Each fiber creates
72    /// its own reader via `create_reader()`.
73    pub source_factory: Option<Arc<dyn polydat::iteration::source::DataSourceFactory>>,
74    /// Suppress the inline stderr progress line (TUI handles
75    /// display). Wrapped in `Arc<AtomicBool>` so the runner can
76    /// flip it at runtime — when the user dismisses the TUI
77    /// mid-run (`q` keypress), this flag drops to `false` and
78    /// the status thread resumes emission, making the
79    /// experience feel like tui=off was set from the start.
80    /// A bare `bool` would have baked the TUI-mode value in at
81    /// activity construction, so post-dismissal there'd be no
82    /// progress display at all.
83    pub suppress_status_line: Arc<std::sync::atomic::AtomicBool>,
84    /// Names of relevancy / live aggregate metrics to surface on
85    /// the inline progress line and the per-phase ✓ DONE summary.
86    /// Empty → no extra metrics are shown (status line carries
87    /// only the universal counters). Set per-phase via the YAML
88    /// `status_metrics: [name]` field; workload-level phases that
89    /// compute relevancy must opt in explicitly — nothing is
90    /// presumed to be present.
91    pub status_metrics: Vec<String>,
92    /// Full root-first coordinate label (e.g.
93    /// `(profile=label_00), (bucket=1, kind=READ)`) for this
94    /// phase's iteration. Used by the ✓ DONE summary line to
95    /// show the same identity the per-phase header would carry,
96    /// so the completed-status line stands alone — no separate
97    /// phase-starting row needed.
98    pub phase_labels: String,
99    /// Pre-map sequence number `[N/total]` for this phase. Same
100    /// numbering the TUI tree row and post-run summary use.
101    /// `None` ⇒ inline-CLI form / pre-map didn't produce a seq.
102    pub phase_seq: Option<(usize, usize)>,
103    /// Resolved `readouts:` slot bindings from the workload
104    /// (SRD-63 §5). Empty → all slots fall through to the
105    /// hard-coded built-in defaults (`phase_outcome` at
106    /// `on_phase_end`, `phase_status` at `on_update`).
107    pub readouts: nmbrs_workload::model::ReadoutsBindings,
108    /// CLI `--readout=<body>` override (SRD-63 §8).
109    /// Applies to the `on_update` slot only; replaces
110    /// (or with `+` prefix, appends to) whatever the
111    /// workload + default path resolved.
112    pub cli_readout_override: Option<String>,
113    /// Per-session SQLite writer. Used by Push 6's snapshot
114    /// store — every binder.fire captures its rendered
115    /// output via `upsert_readout_snapshot` so replay /
116    /// scrollback can reproduce the line later. `None`
117    /// means snapshot capture is skipped (no session db
118    /// — short test fixtures, in-memory sessions).
119    pub snapshot_writer:
120        Option<Arc<std::sync::Mutex<Option<nmbrs_metrics::reporters::sqlite::SqliteReporter>>>>,
121    /// Session-level dryrun mode (`silent` / `emit` / `json`),
122    /// or `None` for a normal run.
123    ///
124    /// `dryrun=cycle` means **full construction of an executable
125    /// cycle path** — real adapter, real cluster connection, real
126    /// prepared statements, real metadata — and then suppression
127    /// of only the outbound `execute()` at cycle time via the
128    /// outermost `DryRunWrapper`. The wrapper is triggered by an
129    /// injected `dryrun:` op-template parameter, and this field is
130    /// the signal that drives that injection. There is no
131    /// substitution of the adapter itself; the adapter lifecycle
132    /// runs end-to-end, so the typed lvalue contract the adapter
133    /// reifies at `map_op` time (CQL: prepare + metadata) is
134    /// available under dryrun exactly as it is under a real run.
135    pub dry_run_mode: Option<String>,
136    /// `dryrun=dispenser`: when true, `run_with_adapters`
137    /// returns cleanly after every op template's dispenser is
138    /// constructed (adapter `map_op` fires, the wrapper plan
139    /// resolves and wraps, the per-template pull plan seals).
140    /// No fiber pool is spawned, no cycles run. Lets the
141    /// operator verify the full construction pipeline ran
142    /// without paying any per-cycle cost.
143    ///
144    /// Set from `ExecCtx::diag.depth` at phase-attach time —
145    /// `< Cycle` flips this on, `>= Cycle` leaves it off.
146    pub stop_after_dispenser_init: bool,
147}
148
149impl Default for ActivityConfig {
150    fn default() -> Self {
151        Self {
152            name: "default".into(),
153            cycles: 1,
154            concurrency: 1,
155            rate: None,
156            sequencer: SequencerType::Bucket,
157            error_spec: ".*:warn,stop".into(),
158            error_rate_max: None,
159            stop_when: Vec::new(),
160            throttle: None,
161            // Default 0 — retries are opt-in via the workload `retries:` param
162            // (runner default is also 0). Matches the effective pre-wrapper
163            // behaviour (retries previously required a policy `retry`
164            // classification, off by default).
165            tries: None,
166            tries_backoff: None,
167            stanza_concurrency: 1,
168            source_factory: None,
169            suppress_status_line: Arc::new(std::sync::atomic::AtomicBool::new(false)),
170            status_metrics: Vec::new(),
171            phase_labels: String::new(),
172            phase_seq: None,
173            readouts: nmbrs_workload::model::ReadoutsBindings::default(),
174            cli_readout_override: None,
175            snapshot_writer: None,
176            dry_run_mode: None,
177            stop_after_dispenser_init: false,
178        }
179    }
180}
181
182/// Standard metrics for an activity. Shared via Arc so the metrics
183/// scheduler can capture snapshots while executor tasks record.
184///
185/// Fields are `Arc<Counter>` / `Arc<Timer>` / `Arc<Histogram>` so
186/// the same instrument is held both here (for per-cycle record
187/// access) and in the activity's `Component` instrument registry
188/// (for the cadence reporter's per-tick capture). Per-cycle code
189/// continues calling `metrics.cycles_total.inc()` etc. through
190/// `Arc`'s `Deref`.
191///
192/// Static instruments (the fields below) register on the component
193/// from [`ActivityMetrics::register_on`] called by
194/// [`Activity::attach_component`]. Dynamic per-error-type counters
195/// and adapter-specific metrics flow through the
196/// [`nmbrs_metrics::component::DynamicCapture`] hook implemented
197/// for [`ActivityMetricsDynamic`].
198///
199/// Late-bound, optional shared dispenser list. Wrapped in a `Mutex`
200/// because it is set once after dispenser creation (post-init) and
201/// read by the dynamic-capture hook thereafter.
202type SharedDispensers = std::sync::Mutex<Option<Arc<Vec<Arc<dyn crate::adapter::OpDispenser>>>>>;
203
204pub struct ActivityMetrics {
205    pub service_time: Arc<Timer>,
206    pub wait_time: Arc<Timer>,
207    pub response_time: Arc<Timer>,
208    /// Number of tries per op (1 = succeeded first try, 2+ = retried).
209    /// Distribution shape reveals incremental saturation.
210    pub tries_histogram: Arc<Histogram>,
211    /// Every op dispatched (incl. skips) — the rate driver. Distinct
212    /// from `result_total`, which excludes skips. SRD-91.
213    pub cycles_total: Arc<Counter>,
214    pub skips_total: Arc<Counter>,
215    // ── SRD-91 op-outcome taxonomy ────────────────────────────────
216    // Two layers that reconcile (the redundancy IS the validation):
217    //   • executor layer — `attempt_*` / `result_*`, counted in the
218    //     stanza hot loop;
219    //   • error-handler layer — `errors_total` + the per-type
220    //     breakdown, counted per failed attempt at error dispatch.
221    // Invariants:
222    //   attempt_total == attempt_success.count + attempt_failure.count
223    //   result_total  == result_success.count  + result_failure.count
224    //   cycles_total  == result_total + skips_total
225    //   errors_total  == Σ per-type == attempt_failure.count
226    //                    (when the policy counts every error)
227    /// Per-ATTEMPT total — one increment per `dispenser.execute`,
228    /// including retries.
229    pub attempt_total: Arc<Counter>,
230    /// Per-ATTEMPT outcomes (+ attempt latency when Timed). The count
231    /// is available in either detail mode — see [`OutcomeInstrument`].
232    pub attempt_success: OutcomeInstrument,
233    pub attempt_failure: OutcomeInstrument,
234    /// Per-OP terminal total — executed results only (success +
235    /// failure; excludes skips). Distinct from `cycles_total` by the
236    /// skip count.
237    pub result_total: Arc<Counter>,
238    /// Per-OP terminal outcomes (+ op latency when Timed).
239    /// `result_success` replaces the former `result_success_time`
240    /// timer and the unexported `successes_total` counter — its
241    /// `count()` IS the terminal-success count.
242    pub result_success: OutcomeInstrument,
243    pub result_failure: OutcomeInstrument,
244    /// Error-handler-layer tally: one increment per failed attempt at
245    /// error dispatch (per-attempt, so retries DO count here), keyed
246    /// by the handler-classified name for the per-type breakdown. The
247    /// per-op error rate uses `result_failure` instead, keeping it in
248    /// [0,1]. SRD-91.
249    pub errors_total: Arc<Counter>,
250    pub stanzas_total: Arc<Counter>,
251    /// Daemon ops that exited cleanly via stop-signal cancellation
252    /// at phase shutdown (the trigger-and-observe happy path).
253    /// Counts increment on `DaemonExit::Cancelled` only — natural
254    /// completions are tracked through `result_success` /
255    /// `result_failure` on the underlying op path. Visibility on
256    /// this counter lets the operator distinguish "phase exited
257    /// with N daemons cancelled" from "phase exited with no
258    /// daemons in flight" without re-reading session.log.
259    pub daemon_cancelled_total: Arc<Counter>,
260    /// Daemon ops whose shutdown failed: returned an error during
261    /// running or shutdown, panicked, or missed the grace window.
262    /// Each increment is paired with the activity's stop_flag
263    /// being set + a stop_reason being recorded.
264    pub daemon_errors_total: Arc<Counter>,
265    /// Number of ops dispatched to adapters (monotonic).
266    pub ops_started: std::sync::atomic::AtomicU64,
267    /// Number of ops returned from adapters (monotonic).
268    pub ops_finished: std::sync::atomic::AtomicU64,
269    pub result_elements: Arc<Counter>,
270    pub result_bytes: Arc<Counter>,
271    /// Derived phase-progress override, in parts-per-million of the
272    /// completion fraction (0..=1_000_000); `u64::MAX` = unset. When
273    /// set, status surfaces render THIS fraction for the phase's
274    /// completion bar / percentage instead of the cycles-based
275    /// `cycles_completed / total_extent` — load-bearing for phases
276    /// whose single long op measures its own progress (e.g. a
277    /// `poll:` await publishing `completion_ratio` from
278    /// `system_views.sstable_tasks`), where the cycle count pins the
279    /// bar at 0% for the whole wait. Stored as integer ppm so the
280    /// producer/consumer handoff stays a lock-free atomic.
281    pub progress_override_ppm: std::sync::atomic::AtomicU64,
282    /// Elapsed milliseconds of the producer that published
283    /// [`Self::progress_override_ppm`] (e.g. the poll's own elapsed),
284    /// `u64::MAX` = unset. Carried so displays can derive an ETA on the
285    /// measured basis — `elapsed × (1−f)/f` — instead of the cycle
286    /// accounting, which stands still for a single long measured op.
287    pub progress_override_elapsed_ms: std::sync::atomic::AtomicU64,
288    /// Per-error-type counters, keyed by error_name.
289    /// Created on demand when a new error type is first seen.
290    /// Captured via the [`DynamicCapture`] hook — the registry on
291    /// `Component` only holds instruments known at init.
292    error_type_counts: std::sync::Mutex<std::collections::HashMap<String, Arc<Counter>>>,
293    labels: Labels,
294    /// Dispensers for adapter-specific metrics capture. Set after dispenser creation.
295    dispensers: SharedDispensers,
296    /// Shared handles to the per-template validation metrics. Populated
297    /// after executor setup so the progress thread can read live
298    /// relevancy aggregates (recall-over-last-N, all-time mean) without
299    /// draining the precision accumulators.
300    validation_metrics:
301        std::sync::Mutex<Option<Arc<Vec<Arc<crate::validation::ValidationMetrics>>>>>,
302}
303
304/// Resolve the SRD-91 op-outcome detail config from the single
305/// `metrics_detail` param. The value is a comma-separated list: a bare
306/// token sets the global default (`counts` / `timers`), and a
307/// `family:mode` token overrides one instrument. Example:
308/// `metrics_detail=timers,attempt_success:counts,attempt_failure:counts`.
309/// Absent or unparseable tokens fall back to the default (timers).
310pub(crate) fn metric_detail_from_params(
311    params: &std::collections::HashMap<String, String>,
312) -> MetricDetailConfig {
313    let Some(spec) = params.get("metrics_detail") else {
314        return MetricDetailConfig::default();
315    };
316    let mut default = MetricDetail::default();
317    let mut overrides: Vec<(String, MetricDetail)> = Vec::new();
318    for tok in spec.split(',') {
319        let tok = tok.trim();
320        if tok.is_empty() {
321            continue;
322        }
323        if let Some((family, mode)) = tok.split_once(':') {
324            if let Some(d) = MetricDetail::parse(mode) {
325                overrides.push((family.trim().to_string(), d));
326            }
327        } else if let Some(d) = MetricDetail::parse(tok) {
328            default = d;
329        }
330    }
331    let mut cfg = MetricDetailConfig::new(default);
332    for (family, detail) in overrides {
333        cfg = cfg.with_override(family, detail);
334    }
335    cfg
336}
337
338impl ActivityMetrics {
339    pub fn new(labels: &Labels) -> Self {
340        Self::with_sigdigs(
341            labels,
342            nmbrs_metrics::instruments::histogram::DEFAULT_HDR_SIGDIGS,
343            &MetricDetailConfig::default(),
344        )
345    }
346
347    /// Construct activity metrics using an explicit HDR
348    /// significant-digits precision for every histogram and
349    /// timer below. The runner resolves `hdr.sigdigs` from the
350    /// session root via
351    /// [`nmbrs_metrics::instruments::histogram::resolve_hdr_sigdigs`]
352    /// once per activity and threads it here (SRD 40 §"HDR
353    /// significant digits — subtree-scoped setting").
354    pub fn with_sigdigs(labels: &Labels, sigdigs: u8, detail: &MetricDetailConfig) -> Self {
355        // Outcome instruments choose counter-vs-timer per family (global
356        // default + override), SRD-91. Default is Timers, preserving the
357        // historical always-on latency distributions.
358        let outcome = |name: &str| {
359            OutcomeInstrument::new(labels.with("name", name), sigdigs, detail.for_family(name))
360        };
361        Self {
362            service_time: Arc::new(Timer::with_sigdigs(
363                labels.with("name", "cycles_servicetime"),
364                sigdigs,
365            )),
366            wait_time: Arc::new(Timer::with_sigdigs(
367                labels.with("name", "cycles_waittime"),
368                sigdigs,
369            )),
370            response_time: Arc::new(Timer::with_sigdigs(
371                labels.with("name", "cycles_responsetime"),
372                sigdigs,
373            )),
374            tries_histogram: Arc::new(
375                nmbrs_metrics::instruments::histogram::Histogram::with_sigdigs(
376                    labels.with("name", "tries"),
377                    sigdigs,
378                ),
379            ),
380            cycles_total: Arc::new(Counter::new(labels.with("name", "cycles_total"))),
381            skips_total: Arc::new(Counter::new(labels.with("name", "skips_total"))),
382            attempt_total: Arc::new(Counter::new(labels.with("name", "attempt_total"))),
383            attempt_success: outcome("attempt_success"),
384            attempt_failure: outcome("attempt_failure"),
385            result_total: Arc::new(Counter::new(labels.with("name", "result_total"))),
386            result_success: outcome("result_success"),
387            result_failure: outcome("result_failure"),
388            errors_total: Arc::new(Counter::new(labels.with("name", "errors_total"))),
389            stanzas_total: Arc::new(Counter::new(labels.with("name", "stanzas_total"))),
390            daemon_cancelled_total: Arc::new(Counter::new(
391                labels.with("name", "daemon_cancelled_total"),
392            )),
393            daemon_errors_total: Arc::new(Counter::new(labels.with("name", "daemon_errors_total"))),
394            ops_started: std::sync::atomic::AtomicU64::new(0),
395            ops_finished: std::sync::atomic::AtomicU64::new(0),
396            result_elements: Arc::new(Counter::new(labels.with("name", "result_elements"))),
397            result_bytes: Arc::new(Counter::new(labels.with("name", "result_bytes"))),
398            progress_override_ppm: std::sync::atomic::AtomicU64::new(u64::MAX),
399            progress_override_elapsed_ms: std::sync::atomic::AtomicU64::new(u64::MAX),
400            error_type_counts: std::sync::Mutex::new(std::collections::HashMap::new()),
401            labels: labels.clone(),
402            dispensers: std::sync::Mutex::new(None),
403            validation_metrics: std::sync::Mutex::new(None),
404        }
405    }
406
407    /// Register every static instrument on `component` and install
408    /// a [`DynamicCapture`] hook for the dynamic surface (per-error-type
409    /// counters and adapter-specific metrics from registered dispensers).
410    ///
411    /// Called once from [`Activity::attach_component`]. After this
412    /// point:
413    /// - The cadence reporter's tree walk picks up every static
414    ///   instrument here through `component.capture_delta`.
415    /// - Per-cycle code continues recording through this struct's
416    ///   typed `Arc` fields — same `Arc` that the registry holds.
417    pub fn register_on(
418        self: &Arc<Self>,
419        component: &mut nmbrs_metrics::component::Component,
420    ) -> Result<(), String> {
421        use nmbrs_metrics::component::InstrumentRef;
422        // Order mirrors the historical capture_delta emission so
423        // metric_family ordering stays stable for downstream
424        // consumers. SRD-91 outcome instruments register via
425        // `instrument_ref()` (a counter or summary family per the
426        // resolved detail mode).
427        component.register_instrument(
428            "cycles_servicetime",
429            InstrumentRef::Timer(self.service_time.clone()),
430        )?;
431        component.register_instrument(
432            "cycles_waittime",
433            InstrumentRef::Timer(self.wait_time.clone()),
434        )?;
435        component.register_instrument(
436            "cycles_responsetime",
437            InstrumentRef::Timer(self.response_time.clone()),
438        )?;
439        component.register_instrument("result_success", self.result_success.instrument_ref())?;
440        component.register_instrument("result_failure", self.result_failure.instrument_ref())?;
441        component.register_instrument(
442            "result_total",
443            InstrumentRef::Counter(self.result_total.clone()),
444        )?;
445        component.register_instrument(
446            "cycles_total",
447            InstrumentRef::Counter(self.cycles_total.clone()),
448        )?;
449        component.register_instrument(
450            "skips_total",
451            InstrumentRef::Counter(self.skips_total.clone()),
452        )?;
453        component.register_instrument(
454            "errors_total",
455            InstrumentRef::Counter(self.errors_total.clone()),
456        )?;
457        component.register_instrument(
458            "attempt_total",
459            InstrumentRef::Counter(self.attempt_total.clone()),
460        )?;
461        component.register_instrument("attempt_success", self.attempt_success.instrument_ref())?;
462        component.register_instrument("attempt_failure", self.attempt_failure.instrument_ref())?;
463        component.register_instrument(
464            "stanzas_total",
465            InstrumentRef::Counter(self.stanzas_total.clone()),
466        )?;
467        component.register_instrument(
468            "daemon_cancelled_total",
469            InstrumentRef::Counter(self.daemon_cancelled_total.clone()),
470        )?;
471        component.register_instrument(
472            "daemon_errors_total",
473            InstrumentRef::Counter(self.daemon_errors_total.clone()),
474        )?;
475        component.register_instrument(
476            "result_elements",
477            InstrumentRef::Counter(self.result_elements.clone()),
478        )?;
479        component.register_instrument(
480            "result_bytes",
481            InstrumentRef::Counter(self.result_bytes.clone()),
482        )?;
483        component.register_instrument(
484            "tries",
485            InstrumentRef::Histogram(self.tries_histogram.clone()),
486        )?;
487
488        component.set_dynamic_capture(Arc::new(ActivityMetricsDynamic {
489            metrics: self.clone(),
490            prev_counters: std::sync::Mutex::new(std::collections::HashMap::new()),
491        }));
492        Ok(())
493    }
494
495    /// Return the number of cycles completed so far.
496    ///
497    /// Reads from the `cycles_total` counter atomically. Used by the
498    /// progress reporter thread to display live throughput.
499    pub fn cycles_completed(&self) -> u64 {
500        self.cycles_total.get()
501    }
502
503    /// Publish (or clear, with `None`) the derived phase-progress
504    /// override. `Some(f)` is clamped to `[0.0, 1.0]` and stored in
505    /// ppm; see the field doc on [`Self::progress_override_ppm`].
506    /// Clearing also clears the producer-elapsed companion.
507    pub fn set_progress_override(&self, fraction: Option<f64>) {
508        let ppm = match fraction {
509            Some(f) => (f.clamp(0.0, 1.0) * 1_000_000.0).round() as u64,
510            None => u64::MAX,
511        };
512        self.progress_override_ppm
513            .store(ppm, std::sync::atomic::Ordering::Relaxed);
514        if fraction.is_none() {
515            self.progress_override_elapsed_ms
516                .store(u64::MAX, std::sync::atomic::Ordering::Relaxed);
517        }
518    }
519
520    /// As [`Self::set_progress_override`], additionally recording the
521    /// producer's own elapsed seconds at publish time. Displays derive
522    /// the measured-basis ETA from the pair: `elapsed × (1−f)/f`.
523    pub fn set_progress_override_with_elapsed(&self, fraction: f64, elapsed_secs: f64) {
524        self.set_progress_override(Some(fraction));
525        let ms = (elapsed_secs.max(0.0) * 1000.0).round() as u64;
526        self.progress_override_elapsed_ms
527            .store(ms.min(u64::MAX - 1), std::sync::atomic::Ordering::Relaxed);
528    }
529
530    /// The derived phase-progress override as a fraction in
531    /// `[0.0, 1.0]`, or `None` when no producer has published one.
532    pub fn progress_override(&self) -> Option<f64> {
533        let ppm = self
534            .progress_override_ppm
535            .load(std::sync::atomic::Ordering::Relaxed);
536        (ppm != u64::MAX).then(|| (ppm.min(1_000_000)) as f64 / 1_000_000.0)
537    }
538
539    /// The producer-elapsed seconds recorded with the override, or
540    /// `None` when unset (no producer, or a producer that publishes
541    /// the fraction only).
542    pub fn progress_override_elapsed_secs(&self) -> Option<f64> {
543        let ms = self
544            .progress_override_elapsed_ms
545            .load(std::sync::atomic::Ordering::Relaxed);
546        (ms != u64::MAX).then(|| ms as f64 / 1000.0)
547    }
548
549    /// Increment counter for a specific error type. Creates the
550    /// counter on first occurrence of each error name. The new
551    /// counter is read by the [`DynamicCapture`] hook on every
552    /// capture tick — registration on `Component` is implicit
553    /// through the hook, not a per-name `register_instrument` call.
554    /// Top-N error types by count, rendered `name=count` comma-joined —
555    /// empty string when no typed errors were recorded. Gives failure
556    /// messages their WHAT (which error families drove the counters)
557    /// without a metrics query.
558    pub fn top_error_types(&self, n: usize) -> String {
559        let map = self
560            .error_type_counts
561            .lock()
562            .unwrap_or_else(|e| e.into_inner());
563        let mut v: Vec<(String, u64)> = map
564            .iter()
565            .map(|(k, c)| (k.clone(), c.get()))
566            .filter(|(_, c)| *c > 0)
567            .collect();
568        v.sort_by(|a, b| b.1.cmp(&a.1));
569        v.truncate(n);
570        v.iter()
571            .map(|(k, c)| format!("{k}={c}"))
572            .collect::<Vec<_>>()
573            .join(", ")
574    }
575
576    pub fn count_error_type(&self, error_name: &str) {
577        let mut map = self
578            .error_type_counts
579            .lock()
580            .unwrap_or_else(|e| e.into_inner());
581        let counter = map.entry(error_name.to_string()).or_insert_with(|| {
582            Arc::new(Counter::new(
583                self.labels.with("name", format!("errors.{error_name}")),
584            ))
585        });
586        counter.inc();
587    }
588
589    /// Capture an absolute snapshot (counters at their current value,
590    /// timer histograms drained as deltas).
591    ///
592    /// Used by the legacy per-activity capture thread. For the component
593    /// tree scheduler, use [`capture_delta`] instead.
594    pub fn capture(&self, interval: std::time::Duration) -> MetricSet {
595        use nmbrs_metrics::snapshot::split_name_label;
596        let service_snap = self.service_time.snapshot();
597        let wait_snap = self.wait_time.snapshot();
598        let response_snap = self.response_time.snapshot();
599        let tries_snap = self.tries_histogram.snapshot();
600        let now = Instant::now();
601        let mut snap = MetricSet::at(now, interval);
602
603        let (n, lbl) = split_name_label(self.service_time.labels());
604        snap.insert_histogram(n, lbl, service_snap.histogram, now);
605        let (n, lbl) = split_name_label(self.wait_time.labels());
606        snap.insert_histogram(n, lbl, wait_snap.histogram, now);
607        let (n, lbl) = split_name_label(self.response_time.labels());
608        snap.insert_histogram(n, lbl, response_snap.histogram, now);
609        // result_success: a histogram when Timed, a plain count when
610        // Counted (SRD-91 detail mode).
611        match &self.result_success {
612            OutcomeInstrument::Timed(t) => {
613                let (n, lbl) = split_name_label(t.labels());
614                snap.insert_histogram(n, lbl, t.snapshot().histogram, now);
615            }
616            OutcomeInstrument::Counted(c) => {
617                let (n, lbl) = split_name_label(c.labels());
618                snap.insert_counter(n, lbl, c.get(), now);
619            }
620        }
621
622        let (n, lbl) = split_name_label(self.cycles_total.labels());
623        snap.insert_counter(n, lbl, self.cycles_total.get(), now);
624        let (n, lbl) = split_name_label(self.skips_total.labels());
625        snap.insert_counter(n, lbl, self.skips_total.get(), now);
626        let (n, lbl) = split_name_label(self.errors_total.labels());
627        snap.insert_counter(n, lbl, self.errors_total.get(), now);
628        let (n, lbl) = split_name_label(self.stanzas_total.labels());
629        snap.insert_counter(n, lbl, self.stanzas_total.get(), now);
630        let (n, lbl) = split_name_label(self.daemon_cancelled_total.labels());
631        snap.insert_counter(n, lbl, self.daemon_cancelled_total.get(), now);
632        let (n, lbl) = split_name_label(self.daemon_errors_total.labels());
633        snap.insert_counter(n, lbl, self.daemon_errors_total.get(), now);
634        let (n, lbl) = split_name_label(self.result_elements.labels());
635        snap.insert_counter(n, lbl, self.result_elements.get(), now);
636        let (n, lbl) = split_name_label(self.result_bytes.labels());
637        snap.insert_counter(n, lbl, self.result_bytes.get(), now);
638        let (n, lbl) = split_name_label(self.tries_histogram.labels());
639        snap.insert_histogram(n, lbl, tries_snap, now);
640
641        let error_counts = self
642            .error_type_counts
643            .lock()
644            .unwrap_or_else(|e| e.into_inner());
645        for counter in error_counts.values() {
646            let (n, lbl) = split_name_label(counter.labels());
647            snap.insert_counter(n, lbl, counter.get(), now);
648        }
649
650        snap
651    }
652
653    /// Register dispensers for adapter-specific metrics capture.
654    pub fn set_dispensers(&self, dispensers: Arc<Vec<Arc<dyn crate::adapter::OpDispenser>>>) {
655        *self.dispensers.lock().unwrap_or_else(|e| e.into_inner()) = Some(dispensers);
656    }
657
658    /// Register the per-template validation metrics so the progress
659    /// thread can read live relevancy aggregates.
660    pub fn set_validation_metrics(&self, vms: Arc<Vec<Arc<crate::validation::ValidationMetrics>>>) {
661        *self
662            .validation_metrics
663            .lock()
664            .unwrap_or_else(|e| e.into_inner()) = Some(vms);
665    }
666
667    /// Snapshot live relevancy aggregates from every registered
668    /// validation-metrics instance (one per op template that declared
669    /// `relevancy:`). Non-destructive — safe to call every frame.
670    pub fn collect_relevancy_live(&self) -> Vec<crate::validation::RelevancyLive> {
671        let mut out = Vec::new();
672        if let Ok(guard) = self.validation_metrics.lock()
673            && let Some(ref vms) = *guard
674        {
675            for vm in vms.iter() {
676                out.extend(vm.live_snapshot());
677            }
678        }
679        out
680    }
681
682    /// Collect every status-line value whose name matches one of
683    /// `patterns`. Patterns are glob-style (`*` for any run of
684    /// characters, `?` for a single character; literal otherwise),
685    /// matched against the canonical names below. Returns formatted
686    /// ` name:value` strings ready to concatenate into the inline
687    /// progress / DONE summary line, in pattern declaration order
688    /// with duplicates suppressed.
689    ///
690    /// Supported metric families:
691    /// - **Relevancy aggregates** — one entry per registered
692    ///   `relevancy.functions:` (e.g. `recall`, `precision`,
693    ///   `f1`). The relevancy cutoff rides on the metric's
694    ///   `k` / `r` labels rather than the family name.
695    ///   Value: `total_mean × 100` as a percent.
696    /// - **Latency** — `latency_p50`, `latency_p99`, `latency_max`,
697    ///   `latency_mean`, sourced from `service_time` (the per-op
698    ///   timer, exclusive of wait time). Value: auto-scaled
699    ///   duration via [`nmbrs_metrics::reporters::summary::format_duration`].
700    pub fn collect_status_values(&self, patterns: &[String]) -> Vec<String> {
701        if patterns.is_empty() {
702            return Vec::new();
703        }
704        // Build the candidate list once. Order is stable so
705        // pattern ordering, not iteration order, drives the
706        // output sequence.
707        let mut candidates: Vec<(String, String)> = Vec::new();
708        for live in self.collect_relevancy_live() {
709            // Zero evaluations (e.g. every op `if:`-skipped) means
710            // there is no measurement — omit the chip rather than
711            // fabricating `recall:0.00%` from an empty aggregate.
712            if live.total_count == 0 {
713                continue;
714            }
715            candidates.push((live.name, format!("{:.2}%", live.total_mean * 100.0)));
716        }
717        let snap = self.service_time.peek_snapshot();
718        let h = &snap.histogram;
719        if !h.is_empty() {
720            let fmt = nmbrs_metrics::reporters::summary::format_duration;
721            candidates.push((
722                "latency_p50".to_string(),
723                fmt(h.value_at_quantile(0.50) as f64),
724            ));
725            candidates.push((
726                "latency_p99".to_string(),
727                fmt(h.value_at_quantile(0.99) as f64),
728            ));
729            candidates.push(("latency_max".to_string(), fmt(h.max() as f64)));
730            candidates.push(("latency_mean".to_string(), fmt(h.mean())));
731        }
732        // KEY-METRIC accent: `status_metrics:` selection is the
733        // workload author saying "this is the number I'm running the
734        // test for" — the chip gets its own bright palette slot
735        // (bold bright magenta; used by nothing else on the line) so
736        // it reads first among the dim bookkeeping counters.
737        let color = crate::observer::use_color();
738        let accent = if color { "\x1b[1;95m" } else { "" };
739        let reset = if color { "\x1b[0m" } else { "" };
740        let mut out: Vec<String> = Vec::new();
741        let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
742        for pat in patterns {
743            for (name, val) in &candidates {
744                if !seen.contains(name.as_str()) && glob_match(pat, name) {
745                    seen.insert(name.clone());
746                    let label = chip_display_label(name);
747                    out.push(format!(" {accent}{label}:{val}{reset}"));
748                }
749            }
750        }
751        out
752    }
753
754    /// The PRIMARY key metric as a raw numeric sample: the first
755    /// `status_metrics:` pattern's first match, in the same
756    /// candidate order [`Self::collect_status_values`] uses
757    /// (relevancy aggregates as percent, then service-time
758    /// quantiles as milliseconds). Feeds the key-metric row's
759    /// gutter cell TREND (SRD-92 R4): the cell shows the metric's
760    /// history as a sparkline — the current value's single
761    /// placement is the chips text in the row body. `None` when
762    /// nothing matches or nothing has been measured yet.
763    pub fn collect_status_primary(&self, patterns: &[String]) -> Option<(String, f64)> {
764        if patterns.is_empty() {
765            return None;
766        }
767        let mut candidates: Vec<(String, f64)> = Vec::new();
768        for live in self.collect_relevancy_live() {
769            if live.total_count == 0 {
770                continue;
771            }
772            candidates.push((live.name, live.total_mean * 100.0));
773        }
774        let snap = self.service_time.peek_snapshot();
775        let h = &snap.histogram;
776        if !h.is_empty() {
777            let ms = |n: f64| n / 1e6;
778            candidates.push((
779                "latency_p50".to_string(),
780                ms(h.value_at_quantile(0.50) as f64),
781            ));
782            candidates.push((
783                "latency_p99".to_string(),
784                ms(h.value_at_quantile(0.99) as f64),
785            ));
786            candidates.push(("latency_max".to_string(), ms(h.max() as f64)));
787            candidates.push(("latency_mean".to_string(), ms(h.mean())));
788        }
789        for pat in patterns {
790            for (name, val) in &candidates {
791                if glob_match(pat, name) {
792                    return Some((name.clone(), *val));
793                }
794            }
795        }
796        None
797    }
798
799    /// Collect status counters from all registered dispensers.
800    pub fn collect_status_counters(&self) -> Vec<(String, u64)> {
801        let mut counters = Vec::new();
802        if let Ok(guard) = self.dispensers.lock()
803            && let Some(ref disps) = *guard
804        {
805            for disp in disps.iter() {
806                for (name, total) in disp.status_counters() {
807                    counters.push((name.to_string(), total));
808                }
809            }
810        }
811        counters
812    }
813}
814
815/// Map a canonical status-chip metric name to its short display
816/// label. Keeps the underlying pattern-matching identity stable
817/// (workloads' `status_metrics: ["latency_*"]` keeps working)
818/// while the operator-facing chip stays terse — `P50` / `P99` /
819/// `Pmax` / `Pmean` instead of `latency_p50` / etc. Identity
820/// passthrough for any name without a shortcut.
821fn chip_display_label(name: &str) -> &str {
822    match name {
823        "latency_p50" => "P50",
824        "latency_p99" => "P99",
825        "latency_max" => "Pmax",
826        "latency_mean" => "Pmean",
827        other => other,
828    }
829}
830
831/// [`DynamicCapture`] adapter for [`ActivityMetrics`]. Captures the
832/// dynamic surface — per-error-type counters and adapter-specific
833/// metrics from registered dispensers — that isn't known at
834/// `register_on` time and therefore can't live in the static
835/// component instrument registry.
836struct ActivityMetricsDynamic {
837    metrics: Arc<ActivityMetrics>,
838    /// Per-counter previous-value baseline for delta emission on
839    /// the `drain=true` path. Keyed by `counter.labels().identity_hash()`.
840    /// Mirrors the per-component baseline that `Component` keeps for
841    /// registered counters; per-error-type counters live outside
842    /// the registry so the baseline travels with the hook.
843    ///
844    /// Why deltas: `MetricSet::combine_into` for Counter is
845    /// `total = a.total.saturating_add(b.total)` — the cascade
846    /// coalesce path treats Counter.total as the per-interval
847    /// delta and SUMS across intervals. Emitting absolutes here
848    /// would inflate as the cascade coalesces.
849    prev_counters: std::sync::Mutex<std::collections::HashMap<u64, u64>>,
850}
851
852impl nmbrs_metrics::component::DynamicCapture for ActivityMetricsDynamic {
853    fn capture_into(&self, out: &mut MetricSet, now: Instant, drain: bool) {
854        use nmbrs_metrics::snapshot::{MetricType, MetricValue, split_name_label};
855
856        // Per-error-type counters.
857        // - drain=true (cadence path): emit deltas vs. the stored
858        //   baseline so cascade coalesce sums across intervals
859        //   without inflation.
860        // - drain=false (peek path): emit absolute totals.
861        let error_counts = self
862            .metrics
863            .error_type_counts
864            .lock()
865            .unwrap_or_else(|e| e.into_inner());
866        if drain {
867            let mut prev = self.prev_counters.lock().unwrap_or_else(|e| e.into_inner());
868            for counter in error_counts.values() {
869                let (name, lbl) = split_name_label(counter.labels());
870                let current = counter.get();
871                let key = counter.labels().identity_hash();
872                let previous = prev.insert(key, current).unwrap_or(0);
873                out.insert_counter(name, lbl, current.saturating_sub(previous), now);
874            }
875        } else {
876            for counter in error_counts.values() {
877                let (name, lbl) = split_name_label(counter.labels());
878                out.insert_counter(name, lbl, counter.get(), now);
879            }
880        }
881
882        // Adapter-specific metrics from each registered dispenser.
883        // Passthrough — the adapter decides delta vs. absolute
884        // semantics for its own metrics.
885        if let Some(ref disps) = *self
886            .metrics
887            .dispensers
888            .lock()
889            .unwrap_or_else(|e| e.into_inner())
890        {
891            for dispenser in disps.iter() {
892                for (family, metric_labels, value) in dispenser.adapter_metrics() {
893                    let mtype = match &value {
894                        MetricValue::Counter(_) => MetricType::Counter,
895                        MetricValue::Gauge(_) => MetricType::Gauge,
896                        MetricValue::Histogram(_) => MetricType::Summary,
897                        MetricValue::BucketedHistogram(_) => MetricType::Histogram,
898                        MetricValue::Info(_) => MetricType::Info,
899                        MetricValue::StateSet(_) => MetricType::StateSet,
900                    };
901                    out.insert_metric(family, mtype, metric_labels, value, now);
902                }
903            }
904        }
905    }
906}
907
908/// A running activity.
909pub struct Activity {
910    pub config: ActivityConfig,
911    pub labels: Labels,
912    pub metrics: Arc<ActivityMetrics>,
913    pub op_sequence: OpSequence,
914    /// SRD-83 — this phase node's own scope kernel (the structural
915    /// walk's `cached_kernel`). Stop-condition predicates bind to THIS
916    /// native scope as it sits, not a conjured root. `None` when the
917    /// phase has no installed kernel (then no conditions evaluate).
918    pub phase_kernel: Option<Arc<crate::scope_kernel::ScopeKernel>>,
919    /// SRD-82 — the phase shell's [`crate::error_policy::ErrorPolicy`]
920    /// (op router + aggregate guard), resolved at scope-init from the
921    /// parent policy so equal configs share one instance. Built
922    /// standalone only on the test/library path ([`Self::with_params`]).
923    pub error_policy: Arc<crate::error_policy::ErrorPolicy>,
924    /// Source factory — creates per-fiber readers. All phases go through
925    /// sources. `cycles: N` desugars to `range(0, N)`.
926    source_factory: Arc<dyn polydat::iteration::source::DataSourceFactory>,
927    /// Resolved workload parameters (constant per run).
928    pub workload_params: Arc<std::collections::HashMap<String, String>>,
929    /// Shared flag: set to true when a `stop` error handler fires.
930    /// All fibers check this and exit their loop when set.
931    pub stop_flag: Arc<std::sync::atomic::AtomicBool>,
932    /// Per-execution walk-stop flag (SRD-82 Part 4), cloned from this
933    /// execution's `WorkloadShell`. Distinct from `stop_flag` (which is
934    /// this phase's OWN stop): set when the scenario WALK halts — a
935    /// sibling phase failed (a fault) or a stop condition tripped — so
936    /// in-flight fibers abort cooperatively and a concurrent
937    /// (`Bounded(N>1)`) sibling phase stops instead of draining. `None`
938    /// outside a walk (tests, the library shim) → never aborts.
939    pub walk_stop: Option<Arc<std::sync::atomic::AtomicBool>>,
940    /// SRD-82 Part 6 — set ONLY for a daemon phase: the daemon-group
941    /// completion flag, latched by the scenario shell once the scope's
942    /// foreground phases finish. A daemon phase's fibers poll it at their
943    /// cooperative boundaries and exit, so the daemon stops when the
944    /// foreground it shadows completes. `None` for foreground phases.
945    pub daemon_stop: Option<Arc<std::sync::atomic::AtomicBool>>,
946    /// First error message that triggered `stop_flag` — captured
947    /// once (the first stopping error wins, subsequent fibers'
948    /// errors don't overwrite). Surfaced in the phase-level
949    /// error so the user doesn't have to grep the per-cycle
950    /// log to learn what actually stopped the run.
951    pub stop_reason: Arc<std::sync::Mutex<Option<String>>>,
952    /// SRD-83 Part 5 — the two-axis Outcome of the FIRST stop-condition
953    /// trip that stopped this phase, latched together with `stop_reason`
954    /// (same first-stopper-wins discipline, written inside the same slot
955    /// win). A `stop` effect latches Interrupted+Succeeded — a clean
956    /// early halt whose partial result the phase keeps; `fail`/`abort`
957    /// latch Interrupted+Failed. The executor reads this at phase end so
958    /// the shell adopts the condition's DECLARED outcome instead of
959    /// deriving failure from the bare stop flag. `None` whenever the
960    /// stop came from any other source (error router `stop` verb, walk
961    /// stop, poll timeout, Ctrl-C) — those keep failure semantics.
962    pub stop_outcome: Arc<std::sync::Mutex<Option<crate::phase_outcome::Outcome>>>,
963    /// SRD-76 — chronologically ordered per-cycle error
964    /// records. Populated by the per-cycle dispatch path
965    /// (alongside the existing `stop_reason` formatted
966    /// string) so the executor can drain a structured
967    /// list into `PhaseOutcome.errors` at phase end. The
968    /// `stop_reason` string stays — it's the single
969    /// load-bearing format the executor reads to compose
970    /// the `phase 'X' stopped by error handler:` log
971    /// line. This buffer is the orthogonal structured
972    /// projection.
973    pub phase_errors: Arc<std::sync::Mutex<Vec<crate::phase_outcome::PhaseErrorDetail>>>,
974    /// Final validation metrics frame, populated after all cycles complete.
975    /// Read by the metrics capture thread after the activity finishes.
976    pub validation_frame: Arc<std::sync::Mutex<Option<MetricSet>>>,
977    /// Optional handle to this activity's component in the session tree.
978    /// Set by the runner via [`Self::attach_component`] before
979    /// execution; when present, the executor declares the
980    /// `concurrency` control on it (SRD 23) and wires the
981    /// [`crate::fiber_pool::ConcurrencyApplier`] so runtime writes
982    /// resize the fiber pool.
983    pub component: Option<Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>>,
984    /// SRD-32a Push 3 — workload-root wrapper-composition
985    /// override. When populated (from the workload's
986    /// `wrappers: { order: [...] }` block), every op
987    /// template that doesn't carry its own
988    /// per-template override uses this innermost-to-outermost
989    /// list as its composition order. Validated against the
990    /// per-op triggered set at cascade time; mismatch is a
991    /// hard error per SRD-32a §"Workload-level override".
992    pub wrappers_override: Option<Vec<String>>,
993    /// SRD-32a Push 3 — CLI `--wrap-default-order` override.
994    /// Replaces the resolver's built-in `DEFAULT_ORDER`
995    /// tiebreaker for this activity. `None` ⇒ resolver uses
996    /// the built-in order. Distinct from
997    /// `wrappers_override`: that pins the per-op stack;
998    /// this changes the tiebreaker used when constraints
999    /// leave order ambiguous.
1000    pub wrap_default_order: Option<Vec<String>>,
1001    /// Shared retry-exemplar sampling config (`exec_events`): every
1002    /// tries wrapper in this activity that does NOT pin its own
1003    /// `retry_exemplar_*` op params samples through this cell, so
1004    /// the `retry_exemplar_rate` / `retry_exemplar_max_hz` dynamic
1005    /// controls (declared in [`Self::attach_component`]) move them
1006    /// all with one atomic store — push-on-set, no per-op control
1007    /// traffic, and the read only happens on the retry path.
1008    pub exemplar_config: Arc<crate::exec_events::ExemplarConfig>,
1009    /// Shared per-phase retry advisory gate (`exec_events`): one
1010    /// first-sighting advisory per error class per phase, capped —
1011    /// the default-on signal that the retry loop started absorbing
1012    /// errors. Ops opt out with `retry_advisory: off`.
1013    pub advisory_gate: Arc<crate::exec_events::AdvisoryGate>,
1014    /// Phase memo — a short operator-visible string that the
1015    /// `memo` wrapper publishes via `before:` / `after:`
1016    /// templates. Read by the inline-status readout and
1017    /// rendered as `[[ <memo> ]]` above the status line when
1018    /// non-empty. Lock-free atomic so the inline thread can
1019    /// load it every tick without blocking the executor.
1020    /// Default empty.
1021    pub memo: Arc<arc_swap::ArcSwap<String>>,
1022    /// Phase gutter — the contextual left-gutter cell content that
1023    /// the `gutter` wrapper publishes (distinct from `memo`, which
1024    /// owns the `[[ ... ]]` header line). `None` ⇒ the display
1025    /// derives the cell automatically (completion bar for metered
1026    /// phases, latency trend for daemons); `Some` overrides that
1027    /// derivation with the workload-declared spec. Lock-free
1028    /// atomic for the same reason as `memo`.
1029    pub gutter: Arc<arc_swap::ArcSwapOption<crate::wrappers::gutter::GutterSpec>>,
1030    /// The op-declared DURING-execution gutter template (kind +
1031    /// template), retained for the guaranteed one-final-update at
1032    /// phase end when no `final:` form is declared.
1033    pub gutter_spec: std::sync::Mutex<Option<(crate::wrappers::gutter::GutterKind, String)>>,
1034    /// The op-declared `final:` gutter template — evaluated once at
1035    /// phase end (wires first, status-metric aggregates as fallback)
1036    /// and rendered as the ✓ outcome detail line's gutter cell.
1037    pub gutter_final_spec: std::sync::Mutex<Option<(crate::wrappers::gutter::GutterKind, String)>>,
1038    /// SRD-75 phase-poll context. When present, the fiber
1039    /// loop checks the predicate after each source-exhaustion
1040    /// event; if false and the timeout hasn't elapsed, the
1041    /// source factory rewinds and the loop continues.
1042    /// `None` ⇒ no phase-poll (standard activity semantics).
1043    /// Set by the executor at run-phase entry; not part of
1044    /// the YAML-derived `ActivityConfig`.
1045    pub phase_poll: Option<PhasePollContext>,
1046}
1047
1048/// Runtime context for SRD-75 phase-level poll. Carried on
1049/// `Activity` when the phase declares a `poll:` block;
1050/// consumed by the fiber loop after each source-exhaustion
1051/// event to decide whether to terminate (predicate satisfied
1052/// or timeout) or rewind and run another iteration.
1053#[derive(Clone)]
1054pub struct PhasePollContext {
1055    /// Handle to the phase scope kernel. After each iteration the fiber loop
1056    /// resolves [`crate::wrappers::condition::UNTIL_BINDING`] through
1057    /// `CycleWires` over the per-fiber kernel — the same call an op-level
1058    /// `poll:` makes — and ends the loop when it reads truthy per
1059    /// [`crate::wrappers::condition::is_truthy`]. `lookup` would not do: the
1060    /// predicate is a DYNAMIC binding fed by per-iteration capture writes, so
1061    /// it has to be pulled or it returns the last-evaluated value forever.
1062    pub kernel: Arc<crate::scope_kernel::ScopeKernel>,
1063    /// Sleep between iterations (after a predicate check
1064    /// returns "not done").
1065    pub interval: std::time::Duration,
1066    /// Wall-clock cap on the whole poll loop. Computed at
1067    /// run-phase entry as `Instant::now() + timeout_ms`.
1068    pub deadline: std::time::Instant,
1069    /// `Instant` the loop started — used to compute the
1070    /// elapsed-time value emitted under `metric_name`
1071    /// (if set) on successful completion.
1072    pub started_at: std::time::Instant,
1073    /// Optional named metric to emit on successful loop
1074    /// completion (predicate fired). Value is the elapsed
1075    /// wall-clock decoded per the existing `_ns` / `_us` /
1076    /// `_ms` / `_s` / `_m` / `_h` suffix convention. `None`
1077    /// ⇒ no metric written.
1078    pub metric_name: Option<String>,
1079    /// Tolerated consecutive retryable inner-op errors
1080    /// before the loop propagates the error. Mirrors the
1081    /// per-op `PollingDispenser` `max_error_retries`
1082    /// semantics. Default `0` (strict).
1083    pub max_error_retries: u32,
1084    /// What to do when the `deadline` fires without
1085    /// satisfying the predicate. SRD-75 §"on_timeout":
1086    /// - `Error` (default) — set the activity's
1087    ///   stop_flag + stop_reason. The phase returns an
1088    ///   error; the scenario walker's error-routing
1089    ///   policy decides whether sibling phases continue.
1090    /// - `Abort` — additionally call
1091    ///   `session_signals::request_stop()` so the whole
1092    ///   scenario terminates. The workload-author
1093    ///   declares the predicate's satisfaction as a
1094    ///   precondition for any downstream phase being
1095    ///   meaningful; a stuck synchronizer invalidates
1096    ///   the rest of the run.
1097    pub on_timeout: PhasePollTimeoutPolicy,
1098    /// SRD-75 (C5) — strict-gate selectors: each must resolve to a
1099    /// registered instrument before the gate's predicate is trusted.
1100    /// While any selector is unresolved the gate HOLDS (the predicate
1101    /// is not consulted — an unregistered family reads 0.0, which
1102    /// could satisfy a `>=`-shaped predicate spuriously); past
1103    /// [`Self::require_grace`] an unresolved selector is a hard
1104    /// `poll_require` failure.
1105    pub require: Vec<String>,
1106    /// Grace deadline for [`Self::require`] resolution: the first
1107    /// poll interval after loop start (late registration is normal —
1108    /// a producing daemon registers its instruments as it spins up).
1109    pub require_grace: std::time::Instant,
1110}
1111
1112/// SRD-75 `on_timeout` policy — see [`PhasePollContext::on_timeout`].
1113#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1114pub enum PhasePollTimeoutPolicy {
1115    /// Phase fails; scenario walker's error-routing policy
1116    /// decides downstream behaviour.
1117    #[default]
1118    Error,
1119    /// Phase fails AND `session_signals::request_stop()`
1120    /// is called — the scenario walker observes the
1121    /// global stop on its next iteration check and
1122    /// terminates the whole run.
1123    Abort,
1124}
1125
1126/// Invoke [`DriverAdapter::declare_controls`] for each unique adapter
1127/// instance against the given parent component, deduping by
1128/// `Arc`-pointer identity. The same adapter `Arc` may be entered
1129/// into the map under multiple alias keys; this guarantees each
1130/// physical instance gets exactly one declaration call per
1131/// invocation of this helper.
1132///
1133/// Called from two sites:
1134///
1135/// 1. The phase executor at component-attach time, so
1136///    `dryrun=controls` walks a populated tree before any
1137///    cycles run.
1138/// 2. [`Activity::run_with_adapters`] at run start, so adapters
1139///    that only ever materialize at run time still get declared.
1140///
1141/// Adapter implementations are expected to be idempotent — calling
1142/// this helper twice against the same parent must not produce
1143/// duplicate subcomponents or duplicate-name control declarations.
1144pub fn declare_adapter_controls(
1145    adapters: &std::collections::HashMap<String, Arc<dyn DriverAdapter>>,
1146    component: &Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
1147) {
1148    let mut seen: Vec<*const dyn DriverAdapter> = Vec::new();
1149    for adapter in adapters.values() {
1150        let ptr = Arc::as_ptr(adapter);
1151        if seen.contains(&ptr) {
1152            continue;
1153        }
1154        seen.push(ptr);
1155        adapter.declare_controls(component);
1156    }
1157}
1158
1159impl Activity {
1160    pub fn new(config: ActivityConfig, parent_labels: &Labels, op_sequence: OpSequence) -> Self {
1161        Self::with_params(
1162            config,
1163            parent_labels,
1164            op_sequence,
1165            std::collections::HashMap::new(),
1166        )
1167    }
1168
1169    pub fn with_params(
1170        config: ActivityConfig,
1171        parent_labels: &Labels,
1172        op_sequence: OpSequence,
1173        params: std::collections::HashMap<String, String>,
1174    ) -> Self {
1175        // Library / test path: no session root policy, so build a
1176        // standalone (un-shared) policy from the config. The real
1177        // execution path resolves the shared instance from the parent
1178        // policy and passes it to `with_params_and_sigdigs`.
1179        let error_policy =
1180            crate::error_policy::ErrorPolicy::standalone(crate::error_policy::PolicyConfig::new(
1181                config.error_spec.clone(),
1182                config.error_rate_max,
1183            ));
1184        let metric_detail = metric_detail_from_params(&params);
1185        Self::with_params_and_sigdigs(
1186            config,
1187            parent_labels,
1188            op_sequence,
1189            params,
1190            nmbrs_metrics::instruments::histogram::DEFAULT_HDR_SIGDIGS,
1191            error_policy,
1192            // This shim is the no-phase-kernel path (tests / library use);
1193            // the executor's run_phase path passes the phase node's kernel.
1194            None,
1195            &metric_detail,
1196        )
1197    }
1198
1199    /// Build an activity with explicit HDR significant-digits
1200    /// precision. Used by the runner after it resolves
1201    /// `hdr.sigdigs` from the session root (SRD 40); every
1202    /// histogram the activity owns is constructed at this
1203    /// precision. Callers that don't resolve from a tree can
1204    /// use [`Self::with_params`] which defaults to
1205    /// [`nmbrs_metrics::instruments::histogram::DEFAULT_HDR_SIGDIGS`].
1206    pub fn with_params_and_sigdigs(
1207        config: ActivityConfig,
1208        parent_labels: &Labels,
1209        op_sequence: OpSequence,
1210        params: std::collections::HashMap<String, String>,
1211        sigdigs: u8,
1212        error_policy: Arc<crate::error_policy::ErrorPolicy>,
1213        phase_kernel: Option<Arc<crate::scope_kernel::ScopeKernel>>,
1214        metric_detail: &MetricDetailConfig,
1215    ) -> Self {
1216        let labels = parent_labels.clone();
1217        // SRD-91 — counter-vs-timer detail for the op-outcome instruments,
1218        // resolved by the caller from the run's effective params (the
1219        // executor passes the CLI-overlaid set; the library shim derives
1220        // from its own params). Default: timers.
1221        let metrics = Arc::new(ActivityMetrics::with_sigdigs(
1222            &labels,
1223            sigdigs,
1224            metric_detail,
1225        ));
1226        // All phases go through sources. cycles: N desugars to range(0, N).
1227        // Named cursors in Polydat provide their own factory via config.source_factory.
1228        let source_factory: Arc<dyn polydat::iteration::source::DataSourceFactory> =
1229            config.source_factory.clone().unwrap_or_else(|| {
1230                Arc::new(polydat::iteration::source::RangeSourceFactory::named(
1231                    "cycles",
1232                    0,
1233                    config.cycles,
1234                ))
1235            });
1236
1237        Self {
1238            config,
1239            labels,
1240            metrics,
1241            op_sequence,
1242            error_policy,
1243            phase_kernel,
1244            source_factory,
1245            workload_params: Arc::new(params),
1246            stop_flag: Arc::new(std::sync::atomic::AtomicBool::new(false)),
1247            walk_stop: None,
1248            daemon_stop: None,
1249            stop_reason: Arc::new(std::sync::Mutex::new(None)),
1250            stop_outcome: Arc::new(std::sync::Mutex::new(None)),
1251            phase_errors: Arc::new(std::sync::Mutex::new(Vec::new())),
1252            validation_frame: Arc::new(std::sync::Mutex::new(None)),
1253            component: None,
1254            wrappers_override: None,
1255            wrap_default_order: None,
1256            exemplar_config: Arc::new(crate::exec_events::ExemplarConfig::new(0.0, 5.0)),
1257            advisory_gate: Arc::new(crate::exec_events::AdvisoryGate::new()),
1258            memo: Arc::new(arc_swap::ArcSwap::from_pointee(String::new())),
1259            gutter: Arc::new(arc_swap::ArcSwapOption::empty()),
1260            gutter_spec: std::sync::Mutex::new(None),
1261            gutter_final_spec: std::sync::Mutex::new(None),
1262            phase_poll: None,
1263        }
1264    }
1265
1266    /// Whether this execution's scenario walk has halted (SRD-82 Part
1267    /// 4). Fibers poll this at their cooperative boundaries, alongside
1268    /// `stop_flag` and `session_signals::stop_requested()`, to abort an
1269    /// in-flight phase when a sibling failed or a stop condition
1270    /// tripped. `false` when no walk-stop flag is wired (tests / shim).
1271    #[inline]
1272    pub fn walk_stop_requested(&self) -> bool {
1273        self.walk_stop
1274            .as_ref()
1275            .is_some_and(|f| f.load(std::sync::atomic::Ordering::Relaxed))
1276    }
1277
1278    /// Whether this (daemon) phase's group has signalled completion
1279    /// (SRD-82 Part 6) — the scenario shell latches `daemon_stop` once
1280    /// the scope's foreground phases finish, and the daemon's fibers poll
1281    /// this to exit. `false` for a foreground phase (no flag wired).
1282    #[inline]
1283    pub fn daemon_stop_requested(&self) -> bool {
1284        self.daemon_stop
1285            .as_ref()
1286            .is_some_and(|f| f.load(std::sync::atomic::Ordering::Relaxed))
1287    }
1288
1289    /// SRD-92 Step 0 — the cooperative-stop view at a loop BREAK boundary:
1290    /// the activity `stop_flag`, the global / per-execution session stop,
1291    /// the SRD-83 `walk_stop`, and the SRD-82 P6 `daemon_stop`. Replaces the
1292    /// scattered per-flag loads at the fiber boundaries. (The
1293    /// failure-determining return deliberately uses a different set that
1294    /// EXCLUDES `daemon_stop` — see
1295    /// [`crate::session_signals::StopView::abnormal`].)
1296    #[inline]
1297    pub fn stopped(&self) -> bool {
1298        self.stop_flag.load(std::sync::atomic::Ordering::Relaxed)
1299            || crate::session_signals::stop_requested()
1300            || self.walk_stop_requested()
1301            || self.daemon_stop_requested()
1302    }
1303
1304    /// The portable [`StopView`](crate::session_signals::StopView) for this
1305    /// activity — handed to the `while:` wrapper (a `Send + 'static`
1306    /// dispenser that cannot borrow the activity) so its loop observes the
1307    /// full stop set, not just `stop_flag`. Built once at wrapper
1308    /// construction, after `walk_stop` / `daemon_stop` are set in `run_phase`.
1309    pub fn stop_view(&self) -> crate::session_signals::StopView {
1310        crate::session_signals::StopView::new(
1311            Some(self.stop_flag.clone()),
1312            self.walk_stop.clone(),
1313            self.daemon_stop.clone(),
1314        )
1315    }
1316
1317    /// SRD-32a Push 3 — set the workload-root wrapper-
1318    /// composition override on this activity. Pass `None` to
1319    /// clear; pass `Some(order)` to install. The order list
1320    /// is innermost-to-outermost; per-op `wrappers:` blocks
1321    /// shadow this entry entirely.
1322    pub fn set_wrappers_override(&mut self, order: Option<Vec<String>>) {
1323        self.wrappers_override = order;
1324    }
1325
1326    /// SRD-32a Push 3 — set the resolver's default-order
1327    /// tiebreaker for this activity (CLI
1328    /// `--wrap-default-order`). `None` ⇒ the resolver uses
1329    /// its built-in `DEFAULT_ORDER` list.
1330    pub fn set_wrap_default_order(&mut self, order: Option<Vec<String>>) {
1331        self.wrap_default_order = order;
1332    }
1333
1334    /// Attach this activity to its component in the session tree.
1335    /// The runner creates the component and installs it here so
1336    /// `run_with_*` can register appliers on the activity's
1337    /// declared controls.
1338    ///
1339    /// Structural control declarations happen here — not at run
1340    /// time — so `dryrun=controls` (and every other pre-execution
1341    /// discovery path) sees the activity's controls without
1342    /// needing to start any cycles. Appliers that depend on
1343    /// run-time state (the fiber pool, the rate limiter) are
1344    /// registered later in `run_with_adapters`.
1345    pub fn attach_component(
1346        &mut self,
1347        component: Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
1348    ) {
1349        use crate::control_catalog::{CONCURRENCY, RATE};
1350        // Derive both controls from their single-source capability descriptors
1351        // (SRD-23) — name, range, and gauge come from the same `ControlDesc`
1352        // that `describe controls` reads, so the discovery surface and the
1353        // live knob can't drift. The instance-specific appliers (fiber-pool
1354        // resize, rate limiter) are registered later in `run_with_adapters`.
1355        let concurrency_control = CONCURRENCY.build_u32(self.config.concurrency as u32);
1356        component
1357            .read()
1358            .unwrap_or_else(|e| e.into_inner())
1359            .controls()
1360            .declare(concurrency_control);
1361
1362        // Declare a `rate` control whenever the activity config has a rate
1363        // set. Its reified gauge projects ops/sec so metric sinks and the
1364        // f64-writable surface (TUI `e` prompt, web POST, Polydat
1365        // `control_set`, `optimize.servo: rate`) all read and write the same
1366        // unit. The [`RateLimiterApplier`] gets registered at run time once
1367        // the limiter exists (see `run_with_adapters`).
1368        if let Some(rate) = self.config.rate {
1369            let rate_control = RATE.build_rate(rate);
1370            component
1371                .read()
1372                .unwrap_or_else(|e| e.into_inner())
1373                .controls()
1374                .declare(rate_control);
1375        }
1376
1377        // Retry-exemplar sampling (SRD-82 Part 3b / `exec_events`):
1378        // both knobs are push-on-set — the applier is one atomic
1379        // store into the activity's shared `ExemplarConfig`, which
1380        // every unpinned tries wrapper reads (only on its retry
1381        // path). Ops that pin `retry_exemplar_*` params hold private
1382        // cells the controls deliberately do not move — authored
1383        // matter wins, the live control moves the rest.
1384        {
1385            use crate::control_catalog::{RETRY_EXEMPLAR_MAX_HZ, RETRY_EXEMPLAR_RATE};
1386            use nmbrs_metrics::controls::SyncApplier;
1387            let cfg = self.exemplar_config.clone();
1388            let rate_control = RETRY_EXEMPLAR_RATE.build_f64(0.0);
1389            rate_control.register_applier(SyncApplier::new(move |v: f64| {
1390                cfg.set_rate(v);
1391                Ok(())
1392            }));
1393            component
1394                .read()
1395                .unwrap_or_else(|e| e.into_inner())
1396                .controls()
1397                .declare(rate_control);
1398
1399            let cfg = self.exemplar_config.clone();
1400            let hz_control = RETRY_EXEMPLAR_MAX_HZ.build_f64(5.0);
1401            hz_control.register_applier(SyncApplier::new(move |v: f64| {
1402                cfg.set_max_hz(v);
1403                Ok(())
1404            }));
1405            component
1406                .read()
1407                .unwrap_or_else(|e| e.into_inner())
1408                .controls()
1409                .declare(hz_control);
1410        }
1411        // Register every static instrument owned by ActivityMetrics
1412        // on this component so the cadence reporter's tree walk
1413        // sees them. Failures here are programming errors
1414        // (duplicate family on the activity's own component) —
1415        // panic so the issue surfaces during init.
1416        {
1417            let mut guard = component.write().unwrap_or_else(|e| e.into_inner());
1418            self.metrics
1419                .register_on(&mut guard)
1420                .expect("ActivityMetrics::register_on failed on a fresh activity component");
1421        }
1422        self.component = Some(component);
1423    }
1424
1425    /// Get a shared reference to the metrics for external capture.
1426    pub fn shared_metrics(&self) -> Arc<ActivityMetrics> {
1427        self.metrics.clone()
1428    }
1429
1430    /// Run the activity with a single adapter for all ops.
1431    pub async fn run_with_driver(
1432        self,
1433        adapter: Arc<dyn DriverAdapter>,
1434        op_builder: Arc<crate::synthesis::OpBuilder>,
1435    ) -> bool {
1436        let mut adapters = std::collections::HashMap::new();
1437        let name = adapter.name().to_string();
1438        adapters.insert(name.clone(), adapter);
1439        self.run_with_adapters(adapters, &name, op_builder).await
1440    }
1441
1442    /// Run the activity with multiple adapters (SRD 38/40).
1443    ///
1444    /// Each op template's `adapter` param selects which adapter to use.
1445    /// Templates without an explicit adapter use `default_adapter`.
1446    /// At init time: maps each template to a dispenser from the
1447    /// appropriate adapter. Per fiber: creates a FiberBuilder. Per
1448    /// cycle: resolves fields via GK, executes via dispenser.
1449    /// Returns true if the activity was stopped by an error handler.
1450    pub async fn run_with_adapters(
1451        self,
1452        adapters: std::collections::HashMap<String, Arc<dyn DriverAdapter>>,
1453        default_adapter: &str,
1454        op_builder: Arc<crate::synthesis::OpBuilder>,
1455    ) -> bool {
1456        let activity = Arc::new(self);
1457        let program = op_builder.program();
1458
1459        // Init time: map each template to a dispenser from its adapter,
1460        // then wrap with result traverser for consumption/capture.
1461        //
1462        // Dryrun injection: when the session is in dryrun mode
1463        // (`config.dry_run_mode` is `Some(mode)`), inject a logical
1464        // `dryrun: <mode>` parameter into every op template's
1465        // `params` map BEFORE the wrapping cascade sees them. This
1466        // triggers the outermost `DryRunWrapper` to install for
1467        // every op; the wrapper short-circuits at cycle time and
1468        // suppresses only the outbound `execute()`.
1469        //
1470        // The real adapter's full lifecycle still runs (connect,
1471        // prepare, metadata) — `dryrun=cycle` means "construct a
1472        // fully-executable cycle path, then suppress only the
1473        // outbound call." So `dry_run_mode` is sourced from the
1474        // session config, NOT from any adapter substitution.
1475        let dryrun_mode: Option<String> = activity.config.dry_run_mode.clone();
1476        let templates_owned: Vec<nmbrs_workload::model::ParsedOp>;
1477        let templates: &[nmbrs_workload::model::ParsedOp] =
1478            if let Some(mode) = dryrun_mode.as_deref() {
1479                templates_owned = activity
1480                    .op_sequence
1481                    .templates()
1482                    .iter()
1483                    .map(|t| {
1484                        let mut clone = t.clone();
1485                        clone
1486                            .params
1487                            .insert("dryrun".into(), serde_json::Value::String(mode.to_string()));
1488                        // dryrun=fields also forces the fields wrapper
1489                        // on so the rendered op text reaches stdout
1490                        // even though DRYRUN short-circuits the
1491                        // adapter call. The fields wrapper is composed
1492                        // OUTER of dryrun (see
1493                        // wrapper_resolver::DEFAULT_ORDER), so its
1494                        // pre-execute render runs first; the
1495                        // subsequent DRYRUN short-circuit suppresses
1496                        // the real adapter call.
1497                        if mode == "fields" {
1498                            clone
1499                                .params
1500                                .insert("fields".into(), serde_json::Value::Bool(true));
1501                        }
1502                        clone
1503                    })
1504                    .collect();
1505                // The user passed `dryrun=<mode>` on the CLI — they
1506                // already know what they asked for. Keep the
1507                // marker-injection record at Debug so step-through /
1508                // session-log audits can still find it without
1509                // narrating it back on stderr every phase.
1510                crate::diag!(
1511                    crate::observer::LogLevel::Debug,
1512                    "dryrun={mode}: injected marker into {n} op template(s)",
1513                    mode = mode,
1514                    n = templates_owned.len()
1515                );
1516                &templates_owned[..]
1517            } else {
1518                activity.op_sequence.templates()
1519            };
1520
1521        // Validate all bind points are resolvable before execution
1522        let program_for_op = |name: &str| op_builder.program_for_op(name);
1523        if let Err(e) = crate::synthesis::validate_bind_points(templates, &program_for_op) {
1524            crate::diag!(crate::observer::LogLevel::Error, "error: {e}");
1525            return true;
1526        }
1527
1528        // Adapter-level dynamic controls (SRD 23). The phase
1529        // executor already declared adapter controls at attach
1530        // time so `dryrun=controls` saw them; calling again here
1531        // is the safety net for adapters that materialize only
1532        // at run time. Adapter `declare_controls` impls are
1533        // contractually idempotent — see `declare_adapter_controls`.
1534        if let Some(component) = activity.component.as_ref() {
1535            declare_adapter_controls(&adapters, component);
1536        }
1537
1538        let traversal_stats = Arc::new(crate::wrappers::TraversalStats {
1539            metrics: activity.metrics.clone(),
1540        });
1541
1542        // SRD-32a — wrapper registry + resolver. The
1543        // registry is fixed at link time (every `inventory::
1544        // submit!` block in the binary contributes one
1545        // entry); the resolver carries the validated
1546        // default-order tiebreaker. Both are built once
1547        // here and reused for every op template in this
1548        // activity.
1549        let wrapper_registry = crate::wrapper_registry::WrapperRegistry::from_inventory();
1550        // SRD-32a Push 3 — CLI `--wrap-default-order` replaces
1551        // the resolver's built-in tiebreaker. When unset, the
1552        // resolver builds with its DEFAULT_ORDER. The CLI list
1553        // is validated against the constraint graph at
1554        // construction; an inconsistent list aborts the run.
1555        let wrapper_resolver = match &activity.wrap_default_order {
1556            Some(order) => {
1557                let names: Vec<&str> = order.iter().map(|s| s.as_str()).collect();
1558                crate::wrapper_resolver::WrapperResolver::from_names(&names, &wrapper_registry)
1559            }
1560            None => crate::wrapper_resolver::WrapperResolver::with_default_order(&wrapper_registry),
1561        };
1562        let wrapper_resolver = match wrapper_resolver {
1563            Ok(r) => r,
1564            Err(e) => {
1565                crate::diag!(
1566                    crate::observer::LogLevel::Error,
1567                    "error: wrapper default-order is inconsistent with the \
1568                     registered wrapper graph: {e}. CLI `--wrap-default-order` \
1569                     and the built-in default both must satisfy every \
1570                     registered constraint."
1571                );
1572                return true;
1573            }
1574        };
1575
1576        let mut dispensers: Vec<Arc<dyn OpDispenser>> = Vec::new();
1577        let mut validation_metrics: Vec<Arc<validation::ValidationMetrics>> = Vec::new();
1578        // SRD-40b §6/§7 — one `Component` per **op dispenser**
1579        // (= per op template), not per op execution. Op
1580        // dispensers are the durable CNS layer of the nmbrs
1581        // runtime; per-cycle op invocations are stack-ephemeral
1582        // and inherit the dispenser's component implicitly via
1583        // the wrapper-stack closure capture. Each component
1584        // carries `op=<template.name>` labels (child of the
1585        // activity component) so SRD-40b §7.2's duplicate-
1586        // family check (`Component::register_instrument`)
1587        // sees one dimensional cell per dispenser, surviving
1588        // for the run's duration. Held here to keep the Arc
1589        // alive.
1590        let mut dispenser_components: Vec<
1591            std::sync::Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
1592        > = Vec::new();
1593        // Per-template wrapper pull plan. Wrapper-side reads
1594        // (validation, conditional, throttle) go through this
1595        // `PullPlan` against the firing fiber's state — see
1596        // SRD 31 §"Pull plan vs bind plan". Adapter-side reads
1597        // moved to the generic `crate::wires::WireSource` surface
1598        // at SRD-68 Push 5; the legacy `field_pulls` /
1599        // `bind_plans` / `batch_configs` lists are retired.
1600        let mut pull_plans_per_template: Vec<crate::fixture::PullPlan> = Vec::new();
1601        for template in templates {
1602            // Resolve adapter: per-template override or default
1603            let adapter_name = template
1604                .params
1605                .get("adapter")
1606                .and_then(|v| v.as_str())
1607                .or_else(|| template.params.get("driver").and_then(|v| v.as_str()))
1608                .unwrap_or(default_adapter);
1609            let adapter = match adapters.get(adapter_name) {
1610                Some(a) => a,
1611                None => {
1612                    let available = adapters.keys().cloned().collect::<Vec<_>>().join(", ");
1613                    crate::diag!(
1614                        crate::observer::LogLevel::Error,
1615                        "error: unknown adapter '{adapter_name}' for op '{}' (available: {available})",
1616                        template.name
1617                    );
1618                    return true; // signal stop — cannot proceed without the adapter
1619                }
1620            };
1621
1622            if template.params.contains_key("batch") {
1623                crate::diag!(
1624                    crate::observer::LogLevel::Debug,
1625                    "[activity] op '{}' has batch param: {:?}",
1626                    template.name,
1627                    template.params.get("batch")
1628                );
1629            }
1630            // SRD 30 §"Core-first field processing": if the adapter
1631            // declares its known op fields, every key in
1632            // `template.op` must be one of them. Core has already
1633            // stripped its own fields during parse (activity_params
1634            // in nmbrs-workload), so anything left is an adapter
1635            // concern. Unknown fields are a typo or a misplaced
1636            // core directive — fail loudly rather than silently
1637            // dropping the field.
1638            if let Some(known) = adapter.known_op_fields() {
1639                let unknown: Vec<&String> = template
1640                    .op
1641                    .keys()
1642                    .filter(|k| !known.contains(&k.as_str()))
1643                    .collect();
1644                if !unknown.is_empty() {
1645                    let list = unknown
1646                        .iter()
1647                        .map(|s| s.as_str())
1648                        .collect::<Vec<_>>()
1649                        .join(", ");
1650                    crate::diag!(
1651                        crate::observer::LogLevel::Error,
1652                        "error: adapter '{}' does not recognize op fields [{list}] on op '{}'; known fields: [{}]",
1653                        adapter.name(),
1654                        template.name,
1655                        known.join(", "),
1656                    );
1657                    return true; // stop — misconfiguration
1658                }
1659            }
1660
1661            // Same idea for `template.params`: validate against
1662            // a closed vocabulary so silent-ignore traps like
1663            // `evaluations: { relevancy: ... }` (wrapper keys
1664            // the runtime never reads) cannot hide a
1665            // misconfigured op. Allowed keys are the union of:
1666            //   1. core op-level params consumed by the runtime
1667            //      (validation, batching, polling, weighting,
1668            //      adapter selection) — `CORE_OP_PARAMS`.
1669            //   2. workload/CLI-level params that the parser
1670            //      blast-merges into every op's params at parse
1671            //      time — `runner::KNOWN_PARAMS`.
1672            //   3. user-declared workload params from the
1673            //      workload's top-level `params:` block (e.g.
1674            //      `table`, `keyspace`, `num_items`). The parser
1675            //      threads these into every op's params during
1676            //      doc → block → op merge, where they're meant
1677            //      for `{name}` interpolation in op templates.
1678            //      Visible here as `activity.workload_params`.
1679            //   4. adapter-specific params declared via
1680            //      `DriverAdapter::known_op_params()`.
1681            // Anything else is a typo / misplaced wrapper / dead
1682            // YAML and is rejected.
1683            {
1684                let allowed_extras = adapter.known_op_params();
1685                let workload_keys = &activity.workload_params;
1686                // A params key is legitimate if it is (1) a non-wrapper
1687                // core param, (2) a wrapper field declared via the
1688                // registry's `owned_fields` — the durable, type-safe
1689                // route that replaced the fragile "it's also a CLI
1690                // param" coincidence for fields like `errors`/`tries`,
1691                // (3) a CLI param blast-merged into every op's params at
1692                // parse time, (4) an adapter-declared extra, or (5) a
1693                // workload-level param.
1694                let unknown_params: Vec<&String> = template
1695                    .params
1696                    .keys()
1697                    .filter(|k| {
1698                        !crate::validation::CORE_OP_PARAMS.contains(&k.as_str())
1699                            && !wrapper_registry.owns_field(k)
1700                            && !crate::runner::is_cli_param(k)
1701                            && !allowed_extras.contains(&k.as_str())
1702                            && !workload_keys.contains_key(k.as_str())
1703                    })
1704                    .collect();
1705                if !unknown_params.is_empty() {
1706                    let list = unknown_params
1707                        .iter()
1708                        .map(|s| s.as_str())
1709                        .collect::<Vec<_>>()
1710                        .join(", ");
1711                    let wrapper_vocab = wrapper_registry
1712                        .all_owned_fields()
1713                        .into_iter()
1714                        .collect::<Vec<_>>()
1715                        .join(", ");
1716                    crate::diag!(
1717                        crate::observer::LogLevel::Error,
1718                        "error: op '{}' has unknown params keys [{list}] — \
1719                         not a core op param, not a wrapper field, not \
1720                         declared by adapter '{}', and not a workload-level \
1721                         param. Known core op params: [{}]. Wrapper fields: \
1722                         [{}]. Adapter extras: [{}]. Did you misspell a \
1723                         wrapper key (e.g. `verify:` / `poll:` / `readout:`) \
1724                         or nest something under a key the runtime doesn't \
1725                         read?",
1726                        template.name,
1727                        adapter.name(),
1728                        crate::validation::CORE_OP_PARAMS.join(", "),
1729                        wrapper_vocab,
1730                        allowed_extras.join(", "),
1731                    );
1732                    return true; // stop — misconfiguration
1733                }
1734            }
1735
1736            // SRD-32a Push 2 — field ownership and misplaced-
1737            // field guard. The wrapper registry knows which
1738            // `params:` keys each wrapper consumes
1739            // (`owned_fields`) and what makes the wrapper
1740            // trigger. A field that's owned by a wrapper that
1741            // ISN'T triggered is misplaced — silently ignoring
1742            // it would mask a typo or a half-applied
1743            // configuration. Example: `poll_interval_ms: 5000`
1744            // on an op without `poll:` is a misconfiguration —
1745            // the operator probably meant to enable polling
1746            // but forgot the trigger.
1747            //
1748            // The closed-vocabulary check above catches "I
1749            // don't recognise this key at all"; THIS check
1750            // catches "I recognise the key but it has no effect
1751            // here." Both surface as hard errors.
1752            //
1753            // Note: we cross-check against `template.params`
1754            // for keys. The registry's owned_fields includes
1755            // a few names that live elsewhere on `ParsedOp`
1756            // (`if` → `template.condition`, `delay` →
1757            // `template.delay`); those happen to BE their
1758            // wrapper's trigger field too, so when set the
1759            // outer `(reg.triggers)(template)` short-circuits
1760            // and we never reach the params-key check for
1761            // them. The set is small enough that we don't
1762            // need a separate "where does this field live"
1763            // helper.
1764            {
1765                let violations = wrapper_registry
1766                    .misplaced_fields(crate::wrapper_registry::WrapperSubject::Op(template));
1767                if !violations.is_empty() {
1768                    for (wrapper, field) in &violations {
1769                        crate::diag!(
1770                            crate::observer::LogLevel::Error,
1771                            "error: op '{}': field `{field}` is owned by wrapper \
1772                             `{wrapper}`, but the trigger condition for `{wrapper}` \
1773                             is not satisfied (no trigger field set on this op). \
1774                             Either remove `{field}` or add the wrapper's trigger \
1775                             field. SRD-32a §\"Field ownership and parse-time \
1776                             validation\".",
1777                            template.name
1778                        );
1779                    }
1780                    return true; // stop — misconfiguration
1781                }
1782            }
1783
1784            // Per-op dispenser-init contract: `map_op` owns the
1785            // typed-binder verification for ITS op as part of
1786            // completing the currying stack — it constructs any
1787            // binders from the adapter's protocol-side metadata,
1788            // verifies them against `parent` via
1789            // `polydat::binder::verify_against_kernel`, and
1790            // surfaces any violation as `Err`. No
1791            // outside-the-dispenser-init phase for binders; the
1792            // map_op return signals the result.
1793            match adapter
1794                .map_op(template, op_builder.canonical_kernel_for_op(&template.name))
1795                .await
1796            {
1797                Ok(d) => {
1798                    let raw: Arc<dyn OpDispenser> = Arc::from(d);
1799
1800                    // Open the per-template scope fixture (SRD 32
1801                    // §"Init-Time Fixture and Consumer Self-
1802                    // Registration"). Each wrapper below registers
1803                    // its own Polydat name dependencies; the fixture is
1804                    // sealed after the wrapper chain is complete and
1805                    // the resulting PullPlan drives cycle-time reads.
1806                    //
1807                    // SRD-13d Phase 9 — when this op-template
1808                    // materialised its own kernel, the fixture
1809                    // builds its plan against THAT program so
1810                    // pulls resolve in the op-template scope.
1811                    // Flattened op-templates fall back to the
1812                    // activity-wide program (same scope as before
1813                    // Phase 9 landed).
1814                    let template_program = op_builder.program_for_op(&template.name);
1815                    let mut fx = crate::fixture::ScopeFixture::new(template_program.clone());
1816
1817                    // SRD-32a — resolve which wrappers fire and
1818                    // in what order. The plan is innermost-first;
1819                    // `traverse` is always inner. The cascade
1820                    // below dispatches to the existing per-
1821                    // wrapper `wrap()` factory based on the plan
1822                    // entries' names. Plan order matches the
1823                    // built-in default order, which mirrors the
1824                    // pre-SRD-32a hand-rolled cascade — existing
1825                    // tests exercise the same composition.
1826                    //
1827                    // SRD-32a Push 3 — override precedence:
1828                    //   1. Per-op `template.wrappers.order` shadows everything else.
1829                    //   2. Else workload-root `activity.wrappers_override`.
1830                    //   3. Else the resolver's default-order tiebreaker.
1831                    // Per-op shadows root entirely (no merge).
1832                    let per_op_override = template
1833                        .wrappers
1834                        .as_ref()
1835                        .filter(|c| !c.order.is_empty())
1836                        .map(|c| c.order.clone());
1837                    let effective_override =
1838                        per_op_override.or_else(|| activity.wrappers_override.clone());
1839                    let plan = match effective_override {
1840                        Some(order) => {
1841                            let order_strs: Vec<&str> = order.iter().map(|s| s.as_str()).collect();
1842                            wrapper_resolver.resolve_with_order(
1843                                crate::wrapper_registry::WrapperSubject::Op(template),
1844                                &wrapper_registry,
1845                                &order_strs,
1846                            )
1847                        }
1848                        None => wrapper_resolver.resolve(
1849                            crate::wrapper_registry::WrapperSubject::Op(template),
1850                            &wrapper_registry,
1851                        ),
1852                    };
1853                    let plan = match plan {
1854                        Ok(p) => p,
1855                        Err(e) => {
1856                            crate::diag!(
1857                                crate::observer::LogLevel::Error,
1858                                "error: op '{}': wrapper resolution failed: {e}",
1859                                template.name
1860                            );
1861                            return true;
1862                        }
1863                    };
1864
1865                    // SRD-32a §"Composition telemetry" — emit one
1866                    // Info-level line per assigned wrapper so
1867                    // operators can see, at session start, exactly
1868                    // which wrappers shape each op and how. Trivial
1869                    // wrappers (e.g. always-on `traverse`) return
1870                    // `None` from `describe_assignment` and are
1871                    // dropped from this list.
1872                    let assignments: Vec<(crate::wrapper_registry::WrapperName, String)> = plan
1873                        .iter_innermost_first()
1874                        .filter_map(|reg| {
1875                            (reg.describe_assignment)(crate::wrapper_registry::WrapperSubject::Op(
1876                                template,
1877                            ))
1878                            .map(|s| (reg.name, s))
1879                        })
1880                        .collect();
1881                    if !assignments.is_empty() {
1882                        // Per-op wrapper assignments are a diagnostic
1883                        // useful when chasing a wrapper-composition
1884                        // bug, not part of normal operator output.
1885                        // `nmbrs describe` renders the same stack on
1886                        // demand (see crates/nmbrs/src/describe.rs); session.log
1887                        // still captures this for postmortem.
1888                        crate::diag!(
1889                            crate::observer::LogLevel::Debug,
1890                            "op '{}' wrappers (innermost → outermost):",
1891                            template.name
1892                        );
1893                        for (i, (_, line)) in assignments.iter().enumerate() {
1894                            crate::diag!(crate::observer::LogLevel::Debug, "  {}. {}", i + 1, line);
1895                        }
1896                    }
1897
1898                    // SRD-82 Part 3b — resolve THIS op's error policy BEFORE
1899                    // the retry activation check, so the policy's `retry` verb
1900                    // can inject a retry budget (the errors→retry bridge). An
1901                    // op-template `errors:` derives a child of the phase
1902                    // policy (value-equality shared across ops declaring the
1903                    // same spec); no override inherits the phase policy by
1904                    // reference. The policy drives the OUTERMOST error-handler
1905                    // wrapper placed after the plan cascade below.
1906                    let op_error_policy =
1907                        match template.params.get("errors").and_then(|v| v.as_str()) {
1908                            Some(spec) => activity.error_policy.resolve_child(Some(
1909                                crate::error_policy::PolicyConfig::new(
1910                                    spec,
1911                                    activity.config.error_rate_max,
1912                                ),
1913                            )),
1914                            None => activity.error_policy.clone(),
1915                        };
1916
1917                    // Tries wrapper — CONDITIONAL innermost wrapper (SRD-82
1918                    // Part 3b): `tries:` is its sigil, the TOTAL attempts the
1919                    // op may make. A budget resolves from, in order: the op's
1920                    // own `tries:` field; the inherited phase/root `tries`
1921                    // (config — the phase FIELD must beat a scope wire, since
1922                    // a workload-root `tries` param also lands in GK scope as
1923                    // a constant and would otherwise shadow an explicit
1924                    // phase-level `tries: 1` pin); a `tries` wire defined in
1925                    // the op's scope (bindings); or a `retry`/`retry(N)` verb
1926                    // in the op's error policy (the injection bridge — N
1927                    // additional attempts → N+1 total; `errors:` and `tries:`
1928                    // stay orthogonal surfaces). No budget anywhere OR
1929                    // `tries: 1` → NO wrapper → single attempt (the
1930                    // error-handler wrapper records the tallies). `tries: 0`
1931                    // constructs the wrapper in its fail-without-executing
1932                    // mode. When constructed it owns the attempt loop, the
1933                    // `attempt_*` counters, and the per-attempt panic catch.
1934                    // Op `tries:` accepts the sugared number/string AND the
1935                    // map form `{count: N, backoff: {...}}` — read `count`
1936                    // from the map. Falls back to the phase/root `tries`, an
1937                    // in-scope `tries` wire, then an `errors:` retry-verb
1938                    // budget.
1939                    let op_tries: Option<u32> = template
1940                        .params
1941                        .get("tries")
1942                        .and_then(|v| match v {
1943                            serde_json::Value::Number(n) => n.as_u64().map(|n| n as u32),
1944                            serde_json::Value::String(s) => s.trim().parse::<u32>().ok(),
1945                            serde_json::Value::Object(m) => {
1946                                m.get("count").and_then(|c| c.as_u64()).map(|n| n as u32)
1947                            }
1948                            _ => None,
1949                        })
1950                        .or(activity.config.tries)
1951                        .or_else(|| {
1952                            raw.canonical_kernel()
1953                                .and_then(|k| {
1954                                    use polydat::kernel::interp::Lookup as _;
1955                                    polydat::kernel::interp::KernelLookup::new(k.as_ref())
1956                                        .lookup("tries")
1957                                })
1958                                .and_then(|v| match v {
1959                                    polydat::ast::Value::U64(n) => Some(n as u32),
1960                                    _ => None,
1961                                })
1962                        })
1963                        .or_else(|| {
1964                            op_error_policy
1965                                .router
1966                                .retry_verb_budget()
1967                                .map(|additional| additional.saturating_add(1))
1968                        });
1969                    let has_tries_wrapper = matches!(op_tries, Some(n) if n != 1);
1970                    // Retry pacing (compaction-demo diagnosis: an immediate
1971                    // `continue` retry loop hammers a dying server). Resolve
1972                    // (min, max, ratio) by precedence, highest first:
1973                    //   1. op `tries:` map `backoff: {ratio, min, max}`
1974                    //   2. op standalone `retry_backoff*` params (back-compat)
1975                    //   3. phase `tries:` map backoff (activity config)
1976                    //   4. built-in defaults: 100ms base, 10s cap, 2.0 ratio.
1977                    // `retry_backoff: 0` (base == 0) disables pacing.
1978                    let json_to_ms = |v: &serde_json::Value| -> Option<u64> {
1979                        match v {
1980                            serde_json::Value::Number(n) => n.as_u64().map(|n| n.to_string()),
1981                            serde_json::Value::String(s) => Some(s.clone()),
1982                            _ => None,
1983                        }
1984                        .and_then(|s| crate::timeval::parse_time_ms(&s).ok())
1985                    };
1986                    let mut backoff_base_ms: u64 = 100;
1987                    let mut backoff_max_ms: u64 = 10_000;
1988                    let mut backoff_ratio: f64 = 2.0;
1989                    // (3) phase-level map backoff
1990                    if let Some(bo) = activity.config.tries_backoff.as_ref() {
1991                        if let Some(r) = bo.ratio {
1992                            backoff_ratio = r;
1993                        }
1994                        if let Some(m) = bo
1995                            .min
1996                            .as_deref()
1997                            .and_then(|s| crate::timeval::parse_time_ms(s).ok())
1998                        {
1999                            backoff_base_ms = m;
2000                        }
2001                        if let Some(m) = bo
2002                            .max
2003                            .as_deref()
2004                            .and_then(|s| crate::timeval::parse_time_ms(s).ok())
2005                        {
2006                            backoff_max_ms = m;
2007                        }
2008                    }
2009                    // (1)/(2) op-level: the `tries:` map backoff wins; absent
2010                    // that, the standalone `retry_backoff*` params.
2011                    let op_map_backoff = template
2012                        .params
2013                        .get("tries")
2014                        .and_then(|v| v.as_object())
2015                        .and_then(|m| m.get("backoff"))
2016                        .and_then(|v| v.as_object());
2017                    if let Some(bo) = op_map_backoff {
2018                        if let Some(r) = bo.get("ratio").and_then(|v| v.as_f64()) {
2019                            backoff_ratio = r;
2020                        }
2021                        if let Some(m) = bo.get("min").and_then(&json_to_ms) {
2022                            backoff_base_ms = m;
2023                        }
2024                        if let Some(m) = bo.get("max").and_then(&json_to_ms) {
2025                            backoff_max_ms = m;
2026                        }
2027                    } else {
2028                        if let Some(m) = template.params.get("retry_backoff").and_then(&json_to_ms)
2029                        {
2030                            backoff_base_ms = m;
2031                        }
2032                        if let Some(m) = template
2033                            .params
2034                            .get("retry_backoff_max")
2035                            .and_then(&json_to_ms)
2036                        {
2037                            backoff_max_ms = m;
2038                        }
2039                        if let Some(r) = template
2040                            .params
2041                            .get("retry_backoff_ratio")
2042                            .and_then(|v| v.as_f64())
2043                        {
2044                            backoff_ratio = r;
2045                        }
2046                    }
2047                    // Retry-error counter-exemplars (`exec_events`): the
2048                    // fraction of caught-and-retried errors sampled onto
2049                    // the structured event sink (default 0.0 = off), and
2050                    // the emission-frequency ceiling that squelches spam.
2051                    // Standalone params, same channel as `retry_backoff*`.
2052                    let json_to_f64 = |v: &serde_json::Value| -> Option<f64> {
2053                        match v {
2054                            serde_json::Value::Number(n) => n.as_f64(),
2055                            serde_json::Value::String(s) => s.trim().parse::<f64>().ok(),
2056                            _ => None,
2057                        }
2058                    };
2059                    let pinned = template.params.contains_key("retry_exemplar_rate")
2060                        || template.params.contains_key("retry_exemplar_max_hz");
2061                    let sampler = if pinned {
2062                        // Authored pin: a private cell the dynamic
2063                        // controls deliberately do not move.
2064                        let exemplar_rate = template
2065                            .params
2066                            .get("retry_exemplar_rate")
2067                            .and_then(&json_to_f64)
2068                            .unwrap_or(0.0);
2069                        let exemplar_max_hz = template
2070                            .params
2071                            .get("retry_exemplar_max_hz")
2072                            .and_then(&json_to_f64)
2073                            .unwrap_or(5.0);
2074                        crate::exec_events::ExemplarSampler::pinned(exemplar_rate, exemplar_max_hz)
2075                    } else {
2076                        // Default: the activity's shared cell — the
2077                        // `retry_exemplar_rate`/`_max_hz` dynamic
2078                        // controls move every sampler built here with
2079                        // one atomic store.
2080                        crate::exec_events::ExemplarSampler::shared(
2081                            activity.exemplar_config.clone(),
2082                        )
2083                    };
2084                    // Default-on advisory; `retry_advisory: off` (or
2085                    // false) opts this op out of the shared gate.
2086                    let advisory_on = template
2087                        .params
2088                        .get("retry_advisory")
2089                        .map(|v| match v {
2090                            serde_json::Value::Bool(b) => *b,
2091                            serde_json::Value::String(s) => {
2092                                !s.eq_ignore_ascii_case("off") && !s.eq_ignore_ascii_case("false")
2093                            }
2094                            _ => true,
2095                        })
2096                        .unwrap_or(true);
2097                    let advisory = advisory_on.then(|| activity.advisory_gate.clone());
2098                    let raw = match op_tries {
2099                        Some(n) if n != 1 => crate::wrappers::TriesDispenser::wrap(
2100                            raw,
2101                            n,
2102                            activity.metrics.clone(),
2103                            backoff_base_ms,
2104                            backoff_max_ms,
2105                            backoff_ratio,
2106                            activity.stop_view(),
2107                            template.name.clone(),
2108                            sampler,
2109                            advisory,
2110                        ),
2111                        _ => raw,
2112                    };
2113
2114                    // Wrap with traversal. Traversal does not read Polydat
2115                    // values; no fixture registration needed. Always present per
2116                    // the registry's always-true trigger.
2117                    let mut current: Arc<dyn OpDispenser> =
2118                        crate::wrappers::TraversingDispenser::wrap(
2119                            raw,
2120                            template,
2121                            traversal_stats.clone(),
2122                        );
2123                    // Late-bound handoff from the metrics arm (outside poll in
2124                    // the cascade, so it runs AFTER the poll arm below) to the
2125                    // poll dispenser: the op's compiled gauge slots, re-published
2126                    // per poll iteration so drains feed the store live.
2127                    let poll_iteration_gauges: Arc<
2128                        arc_swap::ArcSwapOption<Vec<crate::wrappers::metrics::MetricSlot>>,
2129                    > = Arc::new(arc_swap::ArcSwapOption::empty());
2130
2131                    // Apply each remaining wrapper in plan order.
2132                    // Skip `traverse`; it's already constructed.
2133                    for reg in plan.iter_innermost_first() {
2134                        if reg.name == crate::wrappers::traverse::NAME {
2135                            continue;
2136                        }
2137                        let stop = match reg.name {
2138                            crate::wrappers::delay::NAME => {
2139                                let spec = template
2140                                    .delay
2141                                    .as_ref()
2142                                    .expect("delay triggered → delay set");
2143                                let trim = |s: &str| -> String {
2144                                    let t = s.trim();
2145                                    t.strip_prefix('{')
2146                                        .and_then(|s| s.strip_suffix('}'))
2147                                        .unwrap_or(t)
2148                                        .to_string()
2149                                };
2150                                let wrap_result = match spec {
2151                                    nmbrs_workload::model::DelaySpec::Before(name) => {
2152                                        let name = trim(name);
2153                                        crate::wrappers::DelayDispenser::wrap(
2154                                            current.clone(),
2155                                            &name,
2156                                            &mut fx,
2157                                        )
2158                                    }
2159                                    nmbrs_workload::model::DelaySpec::BeforeAfter {
2160                                        before,
2161                                        after,
2162                                    } => {
2163                                        let before = before.as_deref().map(trim);
2164                                        let after = after.as_deref().map(trim);
2165                                        crate::wrappers::DelayDispenser::wrap_before_after(
2166                                            current.clone(),
2167                                            before.as_deref(),
2168                                            after.as_deref(),
2169                                            &mut fx,
2170                                        )
2171                                    }
2172                                };
2173                                match wrap_result {
2174                                    Ok(d) => {
2175                                        current = d;
2176                                        false
2177                                    }
2178                                    Err(e) => {
2179                                        crate::diag!(
2180                                            crate::observer::LogLevel::Error,
2181                                            "error: op '{}': {e}",
2182                                            template.name
2183                                        );
2184                                        true
2185                                    }
2186                                }
2187                            }
2188                            crate::validation::WRAPPER_NAME => {
2189                                match crate::validation::ValidatingDispenser::wrap(
2190                                    current.clone(),
2191                                    template,
2192                                    &activity.labels,
2193                                    Some(&program),
2194                                    &mut fx,
2195                                ) {
2196                                    Ok((d, vm)) => {
2197                                        if let Some(vm) = vm {
2198                                            validation_metrics.push(vm);
2199                                        }
2200                                        current = d;
2201                                        false
2202                                    }
2203                                    Err(e) => {
2204                                        crate::diag!(
2205                                            crate::observer::LogLevel::Error,
2206                                            "error: op '{}': {e}",
2207                                            template.name
2208                                        );
2209                                        true
2210                                    }
2211                                }
2212                            }
2213                            crate::wrappers::poll::NAME => {
2214                                // Poll config reader: `poll:` is either
2215                                // a string (mode only, all-defaults) or
2216                                // a map (`{mode, interval_ms, timeout_ms,
2217                                // max_rows, min_rows, json_path,
2218                                // metric_name, max_error_retries}`). All
2219                                // poll knobs live UNDER `poll:` — no
2220                                // flat `poll_*` prefix keys at op
2221                                // level, so the wrapper's namespace
2222                                // doesn't collide with adapter fields
2223                                // (e.g. HTTP's `request_timeout_ms`).
2224                                let poll_val = template.params.get("poll");
2225                                let cfg = poll_val.and_then(|v| v.as_object());
2226                                let get_u64 = |k: &str, default: u64| -> u64 {
2227                                    cfg.and_then(|m| m.get(k))
2228                                        .and_then(|v| {
2229                                            v.as_u64().or_else(|| {
2230                                                v.as_str().and_then(|s| s.parse::<u64>().ok())
2231                                            })
2232                                        })
2233                                        .unwrap_or(default)
2234                                };
2235                                let get_u32 = |k: &str, default: u32| -> u32 {
2236                                    cfg.and_then(|m| m.get(k))
2237                                        .and_then(|v| {
2238                                            v.as_u64().map(|n| n as u32).or_else(|| {
2239                                                v.as_str().and_then(|s| s.parse::<u32>().ok())
2240                                            })
2241                                        })
2242                                        .unwrap_or(default)
2243                                };
2244                                let get_str = |k: &str| -> Option<String> {
2245                                    cfg.and_then(|m| m.get(k))
2246                                        .and_then(|v| v.as_str())
2247                                        .map(|s| s.to_string())
2248                                };
2249                                let interval = get_u64("interval_ms", 1000);
2250                                let timeout = get_u64("timeout_ms", 300_000);
2251                                // SRD-03 §"Status-Determination
2252                                // Invariant — Retries Within": bounded
2253                                // retry budget for retryable inner
2254                                // errors. Default 0 (strict). Operators
2255                                // raise this when long fixture readiness
2256                                // checks tolerate transient blips.
2257                                let max_error_retries = get_u32("max_error_retries", 0);
2258                                let metric_name = get_str("metric_name");
2259                                // `min_rows` / `max_rows` — completion
2260                                // window: poll considered done when
2261                                // row count is in the closed interval
2262                                // `[min..=max]`. Defaults `min=0, max=0`
2263                                // reproduce `await_empty` (exactly 0 =
2264                                // done). For "settled to N rows" cases
2265                                // (e.g. SAI's `sai_sstable_count == 1`
2266                                // after memtable flush + compaction)
2267                                // use `min=1, max=1`.
2268                                let min_rows = get_u64("min_rows", 0);
2269                                let max_rows = get_u64("max_rows", 0);
2270                                // `json_path` — optional JSON Pointer
2271                                // (RFC 6901, e.g. `/value`) drilled
2272                                // into the body before counting. Lets
2273                                // the count check address a nested
2274                                // field — load-bearing for envelope
2275                                // responses like Jolokia's
2276                                // `{value, status, …}`.
2277                                let json_path = get_str("json_path");
2278                                // `memo` / `progress` — live-status templates
2279                                // re-rendered per poll iteration against the
2280                                // wires: `memo` publishes the measured values
2281                                // to the activity memo (default: append
2282                                // `measured N row(s)…` to the base memo);
2283                                // `progress` parses to an f64 fraction that
2284                                // drives the phase completion bar via the
2285                                // derived-progress override.
2286                                let each_memo = get_str("memo");
2287                                let progress_template = get_str("progress");
2288                                // The op's `gutter:` DURING form rides into the poll
2289                                // loop for per-iteration refresh — the GutterDispenser
2290                                // itself publishes only once, when the (potentially
2291                                // hours-long) drain op completes.
2292                                // `on_done` — values written to the wires on the
2293                                // TERMINATING poll, before its metrics publish.
2294                                // A remote view of in-flight work (a compactions
2295                                // view, a job queue) shows what is RUNNING, so
2296                                // "done" arrives as an empty result and the
2297                                // finished item's attributes are already gone.
2298                                // This proxies the completed measurement we could
2299                                // not observe: `on_done: {completion_ratio: 1.0}`
2300                                // makes the final sample say what the view no
2301                                // longer can.
2302                                let on_done = crate::wrappers::poll::parse_on_done(cfg);
2303                                // `until:` — a polydat predicate replacing the
2304                                // row-count window. Scope synthesis lowered the
2305                                // expression into this node's kernel; the
2306                                // wrapper only needs to know that it exists.
2307                                let until = cfg
2308                                    .and_then(|m| m.get("until"))
2309                                    .and_then(|v| v.as_str())
2310                                    .is_some_and(|s| !s.trim().is_empty());
2311                                let (each_gutter, _) = crate::wrappers::gutter::parse_specs(
2312                                    template.params.get("gutter"),
2313                                );
2314                                let (d, _pm) = crate::wrappers::PollingDispenser::wrap_with_status(
2315                                    current.clone(),
2316                                    interval,
2317                                    timeout,
2318                                    max_error_retries,
2319                                    metric_name,
2320                                    min_rows,
2321                                    max_rows,
2322                                    json_path.clone(),
2323                                    each_memo,
2324                                    Some(activity.memo.clone()),
2325                                    progress_template,
2326                                    Some(activity.metrics.clone()),
2327                                    each_gutter,
2328                                    Some(activity.gutter.clone()),
2329                                    Some(poll_iteration_gauges.clone()),
2330                                    on_done.clone(),
2331                                    until,
2332                                    activity.stop_view(),
2333                                );
2334                                crate::diag!(
2335                                    crate::observer::LogLevel::Debug,
2336                                    "  op '{}': polling enabled (interval={}ms, timeout={}ms, max_error_retries={}, done_on={}, json_path={:?}, on_done={:?})",
2337                                    template.name,
2338                                    interval,
2339                                    timeout,
2340                                    max_error_retries,
2341                                    if until {
2342                                        "until-predicate".to_string()
2343                                    } else {
2344                                        format!("rows=[{min_rows}..={max_rows}]")
2345                                    },
2346                                    json_path,
2347                                    on_done
2348                                        .iter()
2349                                        .map(|(k, v)| format!("{k}={v}"))
2350                                        .collect::<Vec<_>>()
2351                                );
2352                                current = d;
2353                                false
2354                            }
2355                            crate::wrappers::r#if::NAME => {
2356                                // `if:` short-circuits before the
2357                                // inner cascade — load-bearing for
2358                                // the recent fix that pulls polling
2359                                // inside `if`. Resolver order
2360                                // mirrors that.
2361                                let cond = template
2362                                    .condition
2363                                    .as_deref()
2364                                    .expect("if triggered → condition set");
2365                                let cond_name = cond
2366                                    .trim()
2367                                    .strip_prefix('{')
2368                                    .and_then(|s| s.strip_suffix('}'))
2369                                    .unwrap_or(cond.trim());
2370                                match crate::wrappers::ConditionalDispenser::wrap(
2371                                    current.clone(),
2372                                    cond_name,
2373                                    activity.metrics.clone(),
2374                                    &mut fx,
2375                                ) {
2376                                    Ok(d) => {
2377                                        current = d;
2378                                        false
2379                                    }
2380                                    Err(e) => {
2381                                        crate::diag!(
2382                                            crate::observer::LogLevel::Error,
2383                                            "error: op '{}': {e}",
2384                                            template.name
2385                                        );
2386                                        true
2387                                    }
2388                                }
2389                            }
2390                            crate::wrappers::r#while::NAME => {
2391                                // `while:` loops the inner until the
2392                                // synthesised `__while` predicate
2393                                // flips falsy or the activity stops.
2394                                // The op-kernel synthesiser appended
2395                                // `__while := <expr>` to the kernel's
2396                                // result bindings — that's how the
2397                                // expression's free identifiers got
2398                                // their extern slots.
2399                                match crate::wrappers::WhileWrapper::wrap(
2400                                    current.clone(),
2401                                    activity.stop_view(),
2402                                    &mut fx,
2403                                ) {
2404                                    Ok(d) => {
2405                                        current = d;
2406                                        false
2407                                    }
2408                                    Err(e) => {
2409                                        crate::diag!(
2410                                            crate::observer::LogLevel::Error,
2411                                            "error: op '{}': {e}",
2412                                            template.name
2413                                        );
2414                                        true
2415                                    }
2416                                }
2417                            }
2418                            crate::wrappers::rate::NAME => {
2419                                // Per-op rate limiter, independent
2420                                // of the activity-level rate AND of
2421                                // every other op's per-op limiter.
2422                                // Each instance owns its own
2423                                // RateLimiter.
2424                                let rate_spec =
2425                                    template.rate.as_deref().expect("rate triggered → rate set");
2426                                match crate::wrappers::OpRateWrapper::wrap(
2427                                    current.clone(),
2428                                    rate_spec,
2429                                ) {
2430                                    Ok(d) => {
2431                                        current = d;
2432                                        false
2433                                    }
2434                                    Err(e) => {
2435                                        crate::diag!(
2436                                            crate::observer::LogLevel::Error,
2437                                            "error: op '{}': {e}",
2438                                            template.name
2439                                        );
2440                                        true
2441                                    }
2442                                }
2443                            }
2444                            crate::wrappers::fields::NAME => {
2445                                // Capture op_fields so the fields
2446                                // wrapper can render the rendered op
2447                                // text at cycle time. Stable
2448                                // insertion order isn't guaranteed by
2449                                // HashMap, but ParsedOp.op is small
2450                                // enough that a deterministic
2451                                // alphabetical sort keeps the printed
2452                                // output stable across runs.
2453                                let mut op_fields: Vec<(String, serde_json::Value)> = template
2454                                    .op
2455                                    .iter()
2456                                    .map(|(k, v)| (k.clone(), v.clone()))
2457                                    .collect();
2458                                op_fields.sort_by(|a, b| a.0.cmp(&b.0));
2459                                current = crate::wrappers::FieldsDispenser::wrap_with_op_fields(
2460                                    current.clone(),
2461                                    &template.name,
2462                                    op_fields,
2463                                );
2464                                false
2465                            }
2466                            crate::wrappers::result::NAME => {
2467                                // SRD-40b §5: result-as-GK adapter —
2468                                // exposes captured result fields to
2469                                // the op's Polydat scope via
2470                                // `OpResult.captures` so metric
2471                                // expressions (and any later wrappers)
2472                                // can reference them by name. No-op
2473                                // when the op declares no `result:`
2474                                // wires.
2475                                current = crate::wrappers::ResultDispenser::wrap(
2476                                    current.clone(),
2477                                    template.result.as_ref(),
2478                                    template.abstract_interface.as_ref().map(|i| &i.results),
2479                                );
2480                                false
2481                            }
2482                            crate::wrappers::metrics::NAME => {
2483                                // SRD-40b §6/§7 — one `Component`
2484                                // per dispenser carrying
2485                                // `op=<template.name>` so the
2486                                // duplicate-family check sees one
2487                                // dimensional cell per dispenser
2488                                // and child ops collide cleanly on
2489                                // their `op=` label.
2490                                let labels =
2491                                    nmbrs_metrics::labels::Labels::of("op", &template.name);
2492                                let dispenser_component =
2493                                    std::sync::Arc::new(std::sync::RwLock::new(
2494                                        nmbrs_metrics::component::Component::new(
2495                                            labels,
2496                                            std::collections::HashMap::new(),
2497                                        ),
2498                                    ));
2499                                if let Some(parent) = activity.component.as_ref() {
2500                                    nmbrs_metrics::component::attach(parent, &dispenser_component);
2501                                }
2502                                let wrap_result = {
2503                                    let mut guard = dispenser_component
2504                                        .write()
2505                                        .unwrap_or_else(|e| e.into_inner());
2506                                    crate::wrappers::MetricsDispenser::wrap_with_slots(
2507                                        current.clone(),
2508                                        &template.metrics,
2509                                        &mut guard,
2510                                        &dispenser_component,
2511                                        &mut fx,
2512                                    )
2513                                };
2514                                match wrap_result {
2515                                    Ok((d, slots)) => {
2516                                        if let Some(slots) = slots {
2517                                            poll_iteration_gauges.store(Some(slots));
2518                                        }
2519                                        // Mark the dispenser
2520                                        // component Running so the
2521                                        // cadence reporter's
2522                                        // `capture_tree` walk visits
2523                                        // it on every tick.
2524                                        dispenser_component
2525                                            .write()
2526                                            .unwrap_or_else(|e| e.into_inner())
2527                                            .set_state(
2528                                                nmbrs_metrics::component::ComponentState::Running,
2529                                            );
2530                                        dispenser_components.push(dispenser_component);
2531                                        current = d;
2532                                        false
2533                                    }
2534                                    Err(e) => {
2535                                        crate::diag!(
2536                                            crate::observer::LogLevel::Error,
2537                                            "error: op '{}': {e}",
2538                                            template.name
2539                                        );
2540                                        true
2541                                    }
2542                                }
2543                            }
2544                            crate::wrappers::dryrun::NAME => {
2545                                // Outermost short-circuit. Activated
2546                                // by the injected `dryrun:` template
2547                                // parameter (per
2548                                // `run_with_adapters`'s session-
2549                                // startup injection step). The
2550                                // trigger fires only when the
2551                                // template carries the marker, so
2552                                // we know we're in dryrun mode just
2553                                // by being here.
2554                                //
2555                                // Architectural invariant: by sitting
2556                                // outermost in the cascade, the
2557                                // short-circuit returns BEFORE any
2558                                // inner wrapper (verify / metrics /
2559                                // poll / etc.) can observe the
2560                                // empty body the dryrun stand-in
2561                                // produces. The `forbids_outer` set
2562                                // on the DRYRUN registration pins
2563                                // this position structurally.
2564                                // DryRunWrapper has one job: short-circuit
2565                                // the outbound op. No display modes, no
2566                                // op-field snapshot, no extra work — the
2567                                // wrapper is supposed to do NOTHING MORE
2568                                // than wrap the op and not call it.
2569                                current = crate::wrappers::DryRunWrapper::wrap(current.clone());
2570                                false
2571                            }
2572                            crate::wrappers::memo::NAME => {
2573                                // Memo wrapper: parse `memo:` (string
2574                                // shorthand or `{before, after}` map),
2575                                // wrap with cloned ArcSwap handle.
2576                                let (before, after) = match template.params.get("memo") {
2577                                    Some(serde_json::Value::String(s)) => {
2578                                        // Shorthand: same template for
2579                                        // before AND after.
2580                                        (Some(s.clone()), Some(s.clone()))
2581                                    }
2582                                    Some(serde_json::Value::Object(obj)) => {
2583                                        let b = obj
2584                                            .get("before")
2585                                            .and_then(|v| v.as_str())
2586                                            .map(String::from);
2587                                        let a = obj
2588                                            .get("after")
2589                                            .and_then(|v| v.as_str())
2590                                            .map(String::from);
2591                                        (b, a)
2592                                    }
2593                                    _ => (None, None),
2594                                };
2595                                if before.is_none() && after.is_none() {
2596                                    crate::diag!(
2597                                        crate::observer::LogLevel::Warn,
2598                                        "op '{}': memo: requires at least one of \
2599                                         `before` / `after` (or a string shorthand)",
2600                                        template.name
2601                                    );
2602                                    false
2603                                } else {
2604                                    current = crate::wrappers::MemoDispenser::wrap(
2605                                        current.clone(),
2606                                        before,
2607                                        after,
2608                                        activity.memo.clone(),
2609                                    );
2610                                    false
2611                                }
2612                            }
2613                            crate::wrappers::gutter::NAME => {
2614                                // Gutter wrapper: parse `gutter:` (layout-
2615                                // string shorthand or `{bar}` / `{spark}`
2616                                // map), wrap with the activity's shared
2617                                // gutter slot.
2618                                let (during, fin) = crate::wrappers::gutter::parse_specs(
2619                                    template.params.get("gutter"),
2620                                );
2621                                if during.is_none() && fin.is_none() {
2622                                    crate::diag!(
2623                                        crate::observer::LogLevel::Warn,
2624                                        "op '{}': gutter: requires a layout string, one of \
2625                                         `bar:` / `spark:` / `text:`, or a `final:` form",
2626                                        template.name
2627                                    );
2628                                    false
2629                                } else {
2630                                    if let Some((kind, tmpl)) = &during {
2631                                        *activity.gutter_spec.lock().unwrap() =
2632                                            Some((*kind, tmpl.clone()));
2633                                        current = crate::wrappers::GutterDispenser::wrap(
2634                                            current.clone(),
2635                                            *kind,
2636                                            tmpl.clone(),
2637                                            activity.gutter.clone(),
2638                                        );
2639                                    }
2640                                    if let Some((kind, tmpl)) = &fin {
2641                                        *activity.gutter_final_spec.lock().unwrap() =
2642                                            Some((*kind, tmpl.clone()));
2643                                    }
2644                                    false
2645                                }
2646                            }
2647                            crate::wrappers::readout::NAME => {
2648                                // Readout wrapper: op-level status leaf. Wrap the
2649                                // inner op with a lifecycle reporter keyed by the
2650                                // op-template name; the parent phase is resolved at
2651                                // run time from the task-local CURRENT_PHASE. Only
2652                                // `readout: visible` triggers it (see the wrapper's
2653                                // `triggers`), so this arm always wraps when reached.
2654                                current = crate::wrappers::ReadoutDispenser::wrap_with_measure(
2655                                    current.clone(),
2656                                    template.name.clone(),
2657                                    template
2658                                        .params
2659                                        .get("measure")
2660                                        .and_then(|v| v.as_str())
2661                                        .map(|s| s.to_string()),
2662                                );
2663                                crate::diag!(
2664                                    crate::observer::LogLevel::Debug,
2665                                    "  op '{}': readout visible (op-level status line)",
2666                                    template.name
2667                                );
2668                                false
2669                            }
2670                            crate::wrappers::errors::NAME => {
2671                                // SRD-82 Part 3b — hand-placed OUTERMOST after
2672                                // this loop (mirrors the hand-placed innermost
2673                                // tries). The plan entry records presence for
2674                                // telemetry / describe only.
2675                                false
2676                            }
2677                            crate::wrappers::tries::NAME => {
2678                                // SRD-82 Part 3b — hand-placed INNERMOST before
2679                                // this loop (the sigil resolution above). The
2680                                // plan entry records presence for telemetry /
2681                                // describe only.
2682                                false
2683                            }
2684                            other => {
2685                                crate::diag!(
2686                                    crate::observer::LogLevel::Error,
2687                                    "error: op '{}': resolver returned wrapper `{}` \
2688                                     with no dispatch handler in the cascade",
2689                                    template.name,
2690                                    other
2691                                );
2692                                true
2693                            }
2694                        };
2695                        if stop {
2696                            return true;
2697                        }
2698                    }
2699
2700                    // Dryrun short-circuit: when the session is in
2701                    // dryrun mode (`config.dry_run_mode` set), the
2702                    // dryrun template-parameter injection above put
2703                    // a `dryrun:` field on every op template; the
2704                    // wrapper resolver picks up that field and adds
2705                    // `DryRunWrapper` as the OUTERMOST layer. The
2706                    // wrapper never calls its inner — verify /
2707                    // metrics / poll / etc. observers don't fire,
2708                    // and the real adapter's `execute()` is
2709                    // suppressed. The real adapter itself still
2710                    // constructs in full (connect, prepare, gather
2711                    // metadata); only the per-cycle outbound call
2712                    // is short-circuited.
2713                    // SRD-82 Part 3b — the error handler is the OUTERMOST
2714                    // wrapper, hand-placed after the plan cascade (mirroring
2715                    // the hand-placed innermost retry wrapper), driven by the
2716                    // policy resolved BEFORE the retry check above. Only the
2717                    // op-error ROUTER is per-op — the aggregate rate breach
2718                    // stays the phase shell's `error_policy.guard`. The
2719                    // wrapper observes the stack's ONE terminal outcome per
2720                    // cycle: routes it, tallies result-level error counters,
2721                    // captures phase errors, and applies stop/fail effects.
2722                    // Its happy path is a single branch. When no retry
2723                    // wrapper is present it also records the single-attempt
2724                    // `attempt_*` tallies (`records_attempts`).
2725                    let current = crate::wrappers::ErrorHandlerDispenser::wrap(
2726                        current,
2727                        op_error_policy,
2728                        activity.metrics.clone(),
2729                        activity.phase_errors.clone(),
2730                        activity.stop_flag.clone(),
2731                        activity.stop_reason.clone(),
2732                        template.name.clone(),
2733                        /* records_attempts */ !has_tries_wrapper,
2734                    );
2735                    dispensers.push(current);
2736
2737                    // Seal the per-template fixture. The PullPlan
2738                    // drives cycle-time reads for every wrapper that
2739                    // registered (validation ground truth, conditional
2740                    // `if`, throttle `delay`). See SRD 31 §"Pull plan
2741                    // vs bind plan".
2742                    pull_plans_per_template.push(fx.seal());
2743                }
2744                Err(e) => {
2745                    crate::diag!(
2746                        crate::observer::LogLevel::Error,
2747                        "error: adapter.map_op failed for '{}': {e}",
2748                        template.name
2749                    );
2750                    return true;
2751                }
2752            }
2753        }
2754        let dispensers = Arc::new(dispensers);
2755        // Register dispensers for adapter-specific metrics capture
2756        activity.metrics.set_dispensers(dispensers.clone());
2757        let pull_plans_per_template = Arc::new(pull_plans_per_template);
2758
2759        // `dryrun=dispenser` exit point. Every op template's
2760        // dispenser is constructed (map_op succeeded, wrapper
2761        // plan resolved, pull plan sealed). Nothing else needs
2762        // to happen for the operator to know the construction
2763        // pipeline is healthy — return cleanly without spawning
2764        // the fiber pool or the progress thread.
2765        if activity.config.stop_after_dispenser_init {
2766            crate::diag!(
2767                crate::observer::LogLevel::Info,
2768                "dryrun=dispenser: {} op-template dispenser(s) constructed; \
2769                 stopping before cycle execution",
2770                dispensers.len()
2771            );
2772            return false;
2773        }
2774
2775        let validation_metrics = Arc::new(validation_metrics);
2776        // Share the validation-metrics handle with ActivityMetrics so
2777        // the progress thread (below) can read live relevancy aggregates.
2778        activity
2779            .metrics
2780            .set_validation_metrics(validation_metrics.clone());
2781
2782        // Single activity-level rate limiter. One ops-per-sec
2783        // ceiling gates every fiber; there is no separate
2784        // stanza-rate mechanism. Activities with no `rate`
2785        // configured skip construction cleanly.
2786        let rate_limiter = activity
2787            .config
2788            .rate
2789            .map(|r| Arc::new(RateLimiter::start(nmbrs_rate::RateSpec::new(r))));
2790
2791        // Register the [`RateLimiterApplier`] against the
2792        // already-declared `rate` control if both the control
2793        // and the limiter exist. The declaration happens in
2794        // [`Self::attach_component`] — this step only wires the
2795        // applier so a runtime write actually reconfigures the
2796        // running limiter.
2797        if let (Some(ac), Some(rl)) = (activity.component.as_ref(), rate_limiter.as_ref()) {
2798            let existing: Option<nmbrs_metrics::controls::Control<nmbrs_rate::RateSpec>> = ac
2799                .read()
2800                .unwrap_or_else(|e| e.into_inner())
2801                .controls()
2802                .get("rate");
2803            if let Some(ctl) = existing {
2804                ctl.register_applier(nmbrs_rate::RateLimiterApplier::new(Arc::clone(rl)));
2805            }
2806        }
2807
2808        // SRD-100 P2 — the inline-status refresh thread is RETIRED. The
2809        // live phase status is now folded at the display consumer from the
2810        // `active_phases` snapshot (the executor attaches each phase's
2811        // render handle on-task; see `RunObserver::phase_render_attach` and
2812        // `nmbrs_tui::status_fold`). This removes the last-writer race on the
2813        // single status slot, the `std::thread` cross-execution hazard (the
2814        // thread couldn't read the SRD-88 task-local channel), and the
2815        // per-phase-clear-wipes-peers bug. The phase-scoped values below
2816        // still feed the phase-END readout + outcome line.
2817        let activity_name = activity.config.name.clone();
2818        let suppress_progress = adapters.values().any(|a| a.name() == "plotter");
2819        let start_time = Instant::now();
2820        // Use source extent for progress (data-driven), not cycles
2821        let source_for_progress = activity.source_factory.clone();
2822        let total_extent = source_for_progress
2823            .global_extent()
2824            .unwrap_or(activity.config.cycles);
2825        // One Arc<str> shared by every fiber in this phase. The
2826        // Polydat runtime-context `phase()` node clones this per read
2827        // instead of per fiber, keeping the per-cycle cost O(1).
2828        let phase_name_arc: Arc<str> = Arc::from(activity_name.as_str());
2829
2830        // Daemon-op pool. Shared across cycle-pool fibers — each
2831        // fiber's stanza walk dispatches daemon ops by spawning
2832        // a fresh fiber onto this pool instead of running them
2833        // inline. The pool enforces per-op-name fiber caps; an
2834        // overflow is a workload-design error that fails the
2835        // phase. At phase exit the cycle-pool drain runs first,
2836        // then this pool's shutdown signals and waits on every
2837        // still-running daemon (see daemon-pool drain below).
2838        let daemon_pool = Arc::new(crate::daemon_pool::DaemonPool::new());
2839
2840        // SRD 23 §"Fiber executor": fiber lifecycle goes through
2841        // a [`FiberPool`] that the `ConcurrencyApplier` can
2842        // resize via the activity's `concurrency` control. Each
2843        // fiber receives its own stop-flag and exits
2844        // cooperatively at the next cycle boundary when flagged.
2845        let pool_spawner: crate::fiber_pool::FiberSpawner = {
2846            let activity = activity.clone();
2847            let dispensers_outer = dispensers.clone();
2848            let pull_plans_outer = pull_plans_per_template.clone();
2849            let op_builder_outer = op_builder.clone();
2850            let rate_limiter_outer = rate_limiter.clone();
2851            let phase_arc_outer = phase_name_arc.clone();
2852            let daemon_pool_outer = daemon_pool.clone();
2853            // SRD-89 — snapshot this phase's controls ONCE (walk-up from the
2854            // phase component), shared lock-free across all of the phase's
2855            // fibers. Carries live handles, so servo retargets are observed.
2856            let phase_controls_outer = activity
2857                .component
2858                .as_ref()
2859                .map(crate::polydat_nodes::runtime_context::snapshot_controls)
2860                .unwrap_or_else(crate::polydat_nodes::runtime_context::empty_controls);
2861            Box::new(move |stop: crate::fiber_pool::StopFlag| {
2862                let activity = activity.clone();
2863                let dispensers = dispensers_outer.clone();
2864                let pull_plans = pull_plans_outer.clone();
2865                let op_builder = op_builder_outer.clone();
2866                let rate_limiter = rate_limiter_outer.clone();
2867                let phase_arc = phase_arc_outer.clone();
2868                let daemon_pool = daemon_pool_outer.clone();
2869                let phase_controls = phase_controls_outer.clone();
2870                // SRD-88 — carry the per-execution context into the per-cycle
2871                // fiber: the adapter's op-output, log, stop, and exec-identity
2872                // resolve to THIS execution (so concurrent executions sharing a
2873                // session capture their own output / route their own log). A
2874                // no-op on the single-run path (no context scoped — A1).
2875                tokio::spawn(crate::execution_context::propagate(async move {
2876                    // Catch panics inside the fiber so they surface
2877                    // in diagnostics rather than silently terminating
2878                    // the task. Without this, a panic in any cycle's
2879                    // accessor / binder code would leave the fiber
2880                    // gone and the run still "active" from the
2881                    // perspective of the executor, hanging the TUI
2882                    // with no visible cause. The session log line
2883                    // captures location + message; the runtime's
2884                    // own panic reporting (if any) is unchanged.
2885                    use futures::FutureExt as _;
2886                    let activity_for_panic = activity.clone();
2887                    let activity_name_for_log = activity.config.name.clone();
2888                    let phase_arc_for_exec = phase_arc.clone();
2889                    let body = crate::polydat_nodes::runtime_context::with_fiber_context(
2890                        phase_arc,
2891                        phase_controls,
2892                        async move {
2893                            executor_task(
2894                                activity,
2895                                dispensers,
2896                                pull_plans,
2897                                op_builder,
2898                                rate_limiter,
2899                                stop,
2900                                daemon_pool,
2901                                phase_arc_for_exec,
2902                            )
2903                            .await;
2904                        },
2905                    );
2906                    let result = std::panic::AssertUnwindSafe(body).catch_unwind().await;
2907                    match result {
2908                        Ok(()) => {
2909                            // Normal fiber exit is silent — the
2910                            // session log used to record one
2911                            // line per fiber here at Debug, but
2912                            // with concurrency=N, that's N lines
2913                            // per phase boundary in session.log
2914                            // for no diagnostic value (the phase
2915                            // completion + duration already
2916                            // tells the user the fibers
2917                            // completed). Panic exits below
2918                            // remain Error-level.
2919                            let _ = activity_name_for_log;
2920                        }
2921                        Err(panic_payload) => {
2922                            let msg = panic_payload
2923                                .downcast_ref::<&'static str>()
2924                                .map(|s| (*s).to_string())
2925                                .or_else(|| panic_payload.downcast_ref::<String>().cloned())
2926                                .unwrap_or_else(|| "<non-string panic payload>".into());
2927                            crate::diag!(
2928                                crate::observer::LogLevel::Error,
2929                                "fiber panic in activity '{}': {}",
2930                                activity_name_for_log,
2931                                msg
2932                            );
2933                            // A panic is a first-class failure:
2934                            // give the run a headline cause
2935                            // (stop_reason) and a structured
2936                            // PhaseOutcome error, same as the
2937                            // stop-condition failure path.
2938                            // Without these the run dies with
2939                            // only a session-log line and no
2940                            // visible "why" at phase level.
2941                            if let Ok(mut slot) = activity_for_panic.stop_reason.lock()
2942                                && slot.is_none()
2943                            {
2944                                // Headline = first line; the full
2945                                // text lands in phase_errors below.
2946                                let first = msg.lines().next().unwrap_or(&msg);
2947                                *slot = Some(format!(
2948                                    "[panic] fiber panic in activity '{activity_name_for_log}': {first}"
2949                                ));
2950                            }
2951                            if let Ok(mut errs) = activity_for_panic.phase_errors.lock() {
2952                                errs.push(crate::phase_outcome::PhaseErrorDetail {
2953                                    class: "panic".into(),
2954                                    message: msg,
2955                                    op_name: None,
2956                                    cycle: None,
2957                                    op_template: None,
2958                                    op_resolved: None,
2959                                    at_nanos: std::time::SystemTime::now()
2960                                        .duration_since(std::time::UNIX_EPOCH)
2961                                        .map(|d| d.as_nanos() as u64)
2962                                        .unwrap_or(0),
2963                                    retryable: false,
2964                                });
2965                            }
2966                            // Mark stop_flag so other fibers and the
2967                            // executor's main loop see that something
2968                            // went wrong; the run will terminate at
2969                            // the next coordination point rather
2970                            // than continuing in a half-broken state.
2971                            activity_for_panic
2972                                .stop_flag
2973                                .store(true, std::sync::atomic::Ordering::Relaxed);
2974                        }
2975                    }
2976                }))
2977            })
2978        };
2979        let fiber_pool = Arc::new(crate::fiber_pool::FiberPool::new(pool_spawner));
2980
2981        // Register the pool's applier against the already-declared
2982        // `concurrency` control (see [`Self::attach_component`]).
2983        // At run time the applier is what turns a control write
2984        // into an actual fiber-pool resize. Without a component
2985        // attached (library-level tests that call `Activity::new`
2986        // directly) we skip registration — the pool still
2987        // operates, just without the runtime control surface.
2988        if let Some(ac) = activity.component.as_ref() {
2989            let existing: Option<nmbrs_metrics::controls::Control<u32>> = ac
2990                .read()
2991                .unwrap_or_else(|e| e.into_inner())
2992                .controls()
2993                .get("concurrency");
2994            if let Some(ctl) = existing {
2995                ctl.register_applier(crate::fiber_pool::ConcurrencyApplier::new(
2996                    fiber_pool.clone(),
2997                ));
2998            }
2999        }
3000
3001        // SRD-83 Part 9 — the adaptive backpressure governor, built
3002        // BEFORE the pool spawns so the phase OPENS at the slow-start
3003        // offer (default: the floor) instead of assaulting a fragile
3004        // target at the authored ceiling. Built after
3005        // `attach_component` declared the controls it walks.
3006        let mut throttle_governor = activity.config.throttle.as_ref().and_then(|spec| {
3007            crate::throttle::ThrottleGovernor::from_spec(
3008                spec,
3009                activity.component.as_ref(),
3010                &activity.config.name,
3011                activity.config.concurrency,
3012                activity.config.rate,
3013            )
3014        });
3015        let initial_fibers = throttle_governor
3016            .as_ref()
3017            .and_then(|g| g.initial_concurrency())
3018            .unwrap_or(activity.config.concurrency);
3019        fiber_pool.spawn_initial(initial_fibers);
3020        // Wait for fibers to exit by natural exhaustion (source
3021        // drained) or `stop_flag` set by the error router.
3022        // Runtime resize-down flags some of them earlier; those
3023        // exit at the next cycle boundary and the remainder
3024        // drain when the source is done.
3025        let mut last_seen_count = activity.config.concurrency;
3026        let mut last_seen_cycles = activity.metrics.cycles_completed();
3027        let mut stuck_since = std::time::Instant::now();
3028        let mut last_logged_count = activity.config.concurrency;
3029        // SRD-83 — compile this phase's stop conditions (the default
3030        // `error_rate > error_rate_max` plus any declared `stop_when:`
3031        // predicates) as scope-bound `ScopedPredicate`s, evaluated per tick
3032        // below. Fire at most once per phase.
3033        //
3034        // The predicates bind to this phase node's OWN scope kernel
3035        // (`activity.phase_kernel`, the structural walk's cached kernel),
3036        // so they read the phase's wires as they sit — no conjured root.
3037        // (Distribution to other shell levels via `each:` is the broader
3038        // SRD-83 shell-evaluation follow-up; this binds the phase-level
3039        // conditions to their native phase scope.)
3040        let mut policy_tripped = false;
3041        let phase_start = std::time::Instant::now();
3042        let mut stop_conditions = match &activity.phase_kernel {
3043            Some(kernel) => crate::stop_conditions::StopConditionSet::build_for_phase(
3044                kernel,
3045                &activity.config.stop_when,
3046            )
3047            .unwrap_or_else(|e| {
3048                // A predicate that won't compile is the workload author's
3049                // bug; dryrun is where it should be rejected. At runtime,
3050                // log loudly and run with no stop conditions rather than
3051                // abort the phase on a synthesis error.
3052                crate::diag!(
3053                    crate::observer::LogLevel::Error,
3054                    "activity '{}': stop-condition compile failed: {e}",
3055                    activity.config.name
3056                );
3057                crate::stop_conditions::StopConditionSet::empty()
3058            }),
3059            None => crate::stop_conditions::StopConditionSet::empty(),
3060        };
3061        loop {
3062            fiber_pool.reap_finished();
3063            let n = fiber_pool.tracked_count();
3064            if n == 0 {
3065                break;
3066            }
3067            if let Some(governor) = throttle_governor.as_mut() {
3068                governor.tick(
3069                    activity.metrics.attempt_success.count(),
3070                    activity.metrics.attempt_failure.count(),
3071                );
3072            }
3073            // Periodic stall detection. A real stall means
3074            // *neither* signal of progress has moved:
3075            //   - `tracked_count` only changes when a fiber
3076            //     exits. During steady-state rampup every fiber
3077            //     is alive and busy, so this stays constant
3078            //     even when work is flying.
3079            //   - `cycles_completed` increments per finished op,
3080            //     so it reflects actual throughput regardless of
3081            //     whether any fiber has exited yet.
3082            // Either signal moving resets the stuck timer; only
3083            // when both are flat for the full 30 s do we warn.
3084            let cycles = activity.metrics.cycles_completed();
3085            // SRD-83 — evaluate the phase's stop conditions against a
3086            // fresh runtime-state snapshot (the Tick firing event). The
3087            // first predicate that trips stops the shell: fibers drain at
3088            // their next cycle boundary and the phase-end outcome becomes
3089            // Failed. (The per-condition `effect` → two-axis Outcome
3090            // mapping lands with SRD-82 Part 1; for now a trip is Failed.)
3091            if !policy_tripped && !stop_conditions.is_empty() {
3092                // Read `error_count` BEFORE `op_count`, then take a FRESH
3093                // `op_count`: every terminal error also increments
3094                // `cycles_completed`, so a cycles read taken at-or-after the
3095                // errors read is always ≥ it — guaranteeing `error_count ≤
3096                // op_count` and thus `error_rate ≤ 1.0`. The reverse order (the
3097                // stuck-timer's earlier `cycles` snapshot for `op_count`, then a
3098                // later `errors_total` read) let an erroring op completing
3099                // between the two reads push `error_count` past the stale
3100                // `op_count`, momentarily yielding `error_rate > 1.0` and
3101                // SPURIOUSLY tripping the `error_rate > 1.0` guard under
3102                // saturation (where the true rate sits exactly at 1.0).
3103                // SRD-91: the stop-condition error rate is a per-OP
3104                // proportion in [0,1], so it reads the per-op terminal
3105                // failure count (`result_failure`), not the per-attempt
3106                // `errors_total` (which can exceed op_count under retries).
3107                let result_failure = activity.metrics.result_failure.count();
3108                let cycles_total = activity.metrics.cycles_completed();
3109                // Attempt-level tallies (resolved attempts only, per the
3110                // SRD-91 counters): the see-through-retries wires. Each
3111                // wire carries its instrument's name and raw count; any
3112                // rate is derived in the predicate text.
3113                let attempt_success = activity.metrics.attempt_success.count();
3114                let attempt_failure = activity.metrics.attempt_failure.count();
3115                let state = crate::stop_conditions::RuntimeState {
3116                    cycles_total,
3117                    result_failure,
3118                    elapsed_ms: phase_start.elapsed().as_millis() as u64,
3119                    // attempt_total counts at RESOLUTION (tries.rs), so
3120                    // the invariant total == success + failure holds;
3121                    // reading the two parts keeps one consistent view.
3122                    attempt_total: attempt_success + attempt_failure,
3123                    attempt_success,
3124                    attempt_failure,
3125                    ..Default::default()
3126                };
3127                if let Some((outcome, reason, target, cancel_ops)) =
3128                    stop_conditions.evaluate(&state)
3129                {
3130                    policy_tripped = true;
3131                    // SRD-83 Part 5 — honour the condition's effect. A
3132                    // `fail` effect records a phase error (the phase ends
3133                    // Failed); a `stop` effect is a clean halt (no error,
3134                    // the phase ends Completed). Either way the stop_flag
3135                    // drains the fibers at their next cycle boundary.
3136                    // Human-readable detail: the ACTUAL wire values that
3137                    // crossed the threshold (op_count / errors / error_rate /
3138                    // elapsed), so the failure says WHY, not just which
3139                    // predicate. The predicate itself rides the `[{reason}]`
3140                    // class prefix on the slot, so the message no longer
3141                    // repeats it (matches the `[class] message` convention
3142                    // used by the cycle-/daemon-error paths).
3143                    let actual = state.describe();
3144                    // Identity + attribution: the tripped message must say
3145                    // WHICH phase instance (labels carry the sweep cell /
3146                    // partition) and WHAT failed (top error types by
3147                    // count) — not just why the predicate fired.
3148                    let labels = &activity.config.phase_labels;
3149                    let where_part = if labels.is_empty() {
3150                        format!("phase '{}'", activity.config.name)
3151                    } else {
3152                        format!("phase '{}' ({labels})", activity.config.name)
3153                    };
3154                    let top_errs = activity.metrics.top_error_types(3);
3155                    let what_part = if top_errs.is_empty() {
3156                        String::new()
3157                    } else {
3158                        format!(" — top errors: {top_errs}")
3159                    };
3160                    if outcome.is_failure() {
3161                        let msg = format!(
3162                            "stop condition tripped in {where_part} — actual: {actual}{what_part} — failing phase"
3163                        );
3164                        crate::diag!(
3165                            crate::observer::LogLevel::Error,
3166                            "activity '{}': {reason} — {msg}",
3167                            activity.config.name
3168                        );
3169                        if let Ok(mut slot) = activity.stop_reason.lock()
3170                            && slot.is_none()
3171                        {
3172                            *slot = Some(format!("[{reason}] {msg}"));
3173                            // SRD-83 Part 5 — the first stopper owns the
3174                            // outcome as well as the reason: latch the
3175                            // condition's declared effect for the phase
3176                            // shell to adopt.
3177                            if let Ok(mut oc) = activity.stop_outcome.lock() {
3178                                *oc = Some(outcome.clone());
3179                            }
3180                        }
3181                        if let Ok(mut errs) = activity.phase_errors.lock() {
3182                            errs.push(crate::phase_outcome::PhaseErrorDetail {
3183                                class: reason,
3184                                message: msg,
3185                                op_name: None,
3186                                cycle: None,
3187                                op_template: None,
3188                                op_resolved: None,
3189                                at_nanos: std::time::SystemTime::now()
3190                                    .duration_since(std::time::UNIX_EPOCH)
3191                                    .map(|d| d.as_nanos() as u64)
3192                                    .unwrap_or(0),
3193                                retryable: false,
3194                            });
3195                        }
3196                    } else {
3197                        let msg = format!(
3198                            "stop condition tripped in {where_part} — actual: {actual}{what_part} — stopping phase"
3199                        );
3200                        crate::diag!(
3201                            crate::observer::LogLevel::Warn,
3202                            "activity '{}': {reason} — {msg}",
3203                            activity.config.name
3204                        );
3205                        if let Ok(mut slot) = activity.stop_reason.lock()
3206                            && slot.is_none()
3207                        {
3208                            *slot = Some(format!("[{reason}] {msg}"));
3209                            // SRD-83 Part 5 — a graceful `stop` effect:
3210                            // latch Interrupted+Succeeded so the phase
3211                            // shell ends this phase cleanly with its
3212                            // partial result instead of deriving failure
3213                            // from the bare stop flag.
3214                            if let Ok(mut oc) = activity.stop_outcome.lock() {
3215                                *oc = Some(outcome.clone());
3216                            }
3217                        }
3218                    }
3219                    // SRD-83 follow-up — route the action to its target scope.
3220                    // Detection happened here (phase); `target` (from `at:`,
3221                    // default = innermost of `per:`) says WHERE the stop lands.
3222                    // `Phase` halts just this phase (`stop_flag`); `Scenario`/
3223                    // `Workload` latch the workload `walk_stop` so the enclosing
3224                    // shell halts (which also drains this phase via
3225                    // `should_stop()`), leaving the session running. If no
3226                    // workload handle is wired (a standalone phase), fall back
3227                    // to the phase stop.
3228                    match target {
3229                        crate::stop_conditions::StopScope::Phase => {
3230                            activity.stop_flag.store(true, Ordering::Relaxed);
3231                        }
3232                        crate::stop_conditions::StopScope::Scenario
3233                        | crate::stop_conditions::StopScope::Workload => {
3234                            match &activity.walk_stop {
3235                                Some(walk) => walk.store(true, Ordering::Relaxed),
3236                                None => activity.stop_flag.store(true, Ordering::Relaxed),
3237                            }
3238                            // SRD-83 follow-up — a `Workload`-scope halt means
3239                            // "stop the whole run", exactly like a graceful
3240                            // (first) Ctrl-C. Latching `walk_stop` alone only
3241                            // retreats the scope walk lazily, so concurrent
3242                            // siblings (other scenarios, daemon probes, an
3243                            // already-dispatched phase) keep draining until the
3244                            // walk structurally unwinds them. Raising the global
3245                            // session stop makes every live fiber exit at its
3246                            // next cycle boundary — the same cooperative signal
3247                            // Ctrl-C level 1 raises. Cleanup is unchanged: the
3248                            // walk still returns normally and the runner's RAII
3249                            // shutdown guard runs the metrics/WAL/summary
3250                            // teardown. `Scenario` scope stays walk-local — it
3251                            // must halt only its own scenario, not the session.
3252                            // `abort` (cancel_ops) is handled below and
3253                            // supersedes this cooperative global stop.
3254                            if matches!(target, crate::stop_conditions::StopScope::Workload)
3255                                && !cancel_ops
3256                            {
3257                                crate::session_signals::request_stop();
3258                            }
3259                        }
3260                    }
3261                    // SRD-83 follow-up — `action: abort`. Beyond the
3262                    // cooperative halt above, jump STRAIGHT to the cancel-ops
3263                    // rung: `abort_shutdown()` raises the global session stop
3264                    // and drops in-flight op futures NOW — no cooperative
3265                    // drain, no 10s countdown. A stop driven by errors should
3266                    // not wait for doomed ops/phases to finish (a server sick
3267                    // enough to trip this will only ever end them by client
3268                    // timeout). Cleanup is still guaranteed: the walk unwinds
3269                    // normally into the runner's RAII shutdown guard, which
3270                    // runs the metrics/WAL/summary teardown (graceful SESSION
3271                    // shutdown). Not a force-exit — a further Ctrl-C is.
3272                    if cancel_ops {
3273                        crate::session_signals::abort_shutdown(
3274                            crate::session_signals::ShutdownOrigin::StopAction,
3275                        );
3276                    }
3277                }
3278            }
3279            let count_changed = n != last_seen_count;
3280            let cycles_changed = cycles != last_seen_cycles;
3281            if count_changed || cycles_changed {
3282                last_seen_count = n;
3283                last_seen_cycles = cycles;
3284                stuck_since = std::time::Instant::now();
3285                // Log progressing-but-slow drain: every time the
3286                // count changes we re-emit at debug so a stuck
3287                // run's session.log shows the slope (or lack of it)
3288                // without flooding when drain is fast.
3289                if count_changed
3290                    && (last_logged_count.saturating_sub(n) >= 10
3291                        || (n < 10 && n != last_logged_count))
3292                {
3293                    crate::diag!(
3294                        crate::observer::LogLevel::Debug,
3295                        "activity '{}': fiber drain at {n} (from {last_logged_count})",
3296                        activity.config.name
3297                    );
3298                    last_logged_count = n;
3299                }
3300            } else if stuck_since.elapsed() > std::time::Duration::from_secs(30) {
3301                // Distinguish "genuinely stuck" from "fiber is mid-op
3302                // on a long synchronous call." The latter case
3303                // (jolokia compaction, large schema migrations,
3304                // synchronous JMX exec) legitimately blocks one
3305                // fiber for many minutes without being a bug.
3306                // `ops_started > ops_finished` proves the fiber
3307                // is busy in the adapter — log at Debug so the
3308                // session.log timeline still records the slope,
3309                // but don't surface a Warn that operators read as
3310                // "something's wrong."
3311                let started = activity.metrics.ops_started.load(Ordering::Relaxed);
3312                let finished = activity.metrics.ops_finished.load(Ordering::Relaxed);
3313                let in_flight = started.saturating_sub(finished);
3314                if in_flight > 0 {
3315                    crate::diag!(
3316                        crate::observer::LogLevel::Debug,
3317                        "activity '{}': {n} fiber(s), {in_flight} op(s) in flight, \
3318                         {cycles} cycles completed, no fiber-count or cycle-count \
3319                         change for 30s (long-running op in progress)",
3320                        activity.config.name
3321                    );
3322                } else {
3323                    crate::diag!(
3324                        crate::observer::LogLevel::Warn,
3325                        "activity '{}': {n} fibers running, {cycles} cycles completed, \
3326                         no ops in flight, no progress for 30s — likely blocked on \
3327                         lock or IO",
3328                        activity.config.name
3329                    );
3330                }
3331                stuck_since = std::time::Instant::now();
3332            }
3333            tokio::time::sleep(Duration::from_millis(5)).await;
3334        }
3335        crate::diag!(
3336            crate::observer::LogLevel::Debug,
3337            "activity '{}': all fibers drained",
3338            activity.config.name
3339        );
3340
3341        // Daemon-pool drain. Cycle-pool reached zero (cursor
3342        // exhausted or stop signal honoured); now signal each
3343        // still-running daemon, wait its per-op grace window for
3344        // the in-flight future to drop, and aggregate outcomes.
3345        //
3346        // The outcomes split two ways: clean (Completed /
3347        // Cancelled) feed into the phase metrics counters;
3348        // unclean (Errored / TimedOut / Panicked) bubble up as
3349        // phase-stopping errors via the existing stop_flag +
3350        // stop_reason channel that cycle-pool errors use.
3351        if !daemon_pool.is_empty() {
3352            crate::diag!(
3353                crate::observer::LogLevel::Debug,
3354                "activity '{}': draining {} daemon(s)",
3355                activity.config.name,
3356                daemon_pool.len()
3357            );
3358            let outcomes = daemon_pool.shutdown().await;
3359            for (op_name, exit) in &outcomes {
3360                match exit {
3361                    crate::daemon_pool::DaemonExit::Completed => {
3362                        crate::diag!(
3363                            crate::observer::LogLevel::Debug,
3364                            "daemon op '{op_name}': completed"
3365                        );
3366                    }
3367                    crate::daemon_pool::DaemonExit::Cancelled => {
3368                        activity.metrics.daemon_cancelled_total.inc();
3369                        crate::diag!(
3370                            crate::observer::LogLevel::Debug,
3371                            "daemon op '{op_name}': cancelled at phase exit"
3372                        );
3373                    }
3374                    crate::daemon_pool::DaemonExit::Errored(e) => {
3375                        activity.metrics.daemon_errors_total.inc();
3376                        let inner = e.error();
3377                        crate::diag!(
3378                            crate::observer::LogLevel::Error,
3379                            "daemon op '{op_name}' errored: [{}] {}",
3380                            inner.error_name,
3381                            inner.message
3382                        );
3383                        if let Ok(mut slot) = activity.stop_reason.lock()
3384                            && slot.is_none()
3385                        {
3386                            *slot = Some(format!(
3387                                "[{}] daemon op '{op_name}': {}",
3388                                inner.error_name, inner.message,
3389                            ));
3390                        }
3391                    }
3392                    crate::daemon_pool::DaemonExit::TimedOut => {
3393                        activity.metrics.daemon_errors_total.inc();
3394                        crate::diag!(
3395                            crate::observer::LogLevel::Error,
3396                            "daemon op '{op_name}': did not acknowledge stop \
3397                             within grace window — phase fails"
3398                        );
3399                        if let Ok(mut slot) = activity.stop_reason.lock()
3400                            && slot.is_none()
3401                        {
3402                            *slot = Some(format!(
3403                                "[daemon_shutdown_timeout] daemon op \
3404                                 '{op_name}' did not acknowledge stop \
3405                                 within its grace window",
3406                            ));
3407                        }
3408                    }
3409                    crate::daemon_pool::DaemonExit::Panicked(msg) => {
3410                        activity.metrics.daemon_errors_total.inc();
3411                        crate::diag!(
3412                            crate::observer::LogLevel::Error,
3413                            "daemon op '{op_name}' panicked: {msg}"
3414                        );
3415                        if let Ok(mut slot) = activity.stop_reason.lock()
3416                            && slot.is_none()
3417                        {
3418                            *slot = Some(format!("[daemon_panic] daemon op '{op_name}': {msg}",));
3419                        }
3420                    }
3421                }
3422                // SRD-92: the "does this exit fail the phase?" rule lives
3423                // once in the DaemonExit taxonomy's own classifier — gate
3424                // the shared stop-flag latch on it rather than re-encoding
3425                // which variants fail by which arms call store().
3426                if exit.is_phase_error() {
3427                    activity.stop_flag.store(true, Ordering::Relaxed);
3428                }
3429            }
3430        }
3431
3432        // Final completion line — always emitted (one per phase),
3433        // not gated on TTY/extent. Replaces the old executor-side
3434        // `phase 'X' complete (Ns)` line. Honors the live
3435        // `suppress_status_line` flag (TUI takes over rendering)
3436        // and the global `suppress_progress` (e.g. CI / `--quiet`).
3437        if !suppress_progress && !activity.config.suppress_status_line.load(Ordering::Relaxed) {
3438            // Counter snapshots — the readout recomputes
3439            // pct / rate / ok_pct from these primitives, so
3440            // we don't pre-format them here. Retries are
3441            // derived per the existing convention (errors
3442            // minus the skips-adjusted failed-op count).
3443            let consumed = activity.source_factory.global_consumed();
3444            let ops_completed = activity.metrics.cycles_completed();
3445            // SRD-91: terminal-success count = `result_success.count()`;
3446            // `errors_total` is RESULT-level (one inc per terminal
3447            // failure) — it drives `e:` directly. Retries = failed
3448            // attempts that were NOT terminal, i.e. `attempt_failure
3449            // - failed_ops`; the old `errors_total - failed_ops`
3450            // collapsed to ~0 once `errors_total` went result-level
3451            // with the TriesDispenser refactor.
3452            let successes = activity.metrics.result_success.count();
3453            let errors = activity.metrics.errors_total.get();
3454            let elapsed = start_time.elapsed().as_secs_f64();
3455            let failed_ops = ops_completed
3456                .saturating_sub(successes)
3457                .saturating_sub(activity.metrics.skips_total.get());
3458            let retries = activity
3459                .metrics
3460                .attempt_failure
3461                .count()
3462                .saturating_sub(failed_ops);
3463            // Concurrency (fiber count) — the `c:N` tail mirrors
3464            // the live progress line so a completed phase reads
3465            // with the same shape as a running one.
3466            let concurrency = activity.config.concurrency;
3467            // Workload-emphasized metrics — same resolver as the
3468            // inline progress line, glob-matched against the
3469            // declared `status_metrics: [...]`. Empty list ⇒ no
3470            // metrics tail; nothing is presumed to be present.
3471            let relevancy_str: String = activity
3472                .metrics
3473                .collect_status_values(&activity.config.status_metrics)
3474                .concat();
3475            // SRD-100 P2 — no explicit status clear here. The consumer
3476            // folds `active_phases`, and the executor's `PhaseCompleted`
3477            // removes this phase from that map, so the status footer
3478            // self-clears for this phase on the next render tick (while
3479            // any concurrent phase's status survives — the old single-slot
3480            // `status(None)` wiped peers).
3481            // Render the ✓ DONE line via the readout engine.
3482            // SRD-63 / Push 1: the previous inline `format!()`
3483            // is now `phase_outcome.render()` driven by an
3484            // `ActivityReadoutContext` snapshot of the values
3485            // gathered above. Output is byte-equivalent.
3486            let phase_name_bare = activity
3487                .config
3488                .name
3489                .split_once(" (")
3490                .map(|(n, _)| n.to_string())
3491                .unwrap_or_else(|| activity.config.name.clone());
3492            // Phase-end: re-read the source's final extent.
3493            // For static cursors this equals the initial
3494            // `total_extent`; for extending cursors it's the
3495            // last grown value before the policy declined
3496            // further extension.
3497            let final_extent = source_for_progress
3498                .global_extent()
3499                .unwrap_or(activity.config.cycles);
3500            // SRD-76 — the activity-level binder fire happens
3501            // BEFORE the executor records its formal Failed /
3502            // Skipped decision. The activity knows only what it
3503            // measured: a clean completion if it ran to extent,
3504            // a stop-flag trip if the error router fired. Mirror
3505            // the stop_flag into the two-axis Outcome so the
3506            // readout doesn't render ✓ on a stopped phase. The
3507            // executor records the canonical outcome on the scene
3508            // tree; this surface is the realtime display projection.
3509            let outcome = if activity.stop_flag.load(Ordering::Relaxed) {
3510                crate::phase_outcome::Outcome::failed()
3511            } else {
3512                crate::phase_outcome::Outcome::completed()
3513            };
3514            let outcome_errors: Vec<crate::phase_outcome::PhaseErrorDetail> = activity
3515                .phase_errors
3516                .lock()
3517                .ok()
3518                .map(|g| g.clone())
3519                .unwrap_or_default();
3520            let ctx = crate::readout_context::ActivityReadoutContext {
3521                phase_name: phase_name_bare,
3522                phase_seq: activity.config.phase_seq,
3523                phase_labels: activity.config.phase_labels.clone(),
3524                cycles_completed: ops_completed,
3525                cycles_total: final_extent,
3526                ops_ok: successes,
3527                skips: activity.metrics.skips_total.get(),
3528                errors,
3529                retries,
3530                concurrency,
3531                elapsed_secs: elapsed,
3532                consumed,
3533                status_metric_chips: relevancy_str,
3534                depth_indent: crate::scene_tree::running_phase_indent(),
3535                use_color: crate::observer::use_color(),
3536                memo: activity.memo.load().as_str().to_string(),
3537                outcome,
3538                outcome_errors,
3539                outcome_resume_cursor: None,
3540                open_ended: activity.daemon_stop.is_some(),
3541            };
3542            // SRD-63 §6.2 / Push 9c: synthesise one final
3543            // `on_update` tick before the DONE summary. The
3544            // inline thread (if running) fires every 500 ms
3545            // and may have missed the last 100-499 ms of
3546            // counter changes — and for short phases under
3547            // the TTY/extent threshold it never spawned at
3548            // all. This guarantees the snapshot store sees
3549            // the phase's end-of-life on_update render
3550            // matching what the user would have seen if the
3551            // refresh tick had aligned exactly with phase
3552            // termination.
3553            //
3554            // Renders silently (no eprint) — the DONE line
3555            // immediately following carries the visible
3556            // ✓ summary; we just want the snapshot row to
3557            // reflect end-state.
3558            {
3559                let (final_seq, final_depth) =
3560                    crate::readout_context::resolve_phase_coord_by_name(&activity.config.name);
3561                // Row-level cursor progress — only for a DECLARED cursor
3562                // (`config.source_factory` present). Plain `cycles:` phases
3563                // get a synthesized `range(0, cycles)` factory, so gate on
3564                // the config Option, not the resolved factory, to keep them
3565                // on the op-denominated `cycles:` chip.
3566                let (final_rows_consumed, final_rows_total) = match &activity.config.source_factory
3567                {
3568                    Some(_) => (
3569                        activity.source_factory.global_consumed(),
3570                        activity.source_factory.global_extent().unwrap_or(0),
3571                    ),
3572                    None => (0, 0),
3573                };
3574                let final_ctx = crate::readout_context::build_inline_refresh_context(
3575                    &activity.metrics,
3576                    &activity.config.name,
3577                    activity.config.concurrency,
3578                    total_extent,
3579                    final_rows_consumed,
3580                    final_rows_total,
3581                    elapsed,
3582                    u64::MAX, // sentinel: spinner frame doesn't matter at end-of-phase
3583                    &activity.config.status_metrics,
3584                    activity.memo.as_ref(),
3585                    final_seq,
3586                    final_depth,
3587                    // End-of-phase: the subject is closed, never open-ended.
3588                    false,
3589                );
3590                let phase_status_default = {
3591                    let readout = crate::readouts::Registry::lookup("phase_status")
3592                        .expect("phase_status registered");
3593                    crate::readouts::BakedBody::from_single(readout, crate::readouts::Lod::Labeled)
3594                };
3595                if let Ok(mut binder) = crate::readouts::binder::build_event_binder_with_cli(
3596                    &activity.config.readouts,
3597                    crate::lifecycle::EventType::Update,
3598                    phase_status_default,
3599                    activity.config.cli_readout_override.as_deref(),
3600                ) {
3601                    use crate::readouts::ReadoutBinder;
3602                    use crate::readouts::ReadoutContext;
3603                    let mut sink = crate::readouts::StringSink::with_capacity(192);
3604                    binder.fire(crate::lifecycle::EventType::Update, &final_ctx, &mut sink);
3605                    let rendered_final = sink.take();
3606                    crate::readouts::snapshot::capture(
3607                        activity.config.snapshot_writer.as_ref(),
3608                        crate::lifecycle::EventType::Update.slot_name(),
3609                        final_ctx.subject_exec_id(),
3610                        crate::lifecycle::EventType::Update.subject_kind().as_str(),
3611                        &final_ctx.subject_id(),
3612                        "binder",
3613                        crate::readouts::snapshot::lod_str(crate::readouts::Lod::Labeled),
3614                        &rendered_final,
3615                    );
3616                }
3617            }
3618
3619            // gutter/final — the guaranteed ONE final gutter update.
3620            // A declared `final:` template is evaluated here, at phase
3621            // end (wires over the activity chain, status-metric
3622            // aggregates as fallback), and stored into the gutter slot
3623            // the display actor reads when stamping this phase's ✓
3624            // outcome DETAIL line. Without a `final:`, the slot already
3625            // holds the during-form's LAST PUBLISHED value (the
3626            // dispenser publishes on every op completion) — that last
3627            // computed value IS the final update; re-evaluating the
3628            // per-cycle template here would run it against a
3629            // degenerate end-of-phase context.
3630            {
3631                let fin = activity
3632                    .gutter_final_spec
3633                    .lock()
3634                    .ok()
3635                    .and_then(|g| g.clone());
3636                if let Some((kind, tmpl)) = fin {
3637                    if let Some(final_spec) =
3638                        evaluate_final_gutter(&activity, op_builder.source_kernel(), kind, &tmpl)
3639                    {
3640                        activity.gutter.store(Some(Arc::new(final_spec)));
3641                    }
3642                }
3643            }
3644
3645            // Build a one-shot binder for `on_phase_end`:
3646            // workload's `on_phase_end:` overrides if any,
3647            // else the default body — `phase_outcome` for the
3648            // normative ✓/✗ status line, followed by
3649            // `error_readout` which renders the per-error block
3650            // below it. `error_readout` is a no-op (zero bytes)
3651            // when the phase has no recorded errors, so the
3652            // default is safe for both success and failure
3653            // paths; failure paths get the structured error
3654            // block appended without the per-cycle warns
3655            // having to spam the screen mid-phase.
3656            let phase_outcome_default = {
3657                let phase_outcome = crate::readouts::Registry::lookup("phase_outcome")
3658                    .expect("phase_outcome registered");
3659                let error_readout = crate::readouts::Registry::lookup("error_readout")
3660                    .expect("error_readout registered");
3661                crate::readouts::BakedBody::from_steps(vec![
3662                    crate::readouts::binder::RenderStep::Render {
3663                        readout: phase_outcome,
3664                        lod: crate::readouts::Lod::Labeled,
3665                        layout: crate::readouts::binder::LayoutMode::Auto,
3666                        options: crate::readouts::ReadoutOptions::new(),
3667                        color: None,
3668                    },
3669                    crate::readouts::binder::RenderStep::Render {
3670                        readout: error_readout,
3671                        lod: crate::readouts::Lod::Labeled,
3672                        layout: crate::readouts::binder::LayoutMode::Auto,
3673                        options: crate::readouts::ReadoutOptions::new(),
3674                        color: None,
3675                    },
3676                ])
3677            };
3678            let rendered = match crate::readouts::build_event_binder(
3679                &activity.config.readouts,
3680                crate::lifecycle::EventType::PhaseEnd,
3681                phase_outcome_default,
3682            ) {
3683                Ok(mut binder) => {
3684                    use crate::readouts::ReadoutBinder;
3685                    let mut sink = crate::readouts::StringSink::with_capacity(160);
3686                    binder.fire(crate::lifecycle::EventType::PhaseEnd, &ctx, &mut sink);
3687                    sink.take()
3688                }
3689                Err(e) => {
3690                    crate::diag!(
3691                        crate::observer::LogLevel::Error,
3692                        "readouts: failed to bind on_phase_end — {e}"
3693                    );
3694                    String::new()
3695                }
3696            };
3697            // Push 6: capture the on_phase_end render to the
3698            // snapshot store. The DONE line is the canonical
3699            // "what the operator saw at completion" — replay
3700            // returns it byte-for-byte.
3701            if !rendered.is_empty() {
3702                use crate::readouts::ReadoutContext;
3703                crate::readouts::snapshot::capture(
3704                    activity.config.snapshot_writer.as_ref(),
3705                    crate::lifecycle::EventType::PhaseEnd.slot_name(),
3706                    ctx.subject_exec_id(),
3707                    crate::lifecycle::EventType::PhaseEnd
3708                        .subject_kind()
3709                        .as_str(),
3710                    &ctx.subject_id(),
3711                    "binder",
3712                    crate::readouts::snapshot::lod_str(crate::readouts::Lod::Labeled),
3713                    &rendered,
3714                );
3715            }
3716            // `skipped_phases=elide|prune`: a FULLY-SKIPPED phase (every
3717            // cycle `if:`-gated off, nothing measured, no errors) leaves
3718            // no completion line — the snapshot store above still holds
3719            // the render for replay, but the live readout stays silent.
3720            // `mark` (the default) emits phase_outcome's explicit
3721            // `⊘ gated off` form instead.
3722            let fully_skipped = {
3723                let skips = activity.metrics.skips_total.get();
3724                skips > 0
3725                    && ops_completed > 0
3726                    && skips >= ops_completed
3727                    && successes == 0
3728                    && errors == 0
3729            };
3730            let elide_skipped = fully_skipped
3731                && matches!(
3732                    crate::observer::skipped_phase_display(),
3733                    crate::observer::SkippedPhaseDisplay::Elide
3734                        | crate::observer::SkippedPhaseDisplay::Prune
3735                );
3736            if !rendered.is_empty() && !elide_skipped {
3737                // SRD-81 push 1: the per-phase ✓ outcome is a typed
3738                // `PhaseOutcome` projection, not a generic diagnostic.
3739                // The terminal scrollback shows it; the TUI tree /
3740                // active-phase panel render it natively; the TUI log
3741                // panel (diagnostics-only) filters it out instead of
3742                // garbling the multi-line ANSI as one Span. (push 1b
3743                // replaces this pre-rendered string with a structured
3744                // marker the sinks render from the snapshot.)
3745                crate::observer::log_tagged(
3746                    crate::observer::LogLevel::Info,
3747                    crate::observer::EventTag::at(
3748                        crate::lifecycle::EventType::PhaseEnd,
3749                        crate::observer::EventCategory::Outcome,
3750                    ),
3751                    &rendered,
3752                );
3753            }
3754        }
3755
3756        // Print validation summary AND capture to the metrics
3757        // store in one pass. `snapshot()` drains the histogram
3758        // (delta semantics), so we must use the same snapshot
3759        // for both printing and SQLite capture.
3760        if !validation_metrics.is_empty() {
3761            let mut total_passed = 0u64;
3762            let mut total_failed = 0u64;
3763            let now = Instant::now();
3764            let mut final_snapshot = MetricSet::at(now, Duration::ZERO);
3765            let activity_labels = activity.labels.clone();
3766
3767            for vm in validation_metrics.iter() {
3768                total_passed += vm.passed();
3769                total_failed += vm.failed();
3770
3771                for (name, stats) in &vm.relevancy_stats {
3772                    let snap = stats.snapshot();
3773                    if !snap.is_empty() {
3774                        let mean = snap.mean();
3775                        let p50 = snap.p50();
3776                        let p99 = snap.p99();
3777                        let min = snap.min();
3778                        let max = snap.max();
3779                        let n = snap.len();
3780                        // Relevancy stats (recall@k, precision@k, F1@k)
3781                        // are fractions in [0, 1]. Render as percent
3782                        // — the unit operators read these in.
3783                        // Underlying gauges below stay as fractions so
3784                        // downstream consumers (recall_summary,
3785                        // metrics scrapes) keep their existing scale.
3786                        // Indent matches the phase / DONE / complete
3787                        // lines so the relevancy summary nests under
3788                        // the phase row in tui=terminal output.
3789                        let depth_indent = crate::scene_tree::running_phase_indent();
3790                        let color = crate::observer::use_color();
3791                        let dim = if color { "\x1b[2m" } else { "" };
3792                        let bold = if color { "\x1b[1m" } else { "" };
3793                        let reset = if color { "\x1b[0m" } else { "" };
3794                        // SRD-92: a PhaseDetail projection — a detail
3795                        // row of the completion block above it, so the
3796                        // terminal sink renders it under the blank
3797                        // divider margin (no timing triad) and
3798                        // `completed_phases=headers` drops it.
3799                        crate::observer::log_tagged(
3800                            crate::observer::LogLevel::Info,
3801                            crate::observer::EventTag::at(
3802                                crate::lifecycle::EventType::PhaseEnd,
3803                                crate::observer::EventCategory::Evaluation,
3804                            ),
3805                            &format!(
3806                                "{depth_indent}{bold}{name}{reset}: mean={:.2}% {dim}p50={:.2}% p99={:.2}% min={:.2}% max={:.2}% (n={n}){reset}",
3807                                mean * 100.0,
3808                                p50 * 100.0,
3809                                p99 * 100.0,
3810                                min * 100.0,
3811                                max * 100.0,
3812                            ),
3813                        );
3814                        // Pick up `k`/`r` from the F64Stats's
3815                        // labels so per-phase summary gauges
3816                        // remain unique under OpenMetrics §4.5
3817                        // when multiple relevancy configs share
3818                        // a phase but differ in cutoff.
3819                        let stats_labels = stats.labels();
3820                        let k_label = stats_labels.get("k").map(str::to_string);
3821                        let r_label = stats_labels.get("r").map(str::to_string);
3822                        // Generic observability point: a relevancy
3823                        // function's per-phase summary has been
3824                        // computed and is about to be published
3825                        // as `{name}_{stat}` gauges. The trace
3826                        // fires for ANY relevancy function — the
3827                        // labels carry the publishing dimensions
3828                        // (phase, profile, …, k, r, n) from the
3829                        // surrounding scope, not from any
3830                        // workload-specific knowledge.
3831                        if crate::observer::trace_enabled() {
3832                            let mut trace_labels = activity_labels.with("n", n.to_string());
3833                            if let Some(k) = &k_label {
3834                                trace_labels = trace_labels.with("k", k);
3835                            }
3836                            if let Some(r) = &r_label {
3837                                trace_labels = trace_labels.with("r", r);
3838                            }
3839                            crate::observer::trace(
3840                                &trace_labels,
3841                                &format!(
3842                                    "event=relevancy.publish fn={name} n={n} \
3843                                     mean={mean:.6} p50={p50:.6} p99={p99:.6} \
3844                                     min={min:.6} max={max:.6}"
3845                                ),
3846                            );
3847                        }
3848                        for (stat, val) in [
3849                            ("mean", mean),
3850                            ("p50", p50),
3851                            ("p99", p99),
3852                            ("min", min),
3853                            ("max", max),
3854                        ] {
3855                            let mut gauge_labels = activity_labels.with("n", n.to_string());
3856                            if let Some(k) = &k_label {
3857                                gauge_labels = gauge_labels.with("k", k);
3858                            }
3859                            if let Some(r) = &r_label {
3860                                gauge_labels = gauge_labels.with("r", r);
3861                            }
3862                            final_snapshot.insert_gauge(
3863                                format!("{name}_{stat}"),
3864                                gauge_labels,
3865                                val,
3866                                now,
3867                            );
3868                        }
3869                    }
3870                }
3871            }
3872
3873            // Phase-level aggregate counters. One pair per phase
3874            // — `total_passed` / `total_failed` sum across every
3875            // op's `vm` so the metric instance is unique under
3876            // OpenMetrics §4.5 (LabelSets must be unique). The
3877            // earlier per-`vm` insertion path inserted N copies
3878            // with identical labels, which the snapshot
3879            // assembler now rejects as a duplicate. Per-op
3880            // breakdown isn't carried by the validation counters
3881            // anyway — the labels are activity-scope, not
3882            // op-scope.
3883            if total_passed > 0 || total_failed > 0 {
3884                final_snapshot.insert_counter(
3885                    "validations_passed",
3886                    activity_labels.clone(),
3887                    total_passed,
3888                    now,
3889                );
3890                final_snapshot.insert_counter(
3891                    "validations_failed",
3892                    activity_labels.clone(),
3893                    total_failed,
3894                    now,
3895                );
3896            }
3897
3898            // Validation summary line: only emit when there are
3899            // failures. On clean runs the relevancy summary's
3900            // `n=N` already conveys "N validations passed", and
3901            // the `validation: N passed, 0 failed` line was just
3902            // duplicate text on every phase. On failure runs the
3903            // line is signal — promote it to Warn so it stands
3904            // out and route only when failed > 0.
3905            if total_failed > 0 {
3906                let depth_indent = crate::scene_tree::running_phase_indent();
3907                crate::diag!(
3908                    crate::observer::LogLevel::Warn,
3909                    "{depth_indent}validation: {} passed, {} FAILED",
3910                    total_passed,
3911                    total_failed
3912                );
3913            }
3914
3915            if !final_snapshot.is_empty() {
3916                activity
3917                    .validation_frame
3918                    .lock()
3919                    .unwrap_or_else(|e| e.into_inner())
3920                    .replace(final_snapshot);
3921            }
3922        }
3923
3924        // NOTE: this is the `stopped` RETURN that `run_phase` reads to
3925        // decide whether the phase FAILED — so it must reflect only
3926        // abnormal stops (the error-handler `stop_flag`, Ctrl-C, a walk
3927        // fault). A daemon phase's `daemon_stop` is a CLEAN termination
3928        // (Interrupted+Succeeded — the foreground it shadows finished),
3929        // so it drives the loop BREAKS below but is deliberately NOT a
3930        // fault. The daemon exclusion lives once in `StopView::abnormal`
3931        // (session_signals.rs) — this delegates to it rather than
3932        // re-deriving the rule (was a hand-rolled copy; SRD-92 dedup).
3933        activity.stop_view().abnormal()
3934    }
3935}
3936
3937/// Executor task for the tiered DriverAdapter interface.
3938///
3939/// Each fiber has its own FiberBuilder (lock-free Polydat state).
3940/// Ops within a stanza are processed in dependency groups:
3941/// - Groups execute sequentially (captures flow between groups)
3942/// - Ops within a group execute concurrently (join_all)
3943///
3944/// Groups are determined at init time by analyzing capture
3945/// declarations and references across templates.
3946// `pull_plans`: per-template wrapper-side `PullPlan`s, sealed at init.
3947// Drives cycle-time reads for validation / conditional / throttle
3948// wrappers via memoized `PullHandle`s. See SRD 31 §"Pull plan vs bind
3949// plan".
3950/// One-shot daemon dispatch. Mirrors `executor_task`'s setup
3951/// (FiberBuilder + per-op kernel attach) but dispatches the
3952/// daemon's op exactly once, racing the in-flight future
3953/// against the per-daemon stop flag AND the activity-global
3954/// stop flag.
3955///
3956/// Cancellation path: when either flag flips, the
3957/// `dispenser.execute(...)` future is dropped at the next
3958/// await point — for the HTTP adapter that's mid-`send()`,
3959/// which propagates as a clean reqwest cancellation. The
3960/// daemon returns `DaemonExit::Cancelled`. The pool's grace
3961/// window (see `DaemonPool::shutdown`) gives the adapter time
3962/// to observe the drop and exit; deadlines past the grace
3963/// surface as `DaemonExit::TimedOut`.
3964///
3965/// Daemon ops increment `ops_started` / `ops_finished` like
3966/// any other op execution — the operator-visible op-count
3967/// surface stays consistent regardless of fiber kind. Service
3968/// + response timing is recorded the same way the cycle-pool
3969///   records it (one Instant pair around the execute), but no
3970///   rate-limiter acquire — daemons aren't subject to the
3971///   activity's ops-per-second ceiling.
3972async fn daemon_dispatch(
3973    activity: Arc<Activity>,
3974    dispensers: Arc<Vec<Arc<dyn OpDispenser>>>,
3975    pull_plans: Arc<Vec<crate::fixture::PullPlan>>,
3976    op_builder: Arc<crate::synthesis::OpBuilder>,
3977    template_idx: usize,
3978    op_name: String,
3979    stop: crate::daemon_pool::DaemonStopFlag,
3980) -> crate::daemon_pool::DaemonExit {
3981    let mut fiber = op_builder.create_fiber_builder();
3982    fiber.attach_dispenser_kernels(&dispensers);
3983    let dispenser = dispensers[template_idx].clone();
3984    let fields = crate::adapter::ResolvedFields::new(Vec::new(), Vec::new());
3985    // Resolve the wrapper-side pull plan against this daemon fiber's
3986    // kernel — same as the cycle-pool path. Daemon ops carry wrappers
3987    // too (notably `if:`, whose IF_COND wrapper registers a pull for
3988    // its predicate); an empty `ResolvedPulls` would panic when that
3989    // handle resolves. Captures the daemon reads (e.g. a `shared`
3990    // cell written by an earlier stanza op gating `if: sstables > 1`)
3991    // are visible here because the daemon dispatches at its op-walk
3992    // position, after the writer op completed.
3993    let pulls = fiber.resolve_pulls_for_idx(template_idx, &pull_plans[template_idx]);
3994    let cycle_wires = fiber.cycle_wires(template_idx);
3995    let ctx = crate::fixture::ExecCtx::with_wires(&fields, &pulls, &cycle_wires);
3996
3997    activity.metrics.ops_started.fetch_add(1, Ordering::Relaxed);
3998    let started = std::time::Instant::now();
3999    let activity_stop = activity.stop_flag.clone();
4000    let exit = tokio::select! {
4001        result = dispenser.execute(0, &ctx) => match result {
4002            Ok(_) => crate::daemon_pool::DaemonExit::Completed,
4003            Err(e) => crate::daemon_pool::DaemonExit::Errored(e),
4004        },
4005        _ = poll_daemon_stop(&stop, &activity_stop) => {
4006            crate::daemon_pool::DaemonExit::Cancelled
4007        }
4008    };
4009    let service_nanos = started.elapsed().as_nanos() as u64;
4010    activity.metrics.cycles_total.inc();
4011    activity
4012        .metrics
4013        .ops_finished
4014        .fetch_add(1, Ordering::Relaxed);
4015    activity.metrics.service_time.record(service_nanos);
4016    activity.metrics.response_time.record(service_nanos);
4017    // SRD-91 op-outcome taxonomy for daemon dispatch. The `attempt_*`
4018    // counters are owned by the innermost `TriesDispenser` in the op's
4019    // wrapper stack (which this daemon op runs through, same as a foreground
4020    // op), so only the RESULT-level tallies are recorded here — no
4021    // double-count. Cancelled / TimedOut are shutdown outcomes tracked via
4022    // daemon_cancelled_total / daemon_errors_total, not op results.
4023    match &exit {
4024        crate::daemon_pool::DaemonExit::Completed => {
4025            activity.metrics.result_total.inc();
4026            activity.metrics.result_success.observe(service_nanos);
4027        }
4028        crate::daemon_pool::DaemonExit::Errored(_) => {
4029            // `errors_total` / per-type tallies + policy effects already
4030            // ran in the OUTERMOST `ErrorHandlerDispenser` of the daemon
4031            // op's own stack (SRD-82 Part 3b) — only the result-level
4032            // outcome is recorded here.
4033            activity.metrics.result_total.inc();
4034            activity.metrics.result_failure.observe(service_nanos);
4035        }
4036        _ => {}
4037    }
4038    crate::diag!(
4039        crate::observer::LogLevel::Debug,
4040        "daemon op '{op_name}' exit={} elapsed_ms={:.0}",
4041        exit.label(),
4042        service_nanos as f64 / 1_000_000.0
4043    );
4044    exit
4045}
4046
4047/// Polls both the per-daemon stop flag and the activity-global
4048/// stop flag at 50ms granularity. Returns as soon as either is
4049/// set. The 50ms cadence is a compromise: fast enough that
4050/// daemon cancellation lands well under the typical 5-second
4051/// grace window, slow enough that an idle daemon doesn't burn
4052/// CPU. Replacing this with a `tokio::sync::Notify`-backed
4053/// flag would drop the latency to zero but requires touching
4054/// the StopFlag surface and isn't load-bearing for the
4055/// trigger-and-observe pattern.
4056async fn poll_daemon_stop(
4057    daemon_stop: &crate::daemon_pool::DaemonStopFlag,
4058    activity_stop: &Arc<std::sync::atomic::AtomicBool>,
4059) {
4060    loop {
4061        if daemon_stop.load(Ordering::Acquire) || activity_stop.load(Ordering::Acquire) {
4062            return;
4063        }
4064        tokio::time::sleep(Duration::from_millis(50)).await;
4065    }
4066}
4067
4068// reason: cohesive per-fiber executor driver; each argument is a distinct
4069// runtime channel/handle the loop needs — splitting into a struct would only
4070// relocate the same fields with no clarity gain.
4071#[allow(clippy::too_many_arguments)]
4072async fn executor_task(
4073    activity: Arc<Activity>,
4074    dispensers: Arc<Vec<Arc<dyn OpDispenser>>>,
4075    pull_plans: Arc<Vec<crate::fixture::PullPlan>>,
4076    op_builder: Arc<crate::synthesis::OpBuilder>,
4077    // Optional activity-level rate limiter. `acquire` fires
4078    // once per cycle before adapter dispatch. There is no
4079    // separate stanza-rate limiter.
4080    rate_limiter: Option<Arc<RateLimiter>>,
4081    // Per-fiber cooperative-exit flag owned by the activity's
4082    // [`crate::fiber_pool::FiberPool`]. Set to `true` by
4083    // `ConcurrencyApplier` when the pool scales down.
4084    fiber_stop: crate::fiber_pool::StopFlag,
4085    // Daemon-op pool — shared across cycle-pool fibers. Each
4086    // stanza walk that reaches a daemon op dispatches a fresh
4087    // fiber onto this pool via `try_spawn` and continues
4088    // without awaiting; the daemon body runs to completion
4089    // (or until phase-exit drain signals stop) on its own
4090    // tokio task.
4091    daemon_pool: Arc<crate::daemon_pool::DaemonPool>,
4092    // Phase name — used to wrap the daemon body in the same
4093    // runtime-context guard cycle-pool fibers use, so daemon
4094    // ops can read `phase()` / `cycle()` runtime-context
4095    // wires.
4096    phase_name_arc: Arc<str>,
4097) {
4098    let stanza_positions = activity.op_sequence.stanza_length();
4099    // SRD-22 batching cover-once: the phase-cursor stride is the SUM of
4100    // each stanza op's `rows_per_op` (its uniform per-invocation cursor
4101    // consumption), NOT the raw stanza length. A normal op contributes 1
4102    // (identical to the pre-batching model); a batch op contributes its
4103    // fixed stride `N`. Reserving `Σ rows_per_op` per stanza and handing
4104    // each op a contiguous sub-run of its own `rows_per_op` makes
4105    // consecutive stanzas cover DISJOINT ordinal runs — every ordinal is
4106    // inserted exactly once. Precomputed once (rows_per_op is fixed per
4107    // dispenser at map_op) so the hot loop pays no per-stanza virtual
4108    // calls and `Σ per_pos_rows == stanza_stride` by construction.
4109    let per_pos_rows: Vec<usize> = (0..stanza_positions)
4110        .map(|pos| {
4111            let (idx, _) = activity.op_sequence.get_with_index(pos as u64);
4112            dispensers[idx].rows_per_op().max(1)
4113        })
4114        .collect();
4115    let stanza_stride: usize = per_pos_rows.iter().sum::<usize>().max(1);
4116    // Per-fiber `FiberBuilder` carries scope values (per-iteration
4117    // extern inputs) populated by the OpBuilder, so iter-var
4118    // references like `{table}` in op templates resolve to the
4119    // current iteration's value.
4120    let mut fiber = op_builder.create_fiber_builder();
4121
4122    // Shutdown-ladder subscription (session_signals module doc): the op
4123    // dispatch below races the adapter call against the CANCEL rung
4124    // (level 2) so a hung request — one that will only ever end by
4125    // client timeout — can be dropped mid-flight, letting the drain and
4126    // the process-level cleanup (WAL consolidation, summaries) proceed.
4127    // One receiver per fiber; the race is a `select!` per dispatch.
4128    let mut shutdown_rx = crate::session_signals::subscribe_shutdown();
4129
4130    // SRD-68 Push 3 — materialise per-fiber subscope kernels from
4131    // each dispenser's canonical kernel. The fiber holds them as
4132    // `Vec<Option<PolydatKernel>>` indexed parallel to the dispenser
4133    // registry; cycle dispatch reads `fiber.per_op_kernel(template_idx)`
4134    // to populate `ExecCtx::wires` for the firing dispenser.
4135    // Dispensers that return `None` from `canonical_kernel()` get
4136    // a `None` slot and the cycle falls back to `NullWireSource`.
4137    fiber.attach_dispenser_kernels(&dispensers);
4138
4139    // Create per-fiber source reader (used for all phases).
4140    // Source-declared phases will eventually use the advancer model,
4141    // but for now all phases go through the source reader.
4142    // SRD-92 Step 5e — the per-cycle stream flows through the unified
4143    // `ChildSource` contract: a `CursorSource` wraps the reader; `poll_next` IS
4144    // `reserve(stanza_stride)` (per-stanza, not per-cycle → zero per-cycle
4145    // overhead), yielding an ordinal `Range`; `render` is the per-ordinal fetch.
4146    // The level selects `CursorReserve` — this very FiberPool loop.
4147    use crate::child_source::{Child, ChildSource, CursorSource, Drive, select_drive};
4148    let mut source = CursorSource::new(activity.source_factory.create_reader(), stanza_stride);
4149    debug_assert_eq!(select_drive(source.realizability()), Drive::CursorReserve);
4150
4151    loop {
4152        if activity.stopped() {
4153            break;
4154        } // SRD-92 Step 0: one stop view
4155        if fiber_stop.load(std::sync::atomic::Ordering::Acquire) {
4156            break;
4157        } // per-fiber scale-down (distinct)
4158
4159        // Phase 1: RESERVE — CAS on shared cursor, instantaneous.
4160        // Acquires one stanza's worth of ordinals. This is the only
4161        // shared-state interaction per stanza.
4162        let range = match source.poll_next() {
4163            Some(Child::Ordinals(r)) => r,
4164            // CursorSource yields only Ordinals; anything else (None) means the
4165            // source is exhausted → the phase-poll rewind / standard break.
4166            _ => {
4167                // SRD-75 phase-poll: source exhausted ends a poll
4168                // iteration. Check the predicate; if false and
4169                // the deadline hasn't elapsed, sleep, rewind the
4170                // factory's shared cursor, and re-create the
4171                // reader for another iteration. Phase-poll
4172                // mandates concurrency=1 (workload-load
4173                // validation) so this is the only fiber and the
4174                // rewind isn't racing siblings.
4175                if let Some(pp) = activity.phase_poll.clone() {
4176                    // Check predicate first — handles the case
4177                    // where the very first iteration's captures
4178                    // already satisfy the condition.
4179                    //
4180                    // Polydat comparison operators (`==`, `!=`, `<`,
4181                    // …) return u64 (0/1) per SRD-10 §"BinOpKind"
4182                    // — there's no Bool result type for these.
4183                    // Accept either Value::Bool(true) (in case
4184                    // a future Polydat release adds a Bool result
4185                    // path) OR a non-zero numeric value as
4186                    // "satisfied". This mirrors the workload
4187                    // author's expectation that `(a == 1) & (b == 0)`
4188                    // evaluates to "true" when both clauses hold.
4189                    //
4190                    // `__poll_until` is a DYNAMIC binding
4191                    // (SRD-11 §"Two Evaluation Lifecycles") —
4192                    // its value depends on per-iteration capture
4193                    // writes through the phase scope's
4194                    // SharedCells, so a buffer read via
4195                    // `lookup()` returns the LAST-EVALUATED
4196                    // value (None on first iteration, never
4197                    // updated). We MUST trigger re-evaluation
4198                    // via `pull()`. The phase scope kernel is
4199                    // held as `Arc<PolydatKernel>` (immutable
4200                    // handle), so we evaluate via the per-fiber
4201                    // `main_kernel` instead — main_kernel is
4202                    // built from the phase scope program (so
4203                    // it has `__poll_until` as an output) and
4204                    // is wired to the SAME SharedCells the
4205                    // captures wrote to (so its pull returns
4206                    // the live value).
4207                    // SRD-75 (C5) — strict-gate `require:` selectors.
4208                    // While any selector is unresolved the gate HOLDS:
4209                    // the predicate is not trusted, because an
4210                    // unregistered family reads 0.0 silently and a
4211                    // `>=`-shaped predicate could pass spuriously (or
4212                    // a `<`-shaped one hang to timeout). Past the
4213                    // grace window (one poll interval) an unresolved
4214                    // selector is a hard `poll_require` failure —
4215                    // loud and immediate, never a mystery hang.
4216                    let unresolved: Vec<&String> = pp
4217                        .require
4218                        .iter()
4219                        .filter(|s| !nmbrs_metrics::polydat_nodes::metric_selector_resolves(s))
4220                        .collect();
4221                    if !unresolved.is_empty() && std::time::Instant::now() >= pp.require_grace {
4222                        let names = unresolved
4223                            .iter()
4224                            .map(|s| format!("'{s}'"))
4225                            .collect::<Vec<_>>()
4226                            .join(", ");
4227                        let formatted_reason = format!(
4228                            "[poll_require] poll `require:` selector(s) {names} \
4229                             resolved to no registered instrument within the \
4230                             grace window — the gate's `until:` would read 0.0 \
4231                             for them. Check the family name and labels (e.g. \
4232                             phase=<name>), and that the producing phase is \
4233                             running. SRD-75 (C5)."
4234                        );
4235                        crate::diag!(
4236                            crate::observer::LogLevel::Error,
4237                            "activity '{}': {formatted_reason}",
4238                            activity.config.name
4239                        );
4240                        if let Ok(mut slot) = activity.stop_reason.lock()
4241                            && slot.is_none()
4242                        {
4243                            *slot = Some(formatted_reason.clone());
4244                        }
4245                        if let Ok(mut errs) = activity.phase_errors.lock() {
4246                            errs.push(crate::phase_outcome::PhaseErrorDetail {
4247                                class: "poll_require".into(),
4248                                message: formatted_reason,
4249                                op_name: None,
4250                                cycle: None,
4251                                op_template: None,
4252                                op_resolved: None,
4253                                at_nanos: std::time::SystemTime::now()
4254                                    .duration_since(std::time::UNIX_EPOCH)
4255                                    .map(|d| d.as_nanos() as u64)
4256                                    .unwrap_or(0),
4257                                retryable: false,
4258                            });
4259                        }
4260                        activity
4261                            .stop_flag
4262                            .store(true, std::sync::atomic::Ordering::Relaxed);
4263                        break;
4264                    }
4265                    let requires_ok = unresolved.is_empty();
4266                    // SAME call the op-level poll makes. An execution node is
4267                    // an execution node: the predicate is resolved through one
4268                    // interface (`CycleWires`, which PULLS and so re-evaluates
4269                    // a dynamic binding) and judged by one truthiness rule.
4270                    //
4271                    // This was a third, divergent implementation — an inline
4272                    // match whose `_ => false` arm made a non-empty string
4273                    // falsy here and truthy under `if:` / `while:` / op-poll,
4274                    // for the same expression. (The `require:` gate above
4275                    // additionally withholds trust in the predicate until
4276                    // every declared metric selector resolves.)
4277                    let satisfied = requires_ok && {
4278                        let wires = fiber.main_wires();
4279                        match crate::wrappers::condition::holds(
4280                            &wires,
4281                            crate::wrappers::condition::UNTIL_BINDING,
4282                        ) {
4283                            Some(v) => v,
4284                            // Unresolved is a wiring fault, not a false
4285                            // predicate. Keep waiting rather than silently
4286                            // declaring the phase done; the timeout below
4287                            // reports it with the predicate named.
4288                            None => false,
4289                        }
4290                    };
4291                    if satisfied {
4292                        // SRD-75 metric_name emission is wired
4293                        // when the per-fiber mutable kernel handle
4294                        // Elapsed time goes to the metric (when one
4295                        // is configured); operators read it there.
4296                        // Logging it at INFO every loop is just
4297                        // narrating the happy path.
4298                        if let Some(name) = &pp.metric_name {
4299                            let elapsed = pp.started_at.elapsed().as_secs_f64();
4300                            crate::diag!(
4301                                crate::observer::LogLevel::Debug,
4302                                "phase-poll: predicate satisfied; {name}={elapsed:.3}s",
4303                            );
4304                        }
4305                        break;
4306                    }
4307                    if std::time::Instant::now() >= pp.deadline {
4308                        let elapsed = pp.started_at.elapsed().as_secs_f64();
4309                        // Compose the diagnostic once; the `abort`
4310                        // path appends a workload-invalidation
4311                        // note so the operator immediately knows
4312                        // that the whole run is terminating, not
4313                        // just this phase.
4314                        let invalidation_note = match pp.on_timeout {
4315                            PhasePollTimeoutPolicy::Abort => {
4316                                " — `on_timeout: abort` declared by the workload; \
4317                                 requesting session stop (the whole run terminates)"
4318                            }
4319                            PhasePollTimeoutPolicy::Error => "",
4320                        };
4321                        let formatted_reason = format!(
4322                            "[poll_timeout] phase-poll deadline reached after {elapsed:.1}s \
4323                             with predicate '__poll_until' still not Bool(true) \
4324                             (SRD-75 §\"Workload-load validation\" — adjust `timeout_ms` \
4325                             or the `until:` predicate){invalidation_note}"
4326                        );
4327                        if let Ok(mut slot) = activity.stop_reason.lock()
4328                            && slot.is_none()
4329                        {
4330                            *slot = Some(formatted_reason.clone());
4331                        }
4332                        // SRD-76 — push a structured
4333                        // `PhaseErrorDetail` into the
4334                        // activity's phase_errors buffer so
4335                        // the executor's phase-end build of
4336                        // `PhaseOutcome` captures the
4337                        // poll_timeout with its class and
4338                        // message. No op_template /
4339                        // op_resolved because the failure is
4340                        // at the phase level (no specific op
4341                        // dispenser fired the error).
4342                        if let Ok(mut errs) = activity.phase_errors.lock() {
4343                            errs.push(crate::phase_outcome::PhaseErrorDetail {
4344                                class: "poll_timeout".into(),
4345                                message: formatted_reason,
4346                                op_name: None,
4347                                cycle: None,
4348                                op_template: None,
4349                                op_resolved: None,
4350                                at_nanos: std::time::SystemTime::now()
4351                                    .duration_since(std::time::UNIX_EPOCH)
4352                                    .map(|d| d.as_nanos() as u64)
4353                                    .unwrap_or(0),
4354                                retryable: false,
4355                            });
4356                        }
4357                        activity
4358                            .stop_flag
4359                            .store(true, std::sync::atomic::Ordering::Relaxed);
4360                        // SRD-75 `on_timeout: abort` —
4361                        // workload-author declares that an
4362                        // unsatisfied predicate makes the whole
4363                        // run meaningless. Set the session-wide
4364                        // stop signal; the scenario walker
4365                        // observes it on its next iteration check
4366                        // (`session_signals::stop_requested()`)
4367                        // and unwinds without entering the next
4368                        // sweep cell. The phase itself still
4369                        // returns Err for the normal stop-flag
4370                        // path; the session signal is the
4371                        // CROSS-PHASE escalation.
4372                        if matches!(pp.on_timeout, PhasePollTimeoutPolicy::Abort) {
4373                            crate::diag!(
4374                                crate::observer::LogLevel::Error,
4375                                "phase-poll: `on_timeout: abort` triggered after \
4376                                 {elapsed:.1}s; requesting session-wide stop \
4377                                 (SRD-75 §\"on_timeout\")",
4378                            );
4379                            crate::session_signals::request_stop();
4380                        }
4381                        break;
4382                    }
4383                    // Wait, rewind, re-create the reader.
4384                    tokio::time::sleep(pp.interval).await;
4385                    if !activity.source_factory.rewind_for_poll() {
4386                        if let Ok(mut slot) = activity.stop_reason.lock()
4387                            && slot.is_none()
4388                        {
4389                            *slot = Some(
4390                                "[phase_poll] source factory doesn't support \
4391                                 rewind_for_poll(); phase-poll requires a \
4392                                 rewindable source (RangeSourceFactory or \
4393                                 ExtendingRangeSourceFactory). SRD-75."
4394                                    .to_string(),
4395                            );
4396                        }
4397                        activity
4398                            .stop_flag
4399                            .store(true, std::sync::atomic::Ordering::Relaxed);
4400                        break;
4401                    }
4402                    source =
4403                        CursorSource::new(activity.source_factory.create_reader(), stanza_stride);
4404                    continue;
4405                }
4406                break; // source exhausted (standard path)
4407            }
4408        };
4409
4410        activity.metrics.stanzas_total.inc();
4411        // Stanza-boundary `reset_captures()` was historically called
4412        // here to defend against capture leakage between cycles. That
4413        // defence is redundant under the post-closure-binding-economy
4414        // architecture: every reachable wire is either a per-cycle
4415        // kernel output (recomputed each cycle from inputs and the
4416        // current `cycle`), a closed-loop capture (the same op that
4417        // reads the wire is the op that writes it; a failed write
4418        // short-circuits the consumer via OpResult error), a magic
4419        // extern (`body` / `count` / `ok`, rewritten pre-eval by
4420        // ResultDispenser), or a scope-invariant iter-var / workload
4421        // param (constant across the phase activation). None of those
4422        // can hold a stale value that a successful cycle would read.
4423        // The per-cycle reset + re-apply round-trip was 40% of single-
4424        // fiber CPU; removing it leaves end-state semantics identical.
4425
4426        // Phase 2: RENDER + EXECUTE — distribute the reserved run across
4427        // the stanza's ops via `StanzaRuns`. Each op covers a contiguous
4428        // ordinal sub-run `[cycle, cycle + run_len)` of its own
4429        // `rows_per_op` (1 for ordinary ops, N for a batch op); `Σ ==
4430        // stanza_stride`, so consecutive stanzas cover DISJOINT ordinal
4431        // runs (SRD-22 cover-once). The LUT POSITION (not the ordinal)
4432        // selects the op, so a batch op that advances the cursor by N
4433        // still maps to the right stanza slot. At the cursor tail
4434        // `reserve` returned a short range, so the last op's `run_len`
4435        // is the truncated remainder — the partial final batch is still
4436        // inserted, never over-read, never dropped. Sequential in
4437        // declaration order.
4438        for (pos, cycle, run_len) in
4439            crate::child_source::StanzaRuns::new(range.clone(), &per_pos_rows)
4440        {
4441            if activity.stopped() {
4442                break;
4443            } // SRD-92 Step 0: one stop view
4444
4445            // Mark op as active from render through result join.
4446            // "Active" means this fiber is working on an op — resolving
4447            // fields, waiting for the adapter, or recording results.
4448            activity.metrics.ops_started.fetch_add(1, Ordering::Relaxed);
4449
4450            // Render the source item at the sub-run base (fiber-local, no
4451            // shared state). The op reads `[cycle, cycle + run_len)`.
4452            let item = source.render(cycle);
4453            // Publish the cycle to the enclosing fiber-context
4454            // scope so any Polydat node reading `cycle()` or implicitly
4455            // `cycle` inside the DAG sees the same ordinal as
4456            // adapter execution. No-op outside a fiber scope.
4457            crate::polydat_nodes::runtime_context::set_task_cycle(cycle);
4458
4459            let wait_start = Instant::now();
4460            if let Some(ref rl) = rate_limiter {
4461                rl.acquire().await;
4462            }
4463            let wait_nanos = wait_start.elapsed().as_nanos() as u64;
4464
4465            // The stanza POSITION (not the ordinal) selects the op —
4466            // `get_with_index(pos)` returns `lut[pos]`, stable regardless
4467            // of how many ordinals prior batch ops consumed.
4468            let (template_idx, template) = activity.op_sequence.get_with_index(pos as u64);
4469
4470            // Daemon-op dispatch. If the template declares
4471            // `daemon: ...` (non-disabled), spawn a fresh
4472            // fiber onto the daemon pool instead of running
4473            // the op inline. The pool enforces a per-op-name
4474            // fiber cap; an overflow is a workload-design
4475            // error that fails the phase. Cycle-pool moves
4476            // on to the next stanza op as soon as the spawn
4477            // returns (Ok or Err).
4478            if !template.daemon.is_disabled() {
4479                let cap = template
4480                    .daemon
4481                    .max_fibers()
4482                    .expect("non-disabled daemon has cap");
4483                let cancel_grace = template
4484                    .daemon_cancel_grace_ms
4485                    .map(std::time::Duration::from_millis);
4486                let activity_d = activity.clone();
4487                let dispensers_d = dispensers.clone();
4488                let pull_plans_d = pull_plans.clone();
4489                let op_builder_d = op_builder.clone();
4490                let phase_arc_d = phase_name_arc.clone();
4491                let op_name_d = template.name.clone();
4492                let spawn_result =
4493                    daemon_pool.try_spawn(op_name_d.clone(), cap, cancel_grace, move |stop| {
4494                        let activity = activity_d;
4495                        let dispensers = dispensers_d;
4496                        let pull_plans = pull_plans_d;
4497                        let op_builder = op_builder_d;
4498                        let phase_arc = phase_arc_d;
4499                        let op_name = op_name_d;
4500                        async move {
4501                            use futures::FutureExt as _;
4502                            let phase_controls = activity
4503                                .component
4504                                .as_ref()
4505                                .map(crate::polydat_nodes::runtime_context::snapshot_controls)
4506                                .unwrap_or_else(
4507                                    crate::polydat_nodes::runtime_context::empty_controls,
4508                                );
4509                            let body = crate::polydat_nodes::runtime_context::with_fiber_context(
4510                                phase_arc,
4511                                phase_controls,
4512                                daemon_dispatch(
4513                                    activity.clone(),
4514                                    dispensers,
4515                                    pull_plans,
4516                                    op_builder,
4517                                    template_idx,
4518                                    op_name.clone(),
4519                                    stop,
4520                                ),
4521                            );
4522                            match std::panic::AssertUnwindSafe(body).catch_unwind().await {
4523                                Ok(exit) => exit,
4524                                Err(payload) => {
4525                                    let msg = payload
4526                                        .downcast_ref::<&'static str>()
4527                                        .map(|s| (*s).to_string())
4528                                        .or_else(|| payload.downcast_ref::<String>().cloned())
4529                                        .unwrap_or_else(|| "<non-string panic payload>".into());
4530                                    crate::diag!(
4531                                        crate::observer::LogLevel::Error,
4532                                        "daemon op '{op_name}' panicked: {msg}"
4533                                    );
4534                                    activity
4535                                        .stop_flag
4536                                        .store(true, std::sync::atomic::Ordering::Relaxed);
4537                                    crate::daemon_pool::DaemonExit::Panicked(msg)
4538                                }
4539                            }
4540                        }
4541                    });
4542                match spawn_result {
4543                    Ok(()) => {
4544                        // Dispatch succeeded — the daemon fiber
4545                        // (`daemon_dispatch`) owns this op's accounting
4546                        // end to end: it records `ops_started` when it
4547                        // runs and `cycles_total` / `ops_finished` + the
4548                        // SRD-91 result outcome when it completes. Counting
4549                        // the stanza-position here too DOUBLE-counted the
4550                        // op — one extra `cycles_total`/`ops_finished`
4551                        // with no matching result — which dragged ok%
4552                        // below 100% and pushed phase progress above it,
4553                        // and broke the `cycles_total == result_total +
4554                        // skips_total` invariant. `StanzaRuns` already
4555                        // advanced past this op's sub-run (daemon ops have
4556                        // rows_per_op == 1); just skip the inline execute
4557                        // path and let the daemon fiber count.
4558                        continue;
4559                    }
4560                    Err(msg) => {
4561                        crate::diag!(
4562                            crate::observer::LogLevel::Error,
4563                            "daemon op '{}' spawn failed: {msg}",
4564                            template.name
4565                        );
4566                        activity.stop_flag.store(true, Ordering::Release);
4567                        if let Ok(mut slot) = activity.stop_reason.lock()
4568                            && slot.is_none()
4569                        {
4570                            *slot = Some(format!("daemon op '{}' spawn: {msg}", template.name));
4571                        }
4572                        return;
4573                    }
4574                }
4575            }
4576
4577            fiber.set_source_item(&item);
4578            // SRD-68 Push 5: `ctx.fields` is no longer the
4579            // resolution surface for adapters or wrappers — they
4580            // read everything through `ctx.wires` (the bound GK
4581            // context). The empty `ResolvedFields` satisfies the
4582            // `ExecCtx` struct-shape contract until the field is
4583            // removed from the trait surface entirely.
4584            let fields = crate::adapter::ResolvedFields::new(Vec::new(), Vec::new());
4585
4586            // Resolve the wrapper-side pull plan against this
4587            // fiber's GkState (one indexed pull per registered
4588            // name, no name hashing). The resulting `pulls` is
4589            // disjoint from `fields`: adapters see only `fields`,
4590            // wrappers see only `pulls`.
4591            //
4592            // SRD-13d Phase 9 — when this op template materialised
4593            // its own kernel, the plan was sealed against the
4594            // op-template program; resolve_pulls_for_op picks
4595            // that kernel's state. Flattened op-templates fall
4596            // through to the main kernel (the workload program)
4597            // — same call site, the lookup is idempotent.
4598            let pulls = fiber.resolve_pulls_for_idx(template_idx, &pull_plans[template_idx]);
4599            let dispenser = &dispensers[template_idx];
4600            // SRD-68 invariant I-2: cycle-time reads against the
4601            // firing dispenser's per-fiber kernel slot, exposed
4602            // through the narrow `WireSource` trait. `CycleWires`
4603            // wraps the per-fiber kernel handle for the cycle's
4604            // duration so `WireSource::get` can drive output pulls
4605            // (`pull(&mut state, …)` through interior mutability)
4606            // alongside input/constant lookups. Dispensers with
4607            // no canonical kernel (legacy adapters, wrapper
4608            // delegates) fall through to the `NullWireSource`
4609            // baseline `ExecCtx::new` provides.
4610            // SRD-13f / SRD-68: cycle-time wire reads go through a
4611            // single kernel handle — the dispenser's per-fiber
4612            // op-template kernel. Every visible cross-scope wire
4613            // was wired into that kernel at construction (cells
4614            // for shared, folded constants for workload params,
4615            // construction-time slot setup + per-cycle refresh in
4616            // `set_inputs` for other parent outputs). The local
4617            // read API resolves every name; the wires layer never
4618            // composes chains externally.
4619            let cycle_wires = fiber.cycle_wires(template_idx);
4620            let mut exec_ctx = crate::fixture::ExecCtx::with_wires(&fields, &pulls, &cycle_wires);
4621            // Hand the op the ACTUAL reserved sub-run length so a batch op
4622            // inserts exactly `[base, base + run_len)` — the full run, or
4623            // the short tail at cursor exhaustion. Ordinary ops ignore it.
4624            exec_ctx.run_len = run_len;
4625            let service_start = Instant::now();
4626            // The op runs through the wrapper stack ONCE. The innermost
4627            // `TriesDispenser` (when the op has a `tries` budget) owns the
4628            // attempt loop, `attempt_*` counters,
4629            // and the per-attempt panic catch; the OUTERMOST
4630            // `ErrorHandlerDispenser` (SRD-82 Part 3b) owns the whole-stack
4631            // panic backstop and the terminal-error handling — policy
4632            // routing, `errors_total` / per-type tallies, phase-error
4633            // capture, and the stop/fail effects. This loop sees exactly ONE
4634            // terminal outcome per cycle and keeps only the result-level
4635            // accounting. The residual `catch_unwind` guards a panic in the
4636            // error wrapper's OWN routing code (a runtime bug, not an op
4637            // failure): keep the fiber alive, mark the phase stopped.
4638            let outcome: Result<crate::adapter::OpResult, crate::adapter::ExecutionError> = {
4639                use futures::FutureExt as _;
4640                let op_fut = std::panic::AssertUnwindSafe(dispenser.execute(cycle, &exec_ctx))
4641                    .catch_unwind();
4642                // Race the whole op stack against the shutdown ladder's
4643                // CANCEL rung. `biased` polls the op first, so the cancel
4644                // branch costs one extra poll per dispatch on the happy
4645                // path. On cancellation the stack's future is DROPPED —
4646                // that is the cancel — and a synthesised non-retryable
4647                // error stands in as the terminal outcome (the errors
4648                // wrapper was inside the dropped future, so result-level
4649                // accounting below is all that records it).
4650                let raced = tokio::select! {
4651                    biased;
4652                    r = op_fut => Some(r),
4653                    _ = crate::session_signals::ops_cancelled(&mut shutdown_rx) => None,
4654                };
4655                match raced {
4656                    None => Err(crate::adapter::ExecutionError::Op(
4657                        crate::adapter::AdapterError {
4658                            error_name: "cancelled".into(),
4659                            message: "in-flight op cancelled by shutdown \
4660                                      escalation (Ctrl-C)"
4661                                .into(),
4662                            retryable: false,
4663                        },
4664                    )),
4665                    Some(Ok(r)) => r,
4666                    Some(Err(payload)) => {
4667                        let msg = payload
4668                            .downcast_ref::<&'static str>()
4669                            .map(|s| (*s).to_string())
4670                            .or_else(|| payload.downcast_ref::<String>().cloned())
4671                            .unwrap_or_else(|| "<non-string panic payload>".into());
4672                        activity.metrics.errors_total.inc();
4673                        activity.metrics.count_error_type("panic");
4674                        activity.stop_flag.store(true, Ordering::Relaxed);
4675                        if let Ok(mut slot) = activity.stop_reason.lock()
4676                            && slot.is_none()
4677                        {
4678                            // Headline = first line; the full
4679                            // enriched text travels in the
4680                            // AdapterError message and renders
4681                            // once in the phase error list
4682                            // (SRD-82 §one full render).
4683                            let first = msg.lines().next().unwrap_or(&msg);
4684                            *slot = Some(format!(
4685                                "[panic] op '{}' at cycle {}: {first}",
4686                                template.name, cycle,
4687                            ));
4688                        }
4689                        Err(crate::adapter::ExecutionError::Op(
4690                            crate::adapter::AdapterError {
4691                                error_name: "panic".into(),
4692                                message: msg,
4693                                retryable: false,
4694                            },
4695                        ))
4696                    }
4697                }
4698            };
4699            let service_nanos = service_start.elapsed().as_nanos() as u64;
4700            let (success, skipped) = match outcome {
4701                Ok(result) => (true, result.skipped),
4702                Err(_) => (false, false),
4703            };
4704
4705            // Per-OP totals (SRD-91): `cycles_total` counts every op
4706            // dispatched; executed ops go to `result_total`
4707            // (= result_success + result_failure). Skipped ops increment
4708            // `skips_total` in the `if:` wrapper (wrappers/if.rs), so
4709            // `cycles_total == result_total + skips_total` holds without a
4710            // tally here. The per-op error rate reads `result_failure`
4711            // (in [0,1]); the per-attempt `errors_total` was already
4712            // tallied in the loop.
4713            activity.metrics.cycles_total.inc();
4714            if !skipped {
4715                activity.metrics.result_total.inc();
4716                activity.metrics.service_time.record(service_nanos);
4717                activity.metrics.wait_time.record(wait_nanos);
4718                activity
4719                    .metrics
4720                    .response_time
4721                    .record(service_nanos + wait_nanos);
4722                // `tries_histogram` is recorded by the TriesDispenser (it owns
4723                // the attempt count).
4724                if success {
4725                    activity.metrics.result_success.observe(service_nanos);
4726                    // Captures landed on the per-op-template kernel
4727                    // directly via ctx.wires.write inside the
4728                    // dispenser stack — no post-execute pump.
4729                    //
4730                    // Two kernel-side steps remain:
4731                    //
4732                    // 1. Rule 2 write-through commit on the
4733                    //    op-template kernel — pulls every
4734                    //    `__write_<X>` and stores its value through
4735                    //    the cell-bound input slot for `<X>`,
4736                    //    propagating result-binding LHS values up
4737                    //    to parent `shared` cells. No-op when the
4738                    //    kernel carries no write-throughs. A
4739                    //    type-stability violation (scope_model.md
4740                    //    §"Type stability") is a DETERMINISTIC
4741                    //    workload bug — every cycle would repeat it —
4742                    //    so it stops the phase with the write-site
4743                    //    diagnostic (same treatment as the panic arm).
4744                    if let Err(e) = fiber.commit_op_template_write_throughs_for_idx(template_idx) {
4745                        activity.metrics.errors_total.inc();
4746                        activity.metrics.count_error_type("type_mismatch");
4747                        activity.stop_flag.store(true, Ordering::Relaxed);
4748                        if let Ok(mut slot) = activity.stop_reason.lock()
4749                            && slot.is_none()
4750                        {
4751                            *slot = Some(format!(
4752                                "[type_mismatch] op '{}' at cycle {}: {e}",
4753                                template.name, cycle,
4754                            ));
4755                        }
4756                    }
4757                    // 2. Pull every output of the op-template
4758                    //    kernel so side-effecting nodes (log_info,
4759                    //    log_debug) inside result-binding compute
4760                    //    chains actually evaluate. Without this, a
4761                    //    result-binding whose LHS isn't a
4762                    //    write-through stays dormant.
4763                    fiber.pull_all_op_template_outputs_for_idx(template_idx);
4764                } else {
4765                    activity.metrics.result_failure.observe(service_nanos);
4766                }
4767            }
4768
4769            // Op fully processed — render, execute, and metrics all done.
4770            activity
4771                .metrics
4772                .ops_finished
4773                .fetch_add(1, Ordering::Relaxed);
4774        }
4775    }
4776}
4777
4778/// Best-effort terminal column count read off stderr (fd 2) via
4779/// the `TIOCGWINSZ` ioctl. Returns `None` when stderr isn't a
4780/// TTY or the call fails. The status renderer (now hosted by
4781/// `nmbrs-tui::log_only_sink`) uses this to clamp the rendered
4782/// status to a single visual row, since a wrap would leave
4783/// previous-tick text on screen below the cursor — the in-place
4784/// rewrite only erases from the cursor through end of the
4785/// current visual line.
4786#[cfg(unix)]
4787pub fn terminal_cols() -> Option<usize> {
4788    use std::os::raw::c_int;
4789    #[repr(C)]
4790    struct WinSize {
4791        ws_row: u16,
4792        ws_col: u16,
4793        ws_xpixel: u16,
4794        ws_ypixel: u16,
4795    }
4796    let mut ws = WinSize {
4797        ws_row: 0,
4798        ws_col: 0,
4799        ws_xpixel: 0,
4800        ws_ypixel: 0,
4801    };
4802    // SAFETY: `libc::ioctl` is FFI; `TIOCGWINSZ` writes into the
4803    // out-parameter which we own (pinned on the stack for the
4804    // duration of the call). Failure is signalled by negative
4805    // return — we ignore the actual errno.
4806    let rc: c_int = unsafe { libc::ioctl(2, libc::TIOCGWINSZ, &mut ws as *mut _) };
4807    if rc < 0 || ws.ws_col == 0 {
4808        return None;
4809    }
4810    Some(ws.ws_col as usize)
4811}
4812
4813/// Windows variant: `GetConsoleScreenBufferInfo` on the stderr
4814/// handle. The console API is declared by hand rather than via a
4815/// `windows-sys` dependency — two kernel32 imports don't justify
4816/// one. Width is the visible window (srWindow), not the scrollback
4817/// buffer width, matching what the status line can occupy.
4818#[cfg(windows)]
4819pub fn terminal_cols() -> Option<usize> {
4820    use std::ffi::c_void;
4821    #[repr(C)]
4822    struct Coord {
4823        x: i16,
4824        y: i16,
4825    }
4826    #[repr(C)]
4827    struct SmallRect {
4828        left: i16,
4829        top: i16,
4830        right: i16,
4831        bottom: i16,
4832    }
4833    #[repr(C)]
4834    struct ConsoleScreenBufferInfo {
4835        size: Coord,
4836        cursor_position: Coord,
4837        attributes: u16,
4838        window: SmallRect,
4839        maximum_window_size: Coord,
4840    }
4841    #[link(name = "kernel32")]
4842    unsafe extern "system" {
4843        fn GetStdHandle(std_handle: u32) -> *mut c_void;
4844        fn GetConsoleScreenBufferInfo(
4845            console: *mut c_void,
4846            info: *mut ConsoleScreenBufferInfo,
4847        ) -> i32;
4848    }
4849    const STD_ERROR_HANDLE: u32 = -12i32 as u32;
4850    const INVALID_HANDLE_VALUE: *mut c_void = -1isize as *mut c_void;
4851    // SAFETY: both calls only write into the out-parameter we own;
4852    // failure is signalled by null/INVALID handle or zero return.
4853    unsafe {
4854        let handle = GetStdHandle(STD_ERROR_HANDLE);
4855        if handle.is_null() || handle == INVALID_HANDLE_VALUE {
4856            return None;
4857        }
4858        let mut info = std::mem::zeroed::<ConsoleScreenBufferInfo>();
4859        if GetConsoleScreenBufferInfo(handle, &mut info) == 0 {
4860            return None;
4861        }
4862        let cols = i32::from(info.window.right) - i32::from(info.window.left) + 1;
4863        if cols <= 0 { None } else { Some(cols as usize) }
4864    }
4865}
4866
4867/// Glob-style match: `*` matches zero or more characters, `?`
4868/// matches exactly one character, every other byte must match
4869/// literally. Recursive — adequate for the short patterns
4870/// `status_metrics:` accepts (`recall*`, `latency_p99`, etc.).
4871/// Trades worst-case quadratic time for simplicity; the
4872/// candidate set is also tiny (low single-digit count of metric
4873/// Evaluate a gutter template ONCE at phase end (the `final:` form,
4874/// or the during-form's guaranteed last update). Placeholders resolve
4875/// through the wires of a throwaway subscope of the activity's source
4876/// kernel (shared cells and captures visible), then any names still
4877/// unresolved fall back to the phase's STATUS-METRIC aggregates
4878/// (`{recall}`, `{latency_p50}`, …) formatted exactly like the status
4879/// chips. Numeric kinds (`bar`/`spark`) additionally require the fully
4880/// resolved string to parse as f64; failures degrade to None (no
4881/// final cell) — the display must never fail a completed phase.
4882fn evaluate_final_gutter(
4883    activity: &Activity,
4884    source_kernel: &Arc<crate::scope_kernel::ScopeKernel>,
4885    kind: crate::wrappers::gutter::GutterKind,
4886    template: &str,
4887) -> Option<crate::wrappers::gutter::GutterSpec> {
4888    use crate::wrappers::gutter::{GutterKind, GutterSpec};
4889    // Wires pass: a fork of the phase scope (its cells shared) gives
4890    // template names their live end-of-phase values.
4891    let rendered = {
4892        let mut k = source_kernel.fork();
4893        let wires = crate::wires::CycleWires::new(&mut k);
4894        crate::wires::substitute_via_wires(template, &wires).ok()
4895    };
4896    let mut text = rendered.unwrap_or_else(|| template.to_string());
4897
4898    // Status-metric fallback for placeholders the wires didn't know:
4899    // relevancy aggregates and the latency family, formatted like the
4900    // status chips so `{recall}` in a final template reads identically
4901    // to the `recall:` chip beside it.
4902    if text.contains('{') {
4903        let mut candidates: Vec<(String, String)> = Vec::new();
4904        for live in activity.metrics.collect_relevancy_live() {
4905            if live.total_count > 0 {
4906                candidates.push((live.name, format!("{:.2}%", live.total_mean * 100.0)));
4907            }
4908        }
4909        let snap = activity.metrics.service_time.peek_snapshot();
4910        let h = &snap.histogram;
4911        if !h.is_empty() {
4912            let fmt = nmbrs_metrics::reporters::summary::format_duration;
4913            candidates.push(("latency_p50".into(), fmt(h.value_at_quantile(0.50) as f64)));
4914            candidates.push(("latency_p99".into(), fmt(h.value_at_quantile(0.99) as f64)));
4915            candidates.push(("latency_max".into(), fmt(h.max() as f64)));
4916            candidates.push(("latency_mean".into(), fmt(h.mean())));
4917        }
4918        for (name, val) in &candidates {
4919            text = text.replace(&format!("{{{name}}}"), val);
4920        }
4921    }
4922
4923    // A template still carrying unresolved placeholders means the value
4924    // it names was never measured (a gated-off recall phase has no
4925    // relevancy aggregate) — no measurement, no cell. The completion
4926    // detail line keeps its standard stamp instead of showing the
4927    // literal `recall {recall}`.
4928    if text.contains('{') && text.contains('}') {
4929        return None;
4930    }
4931    match kind {
4932        GutterKind::Labeled => {
4933            let (name, value) = text.split_once('\u{1f}').unwrap_or(("", text.as_str()));
4934            Some(GutterSpec::Labeled {
4935                name: name.to_string(),
4936                value: value.to_string(),
4937            })
4938        }
4939        GutterKind::Text => Some(GutterSpec::Text(text)),
4940        GutterKind::Bar => text
4941            .trim()
4942            .parse::<f64>()
4943            .ok()
4944            .map(|v| GutterSpec::Bar(v.clamp(0.0, 1.0))),
4945        GutterKind::Spark => text.trim().parse::<f64>().ok().map(GutterSpec::Spark),
4946    }
4947}
4948
4949/// names per phase).
4950fn glob_match(pattern: &str, candidate: &str) -> bool {
4951    glob_match_bytes(pattern.as_bytes(), candidate.as_bytes())
4952}
4953
4954fn glob_match_bytes(pat: &[u8], s: &[u8]) -> bool {
4955    match (pat.first(), s.first()) {
4956        (None, None) => true,
4957        (Some(b'*'), _) => {
4958            // zero-or-more: try consuming nothing OR consume one
4959            // char of input and re-attempt.
4960            glob_match_bytes(&pat[1..], s) || (!s.is_empty() && glob_match_bytes(pat, &s[1..]))
4961        }
4962        (Some(b'?'), Some(_)) => glob_match_bytes(&pat[1..], &s[1..]),
4963        (Some(p), Some(c)) if p == c => glob_match_bytes(&pat[1..], &s[1..]),
4964        _ => false,
4965    }
4966}
4967
4968// `spinner_frame`, `braille_bar`, `format_eta` moved to
4969// `crate::readouts::format` in Push 2 — the readouts that
4970// consume them now own the helpers. `truncate_to_width`
4971// stays here (it's a surface-level width-clamp concern,
4972// not a readout concern).
4973
4974/// Truncate `s` to at most `max_cols` *visible* columns,
4975/// appending an ellipsis when truncation actually elides
4976/// content. Skips ANSI SGR escape sequences (`\x1b[...m`) when
4977/// counting visible width — they consume characters in the
4978/// string but no terminal columns. The truncation point is
4979/// always at a character boundary that's NOT inside an escape
4980/// sequence, so we never emit a half-broken `\x1b[3` to the
4981/// terminal.
4982pub fn truncate_to_width(s: &str, max_cols: usize) -> String {
4983    if max_cols == 0 {
4984        return String::new();
4985    }
4986    let bytes = s.as_bytes();
4987    let mut visible = 0usize;
4988    let mut byte_pos = 0usize; // last clean truncation point
4989    let mut chars = s.char_indices();
4990    while let Some((i, c)) = chars.next() {
4991        if c == '\x1b' && bytes.get(i + 1) == Some(&b'[') {
4992            // SGR escape: walk until the final byte (`m`,
4993            // `K`, `J`, etc.) so we don't truncate mid-escape.
4994            for (_, ch) in chars.by_ref() {
4995                if ch.is_ascii_alphabetic() {
4996                    break;
4997                }
4998            }
4999            // byte_pos doesn't advance — escape costs no
5000            // visible columns, and the next plain char's
5001            // position is what we'd truncate to.
5002            continue;
5003        }
5004        if visible + 1 > max_cols.saturating_sub(1) {
5005            return format!("{}…", &s[..byte_pos]);
5006        }
5007        visible += 1;
5008        byte_pos = i + c.len_utf8();
5009    }
5010    s.to_string()
5011}
5012
5013#[cfg(test)]
5014mod tests {
5015    use super::*;
5016    use crate::adapter::{AdapterError, ExecutionError, OpResult};
5017    use std::collections::HashMap;
5018    use std::sync::atomic::{AtomicU64, Ordering};
5019
5020    /// SRD-92 R4 — `collect_status_primary` feeds the key-metric
5021    /// gutter cell's trend: first `status_metrics:` pattern's first
5022    /// candidate, as a raw numeric (latency family in milliseconds).
5023    /// The relevancy-first candidate ordering is exercised end-to-end
5024    /// by `crates/nmbrs/tests/srd92_display.rs` (a live relevancy aggregate
5025    /// needs the full validation pipeline).
5026    #[test]
5027    fn status_primary_selects_first_matching_numeric() {
5028        let m = ActivityMetrics::new(&nmbrs_metrics::labels::Labels::empty());
5029        // No patterns → no primary, regardless of measurements.
5030        assert_eq!(m.collect_status_primary(&[]), None);
5031        // Patterns but nothing measured yet → None (no fabricated 0).
5032        assert_eq!(m.collect_status_primary(&["latency_*".into()]), None);
5033
5034        // 5 ms samples land in the service-time histogram.
5035        for _ in 0..10 {
5036            m.service_time.record(5_000_000);
5037        }
5038        let (name, val) = m
5039            .collect_status_primary(&["latency_p50".into()])
5040            .expect("p50 measurable");
5041        assert_eq!(name, "latency_p50");
5042        assert!((val - 5.0).abs() < 0.5, "p50 ≈ 5 ms, got {val}");
5043
5044        // Glob: first pattern's FIRST candidate wins (p50 precedes
5045        // p99 in candidate order).
5046        let (name, _) = m
5047            .collect_status_primary(&["latency_*".into()])
5048            .expect("glob matches");
5049        assert_eq!(name, "latency_p50");
5050
5051        // Non-matching pattern → None.
5052        assert_eq!(m.collect_status_primary(&["recall*".into()]), None);
5053    }
5054
5055    /// A counting DriverAdapter + OpDispenser for testing.
5056    struct CountingDriverAdapter {
5057        count: Arc<AtomicU64>,
5058    }
5059
5060    impl CountingDriverAdapter {
5061        fn new() -> (Self, Arc<AtomicU64>) {
5062            let count = Arc::new(AtomicU64::new(0));
5063            (
5064                Self {
5065                    count: count.clone(),
5066                },
5067                count,
5068            )
5069        }
5070    }
5071
5072    impl DriverAdapter for CountingDriverAdapter {
5073        fn name(&self) -> &str {
5074            "counting"
5075        }
5076        fn map_op<'a>(
5077            &'a self,
5078            _template: &'a nmbrs_workload::model::ParsedOp,
5079            _parent: std::sync::Arc<dyn polydat::Kernel>,
5080        ) -> std::pin::Pin<
5081            Box<dyn std::future::Future<Output = Result<Box<dyn OpDispenser>, String>> + Send + 'a>,
5082        > {
5083            Box::pin(async move {
5084                Ok(Box::new(CountingDispenser {
5085                    count: self.count.clone(),
5086                }) as Box<dyn OpDispenser>)
5087            })
5088        }
5089    }
5090
5091    struct CountingDispenser {
5092        count: Arc<AtomicU64>,
5093    }
5094
5095    impl OpDispenser for CountingDispenser {
5096        fn execute<'a>(
5097            &'a self,
5098            _cycle: u64,
5099            _ctx: &'a crate::fixture::ExecCtx<'a>,
5100        ) -> std::pin::Pin<
5101            Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
5102        > {
5103            self.count.fetch_add(1, Ordering::Relaxed);
5104            Box::pin(async {
5105                Ok(OpResult {
5106                    body: None,
5107                    skipped: false,
5108                })
5109            })
5110        }
5111    }
5112
5113    /// A fail-then-succeed DriverAdapter for retry testing.
5114    struct FailThenSucceedDriverAdapter {
5115        fails_remaining: Arc<AtomicU64>,
5116        total_calls: Arc<AtomicU64>,
5117    }
5118
5119    impl FailThenSucceedDriverAdapter {
5120        fn new(fail_count: u64) -> (Self, Arc<AtomicU64>) {
5121            let total = Arc::new(AtomicU64::new(0));
5122            (
5123                Self {
5124                    fails_remaining: Arc::new(AtomicU64::new(fail_count)),
5125                    total_calls: total.clone(),
5126                },
5127                total,
5128            )
5129        }
5130    }
5131
5132    impl DriverAdapter for FailThenSucceedDriverAdapter {
5133        fn name(&self) -> &str {
5134            "fail-then-succeed"
5135        }
5136        fn map_op<'a>(
5137            &'a self,
5138            _template: &'a nmbrs_workload::model::ParsedOp,
5139            _parent: std::sync::Arc<dyn polydat::Kernel>,
5140        ) -> std::pin::Pin<
5141            Box<dyn std::future::Future<Output = Result<Box<dyn OpDispenser>, String>> + Send + 'a>,
5142        > {
5143            Box::pin(async move {
5144                Ok(Box::new(FailThenSucceedDispenser {
5145                    fails_remaining: self.fails_remaining.clone(),
5146                    total_calls: self.total_calls.clone(),
5147                }) as Box<dyn OpDispenser>)
5148            })
5149        }
5150    }
5151
5152    struct FailThenSucceedDispenser {
5153        fails_remaining: Arc<AtomicU64>,
5154        total_calls: Arc<AtomicU64>,
5155    }
5156
5157    impl OpDispenser for FailThenSucceedDispenser {
5158        fn execute<'a>(
5159            &'a self,
5160            _cycle: u64,
5161            _ctx: &'a crate::fixture::ExecCtx<'a>,
5162        ) -> std::pin::Pin<
5163            Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
5164        > {
5165            self.total_calls.fetch_add(1, Ordering::Relaxed);
5166            let remaining = self.fails_remaining.fetch_sub(1, Ordering::Relaxed);
5167            Box::pin(async move {
5168                if remaining > 0 {
5169                    Err(ExecutionError::Op(AdapterError {
5170                        error_name: "TransientError".into(),
5171                        message: "temporary failure".into(),
5172                        retryable: true,
5173                    }))
5174                } else {
5175                    Ok(OpResult {
5176                        body: None,
5177                        skipped: false,
5178                    })
5179                }
5180            })
5181        }
5182    }
5183
5184    /// Build a minimal Polydat root kernel (single identity node) for tests.
5185    fn test_kernel() -> crate::scope_kernel::ScopeKernel {
5186        use polydat::compile::assembly::{PolydatAssembler, WireRef};
5187        use polydat::library::identity::Identity;
5188        let mut asm = PolydatAssembler::new(vec!["cycle".into()]);
5189        asm.add_node(
5190            "id",
5191            Box::new(Identity::new(polydat::ast::PortType::U64)),
5192            vec![WireRef::input("cycle")],
5193        );
5194        asm.add_output("id", WireRef::node("id"));
5195        asm.compile().unwrap().into()
5196    }
5197
5198    #[tokio::test]
5199    async fn activity_runs_all_cycles() {
5200        let config = ActivityConfig {
5201            name: "test".into(),
5202            cycles: 100,
5203            concurrency: 4,
5204            ..Default::default()
5205        };
5206        let ops = vec![nmbrs_workload::model::ParsedOp::simple("op1", "test")];
5207        let seq = OpSequence::uniform(ops);
5208        let activity = Activity::new(config, &Labels::of("session", "test"), seq);
5209
5210        let (adapter, count) = CountingDriverAdapter::new();
5211        activity
5212            .run_with_driver(
5213                Arc::new(adapter),
5214                Arc::new(crate::synthesis::OpBuilder::new(test_kernel())),
5215            )
5216            .await;
5217
5218        assert_eq!(count.load(Ordering::Relaxed), 100);
5219    }
5220
5221    #[tokio::test]
5222    async fn activity_retries_on_error() {
5223        // Retry backoff keeps this test's op IN FLIGHT for a few
5224        // hundred ms — long enough to overlap the session_signals
5225        // tests, whose bodies legitimately hold the process-global
5226        // cancel rung in force (the raced in-flight cancel would
5227        // kill attempt 3). Serialize with them, same discipline as
5228        // every global-flag test.
5229        let _signals = crate::session_signals::STOP_GLOBAL_TEST_LOCK
5230            .lock()
5231            .unwrap_or_else(|e| e.into_inner());
5232        let config = ActivityConfig {
5233            name: "retrytest".into(),
5234            cycles: 1,
5235            concurrency: 1,
5236            error_spec: "TransientError:retry,warn;.*:stop".into(),
5237            // Total-attempts budget (the `tries` sigil): 6 total ≈ the old
5238            // `retries: 5` additional-attempts budget.
5239            tries: Some(6),
5240            ..Default::default()
5241        };
5242        let ops = vec![nmbrs_workload::model::ParsedOp::simple("op1", "test")];
5243        let seq = OpSequence::uniform(ops);
5244        let activity = Activity::new(config, &Labels::of("session", "s1"), seq);
5245
5246        let (adapter, total_calls) = FailThenSucceedDriverAdapter::new(2);
5247        activity
5248            .run_with_driver(
5249                Arc::new(adapter),
5250                Arc::new(crate::synthesis::OpBuilder::new(test_kernel())),
5251            )
5252            .await;
5253
5254        assert_eq!(total_calls.load(Ordering::Relaxed), 3);
5255    }
5256
5257    #[tokio::test]
5258    async fn daemon_op_dispatches_at_cycle_pool_position() {
5259        // SRD-79 (in-flight): daemon-flagged op spawns onto the
5260        // daemon pool when the cycle-pool fiber's stanza walk
5261        // reaches it. The daemon fiber runs the same dispenser
5262        // as a non-daemon op would, but on its own tokio task,
5263        // and the cycle-pool fiber doesn't await it.
5264        let config = ActivityConfig {
5265            name: "daemon-disp-test".into(),
5266            cycles: 1,
5267            concurrency: 1,
5268            ..Default::default()
5269        };
5270        let mut op = nmbrs_workload::model::ParsedOp::simple("dmn", "test");
5271        op.daemon = nmbrs_workload::model::DaemonSpec::MaxFibers(1);
5272        let seq = OpSequence::uniform(vec![op]);
5273        let activity = Activity::new(config, &Labels::of("session", "s1"), seq);
5274
5275        let (adapter, count) = CountingDriverAdapter::new();
5276        activity
5277            .run_with_driver(
5278                Arc::new(adapter),
5279                Arc::new(crate::synthesis::OpBuilder::new(test_kernel())),
5280            )
5281            .await;
5282
5283        // Daemon dispatched + ran exactly once (cycles=1).
5284        assert_eq!(
5285            count.load(Ordering::Relaxed),
5286            1,
5287            "daemon op should have run via dispatch-time spawn"
5288        );
5289    }
5290
5291    #[tokio::test]
5292    async fn daemon_op_cap_exceeded_fails_phase() {
5293        // Cap=1 with cycles=2: first dispatch succeeds and the
5294        // daemon fiber blocks (the CountingDispenser returns
5295        // instantly, so the daemon should drain before the
5296        // second cycle — but the daemon-pool counter only
5297        // decrements when the body returns, and the second
5298        // cycle may race with the decrement). This test
5299        // primarily checks the no-panic / clean-failure path:
5300        // even if the cap fires, the activity exits cleanly.
5301        let config = ActivityConfig {
5302            name: "daemon-cap-test".into(),
5303            cycles: 50,
5304            concurrency: 1,
5305            ..Default::default()
5306        };
5307        let mut op = nmbrs_workload::model::ParsedOp::simple("dmn", "test");
5308        op.daemon = nmbrs_workload::model::DaemonSpec::MaxFibers(1);
5309        let seq = OpSequence::uniform(vec![op]);
5310        let activity = Activity::new(config, &Labels::of("session", "s2"), seq);
5311
5312        let (adapter, _count) = CountingDriverAdapter::new();
5313        activity
5314            .run_with_driver(
5315                Arc::new(adapter),
5316                Arc::new(crate::synthesis::OpBuilder::new(test_kernel())),
5317            )
5318            .await;
5319        // No assertion on count — the load-bearing behaviour is
5320        // that the activity terminates cleanly even when caps
5321        // bite. Without the cap, this test would hang or panic.
5322    }
5323
5324    #[tokio::test]
5325    async fn shared_metrics_accessible() {
5326        let config = ActivityConfig {
5327            name: "metricstest".into(),
5328            cycles: 50,
5329            concurrency: 2,
5330            ..Default::default()
5331        };
5332        let ops = vec![nmbrs_workload::model::ParsedOp::simple("op1", "test")];
5333        let seq = OpSequence::uniform(ops);
5334        let activity = Activity::new(config, &Labels::of("session", "s1"), seq);
5335
5336        let shared_metrics = activity.shared_metrics();
5337
5338        let (adapter, _count) = CountingDriverAdapter::new();
5339        activity
5340            .run_with_driver(
5341                Arc::new(adapter),
5342                Arc::new(crate::synthesis::OpBuilder::new(test_kernel())),
5343            )
5344            .await;
5345
5346        assert_eq!(shared_metrics.cycles_total.get(), 50);
5347        let frame = shared_metrics.capture(std::time::Duration::from_secs(1));
5348        assert!(!frame.is_empty());
5349    }
5350
5351    #[tokio::test]
5352    async fn per_error_type_counters_emit_deltas_through_dynamic_capture() {
5353        // SRD-40 / cascade coalesce: `MetricSet::combine_into` for
5354        // Counter sums `total` across intervals — so per-cycle
5355        // emissions must be DELTAS, not absolutes. Per-error-type
5356        // counters live on `ActivityMetricsDynamic` (outside the
5357        // static registry), so they need their own delta tracking.
5358        // This test exercises that path directly.
5359        use nmbrs_metrics::component::Component;
5360        use nmbrs_metrics::snapshot::MetricValue;
5361
5362        let metrics = Arc::new(ActivityMetrics::new(&Labels::of("session", "s1")));
5363        let component = Arc::new(std::sync::RwLock::new(Component::new(
5364            Labels::of("activity", "t"),
5365            HashMap::new(),
5366        )));
5367        {
5368            let mut g = component.write().unwrap();
5369            g.set_state(nmbrs_metrics::component::ComponentState::Running);
5370            metrics.register_on(&mut g).unwrap();
5371        }
5372
5373        // Seed two error-type counters with different totals.
5374        for _ in 0..3 {
5375            metrics.count_error_type("net");
5376        }
5377        for _ in 0..7 {
5378            metrics.count_error_type("timeout");
5379        }
5380
5381        // First capture_delta — totals=3 and 7 are the deltas.
5382        let snap1 = component
5383            .read()
5384            .unwrap()
5385            .capture_delta(std::time::Duration::from_secs(1));
5386        let net1 = read_counter(&snap1, "errors.net");
5387        let to1 = read_counter(&snap1, "errors.timeout");
5388        assert_eq!(net1, 3, "first delta for net should be 3, got {net1}");
5389        assert_eq!(to1, 7, "first delta for timeout should be 7, got {to1}");
5390
5391        // Drive the per-error-type counters further.
5392        for _ in 0..2 {
5393            metrics.count_error_type("net");
5394        }
5395        for _ in 0..1 {
5396            metrics.count_error_type("timeout");
5397        }
5398
5399        // Second capture_delta — should report only the new deltas
5400        // (2 and 1), NOT the absolute totals (5 and 8).
5401        let snap2 = component
5402            .read()
5403            .unwrap()
5404            .capture_delta(std::time::Duration::from_secs(1));
5405        let net2 = read_counter(&snap2, "errors.net");
5406        let to2 = read_counter(&snap2, "errors.timeout");
5407        assert_eq!(
5408            net2, 2,
5409            "second delta for net should be 2 (new only), got {net2}"
5410        );
5411        assert_eq!(
5412            to2, 1,
5413            "second delta for timeout should be 1 (new only), got {to2}"
5414        );
5415
5416        // capture_current (drain=false) should still report absolutes.
5417        let cur = component.read().unwrap().capture_current();
5418        let net_abs = read_counter(&cur, "errors.net");
5419        let to_abs = read_counter(&cur, "errors.timeout");
5420        assert_eq!(
5421            net_abs, 5,
5422            "current should be absolute total 5, got {net_abs}"
5423        );
5424        assert_eq!(
5425            to_abs, 8,
5426            "current should be absolute total 8, got {to_abs}"
5427        );
5428
5429        fn read_counter(snap: &nmbrs_metrics::snapshot::MetricSet, family: &str) -> u64 {
5430            let f = snap
5431                .family(family)
5432                .unwrap_or_else(|| panic!("family {family:?} missing from snapshot"));
5433            let m = f.metrics().next().expect("at least one metric");
5434            match m.point().unwrap().value() {
5435                MetricValue::Counter(c) => c.cumulative,
5436                v => panic!("not a counter: {v:?}"),
5437            }
5438        }
5439    }
5440
5441    #[tokio::test]
5442    async fn activity_with_rate() {
5443        let config = ActivityConfig {
5444            name: "ratetest".into(),
5445            cycles: 10,
5446            concurrency: 2,
5447            rate: Some(10000.0),
5448            ..Default::default()
5449        };
5450        let ops = vec![nmbrs_workload::model::ParsedOp::simple("op1", "test")];
5451        let seq = OpSequence::uniform(ops);
5452        let activity = Activity::new(config, &Labels::of("session", "s1"), seq);
5453
5454        let (adapter, count) = CountingDriverAdapter::new();
5455        activity
5456            .run_with_driver(
5457                Arc::new(adapter),
5458                Arc::new(crate::synthesis::OpBuilder::new(test_kernel())),
5459            )
5460            .await;
5461
5462        assert_eq!(count.load(Ordering::Relaxed), 10);
5463    }
5464
5465    #[tokio::test]
5466    async fn activity_with_weighted_ops() {
5467        let config = ActivityConfig {
5468            name: "weighted".into(),
5469            cycles: 12,
5470            concurrency: 1,
5471            ..Default::default()
5472        };
5473        let ops = vec![
5474            nmbrs_workload::model::ParsedOp::simple("read", "SELECT"),
5475            nmbrs_workload::model::ParsedOp::simple("write", "INSERT"),
5476        ];
5477        let seq = OpSequence::build(ops, &[4, 2], SequencerType::Bucket);
5478        let activity = Activity::new(config, &Labels::of("session", "s1"), seq);
5479
5480        let (adapter, count) = CountingDriverAdapter::new();
5481        activity
5482            .run_with_driver(
5483                Arc::new(adapter),
5484                Arc::new(crate::synthesis::OpBuilder::new(test_kernel())),
5485            )
5486            .await;
5487
5488        assert_eq!(count.load(Ordering::Relaxed), 12);
5489    }
5490
5491    #[tokio::test]
5492    async fn rate_control_is_declared_when_rate_configured() {
5493        use nmbrs_metrics::component::Component;
5494        use nmbrs_metrics::labels::Labels as L;
5495        use std::sync::RwLock;
5496
5497        let config = ActivityConfig {
5498            name: "rate_decl".into(),
5499            cycles: 5,
5500            concurrency: 1,
5501            rate: Some(2500.0),
5502            ..Default::default()
5503        };
5504        let ops = vec![nmbrs_workload::model::ParsedOp::simple("op1", "test")];
5505        let seq = OpSequence::uniform(ops);
5506        let mut activity = Activity::new(config, &L::of("session", "s_rate"), seq);
5507        let component = Arc::new(RwLock::new(Component::new(
5508            L::of("session", "s_rate"),
5509            std::collections::HashMap::new(),
5510        )));
5511        activity.attach_component(component.clone());
5512
5513        let (adapter, _count) = CountingDriverAdapter::new();
5514        activity
5515            .run_with_driver(
5516                Arc::new(adapter),
5517                Arc::new(crate::synthesis::OpBuilder::new(test_kernel())),
5518            )
5519            .await;
5520
5521        // After the activity runs, the rate control is on the
5522        // component and reports the configured target via its
5523        // reified gauge.
5524        let guard = component.read().unwrap();
5525        let erased = guard
5526            .controls()
5527            .get_erased("rate")
5528            .expect("rate control should be declared when rate is set");
5529        assert!(erased.accepts_f64_writes());
5530        assert_eq!(erased.gauge_f64(), Some(2500.0));
5531    }
5532
5533    #[tokio::test]
5534    async fn rate_control_is_absent_when_no_rate() {
5535        use nmbrs_metrics::component::Component;
5536        use nmbrs_metrics::labels::Labels as L;
5537        use std::sync::RwLock;
5538
5539        let config = ActivityConfig {
5540            name: "no_rate".into(),
5541            cycles: 3,
5542            concurrency: 1,
5543            rate: None,
5544            ..Default::default()
5545        };
5546        let ops = vec![nmbrs_workload::model::ParsedOp::simple("op1", "test")];
5547        let seq = OpSequence::uniform(ops);
5548        let mut activity = Activity::new(config, &L::of("session", "s_nr"), seq);
5549        let component = Arc::new(RwLock::new(Component::new(
5550            L::of("session", "s_nr"),
5551            std::collections::HashMap::new(),
5552        )));
5553        activity.attach_component(component.clone());
5554
5555        let (adapter, _count) = CountingDriverAdapter::new();
5556        activity
5557            .run_with_driver(
5558                Arc::new(adapter),
5559                Arc::new(crate::synthesis::OpBuilder::new(test_kernel())),
5560            )
5561            .await;
5562
5563        let guard = component.read().unwrap();
5564        assert!(
5565            guard.controls().get_erased("rate").is_none(),
5566            "no rate control should exist without rate configured",
5567        );
5568    }
5569
5570    #[tokio::test]
5571    async fn rate_control_write_retargets_the_running_limiter() {
5572        use nmbrs_metrics::component::Component;
5573        use nmbrs_metrics::controls::ControlOrigin;
5574        use nmbrs_metrics::labels::Labels as L;
5575        use std::sync::RwLock;
5576
5577        // 200 cycles with a low rate + a concurrent writer that
5578        // bumps the rate mid-flight. The committed value on the
5579        // control reflects the write; the limiter carries the
5580        // same target after reconfigure.
5581        let config = ActivityConfig {
5582            name: "rate_live".into(),
5583            cycles: 200,
5584            concurrency: 2,
5585            rate: Some(50.0),
5586            ..Default::default()
5587        };
5588        let ops = vec![nmbrs_workload::model::ParsedOp::simple("op1", "test")];
5589        let seq = OpSequence::uniform(ops);
5590        let mut activity = Activity::new(config, &L::of("session", "s_live"), seq);
5591        let component = Arc::new(RwLock::new(Component::new(
5592            L::of("session", "s_live"),
5593            std::collections::HashMap::new(),
5594        )));
5595        activity.attach_component(component.clone());
5596
5597        // Spawn the activity, wait for the applier to be wired,
5598        // issue a typed write, assert the control value advanced.
5599        let component_for_writer = component.clone();
5600        let writer = tokio::spawn(async move {
5601            for _ in 0..50 {
5602                tokio::time::sleep(std::time::Duration::from_millis(10)).await;
5603                let ctl: Option<nmbrs_metrics::controls::Control<nmbrs_rate::RateSpec>> =
5604                    component_for_writer.read().unwrap().controls().get("rate");
5605                if let Some(c) = ctl {
5606                    // Only attempt once the applier is registered.
5607                    if c.applier_count() > 0 {
5608                        c.set(nmbrs_rate::RateSpec::new(10_000.0), ControlOrigin::Test)
5609                            .await
5610                            .ok();
5611                        return;
5612                    }
5613                }
5614            }
5615        });
5616
5617        let (adapter, _count) = CountingDriverAdapter::new();
5618        activity
5619            .run_with_driver(
5620                Arc::new(adapter),
5621                Arc::new(crate::synthesis::OpBuilder::new(test_kernel())),
5622            )
5623            .await;
5624        let _ = writer.await;
5625
5626        let guard = component.read().unwrap();
5627        let ctl: nmbrs_metrics::controls::Control<nmbrs_rate::RateSpec> =
5628            guard.controls().get("rate").unwrap();
5629        assert_eq!(ctl.value().ops_per_sec, 10_000.0);
5630    }
5631
5632    #[tokio::test]
5633    async fn concurrency_control_is_declared_on_attached_component() {
5634        // SRD 23 integration: the activity declares its
5635        // `concurrency` control on the attached component during
5636        // startup; the control's reified gauge reads the
5637        // configured value.
5638        use nmbrs_metrics::component::Component;
5639        use nmbrs_metrics::labels::Labels as L;
5640        use std::sync::RwLock;
5641
5642        let config = ActivityConfig {
5643            name: "ctrl_decl".into(),
5644            cycles: 10,
5645            concurrency: 3,
5646            ..Default::default()
5647        };
5648        let ops = vec![nmbrs_workload::model::ParsedOp::simple("op1", "test")];
5649        let seq = OpSequence::uniform(ops);
5650        let mut activity = Activity::new(config, &L::of("session", "s_decl"), seq);
5651        let component = Arc::new(RwLock::new(Component::new(
5652            L::of("session", "s_decl"),
5653            std::collections::HashMap::new(),
5654        )));
5655        activity.attach_component(component.clone());
5656
5657        let (adapter, _count) = CountingDriverAdapter::new();
5658        activity
5659            .run_with_driver(
5660                Arc::new(adapter),
5661                Arc::new(crate::synthesis::OpBuilder::new(test_kernel())),
5662            )
5663            .await;
5664
5665        // After run completes the control is still on the
5666        // component (structural declaration survives execution).
5667        let guard = component.read().unwrap();
5668        let erased = guard
5669            .controls()
5670            .get_erased("concurrency")
5671            .expect("concurrency control should be declared on attached component");
5672        assert_eq!(erased.value_string(), "3");
5673        assert!(erased.accepts_f64_writes());
5674        // Gauge projection reads as f64.
5675        assert_eq!(erased.gauge_f64(), Some(3.0));
5676    }
5677}