From c1fae2bb649c694d2db3cdc0aa4e5bd92189270e Mon Sep 17 00:00:00 2001 From: Lukas Hering Date: Sun, 19 Jul 2026 21:23:20 -0400 Subject: [PATCH 1/4] opentelemetry-sdk: add support for advisory attributes on metric instruments --- .../metrics/_internal/__init__.py | 210 +++++++++++++++--- .../metrics/_internal/instrument.py | 91 +++++++- opentelemetry-api/tests/metrics/test_meter.py | 92 +++++++- .../tests/metrics/test_meter_provider.py | 18 +- .../sdk/metrics/_internal/__init__.py | 169 +++++++++++--- .../_internal/_view_instrument_match.py | 13 +- .../sdk/metrics/_internal/aggregation.py | 9 +- .../sdk/metrics/_internal/instrument.py | 29 +-- .../test_attributes_advisory.py | 184 +++++++++++++++ .../tests/metrics/test_aggregation.py | 5 +- .../metrics/test_view_instrument_match.py | 103 ++++++++- 11 files changed, 805 insertions(+), 118 deletions(-) create mode 100644 opentelemetry-sdk/tests/metrics/integration_test/test_attributes_advisory.py diff --git a/opentelemetry-api/src/opentelemetry/metrics/_internal/__init__.py b/opentelemetry-api/src/opentelemetry/metrics/_internal/__init__.py index 9df55291fa5..e60a002976e 100644 --- a/opentelemetry-api/src/opentelemetry/metrics/_internal/__init__.py +++ b/opentelemetry-api/src/opentelemetry/metrics/_internal/__init__.py @@ -1,7 +1,7 @@ # Copyright The OpenTelemetry Authors # SPDX-License-Identifier: Apache-2.0 -# pylint: disable=too-many-ancestors +# pylint: disable=too-many-ancestors,too-many-lines """ The OpenTelemetry metrics API describes the classes used to generate @@ -55,7 +55,8 @@ ObservableGauge, ObservableUpDownCounter, UpDownCounter, - _MetricsHistogramAdvisory, + _MetricsAdvisory, + _normalize_attributes_advisory, _ProxyCounter, _ProxyGauge, _ProxyHistogram, @@ -176,7 +177,7 @@ class _InstrumentRegistrationStatus: instrument_id: str already_registered: bool conflict: bool - current_advisory: _MetricsHistogramAdvisory | None + current_advisory: _MetricsAdvisory | None class Meter(ABC): @@ -196,7 +197,7 @@ def __init__( self._name = name self._version = version self._schema_url = schema_url - self._instrument_ids: dict[str, _MetricsHistogramAdvisory | None] = {} + self._instrument_ids: dict[str, _MetricsAdvisory | None] = {} self._instrument_ids_lock = Lock() @property @@ -226,7 +227,7 @@ def _register_instrument( type_: type, unit: str, description: str, - advisory: _MetricsHistogramAdvisory | None = None, + advisory: _MetricsAdvisory | None = None, ) -> _InstrumentRegistrationStatus: """ Register an instrument with the name, type, unit and description as @@ -289,6 +290,8 @@ def create_counter( name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> Counter: """Creates a `Counter` instrument @@ -305,6 +308,8 @@ def create_up_down_counter( name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> UpDownCounter: """Creates an `UpDownCounter` instrument @@ -322,6 +327,8 @@ def create_observable_counter( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> ObservableCounter: """Creates an `ObservableCounter` instrument @@ -420,6 +427,7 @@ def create_histogram( description: str = "", *, explicit_bucket_boundaries_advisory: Sequence[float] | None = None, + _attributes_advisory: Sequence[str] | None = None, ) -> Histogram: """Creates a :class:`~opentelemetry.metrics.Histogram` instrument @@ -435,6 +443,8 @@ def create_gauge( # type: ignore # pylint: disable=no-self-use name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> Gauge: # pyright: ignore[reportReturnType] """Creates a ``Gauge`` instrument @@ -453,6 +463,8 @@ def create_observable_gauge( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> ObservableGauge: """Creates an `ObservableGauge` instrument @@ -473,6 +485,8 @@ def create_observable_up_down_counter( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> ObservableUpDownCounter: """Creates an `ObservableUpDownCounter` instrument @@ -521,11 +535,23 @@ def create_counter( name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> Counter: with self._lock: if self._real_meter: - return self._real_meter.create_counter(name, unit, description) - proxy = _ProxyCounter(name, unit, description) + return self._real_meter.create_counter( + name, + unit, + description, + _attributes_advisory=_attributes_advisory, + ) + proxy = _ProxyCounter( + name, + unit, + description, + _attributes_advisory=_attributes_advisory, + ) self._instruments.append(proxy) return proxy @@ -534,13 +560,23 @@ def create_up_down_counter( name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> UpDownCounter: with self._lock: if self._real_meter: return self._real_meter.create_up_down_counter( - name, unit, description + name, + unit, + description, + _attributes_advisory=_attributes_advisory, ) - proxy = _ProxyUpDownCounter(name, unit, description) + proxy = _ProxyUpDownCounter( + name, + unit, + description, + _attributes_advisory=_attributes_advisory, + ) self._instruments.append(proxy) return proxy @@ -550,14 +586,24 @@ def create_observable_counter( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> ObservableCounter: with self._lock: if self._real_meter: return self._real_meter.create_observable_counter( - name, callbacks, unit, description + name, + callbacks, + unit, + description, + _attributes_advisory=_attributes_advisory, ) proxy = _ProxyObservableCounter( - name, callbacks, unit=unit, description=description + name, + callbacks, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, ) self._instruments.append(proxy) return proxy @@ -569,6 +615,7 @@ def create_histogram( description: str = "", *, explicit_bucket_boundaries_advisory: Sequence[float] | None = None, + _attributes_advisory: Sequence[str] | None = None, ) -> Histogram: with self._lock: if self._real_meter: @@ -577,9 +624,14 @@ def create_histogram( unit, description, explicit_bucket_boundaries_advisory=explicit_bucket_boundaries_advisory, + _attributes_advisory=_attributes_advisory, ) proxy = _ProxyHistogram( - name, unit, description, explicit_bucket_boundaries_advisory + name, + unit, + description, + explicit_bucket_boundaries_advisory, + _attributes_advisory=_attributes_advisory, ) self._instruments.append(proxy) return proxy @@ -589,11 +641,23 @@ def create_gauge( name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> Gauge: with self._lock: if self._real_meter: - return self._real_meter.create_gauge(name, unit, description) - proxy = _ProxyGauge(name, unit, description) + return self._real_meter.create_gauge( + name, + unit, + description, + _attributes_advisory=_attributes_advisory, + ) + proxy = _ProxyGauge( + name, + unit, + description, + _attributes_advisory=_attributes_advisory, + ) self._instruments.append(proxy) return proxy @@ -603,14 +667,24 @@ def create_observable_gauge( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> ObservableGauge: with self._lock: if self._real_meter: return self._real_meter.create_observable_gauge( - name, callbacks, unit, description + name, + callbacks, + unit, + description, + _attributes_advisory=_attributes_advisory, ) proxy = _ProxyObservableGauge( - name, callbacks, unit=unit, description=description + name, + callbacks, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, ) self._instruments.append(proxy) return proxy @@ -621,6 +695,8 @@ def create_observable_up_down_counter( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> ObservableUpDownCounter: with self._lock: if self._real_meter: @@ -629,9 +705,14 @@ def create_observable_up_down_counter( callbacks, unit, description, + _attributes_advisory=_attributes_advisory, ) proxy = _ProxyObservableUpDownCounter( - name, callbacks, unit=unit, description=description + name, + callbacks, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, ) self._instruments.append(proxy) return proxy @@ -648,10 +729,18 @@ def create_counter( name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> Counter: """Returns a no-op Counter.""" status = self._register_instrument( - name, NoOpCounter, unit, description + name, + NoOpCounter, + unit, + description, + _MetricsAdvisory( + attributes=_normalize_attributes_advisory(_attributes_advisory) + ), ) if status.conflict: self._log_instrument_registration_conflict( @@ -662,16 +751,31 @@ def create_counter( status, ) - return NoOpCounter(name, unit=unit, description=description) + return NoOpCounter( + name, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, + ) def create_gauge( self, name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> Gauge: """Returns a no-op Gauge.""" - status = self._register_instrument(name, NoOpGauge, unit, description) + status = self._register_instrument( + name, + NoOpGauge, + unit, + description, + _MetricsAdvisory( + attributes=_normalize_attributes_advisory(_attributes_advisory) + ), + ) if status.conflict: self._log_instrument_registration_conflict( name, @@ -680,17 +784,30 @@ def create_gauge( description, status, ) - return NoOpGauge(name, unit=unit, description=description) + return NoOpGauge( + name, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, + ) def create_up_down_counter( self, name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> UpDownCounter: """Returns a no-op UpDownCounter.""" status = self._register_instrument( - name, NoOpUpDownCounter, unit, description + name, + NoOpUpDownCounter, + unit, + description, + _MetricsAdvisory( + attributes=_normalize_attributes_advisory(_attributes_advisory) + ), ) if status.conflict: self._log_instrument_registration_conflict( @@ -700,7 +817,12 @@ def create_up_down_counter( description, status, ) - return NoOpUpDownCounter(name, unit=unit, description=description) + return NoOpUpDownCounter( + name, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, + ) def create_observable_counter( self, @@ -708,10 +830,18 @@ def create_observable_counter( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> ObservableCounter: """Returns a no-op ObservableCounter.""" status = self._register_instrument( - name, NoOpObservableCounter, unit, description + name, + NoOpObservableCounter, + unit, + description, + _MetricsAdvisory( + attributes=_normalize_attributes_advisory(_attributes_advisory) + ), ) if status.conflict: self._log_instrument_registration_conflict( @@ -726,6 +856,7 @@ def create_observable_counter( callbacks, unit=unit, description=description, + _attributes_advisory=_attributes_advisory, ) def create_histogram( @@ -735,6 +866,7 @@ def create_histogram( description: str = "", *, explicit_bucket_boundaries_advisory: Sequence[float] | None = None, + _attributes_advisory: Sequence[str] | None = None, ) -> Histogram: """Returns a no-op Histogram.""" status = self._register_instrument( @@ -742,8 +874,11 @@ def create_histogram( NoOpHistogram, unit, description, - _MetricsHistogramAdvisory( - explicit_bucket_boundaries=explicit_bucket_boundaries_advisory + _MetricsAdvisory( + attributes=_normalize_attributes_advisory( + _attributes_advisory + ), + explicit_bucket_boundaries=explicit_bucket_boundaries_advisory, ), ) if status.conflict: @@ -759,6 +894,7 @@ def create_histogram( unit=unit, description=description, explicit_bucket_boundaries_advisory=explicit_bucket_boundaries_advisory, + _attributes_advisory=_attributes_advisory, ) def create_observable_gauge( @@ -767,10 +903,18 @@ def create_observable_gauge( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> ObservableGauge: """Returns a no-op ObservableGauge.""" status = self._register_instrument( - name, NoOpObservableGauge, unit, description + name, + NoOpObservableGauge, + unit, + description, + _MetricsAdvisory( + attributes=_normalize_attributes_advisory(_attributes_advisory) + ), ) if status.conflict: self._log_instrument_registration_conflict( @@ -785,6 +929,7 @@ def create_observable_gauge( callbacks, unit=unit, description=description, + _attributes_advisory=_attributes_advisory, ) def create_observable_up_down_counter( @@ -793,10 +938,18 @@ def create_observable_up_down_counter( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> ObservableUpDownCounter: """Returns a no-op ObservableUpDownCounter.""" status = self._register_instrument( - name, NoOpObservableUpDownCounter, unit, description + name, + NoOpObservableUpDownCounter, + unit, + description, + _MetricsAdvisory( + attributes=_normalize_attributes_advisory(_attributes_advisory) + ), ) if status.conflict: self._log_instrument_registration_conflict( @@ -811,6 +964,7 @@ def create_observable_up_down_counter( callbacks, unit=unit, description=description, + _attributes_advisory=_attributes_advisory, ) diff --git a/opentelemetry-api/src/opentelemetry/metrics/_internal/instrument.py b/opentelemetry-api/src/opentelemetry/metrics/_internal/instrument.py index 79bdfa4e673..ec163c59785 100644 --- a/opentelemetry-api/src/opentelemetry/metrics/_internal/instrument.py +++ b/opentelemetry-api/src/opentelemetry/metrics/_internal/instrument.py @@ -29,10 +29,19 @@ @dataclass(frozen=True) -class _MetricsHistogramAdvisory: +class _MetricsAdvisory: + attributes: frozenset[str] | None = None explicit_bucket_boundaries: Sequence[float] | None = None +def _normalize_attributes_advisory( + attributes_advisory: Sequence[str] | None, +) -> frozenset[str] | None: + return ( + None if attributes_advisory is None else frozenset(attributes_advisory) + ) + + @dataclass(frozen=True) class CallbackOptions: """Options for the callback @@ -62,6 +71,8 @@ def __init__( name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: pass @@ -107,10 +118,13 @@ def __init__( name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: self._name = name self._unit = unit self._description = description + self._attributes_advisory = _attributes_advisory self._real_instrument: InstrumentT | None = None def on_meter_set(self, meter: "metrics.Meter") -> None: @@ -133,8 +147,15 @@ def __init__( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: - super().__init__(name, unit, description) + super().__init__( + name, + unit, + description, + _attributes_advisory=_attributes_advisory, + ) self._callbacks = callbacks @@ -152,8 +173,15 @@ def __init__( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: - super().__init__(name, unit=unit, description=description) + super().__init__( + name, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, + ) class Counter(Synchronous): @@ -184,8 +212,15 @@ def __init__( name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: - super().__init__(name, unit=unit, description=description) + super().__init__( + name, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, + ) def add( self, @@ -211,6 +246,7 @@ def _create_real_instrument(self, meter: "metrics.Meter") -> Counter: self._name, self._unit, self._description, + _attributes_advisory=self._attributes_advisory, ) @@ -246,8 +282,15 @@ def __init__( name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: - super().__init__(name, unit=unit, description=description) + super().__init__( + name, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, + ) def add( self, @@ -273,6 +316,7 @@ def _create_real_instrument(self, meter: "metrics.Meter") -> UpDownCounter: self._name, self._unit, self._description, + _attributes_advisory=self._attributes_advisory, ) @@ -291,12 +335,15 @@ def __init__( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: super().__init__( name, callbacks, unit=unit, description=description, + _attributes_advisory=_attributes_advisory, ) @@ -311,6 +358,7 @@ def _create_real_instrument( self._callbacks, self._unit, self._description, + _attributes_advisory=self._attributes_advisory, ) @@ -330,12 +378,15 @@ def __init__( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: super().__init__( name, callbacks, unit=unit, description=description, + _attributes_advisory=_attributes_advisory, ) @@ -351,6 +402,7 @@ def _create_real_instrument( self._callbacks, self._unit, self._description, + _attributes_advisory=self._attributes_advisory, ) @@ -367,6 +419,8 @@ def __init__( unit: str = "", description: str = "", explicit_bucket_boundaries_advisory: Sequence[float] | None = None, + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: pass @@ -402,12 +456,15 @@ def __init__( unit: str = "", description: str = "", explicit_bucket_boundaries_advisory: Sequence[float] | None = None, + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: super().__init__( name, unit=unit, description=description, explicit_bucket_boundaries_advisory=explicit_bucket_boundaries_advisory, + _attributes_advisory=_attributes_advisory, ) def record( @@ -426,8 +483,15 @@ def __init__( unit: str = "", description: str = "", explicit_bucket_boundaries_advisory: Sequence[float] | None = None, + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: - super().__init__(name, unit=unit, description=description) + super().__init__( + name, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, + ) self._explicit_bucket_boundaries_advisory = ( explicit_bucket_boundaries_advisory ) @@ -447,6 +511,7 @@ def _create_real_instrument(self, meter: "metrics.Meter") -> Histogram: self._unit, self._description, explicit_bucket_boundaries_advisory=self._explicit_bucket_boundaries_advisory, + _attributes_advisory=self._attributes_advisory, ) @@ -466,12 +531,15 @@ def __init__( callbacks: Sequence[CallbackT] | None = None, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: super().__init__( name, callbacks, unit=unit, description=description, + _attributes_advisory=_attributes_advisory, ) @@ -487,6 +555,7 @@ def _create_real_instrument( self._callbacks, self._unit, self._description, + _attributes_advisory=self._attributes_advisory, ) @@ -522,8 +591,15 @@ def __init__( name: str, unit: str = "", description: str = "", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> None: - super().__init__(name, unit=unit, description=description) + super().__init__( + name, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, + ) def set( self, @@ -552,4 +628,5 @@ def _create_real_instrument(self, meter: "metrics.Meter") -> Gauge: self._name, self._unit, self._description, + _attributes_advisory=self._attributes_advisory, ) diff --git a/opentelemetry-api/tests/metrics/test_meter.py b/opentelemetry-api/tests/metrics/test_meter.py index 79b30b6c7f8..30e76bc5276 100644 --- a/opentelemetry-api/tests/metrics/test_meter.py +++ b/opentelemetry-api/tests/metrics/test_meter.py @@ -13,22 +13,41 @@ class ChildMeter(Meter): # pylint: disable=signature-differs - def create_counter(self, name, unit="", description=""): - super().create_counter(name, unit=unit, description=description) + def create_counter( + self, name, unit="", description="", *, _attributes_advisory=None + ): + super().create_counter( + name, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, + ) - def create_up_down_counter(self, name, unit="", description=""): + def create_up_down_counter( + self, name, unit="", description="", *, _attributes_advisory=None + ): super().create_up_down_counter( - name, unit=unit, description=description + name, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, ) def create_observable_counter( - self, name, callbacks, unit="", description="" + self, + name, + callbacks, + unit="", + description="", + *, + _attributes_advisory=None, ): super().create_observable_counter( name, callbacks, unit=unit, description=description, + _attributes_advisory=_attributes_advisory, ) def create_histogram( @@ -38,35 +57,58 @@ def create_histogram( description="", *, explicit_bucket_boundaries_advisory=None, + _attributes_advisory=None, ): super().create_histogram( name, unit=unit, description=description, explicit_bucket_boundaries_advisory=explicit_bucket_boundaries_advisory, + _attributes_advisory=_attributes_advisory, ) - def create_gauge(self, name, unit="", description=""): - super().create_gauge(name, unit=unit, description=description) + def create_gauge( + self, name, unit="", description="", *, _attributes_advisory=None + ): + super().create_gauge( + name, + unit=unit, + description=description, + _attributes_advisory=_attributes_advisory, + ) def create_observable_gauge( - self, name, callbacks, unit="", description="" + self, + name, + callbacks, + unit="", + description="", + *, + _attributes_advisory=None, ): super().create_observable_gauge( name, callbacks, unit=unit, description=description, + _attributes_advisory=_attributes_advisory, ) def create_observable_up_down_counter( - self, name, callbacks, unit="", description="" + self, + name, + callbacks, + unit="", + description="", + *, + _attributes_advisory=None, ): super().create_observable_up_down_counter( name, callbacks, unit=unit, description=description, + _attributes_advisory=_attributes_advisory, ) @@ -131,6 +173,38 @@ def test_repeated_instrument_names_with_different_advisory(self): instrument_name, ) + def test_repeated_instrument_names_with_different_attributes_advisory( + self, + ): + test_meter = NoOpMeter("name") + + for instrument_name in [ + "counter", + "up_down_counter", + "histogram", + "gauge", + "observable_counter", + "observable_up_down_counter", + "observable_gauge", + ]: + with self.subTest(instrument_name=instrument_name): + create_instrument = getattr( + test_meter, f"create_{instrument_name}" + ) + create_instrument( + instrument_name, _attributes_advisory=["a", "b"] + ) + + with self.assertNoLogs(level=WARNING): + create_instrument( + instrument_name, _attributes_advisory=["b", "a"] + ) + + with self.assertLogs(level=WARNING): + create_instrument( + instrument_name, _attributes_advisory=["c"] + ) + def test_create_counter(self): """ Test that the meter provides a function to create a new Counter diff --git a/opentelemetry-api/tests/metrics/test_meter_provider.py b/opentelemetry-api/tests/metrics/test_meter_provider.py index 3ea3e2041bc..8575d38ed1d 100644 --- a/opentelemetry-api/tests/metrics/test_meter_provider.py +++ b/opentelemetry-api/tests/metrics/test_meter_provider.py @@ -257,25 +257,29 @@ def test_proxy_meter(self): real_meter: Mock = real_meter_provider.get_meter() real_meter.create_counter.assert_called_once_with( - name, unit, description + name, unit, description, _attributes_advisory=None ) real_meter.create_up_down_counter.assert_called_once_with( - name, unit, description + name, unit, description, _attributes_advisory=None ) real_meter.create_histogram.assert_called_once_with( - name, unit, description, explicit_bucket_boundaries_advisory=None + name, + unit, + description, + explicit_bucket_boundaries_advisory=None, + _attributes_advisory=None, ) real_meter.create_gauge.assert_called_once_with( - name, unit, description + name, unit, description, _attributes_advisory=None ) real_meter.create_observable_counter.assert_called_once_with( - name, [callback], unit, description + name, [callback], unit, description, _attributes_advisory=None ) real_meter.create_observable_up_down_counter.assert_called_once_with( - name, [callback], unit, description + name, [callback], unit, description, _attributes_advisory=None ) real_meter.create_observable_gauge.assert_called_once_with( - name, [callback], unit, description + name, [callback], unit, description, _attributes_advisory=None ) # The synchronous instrument measurement methods should call through to diff --git a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/__init__.py b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/__init__.py index 61db3dadd39..c03e709fadb 100644 --- a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/__init__.py +++ b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/__init__.py @@ -25,6 +25,10 @@ ) from opentelemetry.metrics import UpDownCounter as APIUpDownCounter from opentelemetry.metrics import _Gauge as APIGauge +from opentelemetry.metrics._internal.instrument import ( + _MetricsAdvisory, + _normalize_attributes_advisory, +) from opentelemetry.sdk.environment_variables import ( OTEL_METRICS_EXEMPLAR_FILTER, OTEL_SDK_DISABLED, @@ -68,6 +72,53 @@ _logger = getLogger(__name__) +def _validate_explicit_bucket_boundaries_advisory( + explicit_bucket_boundaries_advisory: Sequence[float] | None, +) -> Sequence[float] | None: + """Returns the advisory if it is a sequence of numbers, otherwise logs a + warning and returns `None` so the advisory is ignored.""" + + if explicit_bucket_boundaries_advisory is None: + return None + + if isinstance(explicit_bucket_boundaries_advisory, Sequence): + try: + if all( + isinstance(entry, (float, int)) + for entry in explicit_bucket_boundaries_advisory + ): + return explicit_bucket_boundaries_advisory + except (KeyError, TypeError): + pass + + _logger.warning( + "explicit_bucket_boundaries_advisory must be a sequence of numbers" + ) + return None + + +def _validate_attributes_advisory( + attributes_advisory: Sequence[str] | None, +) -> frozenset[str] | None: + """Returns the advisory as a frozenset of attribute keys, otherwise logs a + warning and returns `None` so the advisory is ignored.""" + + if attributes_advisory is None: + return None + + if isinstance(attributes_advisory, Sequence) and not isinstance( + attributes_advisory, str + ): + try: + if all(isinstance(entry, str) for entry in attributes_advisory): + return _normalize_attributes_advisory(attributes_advisory) + except (KeyError, TypeError): + pass + + _logger.warning("_attributes_advisory must be a sequence of strings") + return None + + @dataclass class _MeterConfig: is_enabled: bool = True @@ -118,10 +169,21 @@ def _is_enabled(self) -> bool: def _set_meter_config(self, meter_config: _MeterConfig) -> None: self._meter_config.update(meter_config) - def create_counter(self, name, unit="", description="") -> APICounter: + def create_counter( + self, + name, + unit="", + description="", + *, + _attributes_advisory: Sequence[str] | None = None, + ) -> APICounter: + advisory = _MetricsAdvisory( + attributes=_validate_attributes_advisory(_attributes_advisory) + ) + with self._instrument_registration_lock: status = self._register_instrument( - name, _Counter, unit, description + name, _Counter, unit, description, advisory ) if not status.already_registered: self._instrument_id_instrument[status.instrument_id] = ( @@ -131,6 +193,7 @@ def create_counter(self, name, unit="", description="") -> APICounter: self._measurement_consumer, unit, description, + _advisory=advisory, _meter_config=self._meter_config, ) ) @@ -150,11 +213,20 @@ def create_counter(self, name, unit="", description="") -> APICounter: return instrument def create_up_down_counter( - self, name, unit="", description="" + self, + name, + unit="", + description="", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> APIUpDownCounter: + advisory = _MetricsAdvisory( + attributes=_validate_attributes_advisory(_attributes_advisory) + ) + with self._instrument_registration_lock: status = self._register_instrument( - name, _UpDownCounter, unit, description + name, _UpDownCounter, unit, description, advisory ) if not status.already_registered: self._instrument_id_instrument[status.instrument_id] = ( @@ -164,6 +236,7 @@ def create_up_down_counter( self._measurement_consumer, unit, description, + _advisory=advisory, _meter_config=self._meter_config, ) ) @@ -188,10 +261,16 @@ def create_observable_counter( callbacks=None, unit="", description="", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> APIObservableCounter: + advisory = _MetricsAdvisory( + attributes=_validate_attributes_advisory(_attributes_advisory) + ) + with self._instrument_registration_lock: status = self._register_instrument( - name, _ObservableCounter, unit, description + name, _ObservableCounter, unit, description, advisory ) if not status.already_registered: self._instrument_id_instrument[status.instrument_id] = ( @@ -202,6 +281,7 @@ def create_observable_counter( callbacks, unit, description, + _advisory=advisory, _meter_config=self._meter_config, ) ) @@ -232,27 +312,14 @@ def create_histogram( description: str = "", *, explicit_bucket_boundaries_advisory: Sequence[float] | None = None, + _attributes_advisory: Sequence[str] | None = None, ) -> APIHistogram: - if explicit_bucket_boundaries_advisory is not None: - invalid_advisory = False - if isinstance(explicit_bucket_boundaries_advisory, Sequence): - try: - invalid_advisory = not ( - all( - isinstance(e, (float, int)) - for e in explicit_bucket_boundaries_advisory - ) - ) - except (KeyError, TypeError): - invalid_advisory = True - else: - invalid_advisory = True - - if invalid_advisory: - explicit_bucket_boundaries_advisory = None - _logger.warning( - "explicit_bucket_boundaries_advisory must be a sequence of numbers" - ) + advisory = _MetricsAdvisory( + attributes=_validate_attributes_advisory(_attributes_advisory), + explicit_bucket_boundaries=_validate_explicit_bucket_boundaries_advisory( + explicit_bucket_boundaries_advisory + ), + ) with self._instrument_registration_lock: status = self._register_instrument( @@ -260,7 +327,7 @@ def create_histogram( _Histogram, unit, description, - explicit_bucket_boundaries_advisory, + advisory, ) if not status.already_registered: self._instrument_id_instrument[status.instrument_id] = ( @@ -270,7 +337,7 @@ def create_histogram( self._measurement_consumer, unit, description, - explicit_bucket_boundaries_advisory, + _advisory=advisory, _meter_config=self._meter_config, ) ) @@ -289,9 +356,22 @@ def create_histogram( ) return instrument - def create_gauge(self, name, unit="", description="") -> APIGauge: + def create_gauge( + self, + name, + unit="", + description="", + *, + _attributes_advisory: Sequence[str] | None = None, + ) -> APIGauge: + advisory = _MetricsAdvisory( + attributes=_validate_attributes_advisory(_attributes_advisory) + ) + with self._instrument_registration_lock: - status = self._register_instrument(name, _Gauge, unit, description) + status = self._register_instrument( + name, _Gauge, unit, description, advisory + ) if not status.already_registered: self._instrument_id_instrument[status.instrument_id] = _Gauge( name, @@ -299,6 +379,7 @@ def create_gauge(self, name, unit="", description="") -> APIGauge: self._measurement_consumer, unit, description, + _advisory=advisory, _meter_config=self._meter_config, ) instrument = self._instrument_id_instrument[status.instrument_id] @@ -317,11 +398,21 @@ def create_gauge(self, name, unit="", description="") -> APIGauge: return instrument def create_observable_gauge( - self, name, callbacks=None, unit="", description="" + self, + name, + callbacks=None, + unit="", + description="", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> APIObservableGauge: + advisory = _MetricsAdvisory( + attributes=_validate_attributes_advisory(_attributes_advisory) + ) + with self._instrument_registration_lock: status = self._register_instrument( - name, _ObservableGauge, unit, description + name, _ObservableGauge, unit, description, advisory ) if not status.already_registered: self._instrument_id_instrument[status.instrument_id] = ( @@ -332,6 +423,7 @@ def create_observable_gauge( callbacks, unit, description, + _advisory=advisory, _meter_config=self._meter_config, ) ) @@ -356,11 +448,21 @@ def create_observable_gauge( return instrument def create_observable_up_down_counter( - self, name, callbacks=None, unit="", description="" + self, + name, + callbacks=None, + unit="", + description="", + *, + _attributes_advisory: Sequence[str] | None = None, ) -> APIObservableUpDownCounter: + advisory = _MetricsAdvisory( + attributes=_validate_attributes_advisory(_attributes_advisory) + ) + with self._instrument_registration_lock: status = self._register_instrument( - name, _ObservableUpDownCounter, unit, description + name, _ObservableUpDownCounter, unit, description, advisory ) if not status.already_registered: self._instrument_id_instrument[status.instrument_id] = ( @@ -371,6 +473,7 @@ def create_observable_up_down_counter( callbacks, unit, description, + _advisory=advisory, _meter_config=self._meter_config, ) ) diff --git a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/_view_instrument_match.py b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/_view_instrument_match.py index d9ba05363b1..83b6af99fc3 100644 --- a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/_view_instrument_match.py +++ b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/_view_instrument_match.py @@ -8,6 +8,7 @@ from time import time_ns from typing import cast +from opentelemetry.metrics._internal.instrument import _MetricsAdvisory from opentelemetry.sdk.metrics._internal.aggregation import ( Aggregation, AggregationTemporality, @@ -39,6 +40,14 @@ def __init__( self._description = ( self._view._description or self._instrument.description ) + # A view that specifies attribute keys takes precedence over the + # instrument's attributes advisory. The advisory only applies when no + # view configures this aspect of the stream + self._attribute_keys = self._view._attribute_keys + if self._attribute_keys is None: + advisory = getattr(self._instrument, "_advisory", None) + if isinstance(advisory, _MetricsAdvisory): + self._attribute_keys = advisory.attributes if not isinstance(self._view._aggregation, DefaultAggregation): self._aggregation = self._view._aggregation._create_aggregation( self._instrument, @@ -87,11 +96,11 @@ def conflicts(self, other: "_ViewInstrumentMatch") -> bool: def consume_measurement( self, measurement: Measurement, should_sample_exemplar: bool = True ) -> None: - if self._view._attribute_keys is not None: + if self._attribute_keys is not None: attributes = {} for key, value in (measurement.attributes or {}).items(): - if key in self._view._attribute_keys: + if key in self._attribute_keys: attributes[key] = value elif measurement.attributes is not None: attributes = dict(measurement.attributes) diff --git a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/aggregation.py b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/aggregation.py index a5c5d5571dd..17cb6ab81a9 100644 --- a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/aggregation.py +++ b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/aggregation.py @@ -1409,13 +1409,10 @@ def _create_aggregation( if self._boundaries is not None: boundaries = self._boundaries else: - # guard for usage with instruments without advisory + # guard for usage with instruments without advisory, and for + # instruments whose advisory is not histogram shaped advisory = getattr(instrument, "_advisory", None) - boundaries = ( - advisory.explicit_bucket_boundaries - if advisory is not None - else None - ) + boundaries = getattr(advisory, "explicit_bucket_boundaries", None) return _ExplicitBucketHistogramAggregation( attributes, diff --git a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/instrument.py b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/instrument.py index 548aa8ee1f3..ac2875978cd 100644 --- a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/instrument.py +++ b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/instrument.py @@ -29,7 +29,7 @@ from opentelemetry.metrics import _Gauge as APIGauge from opentelemetry.metrics._internal.instrument import ( CallbackOptions, - _MetricsHistogramAdvisory, + _MetricsAdvisory, ) from opentelemetry.sdk.metrics._internal.measurement import Measurement @@ -68,6 +68,7 @@ def __init__( unit: str = "", description: str = "", *, + _advisory: _MetricsAdvisory | None = None, _meter_config: _ProxyMeterConfig | None = None, ): # pylint: disable=no-member @@ -90,6 +91,7 @@ def __init__( self.description = description self.instrumentation_scope = instrumentation_scope self._measurement_consumer = measurement_consumer + self._advisory = _advisory or _MetricsAdvisory() self._meter_config = _meter_config super().__init__(name, unit=unit, description=description) @@ -107,6 +109,7 @@ def __init__( unit: str = "", description: str = "", *, + _advisory: _MetricsAdvisory | None = None, _meter_config: _ProxyMeterConfig | None = None, ): # pylint: disable=no-member @@ -129,6 +132,7 @@ def __init__( self.description = description self.instrumentation_scope = instrumentation_scope self._measurement_consumer = measurement_consumer + self._advisory = _advisory or _MetricsAdvisory() self._meter_config = _meter_config super().__init__(name, callbacks, unit=unit, description=description) @@ -278,29 +282,6 @@ def __new__(cls, *args, **kwargs): class Histogram(_Synchronous, APIHistogram): - def __init__( - self, - name: str, - instrumentation_scope: InstrumentationScope, - measurement_consumer: MeasurementConsumer, - unit: str = "", - description: str = "", - explicit_bucket_boundaries_advisory: Sequence[float] | None = None, - *, - _meter_config: _ProxyMeterConfig | None = None, - ): - super().__init__( - name, - unit=unit, - description=description, - instrumentation_scope=instrumentation_scope, - measurement_consumer=measurement_consumer, - _meter_config=_meter_config, - ) - self._advisory = _MetricsHistogramAdvisory( - explicit_bucket_boundaries=explicit_bucket_boundaries_advisory - ) - def __new__(cls, *args, **kwargs): if cls is Histogram: raise TypeError("Histogram must be instantiated via a meter.") diff --git a/opentelemetry-sdk/tests/metrics/integration_test/test_attributes_advisory.py b/opentelemetry-sdk/tests/metrics/integration_test/test_attributes_advisory.py new file mode 100644 index 00000000000..baf0e563b2e --- /dev/null +++ b/opentelemetry-sdk/tests/metrics/integration_test/test_attributes_advisory.py @@ -0,0 +1,184 @@ +# Copyright The OpenTelemetry Authors +# SPDX-License-Identifier: Apache-2.0 + +from logging import WARNING +from unittest import TestCase + +from opentelemetry.metrics import Observation +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.metrics.export import InMemoryMetricReader +from opentelemetry.sdk.metrics.view import View + + +class TestAttributesAdvisory(TestCase): + _SYNC_INSTRUMENTS = { + "counter": ("create_counter", "add"), + "up_down_counter": ("create_up_down_counter", "add"), + "histogram": ("create_histogram", "record"), + "gauge": ("create_gauge", "set"), + } + + _ASYNC_INSTRUMENTS = { + "observable_counter": "create_observable_counter", + "observable_up_down_counter": "create_observable_up_down_counter", + "observable_gauge": "create_observable_gauge", + } + + _ALL_INSTRUMENTS = (*_SYNC_INSTRUMENTS, *_ASYNC_INSTRUMENTS) + + def _collect_attributes(self, reader): + metrics = reader.get_metrics_data() + self.assertEqual(len(metrics.resource_metrics), 1) + self.assertEqual(len(metrics.resource_metrics[0].scope_metrics), 1) + self.assertEqual( + len(metrics.resource_metrics[0].scope_metrics[0].metrics), 1 + ) + metric = metrics.resource_metrics[0].scope_metrics[0].metrics[0] + return [ + dict(data_point.attributes) + for data_point in metric.data.data_points + ] + + def _record(self, meter, instrument, attributes, **create_kwargs): + if instrument in self._SYNC_INSTRUMENTS: + create_method, record_method = self._SYNC_INSTRUMENTS[instrument] + synchronous = getattr(meter, create_method)( + "testinstrument", **create_kwargs + ) + getattr(synchronous, record_method)(1, attributes) + else: + create_method = self._ASYNC_INSTRUMENTS[instrument] + getattr(meter, create_method)( + "testinstrument", + callbacks=[ + lambda options, values=attributes: [ + Observation(1, values) + ] + ], + **create_kwargs, + ) + + def test_advisory(self): + for instrument in self._ALL_INSTRUMENTS: + with self.subTest(instrument=instrument): + reader = InMemoryMetricReader() + meter = MeterProvider(metric_readers=[reader]).get_meter("m") + self._record( + meter, + instrument, + {"label": "value", "dropped": "value"}, + _attributes_advisory=["label"], + ) + + self.assertEqual( + self._collect_attributes(reader), [{"label": "value"}] + ) + + def test_no_advisory(self): + for instrument in self._ALL_INSTRUMENTS: + with self.subTest(instrument=instrument): + reader = InMemoryMetricReader() + meter = MeterProvider(metric_readers=[reader]).get_meter("m") + self._record( + meter, + instrument, + {"label": "value", "kept": "value"}, + ) + + self.assertEqual( + self._collect_attributes(reader), + [{"label": "value", "kept": "value"}], + ) + + def test_empty_advisory(self): + for instrument in self._ALL_INSTRUMENTS: + with self.subTest(instrument=instrument): + reader = InMemoryMetricReader() + meter = MeterProvider(metric_readers=[reader]).get_meter("m") + self._record( + meter, + instrument, + {"label": "value"}, + _attributes_advisory=[], + ) + + self.assertEqual(self._collect_attributes(reader), [{}]) + + def test_view_without_attribute_keys(self): + for instrument in self._ALL_INSTRUMENTS: + with self.subTest(instrument=instrument): + reader = InMemoryMetricReader() + meter = MeterProvider( + metric_readers=[reader], + views=[ + View(instrument_name="testinstrument", name="renamed") + ], + ).get_meter("m") + self._record( + meter, + instrument, + {"label": "value", "dropped": "value"}, + _attributes_advisory=["label"], + ) + + self.assertEqual( + self._collect_attributes(reader), [{"label": "value"}] + ) + + def test_view_overrides_advisory(self): + for instrument in self._ALL_INSTRUMENTS: + with self.subTest(instrument=instrument): + reader = InMemoryMetricReader() + meter = MeterProvider( + metric_readers=[reader], + views=[ + View( + instrument_name="testinstrument", + attribute_keys={"other"}, + ) + ], + ).get_meter("m") + self._record( + meter, + instrument, + {"label": "value", "other": "value"}, + _attributes_advisory=["label"], + ) + + self.assertEqual( + self._collect_attributes(reader), [{"other": "value"}] + ) + + def test_invalid_advisory_ignored(self): + for instrument in self._ALL_INSTRUMENTS: + for invalid_advisory in ( + ["label", 1], + [None], + "label", + 123, + ): + with self.subTest( + instrument=instrument, invalid_advisory=invalid_advisory + ): + reader = InMemoryMetricReader() + meter = MeterProvider( + metric_readers=[reader] + ).get_meter("m") + + # the warning is emitted when the instrument is created + with self.assertLogs(level=WARNING) as log: + self._record( + meter, + instrument, + {"label": "value", "kept": "value"}, + _attributes_advisory=invalid_advisory, + ) + self.assertIn( + "_attributes_advisory must be a sequence of strings", + log.records[0].message, + ) + + self.assertEqual( + self._collect_attributes(reader), + [{"label": "value", "kept": "value"}], + ) diff --git a/opentelemetry-sdk/tests/metrics/test_aggregation.py b/opentelemetry-sdk/tests/metrics/test_aggregation.py index 12a42b31771..3b817993b04 100644 --- a/opentelemetry-sdk/tests/metrics/test_aggregation.py +++ b/opentelemetry-sdk/tests/metrics/test_aggregation.py @@ -9,6 +9,7 @@ from unittest.mock import Mock from opentelemetry.context import Context +from opentelemetry.metrics._internal.instrument import _MetricsAdvisory from opentelemetry.sdk.metrics._internal.aggregation import ( _ExplicitBucketHistogramAggregation, _LastValueAggregation, @@ -696,7 +697,9 @@ def test_histogram_with_advisory(self): "name", Mock(), Mock(), - explicit_bucket_boundaries_advisory=boundaries, + _advisory=_MetricsAdvisory( + explicit_bucket_boundaries=boundaries + ), ), Mock(), _default_reservoir_factory, diff --git a/opentelemetry-sdk/tests/metrics/test_view_instrument_match.py b/opentelemetry-sdk/tests/metrics/test_view_instrument_match.py index 73b0e6e24a5..193fff4b9cf 100644 --- a/opentelemetry-sdk/tests/metrics/test_view_instrument_match.py +++ b/opentelemetry-sdk/tests/metrics/test_view_instrument_match.py @@ -10,6 +10,7 @@ from unittest.mock import MagicMock, Mock, patch from opentelemetry.context import Context +from opentelemetry.metrics._internal.instrument import _MetricsAdvisory from opentelemetry.sdk.metrics._internal._view_instrument_match import ( _ViewInstrumentMatch, ) @@ -24,7 +25,15 @@ ExemplarReservoirBuilder, SimpleFixedSizeExemplarReservoir, ) -from opentelemetry.sdk.metrics._internal.instrument import _Counter, _Histogram +from opentelemetry.sdk.metrics._internal.instrument import ( + _Counter, + _Gauge, + _Histogram, + _ObservableCounter, + _ObservableGauge, + _ObservableUpDownCounter, + _UpDownCounter, +) from opentelemetry.sdk.metrics._internal.measurement import Measurement from opentelemetry.sdk.metrics._internal.sdk_configuration import ( SdkConfiguration, @@ -211,6 +220,98 @@ def test_consume_measurement(self): _DropAggregation, ) + def test_attributes_advisory(self): + for instrument_class in ( + _Counter, + _UpDownCounter, + _Histogram, + _Gauge, + _ObservableCounter, + _ObservableUpDownCounter, + _ObservableGauge, + ): + for view_attribute_keys, advisory_attributes, expected_keys in ( + (None, {"a"}, {("a", 1)}), + (None, None, {("a", 1), ("b", 2)}), + (None, (), ()), + ({"b"}, {"a"}, {("b", 2)}), + (set(), {"a"}, set()), + ): + with self.subTest( + instrument_class=instrument_class, + view_attribute_keys=view_attribute_keys, + advisory_attributes=advisory_attributes, + ): + instrument = instrument_class( + "instrument", + self.mock_instrumentation_scope, + Mock(), + _advisory=_MetricsAdvisory( + attributes=None + if advisory_attributes is None + else frozenset(advisory_attributes) + ), + ) + view_instrument_match = _ViewInstrumentMatch( + view=View( + instrument_name="instrument", + name="name", + aggregation=Mock(), + attribute_keys=view_attribute_keys, + ), + instrument=instrument, + instrument_class_aggregation=MagicMock( + **{ + "__getitem__.return_value": DefaultAggregation() + } + ), + ) + + view_instrument_match.consume_measurement( + Measurement( + value=0, + time_unix_nano=time_ns(), + instrument=instrument, + context=Context(), + attributes={"a": 1, "b": 2}, + ) + ) + + self.assertEqual( + list(view_instrument_match._attributes_aggregation), + [frozenset(expected_keys)], + ) + + def test_attributes_advisory_missing(self): + instrument = Mock(name="instrument") + instrument.instrumentation_scope = self.mock_instrumentation_scope + view_instrument_match = _ViewInstrumentMatch( + view=View( + instrument_name="instrument", + name="name", + aggregation=Mock(), + ), + instrument=instrument, + instrument_class_aggregation=MagicMock( + **{"__getitem__.return_value": DefaultAggregation()} + ), + ) + + view_instrument_match.consume_measurement( + Measurement( + value=0, + time_unix_nano=time_ns(), + instrument=instrument, + context=Context(), + attributes={"a": 1, "b": 2}, + ) + ) + + self.assertEqual( + list(view_instrument_match._attributes_aggregation), + [frozenset([("a", 1), ("b", 2)])], + ) + def test_collect(self): instrument1 = _Counter( "instrument1", From e6b4ad27c064793a944005d4fe01d4c0d33116e9 Mon Sep 17 00:00:00 2001 From: Lukas Hering Date: Sun, 19 Jul 2026 21:28:00 -0400 Subject: [PATCH 2/4] add changelog fragment --- .changelog/5444.added | 1 + 1 file changed, 1 insertion(+) create mode 100644 .changelog/5444.added diff --git a/.changelog/5444.added b/.changelog/5444.added new file mode 100644 index 00000000000..fb73aaf737b --- /dev/null +++ b/.changelog/5444.added @@ -0,0 +1 @@ +`opentelemetry-sdk`: add support for advisory attributes on metric instruments From 4b62a31f050a412e90ca37c7000af2acbdb9ce43 Mon Sep 17 00:00:00 2001 From: Lukas Hering Date: Sun, 19 Jul 2026 21:32:40 -0400 Subject: [PATCH 3/4] fix formatting --- .../integration_test/test_attributes_advisory.py | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/opentelemetry-sdk/tests/metrics/integration_test/test_attributes_advisory.py b/opentelemetry-sdk/tests/metrics/integration_test/test_attributes_advisory.py index baf0e563b2e..8cf2da237b8 100644 --- a/opentelemetry-sdk/tests/metrics/integration_test/test_attributes_advisory.py +++ b/opentelemetry-sdk/tests/metrics/integration_test/test_attributes_advisory.py @@ -51,9 +51,7 @@ def _record(self, meter, instrument, attributes, **create_kwargs): getattr(meter, create_method)( "testinstrument", callbacks=[ - lambda options, values=attributes: [ - Observation(1, values) - ] + lambda options, values=attributes: [Observation(1, values)] ], **create_kwargs, ) @@ -161,9 +159,9 @@ def test_invalid_advisory_ignored(self): instrument=instrument, invalid_advisory=invalid_advisory ): reader = InMemoryMetricReader() - meter = MeterProvider( - metric_readers=[reader] - ).get_meter("m") + meter = MeterProvider(metric_readers=[reader]).get_meter( + "m" + ) # the warning is emitted when the instrument is created with self.assertLogs(level=WARNING) as log: From 5be3e3925104f8fc8fd44ea95de2655aa3e69221 Mon Sep 17 00:00:00 2001 From: Lukas Hering Date: Sun, 19 Jul 2026 21:51:55 -0400 Subject: [PATCH 4/4] update '_validate_explicit_bucket_boundaries_advisory' to return a tuple --- .../metrics/_internal/__init__.py | 21 ++++---- .../metrics/_internal/instrument.py | 47 +++++++++++++++-- .../sdk/metrics/_internal/__init__.py | 50 +------------------ .../tests/metrics/test_metrics.py | 2 +- 4 files changed, 57 insertions(+), 63 deletions(-) diff --git a/opentelemetry-api/src/opentelemetry/metrics/_internal/__init__.py b/opentelemetry-api/src/opentelemetry/metrics/_internal/__init__.py index e60a002976e..4f9acf8d3e9 100644 --- a/opentelemetry-api/src/opentelemetry/metrics/_internal/__init__.py +++ b/opentelemetry-api/src/opentelemetry/metrics/_internal/__init__.py @@ -56,7 +56,6 @@ ObservableUpDownCounter, UpDownCounter, _MetricsAdvisory, - _normalize_attributes_advisory, _ProxyCounter, _ProxyGauge, _ProxyHistogram, @@ -64,6 +63,8 @@ _ProxyObservableGauge, _ProxyObservableUpDownCounter, _ProxyUpDownCounter, + _validate_attributes_advisory, + _validate_explicit_bucket_boundaries_advisory, ) from opentelemetry.util._once import Once from opentelemetry.util._providers import _load_provider @@ -739,7 +740,7 @@ def create_counter( unit, description, _MetricsAdvisory( - attributes=_normalize_attributes_advisory(_attributes_advisory) + attributes=_validate_attributes_advisory(_attributes_advisory) ), ) if status.conflict: @@ -773,7 +774,7 @@ def create_gauge( unit, description, _MetricsAdvisory( - attributes=_normalize_attributes_advisory(_attributes_advisory) + attributes=_validate_attributes_advisory(_attributes_advisory) ), ) if status.conflict: @@ -806,7 +807,7 @@ def create_up_down_counter( unit, description, _MetricsAdvisory( - attributes=_normalize_attributes_advisory(_attributes_advisory) + attributes=_validate_attributes_advisory(_attributes_advisory) ), ) if status.conflict: @@ -840,7 +841,7 @@ def create_observable_counter( unit, description, _MetricsAdvisory( - attributes=_normalize_attributes_advisory(_attributes_advisory) + attributes=_validate_attributes_advisory(_attributes_advisory) ), ) if status.conflict: @@ -875,10 +876,10 @@ def create_histogram( unit, description, _MetricsAdvisory( - attributes=_normalize_attributes_advisory( - _attributes_advisory + attributes=_validate_attributes_advisory(_attributes_advisory), + explicit_bucket_boundaries=_validate_explicit_bucket_boundaries_advisory( + explicit_bucket_boundaries_advisory ), - explicit_bucket_boundaries=explicit_bucket_boundaries_advisory, ), ) if status.conflict: @@ -913,7 +914,7 @@ def create_observable_gauge( unit, description, _MetricsAdvisory( - attributes=_normalize_attributes_advisory(_attributes_advisory) + attributes=_validate_attributes_advisory(_attributes_advisory) ), ) if status.conflict: @@ -948,7 +949,7 @@ def create_observable_up_down_counter( unit, description, _MetricsAdvisory( - attributes=_normalize_attributes_advisory(_attributes_advisory) + attributes=_validate_attributes_advisory(_attributes_advisory) ), ) if status.conflict: diff --git a/opentelemetry-api/src/opentelemetry/metrics/_internal/instrument.py b/opentelemetry-api/src/opentelemetry/metrics/_internal/instrument.py index ec163c59785..683702c30c6 100644 --- a/opentelemetry-api/src/opentelemetry/metrics/_internal/instrument.py +++ b/opentelemetry-api/src/opentelemetry/metrics/_internal/instrument.py @@ -31,15 +31,54 @@ @dataclass(frozen=True) class _MetricsAdvisory: attributes: frozenset[str] | None = None - explicit_bucket_boundaries: Sequence[float] | None = None + explicit_bucket_boundaries: tuple[float, ...] | None = None -def _normalize_attributes_advisory( +def _validate_attributes_advisory( attributes_advisory: Sequence[str] | None, ) -> frozenset[str] | None: - return ( - None if attributes_advisory is None else frozenset(attributes_advisory) + """Returns the advisory as a frozenset of attribute keys, otherwise logs a + warning and returns `None` so the advisory is ignored.""" + + if attributes_advisory is None: + return None + + if isinstance(attributes_advisory, Sequence) and not isinstance( + attributes_advisory, str + ): + try: + if all(isinstance(entry, str) for entry in attributes_advisory): + return frozenset(attributes_advisory) + except (KeyError, TypeError): + pass + + _logger.warning("_attributes_advisory must be a sequence of strings") + return None + + +def _validate_explicit_bucket_boundaries_advisory( + explicit_bucket_boundaries_advisory: Sequence[float] | None, +) -> tuple[float, ...] | None: + """Returns the advisory as a tuple of numbers, otherwise logs a warning and + returns `None` so the advisory is ignored.""" + + if explicit_bucket_boundaries_advisory is None: + return None + + if isinstance(explicit_bucket_boundaries_advisory, Sequence): + try: + if all( + isinstance(entry, (float, int)) + for entry in explicit_bucket_boundaries_advisory + ): + return tuple(explicit_bucket_boundaries_advisory) + except (KeyError, TypeError): + pass + + _logger.warning( + "explicit_bucket_boundaries_advisory must be a sequence of numbers" ) + return None @dataclass(frozen=True) diff --git a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/__init__.py b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/__init__.py index c03e709fadb..d7d343bcf5d 100644 --- a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/__init__.py +++ b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/__init__.py @@ -27,7 +27,8 @@ from opentelemetry.metrics import _Gauge as APIGauge from opentelemetry.metrics._internal.instrument import ( _MetricsAdvisory, - _normalize_attributes_advisory, + _validate_attributes_advisory, + _validate_explicit_bucket_boundaries_advisory, ) from opentelemetry.sdk.environment_variables import ( OTEL_METRICS_EXEMPLAR_FILTER, @@ -72,53 +73,6 @@ _logger = getLogger(__name__) -def _validate_explicit_bucket_boundaries_advisory( - explicit_bucket_boundaries_advisory: Sequence[float] | None, -) -> Sequence[float] | None: - """Returns the advisory if it is a sequence of numbers, otherwise logs a - warning and returns `None` so the advisory is ignored.""" - - if explicit_bucket_boundaries_advisory is None: - return None - - if isinstance(explicit_bucket_boundaries_advisory, Sequence): - try: - if all( - isinstance(entry, (float, int)) - for entry in explicit_bucket_boundaries_advisory - ): - return explicit_bucket_boundaries_advisory - except (KeyError, TypeError): - pass - - _logger.warning( - "explicit_bucket_boundaries_advisory must be a sequence of numbers" - ) - return None - - -def _validate_attributes_advisory( - attributes_advisory: Sequence[str] | None, -) -> frozenset[str] | None: - """Returns the advisory as a frozenset of attribute keys, otherwise logs a - warning and returns `None` so the advisory is ignored.""" - - if attributes_advisory is None: - return None - - if isinstance(attributes_advisory, Sequence) and not isinstance( - attributes_advisory, str - ): - try: - if all(isinstance(entry, str) for entry in attributes_advisory): - return _normalize_attributes_advisory(attributes_advisory) - except (KeyError, TypeError): - pass - - _logger.warning("_attributes_advisory must be a sequence of strings") - return None - - @dataclass class _MeterConfig: is_enabled: bool = True diff --git a/opentelemetry-sdk/tests/metrics/test_metrics.py b/opentelemetry-sdk/tests/metrics/test_metrics.py index 176c6a2aef3..e34dac3975e 100644 --- a/opentelemetry-sdk/tests/metrics/test_metrics.py +++ b/opentelemetry-sdk/tests/metrics/test_metrics.py @@ -774,7 +774,7 @@ def test_create_histogram_with_advisory(self): self.assertEqual(histogram.name, "name") self.assertEqual( histogram._advisory.explicit_bucket_boundaries, - [0.0, 1.0, 2], + (0.0, 1.0, 2), ) def test_create_histogram_advisory_validation(self):