1#[allow(unreachable_pub)]
43#[allow(unused)]
44pub(crate) mod aggregation;
45pub mod data;
46mod error;
47pub mod exporter;
48pub(crate) mod instrument;
49pub(crate) mod internal;
50#[cfg(feature = "experimental_metrics_custom_reader")]
51pub(crate) mod manual_reader;
52pub(crate) mod meter;
53mod meter_provider;
54pub(crate) mod noop;
55pub(crate) mod periodic_reader;
56#[cfg(feature = "experimental_metrics_periodicreader_with_async_runtime")]
57pub mod periodic_reader_with_async_runtime;
59pub(crate) mod pipeline;
60#[cfg(feature = "experimental_metrics_custom_reader")]
61pub mod reader;
62#[cfg(not(feature = "experimental_metrics_custom_reader"))]
63pub(crate) mod reader;
64pub(crate) mod view;
65
66#[cfg(any(feature = "testing", test))]
68#[cfg_attr(docsrs, doc(cfg(any(feature = "testing", test))))]
69pub mod in_memory_exporter;
70#[cfg(any(feature = "testing", test))]
71#[cfg_attr(docsrs, doc(cfg(any(feature = "testing", test))))]
72pub use in_memory_exporter::{InMemoryMetricExporter, InMemoryMetricExporterBuilder};
73
74pub use aggregation::*;
75#[cfg(feature = "experimental_metrics_custom_reader")]
76pub use manual_reader::*;
77pub use meter_provider::*;
78pub use periodic_reader::*;
79#[cfg(feature = "experimental_metrics_custom_reader")]
80pub use pipeline::Pipeline;
81
82pub use instrument::{Instrument, InstrumentKind, Stream, StreamBuilder};
83
84use std::hash::Hash;
85use std::str::FromStr;
86
87#[derive(Debug, Copy, Clone, Default, PartialEq, Eq, Hash)]
89#[non_exhaustive]
90pub enum Temporality {
91 #[default]
96 Cumulative,
97
98 Delta,
103
104 LowMemory,
108}
109
110impl FromStr for Temporality {
111 type Err = ();
112
113 fn from_str(s: &str) -> Result<Self, Self::Err> {
114 match s.to_lowercase().as_str() {
115 "cumulative" => Ok(Temporality::Cumulative),
116 "delta" => Ok(Temporality::Delta),
117 "lowmemory" => Ok(Temporality::LowMemory),
118 _ => Err(()),
119 }
120 }
121}
122
123#[cfg(all(test, feature = "testing"))]
124mod tests {
125 #[cfg(feature = "experimental_metrics_bound_instruments")]
126 use self::data::ExponentialHistogramDataPoint;
127 use self::data::{HistogramDataPoint, MetricData, ScopeMetrics, SumDataPoint};
128 use super::internal::Number;
129 use super::*;
130 use crate::metrics::data::ResourceMetrics;
131 use crate::metrics::internal::AggregatedMetricsAccess;
132 use crate::metrics::InMemoryMetricExporter;
133 use crate::metrics::InMemoryMetricExporterBuilder;
134 use data::GaugeDataPoint;
135 use opentelemetry::metrics::{Counter, Meter, UpDownCounter};
136 use opentelemetry::InstrumentationScope;
137 use opentelemetry::Value;
138 use opentelemetry::{metrics::MeterProvider as _, KeyValue};
139 use rand::{rngs, Rng, SeedableRng};
140 use std::cmp::{max, min};
141 use std::sync::atomic::{AtomicBool, Ordering};
142 use std::sync::{Arc, Mutex};
143 use std::thread;
144 use std::time::Duration;
145
146 use rstest::rstest;
147
148 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
155 #[cfg(not(feature = "experimental_metrics_disable_name_validation"))]
156 async fn invalid_instrument_config_noops() {
157 let invalid_instrument_names = vec![
160 "_startWithNoneAlphabet",
161 "utf8char锈",
162 "a".repeat(256).leak(),
163 "invalid name",
164 ];
165 for name in invalid_instrument_names {
166 let test_context = TestContext::new(Temporality::Cumulative);
167 let counter = test_context.meter().u64_counter(name).build();
168 counter.add(1, &[]);
169
170 let up_down_counter = test_context.meter().i64_up_down_counter(name).build();
171 up_down_counter.add(1, &[]);
172
173 let gauge = test_context.meter().f64_gauge(name).build();
174 gauge.record(1.9, &[]);
175
176 let histogram = test_context.meter().f64_histogram(name).build();
177 histogram.record(1.0, &[]);
178
179 let _observable_counter = test_context
180 .meter()
181 .u64_observable_counter(name)
182 .with_callback(move |observer| {
183 observer.observe(1, &[]);
184 })
185 .build();
186
187 let _observable_gauge = test_context
188 .meter()
189 .f64_observable_gauge(name)
190 .with_callback(move |observer| {
191 observer.observe(1.0, &[]);
192 })
193 .build();
194
195 let _observable_up_down_counter = test_context
196 .meter()
197 .i64_observable_up_down_counter(name)
198 .with_callback(move |observer| {
199 observer.observe(1, &[]);
200 })
201 .build();
202
203 test_context.flush_metrics();
204
205 test_context.check_no_metrics();
207 }
208
209 let invalid_bucket_boundaries = vec![
210 vec![1.0, 1.0], vec![1.0, 2.0, 3.0, 2.0], vec![1.0, 2.0, 3.0, 4.0, 2.5], vec![1.0, 2.0, 3.0, f64::INFINITY, 4.0], vec![1.0, 2.0, 3.0, f64::NAN], vec![f64::NEG_INFINITY, 2.0, 3.0], ];
217 for bucket_boundaries in invalid_bucket_boundaries {
218 let test_context = TestContext::new(Temporality::Cumulative);
219 let histogram = test_context
220 .meter()
221 .f64_histogram("test")
222 .with_boundaries(bucket_boundaries)
223 .build();
224 histogram.record(1.9, &[]);
225 test_context.flush_metrics();
226
227 test_context.check_no_metrics();
230 }
231 }
232
233 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
234 #[cfg(feature = "experimental_metrics_disable_name_validation")]
235 async fn valid_instrument_config_with_feature_experimental_metrics_disable_name_validation() {
236 let invalid_instrument_names = vec![
239 "_startWithNoneAlphabet",
240 "utf8char锈",
241 "",
242 "a".repeat(256).leak(),
243 "\\allow\\slash /sec",
244 "\\allow\\$$slash /sec",
245 "Total $ Count",
246 "\\test\\UsagePercent(Total) > 80%",
247 "invalid name",
248 ];
249 for name in invalid_instrument_names {
250 let test_context = TestContext::new(Temporality::Cumulative);
251 let counter = test_context.meter().u64_counter(name).build();
252 counter.add(1, &[]);
253
254 let up_down_counter = test_context.meter().i64_up_down_counter(name).build();
255 up_down_counter.add(1, &[]);
256
257 let gauge = test_context.meter().f64_gauge(name).build();
258 gauge.record(1.9, &[]);
259
260 let histogram = test_context.meter().f64_histogram(name).build();
261 histogram.record(1.0, &[]);
262
263 let _observable_counter = test_context
264 .meter()
265 .u64_observable_counter(name)
266 .with_callback(move |observer| {
267 observer.observe(1, &[]);
268 })
269 .build();
270
271 let _observable_gauge = test_context
272 .meter()
273 .f64_observable_gauge(name)
274 .with_callback(move |observer| {
275 observer.observe(1.0, &[]);
276 })
277 .build();
278
279 let _observable_up_down_counter = test_context
280 .meter()
281 .i64_observable_up_down_counter(name)
282 .with_callback(move |observer| {
283 observer.observe(1, &[]);
284 })
285 .build();
286
287 test_context.flush_metrics();
288
289 let resource_metrics = test_context
291 .exporter
292 .get_finished_metrics()
293 .expect("metrics expected to be exported");
294
295 assert!(!resource_metrics.is_empty(), "metrics should be exported");
296 }
297
298 let invalid_bucket_boundaries = vec![
301 vec![1.0, 1.0], vec![1.0, 2.0, 3.0, 2.0], vec![1.0, 2.0, 3.0, 4.0, 2.5], vec![1.0, 2.0, 3.0, f64::INFINITY, 4.0], vec![1.0, 2.0, 3.0, f64::NAN], vec![f64::NEG_INFINITY, 2.0, 3.0], ];
308 for bucket_boundaries in invalid_bucket_boundaries {
309 let test_context = TestContext::new(Temporality::Cumulative);
310 let histogram = test_context
311 .meter()
312 .f64_histogram("test")
313 .with_boundaries(bucket_boundaries)
314 .build();
315 histogram.record(1.9, &[]);
316 test_context.flush_metrics();
317
318 test_context.check_no_metrics();
321 }
322 }
323
324 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
325 async fn counter_aggregation_delta() {
326 counter_aggregation_helper(Temporality::Delta);
329 }
330
331 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
332 async fn counter_aggregation_cumulative() {
333 counter_aggregation_helper(Temporality::Cumulative);
336 }
337
338 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
339 async fn counter_aggregation_no_attributes_cumulative() {
340 let mut test_context = TestContext::new(Temporality::Cumulative);
341 let counter = test_context.u64_counter("test", "my_counter", None);
342
343 counter.add(50, &[]);
344 test_context.flush_metrics();
345
346 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
347 unreachable!()
348 };
349
350 assert_eq!(sum.data_points.len(), 1, "Expected only one data point");
351 assert!(sum.is_monotonic, "Should produce monotonic.");
352 assert_eq!(
353 sum.temporality,
354 Temporality::Cumulative,
355 "Should produce cumulative"
356 );
357
358 let data_point = &sum.data_points[0];
359 assert!(data_point.attributes.is_empty(), "Non-empty attribute set");
360 assert_eq!(data_point.value, 50, "Unexpected data point value");
361 }
362
363 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
364 async fn counter_aggregation_no_attributes_delta() {
365 let mut test_context = TestContext::new(Temporality::Delta);
366 let counter = test_context.u64_counter("test", "my_counter", None);
367
368 counter.add(50, &[]);
369 test_context.flush_metrics();
370
371 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
372 unreachable!()
373 };
374
375 assert_eq!(sum.data_points.len(), 1, "Expected only one data point");
376 assert!(sum.is_monotonic, "Should produce monotonic.");
377 assert_eq!(sum.temporality, Temporality::Delta, "Should produce delta");
378
379 let data_point = &sum.data_points[0];
380 assert!(data_point.attributes.is_empty(), "Non-empty attribute set");
381 assert_eq!(data_point.value, 50, "Unexpected data point value");
382 }
383
384 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
385 async fn counter_aggregation_overflow_delta() {
386 counter_aggregation_overflow_helper(Temporality::Delta);
387 counter_aggregation_overflow_helper_custom_limit(Temporality::Delta);
388 }
389
390 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
391 async fn counter_aggregation_overflow_cumulative() {
392 counter_aggregation_overflow_helper(Temporality::Cumulative);
393 counter_aggregation_overflow_helper_custom_limit(Temporality::Cumulative);
394 }
395
396 #[rstest]
397 #[case(Temporality::Delta, true)]
398 #[case(Temporality::Delta, false)]
399 #[case(Temporality::Cumulative, true)]
400 #[case(Temporality::Cumulative, false)]
401 fn counter_aggregation(#[case] temporality: Temporality, #[case] start_sorted: bool) {
402 counter_aggregation_attribute_order_helper(temporality, start_sorted);
405 }
406
407 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
408 async fn histogram_aggregation_cumulative() {
409 histogram_aggregation_helper(Temporality::Cumulative);
412 }
413
414 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
415 async fn histogram_aggregation_delta() {
416 histogram_aggregation_helper(Temporality::Delta);
419 }
420
421 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
422 async fn histogram_aggregation_with_custom_bounds() {
423 histogram_aggregation_with_custom_bounds_helper(Temporality::Delta);
426 histogram_aggregation_with_custom_bounds_helper(Temporality::Cumulative);
427 }
428
429 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
430 async fn histogram_aggregation_with_empty_bounds() {
431 histogram_aggregation_with_empty_bounds_helper(Temporality::Delta);
434 histogram_aggregation_with_empty_bounds_helper(Temporality::Cumulative);
435 }
436
437 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
438 async fn histogram_aggregation_with_custom_bounds_and_view() {
439 histogram_aggregation_with_custom_bounds_and_view_helper(Temporality::Delta);
442 histogram_aggregation_with_custom_bounds_and_view_helper(Temporality::Cumulative);
443 }
444
445 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
446 async fn exponential_histogram_aggregation_with_view() {
447 exponential_histogram_aggregation_with_view_helper(Temporality::Delta);
450 exponential_histogram_aggregation_with_view_helper(Temporality::Cumulative);
451 }
452
453 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
454 async fn updown_counter_aggregation_cumulative() {
455 updown_counter_aggregation_helper(Temporality::Cumulative);
458 }
459
460 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
461 async fn updown_counter_aggregation_delta() {
462 updown_counter_aggregation_helper(Temporality::Delta);
465 }
466
467 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
468 async fn gauge_aggregation() {
469 gauge_aggregation_helper(Temporality::Delta);
474 gauge_aggregation_helper(Temporality::Cumulative);
475 }
476
477 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
478 async fn observable_gauge_aggregation() {
479 observable_gauge_aggregation_helper(Temporality::Delta, false);
484 observable_gauge_aggregation_helper(Temporality::Delta, true);
485 observable_gauge_aggregation_helper(Temporality::Cumulative, false);
486 observable_gauge_aggregation_helper(Temporality::Cumulative, true);
487 }
488
489 #[rstest]
490 #[case(Temporality::Cumulative, 10, false)]
491 #[case(Temporality::Cumulative, 10, true)]
492 #[case(Temporality::Delta, 10, false)]
493 #[case(Temporality::Delta, 10, true)]
494 #[case(Temporality::Cumulative, 0, false)]
495 #[case(Temporality::Cumulative, 0, true)]
496 #[case(Temporality::Delta, 0, false)]
497 #[case(Temporality::Delta, 0, true)]
498 fn observable_counter_aggregation(
499 #[case] temporality: Temporality,
500 #[case] increment: u64,
501 #[case] is_empty_attributes: bool,
502 ) {
503 observable_counter_aggregation_helper(temporality, 100, increment, 4, is_empty_attributes);
506 }
507
508 fn observable_counter_aggregation_helper(
509 temporality: Temporality,
510 start: u64,
511 increment: u64,
512 length: u64,
513 is_empty_attributes: bool,
514 ) {
515 let mut test_context = TestContext::new(temporality);
517 let attributes = if is_empty_attributes {
518 vec![]
519 } else {
520 vec![KeyValue::new("key1", "value1")]
521 };
522 let values: Vec<u64> = (0..length).map(|i| start + i * increment).collect();
524 println!("Testing with observable values: {values:?}");
525 let values = Arc::new(values);
526 let values_clone = values.clone();
527 let i = Arc::new(Mutex::new(0));
528 let _observable_counter = test_context
529 .meter()
530 .u64_observable_counter("my_observable_counter")
531 .with_unit("my_unit")
532 .with_callback(move |observer| {
533 let mut index = i.lock().unwrap();
534 if *index < values.len() {
535 observer.observe(values[*index], &attributes);
536 *index += 1;
537 }
538 })
539 .build();
540
541 for (iter, v) in values_clone.iter().enumerate() {
542 test_context.flush_metrics();
543 let MetricData::Sum(sum) =
544 test_context.get_aggregation::<u64>("my_observable_counter", None)
545 else {
546 unreachable!()
547 };
548 assert_eq!(sum.data_points.len(), 1);
549 assert!(sum.is_monotonic, "Counter should produce monotonic.");
550 if let Temporality::Cumulative = temporality {
551 assert_eq!(
552 sum.temporality,
553 Temporality::Cumulative,
554 "Should produce cumulative"
555 );
556 } else {
557 assert_eq!(sum.temporality, Temporality::Delta, "Should produce delta");
558 }
559
560 let data_point = if is_empty_attributes {
562 &sum.data_points[0]
563 } else {
564 find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value1")
565 .expect("datapoint with key1=value1 expected")
566 };
567
568 if let Temporality::Cumulative = temporality {
569 assert_eq!(data_point.value, *v);
571 } else {
572 if iter == 0 {
575 assert_eq!(data_point.value, start);
576 } else {
577 assert_eq!(data_point.value, increment);
578 }
579 }
580
581 test_context.reset_metrics();
582 }
583 }
584
585 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
586 async fn observable_counter_delta_attribute_order_changes() {
587 let mut test_context = TestContext::new(Temporality::Delta);
588
589 let observations = Arc::new(Mutex::new(std::collections::VecDeque::from([
590 (
591 100,
592 vec![
593 KeyValue::new("A", "a"),
594 KeyValue::new("B", "b"),
595 KeyValue::new("C", "c"),
596 ],
597 ),
598 (
599 150,
600 vec![
601 KeyValue::new("B", "b"),
602 KeyValue::new("A", "a"),
603 KeyValue::new("C", "c"),
604 ],
605 ),
606 (
607 175,
608 vec![
609 KeyValue::new("C", "c"),
610 KeyValue::new("B", "b"),
611 KeyValue::new("A", "a"),
612 ],
613 ),
614 ])));
615 let observations_clone = observations.clone();
616
617 let _observable_counter = test_context
618 .meter()
619 .u64_observable_counter("my_observable_counter")
620 .with_callback(move |observer| {
621 let mut observations = observations_clone.lock().unwrap();
622 if let Some((value, attributes)) = observations.pop_front() {
623 observer.observe(value, &attributes);
624 }
625 })
626 .build();
627
628 for expected in [100, 50, 25] {
629 test_context.flush_metrics();
630
631 let MetricData::Sum(sum) =
632 test_context.get_aggregation::<u64>("my_observable_counter", None)
633 else {
634 unreachable!()
635 };
636
637 assert_eq!(sum.data_points.len(), 1);
638 assert_eq!(sum.temporality, Temporality::Delta);
639
640 let data_point = &sum.data_points[0];
641 assert_eq!(data_point.attributes.len(), 3);
642 assert!(data_point
643 .attributes
644 .iter()
645 .any(|kv| kv.key.as_str() == "A" && kv.value.as_str() == "a"));
646 assert!(data_point
647 .attributes
648 .iter()
649 .any(|kv| kv.key.as_str() == "B" && kv.value.as_str() == "b"));
650 assert!(data_point
651 .attributes
652 .iter()
653 .any(|kv| kv.key.as_str() == "C" && kv.value.as_str() == "c"));
654 assert_eq!(data_point.value, expected);
655
656 test_context.reset_metrics();
657 }
658 }
659
660 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
661 async fn observable_counter_delta_attribute_set_reappears_after_gap() {
662 let mut test_context = TestContext::new(Temporality::Delta);
680
681 let callback_state = Arc::new(Mutex::new((0u32, 0u64, Option::<u64>::None)));
684 let callback_state_clone = callback_state.clone();
685
686 let _observable_counter = test_context
687 .meter()
688 .u64_observable_counter("my_observable_counter")
689 .with_callback(move |observer| {
690 let state = callback_state_clone.lock().unwrap();
691 let (_cycle, value_a, value_b_option) = *state;
692
693 observer.observe(value_a, &[KeyValue::new("key", "A")]);
694 if let Some(value_b) = value_b_option {
695 observer.observe(value_b, &[KeyValue::new("key", "B")]);
696 }
697 })
698 .build();
699
700 {
702 *callback_state.lock().unwrap() = (1, 100, Some(50));
703 test_context.flush_metrics();
704
705 let MetricData::Sum(sum) =
706 test_context.get_aggregation::<u64>("my_observable_counter", None)
707 else {
708 unreachable!()
709 };
710
711 assert_eq!(sum.data_points.len(), 2);
712 assert_eq!(sum.temporality, Temporality::Delta);
713
714 let dp_a = find_sum_datapoint_with_key_value(&sum.data_points, "key", "A")
715 .expect("datapoint for A expected");
716 let dp_b = find_sum_datapoint_with_key_value(&sum.data_points, "key", "B")
717 .expect("datapoint for B expected");
718
719 assert_eq!(
721 dp_a.value, 100,
722 "A's delta should be 100 (first collection)"
723 );
724 assert_eq!(dp_b.value, 50, "B's delta should be 50 (first collection)");
725
726 test_context.reset_metrics();
727 }
728
729 {
731 *callback_state.lock().unwrap() = (2, 150, None);
732 test_context.flush_metrics();
733
734 let MetricData::Sum(sum) =
735 test_context.get_aggregation::<u64>("my_observable_counter", None)
736 else {
737 unreachable!()
738 };
739
740 assert_eq!(
742 sum.data_points.len(),
743 1,
744 "Only A should be exported when B is not observed"
745 );
746
747 let dp_a = find_sum_datapoint_with_key_value(&sum.data_points, "key", "A")
748 .expect("datapoint for A expected");
749 assert_eq!(dp_a.value, 50, "A's delta should be 50 (150 - 100)");
750
751 test_context.reset_metrics();
752 }
753
754 {
756 *callback_state.lock().unwrap() = (3, 200, Some(80));
757 test_context.flush_metrics();
758
759 let MetricData::Sum(sum) =
760 test_context.get_aggregation::<u64>("my_observable_counter", None)
761 else {
762 unreachable!()
763 };
764
765 assert_eq!(sum.data_points.len(), 2);
766
767 let dp_a = find_sum_datapoint_with_key_value(&sum.data_points, "key", "A")
768 .expect("datapoint for A expected");
769 let dp_b = find_sum_datapoint_with_key_value(&sum.data_points, "key", "B")
770 .expect("datapoint for B expected");
771
772 assert_eq!(dp_a.value, 50, "A's delta should be 50 (200 - 150)");
773
774 assert_eq!(
779 dp_b.value, 80,
780 "B's delta should be 80 (fresh start after gap, not 30 from last known value)"
781 );
782
783 test_context.reset_metrics();
784 }
785 }
786
787 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
788 async fn empty_meter_name_retained() {
789 async fn meter_name_retained_helper(
790 meter: Meter,
791 provider: SdkMeterProvider,
792 exporter: InMemoryMetricExporter,
793 ) {
794 let counter = meter.u64_counter("my_counter").build();
796
797 counter.add(10, &[]);
798 provider.force_flush().unwrap();
799
800 let resource_metrics = exporter
802 .get_finished_metrics()
803 .expect("metrics are expected to be exported.");
804 assert!(
805 resource_metrics[0].scope_metrics[0].metrics.len() == 1,
806 "There should be a single metric"
807 );
808 let meter_name = resource_metrics[0].scope_metrics[0].scope.name();
809 assert_eq!(meter_name, "");
810 }
811
812 let exporter = InMemoryMetricExporter::default();
813 let meter_provider = SdkMeterProvider::builder()
814 .with_periodic_exporter(exporter.clone())
815 .build();
816
817 let meter1 = meter_provider.meter("");
819 meter_name_retained_helper(meter1, meter_provider.clone(), exporter.clone()).await;
820
821 let meter_scope = InstrumentationScope::builder("").build();
822 let meter2 = meter_provider.meter_with_scope(meter_scope);
823 meter_name_retained_helper(meter2, meter_provider, exporter).await;
824 }
825
826 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
827 async fn counter_duplicate_instrument_merge() {
828 let exporter = InMemoryMetricExporter::default();
830 let meter_provider = SdkMeterProvider::builder()
831 .with_periodic_exporter(exporter.clone())
832 .build();
833
834 let meter = meter_provider.meter("test");
836 let counter = meter
837 .u64_counter("my_counter")
838 .with_unit("my_unit")
839 .with_description("my_description")
840 .build();
841
842 let counter_duplicated = meter
843 .u64_counter("my_counter")
844 .with_unit("my_unit")
845 .with_description("my_description")
846 .build();
847
848 let attribute = vec![KeyValue::new("key1", "value1")];
849 counter.add(10, &attribute);
850 counter_duplicated.add(5, &attribute);
851
852 meter_provider.force_flush().unwrap();
853
854 let resource_metrics = exporter
856 .get_finished_metrics()
857 .expect("metrics are expected to be exported.");
858 assert!(
859 resource_metrics[0].scope_metrics[0].metrics.len() == 1,
860 "There should be single metric merging duplicate instruments"
861 );
862 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
863 assert_eq!(metric.name, "my_counter");
864 assert_eq!(metric.unit, "my_unit");
865 let MetricData::Sum(sum) = u64::extract_metrics_data_ref(&metric.data)
866 .expect("Sum aggregation expected for Counter instruments by default")
867 else {
868 unreachable!()
869 };
870
871 assert_eq!(sum.data_points.len(), 1);
873
874 let datapoint = &sum.data_points[0];
875 assert_eq!(datapoint.value, 15);
876 }
877
878 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
879 async fn counter_duplicate_instrument_different_meter_no_merge() {
880 let exporter = InMemoryMetricExporter::default();
882 let meter_provider = SdkMeterProvider::builder()
883 .with_periodic_exporter(exporter.clone())
884 .build();
885
886 let meter1 = meter_provider.meter("test.meter1");
888 let meter2 = meter_provider.meter("test.meter2");
889 let counter1 = meter1
890 .u64_counter("my_counter")
891 .with_unit("my_unit")
892 .with_description("my_description")
893 .build();
894
895 let counter2 = meter2
896 .u64_counter("my_counter")
897 .with_unit("my_unit")
898 .with_description("my_description")
899 .build();
900
901 let attribute = vec![KeyValue::new("key1", "value1")];
902 counter1.add(10, &attribute);
903 counter2.add(5, &attribute);
904
905 meter_provider.force_flush().unwrap();
906
907 let resource_metrics = exporter
909 .get_finished_metrics()
910 .expect("metrics are expected to be exported.");
911 assert!(
912 resource_metrics[0].scope_metrics.len() == 2,
913 "There should be 2 separate scope"
914 );
915 assert!(
916 resource_metrics[0].scope_metrics[0].metrics.len() == 1,
917 "There should be single metric for the scope"
918 );
919 assert!(
920 resource_metrics[0].scope_metrics[1].metrics.len() == 1,
921 "There should be single metric for the scope"
922 );
923
924 let scope1 = find_scope_metric(&resource_metrics[0].scope_metrics, "test.meter1");
925 let scope2 = find_scope_metric(&resource_metrics[0].scope_metrics, "test.meter2");
926
927 if let Some(scope1) = scope1 {
928 let metric1 = &scope1.metrics[0];
929 assert_eq!(metric1.name, "my_counter");
930 assert_eq!(metric1.unit, "my_unit");
931 assert_eq!(metric1.description, "my_description");
932 let MetricData::Sum(sum1) = u64::extract_metrics_data_ref(&metric1.data)
933 .expect("Sum aggregation expected for Counter instruments by default")
934 else {
935 unreachable!()
936 };
937
938 assert_eq!(sum1.data_points.len(), 1);
940
941 let datapoint1 = &sum1.data_points[0];
942 assert_eq!(datapoint1.value, 10);
943 } else {
944 panic!("No MetricScope found for 'test.meter1'");
945 }
946
947 if let Some(scope2) = scope2 {
948 let metric2 = &scope2.metrics[0];
949 assert_eq!(metric2.name, "my_counter");
950 assert_eq!(metric2.unit, "my_unit");
951 assert_eq!(metric2.description, "my_description");
952
953 let MetricData::Sum(sum2) = u64::extract_metrics_data_ref(&metric2.data)
954 .expect("Sum aggregation expected for Counter instruments by default")
955 else {
956 unreachable!()
957 };
958
959 assert_eq!(sum2.data_points.len(), 1);
961
962 let datapoint2 = &sum2.data_points[0];
963 assert_eq!(datapoint2.value, 5);
964 } else {
965 panic!("No MetricScope found for 'test.meter2'");
966 }
967 }
968
969 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
970 async fn instrumentation_scope_identity_test() {
971 let exporter = InMemoryMetricExporter::default();
973 let meter_provider = SdkMeterProvider::builder()
974 .with_periodic_exporter(exporter.clone())
975 .build();
976
977 let make_scope = |attributes| {
981 InstrumentationScope::builder("test.meter")
982 .with_version("v0.1.0")
983 .with_schema_url("http://example.com")
984 .with_attributes(attributes)
985 .build()
986 };
987
988 let meter1 =
989 meter_provider.meter_with_scope(make_scope(vec![KeyValue::new("key", "value1")]));
990 let meter2 =
991 meter_provider.meter_with_scope(make_scope(vec![KeyValue::new("key", "value1")]));
992
993 let counter1 = meter1
994 .u64_counter("my_counter")
995 .with_unit("my_unit")
996 .with_description("my_description")
997 .build();
998
999 let counter2 = meter2
1000 .u64_counter("my_counter")
1001 .with_unit("my_unit")
1002 .with_description("my_description")
1003 .build();
1004
1005 let attribute = vec![KeyValue::new("key1", "value1")];
1006 counter1.add(10, &attribute);
1007 counter2.add(5, &attribute);
1008
1009 meter_provider.force_flush().unwrap();
1010
1011 let resource_metrics = exporter
1013 .get_finished_metrics()
1014 .expect("metrics are expected to be exported.");
1015 println!("resource_metrics: {resource_metrics:?}");
1016 assert!(
1017 resource_metrics[0].scope_metrics.len() == 1,
1018 "There should be a single scope as the meters are identical"
1019 );
1020 assert!(
1021 resource_metrics[0].scope_metrics[0].metrics.len() == 1,
1022 "There should be single metric for the scope as instruments are identical"
1023 );
1024
1025 let scope = &resource_metrics[0].scope_metrics[0].scope;
1026 assert_eq!(scope.name(), "test.meter");
1027 assert_eq!(scope.version(), Some("v0.1.0"));
1028 assert_eq!(scope.schema_url(), Some("http://example.com"));
1029
1030 assert!(scope.attributes().eq(&[KeyValue::new("key", "value1")]));
1033
1034 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1035 assert_eq!(metric.name, "my_counter");
1036 assert_eq!(metric.unit, "my_unit");
1037 assert_eq!(metric.description, "my_description");
1038
1039 let MetricData::Sum(sum) = u64::extract_metrics_data_ref(&metric.data)
1040 .expect("Sum aggregation expected for Counter instruments by default")
1041 else {
1042 unreachable!()
1043 };
1044
1045 assert_eq!(sum.data_points.len(), 1);
1047
1048 let datapoint = &sum.data_points[0];
1049 assert_eq!(datapoint.value, 15);
1050 }
1051
1052 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1053 async fn histogram_aggregation_with_invalid_aggregation_should_proceed_as_if_view_not_exist() {
1054 let exporter = InMemoryMetricExporter::default();
1059 let view = |i: &Instrument| {
1060 if i.name == "test_histogram" {
1061 Stream::builder()
1062 .with_aggregation(aggregation::Aggregation::ExplicitBucketHistogram {
1063 boundaries: vec![0.9, 1.9, 1.2, 1.3, 1.4, 1.5], record_min_max: false,
1065 })
1066 .with_name("test_histogram_renamed")
1067 .with_unit("test_unit_renamed")
1068 .build()
1069 .ok()
1070 } else {
1071 None
1072 }
1073 };
1074 let meter_provider = SdkMeterProvider::builder()
1075 .with_periodic_exporter(exporter.clone())
1076 .with_view(view)
1077 .build();
1078
1079 let meter = meter_provider.meter("test");
1081 let histogram = meter
1082 .f64_histogram("test_histogram")
1083 .with_unit("test_unit")
1084 .build();
1085
1086 histogram.record(1.5, &[KeyValue::new("key1", "value1")]);
1087 meter_provider.force_flush().unwrap();
1088
1089 let resource_metrics = exporter
1091 .get_finished_metrics()
1092 .expect("metrics are expected to be exported.");
1093 assert!(!resource_metrics.is_empty());
1094 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1095 assert_eq!(
1096 metric.name, "test_histogram",
1097 "View rename should be ignored and original name retained."
1098 );
1099 assert_eq!(
1100 metric.unit, "test_unit",
1101 "View rename of unit should be ignored and original unit retained."
1102 );
1103 }
1104
1105 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1106 async fn counter_with_lastvalue_aggregation_uses_default() {
1107 let exporter = InMemoryMetricExporter::default();
1114 let view = |i: &Instrument| {
1115 if i.name == "my_counter" {
1116 Stream::builder()
1117 .with_aggregation(aggregation::Aggregation::LastValue)
1118 .with_name("my_counter_renamed")
1119 .build()
1120 .ok()
1121 } else {
1122 None
1123 }
1124 };
1125 let meter_provider = SdkMeterProvider::builder()
1126 .with_periodic_exporter(exporter.clone())
1127 .with_view(view)
1128 .build();
1129
1130 let meter = meter_provider.meter("test");
1132 let counter = meter.u64_counter("my_counter").build();
1133 counter.add(10, &[KeyValue::new("key1", "value1")]);
1134 meter_provider.force_flush().unwrap();
1135
1136 let resource_metrics = exporter
1138 .get_finished_metrics()
1139 .expect("metrics are expected to be exported.");
1140 assert!(!resource_metrics.is_empty());
1141 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1142 assert_eq!(
1144 metric.name, "my_counter",
1145 "View rename should be ignored due to incompatible aggregation."
1146 );
1147 assert!(
1149 matches!(
1150 &metric.data,
1151 data::AggregatedMetrics::U64(data::MetricData::Sum(_))
1152 ),
1153 "Counter should use default Sum aggregation when LastValue is incompatible."
1154 );
1155 }
1156
1157 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1158 async fn gauge_with_sum_aggregation_uses_default() {
1159 let exporter = InMemoryMetricExporter::default();
1166 let view = |i: &Instrument| {
1167 if i.name == "my_gauge" {
1168 Stream::builder()
1169 .with_aggregation(aggregation::Aggregation::Sum)
1170 .with_name("my_gauge_renamed")
1171 .build()
1172 .ok()
1173 } else {
1174 None
1175 }
1176 };
1177 let meter_provider = SdkMeterProvider::builder()
1178 .with_periodic_exporter(exporter.clone())
1179 .with_view(view)
1180 .build();
1181
1182 let meter = meter_provider.meter("test");
1184 let gauge = meter.f64_gauge("my_gauge").build();
1185 gauge.record(42.0, &[KeyValue::new("key1", "value1")]);
1186 meter_provider.force_flush().unwrap();
1187
1188 let resource_metrics = exporter
1190 .get_finished_metrics()
1191 .expect("metrics are expected to be exported.");
1192 assert!(!resource_metrics.is_empty());
1193 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1194 assert_eq!(
1196 metric.name, "my_gauge",
1197 "View rename should be ignored due to incompatible aggregation."
1198 );
1199 assert!(
1201 matches!(
1202 &metric.data,
1203 data::AggregatedMetrics::F64(data::MetricData::Gauge(_))
1204 ),
1205 "Gauge should use default LastValue aggregation when Sum is incompatible."
1206 );
1207 }
1208
1209 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1210 async fn updowncounter_with_lastvalue_aggregation_uses_default() {
1211 let exporter = InMemoryMetricExporter::default();
1218 let view = |i: &Instrument| {
1219 if i.name == "my_updown_counter" {
1220 Stream::builder()
1221 .with_aggregation(aggregation::Aggregation::LastValue)
1222 .with_name("my_updown_counter_renamed")
1223 .build()
1224 .ok()
1225 } else {
1226 None
1227 }
1228 };
1229 let meter_provider = SdkMeterProvider::builder()
1230 .with_periodic_exporter(exporter.clone())
1231 .with_view(view)
1232 .build();
1233
1234 let meter = meter_provider.meter("test");
1236 let counter = meter.i64_up_down_counter("my_updown_counter").build();
1237 counter.add(-5, &[KeyValue::new("key1", "value1")]);
1238 meter_provider.force_flush().unwrap();
1239
1240 let resource_metrics = exporter
1242 .get_finished_metrics()
1243 .expect("metrics are expected to be exported.");
1244 assert!(!resource_metrics.is_empty());
1245 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1246 assert_eq!(
1248 metric.name, "my_updown_counter",
1249 "View rename should be ignored due to incompatible aggregation."
1250 );
1251 assert!(
1253 matches!(
1254 &metric.data,
1255 data::AggregatedMetrics::I64(data::MetricData::Sum(_))
1256 ),
1257 "UpDownCounter should use default Sum aggregation when LastValue is incompatible."
1258 );
1259 }
1260
1261 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1262 async fn histogram_with_lastvalue_aggregation_uses_default() {
1263 let exporter = InMemoryMetricExporter::default();
1270 let view = |i: &Instrument| {
1271 if i.name == "my_histogram" {
1272 Stream::builder()
1273 .with_aggregation(aggregation::Aggregation::LastValue)
1274 .with_name("my_histogram_renamed")
1275 .build()
1276 .ok()
1277 } else {
1278 None
1279 }
1280 };
1281 let meter_provider = SdkMeterProvider::builder()
1282 .with_periodic_exporter(exporter.clone())
1283 .with_view(view)
1284 .build();
1285
1286 let meter = meter_provider.meter("test");
1288 let histogram = meter.f64_histogram("my_histogram").build();
1289 histogram.record(42.0, &[KeyValue::new("key1", "value1")]);
1290 meter_provider.force_flush().unwrap();
1291
1292 let resource_metrics = exporter
1294 .get_finished_metrics()
1295 .expect("metrics are expected to be exported.");
1296 assert!(!resource_metrics.is_empty());
1297 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1298 assert_eq!(
1300 metric.name, "my_histogram",
1301 "View rename should be ignored due to incompatible aggregation."
1302 );
1303 assert!(
1305 matches!(
1306 &metric.data,
1307 data::AggregatedMetrics::F64(data::MetricData::Histogram(_))
1308 ),
1309 "Histogram should use default ExplicitBucketHistogram aggregation when LastValue is incompatible."
1310 );
1311 }
1312
1313 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1314 async fn observable_gauge_with_sum_aggregation_uses_default() {
1315 let exporter = InMemoryMetricExporter::default();
1322 let view = |i: &Instrument| {
1323 if i.name == "my_observable_gauge" {
1324 Stream::builder()
1325 .with_aggregation(aggregation::Aggregation::Sum)
1326 .with_name("my_observable_gauge_renamed")
1327 .build()
1328 .ok()
1329 } else {
1330 None
1331 }
1332 };
1333 let meter_provider = SdkMeterProvider::builder()
1334 .with_periodic_exporter(exporter.clone())
1335 .with_view(view)
1336 .build();
1337
1338 let meter = meter_provider.meter("test");
1340 let _observable_gauge = meter
1341 .f64_observable_gauge("my_observable_gauge")
1342 .with_callback(|observer| {
1343 observer.observe(42.0, &[KeyValue::new("key1", "value1")]);
1344 })
1345 .build();
1346 meter_provider.force_flush().unwrap();
1347
1348 let resource_metrics = exporter
1350 .get_finished_metrics()
1351 .expect("metrics are expected to be exported.");
1352 assert!(!resource_metrics.is_empty());
1353 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1354 assert_eq!(
1356 metric.name, "my_observable_gauge",
1357 "View rename should be ignored due to incompatible aggregation."
1358 );
1359 assert!(
1361 matches!(
1362 &metric.data,
1363 data::AggregatedMetrics::F64(data::MetricData::Gauge(_))
1364 ),
1365 "Observable Gauge should use default LastValue aggregation when Sum is incompatible."
1366 );
1367 }
1368
1369 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1370 async fn observable_counter_with_lastvalue_aggregation_uses_default() {
1371 let exporter = InMemoryMetricExporter::default();
1378 let view = |i: &Instrument| {
1379 if i.name == "my_observable_counter" {
1380 Stream::builder()
1381 .with_aggregation(aggregation::Aggregation::LastValue)
1382 .with_name("my_observable_counter_renamed")
1383 .build()
1384 .ok()
1385 } else {
1386 None
1387 }
1388 };
1389 let meter_provider = SdkMeterProvider::builder()
1390 .with_periodic_exporter(exporter.clone())
1391 .with_view(view)
1392 .build();
1393
1394 let meter = meter_provider.meter("test");
1396 let _observable_counter = meter
1397 .u64_observable_counter("my_observable_counter")
1398 .with_callback(|observer| {
1399 observer.observe(100, &[KeyValue::new("key1", "value1")]);
1400 })
1401 .build();
1402 meter_provider.force_flush().unwrap();
1403
1404 let resource_metrics = exporter
1406 .get_finished_metrics()
1407 .expect("metrics are expected to be exported.");
1408 assert!(!resource_metrics.is_empty());
1409 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1410 assert_eq!(
1412 metric.name, "my_observable_counter",
1413 "View rename should be ignored due to incompatible aggregation."
1414 );
1415 assert!(
1417 matches!(
1418 &metric.data,
1419 data::AggregatedMetrics::U64(data::MetricData::Sum(_))
1420 ),
1421 "Observable Counter should use default Sum aggregation when LastValue is incompatible."
1422 );
1423 }
1424
1425 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1426 async fn observable_updowncounter_with_lastvalue_aggregation_uses_default() {
1427 let exporter = InMemoryMetricExporter::default();
1434 let view = |i: &Instrument| {
1435 if i.name == "my_observable_updowncounter" {
1436 Stream::builder()
1437 .with_aggregation(aggregation::Aggregation::LastValue)
1438 .with_name("my_observable_updowncounter_renamed")
1439 .build()
1440 .ok()
1441 } else {
1442 None
1443 }
1444 };
1445 let meter_provider = SdkMeterProvider::builder()
1446 .with_periodic_exporter(exporter.clone())
1447 .with_view(view)
1448 .build();
1449
1450 let meter = meter_provider.meter("test");
1452 let _observable_updowncounter = meter
1453 .i64_observable_up_down_counter("my_observable_updowncounter")
1454 .with_callback(|observer| {
1455 observer.observe(-50, &[KeyValue::new("key1", "value1")]);
1456 })
1457 .build();
1458 meter_provider.force_flush().unwrap();
1459
1460 let resource_metrics = exporter
1462 .get_finished_metrics()
1463 .expect("metrics are expected to be exported.");
1464 assert!(!resource_metrics.is_empty());
1465 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1466 assert_eq!(
1468 metric.name, "my_observable_updowncounter",
1469 "View rename should be ignored due to incompatible aggregation."
1470 );
1471 assert!(
1473 matches!(
1474 &metric.data,
1475 data::AggregatedMetrics::I64(data::MetricData::Sum(_))
1476 ),
1477 "Observable UpDownCounter should use default Sum aggregation when LastValue is incompatible."
1478 );
1479 }
1480
1481 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1482 async fn histogram_with_sum_aggregation_is_valid() {
1483 let exporter = InMemoryMetricExporter::default();
1489 let view = |i: &Instrument| {
1490 if i.name == "my_histogram" {
1491 Stream::builder()
1492 .with_aggregation(aggregation::Aggregation::Sum)
1493 .with_name("my_histogram_renamed")
1494 .build()
1495 .ok()
1496 } else {
1497 None
1498 }
1499 };
1500 let meter_provider = SdkMeterProvider::builder()
1501 .with_periodic_exporter(exporter.clone())
1502 .with_view(view)
1503 .build();
1504
1505 let meter = meter_provider.meter("test");
1507 let histogram = meter.f64_histogram("my_histogram").build();
1508 histogram.record(10.0, &[KeyValue::new("key1", "value1")]);
1509 histogram.record(20.0, &[KeyValue::new("key1", "value1")]);
1510 histogram.record(30.0, &[KeyValue::new("key1", "value1")]);
1511 meter_provider.force_flush().unwrap();
1512
1513 let resource_metrics = exporter
1515 .get_finished_metrics()
1516 .expect("metrics are expected to be exported.");
1517 assert!(!resource_metrics.is_empty());
1518 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1519 assert_eq!(
1521 metric.name, "my_histogram_renamed",
1522 "View rename should be applied for compatible aggregation."
1523 );
1524 let MetricData::Sum(sum) = f64::extract_metrics_data_ref(&metric.data)
1526 .expect("Sum aggregation expected when view specifies Sum")
1527 else {
1528 panic!("Expected Sum aggregation for Histogram with Sum view");
1529 };
1530 assert_eq!(sum.data_points.len(), 1);
1531 assert_eq!(sum.data_points[0].value, 60.0); }
1533
1534 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1535 async fn gauge_with_histogram_aggregation_is_valid() {
1536 let exporter = InMemoryMetricExporter::default();
1542 let view = |i: &Instrument| {
1543 if i.name == "my_gauge" {
1544 Stream::builder()
1545 .with_aggregation(aggregation::Aggregation::ExplicitBucketHistogram {
1546 boundaries: vec![5.0, 10.0, 25.0, 50.0],
1547 record_min_max: true,
1548 })
1549 .with_name("my_gauge_renamed")
1550 .build()
1551 .ok()
1552 } else {
1553 None
1554 }
1555 };
1556 let meter_provider = SdkMeterProvider::builder()
1557 .with_periodic_exporter(exporter.clone())
1558 .with_view(view)
1559 .build();
1560
1561 let meter = meter_provider.meter("test");
1563 let gauge = meter.f64_gauge("my_gauge").build();
1564 gauge.record(3.0, &[KeyValue::new("key1", "value1")]);
1565 gauge.record(7.0, &[KeyValue::new("key1", "value1")]);
1566 gauge.record(15.0, &[KeyValue::new("key1", "value1")]);
1567 gauge.record(30.0, &[KeyValue::new("key1", "value1")]);
1568 meter_provider.force_flush().unwrap();
1569
1570 let resource_metrics = exporter
1572 .get_finished_metrics()
1573 .expect("metrics are expected to be exported.");
1574 assert!(!resource_metrics.is_empty());
1575 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1576 assert_eq!(
1578 metric.name, "my_gauge_renamed",
1579 "View rename should be applied for compatible aggregation."
1580 );
1581 let MetricData::Histogram(histogram) = f64::extract_metrics_data_ref(&metric.data)
1583 .expect("Histogram aggregation expected when view specifies Histogram")
1584 else {
1585 panic!("Expected Histogram aggregation for Gauge with Histogram view");
1586 };
1587 assert_eq!(histogram.data_points.len(), 1);
1588 let dp = &histogram.data_points[0];
1589 assert_eq!(dp.count, 4);
1590 assert_eq!(dp.sum, 55.0); assert_eq!(dp.min, Some(3.0));
1592 assert_eq!(dp.max, Some(30.0));
1593 assert_eq!(dp.bucket_counts, vec![1, 1, 1, 1, 0]);
1596 }
1597
1598 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1599 async fn counter_with_histogram_aggregation_is_valid() {
1600 let exporter = InMemoryMetricExporter::default();
1606 let view = |i: &Instrument| {
1607 if i.name == "my_counter" {
1608 Stream::builder()
1609 .with_aggregation(aggregation::Aggregation::ExplicitBucketHistogram {
1610 boundaries: vec![5.0, 10.0, 25.0, 50.0],
1611 record_min_max: true,
1612 })
1613 .with_name("my_counter_renamed")
1614 .build()
1615 .ok()
1616 } else {
1617 None
1618 }
1619 };
1620 let meter_provider = SdkMeterProvider::builder()
1621 .with_periodic_exporter(exporter.clone())
1622 .with_view(view)
1623 .build();
1624
1625 let meter = meter_provider.meter("test");
1627 let counter = meter.u64_counter("my_counter").build();
1628 counter.add(3, &[KeyValue::new("key1", "value1")]);
1629 counter.add(7, &[KeyValue::new("key1", "value1")]);
1630 counter.add(15, &[KeyValue::new("key1", "value1")]);
1631 counter.add(30, &[KeyValue::new("key1", "value1")]);
1632 meter_provider.force_flush().unwrap();
1633
1634 let resource_metrics = exporter
1636 .get_finished_metrics()
1637 .expect("metrics are expected to be exported.");
1638 assert!(!resource_metrics.is_empty());
1639 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1640 assert_eq!(
1642 metric.name, "my_counter_renamed",
1643 "View rename should be applied for compatible aggregation."
1644 );
1645 let MetricData::Histogram(histogram) = u64::extract_metrics_data_ref(&metric.data)
1647 .expect("Histogram aggregation expected when view specifies Histogram")
1648 else {
1649 panic!("Expected Histogram aggregation for Counter with Histogram view");
1650 };
1651 assert_eq!(histogram.data_points.len(), 1);
1652 let dp = &histogram.data_points[0];
1653 assert_eq!(dp.count, 4);
1654 assert_eq!(dp.sum, 55); assert_eq!(dp.min, Some(3));
1656 assert_eq!(dp.max, Some(30));
1657 assert_eq!(dp.bucket_counts, vec![1, 1, 1, 1, 0]);
1660 }
1661
1662 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1663 async fn updowncounter_with_histogram_aggregation_is_valid() {
1664 let exporter = InMemoryMetricExporter::default();
1670 let view = |i: &Instrument| {
1671 if i.name == "my_updowncounter" {
1672 Stream::builder()
1673 .with_aggregation(aggregation::Aggregation::ExplicitBucketHistogram {
1674 boundaries: vec![0.0, 10.0, 20.0, 50.0],
1675 record_min_max: true,
1676 })
1677 .with_name("my_updowncounter_renamed")
1678 .build()
1679 .ok()
1680 } else {
1681 None
1682 }
1683 };
1684 let meter_provider = SdkMeterProvider::builder()
1685 .with_periodic_exporter(exporter.clone())
1686 .with_view(view)
1687 .build();
1688
1689 let meter = meter_provider.meter("test");
1691 let updowncounter = meter.i64_up_down_counter("my_updowncounter").build();
1692 updowncounter.add(-5, &[KeyValue::new("key1", "value1")]);
1693 updowncounter.add(15, &[KeyValue::new("key1", "value1")]);
1694 updowncounter.add(25, &[KeyValue::new("key1", "value1")]);
1695 meter_provider.force_flush().unwrap();
1696
1697 let resource_metrics = exporter
1699 .get_finished_metrics()
1700 .expect("metrics are expected to be exported.");
1701 assert!(!resource_metrics.is_empty());
1702 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1703 assert_eq!(
1705 metric.name, "my_updowncounter_renamed",
1706 "View rename should be applied for compatible aggregation."
1707 );
1708 let MetricData::Histogram(histogram) = i64::extract_metrics_data_ref(&metric.data)
1710 .expect("Histogram aggregation expected when view specifies Histogram")
1711 else {
1712 panic!("Expected Histogram aggregation for UpDownCounter with Histogram view");
1713 };
1714 assert_eq!(histogram.data_points.len(), 1);
1715 let dp = &histogram.data_points[0];
1716 assert_eq!(dp.count, 3);
1717 }
1719
1720 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1721 async fn counter_with_exponential_histogram_aggregation_is_valid() {
1722 let exporter = InMemoryMetricExporter::default();
1728 let view = |i: &Instrument| {
1729 if i.name == "my_counter" {
1730 Stream::builder()
1731 .with_aggregation(aggregation::Aggregation::Base2ExponentialHistogram {
1732 max_size: 160,
1733 max_scale: 20,
1734 record_min_max: true,
1735 })
1736 .with_name("my_counter_renamed")
1737 .build()
1738 .ok()
1739 } else {
1740 None
1741 }
1742 };
1743 let meter_provider = SdkMeterProvider::builder()
1744 .with_periodic_exporter(exporter.clone())
1745 .with_view(view)
1746 .build();
1747
1748 let meter = meter_provider.meter("test");
1750 let counter = meter.u64_counter("my_counter").build();
1751 counter.add(5, &[KeyValue::new("key1", "value1")]);
1752 counter.add(10, &[KeyValue::new("key1", "value1")]);
1753 counter.add(20, &[KeyValue::new("key1", "value1")]);
1754 meter_provider.force_flush().unwrap();
1755
1756 let resource_metrics = exporter
1758 .get_finished_metrics()
1759 .expect("metrics are expected to be exported.");
1760 assert!(!resource_metrics.is_empty());
1761 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1762 assert_eq!(
1764 metric.name, "my_counter_renamed",
1765 "View rename should be applied for compatible aggregation."
1766 );
1767 let MetricData::ExponentialHistogram(exp_hist) =
1769 u64::extract_metrics_data_ref(&metric.data)
1770 .expect("ExponentialHistogram aggregation expected when view specifies it")
1771 else {
1772 panic!("Expected ExponentialHistogram aggregation for Counter with ExponentialHistogram view");
1773 };
1774 assert_eq!(exp_hist.data_points.len(), 1);
1775 let dp = &exp_hist.data_points[0];
1776 assert_eq!(dp.count(), 3);
1777 assert_eq!(dp.sum(), 35); assert_eq!(dp.min(), Some(5));
1779 assert_eq!(dp.max(), Some(20));
1780 }
1781
1782 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1783 async fn gauge_with_exponential_histogram_aggregation_is_valid() {
1784 let exporter = InMemoryMetricExporter::default();
1790 let view = |i: &Instrument| {
1791 if i.name == "my_gauge" {
1792 Stream::builder()
1793 .with_aggregation(aggregation::Aggregation::Base2ExponentialHistogram {
1794 max_size: 160,
1795 max_scale: 20,
1796 record_min_max: true,
1797 })
1798 .with_name("my_gauge_renamed")
1799 .build()
1800 .ok()
1801 } else {
1802 None
1803 }
1804 };
1805 let meter_provider = SdkMeterProvider::builder()
1806 .with_periodic_exporter(exporter.clone())
1807 .with_view(view)
1808 .build();
1809
1810 let meter = meter_provider.meter("test");
1812 let gauge = meter.f64_gauge("my_gauge").build();
1813 gauge.record(2.5, &[KeyValue::new("key1", "value1")]);
1814 gauge.record(7.5, &[KeyValue::new("key1", "value1")]);
1815 gauge.record(15.0, &[KeyValue::new("key1", "value1")]);
1816 meter_provider.force_flush().unwrap();
1817
1818 let resource_metrics = exporter
1820 .get_finished_metrics()
1821 .expect("metrics are expected to be exported.");
1822 assert!(!resource_metrics.is_empty());
1823 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
1824 assert_eq!(
1826 metric.name, "my_gauge_renamed",
1827 "View rename should be applied for compatible aggregation."
1828 );
1829 let MetricData::ExponentialHistogram(exp_hist) =
1831 f64::extract_metrics_data_ref(&metric.data)
1832 .expect("ExponentialHistogram aggregation expected when view specifies it")
1833 else {
1834 panic!("Expected ExponentialHistogram aggregation for Gauge with ExponentialHistogram view");
1835 };
1836 assert_eq!(exp_hist.data_points.len(), 1);
1837 let dp = &exp_hist.data_points[0];
1838 assert_eq!(dp.count(), 3);
1839 assert_eq!(dp.sum(), 25.0); assert_eq!(dp.min(), Some(2.5));
1841 assert_eq!(dp.max(), Some(15.0));
1842 }
1843
1844 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1845 async fn counter_with_drop_aggregation_is_dropped() {
1846 let exporter = InMemoryMetricExporter::default();
1854 let view = |i: &Instrument| {
1855 if i.name == "my_counter_to_drop" {
1856 Stream::builder()
1857 .with_aggregation(aggregation::Aggregation::Drop)
1858 .build()
1859 .ok()
1860 } else {
1861 None
1862 }
1863 };
1864 let meter_provider = SdkMeterProvider::builder()
1865 .with_periodic_exporter(exporter.clone())
1866 .with_view(view)
1867 .build();
1868
1869 let meter = meter_provider.meter("test");
1871 let counter = meter.u64_counter("my_counter_to_drop").build();
1872 counter.add(10, &[KeyValue::new("key1", "value1")]);
1873 meter_provider.force_flush().unwrap();
1874
1875 let resource_metrics = exporter
1877 .get_finished_metrics()
1878 .expect("metrics result expected");
1879 assert!(
1880 resource_metrics.is_empty()
1881 || resource_metrics[0].scope_metrics.is_empty()
1882 || resource_metrics[0].scope_metrics[0].metrics.is_empty(),
1883 "No metrics should be exported when view uses Aggregation::Drop. Got: {:?}",
1884 resource_metrics
1885 );
1886 }
1887
1888 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1889 async fn histogram_with_drop_aggregation_is_dropped() {
1890 let exporter = InMemoryMetricExporter::default();
1898 let view = |i: &Instrument| {
1899 if i.name == "my_histogram_to_drop" {
1900 Stream::builder()
1901 .with_aggregation(aggregation::Aggregation::Drop)
1902 .build()
1903 .ok()
1904 } else {
1905 None
1906 }
1907 };
1908 let meter_provider = SdkMeterProvider::builder()
1909 .with_periodic_exporter(exporter.clone())
1910 .with_view(view)
1911 .build();
1912
1913 let meter = meter_provider.meter("test");
1915 let histogram = meter.f64_histogram("my_histogram_to_drop").build();
1916 histogram.record(42.0, &[KeyValue::new("key1", "value1")]);
1917 meter_provider.force_flush().unwrap();
1918
1919 let resource_metrics = exporter
1921 .get_finished_metrics()
1922 .expect("metrics result expected");
1923 assert!(
1924 resource_metrics.is_empty()
1925 || resource_metrics[0].scope_metrics.is_empty()
1926 || resource_metrics[0].scope_metrics[0].metrics.is_empty(),
1927 "No metrics should be exported when view uses Aggregation::Drop. Got: {:?}",
1928 resource_metrics
1929 );
1930 }
1931
1932 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1933 async fn gauge_with_drop_aggregation_is_dropped() {
1934 let exporter = InMemoryMetricExporter::default();
1942 let view = |i: &Instrument| {
1943 if i.name == "my_gauge_to_drop" {
1944 Stream::builder()
1945 .with_aggregation(aggregation::Aggregation::Drop)
1946 .build()
1947 .ok()
1948 } else {
1949 None
1950 }
1951 };
1952 let meter_provider = SdkMeterProvider::builder()
1953 .with_periodic_exporter(exporter.clone())
1954 .with_view(view)
1955 .build();
1956
1957 let meter = meter_provider.meter("test");
1959 let gauge = meter.f64_gauge("my_gauge_to_drop").build();
1960 gauge.record(42.0, &[KeyValue::new("key1", "value1")]);
1961 meter_provider.force_flush().unwrap();
1962
1963 let resource_metrics = exporter
1965 .get_finished_metrics()
1966 .expect("metrics result expected");
1967 assert!(
1968 resource_metrics.is_empty()
1969 || resource_metrics[0].scope_metrics.is_empty()
1970 || resource_metrics[0].scope_metrics[0].metrics.is_empty(),
1971 "No metrics should be exported when view uses Aggregation::Drop. Got: {:?}",
1972 resource_metrics
1973 );
1974 }
1975
1976 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1977 async fn counter_with_drop_aggregation_and_rename_should_still_drop() {
1978 let exporter = InMemoryMetricExporter::default();
1988 let view = |i: &Instrument| {
1989 if i.name == "my_counter" {
1990 Stream::builder()
1991 .with_name("dropped_counter") .with_aggregation(aggregation::Aggregation::Drop)
1993 .build()
1994 .ok()
1995 } else {
1996 None
1997 }
1998 };
1999 let meter_provider = SdkMeterProvider::builder()
2000 .with_periodic_exporter(exporter.clone())
2001 .with_view(view)
2002 .build();
2003
2004 let meter = meter_provider.meter("test");
2006 let counter = meter.u64_counter("my_counter").build();
2007 counter.add(10, &[KeyValue::new("key1", "value1")]);
2008 meter_provider.force_flush().unwrap();
2009
2010 let resource_metrics = exporter
2012 .get_finished_metrics()
2013 .expect("metrics result expected");
2014 assert!(
2015 resource_metrics.is_empty()
2016 || resource_metrics[0].scope_metrics.is_empty()
2017 || resource_metrics[0].scope_metrics[0].metrics.is_empty(),
2018 "No metrics should be exported when view uses Aggregation::Drop, even with rename. Got: {:?}",
2019 resource_metrics
2020 );
2021 }
2022
2023 #[cfg(feature = "spec_unstable_metrics_views")]
2024 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2025 #[ignore = "Spatial aggregation is not yet implemented."]
2026 async fn spatial_aggregation_when_view_drops_attributes_observable_counter() {
2027 let exporter = InMemoryMetricExporter::default();
2031 let view = |i: &Instrument| {
2033 if i.name == "my_observable_counter" {
2034 Stream::builder()
2035 .with_allowed_attribute_keys(vec![])
2036 .build()
2037 .ok()
2038 } else {
2039 None
2040 }
2041 };
2042 let meter_provider = SdkMeterProvider::builder()
2043 .with_periodic_exporter(exporter.clone())
2044 .with_view(view)
2045 .build();
2046
2047 let meter = meter_provider.meter("test");
2049 let _observable_counter = meter
2050 .u64_observable_counter("my_observable_counter")
2051 .with_callback(|observer| {
2052 observer.observe(
2053 100,
2054 &[
2055 KeyValue::new("statusCode", "200"),
2056 KeyValue::new("verb", "get"),
2057 ],
2058 );
2059
2060 observer.observe(
2061 100,
2062 &[
2063 KeyValue::new("statusCode", "200"),
2064 KeyValue::new("verb", "post"),
2065 ],
2066 );
2067
2068 observer.observe(
2069 100,
2070 &[
2071 KeyValue::new("statusCode", "500"),
2072 KeyValue::new("verb", "get"),
2073 ],
2074 );
2075 })
2076 .build();
2077
2078 meter_provider.force_flush().unwrap();
2079
2080 let resource_metrics = exporter
2082 .get_finished_metrics()
2083 .expect("metrics are expected to be exported.");
2084 assert!(!resource_metrics.is_empty());
2085 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
2086 assert_eq!(metric.name, "my_observable_counter",);
2087
2088 let MetricData::Sum(sum) = u64::extract_metrics_data_ref(&metric.data)
2089 .expect("Sum aggregation expected for ObservableCounter instruments by default")
2090 else {
2091 unreachable!()
2092 };
2093
2094 assert_eq!(sum.data_points.len(), 1);
2098
2099 let data_point = &sum.data_points[0];
2101 assert_eq!(data_point.value, 300);
2102 }
2103
2104 #[cfg(feature = "spec_unstable_metrics_views")]
2105 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2106 async fn spatial_aggregation_when_view_drops_attributes_counter() {
2107 let exporter = InMemoryMetricExporter::default();
2111 let view = |i: &Instrument| {
2113 if i.name == "my_counter" {
2114 Some(
2115 Stream::builder()
2116 .with_allowed_attribute_keys(vec![])
2117 .build()
2118 .unwrap(),
2119 )
2120 } else {
2121 None
2122 }
2123 };
2124 let meter_provider = SdkMeterProvider::builder()
2125 .with_periodic_exporter(exporter.clone())
2126 .with_view(view)
2127 .build();
2128
2129 let meter = meter_provider.meter("test");
2131 let counter = meter.u64_counter("my_counter").build();
2132
2133 counter.add(
2136 10,
2137 [
2138 KeyValue::new("statusCode", "200"),
2139 KeyValue::new("verb", "Get"),
2140 ]
2141 .as_ref(),
2142 );
2143
2144 counter.add(
2145 10,
2146 [
2147 KeyValue::new("statusCode", "500"),
2148 KeyValue::new("verb", "Get"),
2149 ]
2150 .as_ref(),
2151 );
2152
2153 counter.add(
2154 10,
2155 [
2156 KeyValue::new("statusCode", "200"),
2157 KeyValue::new("verb", "Post"),
2158 ]
2159 .as_ref(),
2160 );
2161
2162 meter_provider.force_flush().unwrap();
2163
2164 let resource_metrics = exporter
2166 .get_finished_metrics()
2167 .expect("metrics are expected to be exported.");
2168 assert!(!resource_metrics.is_empty());
2169 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
2170 assert_eq!(metric.name, "my_counter",);
2171
2172 let MetricData::Sum(sum) = u64::extract_metrics_data_ref(&metric.data)
2173 .expect("Sum aggregation expected for Counter instruments by default")
2174 else {
2175 unreachable!()
2176 };
2177
2178 assert_eq!(sum.data_points.len(), 1);
2182 let data_point = &sum.data_points[0];
2184 assert_eq!(data_point.value, 30);
2185 }
2186
2187 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2188 async fn no_attr_cumulative_up_down_counter() {
2189 let mut test_context = TestContext::new(Temporality::Cumulative);
2190 let counter = test_context.i64_up_down_counter("test", "my_counter", Some("my_unit"));
2191
2192 counter.add(50, &[]);
2193 test_context.flush_metrics();
2194
2195 let MetricData::Sum(sum) =
2196 test_context.get_aggregation::<i64>("my_counter", Some("my_unit"))
2197 else {
2198 unreachable!()
2199 };
2200
2201 assert_eq!(sum.data_points.len(), 1, "Expected only one data point");
2202 assert!(!sum.is_monotonic, "Should not produce monotonic.");
2203 assert_eq!(
2204 sum.temporality,
2205 Temporality::Cumulative,
2206 "Should produce cumulative"
2207 );
2208
2209 let data_point = &sum.data_points[0];
2210 assert!(data_point.attributes.is_empty(), "Non-empty attribute set");
2211 assert_eq!(data_point.value, 50, "Unexpected data point value");
2212 }
2213
2214 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2215 async fn no_attr_up_down_counter_always_cumulative() {
2216 let mut test_context = TestContext::new(Temporality::Delta);
2217 let counter = test_context.i64_up_down_counter("test", "my_counter", Some("my_unit"));
2218
2219 counter.add(50, &[]);
2220 test_context.flush_metrics();
2221
2222 let MetricData::Sum(sum) =
2223 test_context.get_aggregation::<i64>("my_counter", Some("my_unit"))
2224 else {
2225 unreachable!()
2226 };
2227
2228 assert_eq!(sum.data_points.len(), 1, "Expected only one data point");
2229 assert!(!sum.is_monotonic, "Should not produce monotonic.");
2230 assert_eq!(
2231 sum.temporality,
2232 Temporality::Cumulative,
2233 "Should produce Cumulative due to UpDownCounter temporality_preference"
2234 );
2235
2236 let data_point = &sum.data_points[0];
2237 assert!(data_point.attributes.is_empty(), "Non-empty attribute set");
2238 assert_eq!(data_point.value, 50, "Unexpected data point value");
2239 }
2240
2241 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2242 async fn no_attr_cumulative_counter_value_added_after_export() {
2243 let mut test_context = TestContext::new(Temporality::Cumulative);
2244 let counter = test_context.u64_counter("test", "my_counter", None);
2245
2246 counter.add(50, &[]);
2247 test_context.flush_metrics();
2248 let _ = test_context.get_aggregation::<u64>("my_counter", None);
2249 test_context.reset_metrics();
2250
2251 counter.add(5, &[]);
2252 test_context.flush_metrics();
2253 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
2254 unreachable!()
2255 };
2256
2257 assert_eq!(sum.data_points.len(), 1, "Expected only one data point");
2258 assert!(sum.is_monotonic, "Should produce monotonic.");
2259 assert_eq!(
2260 sum.temporality,
2261 Temporality::Cumulative,
2262 "Should produce cumulative"
2263 );
2264
2265 let data_point = &sum.data_points[0];
2266 assert!(data_point.attributes.is_empty(), "Non-empty attribute set");
2267 assert_eq!(data_point.value, 55, "Unexpected data point value");
2268 }
2269
2270 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2271 async fn no_attr_delta_counter_value_reset_after_export() {
2272 let mut test_context = TestContext::new(Temporality::Delta);
2273 let counter = test_context.u64_counter("test", "my_counter", None);
2274
2275 counter.add(50, &[]);
2276 test_context.flush_metrics();
2277 let _ = test_context.get_aggregation::<u64>("my_counter", None);
2278 test_context.reset_metrics();
2279
2280 counter.add(5, &[]);
2281 test_context.flush_metrics();
2282 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
2283 unreachable!()
2284 };
2285
2286 assert_eq!(sum.data_points.len(), 1, "Expected only one data point");
2287 assert!(sum.is_monotonic, "Should produce monotonic.");
2288 assert_eq!(
2289 sum.temporality,
2290 Temporality::Delta,
2291 "Should produce cumulative"
2292 );
2293
2294 let data_point = &sum.data_points[0];
2295 assert!(data_point.attributes.is_empty(), "Non-empty attribute set");
2296 assert_eq!(data_point.value, 5, "Unexpected data point value");
2297 }
2298
2299 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2300 async fn second_delta_export_does_not_give_no_attr_value_if_add_not_called() {
2301 let mut test_context = TestContext::new(Temporality::Delta);
2302 let counter = test_context.u64_counter("test", "my_counter", None);
2303
2304 counter.add(50, &[]);
2305 test_context.flush_metrics();
2306 let _ = test_context.get_aggregation::<u64>("my_counter", None);
2307 test_context.reset_metrics();
2308
2309 counter.add(50, &[KeyValue::new("a", "b")]);
2310 test_context.flush_metrics();
2311 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
2312 unreachable!()
2313 };
2314
2315 let no_attr_data_point = sum.data_points.iter().find(|x| x.attributes.is_empty());
2316
2317 assert!(
2318 no_attr_data_point.is_none(),
2319 "Expected no data points with no attributes"
2320 );
2321 }
2322
2323 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2324 async fn delta_memory_efficiency_test() {
2325 let mut test_context = TestContext::new(Temporality::Delta);
2330 let counter = test_context.u64_counter("test", "my_counter", None);
2331
2332 counter.add(1, &[KeyValue::new("key1", "value1")]);
2334 counter.add(1, &[KeyValue::new("key1", "value1")]);
2335 counter.add(1, &[KeyValue::new("key1", "value1")]);
2336 counter.add(1, &[KeyValue::new("key1", "value1")]);
2337 counter.add(1, &[KeyValue::new("key1", "value1")]);
2338
2339 counter.add(1, &[KeyValue::new("key1", "value2")]);
2340 counter.add(1, &[KeyValue::new("key1", "value2")]);
2341 counter.add(1, &[KeyValue::new("key1", "value2")]);
2342 test_context.flush_metrics();
2343
2344 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
2345 unreachable!()
2346 };
2347
2348 assert_eq!(sum.data_points.len(), 2);
2350
2351 let data_point1 = find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value1")
2353 .expect("datapoint with key1=value1 expected");
2354 assert_eq!(data_point1.value, 5);
2355
2356 let data_point1 = find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value2")
2358 .expect("datapoint with key1=value2 expected");
2359 assert_eq!(data_point1.value, 3);
2360
2361 test_context.exporter.reset();
2362 test_context.flush_metrics();
2365
2366 let resource_metrics = test_context
2367 .exporter
2368 .get_finished_metrics()
2369 .expect("metrics are expected to be exported.");
2370 println!("resource_metrics: {resource_metrics:?}");
2371 assert!(resource_metrics.is_empty(), "No metrics should be exported as no new measurements were recorded since last collect.");
2372 }
2373
2374 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2375 async fn counter_multithreaded() {
2376 counter_multithreaded_aggregation_helper(Temporality::Delta);
2380 counter_multithreaded_aggregation_helper(Temporality::Cumulative);
2381 }
2382
2383 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2384 async fn counter_f64_multithreaded() {
2385 counter_f64_multithreaded_aggregation_helper(Temporality::Delta);
2389 counter_f64_multithreaded_aggregation_helper(Temporality::Cumulative);
2390 }
2391
2392 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2393 async fn histogram_multithreaded() {
2394 histogram_multithreaded_aggregation_helper(Temporality::Delta);
2398 histogram_multithreaded_aggregation_helper(Temporality::Cumulative);
2399 }
2400
2401 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2402 async fn histogram_f64_multithreaded() {
2403 histogram_f64_multithreaded_aggregation_helper(Temporality::Delta);
2407 histogram_f64_multithreaded_aggregation_helper(Temporality::Cumulative);
2408 }
2409 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2410 async fn synchronous_instruments_cumulative_with_gap_in_measurements() {
2411 synchronous_instruments_cumulative_with_gap_in_measurements_helper("counter");
2415 synchronous_instruments_cumulative_with_gap_in_measurements_helper("updown_counter");
2416 synchronous_instruments_cumulative_with_gap_in_measurements_helper("histogram");
2417 synchronous_instruments_cumulative_with_gap_in_measurements_helper("gauge");
2418 }
2419
2420 fn synchronous_instruments_cumulative_with_gap_in_measurements_helper(
2421 instrument_name: &'static str,
2422 ) {
2423 let mut test_context = TestContext::new(Temporality::Cumulative);
2424 let attributes = &[KeyValue::new("key1", "value1")];
2425
2426 match instrument_name {
2428 "counter" => {
2429 let counter = test_context.meter().u64_counter("test_counter").build();
2430 counter.add(5, &[]);
2431 counter.add(10, attributes);
2432 }
2433 "updown_counter" => {
2434 let updown_counter = test_context
2435 .meter()
2436 .i64_up_down_counter("test_updowncounter")
2437 .build();
2438 updown_counter.add(15, &[]);
2439 updown_counter.add(20, attributes);
2440 }
2441 "histogram" => {
2442 let histogram = test_context.meter().u64_histogram("test_histogram").build();
2443 histogram.record(25, &[]);
2444 histogram.record(30, attributes);
2445 }
2446 "gauge" => {
2447 let gauge = test_context.meter().u64_gauge("test_gauge").build();
2448 gauge.record(35, &[]);
2449 gauge.record(40, attributes);
2450 }
2451 _ => panic!("Incorrect instrument kind provided"),
2452 };
2453
2454 test_context.flush_metrics();
2455
2456 assert_correct_export(&mut test_context, instrument_name);
2458
2459 test_context.reset_metrics();
2461
2462 test_context.flush_metrics();
2463
2464 assert_correct_export(&mut test_context, instrument_name);
2466
2467 fn assert_correct_export(test_context: &mut TestContext, instrument_name: &'static str) {
2468 match instrument_name {
2469 "counter" => {
2470 let MetricData::Sum(sum) =
2471 test_context.get_aggregation::<u64>("test_counter", None)
2472 else {
2473 unreachable!()
2474 };
2475 assert_eq!(sum.data_points.len(), 2);
2476 let zero_attribute_datapoint =
2477 find_sum_datapoint_with_no_attributes(&sum.data_points)
2478 .expect("datapoint with no attributes expected");
2479 assert_eq!(zero_attribute_datapoint.value, 5);
2480 let data_point1 =
2481 find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value1")
2482 .expect("datapoint with key1=value1 expected");
2483 assert_eq!(data_point1.value, 10);
2484 }
2485 "updown_counter" => {
2486 let MetricData::Sum(sum) =
2487 test_context.get_aggregation::<i64>("test_updowncounter", None)
2488 else {
2489 unreachable!()
2490 };
2491 assert_eq!(sum.data_points.len(), 2);
2492 let zero_attribute_datapoint =
2493 find_sum_datapoint_with_no_attributes(&sum.data_points)
2494 .expect("datapoint with no attributes expected");
2495 assert_eq!(zero_attribute_datapoint.value, 15);
2496 let data_point1 =
2497 find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value1")
2498 .expect("datapoint with key1=value1 expected");
2499 assert_eq!(data_point1.value, 20);
2500 }
2501 "histogram" => {
2502 let MetricData::Histogram(histogram_data) =
2503 test_context.get_aggregation::<u64>("test_histogram", None)
2504 else {
2505 unreachable!()
2506 };
2507 assert_eq!(histogram_data.data_points.len(), 2);
2508 let zero_attribute_datapoint =
2509 find_histogram_datapoint_with_no_attributes(&histogram_data.data_points)
2510 .expect("datapoint with no attributes expected");
2511 assert_eq!(zero_attribute_datapoint.count, 1);
2512 assert_eq!(zero_attribute_datapoint.sum, 25);
2513 assert_eq!(zero_attribute_datapoint.min, Some(25));
2514 assert_eq!(zero_attribute_datapoint.max, Some(25));
2515 let data_point1 = find_histogram_datapoint_with_key_value(
2516 &histogram_data.data_points,
2517 "key1",
2518 "value1",
2519 )
2520 .expect("datapoint with key1=value1 expected");
2521 assert_eq!(data_point1.count, 1);
2522 assert_eq!(data_point1.sum, 30);
2523 assert_eq!(data_point1.min, Some(30));
2524 assert_eq!(data_point1.max, Some(30));
2525 }
2526 "gauge" => {
2527 let MetricData::Gauge(gauge_data) =
2528 test_context.get_aggregation::<u64>("test_gauge", None)
2529 else {
2530 unreachable!()
2531 };
2532 assert_eq!(gauge_data.data_points.len(), 2);
2533 let zero_attribute_datapoint =
2534 find_gauge_datapoint_with_no_attributes(&gauge_data.data_points)
2535 .expect("datapoint with no attributes expected");
2536 assert_eq!(zero_attribute_datapoint.value, 35);
2537 let data_point1 = find_gauge_datapoint_with_key_value(
2538 &gauge_data.data_points,
2539 "key1",
2540 "value1",
2541 )
2542 .expect("datapoint with key1=value1 expected");
2543 assert_eq!(data_point1.value, 40);
2544 }
2545 _ => panic!("Incorrect instrument kind provided"),
2546 }
2547 }
2548 }
2549
2550 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
2551 async fn asynchronous_instruments_cumulative_data_points_only_from_last_measurement() {
2552 asynchronous_instruments_cumulative_data_points_only_from_last_measurement_helper("gauge");
2556 asynchronous_instruments_cumulative_data_points_only_from_last_measurement_helper(
2557 "counter",
2558 );
2559 asynchronous_instruments_cumulative_data_points_only_from_last_measurement_helper(
2560 "updown_counter",
2561 );
2562 }
2563
2564 #[test]
2565 fn view_test_rename() {
2566 test_view_customization(
2567 |i| {
2568 if i.name == "my_counter" {
2569 Some(
2570 Stream::builder()
2571 .with_name("my_counter_renamed")
2572 .build()
2573 .unwrap(),
2574 )
2575 } else {
2576 None
2577 }
2578 },
2579 "my_counter_renamed",
2580 "my_unit",
2581 "my_description",
2582 )
2583 }
2584
2585 #[test]
2586 fn view_test_change_unit() {
2587 test_view_customization(
2588 |i| {
2589 if i.name == "my_counter" {
2590 Some(Stream::builder().with_unit("my_unit_new").build().unwrap())
2591 } else {
2592 None
2593 }
2594 },
2595 "my_counter",
2596 "my_unit_new",
2597 "my_description",
2598 )
2599 }
2600
2601 #[test]
2602 fn view_test_change_description() {
2603 test_view_customization(
2604 |i| {
2605 if i.name == "my_counter" {
2606 Some(
2607 Stream::builder()
2608 .with_description("my_description_new")
2609 .build()
2610 .unwrap(),
2611 )
2612 } else {
2613 None
2614 }
2615 },
2616 "my_counter",
2617 "my_unit",
2618 "my_description_new",
2619 )
2620 }
2621
2622 #[test]
2623 fn view_test_change_name_unit() {
2624 test_view_customization(
2625 |i| {
2626 if i.name == "my_counter" {
2627 Some(
2628 Stream::builder()
2629 .with_name("my_counter_renamed")
2630 .with_unit("my_unit_new")
2631 .build()
2632 .unwrap(),
2633 )
2634 } else {
2635 None
2636 }
2637 },
2638 "my_counter_renamed",
2639 "my_unit_new",
2640 "my_description",
2641 )
2642 }
2643
2644 #[test]
2645 fn view_test_change_name_unit_desc() {
2646 test_view_customization(
2647 |i| {
2648 if i.name == "my_counter" {
2649 Some(
2650 Stream::builder()
2651 .with_name("my_counter_renamed")
2652 .with_unit("my_unit_new")
2653 .with_description("my_description_new")
2654 .build()
2655 .unwrap(),
2656 )
2657 } else {
2658 None
2659 }
2660 },
2661 "my_counter_renamed",
2662 "my_unit_new",
2663 "my_description_new",
2664 )
2665 }
2666
2667 #[test]
2668 fn view_test_match_unit() {
2669 test_view_customization(
2670 |i| {
2671 if i.unit == "my_unit" {
2672 Some(Stream::builder().with_unit("my_unit_new").build().unwrap())
2673 } else {
2674 None
2675 }
2676 },
2677 "my_counter",
2678 "my_unit_new",
2679 "my_description",
2680 )
2681 }
2682
2683 #[test]
2684 fn view_test_match_none() {
2685 test_view_customization(
2686 |i| {
2687 if i.name == "not_expected_to_match" {
2688 Some(Stream::builder().build().unwrap())
2689 } else {
2690 None
2691 }
2692 },
2693 "my_counter",
2694 "my_unit",
2695 "my_description",
2696 )
2697 }
2698
2699 #[test]
2700 fn view_test_match_multiple() {
2701 test_view_customization(
2702 |i| {
2703 if i.name == "my_counter" && i.unit == "my_unit" {
2704 Some(
2705 Stream::builder()
2706 .with_name("my_counter_renamed")
2707 .build()
2708 .unwrap(),
2709 )
2710 } else {
2711 None
2712 }
2713 },
2714 "my_counter_renamed",
2715 "my_unit",
2716 "my_description",
2717 )
2718 }
2719
2720 fn test_view_customization<F>(
2722 view_function: F,
2723 expected_name: &str,
2724 expected_unit: &str,
2725 expected_description: &str,
2726 ) where
2727 F: Fn(&Instrument) -> Option<Stream> + Send + Sync + 'static,
2728 {
2729 let exporter = InMemoryMetricExporter::default();
2734 let meter_provider = SdkMeterProvider::builder()
2735 .with_periodic_exporter(exporter.clone())
2736 .with_view(view_function)
2737 .build();
2738
2739 let meter = meter_provider.meter("test");
2741 let counter = meter
2742 .f64_counter("my_counter")
2743 .with_unit("my_unit")
2744 .with_description("my_description")
2745 .build();
2746
2747 counter.add(1.5, &[KeyValue::new("key1", "value1")]);
2748 meter_provider.force_flush().unwrap();
2749
2750 let resource_metrics = exporter
2752 .get_finished_metrics()
2753 .expect("metrics are expected to be exported.");
2754 assert_eq!(resource_metrics.len(), 1);
2755 assert_eq!(resource_metrics[0].scope_metrics.len(), 1);
2756 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
2757 assert_eq!(
2758 metric.name, expected_name,
2759 "Expected name: {expected_name}."
2760 );
2761 assert_eq!(
2762 metric.unit, expected_unit,
2763 "Expected unit: {expected_unit}."
2764 );
2765 assert_eq!(
2766 metric.description, expected_description,
2767 "Expected description: {expected_description}."
2768 );
2769 }
2770
2771 #[test]
2778 fn test_view_single_instrument_multiple_stream() {
2779 let view1 = |i: &Instrument| {
2787 if i.name() == "my_counter" {
2788 Some(Stream::builder().with_name("my_counter_1").build().unwrap())
2789 } else {
2790 None
2791 }
2792 };
2793
2794 let view2 = |i: &Instrument| {
2795 if i.name() == "my_counter" {
2796 Some(Stream::builder().with_name("my_counter_2").build().unwrap())
2797 } else {
2798 None
2799 }
2800 };
2801
2802 let exporter = InMemoryMetricExporter::default();
2804 let meter_provider = SdkMeterProvider::builder()
2805 .with_periodic_exporter(exporter.clone())
2806 .with_view(view1)
2807 .with_view(view2)
2808 .build();
2809
2810 let meter = meter_provider.meter("test");
2812 let counter = meter.f64_counter("my_counter").build();
2813
2814 counter.add(1.5, &[KeyValue::new("key1", "value1")]);
2815 meter_provider.force_flush().unwrap();
2816
2817 let resource_metrics = exporter
2819 .get_finished_metrics()
2820 .expect("metrics are expected to be exported.");
2821 assert_eq!(resource_metrics.len(), 1);
2822 assert_eq!(resource_metrics[0].scope_metrics.len(), 1);
2823 let metrics = &resource_metrics[0].scope_metrics[0].metrics;
2824 assert_eq!(metrics.len(), 2);
2825 assert_eq!(metrics[0].name, "my_counter_1");
2826 assert_eq!(metrics[1].name, "my_counter_2");
2827 }
2828
2829 #[test]
2830 fn test_view_multiple_instrument_single_stream() {
2831 let view = |i: &Instrument| {
2838 if i.name() == "my_counter1" || i.name() == "my_counter2" {
2839 Some(Stream::builder().with_name("my_counter").build().unwrap())
2840 } else {
2841 None
2842 }
2843 };
2844
2845 let exporter = InMemoryMetricExporter::default();
2847 let meter_provider = SdkMeterProvider::builder()
2848 .with_periodic_exporter(exporter.clone())
2849 .with_view(view)
2850 .build();
2851
2852 let meter = meter_provider.meter("test");
2854 let counter1 = meter.f64_counter("my_counter1").build();
2855 let counter2 = meter.f64_counter("my_counter2").build();
2856
2857 counter1.add(1.5, &[KeyValue::new("key1", "value1")]);
2858 counter2.add(1.5, &[KeyValue::new("key1", "value1")]);
2859 meter_provider.force_flush().unwrap();
2860
2861 let resource_metrics = exporter
2863 .get_finished_metrics()
2864 .expect("metrics are expected to be exported.");
2865 assert_eq!(resource_metrics.len(), 1);
2866 assert_eq!(resource_metrics[0].scope_metrics.len(), 1);
2867 let metrics = &resource_metrics[0].scope_metrics[0].metrics;
2868 assert_eq!(metrics.len(), 1);
2869 assert_eq!(metrics[0].name, "my_counter");
2870 }
2872
2873 fn asynchronous_instruments_cumulative_data_points_only_from_last_measurement_helper(
2874 instrument_name: &'static str,
2875 ) {
2876 let mut test_context = TestContext::new(Temporality::Cumulative);
2877 let attributes = Arc::new([KeyValue::new("key1", "value1")]);
2878
2879 match instrument_name {
2881 "counter" => {
2882 let has_run = AtomicBool::new(false);
2883 let _observable_counter = test_context
2884 .meter()
2885 .u64_observable_counter("test_counter")
2886 .with_callback(move |observer| {
2887 if !has_run.load(Ordering::SeqCst) {
2888 observer.observe(5, &[]);
2889 observer.observe(10, &*attributes.clone());
2890 has_run.store(true, Ordering::SeqCst);
2891 }
2892 })
2893 .build();
2894 }
2895 "updown_counter" => {
2896 let has_run = AtomicBool::new(false);
2897 let _observable_up_down_counter = test_context
2898 .meter()
2899 .i64_observable_up_down_counter("test_updowncounter")
2900 .with_callback(move |observer| {
2901 if !has_run.load(Ordering::SeqCst) {
2902 observer.observe(15, &[]);
2903 observer.observe(20, &*attributes.clone());
2904 has_run.store(true, Ordering::SeqCst);
2905 }
2906 })
2907 .build();
2908 }
2909 "gauge" => {
2910 let has_run = AtomicBool::new(false);
2911 let _observable_gauge = test_context
2912 .meter()
2913 .u64_observable_gauge("test_gauge")
2914 .with_callback(move |observer| {
2915 if !has_run.load(Ordering::SeqCst) {
2916 observer.observe(25, &[]);
2917 observer.observe(30, &*attributes.clone());
2918 has_run.store(true, Ordering::SeqCst);
2919 }
2920 })
2921 .build();
2922 }
2923 _ => panic!("Incorrect instrument kind provided"),
2924 };
2925
2926 test_context.flush_metrics();
2927
2928 assert_correct_export(&mut test_context, instrument_name);
2930
2931 test_context.reset_metrics();
2933
2934 test_context.flush_metrics();
2935
2936 test_context.check_no_metrics();
2937
2938 fn assert_correct_export(test_context: &mut TestContext, instrument_name: &'static str) {
2939 match instrument_name {
2940 "counter" => {
2941 let MetricData::Sum(sum) =
2942 test_context.get_aggregation::<u64>("test_counter", None)
2943 else {
2944 unreachable!()
2945 };
2946 assert_eq!(sum.data_points.len(), 2);
2947 assert!(sum.is_monotonic);
2948 let zero_attribute_datapoint =
2949 find_sum_datapoint_with_no_attributes(&sum.data_points)
2950 .expect("datapoint with no attributes expected");
2951 assert_eq!(zero_attribute_datapoint.value, 5);
2952 let data_point1 =
2953 find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value1")
2954 .expect("datapoint with key1=value1 expected");
2955 assert_eq!(data_point1.value, 10);
2956 }
2957 "updown_counter" => {
2958 let MetricData::Sum(sum) =
2959 test_context.get_aggregation::<i64>("test_updowncounter", None)
2960 else {
2961 unreachable!()
2962 };
2963 assert_eq!(sum.data_points.len(), 2);
2964 assert!(!sum.is_monotonic);
2965 let zero_attribute_datapoint =
2966 find_sum_datapoint_with_no_attributes(&sum.data_points)
2967 .expect("datapoint with no attributes expected");
2968 assert_eq!(zero_attribute_datapoint.value, 15);
2969 let data_point1 =
2970 find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value1")
2971 .expect("datapoint with key1=value1 expected");
2972 assert_eq!(data_point1.value, 20);
2973 }
2974 "gauge" => {
2975 let MetricData::Gauge(gauge_data) =
2976 test_context.get_aggregation::<u64>("test_gauge", None)
2977 else {
2978 unreachable!()
2979 };
2980 assert_eq!(gauge_data.data_points.len(), 2);
2981 let zero_attribute_datapoint =
2982 find_gauge_datapoint_with_no_attributes(&gauge_data.data_points)
2983 .expect("datapoint with no attributes expected");
2984 assert_eq!(zero_attribute_datapoint.value, 25);
2985 let data_point1 = find_gauge_datapoint_with_key_value(
2986 &gauge_data.data_points,
2987 "key1",
2988 "value1",
2989 )
2990 .expect("datapoint with key1=value1 expected");
2991 assert_eq!(data_point1.value, 30);
2992 }
2993 _ => panic!("Incorrect instrument kind provided"),
2994 }
2995 }
2996 }
2997
2998 fn counter_multithreaded_aggregation_helper(temporality: Temporality) {
2999 let mut test_context = TestContext::new(temporality);
3001 let counter = Arc::new(test_context.u64_counter("test", "my_counter", None));
3002
3003 for i in 0..10 {
3004 thread::scope(|s| {
3005 s.spawn(|| {
3006 counter.add(1, &[]);
3007
3008 counter.add(1, &[KeyValue::new("key1", "value1")]);
3009 counter.add(1, &[KeyValue::new("key1", "value1")]);
3010 counter.add(1, &[KeyValue::new("key1", "value1")]);
3011
3012 if i % 2 == 0 {
3014 test_context.flush_metrics();
3015 thread::sleep(Duration::from_millis(i)); }
3017
3018 counter.add(1, &[KeyValue::new("key1", "value1")]);
3019 counter.add(1, &[KeyValue::new("key1", "value1")]);
3020 });
3021 });
3022 }
3023
3024 test_context.flush_metrics();
3025
3026 let sums = test_context
3029 .get_from_multiple_aggregations::<u64>("my_counter", None, 6)
3030 .into_iter()
3031 .map(|data| {
3032 if let MetricData::Sum(sum) = data {
3033 sum
3034 } else {
3035 unreachable!()
3036 }
3037 })
3038 .collect::<Vec<_>>();
3039
3040 let mut sum_zero_attributes = 0;
3041 let mut sum_key1_value1 = 0;
3042 sums.iter().for_each(|sum| {
3043 assert_eq!(sum.data_points.len(), 2); assert!(sum.is_monotonic, "Counter should produce monotonic.");
3045 assert_eq!(sum.temporality, temporality);
3046
3047 if temporality == Temporality::Delta {
3048 sum_zero_attributes += sum.data_points[0].value;
3049 sum_key1_value1 += sum.data_points[1].value;
3050 } else {
3051 sum_zero_attributes = sum.data_points[0].value;
3052 sum_key1_value1 = sum.data_points[1].value;
3053 };
3054 });
3055
3056 assert_eq!(sum_zero_attributes, 10);
3057 assert_eq!(sum_key1_value1, 50); }
3059
3060 fn counter_f64_multithreaded_aggregation_helper(temporality: Temporality) {
3061 let mut test_context = TestContext::new(temporality);
3063 let counter = Arc::new(test_context.meter().f64_counter("test_counter").build());
3064
3065 for i in 0..10 {
3066 thread::scope(|s| {
3067 s.spawn(|| {
3068 counter.add(1.23, &[]);
3069
3070 counter.add(1.23, &[KeyValue::new("key1", "value1")]);
3071 counter.add(1.23, &[KeyValue::new("key1", "value1")]);
3072 counter.add(1.23, &[KeyValue::new("key1", "value1")]);
3073
3074 if i % 2 == 0 {
3076 test_context.flush_metrics();
3077 thread::sleep(Duration::from_millis(i)); }
3079
3080 counter.add(1.23, &[KeyValue::new("key1", "value1")]);
3081 counter.add(1.23, &[KeyValue::new("key1", "value1")]);
3082 });
3083 });
3084 }
3085
3086 test_context.flush_metrics();
3087
3088 let sums = test_context
3091 .get_from_multiple_aggregations::<f64>("test_counter", None, 6)
3092 .into_iter()
3093 .map(|data| {
3094 if let MetricData::Sum(sum) = data {
3095 sum
3096 } else {
3097 unreachable!()
3098 }
3099 })
3100 .collect::<Vec<_>>();
3101
3102 let mut sum_zero_attributes = 0.0;
3103 let mut sum_key1_value1 = 0.0;
3104 sums.iter().for_each(|sum| {
3105 assert_eq!(sum.data_points.len(), 2); assert!(sum.is_monotonic, "Counter should produce monotonic.");
3107 assert_eq!(sum.temporality, temporality);
3108
3109 if temporality == Temporality::Delta {
3110 sum_zero_attributes += sum.data_points[0].value;
3111 sum_key1_value1 += sum.data_points[1].value;
3112 } else {
3113 sum_zero_attributes = sum.data_points[0].value;
3114 sum_key1_value1 = sum.data_points[1].value;
3115 };
3116 });
3117
3118 assert!(f64::abs(12.3 - sum_zero_attributes) < 0.0001);
3119 assert!(f64::abs(61.5 - sum_key1_value1) < 0.0001); }
3121
3122 fn histogram_multithreaded_aggregation_helper(temporality: Temporality) {
3123 let mut test_context = TestContext::new(temporality);
3125 let histogram = Arc::new(test_context.meter().u64_histogram("test_histogram").build());
3126
3127 for i in 0..10 {
3128 thread::scope(|s| {
3129 s.spawn(|| {
3130 histogram.record(1, &[]);
3131 histogram.record(4, &[]);
3132
3133 histogram.record(5, &[KeyValue::new("key1", "value1")]);
3134 histogram.record(7, &[KeyValue::new("key1", "value1")]);
3135 histogram.record(18, &[KeyValue::new("key1", "value1")]);
3136
3137 if i % 2 == 0 {
3139 test_context.flush_metrics();
3140 thread::sleep(Duration::from_millis(i)); }
3142
3143 histogram.record(35, &[KeyValue::new("key1", "value1")]);
3144 histogram.record(35, &[KeyValue::new("key1", "value1")]);
3145 });
3146 });
3147 }
3148
3149 test_context.flush_metrics();
3150
3151 let histograms = test_context
3154 .get_from_multiple_aggregations::<u64>("test_histogram", None, 6)
3155 .into_iter()
3156 .map(|data| {
3157 if let MetricData::Histogram(hist) = data {
3158 hist
3159 } else {
3160 unreachable!()
3161 }
3162 })
3163 .collect::<Vec<_>>();
3164
3165 let (
3166 mut sum_zero_attributes,
3167 mut count_zero_attributes,
3168 mut min_zero_attributes,
3169 mut max_zero_attributes,
3170 ) = (0, 0, u64::MAX, u64::MIN);
3171 let (mut sum_key1_value1, mut count_key1_value1, mut min_key1_value1, mut max_key1_value1) =
3172 (0, 0, u64::MAX, u64::MIN);
3173
3174 let mut bucket_counts_zero_attributes = vec![0; 16]; let mut bucket_counts_key1_value1 = vec![0; 16];
3176
3177 histograms.iter().for_each(|histogram| {
3178 assert_eq!(histogram.data_points.len(), 2); assert_eq!(histogram.temporality, temporality);
3180
3181 let data_point_zero_attributes =
3182 find_histogram_datapoint_with_no_attributes(&histogram.data_points).unwrap();
3183 let data_point_key1_value1 =
3184 find_histogram_datapoint_with_key_value(&histogram.data_points, "key1", "value1")
3185 .unwrap();
3186
3187 if temporality == Temporality::Delta {
3188 sum_zero_attributes += data_point_zero_attributes.sum;
3189 sum_key1_value1 += data_point_key1_value1.sum;
3190
3191 count_zero_attributes += data_point_zero_attributes.count;
3192 count_key1_value1 += data_point_key1_value1.count;
3193
3194 min_zero_attributes =
3195 min(min_zero_attributes, data_point_zero_attributes.min.unwrap());
3196 min_key1_value1 = min(min_key1_value1, data_point_key1_value1.min.unwrap());
3197
3198 max_zero_attributes =
3199 max(max_zero_attributes, data_point_zero_attributes.max.unwrap());
3200 max_key1_value1 = max(max_key1_value1, data_point_key1_value1.max.unwrap());
3201
3202 assert_eq!(data_point_zero_attributes.bucket_counts.len(), 16);
3203 assert_eq!(data_point_key1_value1.bucket_counts.len(), 16);
3204
3205 for (i, _) in data_point_zero_attributes.bucket_counts.iter().enumerate() {
3206 bucket_counts_zero_attributes[i] += data_point_zero_attributes.bucket_counts[i];
3207 }
3208
3209 for (i, _) in data_point_key1_value1.bucket_counts.iter().enumerate() {
3210 bucket_counts_key1_value1[i] += data_point_key1_value1.bucket_counts[i];
3211 }
3212 } else {
3213 sum_zero_attributes = data_point_zero_attributes.sum;
3214 sum_key1_value1 = data_point_key1_value1.sum;
3215
3216 count_zero_attributes = data_point_zero_attributes.count;
3217 count_key1_value1 = data_point_key1_value1.count;
3218
3219 min_zero_attributes = data_point_zero_attributes.min.unwrap();
3220 min_key1_value1 = data_point_key1_value1.min.unwrap();
3221
3222 max_zero_attributes = data_point_zero_attributes.max.unwrap();
3223 max_key1_value1 = data_point_key1_value1.max.unwrap();
3224
3225 assert_eq!(data_point_zero_attributes.bucket_counts.len(), 16);
3226 assert_eq!(data_point_key1_value1.bucket_counts.len(), 16);
3227
3228 bucket_counts_zero_attributes.clone_from(&data_point_zero_attributes.bucket_counts);
3229 bucket_counts_key1_value1.clone_from(&data_point_key1_value1.bucket_counts);
3230 };
3231 });
3232
3233 assert_eq!(count_zero_attributes, 20); assert_eq!(sum_zero_attributes, 50); assert_eq!(min_zero_attributes, 1);
3240 assert_eq!(max_zero_attributes, 4);
3241
3242 for (i, count) in bucket_counts_zero_attributes.iter().enumerate() {
3243 match i {
3244 1 => assert_eq!(*count, 20), _ => assert_eq!(*count, 0),
3246 }
3247 }
3248
3249 assert_eq!(count_key1_value1, 50); assert_eq!(sum_key1_value1, 1000); assert_eq!(min_key1_value1, 5);
3252 assert_eq!(max_key1_value1, 35);
3253
3254 for (i, count) in bucket_counts_key1_value1.iter().enumerate() {
3255 match i {
3256 1 => assert_eq!(*count, 10), 2 => assert_eq!(*count, 10), 3 => assert_eq!(*count, 10), 4 => assert_eq!(*count, 20), _ => assert_eq!(*count, 0),
3261 }
3262 }
3263 }
3264
3265 fn histogram_f64_multithreaded_aggregation_helper(temporality: Temporality) {
3266 let mut test_context = TestContext::new(temporality);
3268 let histogram = Arc::new(test_context.meter().f64_histogram("test_histogram").build());
3269
3270 for i in 0..10 {
3271 thread::scope(|s| {
3272 s.spawn(|| {
3273 histogram.record(1.5, &[]);
3274 histogram.record(4.6, &[]);
3275
3276 histogram.record(5.0, &[KeyValue::new("key1", "value1")]);
3277 histogram.record(7.3, &[KeyValue::new("key1", "value1")]);
3278 histogram.record(18.1, &[KeyValue::new("key1", "value1")]);
3279
3280 if i % 2 == 0 {
3282 test_context.flush_metrics();
3283 thread::sleep(Duration::from_millis(i)); }
3285
3286 histogram.record(35.1, &[KeyValue::new("key1", "value1")]);
3287 histogram.record(35.1, &[KeyValue::new("key1", "value1")]);
3288 });
3289 });
3290 }
3291
3292 test_context.flush_metrics();
3293
3294 let histograms = test_context
3297 .get_from_multiple_aggregations::<f64>("test_histogram", None, 6)
3298 .into_iter()
3299 .map(|data| {
3300 if let MetricData::Histogram(hist) = data {
3301 hist
3302 } else {
3303 unreachable!()
3304 }
3305 })
3306 .collect::<Vec<_>>();
3307
3308 let (
3309 mut sum_zero_attributes,
3310 mut count_zero_attributes,
3311 mut min_zero_attributes,
3312 mut max_zero_attributes,
3313 ) = (0.0, 0, f64::MAX, f64::MIN);
3314 let (mut sum_key1_value1, mut count_key1_value1, mut min_key1_value1, mut max_key1_value1) =
3315 (0.0, 0, f64::MAX, f64::MIN);
3316
3317 let mut bucket_counts_zero_attributes = vec![0; 16]; let mut bucket_counts_key1_value1 = vec![0; 16];
3319
3320 histograms.iter().for_each(|histogram| {
3321 assert_eq!(histogram.data_points.len(), 2); assert_eq!(histogram.temporality, temporality);
3323
3324 let data_point_zero_attributes =
3325 find_histogram_datapoint_with_no_attributes(&histogram.data_points).unwrap();
3326 let data_point_key1_value1 =
3327 find_histogram_datapoint_with_key_value(&histogram.data_points, "key1", "value1")
3328 .unwrap();
3329
3330 if temporality == Temporality::Delta {
3331 sum_zero_attributes += data_point_zero_attributes.sum;
3332 sum_key1_value1 += data_point_key1_value1.sum;
3333
3334 count_zero_attributes += data_point_zero_attributes.count;
3335 count_key1_value1 += data_point_key1_value1.count;
3336
3337 min_zero_attributes =
3338 min_zero_attributes.min(data_point_zero_attributes.min.unwrap());
3339 min_key1_value1 = min_key1_value1.min(data_point_key1_value1.min.unwrap());
3340
3341 max_zero_attributes =
3342 max_zero_attributes.max(data_point_zero_attributes.max.unwrap());
3343 max_key1_value1 = max_key1_value1.max(data_point_key1_value1.max.unwrap());
3344
3345 assert_eq!(data_point_zero_attributes.bucket_counts.len(), 16);
3346 assert_eq!(data_point_key1_value1.bucket_counts.len(), 16);
3347
3348 for (i, _) in data_point_zero_attributes.bucket_counts.iter().enumerate() {
3349 bucket_counts_zero_attributes[i] += data_point_zero_attributes.bucket_counts[i];
3350 }
3351
3352 for (i, _) in data_point_key1_value1.bucket_counts.iter().enumerate() {
3353 bucket_counts_key1_value1[i] += data_point_key1_value1.bucket_counts[i];
3354 }
3355 } else {
3356 sum_zero_attributes = data_point_zero_attributes.sum;
3357 sum_key1_value1 = data_point_key1_value1.sum;
3358
3359 count_zero_attributes = data_point_zero_attributes.count;
3360 count_key1_value1 = data_point_key1_value1.count;
3361
3362 min_zero_attributes = data_point_zero_attributes.min.unwrap();
3363 min_key1_value1 = data_point_key1_value1.min.unwrap();
3364
3365 max_zero_attributes = data_point_zero_attributes.max.unwrap();
3366 max_key1_value1 = data_point_key1_value1.max.unwrap();
3367
3368 assert_eq!(data_point_zero_attributes.bucket_counts.len(), 16);
3369 assert_eq!(data_point_key1_value1.bucket_counts.len(), 16);
3370
3371 bucket_counts_zero_attributes.clone_from(&data_point_zero_attributes.bucket_counts);
3372 bucket_counts_key1_value1.clone_from(&data_point_key1_value1.bucket_counts);
3373 };
3374 });
3375
3376 assert_eq!(count_zero_attributes, 20); assert!(f64::abs(61.0 - sum_zero_attributes) < 0.0001); assert_eq!(min_zero_attributes, 1.5);
3383 assert_eq!(max_zero_attributes, 4.6);
3384
3385 for (i, count) in bucket_counts_zero_attributes.iter().enumerate() {
3386 match i {
3387 1 => assert_eq!(*count, 20), _ => assert_eq!(*count, 0),
3389 }
3390 }
3391
3392 assert_eq!(count_key1_value1, 50); assert!(f64::abs(1006.0 - sum_key1_value1) < 0.0001); assert_eq!(min_key1_value1, 5.0);
3395 assert_eq!(max_key1_value1, 35.1);
3396
3397 for (i, count) in bucket_counts_key1_value1.iter().enumerate() {
3398 match i {
3399 1 => assert_eq!(*count, 10), 2 => assert_eq!(*count, 10), 3 => assert_eq!(*count, 10), 4 => assert_eq!(*count, 20), _ => assert_eq!(*count, 0),
3404 }
3405 }
3406 }
3407
3408 fn histogram_aggregation_helper(temporality: Temporality) {
3409 let mut test_context = TestContext::new(temporality);
3411 let histogram = test_context.meter().u64_histogram("my_histogram").build();
3412
3413 let mut rand = rngs::SmallRng::from_os_rng();
3415 let values_kv1 = (0..50)
3416 .map(|_| rand.random_range(0..100))
3417 .collect::<Vec<u64>>();
3418 for value in values_kv1.iter() {
3419 histogram.record(*value, &[KeyValue::new("key1", "value1")]);
3420 }
3421
3422 let values_kv2 = (0..30)
3423 .map(|_| rand.random_range(0..100))
3424 .collect::<Vec<u64>>();
3425 for value in values_kv2.iter() {
3426 histogram.record(*value, &[KeyValue::new("key1", "value2")]);
3427 }
3428
3429 test_context.flush_metrics();
3430
3431 let MetricData::Histogram(histogram_data) =
3433 test_context.get_aggregation::<u64>("my_histogram", None)
3434 else {
3435 unreachable!()
3436 };
3437 assert_eq!(histogram_data.data_points.len(), 2);
3439 if let Temporality::Cumulative = temporality {
3440 assert_eq!(
3441 histogram_data.temporality,
3442 Temporality::Cumulative,
3443 "Should produce cumulative"
3444 );
3445 } else {
3446 assert_eq!(
3447 histogram_data.temporality,
3448 Temporality::Delta,
3449 "Should produce delta"
3450 );
3451 }
3452
3453 let data_point1 =
3455 find_histogram_datapoint_with_key_value(&histogram_data.data_points, "key1", "value1")
3456 .expect("datapoint with key1=value1 expected");
3457 assert_eq!(data_point1.count, values_kv1.len() as u64);
3458 assert_eq!(data_point1.sum, values_kv1.iter().sum::<u64>());
3459 assert_eq!(data_point1.min.unwrap(), *values_kv1.iter().min().unwrap());
3460 assert_eq!(data_point1.max.unwrap(), *values_kv1.iter().max().unwrap());
3461
3462 let data_point2 =
3463 find_histogram_datapoint_with_key_value(&histogram_data.data_points, "key1", "value2")
3464 .expect("datapoint with key1=value2 expected");
3465 assert_eq!(data_point2.count, values_kv2.len() as u64);
3466 assert_eq!(data_point2.sum, values_kv2.iter().sum::<u64>());
3467 assert_eq!(data_point2.min.unwrap(), *values_kv2.iter().min().unwrap());
3468 assert_eq!(data_point2.max.unwrap(), *values_kv2.iter().max().unwrap());
3469
3470 test_context.reset_metrics();
3472 for value in values_kv1.iter() {
3473 histogram.record(*value, &[KeyValue::new("key1", "value1")]);
3474 }
3475
3476 for value in values_kv2.iter() {
3477 histogram.record(*value, &[KeyValue::new("key1", "value2")]);
3478 }
3479
3480 test_context.flush_metrics();
3481
3482 let MetricData::Histogram(histogram_data) =
3483 test_context.get_aggregation::<u64>("my_histogram", None)
3484 else {
3485 unreachable!()
3486 };
3487 assert_eq!(histogram_data.data_points.len(), 2);
3488 let data_point1 =
3489 find_histogram_datapoint_with_key_value(&histogram_data.data_points, "key1", "value1")
3490 .expect("datapoint with key1=value1 expected");
3491 if temporality == Temporality::Cumulative {
3492 assert_eq!(data_point1.count, 2 * (values_kv1.len() as u64));
3493 assert_eq!(data_point1.sum, 2 * (values_kv1.iter().sum::<u64>()));
3494 assert_eq!(data_point1.min.unwrap(), *values_kv1.iter().min().unwrap());
3495 assert_eq!(data_point1.max.unwrap(), *values_kv1.iter().max().unwrap());
3496 } else {
3497 assert_eq!(data_point1.count, values_kv1.len() as u64);
3498 assert_eq!(data_point1.sum, values_kv1.iter().sum::<u64>());
3499 assert_eq!(data_point1.min.unwrap(), *values_kv1.iter().min().unwrap());
3500 assert_eq!(data_point1.max.unwrap(), *values_kv1.iter().max().unwrap());
3501 }
3502
3503 let data_point1 =
3504 find_histogram_datapoint_with_key_value(&histogram_data.data_points, "key1", "value2")
3505 .expect("datapoint with key1=value1 expected");
3506 if temporality == Temporality::Cumulative {
3507 assert_eq!(data_point1.count, 2 * (values_kv2.len() as u64));
3508 assert_eq!(data_point1.sum, 2 * (values_kv2.iter().sum::<u64>()));
3509 assert_eq!(data_point1.min.unwrap(), *values_kv2.iter().min().unwrap());
3510 assert_eq!(data_point1.max.unwrap(), *values_kv2.iter().max().unwrap());
3511 } else {
3512 assert_eq!(data_point1.count, values_kv2.len() as u64);
3513 assert_eq!(data_point1.sum, values_kv2.iter().sum::<u64>());
3514 assert_eq!(data_point1.min.unwrap(), *values_kv2.iter().min().unwrap());
3515 assert_eq!(data_point1.max.unwrap(), *values_kv2.iter().max().unwrap());
3516 }
3517 }
3518
3519 fn histogram_aggregation_with_custom_bounds_helper(temporality: Temporality) {
3520 let mut test_context = TestContext::new(temporality);
3521 let histogram = test_context
3522 .meter()
3523 .u64_histogram("test_histogram")
3524 .with_boundaries(vec![1.0, 2.5, 5.5])
3525 .build();
3526 histogram.record(1, &[KeyValue::new("key1", "value1")]);
3527 histogram.record(2, &[KeyValue::new("key1", "value1")]);
3528 histogram.record(3, &[KeyValue::new("key1", "value1")]);
3529 histogram.record(4, &[KeyValue::new("key1", "value1")]);
3530 histogram.record(5, &[KeyValue::new("key1", "value1")]);
3531
3532 test_context.flush_metrics();
3533
3534 let MetricData::Histogram(histogram_data) =
3536 test_context.get_aggregation::<u64>("test_histogram", None)
3537 else {
3538 unreachable!()
3539 };
3540 assert_eq!(histogram_data.data_points.len(), 1);
3542 if let Temporality::Cumulative = temporality {
3543 assert_eq!(
3544 histogram_data.temporality,
3545 Temporality::Cumulative,
3546 "Should produce cumulative"
3547 );
3548 } else {
3549 assert_eq!(
3550 histogram_data.temporality,
3551 Temporality::Delta,
3552 "Should produce delta"
3553 );
3554 }
3555
3556 let data_point =
3558 find_histogram_datapoint_with_key_value(&histogram_data.data_points, "key1", "value1")
3559 .expect("datapoint with key1=value1 expected");
3560
3561 assert_eq!(data_point.count, 5);
3562 assert_eq!(data_point.sum, 15);
3563
3564 assert_eq!(vec![1.0, 2.5, 5.5], data_point.bounds);
3571 assert_eq!(vec![1, 1, 3, 0], data_point.bucket_counts);
3572 }
3573
3574 fn histogram_aggregation_with_empty_bounds_helper(temporality: Temporality) {
3575 let mut test_context = TestContext::new(temporality);
3576 let histogram = test_context
3577 .meter()
3578 .u64_histogram("test_histogram")
3579 .with_boundaries(vec![])
3580 .build();
3581 histogram.record(1, &[KeyValue::new("key1", "value1")]);
3582 histogram.record(2, &[KeyValue::new("key1", "value1")]);
3583 histogram.record(3, &[KeyValue::new("key1", "value1")]);
3584 histogram.record(4, &[KeyValue::new("key1", "value1")]);
3585 histogram.record(5, &[KeyValue::new("key1", "value1")]);
3586
3587 test_context.flush_metrics();
3588
3589 let MetricData::Histogram(histogram_data) =
3591 test_context.get_aggregation::<u64>("test_histogram", None)
3592 else {
3593 unreachable!()
3594 };
3595 assert_eq!(histogram_data.data_points.len(), 1);
3597 if let Temporality::Cumulative = temporality {
3598 assert_eq!(
3599 histogram_data.temporality,
3600 Temporality::Cumulative,
3601 "Should produce cumulative"
3602 );
3603 } else {
3604 assert_eq!(
3605 histogram_data.temporality,
3606 Temporality::Delta,
3607 "Should produce delta"
3608 );
3609 }
3610
3611 let data_point =
3613 find_histogram_datapoint_with_key_value(&histogram_data.data_points, "key1", "value1")
3614 .expect("datapoint with key1=value1 expected");
3615
3616 assert_eq!(data_point.count, 5);
3617 assert_eq!(data_point.sum, 15);
3618 assert!(data_point.bounds.is_empty());
3619 assert!(data_point.bucket_counts.is_empty());
3620 }
3621
3622 fn histogram_aggregation_with_custom_bounds_and_view_helper(temporality: Temporality) {
3623 for specify_boundaries_in_view in [false, true] {
3624 let view = move |_: &Instrument| {
3625 let mut builder = Stream::builder();
3626 if specify_boundaries_in_view {
3627 builder = builder.with_aggregation(Aggregation::ExplicitBucketHistogram {
3628 boundaries: vec![1.5, 4.2, 6.7],
3629 record_min_max: true,
3630 });
3631 }
3632 Some(builder.build().unwrap())
3633 };
3634 let mut test_context = TestContext::new_with_view(temporality, view);
3635 let histogram = test_context
3636 .meter()
3637 .u64_histogram("test_histogram")
3638 .with_boundaries(vec![1.0, 2.5, 5.5])
3639 .build();
3640 histogram.record(1, &[KeyValue::new("key1", "value1")]);
3641 histogram.record(2, &[KeyValue::new("key1", "value1")]);
3642 histogram.record(3, &[KeyValue::new("key1", "value1")]);
3643 histogram.record(4, &[KeyValue::new("key1", "value1")]);
3644 histogram.record(5, &[KeyValue::new("key1", "value1")]);
3645
3646 test_context.flush_metrics();
3647
3648 let MetricData::Histogram(histogram_data) =
3649 test_context.get_aggregation::<u64>("test_histogram", None)
3650 else {
3651 unreachable!()
3652 };
3653 assert_eq!(histogram_data.data_points.len(), 1);
3654 if let Temporality::Cumulative = temporality {
3655 assert_eq!(
3656 histogram_data.temporality,
3657 Temporality::Cumulative,
3658 "Should produce cumulative"
3659 );
3660 } else {
3661 assert_eq!(
3662 histogram_data.temporality,
3663 Temporality::Delta,
3664 "Should produce delta"
3665 );
3666 }
3667
3668 let data_point = find_histogram_datapoint_with_key_value(
3670 &histogram_data.data_points,
3671 "key1",
3672 "value1",
3673 )
3674 .expect("datapoint with key1=value1 expected");
3675
3676 assert_eq!(data_point.count, 5);
3677 assert_eq!(data_point.sum, 15);
3678
3679 if specify_boundaries_in_view {
3681 assert_eq!(vec![1.5, 4.2, 6.7], data_point.bounds);
3683 assert_eq!(vec![1, 3, 1, 0], data_point.bucket_counts);
3684 } else {
3685 assert_eq!(vec![1.0, 2.5, 5.5], data_point.bounds);
3688 assert_eq!(vec![1, 1, 3, 0], data_point.bucket_counts);
3689 }
3690 }
3691 }
3692
3693 fn exponential_histogram_aggregation_with_view_helper(temporality: Temporality) {
3694 let view = |i: &Instrument| {
3696 if i.name == "test_histogram" {
3697 Some(
3698 Stream::builder()
3699 .with_aggregation(Aggregation::Base2ExponentialHistogram {
3700 max_size: 160,
3701 max_scale: 20,
3702 record_min_max: true,
3703 })
3704 .build()
3705 .unwrap(),
3706 )
3707 } else {
3708 None
3709 }
3710 };
3711 let mut test_context = TestContext::new_with_view(temporality, view);
3712 let histogram = test_context.meter().f64_histogram("test_histogram").build();
3713
3714 histogram.record(1.0, &[KeyValue::new("key1", "value1")]);
3716 histogram.record(2.0, &[KeyValue::new("key1", "value1")]);
3717 histogram.record(3.0, &[KeyValue::new("key1", "value1")]);
3718 histogram.record(4.0, &[KeyValue::new("key1", "value1")]);
3719 histogram.record(5.0, &[KeyValue::new("key1", "value1")]);
3720
3721 test_context.flush_metrics();
3722
3723 let exponential_histogram_data =
3725 test_context.get_aggregation::<f64>("test_histogram", None);
3726 let MetricData::ExponentialHistogram(exp_hist) = exponential_histogram_data else {
3727 panic!(
3728 "Expected ExponentialHistogram aggregation, got {:?}",
3729 exponential_histogram_data
3730 );
3731 };
3732
3733 assert_eq!(exp_hist.data_points.len(), 1);
3734 if let Temporality::Cumulative = temporality {
3735 assert_eq!(
3736 exp_hist.temporality,
3737 Temporality::Cumulative,
3738 "Should produce cumulative"
3739 );
3740 } else {
3741 assert_eq!(
3742 exp_hist.temporality,
3743 Temporality::Delta,
3744 "Should produce delta"
3745 );
3746 }
3747
3748 let data_point = &exp_hist.data_points[0];
3750 assert_eq!(data_point.count(), 5);
3751 assert_eq!(data_point.sum(), 15.0);
3752 assert_eq!(data_point.min(), Some(1.0));
3753 assert_eq!(data_point.max(), Some(5.0));
3754
3755 let scale = data_point.scale();
3758 assert!(
3759 (-10..=20).contains(&scale),
3760 "Scale {} should be within valid range [-10, 20]",
3761 scale
3762 );
3763
3764 assert_eq!(
3766 data_point.zero_count(),
3767 0,
3768 "zero_count should be 0 for positive values"
3769 );
3770
3771 let positive_bucket = data_point.positive_bucket();
3773 let positive_counts: Vec<u64> = positive_bucket.counts().collect();
3774 let total_positive_count: u64 = positive_counts.iter().sum();
3775 assert_eq!(
3776 total_positive_count, 5,
3777 "Total count in positive buckets should equal number of recorded values"
3778 );
3779
3780 let negative_bucket = data_point.negative_bucket();
3782 let negative_counts: Vec<u64> = negative_bucket.counts().collect();
3783 let total_negative_count: u64 = negative_counts.iter().sum();
3784 assert_eq!(
3785 total_negative_count, 0,
3786 "Negative bucket should be empty for positive-only values"
3787 );
3788
3789 let attrs: Vec<_> = data_point.attributes().collect();
3791 assert_eq!(attrs.len(), 1);
3792 assert_eq!(attrs[0].key.as_str(), "key1");
3793
3794 test_context.reset_metrics();
3796 histogram.record(10.0, &[KeyValue::new("key1", "value1")]);
3797 histogram.record(20.0, &[KeyValue::new("key1", "value1")]);
3798 histogram.record(30.0, &[KeyValue::new("key1", "value1")]);
3799
3800 test_context.flush_metrics();
3801
3802 let exponential_histogram_data =
3804 test_context.get_aggregation::<f64>("test_histogram", None);
3805 let MetricData::ExponentialHistogram(exp_hist) = exponential_histogram_data else {
3806 panic!(
3807 "Expected ExponentialHistogram aggregation, got {:?}",
3808 exponential_histogram_data
3809 );
3810 };
3811
3812 assert_eq!(exp_hist.data_points.len(), 1);
3813 let data_point = &exp_hist.data_points[0];
3814
3815 if temporality == Temporality::Cumulative {
3816 assert_eq!(data_point.count(), 8);
3818 assert_eq!(data_point.sum(), 75.0);
3819 assert_eq!(data_point.min(), Some(1.0)); assert_eq!(data_point.max(), Some(30.0)); } else {
3822 assert_eq!(data_point.count(), 3);
3824 assert_eq!(data_point.sum(), 60.0);
3825 assert_eq!(data_point.min(), Some(10.0));
3826 assert_eq!(data_point.max(), Some(30.0));
3827 }
3828
3829 let positive_bucket = data_point.positive_bucket();
3831 let positive_counts: Vec<u64> = positive_bucket.counts().collect();
3832 let total_positive_count: u64 = positive_counts.iter().sum();
3833 assert_eq!(
3834 total_positive_count,
3835 data_point.count() as u64,
3836 "Total count in positive buckets should equal count"
3837 );
3838 }
3839
3840 fn gauge_aggregation_helper(temporality: Temporality) {
3841 let mut test_context = TestContext::new(temporality);
3843 let gauge = test_context.meter().i64_gauge("my_gauge").build();
3844
3845 gauge.record(1, &[KeyValue::new("key1", "value1")]);
3847 gauge.record(2, &[KeyValue::new("key1", "value1")]);
3848 gauge.record(1, &[KeyValue::new("key1", "value1")]);
3849 gauge.record(3, &[KeyValue::new("key1", "value1")]);
3850 gauge.record(4, &[KeyValue::new("key1", "value1")]);
3851
3852 gauge.record(11, &[KeyValue::new("key1", "value2")]);
3853 gauge.record(13, &[KeyValue::new("key1", "value2")]);
3854 gauge.record(6, &[KeyValue::new("key1", "value2")]);
3855
3856 test_context.flush_metrics();
3857
3858 let MetricData::Gauge(gauge_data_point) =
3860 test_context.get_aggregation::<i64>("my_gauge", None)
3861 else {
3862 unreachable!()
3863 };
3864 assert_eq!(gauge_data_point.data_points.len(), 2);
3866
3867 let data_point1 =
3869 find_gauge_datapoint_with_key_value(&gauge_data_point.data_points, "key1", "value1")
3870 .expect("datapoint with key1=value1 expected");
3871 assert_eq!(data_point1.value, 4);
3872
3873 let data_point1 =
3874 find_gauge_datapoint_with_key_value(&gauge_data_point.data_points, "key1", "value2")
3875 .expect("datapoint with key1=value2 expected");
3876 assert_eq!(data_point1.value, 6);
3877
3878 test_context.reset_metrics();
3880 gauge.record(1, &[KeyValue::new("key1", "value1")]);
3881 gauge.record(2, &[KeyValue::new("key1", "value1")]);
3882 gauge.record(11, &[KeyValue::new("key1", "value1")]);
3883 gauge.record(3, &[KeyValue::new("key1", "value1")]);
3884 gauge.record(41, &[KeyValue::new("key1", "value1")]);
3885
3886 gauge.record(34, &[KeyValue::new("key1", "value2")]);
3887 gauge.record(12, &[KeyValue::new("key1", "value2")]);
3888 gauge.record(54, &[KeyValue::new("key1", "value2")]);
3889
3890 test_context.flush_metrics();
3891
3892 let MetricData::Gauge(gauge) = test_context.get_aggregation::<i64>("my_gauge", None) else {
3893 unreachable!()
3894 };
3895 assert_eq!(gauge.data_points.len(), 2);
3896 let data_point1 = find_gauge_datapoint_with_key_value(&gauge.data_points, "key1", "value1")
3897 .expect("datapoint with key1=value1 expected");
3898 assert_eq!(data_point1.value, 41);
3899
3900 let data_point1 = find_gauge_datapoint_with_key_value(&gauge.data_points, "key1", "value2")
3901 .expect("datapoint with key1=value2 expected");
3902 assert_eq!(data_point1.value, 54);
3903 }
3904
3905 fn observable_gauge_aggregation_helper(temporality: Temporality, use_empty_attributes: bool) {
3906 let mut test_context = TestContext::new(temporality);
3908 let _observable_gauge = test_context
3909 .meter()
3910 .i64_observable_gauge("test_observable_gauge")
3911 .with_callback(move |observer| {
3912 if use_empty_attributes {
3913 observer.observe(1, &[]);
3914 }
3915 observer.observe(4, &[KeyValue::new("key1", "value1")]);
3916 observer.observe(5, &[KeyValue::new("key2", "value2")]);
3917 })
3918 .build();
3919
3920 test_context.flush_metrics();
3921
3922 let MetricData::Gauge(gauge) =
3924 test_context.get_aggregation::<i64>("test_observable_gauge", None)
3925 else {
3926 unreachable!()
3927 };
3928 let expected_time_series_count = if use_empty_attributes { 3 } else { 2 };
3930 assert_eq!(gauge.data_points.len(), expected_time_series_count);
3931
3932 if use_empty_attributes {
3933 let zero_attribute_datapoint =
3935 find_gauge_datapoint_with_no_attributes(&gauge.data_points)
3936 .expect("datapoint with no attributes expected");
3937 assert_eq!(zero_attribute_datapoint.value, 1);
3938 }
3939
3940 let data_point1 = find_gauge_datapoint_with_key_value(&gauge.data_points, "key1", "value1")
3942 .expect("datapoint with key1=value1 expected");
3943 assert_eq!(data_point1.value, 4);
3944
3945 let data_point2 = find_gauge_datapoint_with_key_value(&gauge.data_points, "key2", "value2")
3947 .expect("datapoint with key2=value2 expected");
3948 assert_eq!(data_point2.value, 5);
3949
3950 test_context.reset_metrics();
3952
3953 test_context.flush_metrics();
3954
3955 let MetricData::Gauge(gauge) =
3956 test_context.get_aggregation::<i64>("test_observable_gauge", None)
3957 else {
3958 unreachable!()
3959 };
3960 assert_eq!(gauge.data_points.len(), expected_time_series_count);
3961
3962 if use_empty_attributes {
3963 let zero_attribute_datapoint =
3964 find_gauge_datapoint_with_no_attributes(&gauge.data_points)
3965 .expect("datapoint with no attributes expected");
3966 assert_eq!(zero_attribute_datapoint.value, 1);
3967 }
3968
3969 let data_point1 = find_gauge_datapoint_with_key_value(&gauge.data_points, "key1", "value1")
3970 .expect("datapoint with key1=value1 expected");
3971 assert_eq!(data_point1.value, 4);
3972
3973 let data_point2 = find_gauge_datapoint_with_key_value(&gauge.data_points, "key2", "value2")
3974 .expect("datapoint with key2=value2 expected");
3975 assert_eq!(data_point2.value, 5);
3976 }
3977
3978 fn counter_aggregation_helper(temporality: Temporality) {
3979 let mut test_context = TestContext::new(temporality);
3981 let counter = test_context.u64_counter("test", "my_counter", None);
3982
3983 counter.add(1, &[KeyValue::new("key1", "value1")]);
3985 counter.add(1, &[KeyValue::new("key1", "value1")]);
3986 counter.add(1, &[KeyValue::new("key1", "value1")]);
3987 counter.add(1, &[KeyValue::new("key1", "value1")]);
3988 counter.add(1, &[KeyValue::new("key1", "value1")]);
3989
3990 counter.add(1, &[KeyValue::new("key1", "value2")]);
3991 counter.add(1, &[KeyValue::new("key1", "value2")]);
3992 counter.add(1, &[KeyValue::new("key1", "value2")]);
3993
3994 test_context.flush_metrics();
3995
3996 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
3998 unreachable!()
3999 };
4000 assert_eq!(sum.data_points.len(), 2);
4002 assert!(sum.is_monotonic, "Counter should produce monotonic.");
4003 if let Temporality::Cumulative = temporality {
4004 assert_eq!(
4005 sum.temporality,
4006 Temporality::Cumulative,
4007 "Should produce cumulative"
4008 );
4009 } else {
4010 assert_eq!(sum.temporality, Temporality::Delta, "Should produce delta");
4011 }
4012
4013 let data_point1 = find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value1")
4015 .expect("datapoint with key1=value1 expected");
4016 assert_eq!(data_point1.value, 5);
4017
4018 let data_point1 = find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value2")
4019 .expect("datapoint with key1=value2 expected");
4020 assert_eq!(data_point1.value, 3);
4021
4022 test_context.reset_metrics();
4024 counter.add(1, &[KeyValue::new("key1", "value1")]);
4025 counter.add(1, &[KeyValue::new("key1", "value1")]);
4026 counter.add(1, &[KeyValue::new("key1", "value1")]);
4027 counter.add(1, &[KeyValue::new("key1", "value1")]);
4028 counter.add(1, &[KeyValue::new("key1", "value1")]);
4029
4030 counter.add(1, &[KeyValue::new("key1", "value2")]);
4031 counter.add(1, &[KeyValue::new("key1", "value2")]);
4032 counter.add(1, &[KeyValue::new("key1", "value2")]);
4033
4034 test_context.flush_metrics();
4035
4036 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4037 unreachable!()
4038 };
4039 assert_eq!(sum.data_points.len(), 2);
4040 let data_point1 = find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value1")
4041 .expect("datapoint with key1=value1 expected");
4042 if temporality == Temporality::Cumulative {
4043 assert_eq!(data_point1.value, 10);
4044 } else {
4045 assert_eq!(data_point1.value, 5);
4046 }
4047
4048 let data_point1 = find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value2")
4049 .expect("datapoint with key1=value2 expected");
4050 if temporality == Temporality::Cumulative {
4051 assert_eq!(data_point1.value, 6);
4052 } else {
4053 assert_eq!(data_point1.value, 3);
4054 }
4055 }
4056
4057 fn counter_aggregation_overflow_helper(temporality: Temporality) {
4058 let mut test_context = TestContext::new(temporality);
4060 let counter = test_context.u64_counter("test", "my_counter", None);
4061
4062 for v in 0..2000 {
4065 counter.add(100, &[KeyValue::new("A", v.to_string())]);
4066 }
4067
4068 counter.add(3, &[]);
4070 counter.add(3, &[]);
4071
4072 counter.add(100, &[KeyValue::new("A", "foo")]);
4074 counter.add(100, &[KeyValue::new("A", "another")]);
4075 counter.add(100, &[KeyValue::new("A", "yet_another")]);
4076 test_context.flush_metrics();
4077
4078 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4079 unreachable!()
4080 };
4081
4082 assert_eq!(sum.data_points.len(), 2002);
4084
4085 let data_point =
4086 find_overflow_sum_datapoint(&sum.data_points).expect("overflow point expected");
4087 assert_eq!(data_point.value, 300);
4088
4089 let empty_attrs_data_point = find_sum_datapoint_with_no_attributes(&sum.data_points)
4091 .expect("Empty attributes point expected");
4092 assert!(
4093 empty_attrs_data_point.attributes.is_empty(),
4094 "Non-empty attribute set"
4095 );
4096 assert_eq!(
4097 empty_attrs_data_point.value, 6,
4098 "Empty attributes value should be 3+3=6"
4099 );
4100
4101 test_context.reset_metrics();
4106 if temporality == Temporality::Delta {
4107 test_context.flush_metrics();
4108 test_context.reset_metrics();
4109 }
4110 counter.add(100, &[KeyValue::new("A", "foo")]);
4113 counter.add(100, &[KeyValue::new("A", "another")]);
4114 counter.add(100, &[KeyValue::new("A", "yet_another")]);
4115 test_context.flush_metrics();
4116
4117 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4118 unreachable!()
4119 };
4120
4121 if temporality == Temporality::Delta {
4122 assert_eq!(sum.data_points.len(), 3);
4123
4124 let data_point = find_sum_datapoint_with_key_value(&sum.data_points, "A", "foo")
4125 .expect("point expected");
4126 assert_eq!(data_point.value, 100);
4127
4128 let data_point = find_sum_datapoint_with_key_value(&sum.data_points, "A", "another")
4129 .expect("point expected");
4130 assert_eq!(data_point.value, 100);
4131
4132 let data_point =
4133 find_sum_datapoint_with_key_value(&sum.data_points, "A", "yet_another")
4134 .expect("point expected");
4135 assert_eq!(data_point.value, 100);
4136 } else {
4137 assert_eq!(sum.data_points.len(), 2002);
4139 let data_point =
4140 find_overflow_sum_datapoint(&sum.data_points).expect("overflow point expected");
4141 assert_eq!(data_point.value, 600);
4142
4143 let data_point = find_sum_datapoint_with_key_value(&sum.data_points, "A", "foo");
4144 assert!(data_point.is_none(), "point should not be present");
4145
4146 let data_point = find_sum_datapoint_with_key_value(&sum.data_points, "A", "another");
4147 assert!(data_point.is_none(), "point should not be present");
4148
4149 let data_point =
4150 find_sum_datapoint_with_key_value(&sum.data_points, "A", "yet_another");
4151 assert!(data_point.is_none(), "point should not be present");
4152 }
4153 }
4154
4155 fn counter_aggregation_overflow_helper_custom_limit(temporality: Temporality) {
4156 let cardinality_limit = 2300;
4158 let view_change_cardinality = move |i: &Instrument| {
4159 if i.name == "my_counter" {
4160 Some(
4161 Stream::builder()
4162 .with_name("my_counter")
4163 .with_cardinality_limit(cardinality_limit)
4164 .build()
4165 .unwrap(),
4166 )
4167 } else {
4168 None
4169 }
4170 };
4171 let mut test_context = TestContext::new_with_view(temporality, view_change_cardinality);
4172 let counter = test_context.u64_counter("test", "my_counter", None);
4173
4174 for v in 0..cardinality_limit {
4177 counter.add(100, &[KeyValue::new("A", v.to_string())]);
4178 }
4179
4180 counter.add(3, &[]);
4182 counter.add(3, &[]);
4183
4184 counter.add(100, &[KeyValue::new("A", "foo")]);
4186 counter.add(100, &[KeyValue::new("A", "another")]);
4187 counter.add(100, &[KeyValue::new("A", "yet_another")]);
4188 test_context.flush_metrics();
4189
4190 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4191 unreachable!()
4192 };
4193
4194 assert_eq!(sum.data_points.len(), cardinality_limit + 1 + 1);
4196
4197 let data_point =
4198 find_overflow_sum_datapoint(&sum.data_points).expect("overflow point expected");
4199 assert_eq!(data_point.value, 300);
4200
4201 let empty_attrs_data_point = find_sum_datapoint_with_no_attributes(&sum.data_points)
4203 .expect("Empty attributes point expected");
4204 assert!(
4205 empty_attrs_data_point.attributes.is_empty(),
4206 "Non-empty attribute set"
4207 );
4208 assert_eq!(
4209 empty_attrs_data_point.value, 6,
4210 "Empty attributes value should be 3+3=6"
4211 );
4212
4213 test_context.reset_metrics();
4218 if temporality == Temporality::Delta {
4219 test_context.flush_metrics();
4220 test_context.reset_metrics();
4221 }
4222 counter.add(100, &[KeyValue::new("A", "foo")]);
4225 counter.add(100, &[KeyValue::new("A", "another")]);
4226 counter.add(100, &[KeyValue::new("A", "yet_another")]);
4227 test_context.flush_metrics();
4228
4229 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4230 unreachable!()
4231 };
4232
4233 if temporality == Temporality::Delta {
4234 assert_eq!(sum.data_points.len(), 3);
4235
4236 let data_point = find_sum_datapoint_with_key_value(&sum.data_points, "A", "foo")
4237 .expect("point expected");
4238 assert_eq!(data_point.value, 100);
4239
4240 let data_point = find_sum_datapoint_with_key_value(&sum.data_points, "A", "another")
4241 .expect("point expected");
4242 assert_eq!(data_point.value, 100);
4243
4244 let data_point =
4245 find_sum_datapoint_with_key_value(&sum.data_points, "A", "yet_another")
4246 .expect("point expected");
4247 assert_eq!(data_point.value, 100);
4248 } else {
4249 assert_eq!(sum.data_points.len(), cardinality_limit + 1 + 1);
4251 let data_point =
4252 find_overflow_sum_datapoint(&sum.data_points).expect("overflow point expected");
4253 assert_eq!(data_point.value, 600);
4254
4255 let data_point = find_sum_datapoint_with_key_value(&sum.data_points, "A", "foo");
4256 assert!(data_point.is_none(), "point should not be present");
4257
4258 let data_point = find_sum_datapoint_with_key_value(&sum.data_points, "A", "another");
4259 assert!(data_point.is_none(), "point should not be present");
4260
4261 let data_point =
4262 find_sum_datapoint_with_key_value(&sum.data_points, "A", "yet_another");
4263 assert!(data_point.is_none(), "point should not be present");
4264 }
4265 }
4266
4267 fn counter_aggregation_attribute_order_helper(temporality: Temporality, start_sorted: bool) {
4268 let mut test_context = TestContext::new(temporality);
4270 let counter = test_context.u64_counter("test", "my_counter", None);
4271
4272 if start_sorted {
4277 counter.add(
4278 1,
4279 &[
4280 KeyValue::new("A", "a"),
4281 KeyValue::new("B", "b"),
4282 KeyValue::new("C", "c"),
4283 ],
4284 );
4285 } else {
4286 counter.add(
4287 1,
4288 &[
4289 KeyValue::new("A", "a"),
4290 KeyValue::new("C", "c"),
4291 KeyValue::new("B", "b"),
4292 ],
4293 );
4294 }
4295
4296 counter.add(
4297 1,
4298 &[
4299 KeyValue::new("A", "a"),
4300 KeyValue::new("C", "c"),
4301 KeyValue::new("B", "b"),
4302 ],
4303 );
4304 counter.add(
4305 1,
4306 &[
4307 KeyValue::new("B", "b"),
4308 KeyValue::new("A", "a"),
4309 KeyValue::new("C", "c"),
4310 ],
4311 );
4312 counter.add(
4313 1,
4314 &[
4315 KeyValue::new("B", "b"),
4316 KeyValue::new("C", "c"),
4317 KeyValue::new("A", "a"),
4318 ],
4319 );
4320 counter.add(
4321 1,
4322 &[
4323 KeyValue::new("C", "c"),
4324 KeyValue::new("B", "b"),
4325 KeyValue::new("A", "a"),
4326 ],
4327 );
4328 counter.add(
4329 1,
4330 &[
4331 KeyValue::new("C", "c"),
4332 KeyValue::new("A", "a"),
4333 KeyValue::new("B", "b"),
4334 ],
4335 );
4336 test_context.flush_metrics();
4337
4338 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4339 unreachable!()
4340 };
4341
4342 assert_eq!(sum.data_points.len(), 1);
4344
4345 let data_point1 = &sum.data_points[0];
4347 assert_eq!(data_point1.value, 6);
4348 }
4349
4350 fn updown_counter_aggregation_helper(temporality: Temporality) {
4351 let mut test_context = TestContext::new(temporality);
4353 let counter = test_context.i64_up_down_counter("test", "my_updown_counter", None);
4354
4355 counter.add(10, &[KeyValue::new("key1", "value1")]);
4357 counter.add(-1, &[KeyValue::new("key1", "value1")]);
4358 counter.add(-5, &[KeyValue::new("key1", "value1")]);
4359 counter.add(0, &[KeyValue::new("key1", "value1")]);
4360 counter.add(1, &[KeyValue::new("key1", "value1")]);
4361
4362 counter.add(10, &[KeyValue::new("key1", "value2")]);
4363 counter.add(0, &[KeyValue::new("key1", "value2")]);
4364 counter.add(-3, &[KeyValue::new("key1", "value2")]);
4365
4366 test_context.flush_metrics();
4367
4368 let MetricData::Sum(sum) = test_context.get_aggregation::<i64>("my_updown_counter", None)
4370 else {
4371 unreachable!()
4372 };
4373 assert_eq!(sum.data_points.len(), 2);
4375 assert!(
4376 !sum.is_monotonic,
4377 "UpDownCounter should produce non-monotonic."
4378 );
4379 assert_eq!(
4380 sum.temporality,
4381 Temporality::Cumulative,
4382 "Should produce Cumulative for UpDownCounter"
4383 );
4384
4385 let data_point1 = find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value1")
4387 .expect("datapoint with key1=value1 expected");
4388 assert_eq!(data_point1.value, 5);
4389
4390 let data_point1 = find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value2")
4391 .expect("datapoint with key1=value2 expected");
4392 assert_eq!(data_point1.value, 7);
4393
4394 test_context.reset_metrics();
4396 counter.add(10, &[KeyValue::new("key1", "value1")]);
4397 counter.add(-1, &[KeyValue::new("key1", "value1")]);
4398 counter.add(-5, &[KeyValue::new("key1", "value1")]);
4399 counter.add(0, &[KeyValue::new("key1", "value1")]);
4400 counter.add(1, &[KeyValue::new("key1", "value1")]);
4401
4402 counter.add(10, &[KeyValue::new("key1", "value2")]);
4403 counter.add(0, &[KeyValue::new("key1", "value2")]);
4404 counter.add(-3, &[KeyValue::new("key1", "value2")]);
4405
4406 test_context.flush_metrics();
4407
4408 let MetricData::Sum(sum) = test_context.get_aggregation::<i64>("my_updown_counter", None)
4409 else {
4410 unreachable!()
4411 };
4412 assert_eq!(sum.data_points.len(), 2);
4413 let data_point1 = find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value1")
4414 .expect("datapoint with key1=value1 expected");
4415 assert_eq!(data_point1.value, 10);
4416
4417 let data_point1 = find_sum_datapoint_with_key_value(&sum.data_points, "key1", "value2")
4418 .expect("datapoint with key1=value2 expected");
4419 assert_eq!(data_point1.value, 14);
4420 }
4421
4422 fn find_sum_datapoint_with_key_value<'a, T>(
4423 data_points: &'a [SumDataPoint<T>],
4424 key: &str,
4425 value: &str,
4426 ) -> Option<&'a SumDataPoint<T>> {
4427 data_points.iter().find(|&datapoint| {
4428 datapoint
4429 .attributes
4430 .iter()
4431 .any(|kv| kv.key.as_str() == key && kv.value.as_str() == value)
4432 })
4433 }
4434
4435 fn find_overflow_sum_datapoint<T>(data_points: &[SumDataPoint<T>]) -> Option<&SumDataPoint<T>> {
4436 data_points.iter().find(|&datapoint| {
4437 datapoint.attributes.iter().any(|kv| {
4438 kv.key.as_str() == "otel.metric.overflow" && kv.value == Value::Bool(true)
4439 })
4440 })
4441 }
4442
4443 fn find_gauge_datapoint_with_key_value<'a, T>(
4444 data_points: &'a [GaugeDataPoint<T>],
4445 key: &str,
4446 value: &str,
4447 ) -> Option<&'a GaugeDataPoint<T>> {
4448 data_points.iter().find(|&datapoint| {
4449 datapoint
4450 .attributes
4451 .iter()
4452 .any(|kv| kv.key.as_str() == key && kv.value.as_str() == value)
4453 })
4454 }
4455
4456 fn find_sum_datapoint_with_no_attributes<T>(
4457 data_points: &[SumDataPoint<T>],
4458 ) -> Option<&SumDataPoint<T>> {
4459 data_points
4460 .iter()
4461 .find(|&datapoint| datapoint.attributes.is_empty())
4462 }
4463
4464 fn find_gauge_datapoint_with_no_attributes<T>(
4465 data_points: &[GaugeDataPoint<T>],
4466 ) -> Option<&GaugeDataPoint<T>> {
4467 data_points
4468 .iter()
4469 .find(|&datapoint| datapoint.attributes.is_empty())
4470 }
4471
4472 fn find_histogram_datapoint_with_key_value<'a, T>(
4473 data_points: &'a [HistogramDataPoint<T>],
4474 key: &str,
4475 value: &str,
4476 ) -> Option<&'a HistogramDataPoint<T>> {
4477 data_points.iter().find(|&datapoint| {
4478 datapoint
4479 .attributes
4480 .iter()
4481 .any(|kv| kv.key.as_str() == key && kv.value.as_str() == value)
4482 })
4483 }
4484
4485 fn find_histogram_datapoint_with_no_attributes<T>(
4486 data_points: &[HistogramDataPoint<T>],
4487 ) -> Option<&HistogramDataPoint<T>> {
4488 data_points
4489 .iter()
4490 .find(|&datapoint| datapoint.attributes.is_empty())
4491 }
4492
4493 #[cfg(feature = "experimental_metrics_bound_instruments")]
4494 fn find_overflow_histogram_datapoint<T>(
4495 data_points: &[HistogramDataPoint<T>],
4496 ) -> Option<&HistogramDataPoint<T>> {
4497 data_points.iter().find(|&datapoint| {
4498 datapoint.attributes.iter().any(|kv| {
4499 kv.key.as_str() == "otel.metric.overflow" && kv.value == Value::Bool(true)
4500 })
4501 })
4502 }
4503
4504 #[cfg(feature = "experimental_metrics_bound_instruments")]
4505 fn find_overflow_exponential_histogram_datapoint<T>(
4506 data_points: &[ExponentialHistogramDataPoint<T>],
4507 ) -> Option<&ExponentialHistogramDataPoint<T>> {
4508 data_points.iter().find(|&datapoint| {
4509 datapoint.attributes.iter().any(|kv| {
4510 kv.key.as_str() == "otel.metric.overflow" && kv.value == Value::Bool(true)
4511 })
4512 })
4513 }
4514
4515 #[cfg(feature = "experimental_metrics_bound_instruments")]
4516 fn find_overflow_gauge_datapoint<T>(
4517 data_points: &[GaugeDataPoint<T>],
4518 ) -> Option<&GaugeDataPoint<T>> {
4519 data_points.iter().find(|&datapoint| {
4520 datapoint.attributes.iter().any(|kv| {
4521 kv.key.as_str() == "otel.metric.overflow" && kv.value == Value::Bool(true)
4522 })
4523 })
4524 }
4525
4526 fn find_scope_metric<'a>(
4527 metrics: &'a [ScopeMetrics],
4528 name: &'a str,
4529 ) -> Option<&'a ScopeMetrics> {
4530 metrics
4531 .iter()
4532 .find(|&scope_metric| scope_metric.scope.name() == name)
4533 }
4534
4535 struct TestContext {
4536 exporter: InMemoryMetricExporter,
4537 meter_provider: SdkMeterProvider,
4538
4539 resource_metrics: Vec<ResourceMetrics>,
4541 }
4542
4543 impl TestContext {
4544 fn new(temporality: Temporality) -> Self {
4545 let exporter = InMemoryMetricExporterBuilder::new().with_temporality(temporality);
4546 let exporter = exporter.build();
4547 let meter_provider = SdkMeterProvider::builder()
4548 .with_periodic_exporter(exporter.clone())
4549 .build();
4550
4551 TestContext {
4552 exporter,
4553 meter_provider,
4554 resource_metrics: vec![],
4555 }
4556 }
4557
4558 fn new_with_view<T>(temporality: Temporality, view: T) -> Self
4559 where
4560 T: Fn(&Instrument) -> Option<Stream> + Send + Sync + 'static,
4561 {
4562 let exporter = InMemoryMetricExporterBuilder::new().with_temporality(temporality);
4563 let exporter = exporter.build();
4564 let meter_provider = SdkMeterProvider::builder()
4565 .with_periodic_exporter(exporter.clone())
4566 .with_view(view)
4567 .build();
4568
4569 TestContext {
4570 exporter,
4571 meter_provider,
4572 resource_metrics: vec![],
4573 }
4574 }
4575
4576 fn u64_counter(
4577 &self,
4578 meter_name: &'static str,
4579 counter_name: &'static str,
4580 unit: Option<&'static str>,
4581 ) -> Counter<u64> {
4582 let meter = self.meter_provider.meter(meter_name);
4583 let mut counter_builder = meter.u64_counter(counter_name);
4584 if let Some(unit_name) = unit {
4585 counter_builder = counter_builder.with_unit(unit_name);
4586 }
4587 counter_builder.build()
4588 }
4589
4590 fn i64_up_down_counter(
4591 &self,
4592 meter_name: &'static str,
4593 counter_name: &'static str,
4594 unit: Option<&'static str>,
4595 ) -> UpDownCounter<i64> {
4596 let meter = self.meter_provider.meter(meter_name);
4597 let mut updown_counter_builder = meter.i64_up_down_counter(counter_name);
4598 if let Some(unit_name) = unit {
4599 updown_counter_builder = updown_counter_builder.with_unit(unit_name);
4600 }
4601 updown_counter_builder.build()
4602 }
4603
4604 fn meter(&self) -> Meter {
4605 self.meter_provider.meter("test")
4606 }
4607
4608 fn flush_metrics(&self) {
4609 self.meter_provider.force_flush().unwrap();
4610 }
4611
4612 fn reset_metrics(&self) {
4613 self.exporter.reset();
4614 }
4615
4616 fn check_no_metrics(&self) {
4617 let resource_metrics = self
4618 .exporter
4619 .get_finished_metrics()
4620 .expect("metrics expected to be exported"); assert!(resource_metrics.is_empty(), "no metrics should be exported");
4623 }
4624
4625 fn get_aggregation<T: Number>(
4626 &mut self,
4627 counter_name: &str,
4628 unit_name: Option<&str>,
4629 ) -> &MetricData<T> {
4630 self.resource_metrics = self
4631 .exporter
4632 .get_finished_metrics()
4633 .expect("metrics expected to be exported");
4634
4635 assert!(
4636 !self.resource_metrics.is_empty(),
4637 "no metrics were exported"
4638 );
4639
4640 assert!(
4641 self.resource_metrics.len() == 1,
4642 "Expected single resource metrics."
4643 );
4644 let resource_metric = self
4645 .resource_metrics
4646 .first()
4647 .expect("This should contain exactly one resource metric, as validated above.");
4648
4649 assert!(
4650 !resource_metric.scope_metrics.is_empty(),
4651 "No scope metrics in latest export"
4652 );
4653 assert!(!resource_metric.scope_metrics[0].metrics.is_empty());
4654
4655 let metric = &resource_metric.scope_metrics[0].metrics[0];
4656 assert_eq!(metric.name, counter_name);
4657 if let Some(expected_unit) = unit_name {
4658 assert_eq!(metric.unit, expected_unit);
4659 }
4660
4661 T::extract_metrics_data_ref(&metric.data)
4662 .expect("Failed to cast aggregation to expected type")
4663 }
4664
4665 fn get_from_multiple_aggregations<T: Number>(
4666 &mut self,
4667 counter_name: &str,
4668 unit_name: Option<&str>,
4669 invocation_count: usize,
4670 ) -> Vec<&MetricData<T>> {
4671 self.resource_metrics = self
4672 .exporter
4673 .get_finished_metrics()
4674 .expect("metrics expected to be exported");
4675
4676 assert!(
4677 !self.resource_metrics.is_empty(),
4678 "no metrics were exported"
4679 );
4680
4681 assert_eq!(
4682 self.resource_metrics.len(),
4683 invocation_count,
4684 "Expected collect to be called {invocation_count} times"
4685 );
4686
4687 let result = self
4688 .resource_metrics
4689 .iter()
4690 .map(|resource_metric| {
4691 assert!(
4692 !resource_metric.scope_metrics.is_empty(),
4693 "An export with no scope metrics occurred"
4694 );
4695
4696 assert!(!resource_metric.scope_metrics[0].metrics.is_empty());
4697
4698 let metric = &resource_metric.scope_metrics[0].metrics[0];
4699 assert_eq!(metric.name, counter_name);
4700
4701 if let Some(expected_unit) = unit_name {
4702 assert_eq!(metric.unit, expected_unit);
4703 }
4704
4705 let aggregation = T::extract_metrics_data_ref(&metric.data)
4706 .expect("Failed to cast aggregation to expected type");
4707 aggregation
4708 })
4709 .collect::<Vec<_>>();
4710
4711 result
4712 }
4713 }
4714
4715 #[test]
4716 fn parse_valid_temporality_values() {
4717 assert_eq!(
4718 "cumulative".parse::<Temporality>(),
4719 Ok(Temporality::Cumulative)
4720 );
4721 assert_eq!("delta".parse::<Temporality>(), Ok(Temporality::Delta));
4722 assert_eq!(
4723 "lowmemory".parse::<Temporality>(),
4724 Ok(Temporality::LowMemory)
4725 );
4726 }
4727
4728 #[test]
4729 fn parse_temporality_case_insensitive() {
4730 assert_eq!(
4731 "Cumulative".parse::<Temporality>(),
4732 Ok(Temporality::Cumulative)
4733 );
4734 assert_eq!("DELTA".parse::<Temporality>(), Ok(Temporality::Delta));
4735 assert_eq!(
4736 "LowMemory".parse::<Temporality>(),
4737 Ok(Temporality::LowMemory)
4738 );
4739 assert_eq!(
4740 "LOWMEMORY".parse::<Temporality>(),
4741 Ok(Temporality::LowMemory)
4742 );
4743 }
4744
4745 #[test]
4746 fn parse_invalid_temporality_returns_err() {
4747 assert!("unknown".parse::<Temporality>().is_err());
4748 assert!("".parse::<Temporality>().is_err());
4749 assert!("cumulativ".parse::<Temporality>().is_err());
4750 }
4751
4752 #[cfg(feature = "experimental_metrics_bound_instruments")]
4753 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
4754 async fn bound_counter_cumulative() {
4755 let mut test_context = TestContext::new(Temporality::Cumulative);
4756 let counter = test_context.u64_counter("test", "my_counter", None);
4757 let attrs = vec![KeyValue::new("key1", "bound_value")];
4758 let bound = counter.bind(&attrs);
4759
4760 bound.add(10);
4761 bound.add(20);
4762 bound.add(30);
4763 test_context.flush_metrics();
4764
4765 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4766 unreachable!()
4767 };
4768
4769 assert_eq!(sum.data_points.len(), 1, "Expected one data point");
4770 assert!(sum.is_monotonic);
4771 assert_eq!(sum.temporality, Temporality::Cumulative);
4772
4773 let data_point = &sum.data_points[0];
4774 assert_eq!(data_point.value, 60);
4775 assert_eq!(
4776 data_point.attributes,
4777 vec![KeyValue::new("key1", "bound_value")]
4778 );
4779 }
4780
4781 #[cfg(feature = "experimental_metrics_bound_instruments")]
4782 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
4783 async fn bound_counter_delta() {
4784 let mut test_context = TestContext::new(Temporality::Delta);
4785 let counter = test_context.u64_counter("test", "my_counter", None);
4786 let attrs = vec![KeyValue::new("key1", "bound_value")];
4787 let bound = counter.bind(&attrs);
4788
4789 bound.add(50);
4790 test_context.flush_metrics();
4791
4792 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4793 unreachable!()
4794 };
4795 assert_eq!(sum.temporality, Temporality::Delta);
4796 assert_eq!(sum.data_points.len(), 1);
4797 assert_eq!(sum.data_points[0].value, 50);
4798
4799 test_context.reset_metrics();
4801 bound.add(25);
4802 test_context.flush_metrics();
4803
4804 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4805 unreachable!()
4806 };
4807 assert_eq!(sum.data_points.len(), 1);
4808 assert_eq!(
4809 sum.data_points[0].value, 25,
4810 "Delta should reset between collections"
4811 );
4812 }
4813
4814 #[cfg(feature = "experimental_metrics_bound_instruments")]
4815 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
4816 async fn bound_histogram_cumulative() {
4817 let mut test_context = TestContext::new(Temporality::Cumulative);
4818 let histogram = test_context
4819 .meter()
4820 .f64_histogram("my_histogram")
4821 .with_boundaries(vec![5.0, 10.0, 25.0, 50.0])
4822 .build();
4823 let attrs = vec![KeyValue::new("key1", "bound_value")];
4824 let bound = histogram.bind(&attrs);
4825
4826 bound.record(1.0);
4827 bound.record(7.5);
4828 bound.record(15.0);
4829 bound.record(30.0);
4830 test_context.flush_metrics();
4831
4832 let MetricData::Histogram(histogram_data) =
4833 test_context.get_aggregation::<f64>("my_histogram", None)
4834 else {
4835 unreachable!()
4836 };
4837
4838 assert_eq!(histogram_data.data_points.len(), 1);
4839 assert_eq!(histogram_data.temporality, Temporality::Cumulative);
4840
4841 let dp = &histogram_data.data_points[0];
4842 assert_eq!(dp.count, 4);
4843 assert_eq!(dp.sum, 53.5);
4844 assert_eq!(dp.min.unwrap(), 1.0);
4845 assert_eq!(dp.max.unwrap(), 30.0);
4846 assert_eq!(dp.attributes, vec![KeyValue::new("key1", "bound_value")]);
4847 }
4848
4849 #[cfg(feature = "experimental_metrics_bound_instruments")]
4850 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
4851 async fn bound_counter_matches_unbound() {
4852 let mut test_context = TestContext::new(Temporality::Cumulative);
4853 let counter = test_context.u64_counter("test", "my_counter", None);
4854 let attrs = vec![KeyValue::new("key1", "shared")];
4855 let bound = counter.bind(&attrs);
4856
4857 counter.add(10, &attrs);
4859 bound.add(20);
4860 counter.add(30, &attrs);
4861 bound.add(40);
4862 test_context.flush_metrics();
4863
4864 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4865 unreachable!()
4866 };
4867
4868 assert_eq!(
4869 sum.data_points.len(),
4870 1,
4871 "Bound and unbound should share the same data point"
4872 );
4873 assert_eq!(sum.data_points[0].value, 100);
4874 }
4875
4876 #[cfg(feature = "experimental_metrics_bound_instruments")]
4877 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
4878 async fn bound_counter_delta_no_update_no_export() {
4879 let mut test_context = TestContext::new(Temporality::Delta);
4880 let counter = test_context.u64_counter("test", "my_counter", None);
4881 let attrs = vec![KeyValue::new("key1", "bound_value")];
4882 let bound = counter.bind(&attrs);
4883
4884 bound.add(10);
4886 test_context.flush_metrics();
4887 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4888 unreachable!()
4889 };
4890 assert_eq!(sum.data_points.len(), 1);
4891 assert_eq!(sum.data_points[0].value, 10);
4892
4893 test_context.reset_metrics();
4895 test_context.flush_metrics();
4896 let resource_metrics = test_context
4897 .exporter
4898 .get_finished_metrics()
4899 .expect("metrics export should succeed");
4900 assert!(
4901 resource_metrics.is_empty(),
4902 "Bound handle with no updates should not export"
4903 );
4904
4905 test_context.reset_metrics();
4907 bound.add(5);
4908 test_context.flush_metrics();
4909 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4910 unreachable!()
4911 };
4912 assert_eq!(sum.data_points.len(), 1);
4913 assert_eq!(
4914 sum.data_points[0].value, 5,
4915 "Bound handle should produce fresh delta after quiet cycle"
4916 );
4917 }
4918
4919 #[cfg(feature = "experimental_metrics_bound_instruments")]
4920 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
4921 async fn bound_counter_at_overflow_attributes_to_overflow_bucket() {
4922 let cardinality_limit = 3;
4923 let view = move |i: &Instrument| {
4924 if i.name() == "my_counter" {
4925 Stream::builder()
4926 .with_name("my_counter")
4927 .with_cardinality_limit(cardinality_limit)
4928 .build()
4929 .ok()
4930 } else {
4931 None
4932 }
4933 };
4934 let mut test_context = TestContext::new_with_view(Temporality::Delta, view);
4935 let counter = test_context.u64_counter("test", "my_counter", None);
4936
4937 for v in 0..cardinality_limit {
4939 counter.add(1, &[KeyValue::new("A", v.to_string())]);
4940 }
4941
4942 let overflow_attrs = vec![KeyValue::new("A", "overflow_bind")];
4944 let bound = counter.bind(&overflow_attrs);
4945 bound.add(42);
4946
4947 test_context.flush_metrics();
4948 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4949 unreachable!()
4950 };
4951
4952 assert_eq!(
4954 sum.data_points.len(),
4955 cardinality_limit + 1,
4956 "Expected {} unique + 1 overflow data points",
4957 cardinality_limit
4958 );
4959
4960 let overflow_dp =
4962 find_overflow_sum_datapoint(&sum.data_points).expect("overflow point expected");
4963 assert_eq!(
4964 overflow_dp.value, 42,
4965 "Bound-at-overflow data should go to overflow bucket"
4966 );
4967 }
4968
4969 #[cfg(feature = "experimental_metrics_bound_instruments")]
4970 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
4971 async fn bound_counter_overflow_recovery_after_delta_eviction() {
4972 let cardinality_limit = 3;
4973 let view = move |i: &Instrument| {
4974 if i.name() == "my_counter" {
4975 Stream::builder()
4976 .with_name("my_counter")
4977 .with_cardinality_limit(cardinality_limit)
4978 .build()
4979 .ok()
4980 } else {
4981 None
4982 }
4983 };
4984 let mut test_context = TestContext::new_with_view(Temporality::Delta, view);
4985 let counter = test_context.u64_counter("test", "my_counter", None);
4986
4987 for v in 0..cardinality_limit {
4989 counter.add(1, &[KeyValue::new("A", v.to_string())]);
4990 }
4991
4992 test_context.flush_metrics();
4994 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
4995 unreachable!()
4996 };
4997 assert_eq!(sum.data_points.len(), cardinality_limit);
4998
4999 test_context.reset_metrics();
5002 test_context.flush_metrics(); let new_attrs = vec![KeyValue::new("A", "recovered")];
5005 let bound = counter.bind(&new_attrs);
5006 bound.add(99);
5007
5008 test_context.reset_metrics();
5009 test_context.flush_metrics();
5010 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
5011 unreachable!()
5012 };
5013
5014 assert_eq!(
5016 sum.data_points.len(),
5017 1,
5018 "Only the bound entry should be exported"
5019 );
5020 let dp = find_sum_datapoint_with_key_value(&sum.data_points, "A", "recovered")
5021 .expect("should find dedicated data point for recovered attrs");
5022 assert_eq!(
5023 dp.value, 99,
5024 "Bound handle after recovery should have dedicated tracker"
5025 );
5026 assert!(
5027 find_overflow_sum_datapoint(&sum.data_points).is_none(),
5028 "Should not have overflow — bind() after eviction should get a dedicated tracker"
5029 );
5030 }
5031
5032 #[cfg(feature = "experimental_metrics_bound_instruments")]
5033 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5034 async fn bound_counter_multiple_overflow_handles_share_overflow_bucket() {
5035 let cardinality_limit = 2;
5036 let view = move |i: &Instrument| {
5037 if i.name() == "my_counter" {
5038 Stream::builder()
5039 .with_name("my_counter")
5040 .with_cardinality_limit(cardinality_limit)
5041 .build()
5042 .ok()
5043 } else {
5044 None
5045 }
5046 };
5047 let mut test_context = TestContext::new_with_view(Temporality::Delta, view);
5048 let counter = test_context.u64_counter("test", "my_counter", None);
5049
5050 counter.add(1, &[KeyValue::new("A", "0")]);
5052 counter.add(1, &[KeyValue::new("A", "1")]);
5053
5054 let bound_a = counter.bind(&[KeyValue::new("A", "overflow_a")]);
5056 let bound_b = counter.bind(&[KeyValue::new("A", "overflow_b")]);
5057
5058 bound_a.add(10);
5059 bound_b.add(20);
5060 bound_a.add(5);
5061
5062 test_context.flush_metrics();
5063 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
5064 unreachable!()
5065 };
5066
5067 let overflow_dp =
5068 find_overflow_sum_datapoint(&sum.data_points).expect("overflow point expected");
5069 assert_eq!(
5070 overflow_dp.value, 35,
5071 "All overflow-bound measurements should accumulate in overflow bucket"
5072 );
5073
5074 test_context.reset_metrics();
5076 bound_a.add(7);
5077 bound_b.add(3);
5078 test_context.flush_metrics();
5079
5080 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
5081 unreachable!()
5082 };
5083
5084 let overflow_dp =
5085 find_overflow_sum_datapoint(&sum.data_points).expect("overflow point expected");
5086 assert_eq!(
5087 overflow_dp.value, 10,
5088 "Overflow-bound handles should continue working across delta cycles"
5089 );
5090 }
5091
5092 #[cfg(feature = "experimental_metrics_bound_instruments")]
5093 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5094 async fn bound_counter_overflow_persists_across_eviction_cycles() {
5095 let cardinality_limit = 3;
5102 let view = move |i: &Instrument| {
5103 if i.name() == "my_counter" {
5104 Stream::builder()
5105 .with_name("my_counter")
5106 .with_cardinality_limit(cardinality_limit)
5107 .build()
5108 .ok()
5109 } else {
5110 None
5111 }
5112 };
5113 let mut test_context = TestContext::new_with_view(Temporality::Delta, view);
5114 let counter = test_context.u64_counter("test", "my_counter", None);
5115
5116 for v in 0..cardinality_limit {
5118 counter.add(1, &[KeyValue::new("A", v.to_string())]);
5119 }
5120 let stuck_attrs = vec![KeyValue::new("A", "stuck_in_overflow")];
5121 let bound = counter.bind(&stuck_attrs);
5122 bound.add(10);
5123
5124 test_context.flush_metrics();
5125 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
5126 unreachable!()
5127 };
5128 let overflow_dp = find_overflow_sum_datapoint(&sum.data_points)
5129 .expect("cycle 1: bound write at overflow should land in overflow bucket");
5130 assert_eq!(overflow_dp.value, 10);
5131
5132 test_context.reset_metrics();
5135 test_context.flush_metrics();
5136
5137 test_context.reset_metrics();
5141 bound.add(99);
5142 test_context.flush_metrics();
5143 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
5144 unreachable!()
5145 };
5146 let overflow_dp = find_overflow_sum_datapoint(&sum.data_points)
5147 .expect("cycle 3: bound write must still land in overflow even after space frees up");
5148 assert_eq!(
5149 overflow_dp.value, 99,
5150 "Bound-at-overflow handle must keep writing to overflow even after delta eviction"
5151 );
5152 assert!(
5153 find_sum_datapoint_with_key_value(&sum.data_points, "A", "stuck_in_overflow").is_none(),
5154 "Bound-at-overflow handle must not silently self-heal to a dedicated tracker"
5155 );
5156 }
5157
5158 #[cfg(feature = "experimental_metrics_bound_instruments")]
5159 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5160 async fn bound_histogram_delta() {
5161 let mut test_context = TestContext::new(Temporality::Delta);
5162 let histogram = test_context
5163 .meter()
5164 .f64_histogram("my_histogram")
5165 .with_boundaries(vec![5.0, 10.0, 25.0, 50.0])
5166 .build();
5167 let attrs = vec![KeyValue::new("key1", "bound_value")];
5168 let bound = histogram.bind(&attrs);
5169
5170 bound.record(3.0);
5172 bound.record(12.0);
5173 test_context.flush_metrics();
5174
5175 let MetricData::Histogram(hist) = test_context.get_aggregation::<f64>("my_histogram", None)
5176 else {
5177 unreachable!()
5178 };
5179 assert_eq!(hist.temporality, Temporality::Delta);
5180 assert_eq!(hist.data_points.len(), 1);
5181 assert_eq!(hist.data_points[0].count, 2);
5182 assert_eq!(hist.data_points[0].sum, 15.0);
5183
5184 test_context.reset_metrics();
5186 bound.record(40.0);
5187 test_context.flush_metrics();
5188
5189 let MetricData::Histogram(hist) = test_context.get_aggregation::<f64>("my_histogram", None)
5190 else {
5191 unreachable!()
5192 };
5193 assert_eq!(hist.data_points.len(), 1);
5194 assert_eq!(
5195 hist.data_points[0].count, 1,
5196 "Delta should reset count between collections"
5197 );
5198 assert_eq!(
5199 hist.data_points[0].sum, 40.0,
5200 "Delta should reset sum between collections"
5201 );
5202 }
5203
5204 #[cfg(feature = "experimental_metrics_bound_instruments")]
5205 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5206 async fn bound_histogram_matches_unbound() {
5207 let mut test_context = TestContext::new(Temporality::Cumulative);
5208 let histogram = test_context
5209 .meter()
5210 .f64_histogram("my_histogram")
5211 .with_boundaries(vec![10.0, 50.0])
5212 .build();
5213 let attrs = vec![KeyValue::new("key1", "shared")];
5214 let bound = histogram.bind(&attrs);
5215
5216 histogram.record(5.0, &attrs);
5218 bound.record(15.0);
5219 histogram.record(25.0, &attrs);
5220 bound.record(35.0);
5221 test_context.flush_metrics();
5222
5223 let MetricData::Histogram(hist) = test_context.get_aggregation::<f64>("my_histogram", None)
5224 else {
5225 unreachable!()
5226 };
5227
5228 assert_eq!(
5229 hist.data_points.len(),
5230 1,
5231 "Bound and unbound should share the same data point"
5232 );
5233 assert_eq!(hist.data_points[0].count, 4);
5234 assert_eq!(hist.data_points[0].sum, 80.0);
5235 assert_eq!(hist.data_points[0].min.unwrap(), 5.0);
5236 assert_eq!(hist.data_points[0].max.unwrap(), 35.0);
5237 }
5238
5239 #[cfg(feature = "experimental_metrics_bound_instruments")]
5240 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5241 async fn bound_histogram_delta_no_update_no_export() {
5242 let mut test_context = TestContext::new(Temporality::Delta);
5243 let histogram = test_context
5244 .meter()
5245 .f64_histogram("my_histogram")
5246 .with_boundaries(vec![10.0])
5247 .build();
5248 let attrs = vec![KeyValue::new("key1", "bound_value")];
5249 let bound = histogram.bind(&attrs);
5250
5251 bound.record(5.0);
5253 test_context.flush_metrics();
5254 let MetricData::Histogram(hist) = test_context.get_aggregation::<f64>("my_histogram", None)
5255 else {
5256 unreachable!()
5257 };
5258 assert_eq!(hist.data_points.len(), 1);
5259 assert_eq!(hist.data_points[0].count, 1);
5260
5261 test_context.reset_metrics();
5263 test_context.flush_metrics();
5264 let resource_metrics = test_context
5265 .exporter
5266 .get_finished_metrics()
5267 .expect("metrics export should succeed");
5268 assert!(
5269 resource_metrics.is_empty(),
5270 "Bound histogram with no updates should not export"
5271 );
5272
5273 test_context.reset_metrics();
5275 bound.record(20.0);
5276 test_context.flush_metrics();
5277 let MetricData::Histogram(hist) = test_context.get_aggregation::<f64>("my_histogram", None)
5278 else {
5279 unreachable!()
5280 };
5281 assert_eq!(hist.data_points.len(), 1);
5282 assert_eq!(
5283 hist.data_points[0].count, 1,
5284 "Fresh delta after quiet cycle"
5285 );
5286 assert_eq!(hist.data_points[0].sum, 20.0);
5287 }
5288
5289 #[cfg(feature = "experimental_metrics_bound_instruments")]
5290 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5291 async fn bound_histogram_at_overflow_attributes_to_overflow_bucket() {
5292 let cardinality_limit = 3;
5293 let view = move |i: &Instrument| {
5294 if i.name() == "my_histogram" {
5295 Stream::builder()
5296 .with_name("my_histogram")
5297 .with_cardinality_limit(cardinality_limit)
5298 .build()
5299 .ok()
5300 } else {
5301 None
5302 }
5303 };
5304 let mut test_context = TestContext::new_with_view(Temporality::Delta, view);
5305 let histogram = test_context
5306 .meter()
5307 .f64_histogram("my_histogram")
5308 .with_boundaries(vec![10.0, 50.0])
5309 .build();
5310
5311 for v in 0..cardinality_limit {
5313 histogram.record(1.0, &[KeyValue::new("A", v.to_string())]);
5314 }
5315
5316 let overflow_attrs = vec![KeyValue::new("A", "overflow_bind")];
5318 let bound = histogram.bind(&overflow_attrs);
5319 bound.record(42.0);
5320
5321 test_context.flush_metrics();
5322 let MetricData::Histogram(hist) = test_context.get_aggregation::<f64>("my_histogram", None)
5323 else {
5324 unreachable!()
5325 };
5326
5327 assert_eq!(
5328 hist.data_points.len(),
5329 cardinality_limit + 1,
5330 "Expected {} unique + 1 overflow data points",
5331 cardinality_limit
5332 );
5333
5334 let overflow_dp =
5335 find_overflow_histogram_datapoint(&hist.data_points).expect("overflow point expected");
5336 assert_eq!(
5337 overflow_dp.sum, 42.0,
5338 "Bound-at-overflow data should go to overflow bucket"
5339 );
5340 assert_eq!(overflow_dp.count, 1);
5341 }
5342
5343 #[cfg(feature = "experimental_metrics_bound_instruments")]
5344 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5345 async fn bound_histogram_overflow_persists_across_eviction_cycles() {
5346 let cardinality_limit = 3;
5350 let view = move |i: &Instrument| {
5351 if i.name() == "my_histogram" {
5352 Stream::builder()
5353 .with_name("my_histogram")
5354 .with_cardinality_limit(cardinality_limit)
5355 .build()
5356 .ok()
5357 } else {
5358 None
5359 }
5360 };
5361 let mut test_context = TestContext::new_with_view(Temporality::Delta, view);
5362 let histogram = test_context
5363 .meter()
5364 .f64_histogram("my_histogram")
5365 .with_boundaries(vec![10.0, 50.0])
5366 .build();
5367
5368 for v in 0..cardinality_limit {
5370 histogram.record(1.0, &[KeyValue::new("A", v.to_string())]);
5371 }
5372 let stuck_attrs = vec![KeyValue::new("A", "stuck_in_overflow")];
5373 let bound = histogram.bind(&stuck_attrs);
5374 bound.record(15.0);
5375
5376 test_context.flush_metrics();
5377 let MetricData::Histogram(hist) = test_context.get_aggregation::<f64>("my_histogram", None)
5378 else {
5379 unreachable!()
5380 };
5381 let overflow_dp = find_overflow_histogram_datapoint(&hist.data_points)
5382 .expect("cycle 1: bound write at overflow should land in overflow bucket");
5383 assert_eq!(overflow_dp.sum, 15.0);
5384
5385 test_context.reset_metrics();
5387 test_context.flush_metrics();
5388
5389 test_context.reset_metrics();
5392 bound.record(99.0);
5393 test_context.flush_metrics();
5394 let MetricData::Histogram(hist) = test_context.get_aggregation::<f64>("my_histogram", None)
5395 else {
5396 unreachable!()
5397 };
5398 let overflow_dp = find_overflow_histogram_datapoint(&hist.data_points)
5399 .expect("cycle 3: bound write must still land in overflow even after space frees up");
5400 assert_eq!(
5401 overflow_dp.sum, 99.0,
5402 "Bound-at-overflow histogram must keep writing to overflow even after delta eviction"
5403 );
5404 assert_eq!(overflow_dp.count, 1);
5405 }
5406
5407 #[cfg(feature = "experimental_metrics_bound_instruments")]
5408 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5409 async fn bound_exponential_histogram_delta() {
5410 let view = |i: &Instrument| {
5414 if i.name() == "my_histogram" {
5415 Stream::builder()
5416 .with_aggregation(Aggregation::Base2ExponentialHistogram {
5417 max_size: 160,
5418 max_scale: 20,
5419 record_min_max: true,
5420 })
5421 .build()
5422 .ok()
5423 } else {
5424 None
5425 }
5426 };
5427 let mut test_context = TestContext::new_with_view(Temporality::Delta, view);
5428 let histogram = test_context.meter().f64_histogram("my_histogram").build();
5429 let attrs = vec![KeyValue::new("key1", "bound_value")];
5430 let bound = histogram.bind(&attrs);
5431
5432 bound.record(2.0);
5433 bound.record(4.0);
5434 bound.record(8.0);
5435 bound.record(f64::NAN);
5437 bound.record(f64::INFINITY);
5438
5439 test_context.flush_metrics();
5440 let MetricData::ExponentialHistogram(hist) =
5441 test_context.get_aggregation::<f64>("my_histogram", None)
5442 else {
5443 panic!("expected ExponentialHistogram aggregation");
5444 };
5445 assert_eq!(hist.data_points.len(), 1);
5446 let dp = &hist.data_points[0];
5447 assert_eq!(dp.count, 3, "NaN and infinity should be dropped");
5448 assert_eq!(dp.sum, 14.0);
5449 assert_eq!(dp.min, Some(2.0));
5450 assert_eq!(dp.max, Some(8.0));
5451 }
5452
5453 #[cfg(feature = "experimental_metrics_bound_instruments")]
5454 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5455 async fn bound_exponential_histogram_at_overflow_attributes_to_overflow_bucket() {
5456 let cardinality_limit = 3;
5457 let view = move |i: &Instrument| {
5458 if i.name() == "my_histogram" {
5459 Stream::builder()
5460 .with_aggregation(Aggregation::Base2ExponentialHistogram {
5461 max_size: 160,
5462 max_scale: 20,
5463 record_min_max: true,
5464 })
5465 .with_cardinality_limit(cardinality_limit)
5466 .build()
5467 .ok()
5468 } else {
5469 None
5470 }
5471 };
5472 let mut test_context = TestContext::new_with_view(Temporality::Delta, view);
5473 let histogram = test_context.meter().f64_histogram("my_histogram").build();
5474
5475 for v in 0..cardinality_limit {
5476 histogram.record(1.0, &[KeyValue::new("A", v.to_string())]);
5477 }
5478
5479 let overflow_attrs = vec![KeyValue::new("A", "overflow_bind")];
5481 let bound = histogram.bind(&overflow_attrs);
5482 bound.record(42.0);
5483
5484 test_context.flush_metrics();
5485 let MetricData::ExponentialHistogram(hist) =
5486 test_context.get_aggregation::<f64>("my_histogram", None)
5487 else {
5488 panic!("expected ExponentialHistogram aggregation");
5489 };
5490 let overflow_dp = find_overflow_exponential_histogram_datapoint(&hist.data_points)
5491 .expect("overflow point expected");
5492 assert_eq!(
5493 overflow_dp.sum, 42.0,
5494 "Bound-at-overflow data should go to overflow bucket"
5495 );
5496 assert_eq!(overflow_dp.count, 1);
5497 }
5498
5499 #[cfg(feature = "experimental_metrics_bound_instruments")]
5500 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5501 async fn bound_exponential_histogram_overflow_persists_across_eviction_cycles() {
5502 let cardinality_limit = 3;
5505 let view = move |i: &Instrument| {
5506 if i.name() == "my_histogram" {
5507 Stream::builder()
5508 .with_aggregation(Aggregation::Base2ExponentialHistogram {
5509 max_size: 160,
5510 max_scale: 20,
5511 record_min_max: true,
5512 })
5513 .with_cardinality_limit(cardinality_limit)
5514 .build()
5515 .ok()
5516 } else {
5517 None
5518 }
5519 };
5520 let mut test_context = TestContext::new_with_view(Temporality::Delta, view);
5521 let histogram = test_context.meter().f64_histogram("my_histogram").build();
5522
5523 for v in 0..cardinality_limit {
5524 histogram.record(1.0, &[KeyValue::new("A", v.to_string())]);
5525 }
5526 let stuck_attrs = vec![KeyValue::new("A", "stuck_in_overflow")];
5527 let bound = histogram.bind(&stuck_attrs);
5528 bound.record(15.0);
5529
5530 test_context.flush_metrics();
5531 let MetricData::ExponentialHistogram(hist) =
5532 test_context.get_aggregation::<f64>("my_histogram", None)
5533 else {
5534 panic!("expected ExponentialHistogram aggregation");
5535 };
5536 let overflow_dp = find_overflow_exponential_histogram_datapoint(&hist.data_points)
5537 .expect("cycle 1: bound write at overflow should land in overflow bucket");
5538 assert_eq!(overflow_dp.sum, 15.0);
5539
5540 test_context.reset_metrics();
5542 test_context.flush_metrics();
5543
5544 test_context.reset_metrics();
5546 bound.record(99.0);
5547 test_context.flush_metrics();
5548 let MetricData::ExponentialHistogram(hist) =
5549 test_context.get_aggregation::<f64>("my_histogram", None)
5550 else {
5551 panic!("expected ExponentialHistogram aggregation");
5552 };
5553 let overflow_dp = find_overflow_exponential_histogram_datapoint(&hist.data_points)
5554 .expect("cycle 3: bound write must still land in overflow even after space frees up");
5555 assert_eq!(
5556 overflow_dp.sum, 99.0,
5557 "Bound-at-overflow ExpoHistogram must keep writing to overflow even after delta eviction"
5558 );
5559 assert_eq!(overflow_dp.count, 1);
5560 }
5561
5562 #[cfg(feature = "experimental_metrics_bound_instruments")]
5563 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5564 async fn bound_counter_drop_enables_eviction() {
5565 let mut test_context = TestContext::new(Temporality::Delta);
5566 let counter = test_context.u64_counter("test", "my_counter", None);
5567 let attrs = vec![KeyValue::new("key1", "ephemeral")];
5568
5569 {
5570 let bound = counter.bind(&attrs);
5571 bound.add(100);
5572 test_context.flush_metrics();
5573
5574 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None)
5575 else {
5576 unreachable!()
5577 };
5578 assert_eq!(sum.data_points.len(), 1);
5579 assert_eq!(sum.data_points[0].value, 100);
5580 }
5582
5583 test_context.reset_metrics();
5585 test_context.flush_metrics();
5586
5587 test_context.reset_metrics();
5590 counter.add(1, &[KeyValue::new("key1", "new_entry")]);
5591 test_context.flush_metrics();
5592
5593 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
5594 unreachable!()
5595 };
5596
5597 assert_eq!(sum.data_points.len(), 1);
5599 let dp = find_sum_datapoint_with_key_value(&sum.data_points, "key1", "new_entry")
5600 .expect("new_entry should be present");
5601 assert_eq!(dp.value, 1);
5602 assert!(
5603 find_sum_datapoint_with_key_value(&sum.data_points, "key1", "ephemeral").is_none(),
5604 "Dropped bound entry should have been evicted"
5605 );
5606 }
5607
5608 #[cfg(feature = "experimental_metrics_bound_instruments")]
5609 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5610 async fn bound_counter_multiple_handles_same_attrs() {
5611 let mut test_context = TestContext::new(Temporality::Delta);
5612 let counter = test_context.u64_counter("test", "my_counter", None);
5613 let attrs = vec![KeyValue::new("key1", "shared")];
5614
5615 let bound1 = counter.bind(&attrs);
5616 let bound2 = counter.bind(&attrs);
5617
5618 bound1.add(10);
5619 bound2.add(20);
5620 test_context.flush_metrics();
5621
5622 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
5623 unreachable!()
5624 };
5625 assert_eq!(
5626 sum.data_points.len(),
5627 1,
5628 "Multiple handles to same attrs should share data point"
5629 );
5630 assert_eq!(sum.data_points[0].value, 30);
5631
5632 drop(bound1);
5634 test_context.reset_metrics();
5635 test_context.flush_metrics(); test_context.reset_metrics();
5639 bound2.add(5);
5640 test_context.flush_metrics();
5641
5642 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
5643 unreachable!()
5644 };
5645 assert_eq!(sum.data_points.len(), 1);
5646 assert_eq!(
5647 sum.data_points[0].value, 5,
5648 "Entry should persist while any handle is alive"
5649 );
5650 }
5651
5652 #[cfg(feature = "experimental_metrics_bound_instruments")]
5653 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5654 async fn bound_counter_empty_attributes() {
5655 let mut test_context = TestContext::new(Temporality::Cumulative);
5656 let counter = test_context.u64_counter("test", "my_counter", None);
5657 let bound = counter.bind(&[]);
5658
5659 bound.add(10);
5660 bound.add(30);
5661 test_context.flush_metrics();
5662
5663 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
5664 unreachable!()
5665 };
5666
5667 assert_eq!(sum.data_points.len(), 1);
5668 assert_eq!(sum.data_points[0].value, 40);
5669 }
5670
5671 #[cfg(feature = "experimental_metrics_bound_instruments")]
5672 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5673 async fn bound_counter_empty_attributes_shares_with_unbound() {
5674 let mut test_context = TestContext::new(Temporality::Cumulative);
5675 let counter = test_context.u64_counter("test", "my_counter", None);
5676 let bound = counter.bind(&[]);
5677
5678 counter.add(10, &[]);
5681 bound.add(20);
5682 counter.add(30, &[]);
5683 bound.add(40);
5684 test_context.flush_metrics();
5685
5686 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
5687 unreachable!()
5688 };
5689
5690 assert_eq!(
5691 sum.data_points.len(),
5692 1,
5693 "Bound and unbound with empty attributes must share the same data point"
5694 );
5695 assert_eq!(sum.data_points[0].value, 100);
5696 }
5697
5698 #[cfg(feature = "experimental_metrics_bound_instruments")]
5699 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5700 async fn bound_histogram_empty_attributes_shares_with_unbound() {
5701 let mut test_context = TestContext::new(Temporality::Cumulative);
5702 let histogram = test_context
5703 .meter()
5704 .u64_histogram("my_histogram")
5705 .with_boundaries(vec![5.0, 10.0, 25.0])
5706 .build();
5707 let bound = histogram.bind(&[]);
5708
5709 histogram.record(3, &[]);
5710 bound.record(7);
5711 histogram.record(20, &[]);
5712 test_context.flush_metrics();
5713
5714 let MetricData::Histogram(hist) = test_context.get_aggregation::<u64>("my_histogram", None)
5715 else {
5716 unreachable!()
5717 };
5718
5719 assert_eq!(
5720 hist.data_points.len(),
5721 1,
5722 "Bound and unbound with empty attributes must share the same data point"
5723 );
5724 let dp = &hist.data_points[0];
5725 assert!(dp.attributes.is_empty());
5726 assert_eq!(dp.count, 3);
5727 assert_eq!(dp.sum, 30);
5728 }
5729
5730 #[cfg(feature = "experimental_metrics_bound_instruments")]
5731 #[cfg(feature = "spec_unstable_metrics_views")]
5732 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5733 async fn bound_counter_view_filters_attributes_at_bind_time() {
5734 use opentelemetry::Key;
5735
5736 let exporter = InMemoryMetricExporter::default();
5737 let view = |i: &Instrument| {
5738 if i.name() == "my_counter" {
5739 Stream::builder()
5740 .with_allowed_attribute_keys(vec![Key::new("k1"), Key::new("k2")])
5741 .build()
5742 .ok()
5743 } else {
5744 None
5745 }
5746 };
5747 let meter_provider = SdkMeterProvider::builder()
5748 .with_periodic_exporter(exporter.clone())
5749 .with_view(view)
5750 .build();
5751 let meter = meter_provider.meter("test");
5752 let counter = meter.u64_counter("my_counter").build();
5753
5754 let bound = counter.bind(&[
5756 KeyValue::new("k1", "v1"),
5757 KeyValue::new("k2", "v2"),
5758 KeyValue::new("k3", "v3"),
5759 ]);
5760 bound.add(10);
5761 bound.add(20);
5762
5763 counter.add(
5766 7,
5767 &[
5768 KeyValue::new("k1", "v1"),
5769 KeyValue::new("k2", "v2"),
5770 KeyValue::new("k3", "different"),
5771 ],
5772 );
5773
5774 meter_provider.force_flush().unwrap();
5775 let resource_metrics = exporter
5776 .get_finished_metrics()
5777 .expect("metrics are expected to be exported.");
5778 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
5779 let data::AggregatedMetrics::U64(MetricData::Sum(sum)) = &metric.data else {
5780 unreachable!()
5781 };
5782
5783 assert_eq!(
5784 sum.data_points.len(),
5785 1,
5786 "view should filter k3, leaving bound+unbound to aggregate together"
5787 );
5788 assert_eq!(sum.data_points[0].value, 37);
5789 let attrs = &sum.data_points[0].attributes;
5790 assert_eq!(attrs.len(), 2);
5791 assert!(attrs.iter().any(|kv| kv.key.as_str() == "k1"));
5792 assert!(attrs.iter().any(|kv| kv.key.as_str() == "k2"));
5793 assert!(!attrs.iter().any(|kv| kv.key.as_str() == "k3"));
5794 }
5795
5796 #[cfg(feature = "experimental_metrics_bound_instruments")]
5797 #[cfg(feature = "spec_unstable_metrics_views")]
5798 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5799 async fn bound_histogram_view_filters_attributes_at_bind_time() {
5800 use opentelemetry::Key;
5801
5802 let exporter = InMemoryMetricExporter::default();
5803 let view = |i: &Instrument| {
5804 if i.name() == "my_hist" {
5805 Stream::builder()
5806 .with_allowed_attribute_keys(vec![Key::new("k1"), Key::new("k2")])
5807 .build()
5808 .ok()
5809 } else {
5810 None
5811 }
5812 };
5813 let meter_provider = SdkMeterProvider::builder()
5814 .with_periodic_exporter(exporter.clone())
5815 .with_view(view)
5816 .build();
5817 let meter = meter_provider.meter("test");
5818 let histogram = meter
5819 .u64_histogram("my_hist")
5820 .with_boundaries(vec![5.0, 10.0, 25.0])
5821 .build();
5822
5823 let bound = histogram.bind(&[
5824 KeyValue::new("k1", "v1"),
5825 KeyValue::new("k2", "v2"),
5826 KeyValue::new("k3", "v3"),
5827 ]);
5828 bound.record(3);
5829 bound.record(20);
5830 histogram.record(
5831 7,
5832 &[
5833 KeyValue::new("k1", "v1"),
5834 KeyValue::new("k2", "v2"),
5835 KeyValue::new("k3", "different"),
5836 ],
5837 );
5838
5839 meter_provider.force_flush().unwrap();
5840 let resource_metrics = exporter
5841 .get_finished_metrics()
5842 .expect("metrics are expected to be exported.");
5843 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
5844 let data::AggregatedMetrics::U64(MetricData::Histogram(hist)) = &metric.data else {
5845 unreachable!()
5846 };
5847
5848 assert_eq!(
5849 hist.data_points.len(),
5850 1,
5851 "view should filter k3, leaving bound+unbound to aggregate together"
5852 );
5853 let dp = &hist.data_points[0];
5854 assert_eq!(dp.count, 3);
5855 assert_eq!(dp.sum, 30);
5856 assert_eq!(dp.attributes.len(), 2);
5857 assert!(dp.attributes.iter().any(|kv| kv.key.as_str() == "k1"));
5858 assert!(dp.attributes.iter().any(|kv| kv.key.as_str() == "k2"));
5859 assert!(!dp.attributes.iter().any(|kv| kv.key.as_str() == "k3"));
5860 }
5861
5862 #[cfg(feature = "experimental_metrics_bound_instruments")]
5863 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5864 async fn bound_counter_at_overflow_attributes_to_overflow_bucket_cumulative() {
5865 let cardinality_limit = 3;
5870 let view = move |i: &Instrument| {
5871 if i.name() == "my_counter" {
5872 Stream::builder()
5873 .with_name("my_counter")
5874 .with_cardinality_limit(cardinality_limit)
5875 .build()
5876 .ok()
5877 } else {
5878 None
5879 }
5880 };
5881 let mut test_context = TestContext::new_with_view(Temporality::Cumulative, view);
5882 let counter = test_context.u64_counter("test", "my_counter", None);
5883
5884 for v in 0..cardinality_limit {
5885 counter.add(1, &[KeyValue::new("A", v.to_string())]);
5886 }
5887 let bound = counter.bind(&[KeyValue::new("A", "overflow_bind")]);
5888 bound.add(10);
5889
5890 test_context.flush_metrics();
5891 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
5892 unreachable!()
5893 };
5894 let overflow_dp = find_overflow_sum_datapoint(&sum.data_points)
5895 .expect("cycle 1: overflow point expected");
5896 assert_eq!(overflow_dp.value, 10);
5897
5898 test_context.reset_metrics();
5901 bound.add(7);
5902 test_context.flush_metrics();
5903 let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
5904 unreachable!()
5905 };
5906 let overflow_dp = find_overflow_sum_datapoint(&sum.data_points)
5907 .expect("cycle 2: overflow point expected");
5908 assert_eq!(
5909 overflow_dp.value, 17,
5910 "cumulative overflow-bound writes must accumulate"
5911 );
5912 }
5913
5914 #[cfg(feature = "experimental_metrics_bound_instruments")]
5915 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5916 async fn bound_histogram_at_overflow_attributes_to_overflow_bucket_cumulative() {
5917 let cardinality_limit = 3;
5918 let view = move |i: &Instrument| {
5919 if i.name() == "my_histogram" {
5920 Stream::builder()
5921 .with_name("my_histogram")
5922 .with_cardinality_limit(cardinality_limit)
5923 .build()
5924 .ok()
5925 } else {
5926 None
5927 }
5928 };
5929 let mut test_context = TestContext::new_with_view(Temporality::Cumulative, view);
5930 let histogram = test_context
5931 .meter()
5932 .f64_histogram("my_histogram")
5933 .with_boundaries(vec![10.0, 50.0])
5934 .build();
5935
5936 for v in 0..cardinality_limit {
5937 histogram.record(1.0, &[KeyValue::new("A", v.to_string())]);
5938 }
5939 let bound = histogram.bind(&[KeyValue::new("A", "overflow_bind")]);
5940 bound.record(20.0);
5941
5942 test_context.flush_metrics();
5943 let MetricData::Histogram(hist) = test_context.get_aggregation::<f64>("my_histogram", None)
5944 else {
5945 unreachable!()
5946 };
5947 let overflow_dp = find_overflow_histogram_datapoint(&hist.data_points)
5948 .expect("cycle 1: overflow point expected");
5949 assert_eq!(overflow_dp.sum, 20.0);
5950 assert_eq!(overflow_dp.count, 1);
5951
5952 test_context.reset_metrics();
5955 bound.record(30.0);
5956 test_context.flush_metrics();
5957 let MetricData::Histogram(hist) = test_context.get_aggregation::<f64>("my_histogram", None)
5958 else {
5959 unreachable!()
5960 };
5961 let overflow_dp = find_overflow_histogram_datapoint(&hist.data_points)
5962 .expect("cycle 2: overflow point expected");
5963 assert_eq!(
5964 overflow_dp.sum, 50.0,
5965 "cumulative overflow-bound writes must accumulate"
5966 );
5967 assert_eq!(overflow_dp.count, 2);
5968 }
5969
5970 #[cfg(feature = "experimental_metrics_bound_instruments")]
5971 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
5972 async fn bound_exponential_histogram_at_overflow_attributes_to_overflow_bucket_cumulative() {
5973 let cardinality_limit = 3;
5974 let view = move |i: &Instrument| {
5975 if i.name() == "my_histogram" {
5976 Stream::builder()
5977 .with_aggregation(Aggregation::Base2ExponentialHistogram {
5978 max_size: 160,
5979 max_scale: 20,
5980 record_min_max: true,
5981 })
5982 .with_cardinality_limit(cardinality_limit)
5983 .build()
5984 .ok()
5985 } else {
5986 None
5987 }
5988 };
5989 let mut test_context = TestContext::new_with_view(Temporality::Cumulative, view);
5990 let histogram = test_context.meter().f64_histogram("my_histogram").build();
5991
5992 for v in 0..cardinality_limit {
5993 histogram.record(1.0, &[KeyValue::new("A", v.to_string())]);
5994 }
5995 let bound = histogram.bind(&[KeyValue::new("A", "overflow_bind")]);
5996 bound.record(20.0);
5997
5998 test_context.flush_metrics();
5999 let MetricData::ExponentialHistogram(hist) =
6000 test_context.get_aggregation::<f64>("my_histogram", None)
6001 else {
6002 panic!("expected ExponentialHistogram aggregation");
6003 };
6004 let overflow_dp = find_overflow_exponential_histogram_datapoint(&hist.data_points)
6005 .expect("cycle 1: overflow point expected");
6006 assert_eq!(overflow_dp.sum, 20.0);
6007 assert_eq!(overflow_dp.count, 1);
6008
6009 test_context.reset_metrics();
6012 bound.record(30.0);
6013 test_context.flush_metrics();
6014 let MetricData::ExponentialHistogram(hist) =
6015 test_context.get_aggregation::<f64>("my_histogram", None)
6016 else {
6017 panic!("expected ExponentialHistogram aggregation");
6018 };
6019 let overflow_dp = find_overflow_exponential_histogram_datapoint(&hist.data_points)
6020 .expect("cycle 2: overflow point expected");
6021 assert_eq!(
6022 overflow_dp.sum, 50.0,
6023 "cumulative overflow-bound writes must accumulate"
6024 );
6025 assert_eq!(overflow_dp.count, 2);
6026 }
6027
6028 #[cfg(feature = "experimental_metrics_bound_instruments")]
6029 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6030 async fn bound_exponential_histogram_view_filters_attributes_at_bind_time() {
6031 use opentelemetry::Key;
6032
6033 let exporter = InMemoryMetricExporter::default();
6034 let view = |i: &Instrument| {
6035 if i.name() == "my_hist" {
6036 Stream::builder()
6037 .with_aggregation(Aggregation::Base2ExponentialHistogram {
6038 max_size: 160,
6039 max_scale: 20,
6040 record_min_max: true,
6041 })
6042 .with_allowed_attribute_keys(vec![Key::new("k1"), Key::new("k2")])
6043 .build()
6044 .ok()
6045 } else {
6046 None
6047 }
6048 };
6049 let meter_provider = SdkMeterProvider::builder()
6050 .with_periodic_exporter(exporter.clone())
6051 .with_view(view)
6052 .build();
6053 let meter = meter_provider.meter("test");
6054 let histogram = meter.f64_histogram("my_hist").build();
6055
6056 let bound = histogram.bind(&[
6057 KeyValue::new("k1", "v1"),
6058 KeyValue::new("k2", "v2"),
6059 KeyValue::new("k3", "v3"),
6060 ]);
6061 bound.record(3.0);
6062 bound.record(20.0);
6063 histogram.record(
6066 7.0,
6067 &[
6068 KeyValue::new("k1", "v1"),
6069 KeyValue::new("k2", "v2"),
6070 KeyValue::new("k3", "different"),
6071 ],
6072 );
6073
6074 meter_provider.force_flush().unwrap();
6075 let resource_metrics = exporter
6076 .get_finished_metrics()
6077 .expect("metrics are expected to be exported.");
6078 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
6079 let data::AggregatedMetrics::F64(MetricData::ExponentialHistogram(hist)) = &metric.data
6080 else {
6081 panic!("expected ExponentialHistogram aggregation");
6082 };
6083
6084 assert_eq!(
6085 hist.data_points.len(),
6086 1,
6087 "view should filter k3, leaving bound+unbound to aggregate together"
6088 );
6089 let dp = &hist.data_points[0];
6090 assert_eq!(dp.count, 3);
6091 assert_eq!(dp.sum, 30.0);
6092 assert_eq!(dp.attributes.len(), 2);
6093 assert!(dp.attributes.iter().any(|kv| kv.key.as_str() == "k1"));
6094 assert!(dp.attributes.iter().any(|kv| kv.key.as_str() == "k2"));
6095 assert!(!dp.attributes.iter().any(|kv| kv.key.as_str() == "k3"));
6096 }
6097
6098 #[cfg(feature = "experimental_metrics_bound_instruments")]
6099 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6100 async fn bound_gauge_cumulative() {
6101 let mut test_context = TestContext::new(Temporality::Cumulative);
6102 let gauge = test_context.meter().u64_gauge("my_gauge").build();
6103 let attrs = vec![KeyValue::new("key1", "bound_value")];
6104 let bound = gauge.bind(&attrs);
6105
6106 bound.record(10);
6107 bound.record(20);
6108 bound.record(30);
6109 test_context.flush_metrics();
6110
6111 let MetricData::Gauge(gauge_data) = test_context.get_aggregation::<u64>("my_gauge", None)
6112 else {
6113 unreachable!()
6114 };
6115
6116 assert_eq!(gauge_data.data_points.len(), 1, "Expected one data point");
6117 let dp = &gauge_data.data_points[0];
6118 assert_eq!(dp.value, 30, "Gauge should report the last recorded value");
6119 assert_eq!(dp.attributes, vec![KeyValue::new("key1", "bound_value")]);
6120 }
6121
6122 #[cfg(feature = "experimental_metrics_bound_instruments")]
6123 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6124 async fn bound_gauge_delta() {
6125 let mut test_context = TestContext::new(Temporality::Delta);
6126 let gauge = test_context.meter().u64_gauge("my_gauge").build();
6127 let attrs = vec![KeyValue::new("key1", "bound_value")];
6128 let bound = gauge.bind(&attrs);
6129
6130 bound.record(50);
6131 bound.record(75);
6132 test_context.flush_metrics();
6133
6134 let MetricData::Gauge(gauge_data) = test_context.get_aggregation::<u64>("my_gauge", None)
6135 else {
6136 unreachable!()
6137 };
6138 assert_eq!(gauge_data.data_points.len(), 1);
6139 assert_eq!(gauge_data.data_points[0].value, 75);
6140
6141 test_context.reset_metrics();
6143 bound.record(25);
6144 test_context.flush_metrics();
6145
6146 let MetricData::Gauge(gauge_data) = test_context.get_aggregation::<u64>("my_gauge", None)
6147 else {
6148 unreachable!()
6149 };
6150 assert_eq!(gauge_data.data_points.len(), 1);
6151 assert_eq!(
6152 gauge_data.data_points[0].value, 25,
6153 "Delta gauge should report only the latest value since the last collection"
6154 );
6155 }
6156
6157 #[cfg(feature = "experimental_metrics_bound_instruments")]
6158 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6159 async fn bound_gauge_matches_unbound() {
6160 let mut test_context = TestContext::new(Temporality::Cumulative);
6161 let gauge = test_context.meter().u64_gauge("my_gauge").build();
6162 let attrs = vec![KeyValue::new("key1", "shared")];
6163 let bound = gauge.bind(&attrs);
6164
6165 gauge.record(10, &attrs);
6168 bound.record(20);
6169 gauge.record(30, &attrs);
6170 bound.record(40);
6171 test_context.flush_metrics();
6172
6173 let MetricData::Gauge(gauge_data) = test_context.get_aggregation::<u64>("my_gauge", None)
6174 else {
6175 unreachable!()
6176 };
6177
6178 assert_eq!(
6179 gauge_data.data_points.len(),
6180 1,
6181 "Bound and unbound should share the same data point"
6182 );
6183 assert_eq!(
6184 gauge_data.data_points[0].value, 40,
6185 "Last write (via bound) should win"
6186 );
6187 }
6188
6189 #[cfg(feature = "experimental_metrics_bound_instruments")]
6190 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6191 async fn bound_gauge_delta_no_update_no_export() {
6192 let mut test_context = TestContext::new(Temporality::Delta);
6193 let gauge = test_context.meter().u64_gauge("my_gauge").build();
6194 let attrs = vec![KeyValue::new("key1", "bound_value")];
6195 let bound = gauge.bind(&attrs);
6196
6197 bound.record(10);
6199 test_context.flush_metrics();
6200 let MetricData::Gauge(gauge_data) = test_context.get_aggregation::<u64>("my_gauge", None)
6201 else {
6202 unreachable!()
6203 };
6204 assert_eq!(gauge_data.data_points.len(), 1);
6205 assert_eq!(gauge_data.data_points[0].value, 10);
6206
6207 test_context.reset_metrics();
6209 test_context.flush_metrics();
6210 let resource_metrics = test_context
6211 .exporter
6212 .get_finished_metrics()
6213 .expect("metrics export should succeed");
6214 assert!(
6215 resource_metrics.is_empty(),
6216 "Bound gauge with no updates should not export"
6217 );
6218
6219 test_context.reset_metrics();
6221 bound.record(99);
6222 test_context.flush_metrics();
6223 let MetricData::Gauge(gauge_data) = test_context.get_aggregation::<u64>("my_gauge", None)
6224 else {
6225 unreachable!()
6226 };
6227 assert_eq!(gauge_data.data_points.len(), 1);
6228 assert_eq!(
6229 gauge_data.data_points[0].value, 99,
6230 "Bound handle should produce fresh value after quiet cycle"
6231 );
6232 }
6233
6234 #[cfg(feature = "experimental_metrics_bound_instruments")]
6235 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6236 async fn bound_gauge_at_overflow_attributes_to_overflow_bucket() {
6237 let cardinality_limit = 3;
6238 let view = move |i: &Instrument| {
6239 if i.name() == "my_gauge" {
6240 Stream::builder()
6241 .with_name("my_gauge")
6242 .with_cardinality_limit(cardinality_limit)
6243 .build()
6244 .ok()
6245 } else {
6246 None
6247 }
6248 };
6249 let mut test_context = TestContext::new_with_view(Temporality::Delta, view);
6250 let gauge = test_context.meter().u64_gauge("my_gauge").build();
6251
6252 for v in 0..cardinality_limit {
6254 gauge.record(1, &[KeyValue::new("A", v.to_string())]);
6255 }
6256
6257 let overflow_attrs = vec![KeyValue::new("A", "overflow_bind")];
6259 let bound = gauge.bind(&overflow_attrs);
6260 bound.record(42);
6261
6262 test_context.flush_metrics();
6263 let MetricData::Gauge(gauge_data) = test_context.get_aggregation::<u64>("my_gauge", None)
6264 else {
6265 unreachable!()
6266 };
6267
6268 assert_eq!(
6269 gauge_data.data_points.len(),
6270 cardinality_limit + 1,
6271 "Expected {} unique + 1 overflow data points",
6272 cardinality_limit
6273 );
6274
6275 let overflow_dp = find_overflow_gauge_datapoint(&gauge_data.data_points)
6276 .expect("overflow point expected");
6277 assert_eq!(
6278 overflow_dp.value, 42,
6279 "Bound-at-overflow data should go to overflow bucket"
6280 );
6281 }
6282
6283 #[cfg(feature = "experimental_metrics_bound_instruments")]
6284 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6285 async fn bound_gauge_multiple_handles_same_attrs() {
6286 let mut test_context = TestContext::new(Temporality::Delta);
6287 let gauge = test_context.meter().u64_gauge("my_gauge").build();
6288 let attrs = vec![KeyValue::new("key1", "shared")];
6289
6290 let bound1 = gauge.bind(&attrs);
6291 let bound2 = gauge.bind(&attrs);
6292
6293 bound1.record(10);
6294 bound2.record(20);
6295 test_context.flush_metrics();
6296
6297 let MetricData::Gauge(gauge_data) = test_context.get_aggregation::<u64>("my_gauge", None)
6298 else {
6299 unreachable!()
6300 };
6301 assert_eq!(
6302 gauge_data.data_points.len(),
6303 1,
6304 "Multiple handles to same attrs should share data point"
6305 );
6306 assert_eq!(gauge_data.data_points[0].value, 20);
6307
6308 drop(bound1);
6310 test_context.reset_metrics();
6311 test_context.flush_metrics();
6312
6313 test_context.reset_metrics();
6314 bound2.record(5);
6315 test_context.flush_metrics();
6316
6317 let MetricData::Gauge(gauge_data) = test_context.get_aggregation::<u64>("my_gauge", None)
6318 else {
6319 unreachable!()
6320 };
6321 assert_eq!(gauge_data.data_points.len(), 1);
6322 assert_eq!(
6323 gauge_data.data_points[0].value, 5,
6324 "Entry should persist while any handle is alive"
6325 );
6326 }
6327
6328 #[cfg(feature = "experimental_metrics_bound_instruments")]
6329 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6330 async fn bound_gauge_empty_attributes_shares_with_unbound() {
6331 let mut test_context = TestContext::new(Temporality::Cumulative);
6332 let gauge = test_context.meter().u64_gauge("my_gauge").build();
6333 let bound = gauge.bind(&[]);
6334
6335 gauge.record(10, &[]);
6338 bound.record(20);
6339 gauge.record(30, &[]);
6340 bound.record(40);
6341 test_context.flush_metrics();
6342
6343 let MetricData::Gauge(gauge_data) = test_context.get_aggregation::<u64>("my_gauge", None)
6344 else {
6345 unreachable!()
6346 };
6347
6348 assert_eq!(
6349 gauge_data.data_points.len(),
6350 1,
6351 "Bound and unbound with empty attributes must share the same data point"
6352 );
6353 let dp = &gauge_data.data_points[0];
6354 assert!(dp.attributes.is_empty());
6355 assert_eq!(dp.value, 40);
6356 }
6357
6358 #[cfg(feature = "experimental_metrics_bound_instruments")]
6359 #[cfg(feature = "spec_unstable_metrics_views")]
6360 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6361 async fn bound_gauge_view_filters_attributes_at_bind_time() {
6362 use opentelemetry::Key;
6363
6364 let exporter = InMemoryMetricExporter::default();
6365 let view = |i: &Instrument| {
6366 if i.name() == "my_gauge" {
6367 Stream::builder()
6368 .with_allowed_attribute_keys(vec![Key::new("k1"), Key::new("k2")])
6369 .build()
6370 .ok()
6371 } else {
6372 None
6373 }
6374 };
6375 let meter_provider = SdkMeterProvider::builder()
6376 .with_periodic_exporter(exporter.clone())
6377 .with_view(view)
6378 .build();
6379 let meter = meter_provider.meter("test");
6380 let gauge = meter.u64_gauge("my_gauge").build();
6381
6382 let bound = gauge.bind(&[
6384 KeyValue::new("k1", "v1"),
6385 KeyValue::new("k2", "v2"),
6386 KeyValue::new("k3", "v3"),
6387 ]);
6388 bound.record(10);
6389 bound.record(20);
6390
6391 gauge.record(
6394 7,
6395 &[
6396 KeyValue::new("k1", "v1"),
6397 KeyValue::new("k2", "v2"),
6398 KeyValue::new("k3", "different"),
6399 ],
6400 );
6401
6402 meter_provider.force_flush().unwrap();
6403 let resource_metrics = exporter
6404 .get_finished_metrics()
6405 .expect("metrics are expected to be exported.");
6406 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
6407 let data::AggregatedMetrics::U64(MetricData::Gauge(gauge_data)) = &metric.data else {
6408 unreachable!()
6409 };
6410
6411 assert_eq!(
6412 gauge_data.data_points.len(),
6413 1,
6414 "view should filter k3, leaving bound+unbound to share the same data point"
6415 );
6416 assert_eq!(
6417 gauge_data.data_points[0].value, 7,
6418 "Last write (the unbound record) should win"
6419 );
6420 let attrs = &gauge_data.data_points[0].attributes;
6421 assert_eq!(attrs.len(), 2);
6422 assert!(attrs.iter().any(|kv| kv.key.as_str() == "k1"));
6423 assert!(attrs.iter().any(|kv| kv.key.as_str() == "k2"));
6424 assert!(!attrs.iter().any(|kv| kv.key.as_str() == "k3"));
6425 }
6426
6427 #[cfg(feature = "experimental_metrics_bound_instruments")]
6428 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6429 async fn bound_gauge_drop_enables_eviction() {
6430 let mut test_context = TestContext::new(Temporality::Delta);
6431 let gauge = test_context.meter().u64_gauge("my_gauge").build();
6432 let attrs = vec![KeyValue::new("key1", "ephemeral")];
6433
6434 {
6435 let bound = gauge.bind(&attrs);
6436 bound.record(100);
6437 test_context.flush_metrics();
6438
6439 let MetricData::Gauge(gauge_data) =
6440 test_context.get_aggregation::<u64>("my_gauge", None)
6441 else {
6442 unreachable!()
6443 };
6444 assert_eq!(gauge_data.data_points.len(), 1);
6445 assert_eq!(gauge_data.data_points[0].value, 100);
6446 }
6448
6449 test_context.reset_metrics();
6451 test_context.flush_metrics();
6452
6453 test_context.reset_metrics();
6456 gauge.record(1, &[KeyValue::new("key1", "new_entry")]);
6457 test_context.flush_metrics();
6458
6459 let MetricData::Gauge(gauge_data) = test_context.get_aggregation::<u64>("my_gauge", None)
6460 else {
6461 unreachable!()
6462 };
6463
6464 assert_eq!(gauge_data.data_points.len(), 1);
6465 let dp = find_gauge_datapoint_with_key_value(&gauge_data.data_points, "key1", "new_entry")
6466 .expect("new_entry should be present");
6467 assert_eq!(dp.value, 1);
6468 assert!(
6469 find_gauge_datapoint_with_key_value(&gauge_data.data_points, "key1", "ephemeral")
6470 .is_none(),
6471 "Dropped bound entry should have been evicted"
6472 );
6473 }
6474
6475 #[cfg(feature = "experimental_metrics_bound_instruments")]
6476 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6477 async fn bound_updown_counter_cumulative() {
6478 let mut test_context = TestContext::new(Temporality::Cumulative);
6479 let counter = test_context.i64_up_down_counter("test", "my_updown", None);
6480 let attrs = vec![KeyValue::new("key1", "bound_value")];
6481 let bound = counter.bind(&attrs);
6482
6483 bound.add(50);
6485 bound.add(-20);
6486 bound.add(30);
6487 test_context.flush_metrics();
6488
6489 let MetricData::Sum(sum) = test_context.get_aggregation::<i64>("my_updown", None) else {
6490 unreachable!()
6491 };
6492
6493 assert_eq!(sum.data_points.len(), 1, "Expected one data point");
6494 assert!(
6495 !sum.is_monotonic,
6496 "UpDownCounter must not produce monotonic Sums"
6497 );
6498 assert_eq!(sum.temporality, Temporality::Cumulative);
6499 assert_eq!(sum.data_points[0].value, 60);
6500 assert_eq!(
6501 sum.data_points[0].attributes,
6502 vec![KeyValue::new("key1", "bound_value")]
6503 );
6504 }
6505
6506 #[cfg(feature = "experimental_metrics_bound_instruments")]
6507 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6508 async fn bound_updown_counter_always_cumulative_even_when_delta_requested() {
6509 let mut test_context = TestContext::new(Temporality::Delta);
6512 let counter = test_context.i64_up_down_counter("test", "my_updown", None);
6513 let attrs = vec![KeyValue::new("key1", "bound_value")];
6514 let bound = counter.bind(&attrs);
6515
6516 bound.add(50);
6517 bound.add(-10);
6518 test_context.flush_metrics();
6519
6520 let MetricData::Sum(sum) = test_context.get_aggregation::<i64>("my_updown", None) else {
6521 unreachable!()
6522 };
6523 assert_eq!(
6524 sum.temporality,
6525 Temporality::Cumulative,
6526 "UpDownCounter must produce Cumulative regardless of reader preference"
6527 );
6528 assert!(!sum.is_monotonic);
6529 assert_eq!(sum.data_points.len(), 1);
6530 assert_eq!(sum.data_points[0].value, 40);
6531 }
6532
6533 #[cfg(feature = "experimental_metrics_bound_instruments")]
6534 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6535 async fn bound_updown_counter_matches_unbound() {
6536 let mut test_context = TestContext::new(Temporality::Cumulative);
6537 let counter = test_context.i64_up_down_counter("test", "my_updown", None);
6538 let attrs = vec![KeyValue::new("key1", "shared")];
6539 let bound = counter.bind(&attrs);
6540
6541 counter.add(10, &attrs);
6543 bound.add(-20);
6544 counter.add(30, &attrs);
6545 bound.add(40);
6546 test_context.flush_metrics();
6547
6548 let MetricData::Sum(sum) = test_context.get_aggregation::<i64>("my_updown", None) else {
6549 unreachable!()
6550 };
6551
6552 assert_eq!(
6553 sum.data_points.len(),
6554 1,
6555 "Bound and unbound should share the same data point"
6556 );
6557 assert!(!sum.is_monotonic);
6558 assert_eq!(sum.data_points[0].value, 60);
6559 }
6560
6561 #[cfg(feature = "experimental_metrics_bound_instruments")]
6562 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6563 async fn bound_updown_counter_at_overflow_attributes_to_overflow_bucket() {
6564 let cardinality_limit = 3;
6565 let view = move |i: &Instrument| {
6566 if i.name() == "my_updown" {
6567 Stream::builder()
6568 .with_name("my_updown")
6569 .with_cardinality_limit(cardinality_limit)
6570 .build()
6571 .ok()
6572 } else {
6573 None
6574 }
6575 };
6576 let mut test_context = TestContext::new_with_view(Temporality::Cumulative, view);
6577 let counter = test_context.i64_up_down_counter("test", "my_updown", None);
6578
6579 for v in 0..cardinality_limit {
6581 counter.add(1, &[KeyValue::new("A", v.to_string())]);
6582 }
6583
6584 let overflow_attrs = vec![KeyValue::new("A", "overflow_bind")];
6586 let bound = counter.bind(&overflow_attrs);
6587 bound.add(42);
6588 bound.add(-2);
6589
6590 test_context.flush_metrics();
6591 let MetricData::Sum(sum) = test_context.get_aggregation::<i64>("my_updown", None) else {
6592 unreachable!()
6593 };
6594
6595 assert_eq!(
6596 sum.data_points.len(),
6597 cardinality_limit + 1,
6598 "Expected {} unique + 1 overflow data points",
6599 cardinality_limit
6600 );
6601
6602 let overflow_dp =
6603 find_overflow_sum_datapoint(&sum.data_points).expect("overflow point expected");
6604 assert_eq!(
6605 overflow_dp.value, 40,
6606 "Bound-at-overflow data should go to overflow bucket"
6607 );
6608 }
6609
6610 #[cfg(feature = "experimental_metrics_bound_instruments")]
6611 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6612 async fn bound_updown_counter_multiple_handles_same_attrs() {
6613 let mut test_context = TestContext::new(Temporality::Cumulative);
6614 let counter = test_context.i64_up_down_counter("test", "my_updown", None);
6615 let attrs = vec![KeyValue::new("key1", "shared")];
6616
6617 let bound1 = counter.bind(&attrs);
6618 let bound2 = counter.bind(&attrs);
6619
6620 bound1.add(10);
6621 bound2.add(-3);
6622 test_context.flush_metrics();
6623
6624 let MetricData::Sum(sum) = test_context.get_aggregation::<i64>("my_updown", None) else {
6625 unreachable!()
6626 };
6627 assert_eq!(
6628 sum.data_points.len(),
6629 1,
6630 "Multiple handles to same attrs should share data point"
6631 );
6632 assert_eq!(sum.data_points[0].value, 7);
6633
6634 drop(bound1);
6637 test_context.reset_metrics();
6638 bound2.add(-5);
6639 test_context.flush_metrics();
6640
6641 let MetricData::Sum(sum) = test_context.get_aggregation::<i64>("my_updown", None) else {
6642 unreachable!()
6643 };
6644 assert_eq!(sum.data_points.len(), 1);
6645 assert_eq!(
6646 sum.data_points[0].value, 2,
6647 "Cumulative tracker should retain prior bound contributions"
6648 );
6649 }
6650
6651 #[cfg(feature = "experimental_metrics_bound_instruments")]
6652 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6653 async fn bound_updown_counter_empty_attributes_shares_with_unbound() {
6654 let mut test_context = TestContext::new(Temporality::Cumulative);
6655 let counter = test_context.i64_up_down_counter("test", "my_updown", None);
6656 let bound = counter.bind(&[]);
6657
6658 counter.add(10, &[]);
6661 bound.add(-3);
6662 counter.add(30, &[]);
6663 bound.add(-2);
6664 test_context.flush_metrics();
6665
6666 let MetricData::Sum(sum) = test_context.get_aggregation::<i64>("my_updown", None) else {
6667 unreachable!()
6668 };
6669
6670 assert_eq!(
6671 sum.data_points.len(),
6672 1,
6673 "Bound and unbound with empty attributes must share the same data point"
6674 );
6675 assert!(sum.data_points[0].attributes.is_empty());
6676 assert_eq!(sum.data_points[0].value, 35);
6677 }
6678
6679 #[cfg(feature = "experimental_metrics_bound_instruments")]
6680 #[cfg(feature = "spec_unstable_metrics_views")]
6681 #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
6682 async fn bound_updown_counter_view_filters_attributes_at_bind_time() {
6683 use opentelemetry::Key;
6684
6685 let exporter = InMemoryMetricExporter::default();
6686 let view = |i: &Instrument| {
6687 if i.name() == "my_updown" {
6688 Stream::builder()
6689 .with_allowed_attribute_keys(vec![Key::new("k1"), Key::new("k2")])
6690 .build()
6691 .ok()
6692 } else {
6693 None
6694 }
6695 };
6696 let meter_provider = SdkMeterProvider::builder()
6697 .with_periodic_exporter(exporter.clone())
6698 .with_view(view)
6699 .build();
6700 let meter = meter_provider.meter("test");
6701 let counter = meter.i64_up_down_counter("my_updown").build();
6702
6703 let bound = counter.bind(&[
6705 KeyValue::new("k1", "v1"),
6706 KeyValue::new("k2", "v2"),
6707 KeyValue::new("k3", "v3"),
6708 ]);
6709 bound.add(10);
6710 bound.add(-3);
6711
6712 counter.add(
6715 7,
6716 &[
6717 KeyValue::new("k1", "v1"),
6718 KeyValue::new("k2", "v2"),
6719 KeyValue::new("k3", "different"),
6720 ],
6721 );
6722
6723 meter_provider.force_flush().unwrap();
6724 let resource_metrics = exporter
6725 .get_finished_metrics()
6726 .expect("metrics are expected to be exported.");
6727 let metric = &resource_metrics[0].scope_metrics[0].metrics[0];
6728 let data::AggregatedMetrics::I64(MetricData::Sum(sum)) = &metric.data else {
6729 unreachable!()
6730 };
6731
6732 assert_eq!(
6733 sum.data_points.len(),
6734 1,
6735 "view should filter k3, leaving bound+unbound to aggregate together"
6736 );
6737 assert!(!sum.is_monotonic);
6738 assert_eq!(sum.data_points[0].value, 14);
6739 let attrs = &sum.data_points[0].attributes;
6740 assert_eq!(attrs.len(), 2);
6741 assert!(attrs.iter().any(|kv| kv.key.as_str() == "k1"));
6742 assert!(attrs.iter().any(|kv| kv.key.as_str() == "k2"));
6743 assert!(!attrs.iter().any(|kv| kv.key.as_str() == "k3"));
6744 }
6745}