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(¶ms);
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}