use std::sync::Arc;
use std::time::{Duration, Instant};
use hdrhistogram::Histogram as HdrHistogram;
use crate::labels::Labels;
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
pub enum CloseReason {
Quiesce,
ScopeClose,
Shutdown,
}
#[derive(Clone, Debug)]
pub struct MetricSet {
captured_at: Instant,
interval: Duration,
partial: bool,
close: Option<CloseReason>,
scheduled_ts: Option<Instant>,
families: Vec<MetricFamily>,
}
impl Default for MetricSet {
fn default() -> Self {
Self {
captured_at: Instant::now(),
interval: Duration::ZERO,
partial: false,
close: None,
scheduled_ts: None,
families: Vec::new(),
}
}
}
impl MetricSet {
pub fn new(interval: Duration) -> Self {
Self {
captured_at: Instant::now(),
interval,
partial: false,
close: None,
scheduled_ts: None,
families: Vec::new(),
}
}
pub fn at(captured_at: Instant, interval: Duration) -> Self {
Self {
captured_at,
interval,
partial: false,
close: None,
scheduled_ts: None,
families: Vec::new(),
}
}
pub fn captured_at(&self) -> Instant {
self.captured_at
}
pub fn actual_ts(&self) -> Instant {
self.captured_at
}
pub fn scheduled_ts(&self) -> Option<Instant> {
self.scheduled_ts
}
pub fn set_scheduled_ts(&mut self, scheduled: Instant) {
self.scheduled_ts = Some(scheduled);
}
pub fn interval(&self) -> Duration {
self.interval
}
pub fn is_partial(&self) -> bool {
self.partial
}
pub fn mark_partial(&mut self) {
self.partial = true;
}
pub fn close_reason(&self) -> Option<CloseReason> {
self.close
}
pub fn mark_close(&mut self, reason: CloseReason) {
self.close = Some(self.close.map_or(reason, |c| c.max(reason)));
}
pub fn set_interval(&mut self, interval: Duration) {
self.interval = interval;
}
pub fn families(&self) -> impl Iterator<Item = &MetricFamily> {
self.families.iter()
}
pub fn family(&self, name: &str) -> Option<&MetricFamily> {
self.families.iter().find(|f| f.name() == name)
}
pub fn len(&self) -> usize {
self.families.len()
}
pub fn is_empty(&self) -> bool {
self.families.is_empty()
}
pub fn has_distributions(&self) -> bool {
self.families.iter().any(|f| {
matches!(
f.r#type(),
MetricType::Histogram | MetricType::GaugeHistogram
)
})
}
pub fn without_distributions(&self) -> MetricSet {
MetricSet {
captured_at: self.captured_at,
interval: self.interval,
partial: self.partial,
close: self.close,
scheduled_ts: self.scheduled_ts,
families: self
.families
.iter()
.filter(|f| {
!matches!(
f.r#type(),
MetricType::Histogram | MetricType::GaugeHistogram
)
})
.cloned()
.collect(),
}
}
pub fn insert(&mut self, family: MetricFamily) {
assert!(
!self.families.iter().any(|f| f.name() == family.name()),
"MetricSet already contains family '{}' — names must be unique (OpenMetrics §4.1)",
family.name(),
);
self.families.push(family);
}
pub fn coalesce(snapshots: &[MetricSet]) -> MetricSet {
Self::coalesce_with_mode(snapshots, CombineMode::Coalesce)
}
pub fn coalesce_with_mode(snapshots: &[MetricSet], mode: CombineMode) -> MetricSet {
if snapshots.is_empty() {
return MetricSet::default();
}
if snapshots.len() == 1 {
return snapshots[0].clone();
}
let captured_at = snapshots.iter().map(|s| s.captured_at).max().unwrap();
let interval: Duration = snapshots.iter().map(|s| s.interval).sum();
let partial = snapshots.iter().any(|s| s.partial);
let scheduled_ts = None;
let close = snapshots.iter().filter_map(|s| s.close).max();
let mut out = MetricSet {
captured_at,
interval,
partial,
close,
scheduled_ts,
families: Vec::new(),
};
let mut seen_family: Vec<String> = Vec::new();
for s in snapshots {
for f in &s.families {
if !seen_family.contains(&f.name) {
seen_family.push(f.name.clone());
}
}
}
for fname in seen_family {
let mut acc: Option<MetricFamily> = None;
for s in snapshots {
let Some(src_family) = s.families.iter().find(|f| f.name == fname) else {
continue;
};
if acc.is_none() {
acc = Some(MetricFamily {
name: src_family.name.clone(),
r#type: src_family.r#type,
unit: src_family.unit.clone(),
help: src_family.help.clone(),
metrics: src_family.metrics.clone(),
});
continue;
}
let dst = acc.as_mut().unwrap();
for m in &src_family.metrics {
let dst_metric = dst.metrics.iter_mut().find(|d| d.labels == m.labels);
match dst_metric {
Some(dm) => {
let (Some(dp), Some(sp)) = (dm.points.first_mut(), m.points.first())
else {
continue;
};
combine_into(dp, sp, mode).expect("matching identity must combine");
}
None => {
dst.metrics.push(m.clone());
}
}
}
}
if let Some(family) = acc {
out.families.push(family);
}
}
out
}
}
#[derive(Clone, Debug)]
pub struct MetricFamily {
name: String,
r#type: MetricType,
unit: Option<String>,
help: Option<String>,
metrics: Vec<Metric>,
}
impl MetricFamily {
pub fn new(name: impl Into<String>, r#type: MetricType) -> Self {
Self {
name: name.into(),
r#type,
unit: None,
help: None,
metrics: Vec::new(),
}
}
pub fn with_unit(mut self, unit: impl Into<String>) -> Self {
let unit_str: String = unit.into();
if !unit_str.is_empty()
&& crate::validation::check_unit_suffix(&self.name, Some(&unit_str)).is_err()
{
self.name = format!("{}_{}", self.name, unit_str);
}
self.unit = Some(unit_str);
self
}
pub fn with_help(mut self, help: impl Into<String>) -> Self {
self.help = Some(help.into());
self
}
pub fn name(&self) -> &str {
&self.name
}
pub fn r#type(&self) -> MetricType {
self.r#type
}
pub fn unit(&self) -> Option<&str> {
self.unit.as_deref()
}
pub fn help(&self) -> Option<&str> {
self.help.as_deref()
}
pub fn metrics(&self) -> impl Iterator<Item = &Metric> {
self.metrics.iter()
}
pub fn len(&self) -> usize {
self.metrics.len()
}
pub fn is_empty(&self) -> bool {
self.metrics.is_empty()
}
pub fn metric_with_labels(&self, labels: &Labels) -> Option<&Metric> {
self.metrics.iter().find(|m| m.labels() == labels)
}
pub fn insert(&mut self, metric: Metric) {
assert!(
!self.metrics.iter().any(|m| m.labels() == metric.labels()),
"MetricFamily '{}' already contains a Metric with labels {:?} — LabelSets must be unique (OpenMetrics §4.5)",
self.name,
metric.labels(),
);
self.metrics.push(metric);
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum MetricType {
Counter,
Gauge,
Histogram,
GaugeHistogram,
Summary,
Info,
StateSet,
Unknown,
}
impl MetricType {
pub fn as_str(&self) -> &'static str {
match self {
Self::Counter => "counter",
Self::Gauge => "gauge",
Self::Histogram => "histogram",
Self::GaugeHistogram => "gaugehistogram",
Self::Summary => "summary",
Self::Info => "info",
Self::StateSet => "stateset",
Self::Unknown => "unknown",
}
}
}
#[derive(Clone, Debug)]
pub struct Metric {
labels: Labels,
points: Vec<MetricPoint>,
}
impl Metric {
pub fn new(labels: Labels, points: Vec<MetricPoint>) -> Self {
Self { labels, points }
}
pub fn single(labels: Labels, point: MetricPoint) -> Self {
Self {
labels,
points: vec![point],
}
}
pub fn labels(&self) -> &Labels {
&self.labels
}
pub fn points(&self) -> impl Iterator<Item = &MetricPoint> {
self.points.iter()
}
pub fn point(&self) -> Option<&MetricPoint> {
self.points.first()
}
}
#[derive(Clone, Debug)]
pub struct MetricPoint {
value: MetricValue,
timestamp: Option<Instant>,
}
impl MetricPoint {
pub fn new(value: MetricValue, timestamp: Instant) -> Self {
Self {
value,
timestamp: Some(timestamp),
}
}
pub fn untimed(value: MetricValue) -> Self {
Self {
value,
timestamp: None,
}
}
pub fn value(&self) -> &MetricValue {
&self.value
}
pub fn timestamp(&self) -> Option<Instant> {
self.timestamp
}
}
#[derive(Clone, Debug)]
pub enum MetricValue {
Counter(CounterValue),
Gauge(GaugeValue),
Histogram(HistogramValue),
BucketedHistogram(BucketedHistogramValue),
Info(InfoValue),
StateSet(StateSetValue),
}
#[derive(Clone, Debug)]
pub struct CounterValue {
pub cumulative: u64,
pub created: Option<Instant>,
pub exemplar: Option<Exemplar>,
}
impl CounterValue {
pub fn new(cumulative: u64) -> Self {
Self {
cumulative,
created: None,
exemplar: None,
}
}
pub fn with_created(mut self, t: Instant) -> Self {
self.created = Some(t);
self
}
pub fn with_exemplar(mut self, e: Exemplar) -> Self {
self.exemplar = Some(e);
self
}
}
#[derive(Clone, Debug)]
pub struct GaugeValue {
pub value: f64,
}
impl GaugeValue {
pub fn new(value: f64) -> Self {
Self { value }
}
}
#[derive(Clone, Debug)]
pub struct HistogramValue {
pub reservoir: Arc<HdrHistogram<u64>>,
pub count: u64,
pub cumulative_count: u64,
pub sum: f64,
pub created: Option<Instant>,
pub bucket_exemplars: Vec<Option<Exemplar>>,
}
impl HistogramValue {
pub fn from_hdr(reservoir: HdrHistogram<u64>) -> Self {
let count = reservoir.len();
let sum = hdr_sum(&reservoir);
Self {
reservoir: Arc::new(reservoir),
count,
cumulative_count: count,
sum,
created: None,
bucket_exemplars: Vec::new(),
}
}
pub fn with_cumulative_count(mut self, cumulative_count: u64) -> Self {
self.cumulative_count = cumulative_count;
self
}
pub fn with_created(mut self, t: Instant) -> Self {
self.created = Some(t);
self
}
pub fn with_bucket_exemplars(mut self, exemplars: Vec<Option<Exemplar>>) -> Self {
self.bucket_exemplars = exemplars;
self
}
pub fn project_buckets(&self, bounds: &[u64]) -> Vec<Bucket> {
let mut out = Vec::with_capacity(bounds.len() + 1);
for &le in bounds {
let cumulative = self.reservoir.count_between(0, le);
out.push(Bucket {
upper_bound: BucketBound::Finite(le),
cumulative_count: cumulative,
exemplar: None,
});
}
out.push(Bucket {
upper_bound: BucketBound::PositiveInfinity,
cumulative_count: self.count,
exemplar: None,
});
out
}
}
#[derive(Clone, Debug)]
pub struct Bucket {
pub upper_bound: BucketBound,
pub cumulative_count: u64,
pub exemplar: Option<Exemplar>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum BucketBound {
Finite(u64),
PositiveInfinity,
}
#[derive(Clone, Debug)]
pub struct BucketedHistogramValue {
pub buckets: Vec<(BucketBound, u64)>,
pub sum: Option<f64>,
pub count: u64,
pub cumulative_count: u64,
pub created: Option<Instant>,
pub bucket_exemplars: Vec<Option<Exemplar>>,
}
impl BucketedHistogramValue {
pub fn new(buckets: Vec<(BucketBound, u64)>) -> Self {
let count = buckets.iter().map(|(_, c)| *c).max().unwrap_or(0);
Self {
buckets,
sum: None,
count,
cumulative_count: count,
created: None,
bucket_exemplars: Vec::new(),
}
}
pub fn with_cumulative_count(mut self, cumulative_count: u64) -> Self {
self.cumulative_count = cumulative_count;
self
}
pub fn with_sum(mut self, sum: f64) -> Self {
self.sum = Some(sum);
self
}
pub fn with_created(mut self, t: Instant) -> Self {
self.created = Some(t);
self
}
pub fn with_bucket_exemplars(mut self, ex: Vec<Option<Exemplar>>) -> Self {
self.bucket_exemplars = ex;
self
}
}
#[derive(Clone, Debug, Default)]
pub struct InfoValue;
impl InfoValue {
pub fn new() -> Self {
Self
}
}
#[derive(Clone, Debug, Default)]
pub struct StateSetValue {
pub states: Vec<(String, bool)>,
}
impl StateSetValue {
pub fn new(states: Vec<(String, bool)>) -> Self {
Self { states }
}
pub fn with_state(mut self, name: impl Into<String>, active: bool) -> Self {
self.states.push((name.into(), active));
self
}
}
#[derive(Clone, Debug)]
pub struct Exemplar {
pub labels: Labels,
pub value: f64,
pub timestamp: Option<Instant>,
}
impl Exemplar {
pub fn new(labels: Labels, value: f64) -> Self {
Self {
labels,
value,
timestamp: None,
}
}
pub fn with_timestamp(mut self, t: Instant) -> Self {
self.timestamp = Some(t);
self
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum CombineMode {
Coalesce,
Aggregate,
}
pub fn combine_into(
dst: &mut MetricPoint,
src: &MetricPoint,
mode: CombineMode,
) -> Result<(), CombineError> {
match (&mut dst.value, &src.value) {
(MetricValue::Counter(a), MetricValue::Counter(b)) => {
a.cumulative = match mode {
CombineMode::Coalesce => {
let take_src = dst.timestamp.is_none()
|| src
.timestamp
.map(|s| Some(s) >= dst.timestamp)
.unwrap_or(false);
if take_src { b.cumulative } else { a.cumulative }
}
CombineMode::Aggregate => a.cumulative.saturating_add(b.cumulative),
};
a.created = match (a.created, b.created) {
(Some(x), Some(y)) => Some(x.min(y)),
(Some(x), None) | (None, Some(x)) => Some(x),
(None, None) => None,
};
a.exemplar = pick_more_recent_exemplar(
a.exemplar.take(),
b.exemplar.clone(),
dst.timestamp,
src.timestamp,
);
}
(MetricValue::Gauge(a), MetricValue::Gauge(b)) => {
if dst.timestamp.is_none()
|| src
.timestamp
.map(|s| Some(s) >= dst.timestamp)
.unwrap_or(false)
{
a.value = b.value;
}
}
(MetricValue::Histogram(a), MetricValue::Histogram(b)) => {
let merged = combine_hdr(&a.reservoir, &b.reservoir)?;
a.count = merged.len();
a.sum = hdr_sum(&merged);
a.reservoir = Arc::new(merged);
a.cumulative_count = match mode {
CombineMode::Coalesce => {
let take_src = dst.timestamp.is_none()
|| src
.timestamp
.map(|s| Some(s) >= dst.timestamp)
.unwrap_or(false);
if take_src {
b.cumulative_count
} else {
a.cumulative_count
}
}
CombineMode::Aggregate => a.cumulative_count.saturating_add(b.cumulative_count),
};
a.created = match (a.created, b.created) {
(Some(x), Some(y)) => Some(x.min(y)),
(Some(x), None) | (None, Some(x)) => Some(x),
(None, None) => None,
};
combine_bucket_exemplars(
&mut a.bucket_exemplars,
&b.bucket_exemplars,
dst.timestamp,
src.timestamp,
);
}
(MetricValue::BucketedHistogram(a), MetricValue::BucketedHistogram(b)) => {
if a.buckets.len() != b.buckets.len()
|| a.buckets
.iter()
.zip(b.buckets.iter())
.any(|((la, _), (lb, _))| la != lb)
{
return Err(CombineError::TypeMismatch);
}
for (i, (_, count_b)) in b.buckets.iter().enumerate() {
a.buckets[i].1 = a.buckets[i].1.saturating_add(*count_b);
}
a.count = a.count.saturating_add(b.count);
a.cumulative_count = match mode {
CombineMode::Coalesce => {
let take_src = dst.timestamp.is_none()
|| src
.timestamp
.map(|s| Some(s) >= dst.timestamp)
.unwrap_or(false);
if take_src {
b.cumulative_count
} else {
a.cumulative_count
}
}
CombineMode::Aggregate => a.cumulative_count.saturating_add(b.cumulative_count),
};
a.sum = match (a.sum, b.sum) {
(Some(sa), Some(sb)) => Some(sa + sb),
(Some(s), None) | (None, Some(s)) => Some(s),
(None, None) => None,
};
a.created = match (a.created, b.created) {
(Some(x), Some(y)) => Some(x.min(y)),
(Some(x), None) | (None, Some(x)) => Some(x),
(None, None) => None,
};
combine_bucket_exemplars(
&mut a.bucket_exemplars,
&b.bucket_exemplars,
dst.timestamp,
src.timestamp,
);
}
(MetricValue::Info(_), MetricValue::Info(_)) => {
}
(MetricValue::StateSet(a), MetricValue::StateSet(b)) => {
for (name, active) in &b.states {
if let Some(slot) = a.states.iter_mut().find(|(n, _)| n == name) {
slot.1 = *active;
} else {
a.states.push((name.clone(), *active));
}
}
}
_ => return Err(CombineError::TypeMismatch),
}
if let Some(src_ts) = src.timestamp {
dst.timestamp = Some(match dst.timestamp {
Some(d) if d >= src_ts => d,
_ => src_ts,
});
}
Ok(())
}
pub fn combine_hdr(
a: &HdrHistogram<u64>,
b: &HdrHistogram<u64>,
) -> Result<HdrHistogram<u64>, CombineError> {
let mut out = a.clone();
out.add(b).map_err(|_| CombineError::HdrAddFailed)?;
Ok(out)
}
fn pick_more_recent_exemplar(
a: Option<Exemplar>,
b: Option<Exemplar>,
a_ts: Option<Instant>,
b_ts: Option<Instant>,
) -> Option<Exemplar> {
match (a, b) {
(None, x) | (x, None) => x,
(Some(ax), Some(bx)) => {
if b_ts.map(|s| Some(s) >= a_ts).unwrap_or(false) {
Some(bx)
} else {
Some(ax)
}
}
}
}
fn combine_bucket_exemplars(
dst: &mut Vec<Option<Exemplar>>,
src: &[Option<Exemplar>],
dst_ts: Option<Instant>,
src_ts: Option<Instant>,
) {
if dst.len() < src.len() {
dst.resize(src.len(), None);
}
for (i, src_ex) in src.iter().enumerate() {
let dst_slot = dst[i].take();
dst[i] = pick_more_recent_exemplar(dst_slot, src_ex.clone(), dst_ts, src_ts);
}
}
#[derive(Debug, PartialEq, Eq)]
pub enum CombineError {
TypeMismatch,
HdrAddFailed,
}
impl std::fmt::Display for CombineError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::TypeMismatch => write!(f, "MetricPoint value type mismatch"),
Self::HdrAddFailed => write!(f, "HDR histogram add failed"),
}
}
}
impl std::error::Error for CombineError {}
pub const QUANTILES: &[f64] = &[0.5, 0.75, 0.90, 0.95, 0.98, 0.99, 0.999];
pub fn split_name_label(labels: &Labels) -> (String, Labels) {
let name = labels
.get("name")
.map(|s| s.to_string())
.expect("every metric must have a 'name' label");
let mut residual = Labels::default();
for (k, v) in labels.iter() {
if k != "name" {
residual = residual.with(k, v);
}
}
(name, residual)
}
impl MetricSet {
pub fn insert_metric(
&mut self,
family_name: impl Into<String>,
family_type: MetricType,
labels: Labels,
value: MetricValue,
timestamp: Instant,
) {
let name = family_name.into();
let point = MetricPoint::new(value, timestamp);
if let Some(fam) = self.families.iter_mut().find(|f| f.name == name) {
assert_eq!(
fam.r#type, family_type,
"family '{}' already exists as {:?}; cannot insert as {:?}",
name, fam.r#type, family_type,
);
fam.insert(Metric::single(labels, point));
} else {
let mut fam = MetricFamily::new(name, family_type);
fam.insert(Metric::single(labels, point));
self.families.push(fam);
}
}
pub fn insert_metric_with_unit(
&mut self,
family_name: impl Into<String>,
family_type: MetricType,
unit: Option<&str>,
labels: Labels,
value: MetricValue,
timestamp: Instant,
) {
let bare = family_name.into();
let template = MetricFamily::new(bare.clone(), family_type);
let template = match unit {
Some(u) => template.with_unit(u),
None => template,
};
let effective_name = template.name().to_string();
let point = MetricPoint::new(value, timestamp);
if let Some(fam) = self.families.iter_mut().find(|f| f.name == effective_name) {
assert_eq!(
fam.r#type, family_type,
"family '{}' already exists as {:?}; cannot insert as {:?}",
effective_name, fam.r#type, family_type,
);
fam.insert(Metric::single(labels, point));
} else {
let mut fam = template;
fam.insert(Metric::single(labels, point));
self.families.push(fam);
}
}
pub fn insert_counter(
&mut self,
family_name: impl Into<String>,
labels: Labels,
cumulative: u64,
timestamp: Instant,
) {
self.insert_metric(
family_name,
MetricType::Counter,
labels,
MetricValue::Counter(CounterValue::new(cumulative)),
timestamp,
);
}
pub fn insert_counter_with_unit(
&mut self,
family_name: impl Into<String>,
unit: Option<&str>,
labels: Labels,
cumulative: u64,
timestamp: Instant,
) {
self.insert_metric_with_unit(
family_name,
MetricType::Counter,
unit,
labels,
MetricValue::Counter(CounterValue::new(cumulative)),
timestamp,
);
}
pub fn insert_gauge(
&mut self,
family_name: impl Into<String>,
labels: Labels,
value: f64,
timestamp: Instant,
) {
self.insert_metric(
family_name,
MetricType::Gauge,
labels,
MetricValue::Gauge(GaugeValue::new(value)),
timestamp,
);
}
pub fn insert_gauge_with_unit(
&mut self,
family_name: impl Into<String>,
unit: Option<&str>,
labels: Labels,
value: f64,
timestamp: Instant,
) {
self.insert_metric_with_unit(
family_name,
MetricType::Gauge,
unit,
labels,
MetricValue::Gauge(GaugeValue::new(value)),
timestamp,
);
}
pub fn insert_histogram(
&mut self,
family_name: impl Into<String>,
labels: Labels,
reservoir: HdrHistogram<u64>,
timestamp: Instant,
) {
self.insert_metric(
family_name,
MetricType::Histogram,
labels,
MetricValue::Histogram(HistogramValue::from_hdr(reservoir)),
timestamp,
);
}
pub fn insert_histogram_with_unit(
&mut self,
family_name: impl Into<String>,
unit: Option<&str>,
labels: Labels,
reservoir: HdrHistogram<u64>,
timestamp: Instant,
) {
self.insert_metric_with_unit(
family_name,
MetricType::Histogram,
unit,
labels,
MetricValue::Histogram(HistogramValue::from_hdr(reservoir)),
timestamp,
);
}
pub fn insert_histogram_with_unit_cumulative(
&mut self,
family_name: impl Into<String>,
unit: Option<&str>,
labels: Labels,
reservoir: HdrHistogram<u64>,
cumulative_count: u64,
timestamp: Instant,
) {
self.insert_metric_with_unit(
family_name,
MetricType::Histogram,
unit,
labels,
MetricValue::Histogram(
HistogramValue::from_hdr(reservoir).with_cumulative_count(cumulative_count),
),
timestamp,
);
}
}
fn hdr_sum(h: &HdrHistogram<u64>) -> f64 {
h.mean() * h.len() as f64
}
pub fn counter_family(
name: impl Into<String>,
labels: Labels,
total: u64,
timestamp: Instant,
) -> MetricFamily {
let mut f = MetricFamily::new(name, MetricType::Counter);
f.insert(Metric::single(
labels,
MetricPoint::new(MetricValue::Counter(CounterValue::new(total)), timestamp),
));
f
}
pub fn gauge_family(
name: impl Into<String>,
labels: Labels,
value: f64,
timestamp: Instant,
) -> MetricFamily {
let mut f = MetricFamily::new(name, MetricType::Gauge);
f.insert(Metric::single(
labels,
MetricPoint::new(MetricValue::Gauge(GaugeValue::new(value)), timestamp),
));
f
}
pub fn histogram_family(
name: impl Into<String>,
labels: Labels,
reservoir: HdrHistogram<u64>,
timestamp: Instant,
) -> MetricFamily {
let mut f = MetricFamily::new(name, MetricType::Histogram);
f.insert(Metric::single(
labels,
MetricPoint::new(
MetricValue::Histogram(HistogramValue::from_hdr(reservoir)),
timestamp,
),
));
f
}
#[cfg(test)]
mod tests {
use super::*;
fn ts() -> Instant {
Instant::now()
}
fn empty_set() -> MetricSet {
MetricSet::new(Duration::from_secs(1))
}
#[test]
fn metric_set_is_empty_by_default() {
let m = empty_set();
assert_eq!(m.len(), 0);
assert!(m.is_empty());
assert!(m.family("anything").is_none());
assert_eq!(m.interval(), Duration::from_secs(1));
}
#[test]
fn metric_set_inserts_and_looks_up_by_name() {
let mut m = empty_set();
m.insert(counter_family(
"cycles",
Labels::of("phase", "load"),
100,
ts(),
));
m.insert(gauge_family(
"temp",
Labels::of("phase", "load"),
42.0,
ts(),
));
assert_eq!(m.len(), 2);
assert!(m.family("cycles").is_some());
assert!(m.family("temp").is_some());
assert!(m.family("missing").is_none());
}
#[test]
#[should_panic(expected = "names must be unique")]
fn metric_set_rejects_duplicate_family_names() {
let mut m = empty_set();
m.insert(counter_family("cycles", Labels::of("a", "1"), 1, ts()));
m.insert(counter_family("cycles", Labels::of("a", "2"), 2, ts()));
}
fn make_counter_set(interval: Duration, value: u64) -> MetricSet {
let mut s = MetricSet::new(interval);
s.insert(counter_family(
"cycles",
Labels::of("name", "ops"),
value,
Instant::now(),
));
s
}
fn make_histogram_set(interval: Duration, values: &[u64]) -> MetricSet {
let mut h = HdrHistogram::<u64>::new_with_bounds(1, 3_600_000_000_000, 3).unwrap();
for v in values {
h.record(*v).unwrap();
}
let mut s = MetricSet::new(interval);
s.insert(histogram_family(
"latency",
Labels::of("name", "rt"),
h,
Instant::now(),
));
s
}
fn make_gauge_set(interval: Duration, value: f64) -> MetricSet {
let mut s = MetricSet::new(interval);
s.insert(gauge_family(
"temp",
Labels::of("name", "x"),
value,
Instant::now(),
));
s
}
#[test]
fn coalesce_empty_returns_empty() {
let merged = MetricSet::coalesce(&[]);
assert!(merged.is_empty());
}
#[test]
fn coalesce_single_clones() {
let s = make_counter_set(Duration::from_secs(1), 10);
let m = MetricSet::coalesce(std::slice::from_ref(&s));
assert_eq!(m.interval(), Duration::from_secs(1));
let f = m.family("cycles").unwrap();
let c = match f.metrics().next().unwrap().point().unwrap().value() {
MetricValue::Counter(c) => c.cumulative,
_ => panic!("wrong type"),
};
assert_eq!(c, 10);
}
#[test]
fn coalesce_counters_keep_latest_cumulative_sum_intervals() {
let merged = MetricSet::coalesce(&[
make_counter_set(Duration::from_secs(1), 10),
make_counter_set(Duration::from_secs(1), 25),
make_counter_set(Duration::from_secs(1), 42),
]);
assert_eq!(merged.interval(), Duration::from_secs(3));
let cumulative = match merged
.family("cycles")
.unwrap()
.metrics()
.next()
.unwrap()
.point()
.unwrap()
.value()
{
MetricValue::Counter(c) => c.cumulative,
_ => panic!("wrong type"),
};
assert_eq!(cumulative, 42, "latest window's cumulative, not the sum");
}
#[test]
fn coalesce_histograms_merge_reservoirs() {
let merged = MetricSet::coalesce(&[
make_histogram_set(Duration::from_secs(1), &[1_000, 2_000, 3_000]),
make_histogram_set(Duration::from_secs(1), &[4_000, 5_000]),
]);
let hv = match merged
.family("latency")
.unwrap()
.metrics()
.next()
.unwrap()
.point()
.unwrap()
.value()
{
MetricValue::Histogram(h) => h.clone(),
_ => panic!("wrong type"),
};
assert_eq!(hv.count, 5);
assert!(hv.reservoir.max() >= 4_900);
}
#[test]
fn coalesce_gauges_last_write_wins() {
let merged = MetricSet::coalesce(&[
make_gauge_set(Duration::from_secs(1), 10.0),
make_gauge_set(Duration::from_secs(2), 20.0),
]);
let v = match merged
.family("temp")
.unwrap()
.metrics()
.next()
.unwrap()
.point()
.unwrap()
.value()
{
MetricValue::Gauge(g) => g.value,
_ => panic!("wrong type"),
};
assert_eq!(v, 20.0, "last written value wins, got {v}");
}
#[test]
fn coalesce_disjoint_label_sets_appended() {
let mut a = MetricSet::new(Duration::from_secs(1));
a.insert(counter_family(
"cycles",
Labels::of("phase", "load"),
100,
ts(),
));
let mut b = MetricSet::new(Duration::from_secs(1));
b.insert(counter_family(
"cycles",
Labels::of("phase", "verify"),
50,
ts(),
));
let merged = MetricSet::coalesce(&[a, b]);
let f = merged.family("cycles").unwrap();
assert_eq!(f.len(), 2);
let load = f.metric_with_labels(&Labels::of("phase", "load")).unwrap();
let verify = f
.metric_with_labels(&Labels::of("phase", "verify"))
.unwrap();
match load.point().unwrap().value() {
MetricValue::Counter(c) => assert_eq!(c.cumulative, 100),
_ => panic!(),
}
match verify.point().unwrap().value() {
MetricValue::Counter(c) => assert_eq!(c.cumulative, 50),
_ => panic!(),
}
}
#[test]
fn metric_family_records_type_and_optional_metadata() {
let f = MetricFamily::new("latency", MetricType::Histogram)
.with_unit("nanoseconds")
.with_help("End-to-end op latency");
assert_eq!(f.name(), "latency_nanoseconds");
assert_eq!(f.r#type(), MetricType::Histogram);
assert_eq!(f.unit(), Some("nanoseconds"));
assert_eq!(f.help(), Some("End-to-end op latency"));
}
#[test]
fn with_unit_preserves_name_when_suffix_already_present() {
let f = MetricFamily::new("memory_bytes", MetricType::Gauge).with_unit("bytes");
assert_eq!(f.name(), "memory_bytes");
assert_eq!(f.unit(), Some("bytes"));
}
#[test]
fn with_unit_preserves_name_when_unit_precedes_exposition_suffix() {
let f = MetricFamily::new("process_cpu_seconds_total", MetricType::Counter)
.with_unit("seconds");
assert_eq!(f.name(), "process_cpu_seconds_total");
assert_eq!(f.unit(), Some("seconds"));
}
#[test]
#[should_panic(expected = "LabelSets must be unique")]
fn metric_family_rejects_duplicate_labelsets() {
let mut f = MetricFamily::new("cycles", MetricType::Counter);
f.insert(Metric::single(
Labels::of("phase", "load"),
MetricPoint::untimed(MetricValue::Counter(CounterValue::new(1))),
));
f.insert(Metric::single(
Labels::of("phase", "load"),
MetricPoint::untimed(MetricValue::Counter(CounterValue::new(2))),
));
}
#[test]
fn metric_lookup_by_labels_matches_identity() {
let mut f = MetricFamily::new("cycles", MetricType::Counter);
f.insert(Metric::single(
Labels::of("phase", "load"),
MetricPoint::untimed(MetricValue::Counter(CounterValue::new(10))),
));
f.insert(Metric::single(
Labels::of("phase", "verify"),
MetricPoint::untimed(MetricValue::Counter(CounterValue::new(20))),
));
let load = f.metric_with_labels(&Labels::of("phase", "load")).unwrap();
match load.point().unwrap().value() {
MetricValue::Counter(c) => assert_eq!(c.cumulative, 10),
_ => panic!("wrong type"),
}
assert!(
f.metric_with_labels(&Labels::of("phase", "missing"))
.is_none()
);
}
#[test]
fn metric_type_strings_match_open_metrics_spec() {
assert_eq!(MetricType::Counter.as_str(), "counter");
assert_eq!(MetricType::Gauge.as_str(), "gauge");
assert_eq!(MetricType::Histogram.as_str(), "histogram");
assert_eq!(MetricType::GaugeHistogram.as_str(), "gaugehistogram");
assert_eq!(MetricType::Summary.as_str(), "summary");
assert_eq!(MetricType::Info.as_str(), "info");
assert_eq!(MetricType::StateSet.as_str(), "stateset");
assert_eq!(MetricType::Unknown.as_str(), "unknown");
}
#[test]
fn counter_aggregate_sums_cumulative_keeps_earlier_created() {
let t1 = Instant::now();
let t0 = t1 - Duration::from_secs(60);
let mut a = MetricPoint::new(
MetricValue::Counter(CounterValue::new(10).with_created(t1)),
t1,
);
let b = MetricPoint::new(
MetricValue::Counter(CounterValue::new(25).with_created(t0)),
t1,
);
combine_into(&mut a, &b, CombineMode::Aggregate).unwrap();
match a.value() {
MetricValue::Counter(c) => {
assert_eq!(
c.cumulative, 35,
"aggregate sums cumulative across components"
);
assert_eq!(c.created, Some(t0), "earliest created wins");
}
_ => panic!("wrong type"),
}
}
#[test]
fn counter_coalesce_keeps_latest_cumulative() {
let t1 = Instant::now();
let t2 = t1 + Duration::from_secs(1);
let mut a = MetricPoint::new(MetricValue::Counter(CounterValue::new(100)), t1);
let b = MetricPoint::new(MetricValue::Counter(CounterValue::new(113)), t2);
combine_into(&mut a, &b, CombineMode::Coalesce).unwrap();
match a.value() {
MetricValue::Counter(c) => assert_eq!(
c.cumulative, 113,
"coalesce keeps the latest cumulative (no summing)"
),
_ => panic!("wrong type"),
}
}
#[test]
fn histogram_combine_adds_reservoirs_re_derives_count_and_sum() {
let mut h1 = HdrHistogram::<u64>::new_with_bounds(1, 3_600_000_000_000, 3).unwrap();
h1.record(1_000_000).unwrap();
h1.record(2_000_000).unwrap();
let mut h2 = HdrHistogram::<u64>::new_with_bounds(1, 3_600_000_000_000, 3).unwrap();
h2.record(3_000_000).unwrap();
let mut a = MetricPoint::new(
MetricValue::Histogram(HistogramValue::from_hdr(h1)),
Instant::now(),
);
let b = MetricPoint::new(
MetricValue::Histogram(HistogramValue::from_hdr(h2)),
Instant::now(),
);
combine_into(&mut a, &b, CombineMode::Coalesce).unwrap();
match a.value() {
MetricValue::Histogram(h) => {
assert_eq!(h.count, 3);
assert!(h.sum > 0.0);
assert!(h.reservoir.max() >= 3_000_000);
}
_ => panic!("wrong type"),
}
}
#[test]
fn gauge_combine_keeps_most_recent_value() {
let t1 = Instant::now();
let t2 = t1 + Duration::from_secs(1);
let mut a = MetricPoint::new(MetricValue::Gauge(GaugeValue::new(5.0)), t1);
let b = MetricPoint::new(MetricValue::Gauge(GaugeValue::new(9.0)), t2);
combine_into(&mut a, &b, CombineMode::Coalesce).unwrap();
match a.value() {
MetricValue::Gauge(g) => assert_eq!(g.value, 9.0),
_ => panic!("wrong type"),
}
}
#[test]
fn combine_type_mismatch_is_hard_error() {
let mut a = MetricPoint::untimed(MetricValue::Counter(CounterValue::new(1)));
let b = MetricPoint::untimed(MetricValue::Gauge(GaugeValue::new(1.0)));
let err = combine_into(&mut a, &b, CombineMode::Coalesce).unwrap_err();
assert_eq!(err, CombineError::TypeMismatch);
}
#[test]
fn exemplar_most_recent_wins_on_combine() {
let t1 = Instant::now();
let t2 = t1 + Duration::from_secs(1);
let ex_old = Exemplar::new(Labels::of("trace_id", "old"), 1.0).with_timestamp(t1);
let ex_new = Exemplar::new(Labels::of("trace_id", "new"), 2.0).with_timestamp(t2);
let mut a = MetricPoint::new(
MetricValue::Counter(CounterValue::new(5).with_exemplar(ex_old)),
t1,
);
let b = MetricPoint::new(
MetricValue::Counter(CounterValue::new(5).with_exemplar(ex_new.clone())),
t2,
);
combine_into(&mut a, &b, CombineMode::Coalesce).unwrap();
match a.value() {
MetricValue::Counter(c) => {
let e = c.exemplar.as_ref().expect("exemplar should survive");
assert_eq!(e.labels.get("trace_id"), Some("new"));
}
_ => panic!("wrong type"),
}
}
#[test]
fn histogram_projects_to_open_metrics_buckets() {
let mut h = HdrHistogram::<u64>::new_with_bounds(1, 1_000_000, 3).unwrap();
for v in [10u64, 50, 100, 500, 1000, 5000, 50_000].iter() {
h.record(*v).unwrap();
}
let hv = HistogramValue::from_hdr(h);
let buckets = hv.project_buckets(&[100, 1000, 10_000]);
assert_eq!(buckets.len(), 4);
assert_eq!(buckets[0].upper_bound, BucketBound::Finite(100));
assert_eq!(buckets[3].upper_bound, BucketBound::PositiveInfinity);
assert!(buckets[0].cumulative_count >= 3);
assert_eq!(buckets[3].cumulative_count, hv.count);
for w in buckets.windows(2) {
assert!(w[0].cumulative_count <= w[1].cumulative_count);
}
}
#[test]
fn metric_point_timestamp_propagates_on_construction() {
let now = Instant::now();
let p = MetricPoint::new(MetricValue::Gauge(GaugeValue::new(1.0)), now);
assert_eq!(p.timestamp(), Some(now));
let untimed = MetricPoint::untimed(MetricValue::Gauge(GaugeValue::new(1.0)));
assert_eq!(untimed.timestamp(), None);
}
}