diff --git a/foundations-metrics-registry/src/encode_metric.rs b/foundations-metrics-registry/src/encode_metric.rs index bc59c54c..f2f423e1 100644 --- a/foundations-metrics-registry/src/encode_metric.rs +++ b/foundations-metrics-registry/src/encode_metric.rs @@ -6,5 +6,8 @@ use crate::proto::MetricFamily; /// metric or series that fails, so an empty `Vec` is a valid result. pub trait EncodeMetric: Send + Sync + 'static { /// Encodes this metric into zero or more [`MetricFamily`] messages. + /// + /// Every returned [`MetricFamily`] must set `name` to a complete, non-empty + /// producer-level name. fn encode(&self) -> Vec; } diff --git a/foundations-metrics-registry/src/metadata.rs b/foundations-metrics-registry/src/metadata.rs index 6ebb8e2e..6656e3c7 100644 --- a/foundations-metrics-registry/src/metadata.rs +++ b/foundations-metrics-registry/src/metadata.rs @@ -4,7 +4,7 @@ /// [`register`](crate::register). Build it from [`default`](Self::default) plus /// the setters, since downstream crates can't use a struct literal. #[non_exhaustive] -#[derive(Clone, Default)] +#[derive(Clone, Debug, Default)] pub struct RegistrationMetadata { /// Whether the metric is exported only when optional metrics are requested. pub optional: bool, diff --git a/foundations-metrics-registry/src/proto/mod.rs b/foundations-metrics-registry/src/proto/mod.rs index 50f85e21..beaa392c 100644 --- a/foundations-metrics-registry/src/proto/mod.rs +++ b/foundations-metrics-registry/src/proto/mod.rs @@ -10,7 +10,8 @@ mod model; pub use model::{ - Bucket, BucketSpan, Counter, Gauge, Histogram, LabelPair, Metric, MetricFamily, MetricType, + Bucket, BucketSpan, Counter, Exemplar, Gauge, Histogram, LabelPair, Metric, MetricFamily, + MetricType, Quantile, Summary, }; #[cfg(test)] diff --git a/foundations-metrics/Cargo.toml b/foundations-metrics/Cargo.toml index 48090317..a57ff8b9 100644 --- a/foundations-metrics/Cargo.toml +++ b/foundations-metrics/Cargo.toml @@ -15,6 +15,8 @@ prometheus-client = { version = "0.25.0", features = [ serde = { workspace = true, features = ["derive"] } ryu = "1.0.23" parking_lot = { workspace = true } +prost = { workspace = true } +prost-types = { workspace = true } [lints] workspace = true diff --git a/foundations-metrics/src/collect.rs b/foundations-metrics/src/collect.rs new file mode 100644 index 00000000..a460099c --- /dev/null +++ b/foundations-metrics/src/collect.rs @@ -0,0 +1,490 @@ +use foundations_metrics_registry::{iter, proto::LabelPair}; + +use crate::MetricFamily; +use crate::diagnostics::report_collect_error; +use crate::validation::{ + NAME_REQUIREMENT, ValidationContext, is_valid_name, sanitize_metric_family, +}; + +/// Options that control which registered metrics are collected and how the +/// service name is represented. +#[derive(Copy, Clone, Debug)] +pub struct CollectionOptions<'a> { + /// Whether metrics registered as optional are included. + pub include_optional: bool, + + /// Service name to add to collected metrics, if any. + pub service_name: Option<&'a str>, + + /// How `service_name` is represented in collected metrics. + pub service_name_format: ServiceNameFormat<'a>, +} + +/// How a service name is represented in collected metrics. +#[derive(Copy, Clone, Debug)] +pub enum ServiceNameFormat<'a> { + /// Prefix metric family names with the service name. + MetricPrefix, + + /// Add the service name to every metric row under the given label name. + LabelWithName(&'a str), +} + +/// Collects the currently registered metrics into the canonical protobuf model. +pub fn collect(options: CollectionOptions) -> Vec { + if options.service_name.is_some() + && let ServiceNameFormat::LabelWithName(label_name) = options.service_name_format + && !is_valid_name(label_name) + { + report_collect_error(format_args!( + "non-fatal error while collecting metrics: invalid configured service label name {label_name:?}; expected {NAME_REQUIREMENT}; skipped all metric families" + )); + return Vec::new(); + } + + let mut collected = Vec::new(); + + for registered in iter() { + let metadata = registered.metadata(); + + if metadata.optional && !options.include_optional { + continue; + } + + let mut families = registered.metric().encode(); + + if let Some(service_name) = options.service_name { + match options.service_name_format { + ServiceNameFormat::MetricPrefix if !metadata.unprefixed => { + for family in &mut families { + if let Some(name) = &mut family.name { + name.insert(0, '_'); + name.insert_str(0, service_name); + } + } + } + ServiceNameFormat::LabelWithName(label_name) => { + let service_label = LabelPair { + name: Some(label_name.to_owned()), + value: Some(service_name.to_owned()), + }; + + for family in &mut families { + let family_name = family.name.as_deref().unwrap_or_default(); + family.metric.retain_mut(|metric| { + let mut has_same_value = false; + for label in &metric.label { + if label.name.as_deref() != Some(label_name) { + continue; + } + + if label.value.as_deref() != Some(service_name) { + report_collect_error(format_args!( + "non-fatal error while collecting metrics: skipped row in metric family {family_name:?}; service label {label_name:?} already has a different value" + )); + return false; + } + has_same_value = true; + } + + if !has_same_value { + metric.label.insert(0, service_label.clone()); + } + true + }); + } + } + ServiceNameFormat::MetricPrefix => {} + } + } + + families.retain_mut(|family| sanitize_metric_family(family, ValidationContext::Collection)); + collected.extend(families); + } + + collected +} + +#[cfg(test)] +mod tests { + use foundations_metrics_registry::proto::{ + Bucket, Counter, Exemplar, Gauge, Histogram, LabelPair, Metric, MetricType, + }; + + use super::*; + use crate::{EncodeMetric, RegistrationMetadata, register}; + + struct TestMetric(&'static str); + + struct TestFamilyMetric(MetricFamily); + + impl EncodeMetric for TestMetric { + fn encode(&self) -> Vec { + vec![MetricFamily { + name: Some(self.0.to_owned()), + help: Some("Test metric.".to_owned()), + r#type: Some(MetricType::Gauge as i32), + metric: vec![Metric::default()], + unit: None, + }] + } + } + + impl EncodeMetric for TestFamilyMetric { + fn encode(&self) -> Vec { + vec![self.0.clone()] + } + } + + fn register_test_metric(name: &'static str, metadata: RegistrationMetadata) { + register( + Box::new(TestMetric(name)) as Box, + metadata, + ); + } + + fn register_test_family(family: MetricFamily) { + register( + Box::new(TestFamilyMetric(family)) as Box, + RegistrationMetadata::default(), + ); + } + + fn label(name: &str, value: &str) -> LabelPair { + LabelPair { + name: Some(name.to_owned()), + value: Some(value.to_owned()), + } + } + + #[test] + fn filters_optional_metrics_and_applies_service_prefix() { + register_test_metric("collect_required_metric", RegistrationMetadata::default()); + register_test_metric( + "collect_optional_metric", + RegistrationMetadata::default().optional(true), + ); + register_test_metric( + "collect_unprefixed_metric", + RegistrationMetadata::default().unprefixed(true), + ); + + let required = collect(CollectionOptions { + include_optional: false, + service_name: Some("test_service"), + service_name_format: ServiceNameFormat::MetricPrefix, + }); + let required_names: Vec<_> = required + .iter() + .filter_map(|family| family.name.as_deref()) + .collect(); + + assert!(required_names.contains(&"test_service_collect_required_metric")); + assert!(!required_names.contains(&"test_service_collect_optional_metric")); + assert!(required_names.contains(&"collect_unprefixed_metric")); + + let with_optional = collect(CollectionOptions { + include_optional: true, + service_name: Some("test_service"), + service_name_format: ServiceNameFormat::MetricPrefix, + }); + + assert!(with_optional.iter().any(|family| { + family.name.as_deref() == Some("test_service_collect_optional_metric") + })); + } + + #[test] + fn service_label_is_added_to_prefixed_and_unprefixed_metrics() { + register_test_metric("collect_label_metric", RegistrationMetadata::default()); + register_test_metric( + "collect_label_unprefixed_metric", + RegistrationMetadata::default().unprefixed(true), + ); + + let families = collect(CollectionOptions { + include_optional: false, + service_name: Some("test_service"), + service_name_format: ServiceNameFormat::LabelWithName("service"), + }); + + for name in ["collect_label_metric", "collect_label_unprefixed_metric"] { + let family = families + .iter() + .find(|family| family.name.as_deref() == Some(name)) + .expect("registered metric should be collected"); + let label = &family.metric[0].label[0]; + + assert_eq!(label.name.as_deref(), Some("service")); + assert_eq!(label.value.as_deref(), Some("test_service")); + } + } + + #[test] + fn keeps_nonstandard_service_prefixed_family_names() { + register_test_metric( + "collect_nonstandard_prefix_metric", + RegistrationMetadata::default(), + ); + + let families = collect(CollectionOptions { + include_optional: false, + service_name: Some("invalid-service"), + service_name_format: ServiceNameFormat::MetricPrefix, + }); + + assert!(families.iter().any(|family| { + family.name.as_deref() == Some("invalid-service_collect_nonstandard_prefix_metric") + })); + } + + #[test] + fn keeps_nonstandard_service_label_names() { + register_test_metric( + "collect_nonstandard_service_label_metric", + RegistrationMetadata::default(), + ); + + let families = collect(CollectionOptions { + include_optional: false, + service_name: Some("test_service"), + service_name_format: ServiceNameFormat::LabelWithName("service:name"), + }); + + let family = families + .iter() + .find(|family| { + family.name.as_deref() == Some("collect_nonstandard_service_label_metric") + }) + .expect("metric with a nonstandard service label name should remain"); + assert_eq!( + family.metric[0].label[0], + label("service:name", "test_service") + ); + } + + #[test] + fn service_label_insertion_is_idempotent_and_drops_different_values() { + let service_label_name = "collect_service_collision_label"; + register_test_family(MetricFamily { + name: Some("collect_service_collision_metric".to_owned()), + help: None, + r#type: Some(MetricType::Gauge as i32), + metric: vec![ + Metric { + label: vec![label("id", "same"), label(service_label_name, "wanted")], + gauge: Some(Gauge { value: Some(1.0) }), + ..Default::default() + }, + Metric { + label: vec![label("id", "different"), label(service_label_name, "other")], + gauge: Some(Gauge { value: Some(2.0) }), + ..Default::default() + }, + Metric { + label: vec![label("id", "absent")], + gauge: Some(Gauge { value: Some(3.0) }), + ..Default::default() + }, + ], + unit: None, + }); + + let families = collect(CollectionOptions { + include_optional: false, + service_name: Some("wanted"), + service_name_format: ServiceNameFormat::LabelWithName(service_label_name), + }); + let family = families + .iter() + .find(|family| family.name.as_deref() == Some("collect_service_collision_metric")) + .expect("test family should be collected"); + + assert_eq!(family.metric.len(), 2); + let same = family + .metric + .iter() + .find(|metric| metric.label[0].value.as_deref() == Some("same")) + .expect("same-value row should remain"); + assert_eq!( + same.label + .iter() + .filter(|label| label.name.as_deref() == Some(service_label_name)) + .count(), + 1 + ); + assert_eq!(same.label[0].name.as_deref(), Some("id")); + + let absent = family + .metric + .iter() + .find(|metric| { + metric + .label + .iter() + .any(|label| label.value.as_deref() == Some("absent")) + }) + .expect("row without a service label should remain"); + assert_eq!( + absent.label[0], + label(service_label_name, "wanted"), + "new service labels remain prepended" + ); + } + + #[test] + fn collection_keeps_nonstandard_names_and_skips_duplicate_and_reserved_labels() { + register_test_family(MetricFamily { + name: Some("collect_row_validation_gauge".to_owned()), + help: None, + r#type: Some(MetricType::Gauge as i32), + metric: vec![ + Metric { + label: vec![label("id", "valid")], + gauge: Some(Gauge { value: Some(1.0) }), + ..Default::default() + }, + Metric { + label: vec![label("bad\nname", "nonstandard")], + gauge: Some(Gauge { value: Some(2.0) }), + ..Default::default() + }, + Metric { + label: vec![label("dup", "a"), label("dup", "b")], + gauge: Some(Gauge { value: Some(3.0) }), + ..Default::default() + }, + ], + unit: None, + }); + register_test_family(MetricFamily { + name: Some("collect_row_validation_histogram".to_owned()), + help: None, + r#type: Some(MetricType::Histogram as i32), + metric: vec![ + Metric { + histogram: Some(Histogram::default()), + ..Default::default() + }, + Metric { + label: vec![label("le", "1")], + histogram: Some(Histogram::default()), + ..Default::default() + }, + ], + unit: None, + }); + let families = collect(CollectionOptions { + include_optional: false, + service_name: None, + service_name_format: ServiceNameFormat::MetricPrefix, + }); + + for (name, expected_rows) in [ + ("collect_row_validation_gauge", 2), + ("collect_row_validation_histogram", 1), + ] { + let family = families + .iter() + .find(|family| family.name.as_deref() == Some(name)) + .expect("valid family should remain"); + assert_eq!(family.metric.len(), expected_rows, "family {name}"); + } + } + + #[test] + fn collection_keeps_nonstandard_exemplar_names_and_drops_duplicates() { + register_test_family(MetricFamily { + name: Some("collect_counter_exemplar_validation".to_owned()), + help: None, + r#type: Some(MetricType::Counter as i32), + metric: vec![Metric { + counter: Some(Counter { + value: Some(1.0), + exemplar: Some(Exemplar { + label: vec![label("trace:id", "bad")], + value: Some(2.0), + timestamp: None, + }), + created_timestamp: None, + }), + ..Default::default() + }], + unit: None, + }); + register_test_family(MetricFamily { + name: Some("collect_histogram_exemplar_validation".to_owned()), + help: None, + r#type: Some(MetricType::Histogram as i32), + metric: vec![Metric { + histogram: Some(Histogram { + bucket: vec![Bucket { + exemplar: Some(Exemplar { + label: vec![label("dup", "a"), label("dup", "b")], + ..Default::default() + }), + ..Default::default() + }], + exemplars: vec![ + Exemplar { + label: vec![label("bad name", "bad")], + timestamp: Some(Default::default()), + ..Default::default() + }, + Exemplar { + label: vec![label("trace_id", "missing_timestamp")], + ..Default::default() + }, + Exemplar { + label: vec![label("trace_id", "good")], + timestamp: Some(Default::default()), + ..Default::default() + }, + ], + ..Default::default() + }), + ..Default::default() + }], + unit: None, + }); + + let families = collect(CollectionOptions { + include_optional: false, + service_name: None, + service_name_format: ServiceNameFormat::MetricPrefix, + }); + let counter = families + .iter() + .find(|family| family.name.as_deref() == Some("collect_counter_exemplar_validation")) + .expect("counter family should remain"); + assert_eq!( + counter.metric[0] + .counter + .as_ref() + .unwrap() + .exemplar + .as_ref() + .unwrap() + .label[0] + .name + .as_deref(), + Some("trace:id") + ); + + let histogram = families + .iter() + .find(|family| family.name.as_deref() == Some("collect_histogram_exemplar_validation")) + .expect("histogram family should remain"); + let histogram = histogram.metric[0].histogram.as_ref().unwrap(); + assert!(histogram.bucket[0].exemplar.is_none()); + assert_eq!(histogram.exemplars.len(), 2); + assert_eq!( + histogram + .exemplars + .iter() + .map(|exemplar| exemplar.label[0].name.as_deref().unwrap()) + .collect::>(), + ["bad name", "trace_id"] + ); + } +} diff --git a/foundations-metrics/src/encoding/mod.rs b/foundations-metrics/src/encoding/mod.rs new file mode 100644 index 00000000..59130a72 --- /dev/null +++ b/foundations-metrics/src/encoding/mod.rs @@ -0,0 +1,333 @@ +mod text; + +use prost::Message; + +use crate::MetricFamily; +use crate::validation::{ValidationContext, sanitized_metric_family}; + +pub use text::{OPENMETRICS_CONTENT_TYPE, encode_to_text}; + +/// Encodes metric families as length-delimited Prometheus protobuf messages. +pub fn encode_to_protobuf(families: &[MetricFamily]) -> Vec { + let mut output = Vec::new(); + for family in families { + if let Some(family) = sanitized_metric_family(family, ValidationContext::ProtobufEncoding) { + family + .encode_length_delimited(&mut output) + .expect("encoding a protobuf message to a Vec cannot fail"); + } + } + output +} + +#[cfg(test)] +mod tests { + use foundations_metrics_registry::proto::{ + Bucket, Counter, Exemplar, Gauge, Histogram, LabelPair, Metric, MetricType, Quantile, + Summary, + }; + + use super::*; + + fn label(name: &str, value: &str) -> LabelPair { + LabelPair { + name: Some(name.to_owned()), + value: Some(value.to_owned()), + } + } + + fn decode_families(mut bytes: &[u8]) -> Vec { + let mut families = Vec::new(); + while !bytes.is_empty() { + families.push( + MetricFamily::decode_length_delimited(&mut bytes) + .expect("encoded family should decode"), + ); + } + families + } + + #[test] + fn utf8_protobuf_output_is_unchanged() { + let families = [MetricFamily { + name: Some("valid counter λ".to_owned()), + help: Some("Valid counter.".to_owned()), + r#type: Some(MetricType::Counter as i32), + metric: vec![Metric { + label: vec![label("_label.name λ", "value")], + counter: Some(Counter { + value: Some(1.0), + exemplar: Some(Exemplar::default()), + created_timestamp: None, + }), + ..Default::default() + }], + unit: None, + }]; + let expected: Vec<_> = families + .iter() + .flat_map(Message::encode_length_delimited_to_vec) + .collect(); + + let encoded = encode_to_protobuf(&families); + assert_eq!(encoded, expected); + assert!( + decode_families(&encoded)[0].metric[0] + .counter + .as_ref() + .unwrap() + .exemplar + .is_some(), + "empty exemplars retain their existing protobuf behavior" + ); + } + + #[test] + fn protobuf_keeps_nonstandard_names_and_omits_empty_duplicate_and_reserved_names() { + let families = [ + MetricFamily { + name: Some(String::new()), + help: None, + r#type: Some(MetricType::Gauge as i32), + metric: vec![Metric { + gauge: Some(Gauge { value: Some(100.0) }), + ..Default::default() + }], + unit: None, + }, + MetricFamily { + name: Some("bad\nfamily".to_owned()), + help: None, + r#type: Some(MetricType::Gauge as i32), + metric: vec![Metric { + gauge: Some(Gauge { value: Some(99.0) }), + ..Default::default() + }], + unit: None, + }, + MetricFamily { + name: Some("protobuf_counter".to_owned()), + help: None, + r#type: Some(MetricType::Counter as i32), + metric: vec![ + Metric { + label: vec![label("id", "kept")], + counter: Some(Counter { + value: Some(1.0), + exemplar: Some(Exemplar { + label: vec![label("trace:id", "nonstandard")], + ..Default::default() + }), + created_timestamp: None, + }), + ..Default::default() + }, + Metric { + label: vec![label("bad name", "kept_nonstandard")], + counter: Some(Counter { + value: Some(2.0), + ..Default::default() + }), + ..Default::default() + }, + Metric { + label: vec![label("dup", "a"), label("dup", "b")], + counter: Some(Counter { + value: Some(3.0), + ..Default::default() + }), + ..Default::default() + }, + ], + unit: None, + }, + MetricFamily { + name: Some("protobuf_histogram".to_owned()), + help: None, + r#type: Some(MetricType::Histogram as i32), + metric: vec![ + Metric { + histogram: Some(Histogram { + bucket: vec![Bucket { + exemplar: Some(Exemplar { + label: vec![label("dup", "a"), label("dup", "b")], + ..Default::default() + }), + ..Default::default() + }], + exemplars: vec![ + Exemplar { + label: vec![label("bad#name", "nonstandard")], + timestamp: Some(Default::default()), + ..Default::default() + }, + Exemplar { + label: vec![label("trace_id", "missing_timestamp")], + ..Default::default() + }, + Exemplar { + label: vec![label("trace_id", "good")], + timestamp: Some(Default::default()), + ..Default::default() + }, + ], + ..Default::default() + }), + ..Default::default() + }, + Metric { + label: vec![label("le", "1")], + histogram: Some(Histogram::default()), + ..Default::default() + }, + ], + unit: None, + }, + MetricFamily { + name: Some("protobuf_sibling".to_owned()), + help: None, + r#type: Some(MetricType::Gauge as i32), + metric: vec![Metric { + gauge: Some(Gauge { value: Some(4.0) }), + ..Default::default() + }], + unit: None, + }, + ]; + + let decoded = decode_families(&encode_to_protobuf(&families)); + assert_eq!( + decoded + .iter() + .filter_map(|family| family.name.as_deref()) + .collect::>(), + [ + "bad\nfamily", + "protobuf_counter", + "protobuf_histogram", + "protobuf_sibling", + ] + ); + + assert_eq!(decoded[1].metric.len(), 2); + assert_eq!( + decoded[1].metric[0] + .counter + .as_ref() + .unwrap() + .exemplar + .as_ref() + .unwrap() + .label[0] + .name + .as_deref(), + Some("trace:id") + ); + assert_eq!( + decoded[1].metric[1].label[0].name.as_deref(), + Some("bad name") + ); + + assert_eq!(decoded[2].metric.len(), 1); + let histogram = decoded[2].metric[0].histogram.as_ref().unwrap(); + assert!(histogram.bucket[0].exemplar.is_none()); + assert_eq!(histogram.exemplars.len(), 2); + assert_eq!( + histogram + .exemplars + .iter() + .map(|exemplar| exemplar.label[0].name.as_deref().unwrap()) + .collect::>(), + ["bad#name", "trace_id"] + ); + + assert_eq!(decoded[3].metric.len(), 1); + } + + #[test] + fn preserves_summary_and_gauge_histogram_families() { + let families = [ + MetricFamily { + name: Some("empty_summary".to_owned()), + help: None, + r#type: Some(MetricType::Summary as i32), + metric: Vec::new(), + unit: None, + }, + MetricFamily { + name: Some("request_size".to_owned()), + help: Some("Request size.".to_owned()), + r#type: Some(MetricType::Summary as i32), + metric: vec![Metric { + summary: Some(Summary { + sample_count: Some(2), + sample_sum: Some(6.0), + quantile: vec![Quantile { + quantile: Some(0.5), + value: Some(3.0), + }], + created_timestamp: None, + }), + ..Default::default() + }], + unit: None, + }, + MetricFamily { + name: Some("empty_gauge_histogram".to_owned()), + help: None, + r#type: Some(MetricType::GaugeHistogram as i32), + metric: Vec::new(), + unit: None, + }, + MetricFamily { + name: Some("queue_item_age".to_owned()), + help: Some("Current age distribution of queued items.".to_owned()), + r#type: Some(MetricType::GaugeHistogram as i32), + metric: vec![Metric { + histogram: Some(Histogram { + sample_count: Some(3), + sample_sum: Some(8.0), + bucket: vec![Bucket { + cumulative_count: Some(1), + upper_bound: Some(1.0), + ..Default::default() + }], + ..Default::default() + }), + ..Default::default() + }], + unit: None, + }, + ]; + + let expected: Vec<_> = families + .iter() + .flat_map(Message::encode_length_delimited_to_vec) + .collect(); + + assert_eq!(encode_to_protobuf(&families), expected); + } + + #[test] + fn preserves_legacy_info_gauge_representation() { + let families = [MetricFamily { + name: Some("build_info".to_owned()), + help: Some("Build information.".to_owned()), + r#type: Some(MetricType::Gauge as i32), + metric: vec![Metric { + label: vec![LabelPair { + name: Some("version".to_owned()), + value: Some("1.2.3".to_owned()), + }], + gauge: Some(Gauge { value: Some(1.0) }), + ..Default::default() + }], + unit: None, + }]; + + let encoded = encode_to_protobuf(&families); + let decoded = MetricFamily::decode_length_delimited(encoded.as_slice()).unwrap(); + + assert_eq!(decoded, families[0]); + } +} diff --git a/foundations-metrics/src/encoding/text.rs b/foundations-metrics/src/encoding/text.rs new file mode 100644 index 00000000..c629794e --- /dev/null +++ b/foundations-metrics/src/encoding/text.rs @@ -0,0 +1,1015 @@ +use std::fmt::Write as _; + +use foundations_metrics_registry::proto::{ + Exemplar, Histogram, LabelPair, Metric, MetricFamily, MetricType, +}; + +use crate::diagnostics::report_collect_error; +use crate::validation::{ValidationContext, sanitized_metric_family}; + +/// Content type for the UTF-8 OpenMetrics text emitted by [`encode_to_text`]. +pub const OPENMETRICS_CONTENT_TYPE: &str = + "application/openmetrics-text; version=1.0.0; charset=utf-8; escaping=allow-utf-8"; + +/// Encodes metric families as UTF-8 OpenMetrics text. +/// +/// Label names are always quoted. Metric names outside the legacy Prometheus +/// grammar use the quoted metric-name form. Serve the output with +/// [`OPENMETRICS_CONTENT_TYPE`] so scrapers retain UTF-8 names. +pub fn encode_to_text(families: &[MetricFamily]) -> String { + let mut output = String::new(); + + for family in families { + if let Some(family) = sanitized_metric_family(family, ValidationContext::TextEncoding) { + encode_family(&mut output, &family); + } + } + + output.push_str("# EOF\n"); + output +} + +fn encode_family(output: &mut String, family: &MetricFamily) { + let name = family + .name + .as_deref() + .expect("metric family names are validated before text encoding"); + let Some(metric_type) = family + .r#type + .and_then(|value| MetricType::try_from(value).ok()) + else { + report_collect_error(format_args!( + "non-fatal error while encoding OpenMetrics text: skipped metric family {name:?} with an unknown type" + )); + return; + }; + let metric_type_name = match metric_type { + MetricType::Counter => "counter", + MetricType::Gauge => "gauge", + MetricType::Summary => "summary", + MetricType::Untyped => "unknown", + MetricType::Histogram => "histogram", + MetricType::GaugeHistogram => "gaugehistogram", + }; + + // An empty help string carries no information, so the line is omitted rather + // than written with a blank value. + if let Some(help) = family.help.as_deref().filter(|help| !help.is_empty()) { + output.push_str("# HELP "); + write_metadata_name(output, name); + output.push(' '); + write_escaped(output, help); + output.push('\n'); + } + + output.push_str("# TYPE "); + write_metadata_name(output, name); + output.push(' '); + output.push_str(metric_type_name); + output.push('\n'); + + if let Some(unit) = &family.unit { + output.push_str("# UNIT "); + write_metadata_name(output, name); + output.push(' '); + write_escaped(output, unit); + output.push('\n'); + } + + for metric in &family.metric { + encode_metric(output, name, metric_type, metric); + } +} + +fn encode_metric(output: &mut String, name: &str, metric_type: MetricType, metric: &Metric) { + match metric_type { + MetricType::Counter => { + let Some(counter) = &metric.counter else { + report_missing_value(name, "counter"); + return; + }; + write_sample( + output, + name, + "", + metric, + None, + SampleValue::Float(counter.value.unwrap_or_default()), + counter.exemplar.as_ref(), + ); + } + MetricType::Gauge => { + let Some(gauge) = &metric.gauge else { + report_missing_value(name, "gauge"); + return; + }; + write_plain_sample( + output, + name, + "", + metric, + SampleValue::Float(gauge.value.unwrap_or_default()), + ); + } + MetricType::Summary => { + let Some(summary) = &metric.summary else { + report_missing_value(name, "summary"); + return; + }; + for quantile in &summary.quantile { + write_sample( + output, + name, + "", + metric, + Some(("quantile", quantile.quantile.unwrap_or_default())), + SampleValue::Float(quantile.value.unwrap_or_default()), + None, + ); + } + write_plain_sample( + output, + name, + "_sum", + metric, + SampleValue::Float(summary.sample_sum.unwrap_or_default()), + ); + write_plain_sample( + output, + name, + "_count", + metric, + SampleValue::Unsigned(summary.sample_count.unwrap_or_default()), + ); + } + MetricType::Untyped => { + let Some(untyped) = &metric.untyped else { + report_missing_value(name, "untyped value"); + return; + }; + write_plain_sample( + output, + name, + "", + metric, + SampleValue::Float(untyped.value.unwrap_or_default()), + ); + } + MetricType::Histogram | MetricType::GaugeHistogram => { + let Some(histogram) = &metric.histogram else { + report_missing_value(name, "histogram"); + return; + }; + encode_histogram(output, name, metric_type, metric, histogram); + } + } +} + +fn encode_histogram( + output: &mut String, + name: &str, + metric_type: MetricType, + metric: &Metric, + histogram: &Histogram, +) { + if histogram.bucket.is_empty() && has_native_buckets(histogram) { + report_collect_error(format_args!( + "non-fatal error while encoding OpenMetrics text: skipped native histogram row for {name:?}; native histograms require protobuf output" + )); + return; + } + + if histogram + .sample_count_float + .is_some_and(|count| count > 0.0) + || histogram.bucket.iter().any(|bucket| { + bucket + .cumulative_count_float + .is_some_and(|count| count > 0.0) + }) + { + report_collect_error(format_args!( + "non-fatal error while encoding OpenMetrics text: skipped histogram row for {name:?} with floating-point bucket counts" + )); + return; + } + + let (sum_suffix, count_suffix) = if metric_type == MetricType::GaugeHistogram { + ("_gsum", "_gcount") + } else { + ("_sum", "_count") + }; + let sample_count = histogram.sample_count.unwrap_or_default(); + + write_plain_sample( + output, + name, + sum_suffix, + metric, + SampleValue::Float(histogram.sample_sum.unwrap_or_default()), + ); + write_plain_sample( + output, + name, + count_suffix, + metric, + SampleValue::Unsigned(sample_count), + ); + + let mut infinity_bucket_seen = false; + for bucket in &histogram.bucket { + let upper_bound = bucket.upper_bound.unwrap_or_default(); + let upper_bound = if upper_bound == f64::MAX { + f64::INFINITY + } else { + upper_bound + }; + infinity_bucket_seen |= upper_bound == f64::INFINITY; + + write_sample( + output, + name, + "_bucket", + metric, + Some(("le", upper_bound)), + SampleValue::Unsigned(bucket.cumulative_count.unwrap_or_default()), + bucket.exemplar.as_ref(), + ); + } + + if !infinity_bucket_seen { + write_sample( + output, + name, + "_bucket", + metric, + Some(("le", f64::INFINITY)), + SampleValue::Unsigned(sample_count), + None, + ); + } +} + +fn has_native_buckets(histogram: &Histogram) -> bool { + histogram.schema.is_some() + || histogram.zero_threshold.is_some() + || histogram.zero_count.is_some() + || histogram.zero_count_float.is_some() + || !histogram.negative_span.is_empty() + || !histogram.negative_delta.is_empty() + || !histogram.negative_count.is_empty() + || !histogram.positive_span.is_empty() + || !histogram.positive_delta.is_empty() + || !histogram.positive_count.is_empty() +} + +#[derive(Clone, Copy)] +enum SampleValue { + Float(f64), + Unsigned(u64), +} + +fn write_plain_sample( + output: &mut String, + name: &str, + suffix: &str, + metric: &Metric, + value: SampleValue, +) { + write_sample(output, name, suffix, metric, None, value, None); +} + +fn write_sample( + output: &mut String, + name: &str, + suffix: &str, + metric: &Metric, + additional_label: Option<(&str, f64)>, + value: SampleValue, + exemplar: Option<&Exemplar>, +) { + write_sample_name_and_labels(output, name, suffix, &metric.label, additional_label); + output.push(' '); + match value { + SampleValue::Float(value) => write_float(output, value), + SampleValue::Unsigned(value) => { + write!(output, "{value}").expect("writing to a String cannot fail"); + } + } + + if let Some(timestamp_ms) = metric.timestamp_ms { + output.push(' '); + write_float(output, timestamp_ms as f64 / 1_000.0); + } + + if let Some(exemplar) = exemplar { + output.push_str(" # "); + if exemplar.label.is_empty() { + output.push_str("{}"); + } else { + write_labels(output, &exemplar.label, None); + } + output.push(' '); + write_float(output, exemplar.value.unwrap_or_default()); + + if let Some(timestamp) = &exemplar.timestamp { + output.push(' '); + write_float( + output, + timestamp.seconds as f64 + f64::from(timestamp.nanos) * 1e-9, + ); + } + } + + output.push('\n'); +} + +fn write_labels(output: &mut String, labels: &[LabelPair], additional_label: Option<(&str, f64)>) { + if labels.is_empty() && additional_label.is_none() { + return; + } + + output.push('{'); + let mut separator = ""; + for label in labels { + output.push_str(separator); + write_label_name(output, label.name.as_deref().unwrap_or_default()); + output.push_str("=\""); + write_escaped(output, label.value.as_deref().unwrap_or_default()); + output.push('"'); + separator = ","; + } + + if let Some((name, value)) = additional_label { + output.push_str(separator); + write_label_name(output, name); + output.push_str("=\""); + write_float(output, value); + output.push('"'); + } + + output.push('}'); +} + +fn write_sample_name_and_labels( + output: &mut String, + name: &str, + suffix: &str, + labels: &[LabelPair], + additional_label: Option<(&str, f64)>, +) { + if is_legacy_metric_name(name) { + output.push_str(name); + output.push_str(suffix); + write_labels(output, labels, additional_label); + return; + } + + output.push('{'); + output.push('"'); + write_escaped(output, name); + output.push_str(suffix); + output.push('"'); + + for label in labels { + output.push(','); + write_label_name(output, label.name.as_deref().unwrap_or_default()); + output.push_str("=\""); + write_escaped(output, label.value.as_deref().unwrap_or_default()); + output.push('"'); + } + + if let Some((name, value)) = additional_label { + output.push(','); + write_label_name(output, name); + output.push_str("=\""); + write_float(output, value); + output.push('"'); + } + + output.push('}'); +} + +fn write_metadata_name(output: &mut String, name: &str) { + if is_legacy_metric_name(name) { + output.push_str(name); + } else { + write_quoted(output, name); + } +} + +fn is_legacy_metric_name(name: &str) -> bool { + let mut bytes = name.bytes(); + bytes + .next() + .is_some_and(|byte| byte.is_ascii_alphabetic() || matches!(byte, b'_' | b':')) + && bytes.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b':')) +} + +/// Whether `name` matches the legacy label name grammar, which unlike metric +/// names does not permit colons. +fn is_legacy_label_name(name: &str) -> bool { + let mut bytes = name.bytes(); + bytes + .next() + .is_some_and(|byte| byte.is_ascii_alphabetic() || byte == b'_') + && bytes.all(|byte| byte.is_ascii_alphanumeric() || byte == b'_') +} + +/// Writes a label name, quoting it only when it needs UTF-8 syntax. +/// +/// Legacy-compatible names are emitted unquoted so that output stays byte +/// identical to the classic Prometheus text format, which collectors that +/// predate UTF-8 name support require. +fn write_label_name(output: &mut String, name: &str) { + if is_legacy_label_name(name) { + output.push_str(name); + } else { + write_quoted(output, name); + } +} + +fn write_quoted(output: &mut String, value: &str) { + output.push('"'); + write_escaped(output, value); + output.push('"'); +} + +fn write_float(output: &mut String, value: f64) { + if value.is_nan() { + output.push_str("NaN"); + } else if value == f64::INFINITY { + output.push_str("+Inf"); + } else if value == f64::NEG_INFINITY { + output.push_str("-Inf"); + } else { + let mut buffer = ryu::Buffer::new(); + let formatted = buffer.format(value); + output.push_str(formatted); + if !formatted.contains(['.', 'e', 'E']) { + output.push_str(".0"); + } + } +} + +fn write_escaped(output: &mut String, value: &str) { + for character in value.chars() { + match character { + '\\' => output.push_str("\\\\"), + '\n' => output.push_str("\\n"), + '"' => output.push_str("\\\""), + _ => output.push(character), + } + } +} + +fn report_missing_value(name: &str, expected: &str) { + report_collect_error(format_args!( + "non-fatal error while encoding OpenMetrics text: skipped row in metric family {name:?}; expected {expected} data" + )); +} + +#[cfg(test)] +mod tests { + use foundations_metrics_registry::proto::{ + Bucket, Counter, Exemplar, Gauge, Histogram, LabelPair, Metric, MetricFamily, MetricType, + Quantile, Summary, + }; + + use super::*; + + fn label(name: &str, value: &str) -> LabelPair { + LabelPair { + name: Some(name.to_owned()), + value: Some(value.to_owned()), + } + } + + #[test] + fn quotes_only_label_names_that_need_utf8_syntax() { + let families = [MetricFamily { + name: Some("requests".to_owned()), + help: None, + r#type: Some(MetricType::Counter as i32), + metric: vec![Metric { + label: vec![ + // Legacy-compatible names stay unquoted, so output remains + // byte identical to the classic Prometheus text format. + label("route", "/test"), + label("_internal9", "yes"), + // Colons are valid in metric names but not in label names. + label("trace:id", "abc"), + label("label.name", "dotted"), + label("indicateur_\u{8017}\u{65f6}", "utf8"), + ], + counter: Some(Counter { + value: Some(1.0), + ..Default::default() + }), + ..Default::default() + }], + unit: None, + }]; + + assert_eq!( + encode_to_text(&families), + "# TYPE requests counter\n\ +requests{route=\"/test\",_internal9=\"yes\",\"trace:id\"=\"abc\",\"label.name\"=\"dotted\",\"indicateur_\u{8017}\u{65f6}\"=\"utf8\"} 1.0\n\ +# EOF\n" + ); + } + + #[test] + fn omits_the_help_line_when_there_is_no_help_text() { + let families = [MetricFamily { + name: Some("requests".to_owned()), + help: Some(String::new()), + r#type: Some(MetricType::Counter as i32), + metric: vec![Metric { + counter: Some(Counter { + value: Some(1.0), + ..Default::default() + }), + ..Default::default() + }], + unit: None, + }]; + + assert_eq!( + encode_to_text(&families), + "# TYPE requests counter\n\ +requests 1.0\n\ +# EOF\n" + ); + } + + #[test] + fn encodes_counter_metadata_labels_timestamps_and_exemplars() { + let families = [MetricFamily { + name: Some("requests".to_owned()), + help: Some("A \"quoted\" help\\line\nnext".to_owned()), + r#type: Some(MetricType::Counter as i32), + metric: vec![Metric { + label: vec![LabelPair { + name: Some("kind".to_owned()), + value: Some("a\"b\\c\nd".to_owned()), + }], + counter: Some(Counter { + value: Some(1.0), + exemplar: Some(Exemplar { + label: vec![LabelPair { + name: Some("trace_id".to_owned()), + value: Some("abc".to_owned()), + }], + value: Some(2.0), + timestamp: None, + }), + created_timestamp: None, + }), + timestamp_ms: Some(1_500), + ..Default::default() + }], + unit: None, + }]; + + assert_eq!( + encode_to_text(&families), + "# HELP requests A \\\"quoted\\\" help\\\\line\\nnext\n\ +# TYPE requests counter\n\ +requests{kind=\"a\\\"b\\\\c\\nd\"} 1.0 1.5 # {trace_id=\"abc\"} 2.0\n\ +# EOF\n" + ); + } + + #[test] + fn encodes_classic_histogram_and_maps_terminal_bucket_to_infinity() { + let families = [MetricFamily { + name: Some("request_duration_seconds".to_owned()), + help: Some("Request duration.".to_owned()), + r#type: Some(MetricType::Histogram as i32), + metric: vec![Metric { + label: vec![LabelPair { + name: Some("route".to_owned()), + value: Some("/test".to_owned()), + }], + histogram: Some(Histogram { + sample_count: Some(3), + sample_sum: Some(4.5), + bucket: vec![ + Bucket { + cumulative_count: Some(1), + upper_bound: Some(1.0), + ..Default::default() + }, + Bucket { + cumulative_count: Some(3), + upper_bound: Some(f64::MAX), + ..Default::default() + }, + ], + ..Default::default() + }), + ..Default::default() + }], + unit: Some("seconds".to_owned()), + }]; + + assert_eq!( + encode_to_text(&families), + "# HELP request_duration_seconds Request duration.\n\ +# TYPE request_duration_seconds histogram\n\ +# UNIT request_duration_seconds seconds\n\ +request_duration_seconds_sum{route=\"/test\"} 4.5\n\ +request_duration_seconds_count{route=\"/test\"} 3\n\ +request_duration_seconds_bucket{route=\"/test\",le=\"1.0\"} 1\n\ +request_duration_seconds_bucket{route=\"/test\",le=\"+Inf\"} 3\n\ +# EOF\n" + ); + } + + #[test] + fn appends_an_infinite_histogram_bucket_when_missing() { + let families = [MetricFamily { + name: Some("values".to_owned()), + help: None, + r#type: Some(MetricType::Histogram as i32), + metric: vec![Metric { + histogram: Some(Histogram { + sample_count: Some(2), + sample_sum: Some(3.0), + bucket: vec![Bucket { + cumulative_count: Some(1), + upper_bound: Some(1.0), + ..Default::default() + }], + ..Default::default() + }), + ..Default::default() + }], + unit: None, + }]; + + let output = encode_to_text(&families); + assert!(output.contains("values_bucket{le=\"+Inf\"} 2\n")); + } + + #[test] + fn encodes_summaries_and_gauge_histograms() { + let families = [ + MetricFamily { + name: Some("request_size".to_owned()), + help: None, + r#type: Some(MetricType::Summary as i32), + metric: vec![Metric { + summary: Some(Summary { + sample_count: Some(2), + sample_sum: Some(6.0), + quantile: vec![Quantile { + quantile: Some(0.5), + value: Some(3.0), + }], + created_timestamp: None, + }), + ..Default::default() + }], + unit: None, + }, + MetricFamily { + name: Some("queue_depth".to_owned()), + help: None, + r#type: Some(MetricType::GaugeHistogram as i32), + metric: vec![Metric { + histogram: Some(Histogram { + sample_count: Some(3), + sample_sum: Some(8.0), + bucket: vec![Bucket { + cumulative_count: Some(1), + upper_bound: Some(1.0), + ..Default::default() + }], + ..Default::default() + }), + ..Default::default() + }], + unit: None, + }, + ]; + + assert_eq!( + encode_to_text(&families), + "# TYPE request_size summary\n\ +request_size{quantile=\"0.5\"} 3.0\n\ +request_size_sum 6.0\n\ +request_size_count 2\n\ +# TYPE queue_depth gaugehistogram\n\ +queue_depth_gsum 8.0\n\ +queue_depth_gcount 3\n\ +queue_depth_bucket{le=\"1.0\"} 1\n\ +queue_depth_bucket{le=\"+Inf\"} 3\n\ +# EOF\n" + ); + } + + #[test] + fn encodes_gauges_and_special_float_values() { + let families = [MetricFamily { + name: Some("temperature".to_owned()), + help: None, + r#type: Some(MetricType::Gauge as i32), + metric: [f64::NAN, f64::INFINITY, f64::NEG_INFINITY] + .into_iter() + .map(|value| Metric { + gauge: Some(Gauge { value: Some(value) }), + ..Default::default() + }) + .collect(), + unit: None, + }]; + + assert_eq!( + encode_to_text(&families), + "# TYPE temperature gauge\n\ +temperature NaN\n\ +temperature +Inf\n\ +temperature -Inf\n\ +# EOF\n" + ); + } + + #[test] + fn encodes_legacy_info_metric_as_gauge() { + let families = [MetricFamily { + name: Some("build_info".to_owned()), + help: Some("Build information.".to_owned()), + r#type: Some(MetricType::Gauge as i32), + metric: vec![Metric { + label: vec![LabelPair { + name: Some("version".to_owned()), + value: Some("1.2.3".to_owned()), + }], + gauge: Some(Gauge { value: Some(1.0) }), + ..Default::default() + }], + unit: None, + }]; + + assert_eq!( + encode_to_text(&families), + "# HELP build_info Build information.\n\ +# TYPE build_info gauge\n\ +build_info{version=\"1.2.3\"} 1.0\n\ +# EOF\n" + ); + } + + #[test] + fn appends_histogram_suffixes_inside_quoted_metric_names() { + let families = [MetricFamily { + name: Some("request.耗时".to_owned()), + help: None, + r#type: Some(MetricType::Histogram as i32), + metric: vec![Metric { + histogram: Some(Histogram { + sample_count: Some(1), + sample_sum: Some(0.5), + ..Default::default() + }), + ..Default::default() + }], + unit: None, + }]; + + assert_eq!( + encode_to_text(&families), + "# TYPE \"request.耗时\" histogram\n\ +{\"request.耗时_sum\"} 0.5\n\ +{\"request.耗时_count\"} 1\n\ +{\"request.耗时_bucket\",le=\"+Inf\"} 1\n\ +# EOF\n" + ); + } + + #[test] + fn skips_empty_metric_names() { + let families = [MetricFamily { + name: Some(String::new()), + help: None, + r#type: Some(MetricType::Gauge as i32), + metric: vec![Metric { + gauge: Some(Gauge { value: Some(1.0) }), + ..Default::default() + }], + unit: None, + }]; + + assert_eq!(encode_to_text(&families), "# EOF\n"); + } + + #[test] + fn quotes_and_escapes_utf8_metric_and_label_names() { + let families = [ + MetricFamily { + name: Some("bad\n# HELP injected metadata".to_owned()), + help: Some("Escaped help.".to_owned()), + r#type: Some(MetricType::Gauge as i32), + metric: vec![Metric { + label: vec![label("路由.name\n", "值")], + gauge: Some(Gauge { value: Some(99.0) }), + ..Default::default() + }], + unit: None, + }, + MetricFamily { + name: Some("valid:metric".to_owned()), + help: None, + r#type: Some(MetricType::Gauge as i32), + metric: vec![Metric { + gauge: Some(Gauge { value: Some(1.0) }), + ..Default::default() + }], + unit: None, + }, + ]; + + assert_eq!( + encode_to_text(&families), + "# HELP \"bad\\n# HELP injected metadata\" Escaped help.\n\ +# TYPE \"bad\\n# HELP injected metadata\" gauge\n\ +{\"bad\\n# HELP injected metadata\",\"路由.name\\n\"=\"值\"} 99.0\n\ +# TYPE valid:metric gauge\n\ +valid:metric 1.0\n\ +# EOF\n" + ); + } + + #[test] + fn nonstandard_labels_are_kept_while_duplicate_and_reserved_labels_are_dropped() { + let families = [ + MetricFamily { + name: Some("row_gauge".to_owned()), + help: None, + r#type: Some(MetricType::Gauge as i32), + metric: vec![ + Metric { + label: vec![label("id", "valid")], + gauge: Some(Gauge { value: Some(1.0) }), + ..Default::default() + }, + Metric { + label: vec![label("bad name", "nonstandard")], + gauge: Some(Gauge { value: Some(99.0) }), + ..Default::default() + }, + Metric { + label: vec![label("dup", "a"), label("dup", "b")], + gauge: Some(Gauge { value: Some(98.0) }), + ..Default::default() + }, + ], + unit: None, + }, + MetricFamily { + name: Some("row_histogram".to_owned()), + help: None, + r#type: Some(MetricType::Histogram as i32), + metric: vec![ + Metric { + histogram: Some(Histogram { + sample_count: Some(1), + sample_sum: Some(2.0), + ..Default::default() + }), + ..Default::default() + }, + Metric { + label: vec![label("le", "1")], + histogram: Some(Histogram { + sample_count: Some(99), + sample_sum: Some(99.0), + ..Default::default() + }), + ..Default::default() + }, + ], + unit: None, + }, + MetricFamily { + name: Some("row_summary".to_owned()), + help: None, + r#type: Some(MetricType::Summary as i32), + metric: vec![ + Metric { + summary: Some(Default::default()), + ..Default::default() + }, + Metric { + label: vec![label("quantile", "0.5")], + summary: Some(Default::default()), + ..Default::default() + }, + ], + unit: None, + }, + MetricFamily { + name: Some("row_gauge_histogram".to_owned()), + help: None, + r#type: Some(MetricType::GaugeHistogram as i32), + metric: vec![ + Metric { + histogram: Some(Histogram::default()), + ..Default::default() + }, + Metric { + label: vec![label("le", "1")], + histogram: Some(Histogram::default()), + ..Default::default() + }, + ], + unit: None, + }, + ]; + + let output = encode_to_text(&families); + assert!(output.contains("row_gauge{id=\"valid\"} 1.0\n")); + assert!(output.contains("row_gauge{\"bad name\"=\"nonstandard\"} 99.0\n")); + assert!(!output.contains("98.0")); + assert_eq!(output.matches("row_histogram_sum").count(), 1); + assert_eq!(output.matches("row_summary_sum").count(), 1); + assert_eq!(output.matches("row_gauge_histogram_gsum").count(), 1); + assert!(output.ends_with("# EOF\n")); + } + + #[test] + fn nonstandard_exemplar_labels_are_kept_while_duplicates_are_dropped() { + let families = [ + MetricFamily { + name: Some("exemplar_counter".to_owned()), + help: None, + r#type: Some(MetricType::Counter as i32), + metric: vec![ + Metric { + label: vec![label("id", "nonstandard")], + counter: Some(Counter { + value: Some(1.0), + exemplar: Some(Exemplar { + label: vec![label("trace:id", "bad")], + value: Some(2.0), + timestamp: None, + }), + created_timestamp: None, + }), + ..Default::default() + }, + Metric { + label: vec![label("id", "valid")], + counter: Some(Counter { + value: Some(3.0), + exemplar: Some(Exemplar { + label: vec![label("trace_id", "good")], + value: Some(4.0), + timestamp: None, + }), + created_timestamp: None, + }), + ..Default::default() + }, + ], + unit: None, + }, + MetricFamily { + name: Some("exemplar_histogram".to_owned()), + help: None, + r#type: Some(MetricType::Histogram as i32), + metric: vec![Metric { + histogram: Some(Histogram { + sample_count: Some(1), + sample_sum: Some(1.0), + bucket: vec![Bucket { + cumulative_count: Some(1), + upper_bound: Some(1.0), + exemplar: Some(Exemplar { + label: vec![label("dup", "a"), label("dup", "b")], + value: Some(5.0), + timestamp: None, + }), + ..Default::default() + }], + ..Default::default() + }), + ..Default::default() + }], + unit: None, + }, + ]; + + let output = encode_to_text(&families); + assert!( + output.contains( + "exemplar_counter{id=\"nonstandard\"} 1.0 # {\"trace:id\"=\"bad\"} 2.0\n" + ) + ); + assert!(output.contains("exemplar_counter{id=\"valid\"} 3.0 # {trace_id=\"good\"} 4.0\n")); + assert!(output.contains("exemplar_histogram_bucket{le=\"1.0\"} 1\n")); + assert!(!output.contains("\"dup\"=")); + } +} diff --git a/foundations-metrics/src/labels/serializer.rs b/foundations-metrics/src/labels/serializer.rs index cf3954e6..4c70fc62 100644 --- a/foundations-metrics/src/labels/serializer.rs +++ b/foundations-metrics/src/labels/serializer.rs @@ -2,9 +2,10 @@ use std::fmt; use foundations_metrics_registry::proto::LabelPair; use serde::Serialize; -use serde::ser::{Impossible, SerializeStruct, Serializer}; +use serde::ser::{Impossible, SerializeSeq, SerializeStruct, SerializeTuple, Serializer}; use super::LabelError; +use crate::validation::{NAME_REQUIREMENT, is_valid_name}; // Adapted from prometools' `serde::top::TopSerializer` // (https://github.com/nox/prometools, licensed MIT OR Apache-2.0). @@ -13,8 +14,8 @@ pub(super) struct LabelSetSerializer; impl Serializer for LabelSetSerializer { type Ok = Vec; type Error = LabelError; - type SerializeSeq = Impossible; - type SerializeTuple = Impossible; + type SerializeSeq = LabelSequenceSerializer; + type SerializeTuple = LabelTupleSerializer; type SerializeTupleStruct = Impossible; type SerializeTupleVariant = Impossible; type SerializeMap = Impossible; @@ -147,12 +148,21 @@ impl Serializer for LabelSetSerializer { Err(invalid_label_set()) } - fn serialize_seq(self, _len: Option) -> Result { - Err(invalid_label_set()) + fn serialize_seq(self, len: Option) -> Result { + Ok(LabelSequenceSerializer { + labels: Vec::with_capacity(len.unwrap_or_default()), + }) } - fn serialize_tuple(self, _len: usize) -> Result { - Err(invalid_label_set()) + fn serialize_tuple(self, len: usize) -> Result { + if len != 2 { + return Err(LabelError::new("label pairs must contain a name and value")); + } + + Ok(LabelTupleSerializer { + name: None, + value: None, + }) } fn serialize_tuple_struct( @@ -192,6 +202,77 @@ impl Serializer for LabelSetSerializer { } } +pub(super) struct LabelSequenceSerializer { + labels: Vec, +} + +impl SerializeSeq for LabelSequenceSerializer { + type Ok = Vec; + type Error = LabelError; + + fn serialize_element(&mut self, value: &T) -> Result<(), Self::Error> + where + T: Serialize + ?Sized, + { + let mut labels = value.serialize(LabelSetSerializer)?; + if labels.len() != 1 { + return Err(LabelError::new( + "label sequences must contain name-value pairs", + )); + } + self.labels + .push(labels.pop().expect("one label was encoded")); + Ok(()) + } + + fn end(self) -> Result { + Ok(self.labels) + } +} + +pub(super) struct LabelTupleSerializer { + name: Option, + value: Option, +} + +impl SerializeTuple for LabelTupleSerializer { + type Ok = Vec; + type Error = LabelError; + + fn serialize_element(&mut self, value: &T) -> Result<(), Self::Error> + where + T: Serialize + ?Sized, + { + if self.name.is_none() { + let name = value.serialize(LabelValueSerializer)?; + validate_label_name(&name)?; + self.name = Some(name); + } else if self.value.is_none() { + self.value = Some(value.serialize(LabelValueSerializer)?); + } else { + return Err(LabelError::new( + "label pairs must contain exactly one name and value", + )); + } + + Ok(()) + } + + fn end(self) -> Result { + let name = self + .name + .ok_or_else(|| LabelError::new("label pair is missing its name"))?; + let value = self + .value + .ok_or_else(|| LabelError::new("label pair is missing its value"))?; + + Ok(vec![LabelPair { + name: Some(name), + value: Some(value), + }]) + } +} + // Adapted from prometools' `serde::top::StructSerializer` // (https://github.com/nox/prometools, licensed MIT OR Apache-2.0). pub(super) struct LabelPairSerializer { @@ -396,16 +477,9 @@ impl Serializer for LabelValueSerializer { // Adapted from prometools' `serde::top::check_key` // (https://github.com/nox/prometools, licensed MIT OR Apache-2.0). fn validate_label_name(name: &str) -> Result<(), LabelError> { - let mut chars = name.chars(); - let valid = chars - .next() - .is_some_and(|character| character.is_ascii_alphabetic() || matches!(character, '_' | ':')) - && chars - .all(|character| character.is_ascii_alphanumeric() || matches!(character, '_' | ':')); - - valid.then_some(()).ok_or_else(|| { + is_valid_name(name).then_some(()).ok_or_else(|| { LabelError::new(format!( - "invalid metric label name {name:?}: expected [a-zA-Z_:][a-zA-Z0-9_:]*" + "invalid metric label name {name:?}: expected {NAME_REQUIREMENT}" )) }) } @@ -474,22 +548,58 @@ mod tests { ); } + #[test] + fn serializes_legacy_name_value_sequences() { + let pairs = to_label_pairs(&vec![ + ("trace_id", "abc"), + ("span_id", "def"), + ("trace:id", "ghi"), + ]) + .unwrap(); + let values: Vec<_> = pairs + .iter() + .map(|pair| { + ( + pair.name.as_deref().unwrap(), + pair.value.as_deref().unwrap(), + ) + }) + .collect(); + + assert_eq!( + values, + [("trace_id", "abc"), ("span_id", "def"), ("trace:id", "ghi"),] + ); + } + #[test] fn unit_is_an_empty_label_set() { assert!(to_label_pairs(&()).unwrap().is_empty()); } #[test] - fn rejects_invalid_label_names() { + fn rejects_empty_label_names() { #[derive(Serialize)] struct Invalid { - #[serde(rename = "not-valid")] + #[serde(rename = "")] value: &'static str, } assert!(to_label_pairs(&Invalid { value: "x" }).is_err()); } + #[test] + fn serializes_utf8_label_names() { + #[derive(Serialize)] + struct Labels { + #[serde(rename = "trace.id λ\n\"")] + value: &'static str, + } + + let labels = to_label_pairs(&Labels { value: "x" }).unwrap(); + assert_eq!(labels[0].name.as_deref(), Some("trace.id λ\n\"")); + } + #[test] fn rejects_compound_label_values() { #[derive(Serialize)] diff --git a/foundations-metrics/src/lib.rs b/foundations-metrics/src/lib.rs index 9e2b2e73..eb33dfe4 100644 --- a/foundations-metrics/src/lib.rs +++ b/foundations-metrics/src/lib.rs @@ -6,20 +6,26 @@ //! shared process-global registry and the stable wire format. #![warn(missing_docs)] +mod collect; mod diagnostics; +mod encoding; mod labels; pub mod metrics; mod registered; +mod validation; mod value; +pub use collect::{CollectionOptions, ServiceNameFormat, collect}; pub use diagnostics::{CollectErrorHookAlreadySet, set_collect_error_hook}; +pub use encoding::{OPENMETRICS_CONTENT_TYPE, encode_to_protobuf, encode_to_text}; pub use foundations_metrics_registry::{ EncodeMetric, IntoMetrics, MetricFamily, RegistrationMetadata, register, }; pub use labels::{LabelError, to_label_pairs}; pub use metrics::{ - Counter, CounterAtomic, Family, FamilyMetricGuard, Gauge, GaugeAtomic, GaugeGuard, Histogram, - HistogramBuilder, HistogramSnapshot, HistogramTimer, MetricConstructor, NativeHistogram, - NativeHistogramBuilder, RangeGauge, TimeHistogram, + Counter, CounterAtomic, CounterWithExemplar, Exemplar, Family, FamilyMetricGuard, Gauge, + GaugeAtomic, GaugeGuard, Histogram, HistogramBuilder, HistogramSnapshot, HistogramTimer, + HistogramWithExemplars, MetricConstructor, NativeHistogram, NativeHistogramBuilder, + NativeHistogramWithExemplars, RangeGauge, TimeHistogram, }; pub use registered::NamedMetric; diff --git a/foundations-metrics/src/metrics/counter.rs b/foundations-metrics/src/metrics/counter.rs index 31aaa84d..3b7d2eb9 100644 --- a/foundations-metrics/src/metrics/counter.rs +++ b/foundations-metrics/src/metrics/counter.rs @@ -135,7 +135,7 @@ impl Default for Counter { /// /// The name/help are left empty here; they are populated at registration and /// encode time. -fn encode_counter(value: f64) -> Vec { +pub(super) fn encode_counter(value: f64) -> Vec { vec![MetricFamily { name: Some(String::new()), help: None, diff --git a/foundations-metrics/src/metrics/exemplar.rs b/foundations-metrics/src/metrics/exemplar.rs new file mode 100644 index 00000000..6d8a84db --- /dev/null +++ b/foundations-metrics/src/metrics/exemplar.rs @@ -0,0 +1,710 @@ +use std::collections::HashMap; +use std::marker::PhantomData; +use std::sync::Arc; +use std::sync::atomic::AtomicU64; +use std::time::SystemTime; + +use foundations_metrics_registry::proto; +use parking_lot::{MappedRwLockReadGuard, RwLock, RwLockReadGuard}; +use prost_types::Timestamp; +use serde::Serialize; + +use super::counter::encode_counter; +use super::histogram::encode_snapshot; +use super::{ + CounterAtomic, Histogram, HistogramBuilder, IntoF64, MetricConstructor, NativeHistogram, + NativeHistogramBuilder, +}; +use crate::diagnostics::report_collect_error; +use crate::labels::to_label_pairs; +use crate::validation::EXEMPLAR_SERIALIZATION_ERROR_LABEL; +use crate::{MetricFamily, value::EncodeMetricValue}; + +/// Labels and a sampled value associated with a metric observation. +/// +/// Exemplar values are exposed through [`CounterWithExemplar::get`]. Their fields +/// remain private; exemplars are created by recording labeled metric updates. +#[derive(Debug)] +pub struct Exemplar { + label_set: Arc, + value: V, + timestamp: Option, +} + +impl Exemplar { + /// Returns the exemplar's label set. + pub fn label_set(&self) -> &S { + &self.label_set + } + + /// Returns the sampled increment or observation. + pub fn value(&self) -> &V { + &self.value + } +} + +impl Exemplar { + fn new(label_set: S, value: V, timestamp: Option) -> Self { + Self { + label_set: Arc::new(label_set), + value, + timestamp, + } + } +} + +impl Clone for Exemplar { + fn clone(&self) -> Self { + Self { + label_set: Arc::clone(&self.label_set), + value: self.value.clone(), + timestamp: self.timestamp, + } + } +} + +impl Exemplar { + fn encode(&self) -> Result + where + S: Serialize, + V: Clone + IntoF64, + { + Ok(proto::Exemplar { + label: to_label_pairs(self.label_set.as_ref()).map_err(|error| error.to_string())?, + value: Some(self.value.clone().into_f64()), + timestamp: self.timestamp, + }) + } +} + +fn finish_exemplar(exemplar: Option>) -> Option { + match exemplar { + Some(Ok(exemplar)) => Some(exemplar), + // Defer reporting to validation, after any enclosing Family lock has + // been released. + Some(Err(error)) => Some(proto::Exemplar { + label: vec![proto::LabelPair { + name: Some(EXEMPLAR_SERIALIZATION_ERROR_LABEL.to_owned()), + value: Some(error), + }], + ..Default::default() + }), + None => None, + } +} + +/// A monotonically increasing counter that records an exemplar with an update. +/// +/// Calling [`inc_by`](Self::inc_by) with a label set replaces the previous +/// exemplar. Calling it without a label set clears the previous exemplar. The +/// exemplar value is the increment. [`get`](Self::get) returns the cumulative +/// counter value and read-only access to the current exemplar. Clones share both +/// values. +/// +/// Exemplar labels are serialized with [`serde::Serialize`] during collection. +#[derive(Debug)] +pub struct CounterWithExemplar { + state: Arc>>, + marker: PhantomData, +} + +#[derive(Debug)] +struct CounterWithExemplarState { + value: A, + exemplar: Option>, +} + +impl CounterWithExemplar +where + N: Clone, + A: CounterAtomic, +{ + /// Increments the counter by `value`, returning the previous total. + /// + /// A provided label set replaces the stored exemplar. `None` clears it. + pub fn inc_by(&self, value: N, label_set: Option) -> N { + let exemplar = label_set.map(|label_set| Exemplar::new(label_set, value.clone(), None)); + let mut state = self.state.write(); + state.exemplar = exemplar; + state.value.inc_by(value) + } + + /// Returns the cumulative counter value and the current exemplar. + /// + /// The exemplar guard keeps the counter read-locked until it is dropped. + pub fn get(&self) -> (N, MappedRwLockReadGuard<'_, Option>>) { + let state = self.state.read(); + let value = state.value.get(); + let exemplar = RwLockReadGuard::map(state, |state| &state.exemplar); + (value, exemplar) + } + + /// Returns read-only access to the underlying counter storage. + /// + /// The returned guard prevents updates until it is dropped. + pub fn inner(&self) -> MappedRwLockReadGuard<'_, A> { + RwLockReadGuard::map(self.state.read(), |state| &state.value) + } +} + +impl Clone for CounterWithExemplar { + fn clone(&self) -> Self { + Self { + state: Arc::clone(&self.state), + marker: PhantomData, + } + } +} + +impl Default for CounterWithExemplar { + fn default() -> Self { + Self { + state: Arc::new(RwLock::new(CounterWithExemplarState { + value: A::default(), + exemplar: None, + })), + marker: PhantomData, + } + } +} + +impl EncodeMetricValue for CounterWithExemplar +where + S: Serialize + Send + Sync + 'static, + N: Clone + IntoF64, + A: CounterAtomic + Send + Sync + 'static, +{ + fn encode_metric_value(&self) -> Vec { + let state = self.state.read(); + let value = state.value.get().into_f64(); + let exemplar = state.exemplar.clone(); + drop(state); + + let mut families = encode_counter(value); + let counter = families[0].metric[0] + .counter + .as_mut() + .expect("counter encoder always creates a counter value"); + counter.exemplar = finish_exemplar(exemplar.as_ref().map(Exemplar::encode)); + families + } +} + +/// A classic fixed-bucket histogram that retains one exemplar per bucket. +/// +/// A labeled observation replaces the exemplar in its bucket. An unlabeled +/// observation updates the histogram without clearing an existing exemplar. +/// Clones share histogram and exemplar storage. +#[derive(Clone, Debug)] +pub struct HistogramWithExemplars { + state: Arc>>, +} + +#[derive(Debug)] +struct HistogramWithExemplarsState { + histogram: Histogram, + exemplars: HashMap>, +} + +impl HistogramWithExemplars { + /// Creates a histogram with the provided inclusive upper bounds. + /// + /// Bounds are sorted and a terminal `f64::MAX` bucket is appended. + pub fn new(bounds: impl IntoIterator) -> Self { + Self { + state: Arc::new(RwLock::new(HistogramWithExemplarsState { + histogram: Histogram::new(bounds), + exemplars: HashMap::new(), + })), + } + } + + /// Records an observation and optionally associates it with an exemplar. + pub fn observe(&self, value: f64, label_set: Option) { + let exemplar = label_set.map(|label_set| Exemplar::new(label_set, value, None)); + let mut state = self.state.write(); + let bucket = state.histogram.observe_and_bucket(value); + + if let (Some(bucket), Some(exemplar)) = (bucket, exemplar) { + state.exemplars.insert(bucket, exemplar); + } + } +} + +impl EncodeMetricValue for HistogramWithExemplars +where + S: Serialize + Send + Sync + 'static, +{ + fn encode_metric_value(&self) -> Vec { + let state = self.state.read(); + let snapshot = state.histogram.snapshot(); + let exemplars: Vec<_> = state + .exemplars + .iter() + .map(|(&index, exemplar)| (index, exemplar.clone())) + .collect(); + drop(state); + + let mut families = encode_snapshot(snapshot); + let buckets = &mut families[0].metric[0] + .histogram + .as_mut() + .expect("histogram encoder always creates a histogram value") + .bucket; + + for (index, exemplar) in exemplars { + if let Some(bucket) = buckets.get_mut(index) { + bucket.exemplar = finish_exemplar(Some(exemplar.encode())); + } + } + + families + } +} + +impl MetricConstructor> for HistogramBuilder { + fn new_metric(&self) -> HistogramWithExemplars { + HistogramWithExemplars::new(self.buckets.iter().copied()) + } +} + +/// A native histogram that retains its latest labeled observation as an exemplar. +/// +/// Native histogram exemplars require timestamps in the Prometheus protobuf +/// model. The timestamp is captured when the labeled observation is recorded. +/// Native histograms require protobuf exposition; OpenMetrics text encoding +/// does not expose their sparse buckets or exemplars. +#[derive(Clone, Debug)] +pub struct NativeHistogramWithExemplars { + state: Arc>>, +} + +#[derive(Debug)] +struct NativeHistogramWithExemplarsState { + histogram: NativeHistogram, + exemplar: Option>, +} + +impl NativeHistogramWithExemplars { + /// Creates a native histogram with the given bucket growth `factor`. + /// + /// # Panics + /// + /// Panics if `factor` is not greater than `1.0`. + pub fn new(factor: f64) -> Self { + NativeHistogramBuilder::new(factor).new_metric() + } + + /// Records an observation and optionally retains it as the latest exemplar. + pub fn observe(&self, value: f64, label_set: Option) { + let exemplar = label_set.map(|label_set| Exemplar::new(label_set, value, None)); + let mut state = self.state.write(); + state.histogram.observe(value); + + if let Some(mut exemplar) = exemplar { + exemplar.timestamp = Some(SystemTime::now().into()); + state.exemplar = Some(exemplar); + } + } +} + +impl EncodeMetricValue for NativeHistogramWithExemplars +where + S: Serialize + Send + Sync + 'static, +{ + fn encode_metric_value(&self) -> Vec { + let state = self.state.read(); + let families = state.histogram.try_encode_metric_value(); + let exemplar = state.exemplar.clone(); + drop(state); + + let mut families = match families { + Ok(families) => families, + Err(error) => { + report_collect_error(format_args!( + "non-fatal error while collecting metrics: skipped a native histogram; protobuf encoding failed: {error}" + )); + return Vec::new(); + } + }; + + if let Some(exemplar) = finish_exemplar(exemplar.as_ref().map(Exemplar::encode)) { + for histogram in families + .iter_mut() + .flat_map(|family| &mut family.metric) + .filter_map(|metric| metric.histogram.as_mut()) + { + histogram.exemplars.push(exemplar.clone()); + } + } + + families + } +} + +impl MetricConstructor> for NativeHistogramBuilder { + fn new_metric(&self) -> NativeHistogramWithExemplars { + NativeHistogramWithExemplars { + state: Arc::new(RwLock::new(NativeHistogramWithExemplarsState { + histogram: NativeHistogramBuilder::new_metric(self), + exemplar: None, + })), + } + } +} + +#[cfg(test)] +mod tests { + use foundations_metrics_registry::proto::MetricType; + use prost::Message; + use serde::Serialize; + + use super::*; + use crate::{ + CollectionOptions, EncodeMetric, Family, NamedMetric, RegistrationMetadata, + ServiceNameFormat, collect, encode_to_protobuf, encode_to_text, register, + }; + + #[derive(Clone, Debug, Serialize)] + struct TraceLabels { + trace_id: &'static str, + } + + fn trace_id(exemplar: &proto::Exemplar) -> Option<&str> { + exemplar + .label + .iter() + .find(|label| label.name.as_deref() == Some("trace_id")) + .and_then(|label| label.value.as_deref()) + } + + #[test] + fn counter_replaces_and_clears_exemplars_and_clones_share_storage() { + let counter = CounterWithExemplar::::default(); + let clone = counter.clone(); + + assert_eq!( + counter.inc_by(2, Some(TraceLabels { trace_id: "first" })), + 0 + ); + assert_eq!(clone.inc_by(3, Some(TraceLabels { trace_id: "latest" })), 2); + + let families = counter.encode_metric_value(); + let encoded = families[0].metric[0].counter.as_ref().unwrap(); + let exemplar = encoded.exemplar.as_ref().unwrap(); + assert_eq!(encoded.value, Some(5.0)); + assert_eq!(exemplar.value, Some(3.0)); + assert_eq!(trace_id(exemplar), Some("latest")); + assert!(exemplar.timestamp.is_none()); + + let (value, exemplar) = counter.get(); + assert_eq!(value, 5); + assert!(exemplar.is_some()); + drop(exemplar); + assert_eq!( + counter.inner().load(std::sync::atomic::Ordering::Relaxed), + 5 + ); + + assert_eq!(counter.inc_by(1, None), 5); + assert_eq!(counter.get().0, 6); + assert!( + counter.encode_metric_value()[0].metric[0] + .counter + .as_ref() + .unwrap() + .exemplar + .is_none() + ); + } + + #[test] + fn classic_histogram_retains_latest_exemplar_per_bucket() { + let histogram = HistogramWithExemplars::new([1.0, 2.0]); + histogram.observe(0.5, Some(TraceLabels { trace_id: "first" })); + histogram.observe( + 0.75, + Some(TraceLabels { + trace_id: "replacement", + }), + ); + histogram.observe(0.8, None); + histogram.observe( + 1.5, + Some(TraceLabels { + trace_id: "second_bucket", + }), + ); + + let families = histogram.encode_metric_value(); + let encoded = families[0].metric[0].histogram.as_ref().unwrap(); + assert_eq!(encoded.sample_count, Some(4)); + assert_eq!(encoded.sample_sum, Some(3.55)); + assert_eq!( + encoded + .bucket + .iter() + .map(|bucket| bucket.cumulative_count) + .collect::>(), + [Some(3), Some(4), Some(4)] + ); + assert_eq!( + trace_id(encoded.bucket[0].exemplar.as_ref().unwrap()), + Some("replacement") + ); + assert_eq!( + encoded.bucket[0].exemplar.as_ref().unwrap().value, + Some(0.75) + ); + assert_eq!( + trace_id(encoded.bucket[1].exemplar.as_ref().unwrap()), + Some("second_bucket") + ); + assert!(encoded.bucket[2].exemplar.is_none()); + } + + #[test] + fn native_histogram_retains_latest_timestamped_exemplar() { + let histogram = NativeHistogramWithExemplars::new(1.1); + let clone = histogram.clone(); + histogram.observe(0.5, Some(TraceLabels { trace_id: "first" })); + clone.observe(2.0, None); + clone.observe(3.0, Some(TraceLabels { trace_id: "latest" })); + + let families = histogram.encode_metric_value(); + let encoded = families[0].metric[0].histogram.as_ref().unwrap(); + assert_eq!(encoded.sample_count, Some(3)); + assert_eq!(encoded.sample_sum, Some(5.5)); + assert!(!encoded.positive_span.is_empty()); + assert_eq!(encoded.exemplars.len(), 1); + assert_eq!(encoded.exemplars[0].value, Some(3.0)); + assert_eq!(trace_id(&encoded.exemplars[0]), Some("latest")); + assert!(encoded.exemplars[0].timestamp.is_some()); + } + + #[test] + fn exemplar_label_failure_drops_only_the_exemplar() { + let counter = CounterWithExemplar::<&'static str>::default(); + counter.inc_by(2, Some("not a label set")); + + let families = NamedMetric::new("serialization_failure", "", counter).encode(); + let sentinel = &families[0].metric[0] + .counter + .as_ref() + .unwrap() + .exemplar + .as_ref() + .unwrap() + .label[0]; + assert_eq!( + sentinel.name.as_deref(), + Some(EXEMPLAR_SERIALIZATION_ERROR_LABEL) + ); + assert_eq!( + sentinel.value.as_deref(), + Some("metric labels must serialize as a struct or unit") + ); + + let payload = encode_to_protobuf(&families); + let encoded = MetricFamily::decode_length_delimited(payload.as_slice()).unwrap(); + let encoded = encoded.metric[0].counter.as_ref().unwrap(); + assert_eq!(encoded.value, Some(2.0)); + assert!(encoded.exemplar.is_none()); + } + + #[test] + fn empty_exemplar_label_sets_are_retained_in_text() { + let counter = CounterWithExemplar::<()>::default(); + counter.inc_by(2, Some(())); + + let families = NamedMetric::new("empty_exemplar", "", counter).encode(); + assert!(encode_to_text(&families).contains("empty_exemplar 2.0 # {} 2.0\n")); + } + + #[test] + fn accepts_legacy_sequence_label_sets() { + let counter = CounterWithExemplar::>::default(); + counter.inc_by(1, Some(vec![("trace_id", "legacy")])); + + let families = counter.encode_metric_value(); + let exemplar = families[0].metric[0] + .counter + .as_ref() + .unwrap() + .exemplar + .as_ref() + .unwrap(); + assert_eq!(trace_id(exemplar), Some("legacy")); + } + + #[test] + fn updates_do_not_require_serializable_label_sets() { + struct OpaqueLabels; + + let counter = CounterWithExemplar::::default(); + counter.inc_by(1, Some(OpaqueLabels)); + assert_eq!(counter.get().0, 1); + + let classic = HistogramWithExemplars::new([1.0]); + classic.observe(0.5, Some(OpaqueLabels)); + + let native = NativeHistogramWithExemplars::new(1.1); + native.observe(0.5, Some(OpaqueLabels)); + } + + #[test] + fn families_keep_series_and_exemplar_labels_separate() { + #[derive(Clone, Debug, Eq, Hash, PartialEq, Serialize)] + struct SeriesLabels { + method: &'static str, + } + + let family = Family::>::default(); + family + .get_or_create(&SeriesLabels { method: "GET" }) + .inc_by(1, Some(TraceLabels { trace_id: "abc" })); + + let families = family.encode_metric_value(); + let metric = &families[0].metric[0]; + assert!(metric.label.iter().any(|label| { + label.name.as_deref() == Some("method") && label.value.as_deref() == Some("GET") + })); + assert_eq!( + trace_id(metric.counter.as_ref().unwrap().exemplar.as_ref().unwrap()), + Some("abc") + ); + } + + #[test] + fn histogram_builders_construct_exemplar_metrics() { + let classic: HistogramWithExemplars = HistogramBuilder { + buckets: &[0.5, 1.0], + } + .new_metric(); + classic.observe( + 0.75, + Some(TraceLabels { + trace_id: "classic", + }), + ); + assert!( + classic.encode_metric_value()[0].metric[0] + .histogram + .as_ref() + .unwrap() + .bucket[1] + .exemplar + .is_some() + ); + + let native: NativeHistogramWithExemplars = NativeHistogramBuilder::new(1.1) + .with_max_buckets(160) + .new_metric(); + native.observe(0.75, Some(TraceLabels { trace_id: "native" })); + assert_eq!( + native.encode_metric_value()[0].metric[0] + .histogram + .as_ref() + .unwrap() + .exemplars + .len(), + 1 + ); + } + + #[test] + fn registered_exemplar_counter_is_collected_and_encoded() { + let counter = CounterWithExemplar::::default(); + counter.inc_by( + 4, + Some(TraceLabels { + trace_id: "registered", + }), + ); + register( + Box::new(NamedMetric::new( + "registered_counter_with_exemplar", + "A registered exemplar counter.", + counter, + )) as Box, + RegistrationMetadata::default(), + ); + + let families = collect(CollectionOptions { + include_optional: false, + service_name: None, + service_name_format: ServiceNameFormat::MetricPrefix, + }); + let family = families + .iter() + .find(|family| family.name.as_deref() == Some("registered_counter_with_exemplar")) + .expect("registered exemplar counter is collected"); + assert_eq!(family.r#type, Some(MetricType::Counter as i32)); + assert_eq!( + trace_id( + family.metric[0] + .counter + .as_ref() + .unwrap() + .exemplar + .as_ref() + .unwrap() + ), + Some("registered") + ); + + let text = encode_to_text(std::slice::from_ref(family)); + assert!( + text.contains("registered_counter_with_exemplar 4.0 # {trace_id=\"registered\"} 4.0\n") + ); + } + + #[test] + fn registered_native_exemplar_family_round_trips_through_protobuf() { + #[derive(Clone, Debug, Eq, Hash, PartialEq, Serialize)] + struct SeriesLabels { + method: &'static str, + } + + let family = Family::< + SeriesLabels, + NativeHistogramWithExemplars, + NativeHistogramBuilder, + >::new_with_constructor(NativeHistogramBuilder::new(1.1)); + family + .get_or_create(&SeriesLabels { method: "GET" }) + .observe(0.25, Some(TraceLabels { trace_id: "native" })); + register( + Box::new(NamedMetric::new( + "registered_native_histogram_with_exemplars", + "A registered native histogram with exemplars.", + family, + )) as Box, + RegistrationMetadata::default(), + ); + + let families = collect(CollectionOptions { + include_optional: false, + service_name: None, + service_name_format: ServiceNameFormat::MetricPrefix, + }); + let family = families + .iter() + .find(|family| { + family.name.as_deref() == Some("registered_native_histogram_with_exemplars") + }) + .expect("registered native histogram family is collected"); + let payload = encode_to_protobuf(std::slice::from_ref(family)); + let mut bytes = payload.as_slice(); + let decoded = MetricFamily::decode_length_delimited(&mut bytes).unwrap(); + let metric = &decoded.metric[0]; + assert!(metric.label.iter().any(|label| { + label.name.as_deref() == Some("method") && label.value.as_deref() == Some("GET") + })); + let exemplars = &metric.histogram.as_ref().unwrap().exemplars; + assert_eq!(exemplars.len(), 1); + assert_eq!(trace_id(&exemplars[0]), Some("native")); + assert!(exemplars[0].timestamp.is_some()); + assert!(bytes.is_empty()); + } +} diff --git a/foundations-metrics/src/metrics/histogram.rs b/foundations-metrics/src/metrics/histogram.rs index e58d71e9..0de302f0 100644 --- a/foundations-metrics/src/metrics/histogram.rs +++ b/foundations-metrics/src/metrics/histogram.rs @@ -66,21 +66,27 @@ impl Histogram { /// Records an observed value. pub fn observe(&self, value: f64) { + self.observe_and_bucket(value); + } + + pub(super) fn observe_and_bucket(&self, value: f64) -> Option { let mut state = self.state.lock(); state.sum += value; state.count = state.count.wrapping_add(1); - let bucket = if value.is_nan() { - None - } else { - let index = state - .buckets - .partition_point(|(upper_bound, _)| *upper_bound < value); - state.buckets.get_mut(index) - }; + if value.is_nan() { + return None; + } - if let Some((_, count)) = bucket { + let index = state + .buckets + .partition_point(|(upper_bound, _)| *upper_bound < value); + + if let Some((_, count)) = state.buckets.get_mut(index) { *count = count.wrapping_add(1); + Some(index) + } else { + None } } @@ -102,7 +108,7 @@ impl EncodeMetricValue for Histogram { } } -fn encode_snapshot(snapshot: HistogramSnapshot) -> Vec { +pub(super) fn encode_snapshot(snapshot: HistogramSnapshot) -> Vec { let mut cumulative_count = 0_u64; let buckets = snapshot .buckets diff --git a/foundations-metrics/src/metrics/mod.rs b/foundations-metrics/src/metrics/mod.rs index d7eda4d8..72a6176b 100644 --- a/foundations-metrics/src/metrics/mod.rs +++ b/foundations-metrics/src/metrics/mod.rs @@ -35,12 +35,16 @@ use std::sync::atomic::{AtomicU64, Ordering}; mod counter; +mod exemplar; mod family; mod gauge; mod histogram; mod native_histogram; pub use counter::{Counter, CounterAtomic}; +pub use exemplar::{ + CounterWithExemplar, Exemplar, HistogramWithExemplars, NativeHistogramWithExemplars, +}; pub use family::{Family, FamilyMetricGuard, MetricConstructor}; pub use gauge::{Gauge, GaugeAtomic, GaugeGuard, RangeGauge}; pub use histogram::{ diff --git a/foundations-metrics/src/metrics/native_histogram.rs b/foundations-metrics/src/metrics/native_histogram.rs index 812be320..91da78ca 100644 --- a/foundations-metrics/src/metrics/native_histogram.rs +++ b/foundations-metrics/src/metrics/native_histogram.rs @@ -56,6 +56,16 @@ impl NativeHistogram { pub fn observe(&self, value: f64) { self.inner.observe(value); } + + pub(super) fn try_encode_metric_value(&self) -> Result, std::fmt::Error> { + // Upstream keeps native bucket state private. A cloned histogram shares + // storage, so a temporary registry can drive its protobuf encoder. + let mut registry = Registry::default(); + registry.register("native_histogram", "", self.inner.clone()); + + prometheus_protobuf::encode(®istry) + .map(|families| families.into_iter().map(convert_native_family).collect()) + } } /// Constructs [`NativeHistogram`]s with a fixed configuration. @@ -164,13 +174,8 @@ impl MetricConstructor for NativeHistogramBuilder { impl EncodeMetricValue for NativeHistogram { fn encode_metric_value(&self) -> Vec { - // Upstream keeps native bucket state private. A cloned histogram shares - // storage, so a temporary registry can drive its protobuf encoder. - let mut registry = Registry::default(); - registry.register("native_histogram", "", self.inner.clone()); - - match prometheus_protobuf::encode(®istry) { - Ok(families) => families.into_iter().map(convert_native_family).collect(), + match self.try_encode_metric_value() { + Ok(families) => families, Err(error) => { report_collect_error(format_args!( "non-fatal error while collecting metrics: skipped a native histogram; protobuf encoding failed: {error}" @@ -325,7 +330,7 @@ mod tests { #[test] fn builder_applies_configuration() { - let histogram = NativeHistogramBuilder::new(1.5) + let histogram: NativeHistogram = NativeHistogramBuilder::new(1.5) .with_zero_threshold(0.001) .with_max_buckets(160) .new_metric(); diff --git a/foundations-metrics/src/validation.rs b/foundations-metrics/src/validation.rs new file mode 100644 index 00000000..243912a8 --- /dev/null +++ b/foundations-metrics/src/validation.rs @@ -0,0 +1,310 @@ +use std::borrow::Cow; + +use foundations_metrics_registry::proto::{Exemplar, LabelPair, Metric, MetricFamily, MetricType}; + +use crate::diagnostics::report_collect_error; + +pub(crate) const NAME_REQUIREMENT: &str = "a non-empty UTF-8 string without NUL bytes"; +pub(crate) const EXEMPLAR_SERIALIZATION_ERROR_LABEL: &str = + "\0foundations_metrics_exemplar_serialization_error"; + +#[derive(Clone, Copy)] +pub(crate) enum ValidationContext { + Collection, + TextEncoding, + ProtobufEncoding, +} + +impl ValidationContext { + fn action(self) -> &'static str { + match self { + Self::Collection => "collecting metrics", + Self::TextEncoding => "encoding OpenMetrics text", + Self::ProtobufEncoding => "encoding Prometheus protobuf", + } + } +} + +pub(crate) fn is_valid_name(name: &str) -> bool { + !name.is_empty() && !name.contains('\0') +} + +pub(crate) fn sanitize_metric_family( + family: &mut MetricFamily, + context: ValidationContext, +) -> bool { + let Some(name) = family.name.as_deref().filter(|name| is_valid_name(name)) else { + report_invalid_family_name(context, family.name.as_deref()); + return false; + }; + + sanitize_rows( + &mut family.metric, + family + .r#type + .and_then(|value| MetricType::try_from(value).ok()), + name, + context, + ); + true +} + +pub(crate) fn sanitized_metric_family<'a>( + family: &'a MetricFamily, + context: ValidationContext, +) -> Option> { + let Some(name) = family.name.as_deref().filter(|name| is_valid_name(name)) else { + report_invalid_family_name(context, family.name.as_deref()); + return None; + }; + let metric_type = family + .r#type + .and_then(|value| MetricType::try_from(value).ok()); + + if rows_are_valid(&family.metric, metric_type) { + return Some(Cow::Borrowed(family)); + } + + let mut sanitized = family.clone(); + sanitize_rows(&mut sanitized.metric, metric_type, name, context); + Some(Cow::Owned(sanitized)) +} + +fn report_invalid_family_name(context: ValidationContext, name: Option<&str>) { + report_collect_error(format_args!( + "non-fatal error while {}: skipped metric family with invalid name {name:?}; expected {NAME_REQUIREMENT}", + context.action() + )); +} + +fn rows_are_valid(metrics: &[Metric], metric_type: Option) -> bool { + let reserved_label = reserved_row_label(metric_type); + metrics.iter().all(|metric| { + find_label_issue(&metric.label, reserved_label).is_none() && exemplars_are_valid(metric) + }) +} + +fn sanitize_rows( + metrics: &mut Vec, + metric_type: Option, + family_name: &str, + context: ValidationContext, +) { + let reserved_label = reserved_row_label(metric_type); + metrics.retain_mut(|metric| { + if let Some(issue) = find_label_issue(&metric.label, reserved_label) { + report_row_drop(context, family_name, issue); + return false; + } + + sanitize_exemplars(metric, family_name, context); + true + }); +} + +fn reserved_row_label(metric_type: Option) -> Option<&'static str> { + match metric_type { + Some(MetricType::Histogram | MetricType::GaugeHistogram) => Some("le"), + Some(MetricType::Summary) => Some("quantile"), + _ => None, + } +} + +#[derive(Clone, Copy)] +enum LabelIssue<'a> { + Invalid(Option<&'a str>), + Duplicate(&'a str), + Reserved(&'a str), +} + +enum ExemplarIssue<'a> { + Label(LabelIssue<'a>), + Serialization(&'a str), +} + +fn find_label_issue<'a>( + labels: &'a [LabelPair], + reserved_label: Option<&str>, +) -> Option> { + for (index, label) in labels.iter().enumerate() { + let Some(name) = label.name.as_deref().filter(|name| is_valid_name(name)) else { + return Some(LabelIssue::Invalid(label.name.as_deref())); + }; + + if labels[..index] + .iter() + .any(|previous| previous.name.as_deref() == Some(name)) + { + return Some(LabelIssue::Duplicate(name)); + } + if reserved_label == Some(name) { + return Some(LabelIssue::Reserved(name)); + } + } + + None +} + +fn report_row_drop(context: ValidationContext, family_name: &str, issue: LabelIssue<'_>) { + match issue { + LabelIssue::Invalid(name) => report_collect_error(format_args!( + "non-fatal error while {}: skipped row in metric family {family_name:?} with invalid label name {name:?}; expected {NAME_REQUIREMENT}", + context.action() + )), + LabelIssue::Duplicate(name) => report_collect_error(format_args!( + "non-fatal error while {}: skipped row in metric family {family_name:?} with duplicate label name {name:?}", + context.action() + )), + LabelIssue::Reserved(name) => report_collect_error(format_args!( + "non-fatal error while {}: skipped row in metric family {family_name:?}; label name {name:?} is reserved for this metric type", + context.action() + )), + } +} + +fn exemplars_are_valid(metric: &Metric) -> bool { + metric + .counter + .as_ref() + .and_then(|counter| counter.exemplar.as_ref()) + .is_none_or(exemplar_is_valid) + && metric.histogram.as_ref().is_none_or(|histogram| { + histogram + .bucket + .iter() + .all(|bucket| bucket.exemplar.as_ref().is_none_or(exemplar_is_valid)) + && histogram.exemplars.iter().all(native_exemplar_is_valid) + }) +} + +fn exemplar_is_valid(exemplar: &Exemplar) -> bool { + find_exemplar_issue(&exemplar.label).is_none() +} + +fn native_exemplar_is_valid(exemplar: &Exemplar) -> bool { + exemplar.timestamp.is_some() && exemplar_is_valid(exemplar) +} + +fn sanitize_exemplars(metric: &mut Metric, family_name: &str, context: ValidationContext) { + if let Some(counter) = &mut metric.counter { + sanitize_exemplar_slot(&mut counter.exemplar, "counter", family_name, context); + } + + if let Some(histogram) = &mut metric.histogram { + for bucket in &mut histogram.bucket { + sanitize_exemplar_slot( + &mut bucket.exemplar, + "classic histogram bucket", + family_name, + context, + ); + } + histogram.exemplars.retain(|exemplar| { + if let Some(issue) = find_exemplar_issue(&exemplar.label) { + report_exemplar_drop(context, family_name, "native histogram", issue); + false + } else if exemplar.timestamp.is_none() { + report_collect_error(format_args!( + "non-fatal error while {}: dropped native histogram exemplar in metric family {family_name:?} without a required timestamp", + context.action() + )); + false + } else { + true + } + }); + } +} + +fn sanitize_exemplar_slot( + exemplar: &mut Option, + kind: &str, + family_name: &str, + context: ValidationContext, +) { + let Some(issue) = exemplar + .as_ref() + .and_then(|exemplar| find_exemplar_issue(&exemplar.label)) + else { + return; + }; + + report_exemplar_drop(context, family_name, kind, issue); + *exemplar = None; +} + +fn report_exemplar_drop( + context: ValidationContext, + family_name: &str, + kind: &str, + issue: ExemplarIssue<'_>, +) { + match issue { + ExemplarIssue::Label(LabelIssue::Invalid(name)) => report_collect_error(format_args!( + "non-fatal error while {}: dropped {kind} exemplar in metric family {family_name:?} with invalid label name {name:?}; expected {NAME_REQUIREMENT}", + context.action() + )), + ExemplarIssue::Label(LabelIssue::Duplicate(name)) => report_collect_error(format_args!( + "non-fatal error while {}: dropped {kind} exemplar in metric family {family_name:?} with duplicate label name {name:?}", + context.action() + )), + ExemplarIssue::Serialization(error) => report_collect_error(format_args!( + "non-fatal error while {}: dropped {kind} exemplar in metric family {family_name:?}; label serialization failed: {error}", + context.action() + )), + ExemplarIssue::Label(LabelIssue::Reserved(_)) => { + unreachable!("exemplar labels have no reserved names") + } + } +} + +fn find_exemplar_issue(labels: &[LabelPair]) -> Option> { + if labels.len() == 1 + && labels[0].name.as_deref() == Some(EXEMPLAR_SERIALIZATION_ERROR_LABEL) + && let Some(error) = labels[0].value.as_deref() + { + return Some(ExemplarIssue::Serialization(error)); + } + + find_label_issue(labels, None).map(ExemplarIssue::Label) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn accepts_non_empty_utf8_metric_and_label_names() { + for name in [ + "é", + "aλ", + "metric name", + "metric\nname", + "metric#name", + "metric\"name", + "指标.名称", + ] { + assert!(is_valid_name(name), "metric name {name:?}"); + } + } + + #[test] + fn rejects_empty_and_nul_names() { + for name in ["", "nul\0name"] { + assert!(!is_valid_name(name)); + } + } + + #[test] + fn preserves_exemplar_serialization_errors_for_diagnostics() { + let labels = [LabelPair { + name: Some(EXEMPLAR_SERIALIZATION_ERROR_LABEL.to_owned()), + value: Some("root cause".to_owned()), + }]; + + assert!(matches!( + find_exemplar_issue(&labels), + Some(ExemplarIssue::Serialization("root cause")) + )); + } +}