1use 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
17pub const NAME: WrapperName = WrapperName::new("metrics");
19
20fn 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
38const 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 super::result::NAME,
57];
58
59inventory::submit! {
60 WrapperRegistration {
61 name: NAME,
62 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
75pub struct MetricsDispenser {
107 inner: Arc<dyn OpDispenser>,
108 slots: Arc<Vec<MetricSlot>>,
112}
113
114pub(crate) fn publish_gauges_lenient(slots: &[MetricSlot], wires: &dyn crate::wires::WireSource) {
124 for slot in slots {
125 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
154pub(crate) struct MetricSlot {
157 family: String,
160 value_expr: String,
165 binding_name: String,
170 format: Option<nmbrs_workload::metric_format::FormatSpec>,
173 instrument: Option<MetricInstrument>,
180 placement: Option<CellPlacement>,
182}
183
184pub(crate) struct CellPlacement {
190 dims: Vec<(String, String)>,
194 parent: Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
196 kind: nmbrs_workload::model::MetricKind,
197 unit: Option<String>,
198 instances: std::sync::Mutex<std::collections::HashMap<String, MetricInstrument>>,
202}
203
204#[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 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
235fn 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 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 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 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 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 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 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 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 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 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 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#[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 if result.skipped {
561 return Ok(result);
562 }
563 for slot in self.slots.iter() {
564 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 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 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 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 #[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 fn fresh_component_arc() -> Arc<std::sync::RwLock<nmbrs_metrics::component::Component>> {
734 Arc::new(std::sync::RwLock::new(fresh_component()))
735 }
736
737 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 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 #[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}