Skip to main content

nmbrs_runtime/wrappers/
metrics.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Synthetic-metric recorder (SRD-40b §6). After the inner
5//! adapter returns, pulls each declared metric's value through
6//! the per-fiber op-template kernel via `ctx.wires.get` and
7//! records it onto the kind-specialised instrument
8//! (gauge / histogram / counter).
9
10use std::collections::HashMap;
11use std::sync::Arc;
12
13use crate::adapter::WrappingDispenser;
14use crate::adapter::{ExecutionError, OpDispenser, OpResult};
15use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
16
17/// SRD-32a wrapper name.
18pub const NAME: WrapperName = WrapperName::new("metrics");
19
20/// Trigger: op declares a non-empty `metrics:` map.
21fn triggers(s: WrapperSubject) -> bool {
22    let Some(template) = s.op() else {
23        return false;
24    };
25    !template.metrics.is_empty()
26}
27
28fn describe_assignment(s: WrapperSubject) -> Option<String> {
29    let template = s.op()?;
30    if template.metrics.is_empty() {
31        return None;
32    }
33    let mut names: Vec<&str> = template.metrics.keys().map(|s| s.as_str()).collect();
34    names.sort();
35    Some(format!("metrics: emits {}", names.join(", ")))
36}
37
38/// `metrics` must sit outside every other non-dryrun wrapper —
39/// it observes their per-cycle effects on `ctx.wires`. Listed
40/// as the `forbids_outer` set on metrics' registration so the
41/// constraint validator hard-errors if anything tries to slip
42/// outside it (except DRYRUN, which is strictly outermost via
43/// its own `forbids_outer` declaration).
44const FORBIDS_OUTER: &[WrapperName] = &[
45    super::traverse::NAME,
46    super::delay::NAME,
47    crate::validation::WRAPPER_NAME,
48    super::poll::NAME,
49    super::r#if::NAME,
50    // `emit` is INTENTIONALLY allowed outer of metrics — the
51    // emit wrapper is the operator-visible render surface and
52    // under `dryrun=emit` it sits outer of DRYRUN (and therefore
53    // outer of metrics) so the rendered op text reaches stdout
54    // before the short-circuit. Metrics fires inner of emit;
55    // the per-cycle measurement is unaffected.
56    super::result::NAME,
57];
58
59inventory::submit! {
60    WrapperRegistration {
61        name: NAME,
62        // `metrics:` is parsed into ParsedOp.metrics, not into
63        // params, so there's no owned `params`-key. Trigger
64        // fires only when the metrics map is non-empty.
65        owned_fields: &[],
66        triggers,
67        requires_inner: &[],
68        forbids_outer: FORBIDS_OUTER,
69        mutually_exclusive_with: &[],
70        describe_assignment,
71        levels: &[crate::wrapper_registry::WrapperLevel::Op],
72    }
73}
74
75/// Wraps an inner OpDispenser to publish per-cycle synthetic
76/// metrics declared in the op template's `metrics:` map
77/// (SRD-40b §6).
78///
79/// Per-cycle responsibilities (in order, matching SRD-40b §5.2 →
80/// §6 pipeline):
81///
82/// 1. Await the inner dispenser's `execute`. With a
83///    [`crate::wrappers::ResultDispenser`] in the wrapper stack
84///    between the inner adapter and this one, declared `result:`
85///    wires are already written through `ctx.wires.write` to the
86///    per-fiber kernel by the time we run.
87/// 2. For each declared metric, read the value through
88///    `ctx.wires.get(name)` (bare-binding-name canonical form per
89///    SRD-40b §1). The read pulls fresh through the eval cone
90///    so any computed output (e.g. `row_count := count`) reflects
91///    this cycle's value.
92/// 3. Apply the optional [`metric_format::FormatSpec`] sanitiser
93///    to round to the configured precision (Phase B).
94/// 4. Dispatch to the kind-specific instrument record method
95///    (SRD-40b §6.1):
96///    - [`MetricKind::Gauge`] → [`ValueGauge::set`] (f64).
97///    - [`MetricKind::Histogram`] → [`Histogram::record`]
98///      (truncated to u64 after format rounding).
99///    - [`MetricKind::Counter`] → [`Counter::inc_by`] (u64);
100///      non-positive values warn and skip — counters are
101///      monotonic by definition.
102///
103/// Non-bare-name expressions (`factor * 2.0`, `if(...)`, …) are
104/// **deferred** — the wrap step errors when `spec.value` is not
105/// a bare identifier.
106pub struct MetricsDispenser {
107    inner: Arc<dyn OpDispenser>,
108    /// One slot per declared metric. Stable ordering by metric
109    /// name keeps per-cycle dispatch deterministic for tests and
110    /// makes any per-cycle warning sequence reproducible.
111    slots: Arc<Vec<MetricSlot>>,
112}
113
114/// Lenient per-iteration GAUGE publication for poll drains: the
115/// dispenser above only fires when the (potentially hours-long)
116/// drain op completes, so the poll wrapper re-publishes gauge slots
117/// per poll iteration against that iteration's wires — the store
118/// then carries live samples (`compaction_progress_parts` etc.)
119/// instead of a single end-of-drain point. Gauges only: counters
120/// and histograms keep their one-record-per-op semantics (the
121/// end-of-op publish would double-count them). Resolution misses
122/// degrade to a debug log — status must never fail the poll.
123pub(crate) fn publish_gauges_lenient(slots: &[MetricSlot], wires: &dyn crate::wires::WireSource) {
124    for slot in slots {
125        // Cell-placed metrics are skipped: their instrument depends on a
126        // coordinate resolved from the cycle's wires, which this lenient
127        // path has no basis to choose.
128        let Some(MetricInstrument::Gauge(g)) = &slot.instrument else {
129            continue;
130        };
131        let Some(value) = wires.get(&slot.binding_name) else {
132            crate::diag!(
133                crate::observer::LogLevel::Debug,
134                "poll gauge '{}': binding '{}' unresolved this iteration",
135                slot.family,
136                slot.binding_name
137            );
138            continue;
139        };
140        let Some(raw) = value_to_f64(&value) else {
141            crate::diag!(
142                crate::observer::LogLevel::Debug,
143                "poll gauge '{}': '{}' non-numeric this iteration",
144                slot.family,
145                slot.value_expr
146            );
147            continue;
148        };
149        let sanitised = slot.format.as_ref().map(|f| f.apply(raw)).unwrap_or(raw);
150        g.set(sanitised);
151    }
152}
153
154/// One compiled metric slot: instrument storage + sanitiser +
155/// pre-bound Polydat pull handle.
156pub(crate) struct MetricSlot {
157    /// Family name registered with the [`Component`]. Used in
158    /// diagnostic messages (e.g. the counter non-positive warning).
159    family: String,
160    /// The original `value:` text from the workload, kept for
161    /// diagnostics. Per-cycle resolution reads through the
162    /// internal `binding_name` below; the user's original text
163    /// surfaces in error messages so operators see what they wrote.
164    value_expr: String,
165    /// Internal kernel-output name (`__metric_<name>`) the
166    /// op-template synthesiser created from this metric's
167    /// `value:` expression. Cycle-time reads go through
168    /// `ctx.wires.get(&binding_name)`.
169    binding_name: String,
170    /// Optional value sanitiser. Applied after the value is
171    /// pulled, before the instrument record.
172    format: Option<nmbrs_workload::metric_format::FormatSpec>,
173    /// Resolved instrument storage — exactly one variant is
174    /// populated per slot, matching `MetricSpec.kind`.
175    /// The instrument, for a metric that registers ONCE on the dispenser's
176    /// own component. `None` when the metric is cell-placed: its instruments
177    /// are materialised per coordinate in [`CellPlacement`], because each cell
178    /// is a distinct identity and therefore a distinct instrument.
179    instrument: Option<MetricInstrument>,
180    /// Dimensional placement, when the metric declared `cell:`.
181    placement: Option<CellPlacement>,
182}
183
184/// Per-coordinate instrument materialisation for a cell-placed metric.
185///
186/// The parent is the component this metric would otherwise have registered on,
187/// so a cell REFINES that identity rather than replacing part of it. Nothing
188/// ambient is consulted.
189pub(crate) struct CellPlacement {
190    /// `(dimension, synthesised coordinate wire)`, in the workload's declared
191    /// order. The wire is `__cell_<metric>__<dim>` — see
192    /// `crate::scope::synthesize_cell_binding_name`.
193    dims: Vec<(String, String)>,
194    /// The registration site whose identity each cell refines.
195    parent: Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
196    kind: nmbrs_workload::model::MetricKind,
197    unit: Option<String>,
198    /// Instruments already materialised, keyed by the coordinate's canonical
199    /// rendering. Steady state is a hash lookup: the registry write and the
200    /// component attach happen once per distinct coordinate, never per cycle.
201    instances: std::sync::Mutex<std::collections::HashMap<String, MetricInstrument>>,
202}
203
204/// Kind-specialised instrument storage owned by a [`MetricSlot`].
205///
206/// The same `Arc<...>` is shared with the dispenser's `Component`
207/// instrument registry — registered via
208/// `Component::register_instrument` at [`MetricsDispenser::wrap`] time.
209/// Per-cycle code records through the slot's typed `Arc`; the cadence
210/// reporter snapshots through the registry. One source of truth, two
211/// access paths.
212#[derive(Clone)]
213enum MetricInstrument {
214    Gauge(Arc<nmbrs_metrics::instruments::gauge::ValueGauge>),
215    Histogram(Arc<nmbrs_metrics::instruments::histogram::Histogram>),
216    Counter(Arc<nmbrs_metrics::instruments::counter::Counter>),
217}
218
219impl MetricInstrument {
220    /// Promote the kind-erased slot value into the canonical
221    /// [`InstrumentRef`] for registry storage.
222    fn as_ref(&self) -> nmbrs_metrics::component::InstrumentRef {
223        match self {
224            MetricInstrument::Gauge(g) => nmbrs_metrics::component::InstrumentRef::Gauge(g.clone()),
225            MetricInstrument::Histogram(h) => {
226                nmbrs_metrics::component::InstrumentRef::Histogram(h.clone())
227            }
228            MetricInstrument::Counter(c) => {
229                nmbrs_metrics::component::InstrumentRef::Counter(c.clone())
230            }
231        }
232    }
233}
234
235/// Numeric coercion for capture-map lookups. Returns `None`
236/// for non-numeric variants (string, vector, none) so the
237/// MetricsDispenser slot path logs + skips rather than panicking
238/// through `Value::as_f64`'s strict matcher.
239fn value_to_f64(v: &polydat::ast::Value) -> Option<f64> {
240    match v {
241        polydat::ast::Value::F64(f) => Some(*f),
242        polydat::ast::Value::U64(u) => Some(*u as f64),
243        polydat::ast::Value::Bool(b) => Some(if *b { 1.0 } else { 0.0 }),
244        _ => None,
245    }
246}
247
248impl MetricsDispenser {
249    /// Wrap an inner dispenser with synthetic-metric publication
250    /// for the op template's `metrics:` declarations.
251    ///
252    /// Init steps (SRD-40b §6 init):
253    /// 1. Empty declaration → return `inner` unchanged. No
254    ///    overhead for ops that don't publish synthetic metrics.
255    /// 2. For each `(name, spec)`, allocate the kind-specific
256    ///    instrument and register it on the component via
257    ///    `Component::register_instrument`. A duplicate-family
258    ///    collision (§7.2) errors here, before any cycle runs.
259    /// 3. Pre-parse the optional `format:` string into a
260    ///    [`FormatSpec`].
261    ///
262    /// `component` is borrowed mutably so `register_instrument`
263    /// can claim the family slot atomically with the instrument
264    /// allocation. The same `Arc<...>` is held both on the
265    /// component (for cadence-reporter capture) and in the
266    /// returned dispenser's slots (for per-cycle record).
267    pub fn wrap(
268        inner: Arc<dyn OpDispenser>,
269        metrics: &HashMap<String, nmbrs_workload::model::MetricSpec>,
270        component: &mut nmbrs_metrics::component::Component,
271        component_arc: &Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
272        fx: &mut crate::fixture::ScopeFixture,
273    ) -> Result<Arc<dyn OpDispenser>, String> {
274        Self::wrap_with_slots(inner, metrics, component, component_arc, fx).map(|(d, _)| d)
275    }
276
277    /// As [`Self::wrap`], additionally returning the compiled slot
278    /// list (shared `Arc`) so the poll wrapper can re-publish the
279    /// GAUGE slots per poll iteration (see
280    /// [`publish_gauges_lenient`]). `None` when the op declares no
281    /// metrics.
282    pub(crate) fn wrap_with_slots(
283        inner: Arc<dyn OpDispenser>,
284        metrics: &HashMap<String, nmbrs_workload::model::MetricSpec>,
285        component: &mut nmbrs_metrics::component::Component,
286        // The same component, as the handle cells attach under. A cell refines
287        // THIS identity, so the parent comes from the registration site rather
288        // than from anything ambient.
289        component_arc: &Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
290        fx: &mut crate::fixture::ScopeFixture,
291    ) -> Result<(Arc<dyn OpDispenser>, Option<Arc<Vec<MetricSlot>>>), String> {
292        if metrics.is_empty() {
293            return Ok((inner, None));
294        }
295        // Stable ordering on metric names so init-time
296        // diagnostics + per-cycle dispatch are reproducible.
297        let mut entries: Vec<_> = metrics.iter().collect();
298        entries.sort_by(|a, b| a.0.cmp(b.0));
299
300        let component_labels = component.effective_labels().clone();
301        let mut slots = Vec::with_capacity(entries.len());
302        for (name, spec) in entries {
303            let family = spec.family.clone().unwrap_or_else(|| name.clone());
304
305            let format = match &spec.format {
306                Some(s) => Some(
307                    nmbrs_workload::metric_format::parse_format_spec(s)
308                        .map_err(|e| format!("metric '{name}' format: {e}"))?,
309                ),
310                None => None,
311            };
312
313            let kind = spec.kind.unwrap_or_default();
314            // Instrument labels carry the family name as a label
315            // alongside the component's effective labels. The
316            // `family` argument to `register_instrument` is the
317            // canonical family-name string; the labels are the
318            // dimensional cell.
319            let instr_labels = component_labels.with("family", family.clone());
320            let instrument = match kind {
321                nmbrs_workload::model::MetricKind::Gauge => MetricInstrument::Gauge(Arc::new(
322                    nmbrs_metrics::instruments::gauge::ValueGauge::new(instr_labels),
323                )),
324                nmbrs_workload::model::MetricKind::Histogram => {
325                    MetricInstrument::Histogram(Arc::new(
326                        nmbrs_metrics::instruments::histogram::Histogram::new(instr_labels),
327                    ))
328                }
329                nmbrs_workload::model::MetricKind::Counter => MetricInstrument::Counter(Arc::new(
330                    nmbrs_metrics::instruments::counter::Counter::new(instr_labels),
331                )),
332            };
333
334            // Resolve the metric's value expression against the
335            // Polydat Kernel up front. The op-template synthesiser
336            // appended each metric's `value:` expression as a
337            // `__metric_<name> := <expr>` binding on the kernel
338            // (see `crate::scope::synthesize_metric_binding_name`),
339            // so cycle-time reads go through that internal output —
340            // arbitrary Polydat expressions work, not just bare names.
341            // The closure-binding-economy walker injected magic
342            // externs (body/count/ok) for any of those names this
343            // expression referenced, so a workload that writes
344            // `value: count` no longer needs a fake result-binding
345            // to wedge the slot open.
346            let binding_name = crate::scope::synthesize_metric_binding_name(name);
347            let _ = fx.register_pull(&binding_name).map_err(|e| {
348                format!(
349                    "metric '{name}' value '{value}': {e} (synthesised binding \
350                     '{binding_name}' should have been registered by the \
351                     op-template kernel synthesiser — this is a bug)",
352                    value = spec.value,
353                )
354            })?;
355
356            // SRD-40b §7.2 — collide-on-duplicate at init. The
357            // single registry on `Component` is the canonical
358            // store; the slot's `Arc<...>` shares the same
359            // instrument for the per-cycle hot path. The
360            // optional `unit` rides through to drive the
361            // `_<unit>` suffix on `metric_family.name` and the
362            // `unit` column at capture time (SRD-40a §4.3).
363            // Cell-placed metrics register PER CELL, at first sight of each
364            // coordinate — not here. Registering on the dispenser's own
365            // component as well would claim the family for the un-refined
366            // identity, and the first cell to materialise would then collide
367            // with it on this same duplicate-family check.
368            let placement = if spec.cell.is_empty() {
369                component.register_instrument_with_unit(
370                    family.clone(),
371                    spec.unit.clone(),
372                    instrument.as_ref(),
373                )?;
374                None
375            } else {
376                let mut dims = Vec::with_capacity(spec.cell.len());
377                for dim in spec.cell.keys() {
378                    let wire = crate::scope::synthesize_cell_binding_name(name, dim);
379                    let _ = fx.register_pull(&wire).map_err(|e| {
380                        format!(
381                            "metric '{name}' cell '{dim}': {e} (synthesised \
382                             coordinate binding '{wire}' should have been \
383                             registered by the op-template kernel synthesiser \
384                             — this is a bug)"
385                        )
386                    })?;
387                    dims.push((dim.clone(), wire));
388                }
389                Some(CellPlacement {
390                    dims,
391                    parent: component_arc.clone(),
392                    kind,
393                    unit: spec.unit.clone(),
394                    instances: std::sync::Mutex::new(std::collections::HashMap::new()),
395                })
396            };
397
398            slots.push(MetricSlot {
399                family,
400                value_expr: spec.value.clone(),
401                binding_name,
402                format,
403                instrument: if placement.is_some() {
404                    None
405                } else {
406                    Some(instrument)
407                },
408                placement,
409            });
410        }
411
412        let slots = Arc::new(slots);
413        Ok((
414            Arc::new(Self {
415                inner,
416                slots: slots.clone(),
417            }),
418            Some(slots),
419        ))
420    }
421}
422
423impl CellPlacement {
424    /// The instrument for this cycle's coordinate, materialising the cell and
425    /// registering the family on first sight of each distinct coordinate.
426    ///
427    /// Steady state is a hash lookup. The component attach and the registry
428    /// write happen once per coordinate, never per cycle — which is what keeps
429    /// a per-row metric off the component write lock.
430    fn resolve(
431        &self,
432        wires: &dyn crate::wires::WireSource,
433        family: &str,
434        cycle: u64,
435    ) -> Result<MetricInstrument, ExecutionError> {
436        let mut coord = nmbrs_metrics::labels::Labels::default();
437        for (dim, wire) in &self.dims {
438            let Some(value) = wires.get(wire) else {
439                return Err(ExecutionError::Op(crate::adapter::AdapterError {
440                    error_name: "metric_cell_unresolved".into(),
441                    message: format!(
442                        "metric '{family}' on cycle {cycle}: coordinate binding \
443                         '{wire}' for dimension '{dim}' did not resolve through \
444                         ctx.wires — this is a wiring bug between scope \
445                         synthesis and the metrics wrapper"
446                    ),
447                    retryable: false,
448                }));
449            };
450            // A label value is a string. Anything else would key cells on
451            // formatting rather than on identity, so it is refused here rather
452            // than stringified behind the author's back.
453            let polydat::ast::Value::Str(text) = &value else {
454                return Err(ExecutionError::Op(crate::adapter::AdapterError {
455                    error_name: "metric_cell_not_a_string".into(),
456                    message: format!(
457                        "metric '{family}' cell '{dim}' on cycle {cycle}: \
458                         coordinate resolved to a non-string {disc:?}. A \
459                         dimension's values are label values, which are \
460                         strings — convert the expression explicitly.",
461                        disc = std::mem::discriminant(&value)
462                    ),
463                    retryable: false,
464                }));
465            };
466            coord = coord.with(dim.clone(), text.to_string());
467        }
468
469        let key = coord.to_prometheus();
470        {
471            let cache = self.instances.lock().unwrap_or_else(|e| e.into_inner());
472            if let Some(found) = cache.get(&key) {
473                return Ok(found.clone());
474            }
475        }
476
477        // First sight of this coordinate: materialise the cell, build the
478        // instrument carrying the cell's FULL effective labels, and register
479        // the family there. The duplicate-family check runs per cell, which is
480        // exactly where it belongs.
481        let cell = nmbrs_metrics::cells::resolve_under(&self.parent, &coord);
482        let cell_labels = {
483            let g = cell.read().unwrap_or_else(|e| e.into_inner());
484            g.effective_labels().clone()
485        };
486        let instr_labels = cell_labels.with("family", family.to_string());
487        let instrument = match self.kind {
488            nmbrs_workload::model::MetricKind::Gauge => MetricInstrument::Gauge(Arc::new(
489                nmbrs_metrics::instruments::gauge::ValueGauge::new(instr_labels),
490            )),
491            nmbrs_workload::model::MetricKind::Histogram => MetricInstrument::Histogram(Arc::new(
492                nmbrs_metrics::instruments::histogram::Histogram::new(instr_labels),
493            )),
494            nmbrs_workload::model::MetricKind::Counter => MetricInstrument::Counter(Arc::new(
495                nmbrs_metrics::instruments::counter::Counter::new(instr_labels),
496            )),
497        };
498        {
499            let mut g = cell.write().unwrap_or_else(|e| e.into_inner());
500            g.register_instrument_with_unit(
501                family.to_string(),
502                self.unit.clone(),
503                instrument.as_ref(),
504            )
505            .map_err(|e| {
506                ExecutionError::Op(crate::adapter::AdapterError {
507                    error_name: "metric_cell_family_collision".into(),
508                    message: format!("metric '{family}' cell {key}: {e}"),
509                    retryable: false,
510                })
511            })?;
512        }
513        let mut cache = self.instances.lock().unwrap_or_else(|e| e.into_inner());
514        Ok(cache.entry(key).or_insert(instrument).clone())
515    }
516}
517
518/// A gauge slot wired straight to a wire name — the metric path a poll
519/// publishes through, without the workload/kernel machinery around it.
520/// Lets other modules' tests assert on what actually reaches the
521/// instrument, rather than on a display-side proxy.
522#[cfg(test)]
523pub(crate) fn test_gauge_slot(
524    family: &str,
525    binding_name: &str,
526) -> (
527    MetricSlot,
528    Arc<nmbrs_metrics::instruments::gauge::ValueGauge>,
529) {
530    let g = Arc::new(nmbrs_metrics::instruments::gauge::ValueGauge::new(
531        nmbrs_metrics::labels::Labels::default(),
532    ));
533    (
534        MetricSlot {
535            family: family.to_string(),
536            value_expr: binding_name.to_string(),
537            binding_name: binding_name.to_string(),
538            format: None,
539            instrument: Some(MetricInstrument::Gauge(g.clone())),
540            placement: None,
541        },
542        g,
543    )
544}
545
546impl WrappingDispenser for MetricsDispenser {}
547
548impl OpDispenser for MetricsDispenser {
549    fn execute<'a>(
550        &'a self,
551        cycle: u64,
552        ctx: &'a crate::fixture::ExecCtx<'a>,
553    ) -> std::pin::Pin<
554        Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
555    > {
556        Box::pin(async move {
557            let result = self.inner.execute(cycle, ctx).await?;
558            // Skipped ops produce no measurement — SRD-40b §5.2 /
559            // §6 pipeline only fires on a successfully-executed op.
560            if result.skipped {
561                return Ok(result);
562            }
563            for slot in self.slots.iter() {
564                // Sole resolution path: ctx.wires reads through the
565                // live per-fiber kernel handle (project rule
566                // "GK Is Canonical Scope"). The op-template synthesiser
567                // compiled this metric's `value:` expression into the
568                // kernel as `__metric_<name> := <expr>`; pulling that
569                // wire fires the expression's eval cone — including
570                // magic-extern reads (body/count/ok) ResultDispenser
571                // wrote earlier in the same cycle. ctx.wires.get is
572                // the live read; no pre-stack snapshot.
573                let Some(value) = ctx.wires.get(&slot.binding_name) else {
574                    return Err(ExecutionError::Op(crate::adapter::AdapterError {
575                        error_name: "metric_value_unresolved".into(),
576                        message: format!(
577                            "metric '{family}' on cycle {cycle}: synthesised \
578                             binding '{binding}' (from `value: {expr}`) did not \
579                             resolve through ctx.wires — this is a wiring bug \
580                             between scope synthesis and the metrics wrapper",
581                            family = slot.family,
582                            binding = slot.binding_name,
583                            expr = slot.value_expr,
584                        ),
585                        retryable: false,
586                    }));
587                };
588                // `None` is the absence of a measurement, not a bad one. A
589                // conditional aggregate that matched nothing this cycle
590                // (`:sum(x where kind='y')` with no matching rows) resolves to
591                // None by design — "absent, not zero". Publishing no sample is
592                // the honest response; erroring turns every conditional
593                // aggregate into a run-ending fault the first cycle its subject
594                // is idle, which is exactly when a drain is being watched.
595                //
596                // Genuinely non-numeric TYPES (Str, vector, handle) still fail
597                // loudly below — those are wiring mistakes, not absences.
598                if matches!(value, polydat::ast::Value::None) {
599                    continue;
600                }
601                let raw = match value_to_f64(&value) {
602                    Some(v) => v,
603                    None => {
604                        // The wire resolved but its type can't
605                        // coerce to a numeric metric value (Str,
606                        // vector, handle, etc.). Surface as a
607                        // hard ExecutionError so the activity's
608                        // `errors:` policy decides — by default
609                        // (errors=stop) the phase + run halt.
610                        return Err(ExecutionError::Op(crate::adapter::AdapterError {
611                            error_name: "metric_value_non_numeric".into(),
612                            message: format!(
613                                "metric '{family}' on cycle {cycle}: \
614                                 binding '{expr}' is not coercible to f64 \
615                                 (got value variant {disc:?}); metric \
616                                 values must be numeric (U64 / F64 / Bool)",
617                                family = slot.family,
618                                expr = slot.value_expr,
619                                disc = std::mem::discriminant(&value),
620                            ),
621                            retryable: false,
622                        }));
623                    }
624                };
625                let sanitised = slot.format.as_ref().map(|f| f.apply(raw)).unwrap_or(raw);
626                // A cell-placed metric resolves its coordinate from this
627                // cycle's wires and lands on the instrument for THAT cell.
628                let instrument = match (&slot.instrument, &slot.placement) {
629                    (Some(i), _) => std::borrow::Cow::Borrowed(i),
630                    (None, Some(p)) => match p.resolve(ctx.wires, &slot.family, cycle) {
631                        Ok(i) => std::borrow::Cow::Owned(i),
632                        Err(e) => return Err(e),
633                    },
634                    (None, None) => unreachable!("a slot has either an instrument or a placement"),
635                };
636                match instrument.as_ref() {
637                    MetricInstrument::Gauge(g) => g.set(sanitised),
638                    MetricInstrument::Histogram(h) => h.record(sanitised as u64),
639                    MetricInstrument::Counter(c) => {
640                        if sanitised <= 0.0 {
641                            crate::diag!(
642                                crate::observer::LogLevel::Warn,
643                                "counter '{}' got non-positive value {sanitised}; skipping",
644                                slot.family,
645                            );
646                        } else {
647                            c.inc_by(sanitised as u64);
648                        }
649                    }
650                }
651            }
652            Ok(result)
653        })
654    }
655    fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
656        Some(self.inner.as_ref())
657    }
658}
659
660#[cfg(test)]
661mod absent_value_tests {
662    /// A gauge whose binding resolved to `None` must publish NO SAMPLE, not
663    /// fail the run. `None` is the absence of a measurement: a conditional
664    /// aggregate like `:sum(progress where kind='secondary index build')`
665    /// yields it by design whenever no row has that kind, which is precisely
666    /// when a drain is idle and being watched. Erroring there killed a
667    /// 58-minute run at tier 19 with
668    /// `metric_value_non_numeric ... variant Discriminant(19)`.
669    ///
670    /// Genuinely non-numeric TYPES must still fail loudly — those are wiring
671    /// mistakes, not absences — so the skip is matched on `Value::None`
672    /// specifically, never on "failed to coerce".
673    #[test]
674    fn absent_binding_skips_the_sample_but_bad_types_still_fail() {
675        let src = std::fs::read_to_string(concat!(
676            env!("CARGO_MANIFEST_DIR"),
677            "/src/wrappers/metrics.rs"
678        ))
679        .expect("read own source");
680        let skip = src
681            .find("if matches!(value, polydat::ast::Value::None) {")
682            .expect("None must be skipped explicitly");
683        let coerce = src
684            .find("let raw = match value_to_f64(&value) {")
685            .expect("the coercion site must still exist");
686        assert!(
687            skip < coerce,
688            "the None skip must come BEFORE the coercion, or an absent value \
689             still reaches the error path"
690        );
691        assert!(
692            src.contains("metric_value_non_numeric"),
693            "non-numeric TYPES must still raise metric_value_non_numeric — \
694             the skip is for absence, not for bad wiring"
695        );
696    }
697}
698
699#[cfg(test)]
700mod tests {
701    use super::*;
702    use crate::adapter::{ExecutionError, OpResult};
703    use crate::fixture::ExecCtx;
704    use nmbrs_workload::model::{MetricKind, MetricSpec};
705
706    struct CapturesInner;
707    impl OpDispenser for CapturesInner {
708        fn execute<'a>(
709            &'a self,
710            _cycle: u64,
711            _ctx: &'a ExecCtx<'a>,
712        ) -> std::pin::Pin<
713            Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
714        > {
715            Box::pin(async move {
716                Ok(OpResult {
717                    body: None,
718                    skipped: false,
719                })
720            })
721        }
722    }
723
724    fn fresh_component() -> nmbrs_metrics::component::Component {
725        nmbrs_metrics::component::Component::new(
726            nmbrs_metrics::labels::Labels::empty(),
727            HashMap::new(),
728        )
729    }
730
731    /// A component as the production path holds it: an `Arc` (cells attach
732    /// under it) plus the `&mut` borrow registration needs.
733    fn fresh_component_arc() -> Arc<std::sync::RwLock<nmbrs_metrics::component::Component>> {
734        Arc::new(std::sync::RwLock::new(fresh_component()))
735    }
736
737    /// `MetricsDispenser::wrap` with the component handled the way the
738    /// activity does it.
739    fn wrap_on(
740        inner: Arc<dyn OpDispenser>,
741        decl: &HashMap<String, MetricSpec>,
742        comp: &Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
743        fx: &mut crate::fixture::ScopeFixture,
744    ) -> Result<Arc<dyn OpDispenser>, String> {
745        let mut guard = comp.write().unwrap();
746        MetricsDispenser::wrap(inner, decl, &mut guard, comp, fx)
747    }
748
749    fn fresh_fixture() -> crate::fixture::ScopeFixture {
750        use polydat::compile::assembly::{PolydatAssembler, WireRef};
751        use polydat::library::identity::Identity;
752        let mut asm = PolydatAssembler::new(vec!["cycle".into()]);
753        asm.add_node(
754            "cycle_id",
755            Box::new(Identity::new(polydat::ast::PortType::U64)),
756            vec![WireRef::input("cycle")],
757        );
758        asm.add_output("cycle_id", WireRef::node("cycle_id"));
759        let kernel = asm.compile().expect("test fixture asm.compile");
760        crate::fixture::ScopeFixture::new(kernel.program().clone())
761    }
762
763    fn make_spec(value: &str, kind: MetricKind, format: Option<&str>) -> MetricSpec {
764        MetricSpec {
765            cell: Default::default(),
766            value: value.to_string(),
767            family: None,
768            kind: Some(kind),
769            unit: None,
770            format: format.map(|s| s.to_string()),
771        }
772    }
773
774    #[test]
775    fn metrics_dispenser_empty_returns_inner_unchanged() {
776        let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
777        let inner_ptr = Arc::as_ptr(&inner);
778        let comp = fresh_component_arc();
779        let mut fx = fresh_fixture();
780        let wrapped = wrap_on(inner.clone(), &HashMap::new(), &comp, &mut fx).unwrap();
781        assert_eq!(Arc::as_ptr(&wrapped), inner_ptr);
782    }
783
784    /// Test-only introspection: peek at allocated instrument
785    /// `Arc`s by family name. Tests need this because `wrap`
786    /// returns `Arc<dyn OpDispenser>` and we want to assert
787    /// against the same `ValueGauge` / `Histogram` / `Counter`
788    /// the wrapper writes through.
789    impl MetricsDispenser {
790        fn slot_gauge(
791            &self,
792            family: &str,
793        ) -> Option<Arc<nmbrs_metrics::instruments::gauge::ValueGauge>> {
794            self.slots
795                .iter()
796                .find(|s| s.family == family)
797                .and_then(|s| match s.instrument.as_ref()? {
798                    MetricInstrument::Gauge(g) => Some(g.clone()),
799                    _ => None,
800                })
801        }
802        fn slot_histogram(
803            &self,
804            family: &str,
805        ) -> Option<Arc<nmbrs_metrics::instruments::histogram::Histogram>> {
806            self.slots
807                .iter()
808                .find(|s| s.family == family)
809                .and_then(|s| match s.instrument.as_ref()? {
810                    MetricInstrument::Histogram(h) => Some(h.clone()),
811                    _ => None,
812                })
813        }
814        fn slot_counter(
815            &self,
816            family: &str,
817        ) -> Option<Arc<nmbrs_metrics::instruments::counter::Counter>> {
818            self.slots
819                .iter()
820                .find(|s| s.family == family)
821                .and_then(|s| match s.instrument.as_ref()? {
822                    MetricInstrument::Counter(c) => Some(c.clone()),
823                    _ => None,
824                })
825        }
826    }
827
828    fn kernel_with_const_outputs(
829        consts: &[(&str, f64)],
830    ) -> (
831        crate::scope_kernel::ScopeKernel,
832        crate::fixture::ScopeFixture,
833    ) {
834        use polydat::compile::assembly::{PolydatAssembler, WireRef};
835        use polydat::library::fixed::ConstF64;
836        let mut asm = PolydatAssembler::new(vec!["cycle".into()]);
837        for (name, val) in consts {
838            let binding = crate::scope::synthesize_metric_binding_name(name);
839            asm.add_node(&binding, Box::new(ConstF64::new(*val)), vec![]);
840            asm.add_output(&binding, WireRef::node(&binding));
841        }
842        let kernel =
843            crate::scope_kernel::ScopeKernel::from(asm.compile().expect("test kernel asm.compile"));
844        let fx = crate::fixture::ScopeFixture::new(kernel.program().clone());
845        (kernel, fx)
846    }
847
848    fn typed_wrap_with_kernel(
849        inner: Arc<dyn OpDispenser>,
850        decls: &HashMap<String, MetricSpec>,
851        consts: &[(&str, f64)],
852    ) -> Result<
853        (
854            Arc<MetricsDispenser>,
855            crate::fixture::ResolvedPulls,
856            crate::scope_kernel::ScopeKernel,
857        ),
858        String,
859    > {
860        let (mut kernel, mut fx) = kernel_with_const_outputs(consts);
861        let mut comp = fresh_component();
862
863        if decls.is_empty() {
864            return Err("typed_wrap_with_kernel requires non-empty decls".into());
865        }
866        let mut entries: Vec<_> = decls.iter().collect();
867        entries.sort_by(|a, b| a.0.cmp(b.0));
868        let component_labels = comp.effective_labels().clone();
869        let mut slots = Vec::with_capacity(entries.len());
870        for (name, spec) in entries {
871            let family = spec.family.clone().unwrap_or_else(|| name.clone());
872            let format = match &spec.format {
873                Some(s) => Some(
874                    nmbrs_workload::metric_format::parse_format_spec(s)
875                        .map_err(|e| format!("metric '{name}' format: {e}"))?,
876                ),
877                None => None,
878            };
879            let kind = spec.kind.unwrap_or_default();
880            let instr_labels = component_labels.with("family", family.clone());
881            let instrument = match kind {
882                MetricKind::Gauge => MetricInstrument::Gauge(Arc::new(
883                    nmbrs_metrics::instruments::gauge::ValueGauge::new(instr_labels),
884                )),
885                MetricKind::Histogram => MetricInstrument::Histogram(Arc::new(
886                    nmbrs_metrics::instruments::histogram::Histogram::new(instr_labels),
887                )),
888                MetricKind::Counter => MetricInstrument::Counter(Arc::new(
889                    nmbrs_metrics::instruments::counter::Counter::new(instr_labels),
890                )),
891            };
892            comp.register_instrument_with_unit(
893                family.clone(),
894                spec.unit.clone(),
895                instrument.as_ref(),
896            )?;
897            let binding_name = crate::scope::synthesize_metric_binding_name(name);
898            let _ = fx.register_pull(&binding_name)?;
899            slots.push(MetricSlot {
900                family,
901                value_expr: spec.value.clone(),
902                binding_name,
903                format,
904                instrument: Some(instrument),
905                placement: None,
906            });
907        }
908        let typed = Arc::new(MetricsDispenser {
909            inner,
910            slots: Arc::new(slots),
911        });
912
913        let plan = fx.seal();
914        kernel.set_inputs(&[0]);
915        let pulls = plan.resolve_with(&mut kernel);
916        Ok((typed, pulls, kernel))
917    }
918
919    fn run_dispenser(
920        dispenser: Arc<dyn OpDispenser>,
921        pulls: &crate::fixture::ResolvedPulls,
922        kernel: &mut crate::scope_kernel::ScopeKernel,
923    ) -> Result<OpResult, ExecutionError> {
924        let fields = crate::adapter::ResolvedFields::new(vec![], vec![]);
925        let cw = crate::wires::CycleWires::new(kernel);
926        let ctx = ExecCtx::with_wires(&fields, pulls, &cw);
927        let rt = tokio::runtime::Builder::new_current_thread()
928            .build()
929            .unwrap();
930        rt.block_on(dispenser.execute(0, &ctx))
931    }
932
933    #[test]
934    fn metrics_dispenser_gauge_records_f64() {
935        let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
936        let mut decl = HashMap::new();
937        decl.insert(
938            "my_factor".into(),
939            make_spec("my_factor", MetricKind::Gauge, None),
940        );
941
942        let (typed, pulls, mut kernel) =
943            typed_wrap_with_kernel(inner, &decl, &[("my_factor", 3.5)]).unwrap();
944        let gauge = typed.slot_gauge("my_factor").unwrap();
945        run_dispenser(typed.clone() as Arc<dyn OpDispenser>, &pulls, &mut kernel).unwrap();
946
947        assert!((gauge.get() - 3.5).abs() < 1e-9);
948    }
949
950    #[test]
951    fn metrics_dispenser_histogram_truncates_to_u64() {
952        let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
953        let mut decl = HashMap::new();
954        decl.insert(
955            "latency_ms".into(),
956            make_spec("latency_ms", MetricKind::Histogram, None),
957        );
958
959        let (typed, pulls, mut kernel) =
960            typed_wrap_with_kernel(inner, &decl, &[("latency_ms", 7.9)]).unwrap();
961        let hist = typed.slot_histogram("latency_ms").unwrap();
962        run_dispenser(typed.clone() as Arc<dyn OpDispenser>, &pulls, &mut kernel).unwrap();
963
964        let snap = hist.peek_snapshot();
965        assert_eq!(snap.max(), 7);
966        assert_eq!(snap.len(), 1);
967    }
968
969    #[test]
970    fn metrics_dispenser_counter_positive_inc_and_skip_non_positive() {
971        let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
972        let mut decl = HashMap::new();
973        decl.insert(
974            "ok_inc".into(),
975            make_spec("ok_inc", MetricKind::Counter, None),
976        );
977        decl.insert(
978            "skip_inc".into(),
979            make_spec("skip_inc", MetricKind::Counter, None),
980        );
981
982        let (typed, pulls, mut kernel) =
983            typed_wrap_with_kernel(inner, &decl, &[("ok_inc", 5.0), ("skip_inc", 0.0)]).unwrap();
984        let ok_counter = typed.slot_counter("ok_inc").unwrap();
985        let skip_counter = typed.slot_counter("skip_inc").unwrap();
986        run_dispenser(typed.clone() as Arc<dyn OpDispenser>, &pulls, &mut kernel).unwrap();
987
988        assert_eq!(ok_counter.get(), 5);
989        assert_eq!(skip_counter.get(), 0);
990    }
991
992    #[test]
993    fn metrics_dispenser_format_rounds_value() {
994        let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
995        let mut decl = HashMap::new();
996        decl.insert(
997            "ratio".into(),
998            make_spec("ratio", MetricKind::Gauge, Some("#.##")),
999        );
1000
1001        let (typed, pulls, mut kernel) =
1002            typed_wrap_with_kernel(inner, &decl, &[("ratio", 1.234)]).unwrap();
1003        let gauge = typed.slot_gauge("ratio").unwrap();
1004        run_dispenser(typed.clone() as Arc<dyn OpDispenser>, &pulls, &mut kernel).unwrap();
1005
1006        assert!((gauge.get() - 1.23).abs() < 1e-9);
1007    }
1008
1009    #[test]
1010    fn metrics_dispenser_duplicate_family_errors() {
1011        let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
1012        let comp = fresh_component_arc();
1013        comp.write()
1014            .unwrap()
1015            .register_instrument(
1016                "recall_at_10",
1017                nmbrs_metrics::component::InstrumentRef::Counter(Arc::new(
1018                    nmbrs_metrics::instruments::counter::Counter::new(
1019                        nmbrs_metrics::labels::Labels::of("name", "recall_at_10"),
1020                    ),
1021                )),
1022            )
1023            .unwrap();
1024
1025        let mut decl = HashMap::new();
1026        decl.insert(
1027            "recall_at_10".into(),
1028            make_spec("recall_at_10", MetricKind::Gauge, None),
1029        );
1030
1031        let (_kernel, mut fx) = kernel_with_const_outputs(&[("recall_at_10", 0.0)]);
1032        let err = match wrap_on(inner, &decl, &comp, &mut fx) {
1033            Ok(_) => panic!("expected duplicate-family error, got Ok"),
1034            Err(e) => e,
1035        };
1036        assert!(
1037            err.contains("duplicate family name"),
1038            "unexpected error: {err}"
1039        );
1040    }
1041
1042    #[test]
1043    fn metrics_dispenser_skipped_op_records_nothing() {
1044        struct SkipInner;
1045        impl OpDispenser for SkipInner {
1046            fn execute<'a>(
1047                &'a self,
1048                _cycle: u64,
1049                _ctx: &'a ExecCtx<'a>,
1050            ) -> std::pin::Pin<
1051                Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
1052            > {
1053                Box::pin(async move { Ok(OpResult::skipped()) })
1054            }
1055        }
1056        let mut decl = HashMap::new();
1057        decl.insert("g".into(), make_spec("g", MetricKind::Gauge, None));
1058
1059        let (typed, pulls, mut kernel) =
1060            typed_wrap_with_kernel(Arc::new(SkipInner), &decl, &[("g", 1.0)]).unwrap();
1061        let gauge = typed.slot_gauge("g").unwrap();
1062
1063        let res =
1064            run_dispenser(typed.clone() as Arc<dyn OpDispenser>, &pulls, &mut kernel).unwrap();
1065        assert!(res.skipped);
1066        assert_eq!(gauge.get(), 0.0);
1067    }
1068
1069    #[test]
1070    fn metrics_dispenser_accepts_arbitrary_polydat_expression() {
1071        let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
1072        let mut decl = HashMap::new();
1073        decl.insert(
1074            "computed".into(),
1075            make_spec("factor * 2.0", MetricKind::Gauge, None),
1076        );
1077
1078        let (mut kernel, mut fx) = kernel_with_const_outputs(&[("computed", 6.0)]);
1079        let comp = fresh_component_arc();
1080        let _ = wrap_on(inner, &decl, &comp, &mut fx)
1081            .expect("arbitrary Polydat expression should wrap cleanly");
1082        let plan = fx.seal();
1083        kernel.set_inputs(&[0]);
1084        let _pulls = plan.resolve_with(&mut kernel);
1085    }
1086
1087    #[test]
1088    fn metrics_dispenser_missing_wire_errors_at_init() {
1089        let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
1090        let mut decl = HashMap::new();
1091        decl.insert(
1092            "missing_metric".into(),
1093            make_spec("absent_wire", MetricKind::Gauge, None),
1094        );
1095
1096        let (_kernel, mut fx) = kernel_with_const_outputs(&[("present", 1.0)]);
1097        let comp = fresh_component_arc();
1098        let err = wrap_on(inner, &decl, &comp, &mut fx)
1099            .err()
1100            .expect("missing-wire metric should error at init");
1101        assert!(err.contains("absent_wire"), "msg: {err}");
1102        assert!(err.contains("Available"), "msg: {err}");
1103    }
1104
1105    /// Silence dead-code warning for value_to_f64 (used inside
1106    /// wrapper but private to this module otherwise).
1107    #[test]
1108    fn value_to_f64_smoke() {
1109        assert_eq!(value_to_f64(&polydat::ast::Value::U64(5)), Some(5.0));
1110        assert_eq!(value_to_f64(&polydat::ast::Value::Bool(true)), Some(1.0));
1111        assert_eq!(value_to_f64(&polydat::ast::Value::Str("x".into())), None);
1112    }
1113}