nmbrs_metrics/snapshot.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! OpenMetrics-aligned snapshot data model (SRD-42 §"Snapshot data
5//! model"). Mirrors the OpenMetrics specification 1:1 so external
6//! consumers (Prometheus scrape, OTel translation, third-party
7//! dashboards) get a near-trivial projection.
8//!
9//! ## Container hierarchy (spec terms verbatim)
10//!
11//! | Layer | OpenMetrics § | Type |
12//! |---|---|---|
13//! | top-level | §4.1 `MetricSet` | [`MetricSet`] |
14//! | family | §4.4 `MetricFamily` | [`MetricFamily`] |
15//! | series | §4.5 `Metric` | [`Metric`] |
16//! | point | §4.6 `MetricPoint` | [`MetricPoint`] |
17//!
18//! A time series is identified by `(MetricFamily.name, LabelSet)`
19//! per spec §4.5.1 — the same identity used by cascade-time combine
20//! (matching identity → matching reservoir / counter / gauge → combine
21//! permitted).
22//!
23//! ## Histograms
24//!
25//! Internally we keep the HDR reservoir as the source of truth on
26//! [`HistogramValue`] — that's what combines correctly across
27//! cascade folds and ephemeral merges. The OpenMetrics-shaped
28//! cumulative `Bucket` list is **derived on demand at exposition
29//! time** against the consumer-requested bucket layout. `sum` /
30//! `count` are also derivable but maintained alongside the reservoir
31//! for O(1) access.
32//!
33//! ## Naming convention
34//!
35//! Suffix rules from spec §4.4.1 / §5.x (`_total`, `_count`, `_sum`,
36//! `_bucket`, `_created`, `_info`) are **exposition-time concerns,
37//! not stored**. A counter is named `cycles` in memory; the
38//! exposition layer appends `_total` per spec.
39//!
40//! ## Initial coverage
41//!
42//! `Counter`, `Gauge`, `Histogram` are implemented. `Summary`,
43//! `Info`, `StateSet`, `Unknown`, `GaugeHistogram` are listed in
44//! [`MetricType`] but their value variants are added when a real
45//! consumer needs them.
46
47use std::sync::Arc;
48use std::time::{Duration, Instant};
49
50use hdrhistogram::Histogram as HdrHistogram;
51
52use crate::labels::Labels;
53
54// =========================================================================
55// MetricSet — top-level snapshot (OpenMetrics §4.1)
56// =========================================================================
57
58/// Top-level snapshot container per OpenMetrics §4.1. Holds zero or
59/// more [`MetricFamily`] entries, names unique within the set.
60///
61/// Snapshots are immutable once published. Producers build a new
62/// `MetricSet` per cadence-window close; consumers read the
63/// `Arc<MetricSet>` published into the cadence reporter's store.
64///
65/// `MetricSet` carries two pieces of nmbrs internal metadata that
66/// don't appear in the OpenMetrics spec but are needed by the
67/// scheduler's coalesce path:
68///
69/// - `captured_at` — wall-clock instant the snapshot was sealed.
70/// - `interval` — duration the snapshot represents (cadence window
71/// length for cadence-window snapshots; `Duration::ZERO` for
72/// instantaneous `now` reads).
73///
74/// These are intentionally not in any `MetricPoint`; consumers that
75/// project to OpenMetrics on-wire format should read the per-point
76/// timestamps and `_created` fields instead.
77/// SRD-93 M4/A6 — why a snapshot's window was sealed by a lifecycle
78/// boundary rather than closing naturally on its cadence. Typed so
79/// durable sinks act on the *reason*, never inferring lifecycle from
80/// the `partial` flag (`Quiesce` seals partials without ending scope
81/// and MUST NOT produce exit events). Variant order is severity —
82/// coalesce keeps the strongest reason across inputs.
83#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
84pub enum CloseReason {
85 /// Mid-session quiesce: windows sealed so readers see complete
86 /// data; the owning components continue. Never an exit signal.
87 Quiesce,
88 /// The owning component is being torn down (`scope_close`).
89 ScopeClose,
90 /// Session shutdown flush: every path seals for the last time.
91 Shutdown,
92}
93
94#[derive(Clone, Debug)]
95pub struct MetricSet {
96 captured_at: Instant,
97 interval: Duration,
98 /// SRD-40b §11 / SRD-42 §"Component lifecycle: scope_close flush":
99 /// a snapshot is `partial=true` when it was sealed before its
100 /// cadence window naturally closed — typically because the
101 /// component owning the contributing instruments is being torn
102 /// down between cadence pulses. Partial snapshots fold into the
103 /// next full window via the same `coalesce` rules as normal
104 /// pulse-flushed samples (Counter latest-cumulative, Gauge
105 /// last-write-wins, Histogram HDR-merge); the flag is
106 /// preserved so downstream tooling can distinguish the
107 /// scope-close contribution if needed. Coalesce result is
108 /// partial whenever **any** contributing input was partial.
109 partial: bool,
110 /// SRD-93 M4 — the lifecycle reason this window was sealed, when
111 /// sealed by a boundary (`None` for naturally-closed windows).
112 /// Rides beside `partial` rather than replacing it: `partial`
113 /// answers "is this window shorter than its cadence", `close`
114 /// answers "what lifecycle event sealed it".
115 close: Option<CloseReason>,
116 /// SRD-102 §6: the nominal cadence deadline this snapshot was
117 /// *scheduled* to fire at, when produced by the cadence scheduler.
118 /// `captured_at` is the *actual* fire instant; the pair lets
119 /// downstream tooling see schedule-vs-reality. `None` for
120 /// event-driven (non-cadence) snapshots. Logical/cadence
121 /// processing keeps using the prescribed `interval`/cadence — this
122 /// is a record of divergence, not an input to windowing.
123 scheduled_ts: Option<Instant>,
124 families: Vec<MetricFamily>,
125}
126
127impl Default for MetricSet {
128 fn default() -> Self {
129 Self {
130 captured_at: Instant::now(),
131 interval: Duration::ZERO,
132 partial: false,
133 close: None,
134 scheduled_ts: None,
135 families: Vec::new(),
136 }
137 }
138}
139
140impl MetricSet {
141 /// Construct an empty snapshot stamped with the current instant
142 /// and the given window interval.
143 pub fn new(interval: Duration) -> Self {
144 Self {
145 captured_at: Instant::now(),
146 interval,
147 partial: false,
148 close: None,
149 scheduled_ts: None,
150 families: Vec::new(),
151 }
152 }
153
154 /// Construct an empty snapshot stamped with an explicit
155 /// `captured_at` (e.g., for tests reproducing a known instant).
156 pub fn at(captured_at: Instant, interval: Duration) -> Self {
157 Self {
158 captured_at,
159 interval,
160 partial: false,
161 close: None,
162 scheduled_ts: None,
163 families: Vec::new(),
164 }
165 }
166
167 pub fn captured_at(&self) -> Instant {
168 self.captured_at
169 }
170
171 /// The *actual* instant this snapshot fired/was captured (alias of
172 /// [`captured_at`](Self::captured_at), named for the SRD-102
173 /// scheduled/actual timestamp pair).
174 pub fn actual_ts(&self) -> Instant {
175 self.captured_at
176 }
177
178 /// The nominal cadence deadline this snapshot was scheduled for, if
179 /// produced by the cadence scheduler (SRD-102 §6). `None` for
180 /// event-driven snapshots.
181 pub fn scheduled_ts(&self) -> Option<Instant> {
182 self.scheduled_ts
183 }
184
185 /// Stamp the nominal cadence deadline (the scheduler sets this after
186 /// coalescing a tick's component snapshots).
187 pub fn set_scheduled_ts(&mut self, scheduled: Instant) {
188 self.scheduled_ts = Some(scheduled);
189 }
190 pub fn interval(&self) -> Duration {
191 self.interval
192 }
193
194 /// True if this snapshot represents a partial cadence window
195 /// (sealed before the window naturally closed — typically a
196 /// `scope_close` flush at component teardown). See SRD-42
197 /// §"Component lifecycle: scope_close flush" for the
198 /// semantics.
199 pub fn is_partial(&self) -> bool {
200 self.partial
201 }
202
203 /// Mark this snapshot as a partial-window contribution.
204 /// Idempotent. Once set, `coalesce` carries the flag forward —
205 /// any output that includes this snapshot will be partial too.
206 pub fn mark_partial(&mut self) {
207 self.partial = true;
208 }
209
210 /// SRD-93 M4 — the lifecycle reason this window was sealed, if a
211 /// boundary (rather than the cadence) sealed it.
212 pub fn close_reason(&self) -> Option<CloseReason> {
213 self.close
214 }
215
216 /// Stamp the lifecycle close reason. Monotone by severity: a
217 /// stronger reason overwrites a weaker one, never the reverse
218 /// (matching the coalesce fold, so stamp order can't matter).
219 pub fn mark_close(&mut self, reason: CloseReason) {
220 self.close = Some(self.close.map_or(reason, |c| c.max(reason)));
221 }
222
223 /// Set the represented interval. Used when a coalesce path
224 /// promotes a snapshot to a coarser cadence's window.
225 pub fn set_interval(&mut self, interval: Duration) {
226 self.interval = interval;
227 }
228
229 /// Iterator over all families in this set.
230 pub fn families(&self) -> impl Iterator<Item = &MetricFamily> {
231 self.families.iter()
232 }
233
234 /// Lookup a family by name. `None` if no family has that name.
235 pub fn family(&self, name: &str) -> Option<&MetricFamily> {
236 self.families.iter().find(|f| f.name() == name)
237 }
238
239 /// Number of families in this set.
240 pub fn len(&self) -> usize {
241 self.families.len()
242 }
243 pub fn is_empty(&self) -> bool {
244 self.families.is_empty()
245 }
246
247 /// True if this set carries any HDR-reservoir distribution family
248 /// (`Histogram` / `GaugeHistogram`) — the memory-heavy part of a snapshot.
249 ///
250 /// SRD-90: cumulative counters/gauges cost ~nothing to retain across the
251 /// sub-interval history (a point is just `(timestamp, value)`), but a
252 /// distribution point carries a whole reservoir snapshot. The cadence
253 /// retention uses this to keep counter/gauge sub-windows over a long
254 /// horizon while distributions are kept only over a short one.
255 pub fn has_distributions(&self) -> bool {
256 self.families.iter().any(|f| {
257 matches!(
258 f.r#type(),
259 MetricType::Histogram | MetricType::GaugeHistogram
260 )
261 })
262 }
263
264 /// A clone of this set with the heavy distribution families
265 /// (`Histogram` / `GaugeHistogram`) dropped, keeping the cheap cumulative
266 /// `Counter`/`Gauge` families. Used by the cadence retention to compact an
267 /// aged sub-window past the distribution horizon without losing the
268 /// counter history a windowed `rate()`/`increase()` derives from
269 /// (SRD-90 §M1; counters are cheap, histograms are "a different matter").
270 pub fn without_distributions(&self) -> MetricSet {
271 MetricSet {
272 captured_at: self.captured_at,
273 interval: self.interval,
274 partial: self.partial,
275 close: self.close,
276 scheduled_ts: self.scheduled_ts,
277 families: self
278 .families
279 .iter()
280 .filter(|f| {
281 !matches!(
282 f.r#type(),
283 MetricType::Histogram | MetricType::GaugeHistogram
284 )
285 })
286 .cloned()
287 .collect(),
288 }
289 }
290
291 /// Insert a family. Panics if a family of the same name already
292 /// exists — spec §4.1 requires unique names within a `MetricSet`.
293 pub fn insert(&mut self, family: MetricFamily) {
294 assert!(
295 !self.families.iter().any(|f| f.name() == family.name()),
296 "MetricSet already contains family '{}' — names must be unique (OpenMetrics §4.1)",
297 family.name(),
298 );
299 self.families.push(family);
300 }
301
302 /// Coalesce multiple snapshots into one. Used by the scheduler
303 /// to fold smaller-cadence snapshots into a larger-cadence
304 /// window per SRD-42 §"Streaming coalesce semantics".
305 ///
306 /// Combine rules per SRD-42 §"Combine semantics — algebraic
307 /// uniformity":
308 ///
309 /// - Counter `total` sums; `created` keeps the earliest;
310 /// exemplar most-recent-wins.
311 /// - Gauge values are LAST-WRITE-WINS by timestamp — the
312 /// OpenMetrics/Prometheus/OTel gauge contract: a sample is the
313 /// last-written scalar as of window end, and any summarization
314 /// (avg/min/max over time) happens at the query point via
315 /// `*_over_time` rollups. (Until 2026-08-08 gauges coalesced as
316 /// interval-weighted means, which silently redefined what every
317 /// metricsql rollup meant against this store — and turned
318 /// set-once facts stored as gauges into fractions.)
319 /// - Histogram reservoirs add (`HdrHistogram::add`); `count`/`sum`
320 /// re-derive from the merged reservoir; bucket exemplars
321 /// most-recent-wins per index.
322 /// - Identity: `(family.name, LabelSet)` — matching identity
323 /// combines, others append. Type mismatch on matching identity
324 /// is a hard error (panic).
325 ///
326 /// `captured_at` of the result is the latest contributing
327 /// snapshot's `captured_at`; `interval` is the sum of contributing
328 /// intervals.
329 pub fn coalesce(snapshots: &[MetricSet]) -> MetricSet {
330 // Time-coalesce of consecutive windows (the cadence cascade):
331 // counter `cumulative` keeps the latest. Cross-component
332 // aggregation uses [`Self::coalesce_with_mode`] with `Aggregate`.
333 Self::coalesce_with_mode(snapshots, CombineMode::Coalesce)
334 }
335
336 /// [`coalesce`](Self::coalesce) with an explicit [`CombineMode`] —
337 /// the only difference is how counter `cumulative` folds (latest vs.
338 /// sum). See the cumulative-counter note.
339 pub fn coalesce_with_mode(snapshots: &[MetricSet], mode: CombineMode) -> MetricSet {
340 if snapshots.is_empty() {
341 return MetricSet::default();
342 }
343 if snapshots.len() == 1 {
344 return snapshots[0].clone();
345 }
346
347 let captured_at = snapshots.iter().map(|s| s.captured_at).max().unwrap();
348 let interval: Duration = snapshots.iter().map(|s| s.interval).sum();
349 // Partial flag is sticky — if any contributing input was
350 // a scope_close partial, the merged result is too. SRD-42
351 // §"Component lifecycle: scope_close flush" / SRD-40b §11.2.
352 let partial = snapshots.iter().any(|s| s.partial);
353 // A coalesced set spans multiple inputs, so no single nominal
354 // deadline applies; the scheduler stamps `scheduled_ts` on the
355 // combined tick snapshot after coalescing (SRD-102 §6).
356 let scheduled_ts = None;
357
358 // SRD-93 M4 — the strongest lifecycle reason across inputs
359 // survives the fold (severity order on the enum).
360 let close = snapshots.iter().filter_map(|s| s.close).max();
361 let mut out = MetricSet {
362 captured_at,
363 interval,
364 partial,
365 close,
366 scheduled_ts,
367 families: Vec::new(),
368 };
369
370 // Family identity is `name`. For each unique family name
371 // across inputs, fold its metrics in identity order.
372 let mut seen_family: Vec<String> = Vec::new();
373 for s in snapshots {
374 for f in &s.families {
375 if !seen_family.contains(&f.name) {
376 seen_family.push(f.name.clone());
377 }
378 }
379 }
380
381 for fname in seen_family {
382 let mut acc: Option<MetricFamily> = None;
383
384 for s in snapshots {
385 let Some(src_family) = s.families.iter().find(|f| f.name == fname) else {
386 continue;
387 };
388 if acc.is_none() {
389 // Every kind takes the uniform fold path. Gauges
390 // included: `combine_into`'s gauge arm is
391 // most-recent-wins by timestamp, which — folded in
392 // snapshot order — IS last-write-wins over the
393 // coalesced window. A series absent from later
394 // snapshots keeps its previously-written value,
395 // exactly as a Prometheus scrape would report it.
396 acc = Some(MetricFamily {
397 name: src_family.name.clone(),
398 r#type: src_family.r#type,
399 unit: src_family.unit.clone(),
400 help: src_family.help.clone(),
401 metrics: src_family.metrics.clone(),
402 });
403 continue;
404 }
405 let dst = acc.as_mut().unwrap();
406 for m in &src_family.metrics {
407 let dst_metric = dst.metrics.iter_mut().find(|d| d.labels == m.labels);
408 match dst_metric {
409 Some(dm) => {
410 let (Some(dp), Some(sp)) = (dm.points.first_mut(), m.points.first())
411 else {
412 continue;
413 };
414 combine_into(dp, sp, mode).expect("matching identity must combine");
415 }
416 None => {
417 dst.metrics.push(m.clone());
418 }
419 }
420 }
421 }
422
423 if let Some(family) = acc {
424 out.families.push(family);
425 }
426 }
427
428 out
429 }
430}
431
432// =========================================================================
433// MetricFamily (OpenMetrics §4.4)
434// =========================================================================
435
436/// One metric family — a set of [`Metric`] series sharing a name,
437/// type, optional unit, and optional help text. OpenMetrics §4.4.
438#[derive(Clone, Debug)]
439pub struct MetricFamily {
440 name: String,
441 r#type: MetricType,
442 unit: Option<String>,
443 help: Option<String>,
444 metrics: Vec<Metric>,
445}
446
447impl MetricFamily {
448 /// Construct an empty family.
449 pub fn new(name: impl Into<String>, r#type: MetricType) -> Self {
450 Self {
451 name: name.into(),
452 r#type,
453 unit: None,
454 help: None,
455 metrics: Vec::new(),
456 }
457 }
458
459 /// Attach a unit to this family.
460 ///
461 /// SRD-40b §1 / SRD-40a §4.3: when a unit is set it lands in
462 /// **two** surfaces — concatenated onto the family name as an
463 /// `_<unit>` suffix (per OpenMetrics §4.4) **and** stored in
464 /// the `unit` field for structured access. Both surfaces flow
465 /// from this single declaration so they cannot drift.
466 ///
467 /// If the family name already ends with `_<unit>` (or has
468 /// `_<unit>` immediately before a known exposition suffix
469 /// like `_total` / `_count` per [`crate::validation::check_unit_suffix`]),
470 /// the name is left unchanged — the invariant is already met.
471 /// Otherwise the suffix is appended: `overscan` + `ratio` →
472 /// `overscan_ratio`.
473 ///
474 /// Empty unit is treated as no-op for the name; the unit is
475 /// still stored as `Some("")` so callers can distinguish
476 /// "explicitly empty" from "unset" if they care.
477 pub fn with_unit(mut self, unit: impl Into<String>) -> Self {
478 let unit_str: String = unit.into();
479 if !unit_str.is_empty()
480 && crate::validation::check_unit_suffix(&self.name, Some(&unit_str)).is_err()
481 {
482 self.name = format!("{}_{}", self.name, unit_str);
483 }
484 self.unit = Some(unit_str);
485 self
486 }
487
488 pub fn with_help(mut self, help: impl Into<String>) -> Self {
489 self.help = Some(help.into());
490 self
491 }
492
493 pub fn name(&self) -> &str {
494 &self.name
495 }
496 pub fn r#type(&self) -> MetricType {
497 self.r#type
498 }
499 pub fn unit(&self) -> Option<&str> {
500 self.unit.as_deref()
501 }
502 pub fn help(&self) -> Option<&str> {
503 self.help.as_deref()
504 }
505
506 /// Iterator over the family's series.
507 pub fn metrics(&self) -> impl Iterator<Item = &Metric> {
508 self.metrics.iter()
509 }
510
511 pub fn len(&self) -> usize {
512 self.metrics.len()
513 }
514 pub fn is_empty(&self) -> bool {
515 self.metrics.is_empty()
516 }
517
518 /// Look up the series with the given LabelSet, if any. Identity
519 /// per spec §4.5.1: `(family.name, label_set)`.
520 pub fn metric_with_labels(&self, labels: &Labels) -> Option<&Metric> {
521 self.metrics.iter().find(|m| m.labels() == labels)
522 }
523
524 /// Insert a series. Panics if a series with the same LabelSet
525 /// already exists — spec §4.5 requires unique LabelSets within
526 /// a family.
527 pub fn insert(&mut self, metric: Metric) {
528 assert!(
529 !self.metrics.iter().any(|m| m.labels() == metric.labels()),
530 "MetricFamily '{}' already contains a Metric with labels {:?} — LabelSets must be unique (OpenMetrics §4.5)",
531 self.name,
532 metric.labels(),
533 );
534 self.metrics.push(metric);
535 }
536}
537
538/// OpenMetrics metric type per spec §4.4. Stored on every
539/// [`MetricFamily`] and projected verbatim to exposition.
540#[derive(Clone, Copy, Debug, PartialEq, Eq)]
541pub enum MetricType {
542 Counter,
543 Gauge,
544 Histogram,
545 GaugeHistogram,
546 Summary,
547 Info,
548 StateSet,
549 Unknown,
550}
551
552impl MetricType {
553 /// The exposition-format token for this type (spec §4.4).
554 pub fn as_str(&self) -> &'static str {
555 match self {
556 Self::Counter => "counter",
557 Self::Gauge => "gauge",
558 Self::Histogram => "histogram",
559 Self::GaugeHistogram => "gaugehistogram",
560 Self::Summary => "summary",
561 Self::Info => "info",
562 Self::StateSet => "stateset",
563 Self::Unknown => "unknown",
564 }
565 }
566}
567
568// =========================================================================
569// Metric (OpenMetrics §4.5)
570// =========================================================================
571
572/// One labeled time series within a [`MetricFamily`]. Identity is
573/// `(family.name, labels)` per spec §4.5.1.
574///
575/// Carries an ordered list of [`MetricPoint`]s — typically one in a
576/// snapshot, but the spec permits multiple when several observations
577/// belong to the same series.
578#[derive(Clone, Debug)]
579pub struct Metric {
580 labels: Labels,
581 points: Vec<MetricPoint>,
582}
583
584impl Metric {
585 pub fn new(labels: Labels, points: Vec<MetricPoint>) -> Self {
586 Self { labels, points }
587 }
588
589 /// Convenience for the common single-point case.
590 pub fn single(labels: Labels, point: MetricPoint) -> Self {
591 Self {
592 labels,
593 points: vec![point],
594 }
595 }
596
597 pub fn labels(&self) -> &Labels {
598 &self.labels
599 }
600
601 /// Iterator over the series' points.
602 pub fn points(&self) -> impl Iterator<Item = &MetricPoint> {
603 self.points.iter()
604 }
605
606 /// First point — convenience for the typical single-point case.
607 /// Returns `None` if the series is empty (which violates the
608 /// spec but is permitted at construction time so consumers
609 /// don't have to handle Result).
610 pub fn point(&self) -> Option<&MetricPoint> {
611 self.points.first()
612 }
613}
614
615// =========================================================================
616// MetricPoint (OpenMetrics §4.6)
617// =========================================================================
618
619/// One observation for a [`Metric`], plus an optional timestamp.
620/// OpenMetrics §4.6.
621///
622/// `timestamp` is **always populated** in nmbrs snapshots (the
623/// cadence-window-close instant for cadence-window points, the live
624/// read instant for `now` points, the merge instant for ephemeral
625/// `increase_over` / `session_lifetime` points) — even though the
626/// spec marks it optional.
627#[derive(Clone, Debug)]
628pub struct MetricPoint {
629 value: MetricValue,
630 timestamp: Option<Instant>,
631}
632
633impl MetricPoint {
634 pub fn new(value: MetricValue, timestamp: Instant) -> Self {
635 Self {
636 value,
637 timestamp: Some(timestamp),
638 }
639 }
640
641 /// Construct without a timestamp — for cases (typically tests)
642 /// where the timestamp is unknown or irrelevant. Production
643 /// snapshots always set one.
644 pub fn untimed(value: MetricValue) -> Self {
645 Self {
646 value,
647 timestamp: None,
648 }
649 }
650
651 pub fn value(&self) -> &MetricValue {
652 &self.value
653 }
654 pub fn timestamp(&self) -> Option<Instant> {
655 self.timestamp
656 }
657}
658
659/// The typed value carried by a [`MetricPoint`]. Variants mirror
660/// OpenMetrics §5.x. `Histogram` is the HDR-reservoir-backed
661/// summary shape (percentile-based — semantically a Summary
662/// per OpenMetrics §5.5). [`MetricValue::BucketedHistogram`]
663/// is the explicit-`le`-bucket Histogram shape (§5.3) and
664/// [`MetricType::GaugeHistogram`] (§5.4) reuses it under a
665/// different family-type tag.
666#[derive(Clone, Debug)]
667pub enum MetricValue {
668 /// OpenMetrics §5.1: monotonic counter.
669 Counter(CounterValue),
670 /// OpenMetrics §5.2: instantaneous gauge.
671 Gauge(GaugeValue),
672 /// OpenMetrics §5.5: φ-quantile summary backed by an
673 /// HDR reservoir. (Note: the variant name is historical;
674 /// per OpenMetrics taxonomy this shape is a Summary.)
675 Histogram(HistogramValue),
676 /// OpenMetrics §5.3 (Histogram) / §5.4 (GaugeHistogram):
677 /// explicit cumulative `le`-keyed buckets. The owning
678 /// [`MetricFamily::r#type()`] tag (`Histogram` vs
679 /// `GaugeHistogram`) decides whether bucket counts are
680 /// constrained monotonic — see `nmbrs-metrics::validation`.
681 BucketedHistogram(BucketedHistogramValue),
682 /// OpenMetrics §5.6: descriptive metadata. The label
683 /// set carries the data; the value is conceptually 1.
684 Info(InfoValue),
685 /// OpenMetrics §5.7: named-state indicator set. Each
686 /// state renders as its own Metric in exposition.
687 StateSet(StateSetValue),
688}
689
690// ---- CounterValue (§5.1.1) ---------------------------------------------
691
692/// OpenMetrics §5.1.1 counter point. Sample name carries the
693/// spec-required `_total` suffix on exposition (not stored here).
694#[derive(Clone, Debug)]
695pub struct CounterValue {
696 /// The counter's **cumulative** running total at this point's
697 /// timestamp — the single canonical, Prometheus/VM-schematic
698 /// monotonic value. Per-interval deltas are DERIVED by differencing
699 /// consecutive samples (the metricsql engine, the sqlite `_rate`
700 /// suffix, windowed-throughput readers), never stored. Time-coalesce
701 /// keeps the latest (monotonic ⇒ window-end); cross-component
702 /// aggregate sums it. See `docs/SRD/notes/cumulative_counter_model.md`.
703 pub cumulative: u64,
704 /// Series start time per spec §5.1; lets external consumers
705 /// detect counter resets. Optional in spec.
706 pub created: Option<Instant>,
707 /// Optional exemplar per spec §4.6.1. At most one per
708 /// `CounterValue`.
709 pub exemplar: Option<Exemplar>,
710}
711
712impl CounterValue {
713 /// A counter point at the given cumulative running total.
714 pub fn new(cumulative: u64) -> Self {
715 Self {
716 cumulative,
717 created: None,
718 exemplar: None,
719 }
720 }
721
722 pub fn with_created(mut self, t: Instant) -> Self {
723 self.created = Some(t);
724 self
725 }
726
727 pub fn with_exemplar(mut self, e: Exemplar) -> Self {
728 self.exemplar = Some(e);
729 self
730 }
731}
732
733// ---- GaugeValue (§5.2.1) -----------------------------------------------
734
735/// OpenMetrics §5.2.1 gauge point.
736#[derive(Clone, Debug)]
737pub struct GaugeValue {
738 pub value: f64,
739}
740
741impl GaugeValue {
742 pub fn new(value: f64) -> Self {
743 Self { value }
744 }
745}
746
747// ---- HistogramValue (§5.3.1) -------------------------------------------
748
749/// OpenMetrics §5.3.1 histogram point. Internally carries the HDR
750/// reservoir as the source of truth — OpenMetrics-shaped cumulative
751/// buckets are derived on demand at exposition time.
752///
753/// `sum` / `count` are derivable from the reservoir but cached for
754/// O(1) access. `created` is the series start time (component start)
755/// per spec.
756///
757/// Per-bucket exemplars are NOT stored on the reservoir — they live
758/// in [`HistogramValue::bucket_exemplars`] keyed by bucket index of
759/// the consumer's eventual bucket layout. (Combine semantics:
760/// most-recent-wins by `MetricPoint.timestamp`, per SRD-42
761/// §"Exemplars → Combine semantics".)
762#[derive(Clone, Debug)]
763pub struct HistogramValue {
764 /// HDR reservoir — the lossless source of truth for combining
765 /// across cascade folds and ephemeral merges.
766 pub reservoir: Arc<HdrHistogram<u64>>,
767 /// Cached observation count. Equal to `reservoir.len()` — the
768 /// per-window count in a delta snapshot.
769 pub count: u64,
770 /// The **cumulative** total observation count at this point's
771 /// timestamp (monotonic, Prometheus/VM-schematic), carried
772 /// alongside the per-window `count` exactly as `CounterValue`
773 /// carries `cumulative` alongside `total`. The queryapi exposes
774 /// this so MetricsQL `rate()`/`increase`/`*_over_time` over a
775 /// histogram's count are PromQL-correct. Time-coalesce keeps the
776 /// latest; cross-component aggregate sums it. (The percentile
777 /// reservoir stays windowed — see the cumulative-counter note.)
778 pub cumulative_count: u64,
779 /// Cached observation sum (nanoseconds for latency timers).
780 pub sum: f64,
781 /// Series start time per spec §5.3; optional.
782 pub created: Option<Instant>,
783 /// Sampled exemplars, one per bucket of an eventual exposition
784 /// layout. Sparse: an empty slot means "no exemplar for that
785 /// bucket". See SRD-42 §"Exemplars".
786 pub bucket_exemplars: Vec<Option<Exemplar>>,
787}
788
789impl HistogramValue {
790 pub fn from_hdr(reservoir: HdrHistogram<u64>) -> Self {
791 let count = reservoir.len();
792 let sum = hdr_sum(&reservoir);
793 Self {
794 reservoir: Arc::new(reservoir),
795 count,
796 // Defaults to the per-window count (the value is both for a
797 // one-shot snapshot); the cadence capture path overrides it
798 // with the accumulated absolute via `with_cumulative_count`.
799 cumulative_count: count,
800 sum,
801 created: None,
802 bucket_exemplars: Vec::new(),
803 }
804 }
805
806 /// Override the cumulative observation count (the cadence capture
807 /// path sets the accumulated absolute total here).
808 pub fn with_cumulative_count(mut self, cumulative_count: u64) -> Self {
809 self.cumulative_count = cumulative_count;
810 self
811 }
812
813 pub fn with_created(mut self, t: Instant) -> Self {
814 self.created = Some(t);
815 self
816 }
817
818 pub fn with_bucket_exemplars(mut self, exemplars: Vec<Option<Exemplar>>) -> Self {
819 self.bucket_exemplars = exemplars;
820 self
821 }
822
823 /// Project to OpenMetrics-shaped cumulative buckets at the given
824 /// upper bounds. Per spec §5.3, the final bucket MUST have
825 /// `upper_bound = +Inf`; this helper appends it automatically
826 /// if not present.
827 ///
828 /// `bounds` should be sorted ascending. Returns `(upper_bound,
829 /// cumulative_count)` pairs.
830 pub fn project_buckets(&self, bounds: &[u64]) -> Vec<Bucket> {
831 let mut out = Vec::with_capacity(bounds.len() + 1);
832 for &le in bounds {
833 let cumulative = self.reservoir.count_between(0, le);
834 out.push(Bucket {
835 upper_bound: BucketBound::Finite(le),
836 cumulative_count: cumulative,
837 exemplar: None,
838 });
839 }
840 out.push(Bucket {
841 upper_bound: BucketBound::PositiveInfinity,
842 cumulative_count: self.count,
843 exemplar: None,
844 });
845 out
846 }
847}
848
849/// One cumulative bucket projected from a [`HistogramValue`] for
850/// exposition. Per spec §5.3, the final bucket MUST have
851/// `upper_bound = +Inf`.
852#[derive(Clone, Debug)]
853pub struct Bucket {
854 pub upper_bound: BucketBound,
855 pub cumulative_count: u64,
856 pub exemplar: Option<Exemplar>,
857}
858
859/// Bucket upper bound — `+Inf` is the spec-required final bucket.
860#[derive(Clone, Copy, Debug, PartialEq, Eq)]
861pub enum BucketBound {
862 Finite(u64),
863 PositiveInfinity,
864}
865
866// ---- BucketedHistogramValue (§5.3 / §5.4) -------------------------------
867
868/// OpenMetrics §5.3 (Histogram) and §5.4 (GaugeHistogram)
869/// explicit-bucket value: cumulative observation counts at
870/// producer-chosen `le` boundaries, plus optional
871/// `_sum` / `_count` / `_created` siblings.
872///
873/// Histogram (§5.3) bucket counts are required to be
874/// monotonically non-decreasing; GaugeHistogram (§5.4)
875/// permits decreases. The owning [`MetricFamily::r#type()`]
876/// tag distinguishes which constraint applies; the validation
877/// helper in `crate::validation::check_bucket_monotonicity`
878/// enforces it for `Histogram` families.
879///
880/// Difference from [`HistogramValue`]: that variant is the
881/// HDR-reservoir-backed Summary (§5.5), exposing percentile
882/// columns. This variant is the bucket-keyed form — what
883/// OpenMetrics calls a Histogram strictly.
884#[derive(Clone, Debug)]
885pub struct BucketedHistogramValue {
886 /// Cumulative `(le, count)` pairs in ascending `le`
887 /// order. The final pair SHOULD have `le = +Inf` and
888 /// `count == self.count` per spec §5.3; producers that
889 /// omit `+Inf` are tolerated but exposition is
890 /// permitted to synthesize one.
891 pub buckets: Vec<(BucketBound, u64)>,
892 /// Optional sum of all observations.
893 pub sum: Option<f64>,
894 /// Total observation count (== last bucket count) — the
895 /// per-window total in a delta snapshot.
896 pub count: u64,
897 /// The **cumulative** total observation count (monotonic,
898 /// Prometheus/VM-schematic) carried alongside the per-window
899 /// `count`, mirroring [`HistogramValue::cumulative_count`] and
900 /// `CounterValue::cumulative`. (The per-bucket `le` counts in
901 /// `buckets` are the separate OpenMetrics bucket-cumulative
902 /// thing and are unchanged.)
903 pub cumulative_count: u64,
904 /// Series start time per spec §5.3 / §5.4.
905 pub created: Option<Instant>,
906 /// Optional per-bucket exemplars, parallel to `buckets`.
907 /// Sparse: missing slot ⇒ no exemplar for that bucket.
908 pub bucket_exemplars: Vec<Option<Exemplar>>,
909}
910
911impl BucketedHistogramValue {
912 /// Construct from cumulative `(le, count)` pairs. `count`
913 /// is the final bucket's count when present, else the
914 /// max across the supplied pairs.
915 pub fn new(buckets: Vec<(BucketBound, u64)>) -> Self {
916 let count = buckets.iter().map(|(_, c)| *c).max().unwrap_or(0);
917 Self {
918 buckets,
919 sum: None,
920 count,
921 cumulative_count: count,
922 created: None,
923 bucket_exemplars: Vec::new(),
924 }
925 }
926
927 /// Override the cumulative total observation count.
928 pub fn with_cumulative_count(mut self, cumulative_count: u64) -> Self {
929 self.cumulative_count = cumulative_count;
930 self
931 }
932
933 pub fn with_sum(mut self, sum: f64) -> Self {
934 self.sum = Some(sum);
935 self
936 }
937
938 pub fn with_created(mut self, t: Instant) -> Self {
939 self.created = Some(t);
940 self
941 }
942
943 pub fn with_bucket_exemplars(mut self, ex: Vec<Option<Exemplar>>) -> Self {
944 self.bucket_exemplars = ex;
945 self
946 }
947}
948
949// ---- InfoValue (§5.6) ---------------------------------------------------
950
951/// OpenMetrics §5.6: descriptive info metric. The label set
952/// of the owning [`Metric`] carries the data; the value is
953/// conceptually always `1`. This struct has no payload —
954/// the variant tag itself is the marker.
955#[derive(Clone, Debug, Default)]
956pub struct InfoValue;
957
958impl InfoValue {
959 pub fn new() -> Self {
960 Self
961 }
962}
963
964// ---- StateSetValue (§5.7) -----------------------------------------------
965
966/// OpenMetrics §5.7: named-state indicator set. Each
967/// `(state, active)` pair renders as a separate Metric in
968/// exposition with the state name carried as a label.
969///
970/// Spec §5.7 requires that exactly one state in a StateSet
971/// MAY be true at a time when the StateSet encodes an enum
972/// (vs a free bitset). This struct doesn't enforce the
973/// "exactly one" constraint — that's a producer-side
974/// convention.
975#[derive(Clone, Debug, Default)]
976pub struct StateSetValue {
977 /// Each entry is `(state_name, active_bool)`. State
978 /// names are arbitrary strings (no ABNF restriction
979 /// beyond label-value rules).
980 pub states: Vec<(String, bool)>,
981}
982
983impl StateSetValue {
984 pub fn new(states: Vec<(String, bool)>) -> Self {
985 Self { states }
986 }
987
988 pub fn with_state(mut self, name: impl Into<String>, active: bool) -> Self {
989 self.states.push((name.into(), active));
990 self
991 }
992}
993
994// =========================================================================
995// Exemplar (OpenMetrics §4.6.1, §4.7)
996// =========================================================================
997
998/// OpenMetrics §4.6.1 exemplar: a labeled link from a metric
999/// observation to an external context (typically a trace/span ID,
1000/// workload cycle number, or sample identifier).
1001///
1002/// Per spec §4.7 the serialized LabelSet MUST be ≤ 128 UTF-8
1003/// characters. Validation lives at exposition (a stored exemplar
1004/// that exceeds the limit is dropped from the wire; the recording
1005/// path is allowed to be permissive).
1006#[derive(Clone, Debug)]
1007pub struct Exemplar {
1008 pub labels: Labels,
1009 pub value: f64,
1010 pub timestamp: Option<Instant>,
1011}
1012
1013impl Exemplar {
1014 pub fn new(labels: Labels, value: f64) -> Self {
1015 Self {
1016 labels,
1017 value,
1018 timestamp: None,
1019 }
1020 }
1021
1022 pub fn with_timestamp(mut self, t: Instant) -> Self {
1023 self.timestamp = Some(t);
1024 self
1025 }
1026}
1027
1028// =========================================================================
1029// Combine — algebraic uniformity (SRD-42 §"Combine semantics")
1030// =========================================================================
1031
1032/// How two matching metric points combine. The counter `cumulative`
1033/// field is the only thing that differs (see the cumulative-counter
1034/// note): time-coalesce keeps the latest cumulative (monotonic), a
1035/// cross-component aggregate sums it. Everything else — counter `total`
1036/// (delta), gauges, histograms — is identical in both modes.
1037#[derive(Clone, Copy, Debug, PartialEq, Eq)]
1038pub enum CombineMode {
1039 /// Consecutive windows of one series (the cadence cascade): counter
1040 /// `cumulative` = most-recent by timestamp.
1041 Coalesce,
1042 /// Same family+labels across components (cross-component totals):
1043 /// counter `cumulative` = sum.
1044 Aggregate,
1045}
1046
1047/// In-place combine of `other` into `self`. Both must have the same
1048/// identity `(family.name, labels)` and the same value variant —
1049/// otherwise this is a hard error (panic) per the SRD's "matching
1050/// identity → matching combine" rule.
1051///
1052/// Combine rules per SRD-42 §"Combine semantics":
1053/// - Counter `total` (delta) sums in both modes; `cumulative` folds per
1054/// `mode` (latest on `Coalesce`, sum on `Aggregate`); `created` keeps
1055/// the earliest; exemplar most-recent-wins by `MetricPoint.timestamp`.
1056/// - Gauge values are last-write-wins (newer timestamp wins) — the
1057/// OpenMetrics/Prometheus gauge contract; summarization belongs at
1058/// the query point (`*_over_time`).
1059/// - Histogram reservoirs add via `HdrHistogram::add`; sum/count
1060/// re-derive; bucket exemplars most-recent-wins per index.
1061pub fn combine_into(
1062 dst: &mut MetricPoint,
1063 src: &MetricPoint,
1064 mode: CombineMode,
1065) -> Result<(), CombineError> {
1066 match (&mut dst.value, &src.value) {
1067 (MetricValue::Counter(a), MetricValue::Counter(b)) => {
1068 // A counter is its monotonic running total. Time-coalescing
1069 // consecutive windows of one series keeps the LATEST (the
1070 // window-end cumulative — never a sum); aggregating the same
1071 // series across components sums the running totals. Per-window
1072 // deltas are derived downstream by differencing samples.
1073 a.cumulative = match mode {
1074 CombineMode::Coalesce => {
1075 let take_src = dst.timestamp.is_none()
1076 || src
1077 .timestamp
1078 .map(|s| Some(s) >= dst.timestamp)
1079 .unwrap_or(false);
1080 if take_src { b.cumulative } else { a.cumulative }
1081 }
1082 CombineMode::Aggregate => a.cumulative.saturating_add(b.cumulative),
1083 };
1084 a.created = match (a.created, b.created) {
1085 (Some(x), Some(y)) => Some(x.min(y)),
1086 (Some(x), None) | (None, Some(x)) => Some(x),
1087 (None, None) => None,
1088 };
1089 a.exemplar = pick_more_recent_exemplar(
1090 a.exemplar.take(),
1091 b.exemplar.clone(),
1092 dst.timestamp,
1093 src.timestamp,
1094 );
1095 }
1096 (MetricValue::Gauge(a), MetricValue::Gauge(b)) => {
1097 // Most-recent-wins by timestamp; same-or-missing → src.
1098 if dst.timestamp.is_none()
1099 || src
1100 .timestamp
1101 .map(|s| Some(s) >= dst.timestamp)
1102 .unwrap_or(false)
1103 {
1104 a.value = b.value;
1105 }
1106 }
1107 (MetricValue::Histogram(a), MetricValue::Histogram(b)) => {
1108 let merged = combine_hdr(&a.reservoir, &b.reservoir)?;
1109 a.count = merged.len();
1110 a.sum = hdr_sum(&merged);
1111 a.reservoir = Arc::new(merged);
1112 // `cumulative_count` (the monotonic lifetime count) follows the
1113 // counter rule: keep the latest when coalescing consecutive
1114 // windows of one series; sum when aggregating across components.
1115 a.cumulative_count = match mode {
1116 CombineMode::Coalesce => {
1117 let take_src = dst.timestamp.is_none()
1118 || src
1119 .timestamp
1120 .map(|s| Some(s) >= dst.timestamp)
1121 .unwrap_or(false);
1122 if take_src {
1123 b.cumulative_count
1124 } else {
1125 a.cumulative_count
1126 }
1127 }
1128 CombineMode::Aggregate => a.cumulative_count.saturating_add(b.cumulative_count),
1129 };
1130 a.created = match (a.created, b.created) {
1131 (Some(x), Some(y)) => Some(x.min(y)),
1132 (Some(x), None) | (None, Some(x)) => Some(x),
1133 (None, None) => None,
1134 };
1135 combine_bucket_exemplars(
1136 &mut a.bucket_exemplars,
1137 &b.bucket_exemplars,
1138 dst.timestamp,
1139 src.timestamp,
1140 );
1141 }
1142 (MetricValue::BucketedHistogram(a), MetricValue::BucketedHistogram(b)) => {
1143 // Bucket layouts must be compatible. If they
1144 // share boundaries we sum element-wise; mismatched
1145 // layouts are a type error (the producer is
1146 // responsible for emitting consistent buckets per
1147 // series).
1148 if a.buckets.len() != b.buckets.len()
1149 || a.buckets
1150 .iter()
1151 .zip(b.buckets.iter())
1152 .any(|((la, _), (lb, _))| la != lb)
1153 {
1154 return Err(CombineError::TypeMismatch);
1155 }
1156 for (i, (_, count_b)) in b.buckets.iter().enumerate() {
1157 a.buckets[i].1 = a.buckets[i].1.saturating_add(*count_b);
1158 }
1159 a.count = a.count.saturating_add(b.count);
1160 // `cumulative_count` (the monotonic lifetime total) follows the
1161 // counter rule, like the HDR histogram above: latest when
1162 // coalescing one series' consecutive windows, summed when
1163 // aggregating across components. (The per-bucket `le` counts are
1164 // the separate OpenMetrics bucket-cumulative thing, summed above.)
1165 a.cumulative_count = match mode {
1166 CombineMode::Coalesce => {
1167 let take_src = dst.timestamp.is_none()
1168 || src
1169 .timestamp
1170 .map(|s| Some(s) >= dst.timestamp)
1171 .unwrap_or(false);
1172 if take_src {
1173 b.cumulative_count
1174 } else {
1175 a.cumulative_count
1176 }
1177 }
1178 CombineMode::Aggregate => a.cumulative_count.saturating_add(b.cumulative_count),
1179 };
1180 a.sum = match (a.sum, b.sum) {
1181 (Some(sa), Some(sb)) => Some(sa + sb),
1182 (Some(s), None) | (None, Some(s)) => Some(s),
1183 (None, None) => None,
1184 };
1185 a.created = match (a.created, b.created) {
1186 (Some(x), Some(y)) => Some(x.min(y)),
1187 (Some(x), None) | (None, Some(x)) => Some(x),
1188 (None, None) => None,
1189 };
1190 combine_bucket_exemplars(
1191 &mut a.bucket_exemplars,
1192 &b.bucket_exemplars,
1193 dst.timestamp,
1194 src.timestamp,
1195 );
1196 }
1197 (MetricValue::Info(_), MetricValue::Info(_)) => {
1198 // Info is always-1; combining is a no-op apart
1199 // from the timestamp update at the bottom.
1200 }
1201 (MetricValue::StateSet(a), MetricValue::StateSet(b)) => {
1202 // Most-recent-wins on state values: walk `b`'s
1203 // states and overwrite `a`'s entries by name,
1204 // appending unknown states.
1205 for (name, active) in &b.states {
1206 if let Some(slot) = a.states.iter_mut().find(|(n, _)| n == name) {
1207 slot.1 = *active;
1208 } else {
1209 a.states.push((name.clone(), *active));
1210 }
1211 }
1212 }
1213 _ => return Err(CombineError::TypeMismatch),
1214 }
1215 if let Some(src_ts) = src.timestamp {
1216 dst.timestamp = Some(match dst.timestamp {
1217 Some(d) if d >= src_ts => d,
1218 _ => src_ts,
1219 });
1220 }
1221 Ok(())
1222}
1223
1224/// Combine two `HistogramValue` reservoirs into a new owned HDR
1225/// histogram. Used both by `combine_into` and by ephemeral
1226/// `increase_over` / `session_lifetime` queries that fold many
1227/// reservoirs without mutating any of them.
1228pub fn combine_hdr(
1229 a: &HdrHistogram<u64>,
1230 b: &HdrHistogram<u64>,
1231) -> Result<HdrHistogram<u64>, CombineError> {
1232 let mut out = a.clone();
1233 out.add(b).map_err(|_| CombineError::HdrAddFailed)?;
1234 Ok(out)
1235}
1236
1237fn pick_more_recent_exemplar(
1238 a: Option<Exemplar>,
1239 b: Option<Exemplar>,
1240 a_ts: Option<Instant>,
1241 b_ts: Option<Instant>,
1242) -> Option<Exemplar> {
1243 match (a, b) {
1244 (None, x) | (x, None) => x,
1245 (Some(ax), Some(bx)) => {
1246 if b_ts.map(|s| Some(s) >= a_ts).unwrap_or(false) {
1247 Some(bx)
1248 } else {
1249 Some(ax)
1250 }
1251 }
1252 }
1253}
1254
1255fn combine_bucket_exemplars(
1256 dst: &mut Vec<Option<Exemplar>>,
1257 src: &[Option<Exemplar>],
1258 dst_ts: Option<Instant>,
1259 src_ts: Option<Instant>,
1260) {
1261 if dst.len() < src.len() {
1262 dst.resize(src.len(), None);
1263 }
1264 for (i, src_ex) in src.iter().enumerate() {
1265 let dst_slot = dst[i].take();
1266 dst[i] = pick_more_recent_exemplar(dst_slot, src_ex.clone(), dst_ts, src_ts);
1267 }
1268}
1269
1270/// Errors from [`combine_into`] / [`combine_hdr`].
1271#[derive(Debug, PartialEq, Eq)]
1272pub enum CombineError {
1273 /// The two `MetricPoint`s have different value variants
1274 /// (e.g., Counter vs Gauge). Indicates a programming error —
1275 /// only matching identity should ever combine.
1276 TypeMismatch,
1277 /// HDR `add` failed (typically because reservoirs have
1278 /// incompatible bounds).
1279 HdrAddFailed,
1280}
1281
1282impl std::fmt::Display for CombineError {
1283 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1284 match self {
1285 Self::TypeMismatch => write!(f, "MetricPoint value type mismatch"),
1286 Self::HdrAddFailed => write!(f, "HDR histogram add failed"),
1287 }
1288 }
1289}
1290
1291impl std::error::Error for CombineError {}
1292
1293// =========================================================================
1294// Quantiles
1295// =========================================================================
1296
1297/// Standard quantiles reported for histogram samples by reporters
1298/// that emit a fixed quantile set (Prometheus summary-shape, CSV
1299/// percentile columns, etc.).
1300pub const QUANTILES: &[f64] = &[0.5, 0.75, 0.90, 0.95, 0.98, 0.99, 0.999];
1301
1302// =========================================================================
1303// Migration helpers — extract `name` label, build single-point families
1304// =========================================================================
1305
1306/// Split a [`Labels`] value into `(family_name, residual_labels)` by
1307/// extracting the `name` label. Producers that historically embedded
1308/// the metric name in `Labels` (the pre-snapshot pattern) use this
1309/// to feed [`MetricSet::insert_metric`].
1310///
1311/// Panics if `name` is missing — every metric MUST have a family
1312/// name per OpenMetrics §4.4.
1313pub fn split_name_label(labels: &Labels) -> (String, Labels) {
1314 let name = labels
1315 .get("name")
1316 .map(|s| s.to_string())
1317 .expect("every metric must have a 'name' label");
1318 let mut residual = Labels::default();
1319 for (k, v) in labels.iter() {
1320 if k != "name" {
1321 residual = residual.with(k, v);
1322 }
1323 }
1324 (name, residual)
1325}
1326
1327impl MetricSet {
1328 /// Insert a single observation, looking up or creating its
1329 /// `MetricFamily` by name and appending one [`Metric`]/[`MetricPoint`]
1330 /// pair. Convenience for producers that build snapshots one
1331 /// observation at a time.
1332 ///
1333 /// Panics if the family already exists with a different
1334 /// [`MetricType`] (per identity rules) or if the LabelSet
1335 /// already exists within the family (spec §4.5).
1336 pub fn insert_metric(
1337 &mut self,
1338 family_name: impl Into<String>,
1339 family_type: MetricType,
1340 labels: Labels,
1341 value: MetricValue,
1342 timestamp: Instant,
1343 ) {
1344 let name = family_name.into();
1345 let point = MetricPoint::new(value, timestamp);
1346 if let Some(fam) = self.families.iter_mut().find(|f| f.name == name) {
1347 assert_eq!(
1348 fam.r#type, family_type,
1349 "family '{}' already exists as {:?}; cannot insert as {:?}",
1350 name, fam.r#type, family_type,
1351 );
1352 fam.insert(Metric::single(labels, point));
1353 } else {
1354 let mut fam = MetricFamily::new(name, family_type);
1355 fam.insert(Metric::single(labels, point));
1356 self.families.push(fam);
1357 }
1358 }
1359
1360 /// Insert one observation, looking up or creating the
1361 /// [`MetricFamily`] by name with the given unit. When `unit`
1362 /// is set, the family name picks up the OpenMetrics
1363 /// `_<unit>` suffix at creation (per
1364 /// [`MetricFamily::with_unit`]) and the unit is also stored
1365 /// in the `unit` field. Subsequent inserts for the same bare
1366 /// `(family_name, unit)` pair find the same family by the
1367 /// suffixed name.
1368 ///
1369 /// Equivalent to [`insert_metric`] when `unit` is `None`.
1370 pub fn insert_metric_with_unit(
1371 &mut self,
1372 family_name: impl Into<String>,
1373 family_type: MetricType,
1374 unit: Option<&str>,
1375 labels: Labels,
1376 value: MetricValue,
1377 timestamp: Instant,
1378 ) {
1379 let bare = family_name.into();
1380 let template = MetricFamily::new(bare.clone(), family_type);
1381 let template = match unit {
1382 Some(u) => template.with_unit(u),
1383 None => template,
1384 };
1385 let effective_name = template.name().to_string();
1386 let point = MetricPoint::new(value, timestamp);
1387 if let Some(fam) = self.families.iter_mut().find(|f| f.name == effective_name) {
1388 assert_eq!(
1389 fam.r#type, family_type,
1390 "family '{}' already exists as {:?}; cannot insert as {:?}",
1391 effective_name, fam.r#type, family_type,
1392 );
1393 fam.insert(Metric::single(labels, point));
1394 } else {
1395 let mut fam = template;
1396 fam.insert(Metric::single(labels, point));
1397 self.families.push(fam);
1398 }
1399 }
1400
1401 /// Insert a counter observation. Convenience over
1402 /// [`insert_metric`] that builds the [`CounterValue`] for you.
1403 pub fn insert_counter(
1404 &mut self,
1405 family_name: impl Into<String>,
1406 labels: Labels,
1407 cumulative: u64,
1408 timestamp: Instant,
1409 ) {
1410 self.insert_metric(
1411 family_name,
1412 MetricType::Counter,
1413 labels,
1414 MetricValue::Counter(CounterValue::new(cumulative)),
1415 timestamp,
1416 );
1417 }
1418
1419 /// Counter variant of [`insert_metric_with_unit`]. `cumulative` is the
1420 /// counter's running total (the cadence capture path passes the
1421 /// instrument's absolute); per-interval deltas are derived downstream.
1422 pub fn insert_counter_with_unit(
1423 &mut self,
1424 family_name: impl Into<String>,
1425 unit: Option<&str>,
1426 labels: Labels,
1427 cumulative: u64,
1428 timestamp: Instant,
1429 ) {
1430 self.insert_metric_with_unit(
1431 family_name,
1432 MetricType::Counter,
1433 unit,
1434 labels,
1435 MetricValue::Counter(CounterValue::new(cumulative)),
1436 timestamp,
1437 );
1438 }
1439
1440 /// Insert a gauge observation. Convenience over [`insert_metric`].
1441 pub fn insert_gauge(
1442 &mut self,
1443 family_name: impl Into<String>,
1444 labels: Labels,
1445 value: f64,
1446 timestamp: Instant,
1447 ) {
1448 self.insert_metric(
1449 family_name,
1450 MetricType::Gauge,
1451 labels,
1452 MetricValue::Gauge(GaugeValue::new(value)),
1453 timestamp,
1454 );
1455 }
1456
1457 /// Gauge variant of [`insert_metric_with_unit`].
1458 pub fn insert_gauge_with_unit(
1459 &mut self,
1460 family_name: impl Into<String>,
1461 unit: Option<&str>,
1462 labels: Labels,
1463 value: f64,
1464 timestamp: Instant,
1465 ) {
1466 self.insert_metric_with_unit(
1467 family_name,
1468 MetricType::Gauge,
1469 unit,
1470 labels,
1471 MetricValue::Gauge(GaugeValue::new(value)),
1472 timestamp,
1473 );
1474 }
1475
1476 /// Insert a histogram observation. Convenience over
1477 /// [`insert_metric`] that wraps the HDR reservoir into a
1478 /// [`HistogramValue`] and computes `count`/`sum` from it.
1479 pub fn insert_histogram(
1480 &mut self,
1481 family_name: impl Into<String>,
1482 labels: Labels,
1483 reservoir: HdrHistogram<u64>,
1484 timestamp: Instant,
1485 ) {
1486 self.insert_metric(
1487 family_name,
1488 MetricType::Histogram,
1489 labels,
1490 MetricValue::Histogram(HistogramValue::from_hdr(reservoir)),
1491 timestamp,
1492 );
1493 }
1494
1495 /// Histogram variant of [`insert_metric_with_unit`].
1496 pub fn insert_histogram_with_unit(
1497 &mut self,
1498 family_name: impl Into<String>,
1499 unit: Option<&str>,
1500 labels: Labels,
1501 reservoir: HdrHistogram<u64>,
1502 timestamp: Instant,
1503 ) {
1504 self.insert_metric_with_unit(
1505 family_name,
1506 MetricType::Histogram,
1507 unit,
1508 labels,
1509 MetricValue::Histogram(HistogramValue::from_hdr(reservoir)),
1510 timestamp,
1511 );
1512 }
1513
1514 /// Like [`insert_histogram_with_unit`] but stamps the **cumulative**
1515 /// observation count (the instrument's lifetime total) onto the
1516 /// value, so `HistogramValue::cumulative_count` is the authoritative
1517 /// running total — like the counter's `cumulative`. The cadence
1518 /// capture path uses this; the queryapi then exposes a cumulative
1519 /// histogram count over which `rate()` is PromQL-correct.
1520 pub fn insert_histogram_with_unit_cumulative(
1521 &mut self,
1522 family_name: impl Into<String>,
1523 unit: Option<&str>,
1524 labels: Labels,
1525 reservoir: HdrHistogram<u64>,
1526 cumulative_count: u64,
1527 timestamp: Instant,
1528 ) {
1529 self.insert_metric_with_unit(
1530 family_name,
1531 MetricType::Histogram,
1532 unit,
1533 labels,
1534 MetricValue::Histogram(
1535 HistogramValue::from_hdr(reservoir).with_cumulative_count(cumulative_count),
1536 ),
1537 timestamp,
1538 );
1539 }
1540}
1541
1542// =========================================================================
1543// Helpers
1544// =========================================================================
1545
1546/// Approximate observation sum from an HDR histogram. HDR doesn't
1547/// store a true sum — we estimate by `mean × count`. Sufficient for
1548/// OpenMetrics `_sum` exposition; consumers who need an exact sum
1549/// have to record it independently.
1550fn hdr_sum(h: &HdrHistogram<u64>) -> f64 {
1551 h.mean() * h.len() as f64
1552}
1553
1554/// Convenience: build a single-point Counter family.
1555pub fn counter_family(
1556 name: impl Into<String>,
1557 labels: Labels,
1558 total: u64,
1559 timestamp: Instant,
1560) -> MetricFamily {
1561 let mut f = MetricFamily::new(name, MetricType::Counter);
1562 f.insert(Metric::single(
1563 labels,
1564 MetricPoint::new(MetricValue::Counter(CounterValue::new(total)), timestamp),
1565 ));
1566 f
1567}
1568
1569/// Convenience: build a single-point Gauge family.
1570pub fn gauge_family(
1571 name: impl Into<String>,
1572 labels: Labels,
1573 value: f64,
1574 timestamp: Instant,
1575) -> MetricFamily {
1576 let mut f = MetricFamily::new(name, MetricType::Gauge);
1577 f.insert(Metric::single(
1578 labels,
1579 MetricPoint::new(MetricValue::Gauge(GaugeValue::new(value)), timestamp),
1580 ));
1581 f
1582}
1583
1584/// Convenience: build a single-point Histogram family from an HDR
1585/// reservoir.
1586pub fn histogram_family(
1587 name: impl Into<String>,
1588 labels: Labels,
1589 reservoir: HdrHistogram<u64>,
1590 timestamp: Instant,
1591) -> MetricFamily {
1592 let mut f = MetricFamily::new(name, MetricType::Histogram);
1593 f.insert(Metric::single(
1594 labels,
1595 MetricPoint::new(
1596 MetricValue::Histogram(HistogramValue::from_hdr(reservoir)),
1597 timestamp,
1598 ),
1599 ));
1600 f
1601}
1602
1603// =========================================================================
1604// Tests
1605// =========================================================================
1606
1607#[cfg(test)]
1608mod tests {
1609 use super::*;
1610
1611 fn ts() -> Instant {
1612 Instant::now()
1613 }
1614 fn empty_set() -> MetricSet {
1615 MetricSet::new(Duration::from_secs(1))
1616 }
1617
1618 #[test]
1619 fn metric_set_is_empty_by_default() {
1620 let m = empty_set();
1621 assert_eq!(m.len(), 0);
1622 assert!(m.is_empty());
1623 assert!(m.family("anything").is_none());
1624 assert_eq!(m.interval(), Duration::from_secs(1));
1625 }
1626
1627 #[test]
1628 fn metric_set_inserts_and_looks_up_by_name() {
1629 let mut m = empty_set();
1630 m.insert(counter_family(
1631 "cycles",
1632 Labels::of("phase", "load"),
1633 100,
1634 ts(),
1635 ));
1636 m.insert(gauge_family(
1637 "temp",
1638 Labels::of("phase", "load"),
1639 42.0,
1640 ts(),
1641 ));
1642
1643 assert_eq!(m.len(), 2);
1644 assert!(m.family("cycles").is_some());
1645 assert!(m.family("temp").is_some());
1646 assert!(m.family("missing").is_none());
1647 }
1648
1649 #[test]
1650 #[should_panic(expected = "names must be unique")]
1651 fn metric_set_rejects_duplicate_family_names() {
1652 let mut m = empty_set();
1653 m.insert(counter_family("cycles", Labels::of("a", "1"), 1, ts()));
1654 m.insert(counter_family("cycles", Labels::of("a", "2"), 2, ts()));
1655 }
1656
1657 fn make_counter_set(interval: Duration, value: u64) -> MetricSet {
1658 let mut s = MetricSet::new(interval);
1659 s.insert(counter_family(
1660 "cycles",
1661 Labels::of("name", "ops"),
1662 value,
1663 Instant::now(),
1664 ));
1665 s
1666 }
1667
1668 fn make_histogram_set(interval: Duration, values: &[u64]) -> MetricSet {
1669 let mut h = HdrHistogram::<u64>::new_with_bounds(1, 3_600_000_000_000, 3).unwrap();
1670 for v in values {
1671 h.record(*v).unwrap();
1672 }
1673 let mut s = MetricSet::new(interval);
1674 s.insert(histogram_family(
1675 "latency",
1676 Labels::of("name", "rt"),
1677 h,
1678 Instant::now(),
1679 ));
1680 s
1681 }
1682
1683 fn make_gauge_set(interval: Duration, value: f64) -> MetricSet {
1684 let mut s = MetricSet::new(interval);
1685 s.insert(gauge_family(
1686 "temp",
1687 Labels::of("name", "x"),
1688 value,
1689 Instant::now(),
1690 ));
1691 s
1692 }
1693
1694 #[test]
1695 fn coalesce_empty_returns_empty() {
1696 let merged = MetricSet::coalesce(&[]);
1697 assert!(merged.is_empty());
1698 }
1699
1700 #[test]
1701 fn coalesce_single_clones() {
1702 let s = make_counter_set(Duration::from_secs(1), 10);
1703 let m = MetricSet::coalesce(std::slice::from_ref(&s));
1704 assert_eq!(m.interval(), Duration::from_secs(1));
1705 let f = m.family("cycles").unwrap();
1706 let c = match f.metrics().next().unwrap().point().unwrap().value() {
1707 MetricValue::Counter(c) => c.cumulative,
1708 _ => panic!("wrong type"),
1709 };
1710 assert_eq!(c, 10);
1711 }
1712
1713 #[test]
1714 fn coalesce_counters_keep_latest_cumulative_sum_intervals() {
1715 // Time-coalesce of one series keeps the latest (window-end)
1716 // cumulative — never a sum — and sums the intervals. (Monotonic
1717 // cumulatives in window order; the last one wins.)
1718 let merged = MetricSet::coalesce(&[
1719 make_counter_set(Duration::from_secs(1), 10),
1720 make_counter_set(Duration::from_secs(1), 25),
1721 make_counter_set(Duration::from_secs(1), 42),
1722 ]);
1723 assert_eq!(merged.interval(), Duration::from_secs(3));
1724 let cumulative = match merged
1725 .family("cycles")
1726 .unwrap()
1727 .metrics()
1728 .next()
1729 .unwrap()
1730 .point()
1731 .unwrap()
1732 .value()
1733 {
1734 MetricValue::Counter(c) => c.cumulative,
1735 _ => panic!("wrong type"),
1736 };
1737 assert_eq!(cumulative, 42, "latest window's cumulative, not the sum");
1738 }
1739
1740 #[test]
1741 fn coalesce_histograms_merge_reservoirs() {
1742 let merged = MetricSet::coalesce(&[
1743 make_histogram_set(Duration::from_secs(1), &[1_000, 2_000, 3_000]),
1744 make_histogram_set(Duration::from_secs(1), &[4_000, 5_000]),
1745 ]);
1746 let hv = match merged
1747 .family("latency")
1748 .unwrap()
1749 .metrics()
1750 .next()
1751 .unwrap()
1752 .point()
1753 .unwrap()
1754 .value()
1755 {
1756 MetricValue::Histogram(h) => h.clone(),
1757 _ => panic!("wrong type"),
1758 };
1759 assert_eq!(hv.count, 5);
1760 assert!(hv.reservoir.max() >= 4_900);
1761 }
1762
1763 #[test]
1764 fn coalesce_gauges_last_write_wins() {
1765 // Two snapshots: (1s @ 10.0) then (2s @ 20.0). The coalesced
1766 // window's sample is the LAST-WRITTEN value — the
1767 // OpenMetrics/Prometheus/OTel gauge contract: a sample is the
1768 // scalar as of window end; summarization happens at the query
1769 // point via *_over_time. (This replaced interval-weighted
1770 // averaging, which stored 16.67 here — a value never written —
1771 // and turned set-once facts stored as gauges into fractions,
1772 // 2026-08-08.)
1773 let merged = MetricSet::coalesce(&[
1774 make_gauge_set(Duration::from_secs(1), 10.0),
1775 make_gauge_set(Duration::from_secs(2), 20.0),
1776 ]);
1777 let v = match merged
1778 .family("temp")
1779 .unwrap()
1780 .metrics()
1781 .next()
1782 .unwrap()
1783 .point()
1784 .unwrap()
1785 .value()
1786 {
1787 MetricValue::Gauge(g) => g.value,
1788 _ => panic!("wrong type"),
1789 };
1790 assert_eq!(v, 20.0, "last written value wins, got {v}");
1791 }
1792
1793 #[test]
1794 fn coalesce_disjoint_label_sets_appended() {
1795 // Same family name, different LabelSets — should NOT combine.
1796 let mut a = MetricSet::new(Duration::from_secs(1));
1797 a.insert(counter_family(
1798 "cycles",
1799 Labels::of("phase", "load"),
1800 100,
1801 ts(),
1802 ));
1803 let mut b = MetricSet::new(Duration::from_secs(1));
1804 b.insert(counter_family(
1805 "cycles",
1806 Labels::of("phase", "verify"),
1807 50,
1808 ts(),
1809 ));
1810
1811 let merged = MetricSet::coalesce(&[a, b]);
1812 let f = merged.family("cycles").unwrap();
1813 assert_eq!(f.len(), 2);
1814 let load = f.metric_with_labels(&Labels::of("phase", "load")).unwrap();
1815 let verify = f
1816 .metric_with_labels(&Labels::of("phase", "verify"))
1817 .unwrap();
1818 match load.point().unwrap().value() {
1819 MetricValue::Counter(c) => assert_eq!(c.cumulative, 100),
1820 _ => panic!(),
1821 }
1822 match verify.point().unwrap().value() {
1823 MetricValue::Counter(c) => assert_eq!(c.cumulative, 50),
1824 _ => panic!(),
1825 }
1826 }
1827
1828 #[test]
1829 fn metric_family_records_type_and_optional_metadata() {
1830 // SRD-40b §1 / SRD-40a §4.3: `with_unit` appends the
1831 // `_<unit>` suffix to the family name when the invariant
1832 // is not already met. Both surfaces (name + unit) derive
1833 // from this single declaration.
1834 let f = MetricFamily::new("latency", MetricType::Histogram)
1835 .with_unit("nanoseconds")
1836 .with_help("End-to-end op latency");
1837 assert_eq!(f.name(), "latency_nanoseconds");
1838 assert_eq!(f.r#type(), MetricType::Histogram);
1839 assert_eq!(f.unit(), Some("nanoseconds"));
1840 assert_eq!(f.help(), Some("End-to-end op latency"));
1841 }
1842
1843 #[test]
1844 fn with_unit_preserves_name_when_suffix_already_present() {
1845 // No double-suffixing: if the caller already wrote the
1846 // canonical name, `with_unit` is a no-op for the name.
1847 let f = MetricFamily::new("memory_bytes", MetricType::Gauge).with_unit("bytes");
1848 assert_eq!(f.name(), "memory_bytes");
1849 assert_eq!(f.unit(), Some("bytes"));
1850 }
1851
1852 #[test]
1853 fn with_unit_preserves_name_when_unit_precedes_exposition_suffix() {
1854 // OpenMetrics §4.4 / §5.x: the unit may sit before a
1855 // known exposition suffix (e.g. `_total`).
1856 let f = MetricFamily::new("process_cpu_seconds_total", MetricType::Counter)
1857 .with_unit("seconds");
1858 assert_eq!(f.name(), "process_cpu_seconds_total");
1859 assert_eq!(f.unit(), Some("seconds"));
1860 }
1861
1862 #[test]
1863 #[should_panic(expected = "LabelSets must be unique")]
1864 fn metric_family_rejects_duplicate_labelsets() {
1865 let mut f = MetricFamily::new("cycles", MetricType::Counter);
1866 f.insert(Metric::single(
1867 Labels::of("phase", "load"),
1868 MetricPoint::untimed(MetricValue::Counter(CounterValue::new(1))),
1869 ));
1870 f.insert(Metric::single(
1871 Labels::of("phase", "load"),
1872 MetricPoint::untimed(MetricValue::Counter(CounterValue::new(2))),
1873 ));
1874 }
1875
1876 #[test]
1877 fn metric_lookup_by_labels_matches_identity() {
1878 let mut f = MetricFamily::new("cycles", MetricType::Counter);
1879 f.insert(Metric::single(
1880 Labels::of("phase", "load"),
1881 MetricPoint::untimed(MetricValue::Counter(CounterValue::new(10))),
1882 ));
1883 f.insert(Metric::single(
1884 Labels::of("phase", "verify"),
1885 MetricPoint::untimed(MetricValue::Counter(CounterValue::new(20))),
1886 ));
1887
1888 let load = f.metric_with_labels(&Labels::of("phase", "load")).unwrap();
1889 match load.point().unwrap().value() {
1890 MetricValue::Counter(c) => assert_eq!(c.cumulative, 10),
1891 _ => panic!("wrong type"),
1892 }
1893 assert!(
1894 f.metric_with_labels(&Labels::of("phase", "missing"))
1895 .is_none()
1896 );
1897 }
1898
1899 #[test]
1900 fn metric_type_strings_match_open_metrics_spec() {
1901 assert_eq!(MetricType::Counter.as_str(), "counter");
1902 assert_eq!(MetricType::Gauge.as_str(), "gauge");
1903 assert_eq!(MetricType::Histogram.as_str(), "histogram");
1904 assert_eq!(MetricType::GaugeHistogram.as_str(), "gaugehistogram");
1905 assert_eq!(MetricType::Summary.as_str(), "summary");
1906 assert_eq!(MetricType::Info.as_str(), "info");
1907 assert_eq!(MetricType::StateSet.as_str(), "stateset");
1908 assert_eq!(MetricType::Unknown.as_str(), "unknown");
1909 }
1910
1911 #[test]
1912 fn counter_aggregate_sums_cumulative_keeps_earlier_created() {
1913 let t1 = Instant::now();
1914 let t0 = t1 - Duration::from_secs(60);
1915
1916 let mut a = MetricPoint::new(
1917 MetricValue::Counter(CounterValue::new(10).with_created(t1)),
1918 t1,
1919 );
1920 let b = MetricPoint::new(
1921 MetricValue::Counter(CounterValue::new(25).with_created(t0)),
1922 t1,
1923 );
1924 // Aggregate (cross-component): the running totals sum across the
1925 // matching series from different components.
1926 combine_into(&mut a, &b, CombineMode::Aggregate).unwrap();
1927 match a.value() {
1928 MetricValue::Counter(c) => {
1929 assert_eq!(
1930 c.cumulative, 35,
1931 "aggregate sums cumulative across components"
1932 );
1933 assert_eq!(c.created, Some(t0), "earliest created wins");
1934 }
1935 _ => panic!("wrong type"),
1936 }
1937 }
1938
1939 #[test]
1940 fn counter_coalesce_keeps_latest_cumulative() {
1941 let t1 = Instant::now();
1942 let t2 = t1 + Duration::from_secs(1);
1943 // Two consecutive windows of ONE series: cumulative 100 then 113.
1944 // Time-coalesce keeps the latest (window-end) running total — never
1945 // a sum. Per-window deltas are derived downstream by differencing.
1946 let mut a = MetricPoint::new(MetricValue::Counter(CounterValue::new(100)), t1);
1947 let b = MetricPoint::new(MetricValue::Counter(CounterValue::new(113)), t2);
1948 combine_into(&mut a, &b, CombineMode::Coalesce).unwrap();
1949 match a.value() {
1950 MetricValue::Counter(c) => assert_eq!(
1951 c.cumulative, 113,
1952 "coalesce keeps the latest cumulative (no summing)"
1953 ),
1954 _ => panic!("wrong type"),
1955 }
1956 }
1957
1958 #[test]
1959 fn histogram_combine_adds_reservoirs_re_derives_count_and_sum() {
1960 let mut h1 = HdrHistogram::<u64>::new_with_bounds(1, 3_600_000_000_000, 3).unwrap();
1961 h1.record(1_000_000).unwrap();
1962 h1.record(2_000_000).unwrap();
1963 let mut h2 = HdrHistogram::<u64>::new_with_bounds(1, 3_600_000_000_000, 3).unwrap();
1964 h2.record(3_000_000).unwrap();
1965
1966 let mut a = MetricPoint::new(
1967 MetricValue::Histogram(HistogramValue::from_hdr(h1)),
1968 Instant::now(),
1969 );
1970 let b = MetricPoint::new(
1971 MetricValue::Histogram(HistogramValue::from_hdr(h2)),
1972 Instant::now(),
1973 );
1974 combine_into(&mut a, &b, CombineMode::Coalesce).unwrap();
1975 match a.value() {
1976 MetricValue::Histogram(h) => {
1977 assert_eq!(h.count, 3);
1978 assert!(h.sum > 0.0);
1979 assert!(h.reservoir.max() >= 3_000_000);
1980 }
1981 _ => panic!("wrong type"),
1982 }
1983 }
1984
1985 #[test]
1986 fn gauge_combine_keeps_most_recent_value() {
1987 let t1 = Instant::now();
1988 let t2 = t1 + Duration::from_secs(1);
1989 let mut a = MetricPoint::new(MetricValue::Gauge(GaugeValue::new(5.0)), t1);
1990 let b = MetricPoint::new(MetricValue::Gauge(GaugeValue::new(9.0)), t2);
1991 combine_into(&mut a, &b, CombineMode::Coalesce).unwrap();
1992 match a.value() {
1993 MetricValue::Gauge(g) => assert_eq!(g.value, 9.0),
1994 _ => panic!("wrong type"),
1995 }
1996 }
1997
1998 #[test]
1999 fn combine_type_mismatch_is_hard_error() {
2000 let mut a = MetricPoint::untimed(MetricValue::Counter(CounterValue::new(1)));
2001 let b = MetricPoint::untimed(MetricValue::Gauge(GaugeValue::new(1.0)));
2002 let err = combine_into(&mut a, &b, CombineMode::Coalesce).unwrap_err();
2003 assert_eq!(err, CombineError::TypeMismatch);
2004 }
2005
2006 #[test]
2007 fn exemplar_most_recent_wins_on_combine() {
2008 let t1 = Instant::now();
2009 let t2 = t1 + Duration::from_secs(1);
2010 let ex_old = Exemplar::new(Labels::of("trace_id", "old"), 1.0).with_timestamp(t1);
2011 let ex_new = Exemplar::new(Labels::of("trace_id", "new"), 2.0).with_timestamp(t2);
2012
2013 let mut a = MetricPoint::new(
2014 MetricValue::Counter(CounterValue::new(5).with_exemplar(ex_old)),
2015 t1,
2016 );
2017 let b = MetricPoint::new(
2018 MetricValue::Counter(CounterValue::new(5).with_exemplar(ex_new.clone())),
2019 t2,
2020 );
2021 combine_into(&mut a, &b, CombineMode::Coalesce).unwrap();
2022 match a.value() {
2023 MetricValue::Counter(c) => {
2024 let e = c.exemplar.as_ref().expect("exemplar should survive");
2025 assert_eq!(e.labels.get("trace_id"), Some("new"));
2026 }
2027 _ => panic!("wrong type"),
2028 }
2029 }
2030
2031 #[test]
2032 fn histogram_projects_to_open_metrics_buckets() {
2033 let mut h = HdrHistogram::<u64>::new_with_bounds(1, 1_000_000, 3).unwrap();
2034 for v in [10u64, 50, 100, 500, 1000, 5000, 50_000].iter() {
2035 h.record(*v).unwrap();
2036 }
2037 let hv = HistogramValue::from_hdr(h);
2038
2039 let buckets = hv.project_buckets(&[100, 1000, 10_000]);
2040 // 3 finite + 1 +Inf = 4 total
2041 assert_eq!(buckets.len(), 4);
2042 assert_eq!(buckets[0].upper_bound, BucketBound::Finite(100));
2043 assert_eq!(buckets[3].upper_bound, BucketBound::PositiveInfinity);
2044
2045 // Cumulative: ≤100 should include 10, 50, 100 — count ≥ 3
2046 assert!(buckets[0].cumulative_count >= 3);
2047 // Final +Inf bucket equals total count
2048 assert_eq!(buckets[3].cumulative_count, hv.count);
2049 // Cumulative is monotonically non-decreasing
2050 for w in buckets.windows(2) {
2051 assert!(w[0].cumulative_count <= w[1].cumulative_count);
2052 }
2053 }
2054
2055 #[test]
2056 fn metric_point_timestamp_propagates_on_construction() {
2057 let now = Instant::now();
2058 let p = MetricPoint::new(MetricValue::Gauge(GaugeValue::new(1.0)), now);
2059 assert_eq!(p.timestamp(), Some(now));
2060
2061 let untimed = MetricPoint::untimed(MetricValue::Gauge(GaugeValue::new(1.0)));
2062 assert_eq!(untimed.timestamp(), None);
2063 }
2064}