Skip to main content

opentelemetry_sdk/metrics/
mod.rs

1//! The crust of the OpenTelemetry metrics SDK.
2//!
3//! ## Configuration
4//!
5//! The metrics SDK configuration is stored with each [SdkMeterProvider].
6//! Configuration for [Resource]s, views, and `ManualReader` or
7//! [PeriodicReader] instances can be specified.
8//!
9//! ### Example
10//!
11//! ```
12//! use opentelemetry::global;
13//! use opentelemetry::KeyValue;
14//! use opentelemetry_sdk::{metrics::SdkMeterProvider, Resource};
15//!
16//! // Generate SDK configuration, resource, views, etc
17//! let resource = Resource::builder().build(); // default attributes about the current process
18//!
19//! // Create a meter provider with the desired config
20//! let meter_provider = SdkMeterProvider::builder().with_resource(resource).build();
21//! global::set_meter_provider(meter_provider.clone());
22//!
23//! // Use the meter provider to create meter instances
24//! let meter = global::meter("my_app");
25//!
26//! // Create instruments scoped to the meter
27//! let counter = meter
28//!     .u64_counter("power_consumption")
29//!     .with_unit("kWh")
30//!     .build();
31//!
32//! // use instruments to record measurements
33//! counter.add(10, &[KeyValue::new("rate", "standard")]);
34//!
35//! // shutdown the provider at the end of the application to ensure any metrics not yet
36//! // exported are flushed.
37//! meter_provider.shutdown().unwrap();
38//! ```
39//!
40//! [Resource]: crate::Resource
41
42#[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")]
57/// Module for periodic reader with async runtime.
58pub 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/// In-Memory metric exporter for testing purpose.
67#[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/// Defines the window that an aggregation was calculated over.
88#[derive(Debug, Copy, Clone, Default, PartialEq, Eq, Hash)]
89#[non_exhaustive]
90pub enum Temporality {
91    /// A measurement interval that continues to expand forward in time from a
92    /// starting point.
93    ///
94    /// New measurements are added to all previous measurements since a start time.
95    #[default]
96    Cumulative,
97
98    /// A measurement interval that resets each cycle.
99    ///
100    /// Measurements from one cycle are recorded independently, measurements from
101    /// other cycles do not affect them.
102    Delta,
103
104    /// Configures Synchronous Counter and Histogram instruments to use
105    /// Delta aggregation temporality, which allows them to shed memory
106    /// following a cardinality explosion, thus use less memory.
107    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    // Run all tests in this mod
149    // cargo test metrics::tests --features=testing,spec_unstable_metrics_views
150    // Note for all tests from this point onwards in this mod:
151    // "multi_thread" tokio flavor must be used else flush won't
152    // be able to make progress!
153
154    #[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        // Run this test with stdout enabled to see output.
158        // cargo test invalid_instrument_config_noops --features=testing,spec_unstable_metrics_views -- --nocapture
159        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            // As instrument name is invalid, no metrics should be exported
206            test_context.check_no_metrics();
207        }
208
209        let invalid_bucket_boundaries = vec![
210            vec![1.0, 1.0],                          // duplicate boundaries
211            vec![1.0, 2.0, 3.0, 2.0],                // duplicate non consequent boundaries
212            vec![1.0, 2.0, 3.0, 4.0, 2.5],           // unsorted boundaries
213            vec![1.0, 2.0, 3.0, f64::INFINITY, 4.0], // boundaries with positive infinity
214            vec![1.0, 2.0, 3.0, f64::NAN],           // boundaries with NaNs
215            vec![f64::NEG_INFINITY, 2.0, 3.0],       // boundaries with negative infinity
216        ];
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            // As bucket boundaries provided via advisory params are invalid,
228            // no metrics should be exported
229            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        // Run this test with stdout enabled to see output.
237        // cargo test valid_instrument_config_with_feature_experimental_metrics_disable_name_validation --all-features -- --nocapture
238        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            // As instrument name are valid because of the feature flag, metrics should be exported
290            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        // Ensuring that the Histograms with invalid bucket boundaries are not exported
299        // when using the feature flag
300        let invalid_bucket_boundaries = vec![
301            vec![1.0, 1.0],                          // duplicate boundaries
302            vec![1.0, 2.0, 3.0, 2.0],                // duplicate non consequent boundaries
303            vec![1.0, 2.0, 3.0, 4.0, 2.5],           // unsorted boundaries
304            vec![1.0, 2.0, 3.0, f64::INFINITY, 4.0], // boundaries with positive infinity
305            vec![1.0, 2.0, 3.0, f64::NAN],           // boundaries with NaNs
306            vec![f64::NEG_INFINITY, 2.0, 3.0],       // boundaries with negative infinity
307        ];
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            // As bucket boundaries provided via advisory params are invalid,
319            // no metrics should be exported
320            test_context.check_no_metrics();
321        }
322    }
323
324    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
325    async fn counter_aggregation_delta() {
326        // Run this test with stdout enabled to see output.
327        // cargo test counter_aggregation_delta --features=testing -- --nocapture
328        counter_aggregation_helper(Temporality::Delta);
329    }
330
331    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
332    async fn counter_aggregation_cumulative() {
333        // Run this test with stdout enabled to see output.
334        // cargo test counter_aggregation_cumulative --features=testing -- --nocapture
335        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        // Run this test with stdout enabled to see output.
403        // cargo test counter_aggregation_attribute --features=testing -- --nocapture
404        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        // Run this test with stdout enabled to see output.
410        // cargo test histogram_aggregation_cumulative --features=testing -- --nocapture
411        histogram_aggregation_helper(Temporality::Cumulative);
412    }
413
414    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
415    async fn histogram_aggregation_delta() {
416        // Run this test with stdout enabled to see output.
417        // cargo test histogram_aggregation_delta --features=testing -- --nocapture
418        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        // Run this test with stdout enabled to see output.
424        // cargo test histogram_aggregation_with_custom_bounds --features=testing -- --nocapture
425        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        // Run this test with stdout enabled to see output.
432        // cargo test histogram_aggregation_with_empty_bounds --features=testing -- --nocapture
433        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        // Run this test with stdout enabled to see output.
440        // cargo test histogram_aggregation_with_custom_bounds_and_view --features=testing -- --nocapture
441        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        // Run this test with stdout enabled to see output.
448        // cargo test exponential_histogram_aggregation_with_view --features=testing -- --nocapture
449        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        // Run this test with stdout enabled to see output.
456        // cargo test updown_counter_aggregation_cumulative --features=testing -- --nocapture
457        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        // Run this test with stdout enabled to see output.
463        // cargo test updown_counter_aggregation_delta --features=testing -- --nocapture
464        updown_counter_aggregation_helper(Temporality::Delta);
465    }
466
467    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
468    async fn gauge_aggregation() {
469        // Run this test with stdout enabled to see output.
470        // cargo test gauge_aggregation --features=testing -- --nocapture
471
472        // Gauge should use last value aggregation regardless of the aggregation temporality used.
473        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        // Run this test with stdout enabled to see output.
480        // cargo test observable_gauge_aggregation --features=testing -- --nocapture
481
482        // Gauge should use last value aggregation regardless of the aggregation temporality used.
483        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        // Run this test with stdout enabled to see output.
504        // cargo test observable_counter_aggregation --features=testing -- --nocapture
505        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        // Arrange
516        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        // The Observable counter reports values[0], values[1],....values[n] on each flush.
523        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            // find and validate datapoint
561            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                // Cumulative counter should have the value as is.
570                assert_eq!(data_point.value, *v);
571            } else {
572                // Delta counter should have the increment value.
573                // Except for the first value which should be the start value.
574                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        // Run this test with stdout enabled to see output.
663        // cargo test observable_counter_delta_attribute_set_reappears_after_gap --features=testing -- --nocapture
664
665        // This test verifies the behavior when an attribute set is not reported
666        // for one collection cycle and then reappears.
667        // See: https://github.com/open-telemetry/opentelemetry-specification/issues/4861
668        //
669        // Scenario (Observable Counter with Delta temporality):
670        // | Collection | Callback Reports  | Expected Delta Export  |
671        // |------------|-------------------|------------------------|
672        // | 1          | A=100, B=50       | A=100, B=50            |
673        // | 2          | A=150 (B missing) | A=50 (B not exported)  |
674        // | 3          | A=200, B=80       | A=50, B=80             |
675        //
676        // Current implementation: When B reappears, its delta is calculated from zero
677        // (fresh start), not from the last known value. This is Option 1 from the spec issue.
678
679        let mut test_context = TestContext::new(Temporality::Delta);
680
681        // Shared state for callback: (collection_cycle, value_a, value_b_option)
682        // value_b_option is None when B should not be reported
683        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        // Collection 1: A=100, B=50
701        {
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            // First collection: delta = value - 0
720            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        // Collection 2: A=150, B missing
730        {
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            // Only A should be exported, B is not observed so not exported (per spec)
741            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        // Collection 3: A=200, B=80 (B reappears)
755        {
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            // B reappears after a gap. Current implementation uses "delta from zero" (Option 1).
775            // This means B's delta = 80 - 0 = 80, not 80 - 50 = 30.
776            // See: https://github.com/open-telemetry/opentelemetry-specification/issues/4861
777            // TODO: Watch for spec clarification on this behavior.
778            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            // Act
795            let counter = meter.u64_counter("my_counter").build();
796
797            counter.add(10, &[]);
798            provider.force_flush().unwrap();
799
800            // Assert
801            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        // Test Meter creation in 2 ways, both with empty string as meter name
818        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        // Arrange
829        let exporter = InMemoryMetricExporter::default();
830        let meter_provider = SdkMeterProvider::builder()
831            .with_periodic_exporter(exporter.clone())
832            .build();
833
834        // Act
835        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        // Assert
855        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        // Expecting 1 time-series.
872        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        // Arrange
881        let exporter = InMemoryMetricExporter::default();
882        let meter_provider = SdkMeterProvider::builder()
883            .with_periodic_exporter(exporter.clone())
884            .build();
885
886        // Act
887        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        // Assert
908        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            // Expecting 1 time-series.
939            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            // Expecting 1 time-series.
960            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        // Arrange
972        let exporter = InMemoryMetricExporter::default();
973        let meter_provider = SdkMeterProvider::builder()
974            .with_periodic_exporter(exporter.clone())
975            .build();
976
977        // Act
978        // Meters are identical.
979        // Hence there should be a single metric stream output for this test.
980        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        // Assert
1012        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        // This is validating current behavior, but it is not guaranteed to be the case in the future,
1031        // as this is a user error and SDK reserves right to change this behavior.
1032        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        // Expecting 1 time-series.
1046        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        // Run this test with stdout enabled to see output.
1055        // cargo test histogram_aggregation_with_invalid_aggregation_should_proceed_as_if_view_not_exist --features=testing -- --nocapture
1056
1057        // Arrange
1058        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], // invalid boundaries
1064                        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        // Act
1080        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        // Assert
1090        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        // LastValue aggregation is only valid for Gauge instruments.
1108        // When applied to a Counter via a view, the view is ignored and
1109        // the default aggregation (Sum) is used per the spec:
1110        // "proceed as if the View did not match"
1111
1112        // Arrange
1113        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        // Act
1131        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        // Assert - view is ignored, default aggregation is used
1137        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        // Original name is used (view rename is ignored)
1143        assert_eq!(
1144            metric.name, "my_counter",
1145            "View rename should be ignored due to incompatible aggregation."
1146        );
1147        // Default Sum aggregation is used
1148        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        // Sum aggregation is not valid for Gauge instruments.
1160        // When applied to a Gauge via a view, the view is ignored and
1161        // the default aggregation (LastValue/Gauge) is used per the spec:
1162        // "proceed as if the View did not match"
1163
1164        // Arrange
1165        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        // Act
1183        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        // Assert - view is ignored, default aggregation is used
1189        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        // Original name is used (view rename is ignored)
1195        assert_eq!(
1196            metric.name, "my_gauge",
1197            "View rename should be ignored due to incompatible aggregation."
1198        );
1199        // Default Gauge (LastValue) aggregation is used
1200        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        // LastValue aggregation is only valid for Gauge instruments.
1212        // When applied to an UpDownCounter via a view, the view is ignored and
1213        // the default aggregation (Sum) is used per the spec:
1214        // "proceed as if the View did not match"
1215
1216        // Arrange
1217        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        // Act
1235        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        // Assert - view is ignored, default aggregation is used
1241        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        // Original name is used (view rename is ignored)
1247        assert_eq!(
1248            metric.name, "my_updown_counter",
1249            "View rename should be ignored due to incompatible aggregation."
1250        );
1251        // Default Sum aggregation is used
1252        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        // LastValue aggregation is only valid for Gauge instruments.
1264        // When applied to a Histogram via a view, the view is ignored and
1265        // the default aggregation (ExplicitBucketHistogram) is used per the spec:
1266        // "proceed as if the View did not match"
1267
1268        // Arrange
1269        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        // Act
1287        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        // Assert - view is ignored, default aggregation is used
1293        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        // Original name is used (view rename is ignored)
1299        assert_eq!(
1300            metric.name, "my_histogram",
1301            "View rename should be ignored due to incompatible aggregation."
1302        );
1303        // Default Histogram aggregation is used
1304        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        // Sum aggregation is not valid for Observable Gauge instruments.
1316        // When applied to an Observable Gauge via a view, the view is ignored and
1317        // the default aggregation (LastValue/Gauge) is used per the spec:
1318        // "proceed as if the View did not match"
1319
1320        // Arrange
1321        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        // Act
1339        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        // Assert - view is ignored, default aggregation is used
1349        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        // Original name is used (view rename is ignored)
1355        assert_eq!(
1356            metric.name, "my_observable_gauge",
1357            "View rename should be ignored due to incompatible aggregation."
1358        );
1359        // Default Gauge (LastValue) aggregation is used
1360        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        // LastValue aggregation is only valid for Gauge instruments.
1372        // When applied to an Observable Counter via a view, the view is ignored and
1373        // the default aggregation (Sum) is used per the spec:
1374        // "proceed as if the View did not match"
1375
1376        // Arrange
1377        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        // Act
1395        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        // Assert - view is ignored, default aggregation is used
1405        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        // Original name is used (view rename is ignored)
1411        assert_eq!(
1412            metric.name, "my_observable_counter",
1413            "View rename should be ignored due to incompatible aggregation."
1414        );
1415        // Default Sum aggregation is used
1416        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        // LastValue aggregation is only valid for Gauge instruments.
1428        // When applied to an Observable UpDownCounter via a view, the view is ignored and
1429        // the default aggregation (Sum) is used per the spec:
1430        // "proceed as if the View did not match"
1431
1432        // Arrange
1433        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        // Act
1451        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        // Assert - view is ignored, default aggregation is used
1461        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        // Original name is used (view rename is ignored)
1467        assert_eq!(
1468            metric.name, "my_observable_updowncounter",
1469            "View rename should be ignored due to incompatible aggregation."
1470        );
1471        // Default Sum aggregation is used
1472        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        // Sum aggregation is valid for Histogram instruments.
1484        // When applied via a view, the aggregation should change from
1485        // ExplicitBucketHistogram to Sum.
1486
1487        // Arrange
1488        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        // Act
1506        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        // Assert - view is applied, Sum aggregation is used
1514        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        // View rename is applied
1520        assert_eq!(
1521            metric.name, "my_histogram_renamed",
1522            "View rename should be applied for compatible aggregation."
1523        );
1524        // Sum aggregation is used instead of default Histogram
1525        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); // 10 + 20 + 30
1532    }
1533
1534    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1535    async fn gauge_with_histogram_aggregation_is_valid() {
1536        // Histogram aggregation is valid for Gauge instruments.
1537        // When applied via a view, the aggregation should change from
1538        // LastValue to ExplicitBucketHistogram.
1539
1540        // Arrange
1541        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        // Act
1562        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        // Assert - view is applied, Histogram aggregation is used
1571        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        // View rename is applied
1577        assert_eq!(
1578            metric.name, "my_gauge_renamed",
1579            "View rename should be applied for compatible aggregation."
1580        );
1581        // Histogram aggregation is used instead of default LastValue
1582        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); // 3 + 7 + 15 + 30
1591        assert_eq!(dp.min, Some(3.0));
1592        assert_eq!(dp.max, Some(30.0));
1593        // Bucket boundaries: [5.0, 10.0, 25.0, 50.0]
1594        // Values: 3.0 (bucket 0), 7.0 (bucket 1), 15.0 (bucket 2), 30.0 (bucket 3)
1595        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        // Histogram aggregation is valid for Counter instruments.
1601        // When applied via a view, the aggregation should change from
1602        // Sum to ExplicitBucketHistogram.
1603
1604        // Arrange
1605        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        // Act
1626        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        // Assert - view is applied, Histogram aggregation is used
1635        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        // View rename is applied
1641        assert_eq!(
1642            metric.name, "my_counter_renamed",
1643            "View rename should be applied for compatible aggregation."
1644        );
1645        // Histogram aggregation is used instead of default Sum
1646        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); // 3 + 7 + 15 + 30
1655        assert_eq!(dp.min, Some(3));
1656        assert_eq!(dp.max, Some(30));
1657        // Bucket boundaries: [5.0, 10.0, 25.0, 50.0]
1658        // Values: 3 (bucket 0), 7 (bucket 1), 15 (bucket 2), 30 (bucket 3)
1659        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        // Histogram aggregation is valid for UpDownCounter instruments.
1665        // When applied via a view, the aggregation should change from
1666        // Sum to ExplicitBucketHistogram.
1667
1668        // Arrange
1669        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        // Act
1690        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        // Assert - view is applied, Histogram aggregation is used
1698        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        // View rename is applied
1704        assert_eq!(
1705            metric.name, "my_updowncounter_renamed",
1706            "View rename should be applied for compatible aggregation."
1707        );
1708        // Histogram aggregation is used instead of default Sum
1709        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        // Note: Sum is not recorded for UpDownCounter histogram per the spec
1718    }
1719
1720    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
1721    async fn counter_with_exponential_histogram_aggregation_is_valid() {
1722        // ExponentialHistogram aggregation is valid for Counter instruments.
1723        // When applied via a view, the aggregation should change from
1724        // Sum to Base2ExponentialHistogram.
1725
1726        // Arrange
1727        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        // Act
1749        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        // Assert - view is applied, ExponentialHistogram aggregation is used
1757        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        // View rename is applied
1763        assert_eq!(
1764            metric.name, "my_counter_renamed",
1765            "View rename should be applied for compatible aggregation."
1766        );
1767        // ExponentialHistogram aggregation is used instead of default Sum
1768        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); // 5 + 10 + 20
1778        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        // ExponentialHistogram aggregation is valid for Gauge instruments.
1785        // When applied via a view, the aggregation should change from
1786        // LastValue to Base2ExponentialHistogram.
1787
1788        // Arrange
1789        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        // Act
1811        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        // Assert - view is applied, ExponentialHistogram aggregation is used
1819        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        // View rename is applied
1825        assert_eq!(
1826            metric.name, "my_gauge_renamed",
1827            "View rename should be applied for compatible aggregation."
1828        );
1829        // ExponentialHistogram aggregation is used instead of default LastValue
1830        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); // 2.5 + 7.5 + 15.0
1840        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        // Run this test with stdout enabled to see output.
1847        // cargo test counter_with_drop_aggregation_is_dropped --features=testing -- --nocapture
1848
1849        // When a view matches an instrument and specifies Aggregation::Drop,
1850        // the instrument should be dropped and no metrics should be exported.
1851
1852        // Arrange
1853        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        // Act
1870        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        // Assert - no metrics should be exported because the view drops the instrument
1876        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        // Run this test with stdout enabled to see output.
1891        // cargo test histogram_with_drop_aggregation_is_dropped --features=testing -- --nocapture
1892
1893        // When a view matches a histogram and specifies Aggregation::Drop,
1894        // the instrument should be dropped and no metrics should be exported.
1895
1896        // Arrange
1897        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        // Act
1914        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        // Assert - no metrics should be exported because the view drops the instrument
1920        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        // Run this test with stdout enabled to see output.
1935        // cargo test gauge_with_drop_aggregation_is_dropped --features=testing -- --nocapture
1936
1937        // When a view matches a gauge and specifies Aggregation::Drop,
1938        // the instrument should be dropped and no metrics should be exported.
1939
1940        // Arrange
1941        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        // Act
1958        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        // Assert - no metrics should be exported because the view drops the instrument
1964        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        // Run this test with stdout enabled to see output.
1979        // cargo test counter_with_drop_aggregation_and_rename_should_still_drop --features=testing -- --nocapture
1980
1981        // When a view matches and specifies Aggregation::Drop, the instrument should be
1982        // dropped even if the view also specifies other customizations (like name).
1983        // The name customization is meaningless for a dropped instrument, but the Drop
1984        // intent should still be honored.
1985
1986        // Arrange
1987        let exporter = InMemoryMetricExporter::default();
1988        let view = |i: &Instrument| {
1989            if i.name == "my_counter" {
1990                Stream::builder()
1991                    .with_name("dropped_counter") // Meaningless but shouldn't break Drop
1992                    .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        // Act
2005        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        // Assert - no metrics should be exported because the view drops the instrument
2011        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        // cargo test metrics::tests::spatial_aggregation_when_view_drops_attributes_observable_counter --features=testing
2028
2029        // Arrange
2030        let exporter = InMemoryMetricExporter::default();
2031        // View drops all attributes.
2032        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        // Act
2048        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        // Assert
2081        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        // Expecting 1 time-series only, as the view drops all attributes resulting
2095        // in a single time-series.
2096        // This is failing today, due to lack of support for spatial aggregation.
2097        assert_eq!(sum.data_points.len(), 1);
2098
2099        // find and validate the single datapoint
2100        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        // cargo test spatial_aggregation_when_view_drops_attributes_counter --features=testing
2108
2109        // Arrange
2110        let exporter = InMemoryMetricExporter::default();
2111        // View drops all attributes.
2112        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        // Act
2130        let meter = meter_provider.meter("test");
2131        let counter = meter.u64_counter("my_counter").build();
2132
2133        // Normally, this would generate 3 time-series, but since the view
2134        // drops all attributes, we expect only 1 time-series.
2135        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        // Assert
2165        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        // Expecting 1 time-series only, as the view drops all attributes resulting
2179        // in a single time-series.
2180        // This is failing today, due to lack of support for spatial aggregation.
2181        assert_eq!(sum.data_points.len(), 1);
2182        // find and validate the single datapoint
2183        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        // Run this test with stdout enabled to see output.
2326        // cargo test delta_memory_efficiency_test --features=testing -- --nocapture
2327
2328        // Arrange
2329        let mut test_context = TestContext::new(Temporality::Delta);
2330        let counter = test_context.u64_counter("test", "my_counter", None);
2331
2332        // Act
2333        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        // Expecting 2 time-series.
2349        assert_eq!(sum.data_points.len(), 2);
2350
2351        // find and validate key1=value1 datapoint
2352        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        // find and validate key1=value2 datapoint
2357        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        // flush again, and validate that nothing is flushed
2363        // as delta temporality.
2364        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        // Run this test with stdout enabled to see output.
2377        // cargo test counter_multithreaded --features=testing -- --nocapture
2378
2379        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        // Run this test with stdout enabled to see output.
2386        // cargo test counter_f64_multithreaded --features=testing -- --nocapture
2387
2388        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        // Run this test with stdout enabled to see output.
2395        // cargo test histogram_multithreaded --features=testing -- --nocapture
2396
2397        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        // Run this test with stdout enabled to see output.
2404        // cargo test histogram_f64_multithreaded --features=testing -- --nocapture
2405
2406        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        // Run this test with stdout enabled to see output.
2412        // cargo test synchronous_instruments_cumulative_with_gap_in_measurements --features=testing -- --nocapture
2413
2414        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        // Create instrument and emit measurements
2427        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        // Test the first export
2457        assert_correct_export(&mut test_context, instrument_name);
2458
2459        // Reset and export again without making any measurements
2460        test_context.reset_metrics();
2461
2462        test_context.flush_metrics();
2463
2464        // Test that latest export has the same data as the previous one
2465        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        // Run this test with stdout enabled to see output.
2553        // cargo test asynchronous_instruments_cumulative_data_points_only_from_last_measurement --features=testing -- --nocapture
2554
2555        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    /// Helper function to test view customizations
2721    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        // Run this test with stdout enabled to see output.
2730        // cargo test view_test_* --all-features -- --nocapture
2731
2732        // Arrange
2733        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        // Act
2740        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        // Assert
2751        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    // Following are just a basic set of advanced View tests - Views bring a lot
2772    // of permutations and combinations, and we need
2773    // to expand coverage for more scenarios in future.
2774    // It is best to first split this file into multiple files
2775    // based on scenarios (eg: regular aggregation, cardinality, views, view_advanced, etc)
2776    // and then add more tests for each of the scenarios.
2777    #[test]
2778    fn test_view_single_instrument_multiple_stream() {
2779        // Run this test with stdout enabled to see output.
2780        // cargo test test_view_multiple_stream --all-features
2781
2782        // Each of the views match the instrument name "my_counter" and create a
2783        // new stream with a different name. In other words, View can be used to
2784        // create multiple streams for the same instrument.
2785
2786        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        // Arrange
2803        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        // Act
2811        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        // Assert
2818        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        // Run this test with stdout enabled to see output.
2832        // cargo test test_view_multiple_instrument_single_stream --all-features
2833
2834        // The view matches the instrument name "my_counter1" and "my_counter1"
2835        // and create a single new stream for both. In other words, View can be used to
2836        // "merge" multiple instruments into a single stream.
2837        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        // Arrange
2846        let exporter = InMemoryMetricExporter::default();
2847        let meter_provider = SdkMeterProvider::builder()
2848            .with_periodic_exporter(exporter.clone())
2849            .with_view(view)
2850            .build();
2851
2852        // Act
2853        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        // Assert
2862        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        // TODO: Assert that the data points are aggregated correctly.
2871    }
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        // Create instrument and emit measurements once
2880        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        // Test the first export
2929        assert_correct_export(&mut test_context, instrument_name);
2930
2931        // Reset and export again without making any measurements
2932        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        // Arrange
3000        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                    // Test concurrent collection by forcing half of the update threads to `force_flush` metrics and sleep for some time.
3013                    if i % 2 == 0 {
3014                        test_context.flush_metrics();
3015                        thread::sleep(Duration::from_millis(i)); // Make each thread sleep for some time duration for better testing
3016                    }
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        // Assert
3027        // We invoke `test_context.flush_metrics()` six times.
3028        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); // Expecting 1 time-series.
3044            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); // Each of the 10 update threads record measurements summing up to 5.
3058    }
3059
3060    fn counter_f64_multithreaded_aggregation_helper(temporality: Temporality) {
3061        // Arrange
3062        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                    // Test concurrent collection by forcing half of the update threads to `force_flush` metrics and sleep for some time.
3075                    if i % 2 == 0 {
3076                        test_context.flush_metrics();
3077                        thread::sleep(Duration::from_millis(i)); // Make each thread sleep for some time duration for better testing
3078                    }
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        // Assert
3089        // We invoke `test_context.flush_metrics()` six times.
3090        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); // Expecting 1 time-series.
3106            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); // Each of the 10 update threads record measurements 5 times = 10 * 5 * 1.23 = 61.5
3120    }
3121
3122    fn histogram_multithreaded_aggregation_helper(temporality: Temporality) {
3123        // Arrange
3124        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                    // Test concurrent collection by forcing half of the update threads to `force_flush` metrics and sleep for some time.
3138                    if i % 2 == 0 {
3139                        test_context.flush_metrics();
3140                        thread::sleep(Duration::from_millis(i)); // Make each thread sleep for some time duration for better testing
3141                    }
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        // Assert
3152        // We invoke `test_context.flush_metrics()` six times.
3153        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]; // There are 16 buckets for the default configuration
3175        let mut bucket_counts_key1_value1 = vec![0; 16];
3176
3177        histograms.iter().for_each(|histogram| {
3178            assert_eq!(histogram.data_points.len(), 2); // Expecting 1 time-series.
3179            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        // Default buckets:
3234        // (-∞, 0], (0, 5.0], (5.0, 10.0], (10.0, 25.0], (25.0, 50.0], (50.0, 75.0], (75.0, 100.0], (100.0, 250.0], (250.0, 500.0],
3235        // (500.0, 750.0], (750.0, 1000.0], (1000.0, 2500.0], (2500.0, 5000.0], (5000.0, 7500.0], (7500.0, 10000.0], (10000.0, +∞).
3236
3237        assert_eq!(count_zero_attributes, 20); // Each of the 10 update threads record two measurements.
3238        assert_eq!(sum_zero_attributes, 50); // Each of the 10 update threads record measurements summing up to 5.
3239        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), // For each of the 10 update threads, both the recorded values 1 and 4 fall under the bucket (0, 5].
3245                _ => assert_eq!(*count, 0),
3246            }
3247        }
3248
3249        assert_eq!(count_key1_value1, 50); // Each of the 10 update threads record 5 measurements.
3250        assert_eq!(sum_key1_value1, 1000); // Each of the 10 update threads record measurements summing up to 100 (5 + 7 + 18 + 35 + 35).
3251        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), // For each of the 10 update threads, the recorded value 5 falls under the bucket (0, 5].
3257                2 => assert_eq!(*count, 10), // For each of the 10 update threads, the recorded value 7 falls under the bucket (5, 10].
3258                3 => assert_eq!(*count, 10), // For each of the 10 update threads, the recorded value 18 falls under the bucket (10, 25].
3259                4 => assert_eq!(*count, 20), // For each of the 10 update threads, the recorded value 35 (recorded twice) falls under the bucket (25, 50].
3260                _ => assert_eq!(*count, 0),
3261            }
3262        }
3263    }
3264
3265    fn histogram_f64_multithreaded_aggregation_helper(temporality: Temporality) {
3266        // Arrange
3267        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                    // Test concurrent collection by forcing half of the update threads to `force_flush` metrics and sleep for some time.
3281                    if i % 2 == 0 {
3282                        test_context.flush_metrics();
3283                        thread::sleep(Duration::from_millis(i)); // Make each thread sleep for some time duration for better testing
3284                    }
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        // Assert
3295        // We invoke `test_context.flush_metrics()` six times.
3296        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]; // There are 16 buckets for the default configuration
3318        let mut bucket_counts_key1_value1 = vec![0; 16];
3319
3320        histograms.iter().for_each(|histogram| {
3321            assert_eq!(histogram.data_points.len(), 2); // Expecting 1 time-series.
3322            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        // Default buckets:
3377        // (-∞, 0], (0, 5.0], (5.0, 10.0], (10.0, 25.0], (25.0, 50.0], (50.0, 75.0], (75.0, 100.0], (100.0, 250.0], (250.0, 500.0],
3378        // (500.0, 750.0], (750.0, 1000.0], (1000.0, 2500.0], (2500.0, 5000.0], (5000.0, 7500.0], (7500.0, 10000.0], (10000.0, +∞).
3379
3380        assert_eq!(count_zero_attributes, 20); // Each of the 10 update threads record two measurements.
3381        assert!(f64::abs(61.0 - sum_zero_attributes) < 0.0001); // Each of the 10 update threads record measurements summing up to 6.1 (1.5 + 4.6)
3382        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), // For each of the 10 update threads, both the recorded values 1.5 and 4.6 fall under the bucket (0, 5.0].
3388                _ => assert_eq!(*count, 0),
3389            }
3390        }
3391
3392        assert_eq!(count_key1_value1, 50); // Each of the 10 update threads record 5 measurements.
3393        assert!(f64::abs(1006.0 - sum_key1_value1) < 0.0001); // Each of the 10 update threads record measurements summing up to 100.4 (5.0 + 7.3 + 18.1 + 35.1 + 35.1).
3394        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), // For each of the 10 update threads, the recorded value 5.0 falls under the bucket (0, 5.0].
3400                2 => assert_eq!(*count, 10), // For each of the 10 update threads, the recorded value 7.3 falls under the bucket (5.0, 10.0].
3401                3 => assert_eq!(*count, 10), // For each of the 10 update threads, the recorded value 18.1 falls under the bucket (10.0, 25.0].
3402                4 => assert_eq!(*count, 20), // For each of the 10 update threads, the recorded value 35.1 (recorded twice) falls under the bucket (25.0, 50.0].
3403                _ => assert_eq!(*count, 0),
3404            }
3405        }
3406    }
3407
3408    fn histogram_aggregation_helper(temporality: Temporality) {
3409        // Arrange
3410        let mut test_context = TestContext::new(temporality);
3411        let histogram = test_context.meter().u64_histogram("my_histogram").build();
3412
3413        // Act
3414        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        // Assert
3432        let MetricData::Histogram(histogram_data) =
3433            test_context.get_aggregation::<u64>("my_histogram", None)
3434        else {
3435            unreachable!()
3436        };
3437        // Expecting 2 time-series.
3438        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        // find and validate key1=value2 datapoint
3454        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        // Reset and report more measurements
3471        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        // Assert
3535        let MetricData::Histogram(histogram_data) =
3536            test_context.get_aggregation::<u64>("test_histogram", None)
3537        else {
3538            unreachable!()
3539        };
3540        // Expecting 2 time-series.
3541        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        // find and validate key1=value1 datapoint
3557        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        // Check the bucket counts
3565        // -∞ to 1.0: 1
3566        // 1.0 to 2.5: 1
3567        // 2.5 to 5.5: 3
3568        // 5.5 to +∞: 0
3569
3570        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        // Assert
3590        let MetricData::Histogram(histogram_data) =
3591            test_context.get_aggregation::<u64>("test_histogram", None)
3592        else {
3593            unreachable!()
3594        };
3595        // Expecting 1 time-series.
3596        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        // find and validate key1=value1 datapoint
3612        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            // find and validate key1=value1 datapoint
3669            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            // Check the bucket counts
3680            if specify_boundaries_in_view {
3681                // If boundaries are specified in the view, they should take precedence
3682                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                // If boundaries are not specified in the view, the ones from the instrument
3686                // should be used
3687                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        // Arrange: Create a view that converts a regular histogram to Base2ExponentialHistogram
3695        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        // Act: Record some values
3715        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        // Assert: Verify we get an ExponentialHistogram instead of regular Histogram
3724        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        // Validate the data point
3749        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        // Validate exponential histogram specific fields
3756        // Scale should be within valid range (-10 to 20)
3757        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        // zero_count should be 0 since we only recorded positive values > 0
3765        assert_eq!(
3766            data_point.zero_count(),
3767            0,
3768            "zero_count should be 0 for positive values"
3769        );
3770
3771        // Positive bucket should have counts (we recorded positive values)
3772        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        // Negative bucket should be empty (we only recorded positive values)
3781        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        // Verify the attribute is present
3790        let attrs: Vec<_> = data_point.attributes().collect();
3791        assert_eq!(attrs.len(), 1);
3792        assert_eq!(attrs[0].key.as_str(), "key1");
3793
3794        // Reset and report more measurements to verify Delta vs Cumulative behavior
3795        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        // Assert second collect
3803        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            // Cumulative: values accumulate (5 original + 3 new = 8 count, 15 + 60 = 75 sum)
3817            assert_eq!(data_point.count(), 8);
3818            assert_eq!(data_point.sum(), 75.0);
3819            assert_eq!(data_point.min(), Some(1.0)); // min from first batch
3820            assert_eq!(data_point.max(), Some(30.0)); // max from second batch
3821        } else {
3822            // Delta: only new values (3 count, 60 sum)
3823            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        // Verify positive bucket counts match the count
3830        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        // Arrange
3842        let mut test_context = TestContext::new(temporality);
3843        let gauge = test_context.meter().i64_gauge("my_gauge").build();
3844
3845        // Act
3846        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        // Assert
3859        let MetricData::Gauge(gauge_data_point) =
3860            test_context.get_aggregation::<i64>("my_gauge", None)
3861        else {
3862            unreachable!()
3863        };
3864        // Expecting 2 time-series.
3865        assert_eq!(gauge_data_point.data_points.len(), 2);
3866
3867        // find and validate key1=value2 datapoint
3868        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        // Reset and report more measurements
3879        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        // Arrange
3907        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        // Assert
3923        let MetricData::Gauge(gauge) =
3924            test_context.get_aggregation::<i64>("test_observable_gauge", None)
3925        else {
3926            unreachable!()
3927        };
3928        // Expecting 2 time-series.
3929        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            // find and validate zero attribute datapoint
3934            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        // find and validate key1=value1 datapoint
3941        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        // find and validate key2=value2 datapoint
3946        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        // Reset and report more measurements
3951        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        // Arrange
3980        let mut test_context = TestContext::new(temporality);
3981        let counter = test_context.u64_counter("test", "my_counter", None);
3982
3983        // Act
3984        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        // Assert
3997        let MetricData::Sum(sum) = test_context.get_aggregation::<u64>("my_counter", None) else {
3998            unreachable!()
3999        };
4000        // Expecting 2 time-series.
4001        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        // find and validate key1=value2 datapoint
4014        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        // Reset and report more measurements
4023        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        // Arrange
4059        let mut test_context = TestContext::new(temporality);
4060        let counter = test_context.u64_counter("test", "my_counter", None);
4061
4062        // Act
4063        // Record measurements with A:0, A:1,.......A:1999, which just fits in the 2000 limit
4064        for v in 0..2000 {
4065            counter.add(100, &[KeyValue::new("A", v.to_string())]);
4066        }
4067
4068        // Empty attributes is specially treated and does not count towards the limit.
4069        counter.add(3, &[]);
4070        counter.add(3, &[]);
4071
4072        // All of the below will now go into overflow.
4073        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        // Expecting 2002 metric points. (2000 + 1 overflow + Empty attributes)
4083        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 = &sum.data_points[0];
4090        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        // Phase 2 - for delta temporality, collect_and_reset uses in-place eviction:
4102        // the first collect marks entries as not-updated, and the second collect evicts
4103        // those still-stale entries. We need an extra flush to trigger that eviction
4104        // before adding new measurements that should fit under the cardinality limit.
4105        test_context.reset_metrics();
4106        if temporality == Temporality::Delta {
4107            test_context.flush_metrics();
4108            test_context.reset_metrics();
4109        }
4110        // The following should be aggregated normally for Delta,
4111        // and should go into overflow for Cumulative.
4112        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            // For cumulative, overflow should still be there, and new points should not be added.
4138            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        // Arrange
4157        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        // Act
4175        // Record measurements with A:0, A:1,.......A:cardinality_limit, which just fits in the cardinality_limit
4176        for v in 0..cardinality_limit {
4177            counter.add(100, &[KeyValue::new("A", v.to_string())]);
4178        }
4179
4180        // Empty attributes is specially treated and does not count towards the limit.
4181        counter.add(3, &[]);
4182        counter.add(3, &[]);
4183
4184        // All of the below will now go into overflow.
4185        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        // Expecting (cardinality_limit + 1 overflow + Empty attributes) data points.
4195        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 = &sum.data_points[0];
4202        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        // Phase 2 - for delta temporality, collect_and_reset uses in-place eviction:
4214        // the first collect marks entries as not-updated, and the second collect evicts
4215        // those still-stale entries. We need an extra flush to trigger that eviction
4216        // before adding new measurements that should fit under the cardinality limit.
4217        test_context.reset_metrics();
4218        if temporality == Temporality::Delta {
4219            test_context.flush_metrics();
4220            test_context.reset_metrics();
4221        }
4222        // The following should be aggregated normally for Delta,
4223        // and should go into overflow for Cumulative.
4224        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            // For cumulative, overflow should still be there, and new points should not be added.
4250            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        // Arrange
4269        let mut test_context = TestContext::new(temporality);
4270        let counter = test_context.u64_counter("test", "my_counter", None);
4271
4272        // Act
4273        // Add the same set of attributes in different order. (they are expected
4274        // to be treated as same attributes)
4275        // start with sorted order
4276        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        // Expecting 1 time-series.
4343        assert_eq!(sum.data_points.len(), 1);
4344
4345        // validate the sole datapoint
4346        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        // Arrange
4352        let mut test_context = TestContext::new(temporality);
4353        let counter = test_context.i64_up_down_counter("test", "my_updown_counter", None);
4354
4355        // Act
4356        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        // Assert
4369        let MetricData::Sum(sum) = test_context.get_aggregation::<i64>("my_updown_counter", None)
4370        else {
4371            unreachable!()
4372        };
4373        // Expecting 2 time-series.
4374        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        // find and validate key1=value2 datapoint
4386        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        // Reset and report more measurements
4395        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        // Saving this on the test context for lifetime simplicity
4540        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"); // TODO: Need to fix InMemoryMetricExporter to return None.
4621
4622            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        // After delta collect, add more and collect again
4800        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        // Mix bound and unbound additions to the same attribute set
4858        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        // Cycle 1: add and collect
4885        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        // Cycle 2: no add, collect — should export nothing
4894        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        // Cycle 3: add again — handle is still alive, produces fresh delta
4906        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        // Fill to cardinality limit with unbound calls
4938        for v in 0..cardinality_limit {
4939            counter.add(1, &[KeyValue::new("A", v.to_string())]);
4940        }
4941
4942        // bind() at overflow — handle binds directly to the overflow tracker
4943        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        // Expect: cardinality_limit unique + 1 overflow = cardinality_limit + 1
4953        assert_eq!(
4954            sum.data_points.len(),
4955            cardinality_limit + 1,
4956            "Expected {} unique + 1 overflow data points",
4957            cardinality_limit
4958        );
4959
4960        // The bound handle's value should appear in the overflow bucket
4961        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        // Fill to cardinality limit with unbound calls (these are one-shot, not bound)
4988        for v in 0..cardinality_limit {
4989            counter.add(1, &[KeyValue::new("A", v.to_string())]);
4990        }
4991
4992        // Collect cycle 1: exports the 3 unique entries, then evicts them (no new updates)
4993        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        // Cycle 2: no unbound adds, so the stale entries get evicted.
5000        // Space is now open. A new bind() should get a dedicated tracker.
5001        test_context.reset_metrics();
5002        test_context.flush_metrics(); // triggers eviction of stale entries
5003
5004        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        // The bound handle should have a dedicated tracker, NOT overflow
5015        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        // Fill to limit
5051        counter.add(1, &[KeyValue::new("A", "0")]);
5052        counter.add(1, &[KeyValue::new("A", "1")]);
5053
5054        // Bind two distinct attribute sets at overflow
5055        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        // Cycle 2: bound handles still work after delta collect
5075        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        // Once a bind() lands in overflow, the handle's writes must continue
5096        // landing in overflow for the lifetime of the handle — even after
5097        // delta eviction frees space. This is the predictability guarantee:
5098        // a user inspecting their data should see the bound handle's
5099        // attribution as a single, stable bucket. The recovery story is
5100        // explicit (drop and re-bind), not implicit (silent self-healing).
5101        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        // Cycle 1: fill cardinality with unbound calls, then bind at overflow.
5117        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        // Cycle 2: no calls. The 3 unbound entries become stale and are evicted,
5133        // freeing all of the cardinality budget.
5134        test_context.reset_metrics();
5135        test_context.flush_metrics();
5136
5137        // Cycle 3: the SAME bound handle is used again. Even though space is
5138        // available, its writes must still land in overflow — the handle is
5139        // permanently bound to overflow, not silently re-resolved.
5140        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        // Cycle 1: record and collect
5171        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        // Cycle 2: delta resets, new values
5185        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        // Mix bound and unbound recordings to the same attribute set
5217        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        // Cycle 1: record and collect
5252        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        // Cycle 2: no recordings — should export nothing
5262        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        // Cycle 3: record again — handle still alive
5274        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        // Fill to cardinality limit with unbound calls
5312        for v in 0..cardinality_limit {
5313            histogram.record(1.0, &[KeyValue::new("A", v.to_string())]);
5314        }
5315
5316        // bind() at overflow — handle binds directly to the overflow tracker
5317        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        // Mirror of bound_counter_overflow_persists_across_eviction_cycles for
5347        // histograms: a bound-at-overflow handle must keep landing in overflow
5348        // even after delta eviction frees space, for the lifetime of the handle.
5349        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        // Cycle 1: fill cardinality with unbound calls, then bind at overflow.
5369        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        // Cycle 2: no calls. Stale unbound entries get evicted, freeing space.
5386        test_context.reset_metrics();
5387        test_context.flush_metrics();
5388
5389        // Cycle 3: same bound handle. Even though space is free, writes must
5390        // still land in overflow.
5391        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        // Histogram configured with Base2ExponentialHistogram aggregation goes
5411        // through ExpoHistogram internally. Verify bind() returns a handle
5412        // whose direct writes accumulate correctly.
5413        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        // NaN/inf must be filtered just like the unbound path.
5436        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        // bind() at overflow — handle binds directly to the overflow tracker
5480        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        // Same predictability invariant the counter/histogram tests assert,
5503        // verified for the exponential histogram aggregator path.
5504        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        // Cycle 2: evict stale entries, freeing space.
5541        test_context.reset_metrics();
5542        test_context.flush_metrics();
5543
5544        // Cycle 3: bound handle keeps writing to overflow.
5545        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            // bound drops here
5581        }
5582
5583        // Cycle 2: no updates, bound handle dropped — entry becomes stale and evictable
5584        test_context.reset_metrics();
5585        test_context.flush_metrics();
5586
5587        // Cycle 3: the stale entry should have been evicted, so a new unbound add
5588        // should be the only data point
5589        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        // Only the new entry should be present — the old "ephemeral" entry was evicted
5598        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 one handle — entry should NOT be evicted
5633        drop(bound1);
5634        test_context.reset_metrics();
5635        test_context.flush_metrics(); // idle cycle, but bound2 still holds it
5636
5637        // bound2 still works
5638        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        // Mix bound and unbound calls with empty attributes — they must share
5679        // the same data point (both route to no_attribute_tracker).
5680        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        // bind with k3 included — view should drop it at bind time
5755        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        // unbound call with a *different* k3 value: after view filtering both
5764        // bound and unbound must collapse into the same data point.
5765        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        // Cumulative: cardinality only grows, never evicts. A bind() at the
5866        // limit lands in overflow and accumulates there forever. Verifies the
5867        // bound handle's cumulative writes converge in the overflow bucket
5868        // across multiple collection cycles.
5869        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        // Cycle 2: cumulative state accumulates internally; reset the exporter
5899        // so the assertion sees a single export rather than two appended ones.
5900        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        // Cycle 2: cumulative accumulates internally; reset exporter to see a
5953        // single fresh export rather than two appended ones.
5954        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        // Cycle 2: cumulative accumulates internally; reset exporter to see a
6010        // single fresh export rather than two appended ones.
6011        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        // Unbound call with a different k3: after view filtering, bound and unbound
6064        // must collapse into the same exponential histogram data point.
6065        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        // Next cycle: a new record produces the new last value.
6142        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        // Bound and unbound writes target the same data point. The final
6166        // exported value is whichever write happened last.
6167        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        // Cycle 1: record and collect.
6198        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        // Cycle 2: no record - should export nothing.
6208        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        // Cycle 3: record again - handle still alive.
6220        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        // Fill to cardinality limit with unbound calls.
6253        for v in 0..cardinality_limit {
6254            gauge.record(1, &[KeyValue::new("A", v.to_string())]);
6255        }
6256
6257        // bind() at overflow - handle binds directly to the overflow tracker.
6258        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 one handle - entry should NOT be evicted while the other handle is alive.
6309        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        // Mix bound and unbound calls with empty attributes - they must share
6336        // the same data point. Last write wins.
6337        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        // bind with k3 included - view should drop it at bind time.
6383        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        // Unbound call with a *different* k3 value: after view filtering both
6392        // bound and unbound must collapse into the same data point.
6393        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            // bound drops here
6447        }
6448
6449        // Cycle 2: no updates, bound handle dropped - entry becomes stale and evictable.
6450        test_context.reset_metrics();
6451        test_context.flush_metrics();
6452
6453        // Cycle 3: the stale entry should have been evicted, so a new unbound
6454        // record should be the only data point.
6455        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        // UpDownCounter supports both positive and negative values.
6484        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        // UpDownCounter always uses Cumulative temporality regardless of the
6510        // reader's preference. Confirm the bound API observes the same rule.
6511        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        // Mix bound and unbound additions to the same attribute set.
6542        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        // Fill to cardinality limit with unbound calls.
6580        for v in 0..cardinality_limit {
6581            counter.add(1, &[KeyValue::new("A", v.to_string())]);
6582        }
6583
6584        // bind() at overflow - handle binds directly to the overflow tracker.
6585        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 one handle. Cumulative state still includes its contributions
6635        // and the surviving handle continues to write to the same tracker.
6636        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        // Mix bound and unbound calls with empty attributes - they must share
6659        // the same data point (both route to no_attribute_tracker).
6660        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        // bind with k3 included - view should drop it at bind time.
6704        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        // unbound call with a *different* k3 value: after view filtering both
6713        // bound and unbound must collapse into the same data point.
6714        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}