use std::collections::HashMap;
use std::sync::Arc;
use crate::adapter::WrappingDispenser;
use crate::adapter::{ExecutionError, OpDispenser, OpResult};
use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
pub const NAME: WrapperName = WrapperName::new("metrics");
fn triggers(s: WrapperSubject) -> bool {
let Some(template) = s.op() else {
return false;
};
!template.metrics.is_empty()
}
fn describe_assignment(s: WrapperSubject) -> Option<String> {
let template = s.op()?;
if template.metrics.is_empty() {
return None;
}
let mut names: Vec<&str> = template.metrics.keys().map(|s| s.as_str()).collect();
names.sort();
Some(format!("metrics: emits {}", names.join(", ")))
}
const FORBIDS_OUTER: &[WrapperName] = &[
super::traverse::NAME,
super::delay::NAME,
crate::validation::WRAPPER_NAME,
super::poll::NAME,
super::r#if::NAME,
super::result::NAME,
];
inventory::submit! {
WrapperRegistration {
name: NAME,
owned_fields: &[],
triggers,
requires_inner: &[],
forbids_outer: FORBIDS_OUTER,
mutually_exclusive_with: &[],
describe_assignment,
levels: &[crate::wrapper_registry::WrapperLevel::Op],
}
}
pub struct MetricsDispenser {
inner: Arc<dyn OpDispenser>,
slots: Arc<Vec<MetricSlot>>,
}
pub(crate) fn publish_gauges_lenient(slots: &[MetricSlot], wires: &dyn crate::wires::WireSource) {
for slot in slots {
let Some(MetricInstrument::Gauge(g)) = &slot.instrument else {
continue;
};
let Some(value) = wires.get(&slot.binding_name) else {
crate::diag!(
crate::observer::LogLevel::Debug,
"poll gauge '{}': binding '{}' unresolved this iteration",
slot.family,
slot.binding_name
);
continue;
};
let Some(raw) = value_to_f64(&value) else {
crate::diag!(
crate::observer::LogLevel::Debug,
"poll gauge '{}': '{}' non-numeric this iteration",
slot.family,
slot.value_expr
);
continue;
};
let sanitised = slot.format.as_ref().map(|f| f.apply(raw)).unwrap_or(raw);
g.set(sanitised);
}
}
pub(crate) struct MetricSlot {
family: String,
value_expr: String,
binding_name: String,
format: Option<nmbrs_workload::metric_format::FormatSpec>,
instrument: Option<MetricInstrument>,
placement: Option<CellPlacement>,
}
pub(crate) struct CellPlacement {
dims: Vec<(String, String)>,
parent: Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
kind: nmbrs_workload::model::MetricKind,
unit: Option<String>,
instances: std::sync::Mutex<std::collections::HashMap<String, MetricInstrument>>,
}
#[derive(Clone)]
enum MetricInstrument {
Gauge(Arc<nmbrs_metrics::instruments::gauge::ValueGauge>),
Histogram(Arc<nmbrs_metrics::instruments::histogram::Histogram>),
Counter(Arc<nmbrs_metrics::instruments::counter::Counter>),
}
impl MetricInstrument {
fn as_ref(&self) -> nmbrs_metrics::component::InstrumentRef {
match self {
MetricInstrument::Gauge(g) => nmbrs_metrics::component::InstrumentRef::Gauge(g.clone()),
MetricInstrument::Histogram(h) => {
nmbrs_metrics::component::InstrumentRef::Histogram(h.clone())
}
MetricInstrument::Counter(c) => {
nmbrs_metrics::component::InstrumentRef::Counter(c.clone())
}
}
}
}
fn value_to_f64(v: &polydat::ast::Value) -> Option<f64> {
match v {
polydat::ast::Value::F64(f) => Some(*f),
polydat::ast::Value::U64(u) => Some(*u as f64),
polydat::ast::Value::Bool(b) => Some(if *b { 1.0 } else { 0.0 }),
_ => None,
}
}
impl MetricsDispenser {
pub fn wrap(
inner: Arc<dyn OpDispenser>,
metrics: &HashMap<String, nmbrs_workload::model::MetricSpec>,
component: &mut nmbrs_metrics::component::Component,
component_arc: &Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
fx: &mut crate::fixture::ScopeFixture,
) -> Result<Arc<dyn OpDispenser>, String> {
Self::wrap_with_slots(inner, metrics, component, component_arc, fx).map(|(d, _)| d)
}
pub(crate) fn wrap_with_slots(
inner: Arc<dyn OpDispenser>,
metrics: &HashMap<String, nmbrs_workload::model::MetricSpec>,
component: &mut nmbrs_metrics::component::Component,
component_arc: &Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
fx: &mut crate::fixture::ScopeFixture,
) -> Result<(Arc<dyn OpDispenser>, Option<Arc<Vec<MetricSlot>>>), String> {
if metrics.is_empty() {
return Ok((inner, None));
}
let mut entries: Vec<_> = metrics.iter().collect();
entries.sort_by(|a, b| a.0.cmp(b.0));
let component_labels = component.effective_labels().clone();
let mut slots = Vec::with_capacity(entries.len());
for (name, spec) in entries {
let family = spec.family.clone().unwrap_or_else(|| name.clone());
let format = match &spec.format {
Some(s) => Some(
nmbrs_workload::metric_format::parse_format_spec(s)
.map_err(|e| format!("metric '{name}' format: {e}"))?,
),
None => None,
};
let kind = spec.kind.unwrap_or_default();
let instr_labels = component_labels.with("family", family.clone());
let instrument = match kind {
nmbrs_workload::model::MetricKind::Gauge => MetricInstrument::Gauge(Arc::new(
nmbrs_metrics::instruments::gauge::ValueGauge::new(instr_labels),
)),
nmbrs_workload::model::MetricKind::Histogram => {
MetricInstrument::Histogram(Arc::new(
nmbrs_metrics::instruments::histogram::Histogram::new(instr_labels),
))
}
nmbrs_workload::model::MetricKind::Counter => MetricInstrument::Counter(Arc::new(
nmbrs_metrics::instruments::counter::Counter::new(instr_labels),
)),
};
let binding_name = crate::scope::synthesize_metric_binding_name(name);
let _ = fx.register_pull(&binding_name).map_err(|e| {
format!(
"metric '{name}' value '{value}': {e} (synthesised binding \
'{binding_name}' should have been registered by the \
op-template kernel synthesiser โ this is a bug)",
value = spec.value,
)
})?;
let placement = if spec.cell.is_empty() {
component.register_instrument_with_unit(
family.clone(),
spec.unit.clone(),
instrument.as_ref(),
)?;
None
} else {
let mut dims = Vec::with_capacity(spec.cell.len());
for dim in spec.cell.keys() {
let wire = crate::scope::synthesize_cell_binding_name(name, dim);
let _ = fx.register_pull(&wire).map_err(|e| {
format!(
"metric '{name}' cell '{dim}': {e} (synthesised \
coordinate binding '{wire}' should have been \
registered by the op-template kernel synthesiser \
โ this is a bug)"
)
})?;
dims.push((dim.clone(), wire));
}
Some(CellPlacement {
dims,
parent: component_arc.clone(),
kind,
unit: spec.unit.clone(),
instances: std::sync::Mutex::new(std::collections::HashMap::new()),
})
};
slots.push(MetricSlot {
family,
value_expr: spec.value.clone(),
binding_name,
format,
instrument: if placement.is_some() {
None
} else {
Some(instrument)
},
placement,
});
}
let slots = Arc::new(slots);
Ok((
Arc::new(Self {
inner,
slots: slots.clone(),
}),
Some(slots),
))
}
}
impl CellPlacement {
fn resolve(
&self,
wires: &dyn crate::wires::WireSource,
family: &str,
cycle: u64,
) -> Result<MetricInstrument, ExecutionError> {
let mut coord = nmbrs_metrics::labels::Labels::default();
for (dim, wire) in &self.dims {
let Some(value) = wires.get(wire) else {
return Err(ExecutionError::Op(crate::adapter::AdapterError {
error_name: "metric_cell_unresolved".into(),
message: format!(
"metric '{family}' on cycle {cycle}: coordinate binding \
'{wire}' for dimension '{dim}' did not resolve through \
ctx.wires โ this is a wiring bug between scope \
synthesis and the metrics wrapper"
),
retryable: false,
}));
};
let polydat::ast::Value::Str(text) = &value else {
return Err(ExecutionError::Op(crate::adapter::AdapterError {
error_name: "metric_cell_not_a_string".into(),
message: format!(
"metric '{family}' cell '{dim}' on cycle {cycle}: \
coordinate resolved to a non-string {disc:?}. A \
dimension's values are label values, which are \
strings โ convert the expression explicitly.",
disc = std::mem::discriminant(&value)
),
retryable: false,
}));
};
coord = coord.with(dim.clone(), text.to_string());
}
let key = coord.to_prometheus();
{
let cache = self.instances.lock().unwrap_or_else(|e| e.into_inner());
if let Some(found) = cache.get(&key) {
return Ok(found.clone());
}
}
let cell = nmbrs_metrics::cells::resolve_under(&self.parent, &coord);
let cell_labels = {
let g = cell.read().unwrap_or_else(|e| e.into_inner());
g.effective_labels().clone()
};
let instr_labels = cell_labels.with("family", family.to_string());
let instrument = match self.kind {
nmbrs_workload::model::MetricKind::Gauge => MetricInstrument::Gauge(Arc::new(
nmbrs_metrics::instruments::gauge::ValueGauge::new(instr_labels),
)),
nmbrs_workload::model::MetricKind::Histogram => MetricInstrument::Histogram(Arc::new(
nmbrs_metrics::instruments::histogram::Histogram::new(instr_labels),
)),
nmbrs_workload::model::MetricKind::Counter => MetricInstrument::Counter(Arc::new(
nmbrs_metrics::instruments::counter::Counter::new(instr_labels),
)),
};
{
let mut g = cell.write().unwrap_or_else(|e| e.into_inner());
g.register_instrument_with_unit(
family.to_string(),
self.unit.clone(),
instrument.as_ref(),
)
.map_err(|e| {
ExecutionError::Op(crate::adapter::AdapterError {
error_name: "metric_cell_family_collision".into(),
message: format!("metric '{family}' cell {key}: {e}"),
retryable: false,
})
})?;
}
let mut cache = self.instances.lock().unwrap_or_else(|e| e.into_inner());
Ok(cache.entry(key).or_insert(instrument).clone())
}
}
#[cfg(test)]
pub(crate) fn test_gauge_slot(
family: &str,
binding_name: &str,
) -> (
MetricSlot,
Arc<nmbrs_metrics::instruments::gauge::ValueGauge>,
) {
let g = Arc::new(nmbrs_metrics::instruments::gauge::ValueGauge::new(
nmbrs_metrics::labels::Labels::default(),
));
(
MetricSlot {
family: family.to_string(),
value_expr: binding_name.to_string(),
binding_name: binding_name.to_string(),
format: None,
instrument: Some(MetricInstrument::Gauge(g.clone())),
placement: None,
},
g,
)
}
impl WrappingDispenser for MetricsDispenser {}
impl OpDispenser for MetricsDispenser {
fn execute<'a>(
&'a self,
cycle: u64,
ctx: &'a crate::fixture::ExecCtx<'a>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
> {
Box::pin(async move {
let result = self.inner.execute(cycle, ctx).await?;
if result.skipped {
return Ok(result);
}
for slot in self.slots.iter() {
let Some(value) = ctx.wires.get(&slot.binding_name) else {
return Err(ExecutionError::Op(crate::adapter::AdapterError {
error_name: "metric_value_unresolved".into(),
message: format!(
"metric '{family}' on cycle {cycle}: synthesised \
binding '{binding}' (from `value: {expr}`) did not \
resolve through ctx.wires โ this is a wiring bug \
between scope synthesis and the metrics wrapper",
family = slot.family,
binding = slot.binding_name,
expr = slot.value_expr,
),
retryable: false,
}));
};
if matches!(value, polydat::ast::Value::None) {
continue;
}
let raw = match value_to_f64(&value) {
Some(v) => v,
None => {
return Err(ExecutionError::Op(crate::adapter::AdapterError {
error_name: "metric_value_non_numeric".into(),
message: format!(
"metric '{family}' on cycle {cycle}: \
binding '{expr}' is not coercible to f64 \
(got value variant {disc:?}); metric \
values must be numeric (U64 / F64 / Bool)",
family = slot.family,
expr = slot.value_expr,
disc = std::mem::discriminant(&value),
),
retryable: false,
}));
}
};
let sanitised = slot.format.as_ref().map(|f| f.apply(raw)).unwrap_or(raw);
let instrument = match (&slot.instrument, &slot.placement) {
(Some(i), _) => std::borrow::Cow::Borrowed(i),
(None, Some(p)) => match p.resolve(ctx.wires, &slot.family, cycle) {
Ok(i) => std::borrow::Cow::Owned(i),
Err(e) => return Err(e),
},
(None, None) => unreachable!("a slot has either an instrument or a placement"),
};
match instrument.as_ref() {
MetricInstrument::Gauge(g) => g.set(sanitised),
MetricInstrument::Histogram(h) => h.record(sanitised as u64),
MetricInstrument::Counter(c) => {
if sanitised <= 0.0 {
crate::diag!(
crate::observer::LogLevel::Warn,
"counter '{}' got non-positive value {sanitised}; skipping",
slot.family,
);
} else {
c.inc_by(sanitised as u64);
}
}
}
}
Ok(result)
})
}
fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
Some(self.inner.as_ref())
}
}
#[cfg(test)]
mod absent_value_tests {
#[test]
fn absent_binding_skips_the_sample_but_bad_types_still_fail() {
let src = std::fs::read_to_string(concat!(
env!("CARGO_MANIFEST_DIR"),
"/src/wrappers/metrics.rs"
))
.expect("read own source");
let skip = src
.find("if matches!(value, polydat::ast::Value::None) {")
.expect("None must be skipped explicitly");
let coerce = src
.find("let raw = match value_to_f64(&value) {")
.expect("the coercion site must still exist");
assert!(
skip < coerce,
"the None skip must come BEFORE the coercion, or an absent value \
still reaches the error path"
);
assert!(
src.contains("metric_value_non_numeric"),
"non-numeric TYPES must still raise metric_value_non_numeric โ \
the skip is for absence, not for bad wiring"
);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::adapter::{ExecutionError, OpResult};
use crate::fixture::ExecCtx;
use nmbrs_workload::model::{MetricKind, MetricSpec};
struct CapturesInner;
impl OpDispenser for CapturesInner {
fn execute<'a>(
&'a self,
_cycle: u64,
_ctx: &'a ExecCtx<'a>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
> {
Box::pin(async move {
Ok(OpResult {
body: None,
skipped: false,
})
})
}
}
fn fresh_component() -> nmbrs_metrics::component::Component {
nmbrs_metrics::component::Component::new(
nmbrs_metrics::labels::Labels::empty(),
HashMap::new(),
)
}
fn fresh_component_arc() -> Arc<std::sync::RwLock<nmbrs_metrics::component::Component>> {
Arc::new(std::sync::RwLock::new(fresh_component()))
}
fn wrap_on(
inner: Arc<dyn OpDispenser>,
decl: &HashMap<String, MetricSpec>,
comp: &Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
fx: &mut crate::fixture::ScopeFixture,
) -> Result<Arc<dyn OpDispenser>, String> {
let mut guard = comp.write().unwrap();
MetricsDispenser::wrap(inner, decl, &mut guard, comp, fx)
}
fn fresh_fixture() -> crate::fixture::ScopeFixture {
use polydat::compile::assembly::{PolydatAssembler, WireRef};
use polydat::library::identity::Identity;
let mut asm = PolydatAssembler::new(vec!["cycle".into()]);
asm.add_node(
"cycle_id",
Box::new(Identity::new(polydat::ast::PortType::U64)),
vec![WireRef::input("cycle")],
);
asm.add_output("cycle_id", WireRef::node("cycle_id"));
let kernel = asm.compile().expect("test fixture asm.compile");
crate::fixture::ScopeFixture::new(kernel.program().clone())
}
fn make_spec(value: &str, kind: MetricKind, format: Option<&str>) -> MetricSpec {
MetricSpec {
cell: Default::default(),
value: value.to_string(),
family: None,
kind: Some(kind),
unit: None,
format: format.map(|s| s.to_string()),
}
}
#[test]
fn metrics_dispenser_empty_returns_inner_unchanged() {
let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
let inner_ptr = Arc::as_ptr(&inner);
let comp = fresh_component_arc();
let mut fx = fresh_fixture();
let wrapped = wrap_on(inner.clone(), &HashMap::new(), &comp, &mut fx).unwrap();
assert_eq!(Arc::as_ptr(&wrapped), inner_ptr);
}
impl MetricsDispenser {
fn slot_gauge(
&self,
family: &str,
) -> Option<Arc<nmbrs_metrics::instruments::gauge::ValueGauge>> {
self.slots
.iter()
.find(|s| s.family == family)
.and_then(|s| match s.instrument.as_ref()? {
MetricInstrument::Gauge(g) => Some(g.clone()),
_ => None,
})
}
fn slot_histogram(
&self,
family: &str,
) -> Option<Arc<nmbrs_metrics::instruments::histogram::Histogram>> {
self.slots
.iter()
.find(|s| s.family == family)
.and_then(|s| match s.instrument.as_ref()? {
MetricInstrument::Histogram(h) => Some(h.clone()),
_ => None,
})
}
fn slot_counter(
&self,
family: &str,
) -> Option<Arc<nmbrs_metrics::instruments::counter::Counter>> {
self.slots
.iter()
.find(|s| s.family == family)
.and_then(|s| match s.instrument.as_ref()? {
MetricInstrument::Counter(c) => Some(c.clone()),
_ => None,
})
}
}
fn kernel_with_const_outputs(
consts: &[(&str, f64)],
) -> (
crate::scope_kernel::ScopeKernel,
crate::fixture::ScopeFixture,
) {
use polydat::compile::assembly::{PolydatAssembler, WireRef};
use polydat::library::fixed::ConstF64;
let mut asm = PolydatAssembler::new(vec!["cycle".into()]);
for (name, val) in consts {
let binding = crate::scope::synthesize_metric_binding_name(name);
asm.add_node(&binding, Box::new(ConstF64::new(*val)), vec![]);
asm.add_output(&binding, WireRef::node(&binding));
}
let kernel =
crate::scope_kernel::ScopeKernel::from(asm.compile().expect("test kernel asm.compile"));
let fx = crate::fixture::ScopeFixture::new(kernel.program().clone());
(kernel, fx)
}
fn typed_wrap_with_kernel(
inner: Arc<dyn OpDispenser>,
decls: &HashMap<String, MetricSpec>,
consts: &[(&str, f64)],
) -> Result<
(
Arc<MetricsDispenser>,
crate::fixture::ResolvedPulls,
crate::scope_kernel::ScopeKernel,
),
String,
> {
let (mut kernel, mut fx) = kernel_with_const_outputs(consts);
let mut comp = fresh_component();
if decls.is_empty() {
return Err("typed_wrap_with_kernel requires non-empty decls".into());
}
let mut entries: Vec<_> = decls.iter().collect();
entries.sort_by(|a, b| a.0.cmp(b.0));
let component_labels = comp.effective_labels().clone();
let mut slots = Vec::with_capacity(entries.len());
for (name, spec) in entries {
let family = spec.family.clone().unwrap_or_else(|| name.clone());
let format = match &spec.format {
Some(s) => Some(
nmbrs_workload::metric_format::parse_format_spec(s)
.map_err(|e| format!("metric '{name}' format: {e}"))?,
),
None => None,
};
let kind = spec.kind.unwrap_or_default();
let instr_labels = component_labels.with("family", family.clone());
let instrument = match kind {
MetricKind::Gauge => MetricInstrument::Gauge(Arc::new(
nmbrs_metrics::instruments::gauge::ValueGauge::new(instr_labels),
)),
MetricKind::Histogram => MetricInstrument::Histogram(Arc::new(
nmbrs_metrics::instruments::histogram::Histogram::new(instr_labels),
)),
MetricKind::Counter => MetricInstrument::Counter(Arc::new(
nmbrs_metrics::instruments::counter::Counter::new(instr_labels),
)),
};
comp.register_instrument_with_unit(
family.clone(),
spec.unit.clone(),
instrument.as_ref(),
)?;
let binding_name = crate::scope::synthesize_metric_binding_name(name);
let _ = fx.register_pull(&binding_name)?;
slots.push(MetricSlot {
family,
value_expr: spec.value.clone(),
binding_name,
format,
instrument: Some(instrument),
placement: None,
});
}
let typed = Arc::new(MetricsDispenser {
inner,
slots: Arc::new(slots),
});
let plan = fx.seal();
kernel.set_inputs(&[0]);
let pulls = plan.resolve_with(&mut kernel);
Ok((typed, pulls, kernel))
}
fn run_dispenser(
dispenser: Arc<dyn OpDispenser>,
pulls: &crate::fixture::ResolvedPulls,
kernel: &mut crate::scope_kernel::ScopeKernel,
) -> Result<OpResult, ExecutionError> {
let fields = crate::adapter::ResolvedFields::new(vec![], vec![]);
let cw = crate::wires::CycleWires::new(kernel);
let ctx = ExecCtx::with_wires(&fields, pulls, &cw);
let rt = tokio::runtime::Builder::new_current_thread()
.build()
.unwrap();
rt.block_on(dispenser.execute(0, &ctx))
}
#[test]
fn metrics_dispenser_gauge_records_f64() {
let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
let mut decl = HashMap::new();
decl.insert(
"my_factor".into(),
make_spec("my_factor", MetricKind::Gauge, None),
);
let (typed, pulls, mut kernel) =
typed_wrap_with_kernel(inner, &decl, &[("my_factor", 3.5)]).unwrap();
let gauge = typed.slot_gauge("my_factor").unwrap();
run_dispenser(typed.clone() as Arc<dyn OpDispenser>, &pulls, &mut kernel).unwrap();
assert!((gauge.get() - 3.5).abs() < 1e-9);
}
#[test]
fn metrics_dispenser_histogram_truncates_to_u64() {
let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
let mut decl = HashMap::new();
decl.insert(
"latency_ms".into(),
make_spec("latency_ms", MetricKind::Histogram, None),
);
let (typed, pulls, mut kernel) =
typed_wrap_with_kernel(inner, &decl, &[("latency_ms", 7.9)]).unwrap();
let hist = typed.slot_histogram("latency_ms").unwrap();
run_dispenser(typed.clone() as Arc<dyn OpDispenser>, &pulls, &mut kernel).unwrap();
let snap = hist.peek_snapshot();
assert_eq!(snap.max(), 7);
assert_eq!(snap.len(), 1);
}
#[test]
fn metrics_dispenser_counter_positive_inc_and_skip_non_positive() {
let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
let mut decl = HashMap::new();
decl.insert(
"ok_inc".into(),
make_spec("ok_inc", MetricKind::Counter, None),
);
decl.insert(
"skip_inc".into(),
make_spec("skip_inc", MetricKind::Counter, None),
);
let (typed, pulls, mut kernel) =
typed_wrap_with_kernel(inner, &decl, &[("ok_inc", 5.0), ("skip_inc", 0.0)]).unwrap();
let ok_counter = typed.slot_counter("ok_inc").unwrap();
let skip_counter = typed.slot_counter("skip_inc").unwrap();
run_dispenser(typed.clone() as Arc<dyn OpDispenser>, &pulls, &mut kernel).unwrap();
assert_eq!(ok_counter.get(), 5);
assert_eq!(skip_counter.get(), 0);
}
#[test]
fn metrics_dispenser_format_rounds_value() {
let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
let mut decl = HashMap::new();
decl.insert(
"ratio".into(),
make_spec("ratio", MetricKind::Gauge, Some("#.##")),
);
let (typed, pulls, mut kernel) =
typed_wrap_with_kernel(inner, &decl, &[("ratio", 1.234)]).unwrap();
let gauge = typed.slot_gauge("ratio").unwrap();
run_dispenser(typed.clone() as Arc<dyn OpDispenser>, &pulls, &mut kernel).unwrap();
assert!((gauge.get() - 1.23).abs() < 1e-9);
}
#[test]
fn metrics_dispenser_duplicate_family_errors() {
let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
let comp = fresh_component_arc();
comp.write()
.unwrap()
.register_instrument(
"recall_at_10",
nmbrs_metrics::component::InstrumentRef::Counter(Arc::new(
nmbrs_metrics::instruments::counter::Counter::new(
nmbrs_metrics::labels::Labels::of("name", "recall_at_10"),
),
)),
)
.unwrap();
let mut decl = HashMap::new();
decl.insert(
"recall_at_10".into(),
make_spec("recall_at_10", MetricKind::Gauge, None),
);
let (_kernel, mut fx) = kernel_with_const_outputs(&[("recall_at_10", 0.0)]);
let err = match wrap_on(inner, &decl, &comp, &mut fx) {
Ok(_) => panic!("expected duplicate-family error, got Ok"),
Err(e) => e,
};
assert!(
err.contains("duplicate family name"),
"unexpected error: {err}"
);
}
#[test]
fn metrics_dispenser_skipped_op_records_nothing() {
struct SkipInner;
impl OpDispenser for SkipInner {
fn execute<'a>(
&'a self,
_cycle: u64,
_ctx: &'a ExecCtx<'a>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
> {
Box::pin(async move { Ok(OpResult::skipped()) })
}
}
let mut decl = HashMap::new();
decl.insert("g".into(), make_spec("g", MetricKind::Gauge, None));
let (typed, pulls, mut kernel) =
typed_wrap_with_kernel(Arc::new(SkipInner), &decl, &[("g", 1.0)]).unwrap();
let gauge = typed.slot_gauge("g").unwrap();
let res =
run_dispenser(typed.clone() as Arc<dyn OpDispenser>, &pulls, &mut kernel).unwrap();
assert!(res.skipped);
assert_eq!(gauge.get(), 0.0);
}
#[test]
fn metrics_dispenser_accepts_arbitrary_polydat_expression() {
let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
let mut decl = HashMap::new();
decl.insert(
"computed".into(),
make_spec("factor * 2.0", MetricKind::Gauge, None),
);
let (mut kernel, mut fx) = kernel_with_const_outputs(&[("computed", 6.0)]);
let comp = fresh_component_arc();
let _ = wrap_on(inner, &decl, &comp, &mut fx)
.expect("arbitrary Polydat expression should wrap cleanly");
let plan = fx.seal();
kernel.set_inputs(&[0]);
let _pulls = plan.resolve_with(&mut kernel);
}
#[test]
fn metrics_dispenser_missing_wire_errors_at_init() {
let inner: Arc<dyn OpDispenser> = Arc::new(CapturesInner);
let mut decl = HashMap::new();
decl.insert(
"missing_metric".into(),
make_spec("absent_wire", MetricKind::Gauge, None),
);
let (_kernel, mut fx) = kernel_with_const_outputs(&[("present", 1.0)]);
let comp = fresh_component_arc();
let err = wrap_on(inner, &decl, &comp, &mut fx)
.err()
.expect("missing-wire metric should error at init");
assert!(err.contains("absent_wire"), "msg: {err}");
assert!(err.contains("Available"), "msg: {err}");
}
#[test]
fn value_to_f64_smoke() {
assert_eq!(value_to_f64(&polydat::ast::Value::U64(5)), Some(5.0));
assert_eq!(value_to_f64(&polydat::ast::Value::Bool(true)), Some(1.0));
assert_eq!(value_to_f64(&polydat::ast::Value::Str("x".into())), None);
}
}