diff --git a/.changelog/5419.added b/.changelog/5419.added new file mode 100644 index 00000000000..1cead81a8a9 --- /dev/null +++ b/.changelog/5419.added @@ -0,0 +1 @@ +`opentelemetry-sdk`: add support for the `MetricProducer` interface diff --git a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/export/__init__.py b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/export/__init__.py index aee356f78d2..6ee77234424 100644 --- a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/export/__init__.py +++ b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/export/__init__.py @@ -6,7 +6,7 @@ import os import weakref from abc import ABC, abstractmethod -from collections.abc import Callable, Iterable +from collections.abc import Callable, Iterable, Sequence from enum import Enum from logging import getLogger from os import environ, linesep @@ -55,7 +55,10 @@ _ObservableUpDownCounter, _UpDownCounter, ) -from opentelemetry.sdk.metrics._internal.point import MetricsData +from opentelemetry.sdk.metrics._internal.point import ( + MetricsData, + ScopeMetrics, +) from opentelemetry.semconv._incubating.attributes.otel_attributes import ( OtelComponentTypeValues, ) @@ -178,6 +181,34 @@ def force_flush(self, timeout_millis: float = 10_000) -> bool: return True +class MetricProducer(ABC): + """Interface bridging third-party metric sources into a :class:`MetricReader`. + + An implementation is registered on a ``MetricReader`` (via its + ``metric_producers`` argument) and its metrics are collected alongside the + SDK's own internal state whenever the reader collects. See the OpenTelemetry + metrics SDK specification on + `MetricProducer `__. + """ + + @abstractmethod + def produce( + self, timeout_millis: float = 10_000 + ) -> Iterable[ScopeMetrics]: + """Returns the producer's metrics as an iterable of ``ScopeMetrics``. + + A producer SHOULD emit a single ``InstrumentationScope`` that identifies + the producer itself. + + Args: + timeout_millis: Amount of time in milliseconds before the produce + operation should time out. + + Returns: + An iterable of :class:`~opentelemetry.sdk.metrics.export.ScopeMetrics`. + """ + + class MetricReader(ABC): # pylint: disable=too-many-branches,broad-exception-raised """ @@ -207,6 +238,9 @@ class MetricReader(ABC): default aggregations. The aggregation defined here will be overridden by an aggregation defined by a view that is not `DefaultAggregation`. + metric_producers: A sequence of `MetricProducer` instances that bridge + third-party metric sources. Their metrics are collected alongside + the SDK's own metrics whenever this reader collects. .. document protected _receive_metrics which is a intended to be overridden by subclass .. automethod:: _receive_metrics @@ -221,8 +255,12 @@ def __init__( ] | None = None, *, + metric_producers: Sequence[MetricProducer] = (), otel_component_type: OtelComponentTypeValues | None = None, ) -> None: + self._metric_producers: tuple[MetricProducer, ...] = tuple( + metric_producers + ) self._collect: ( Callable[ [ @@ -432,10 +470,13 @@ def __init__( type, opentelemetry.sdk.metrics.view.Aggregation ] | None = None, + *, + metric_producers: Sequence[MetricProducer] = (), ) -> None: super().__init__( preferred_temporality=preferred_temporality, preferred_aggregation=preferred_aggregation, + metric_producers=metric_producers, ) self._lock = RLock() self._metrics_data: MetricsData | None = None @@ -478,11 +519,14 @@ def __init__( exporter: MetricExporter, export_interval_millis: float | None = None, export_timeout_millis: float | None = None, + *, + metric_producers: Sequence[MetricProducer] = (), ) -> None: # PeriodicExportingMetricReader defers to exporter for configuration super().__init__( preferred_temporality=exporter._preferred_temporality, preferred_aggregation=exporter._preferred_aggregation, + metric_producers=metric_producers, otel_component_type=OtelComponentTypeValues.PERIODIC_METRIC_READER, ) diff --git a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/measurement_consumer.py b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/measurement_consumer.py index b17475bc551..8e55d43ec20 100644 --- a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/measurement_consumer.py +++ b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/measurement_consumer.py @@ -5,6 +5,7 @@ from abc import ABC, abstractmethod from collections.abc import Iterable, Mapping +from logging import getLogger from threading import Lock from time import time_ns @@ -17,7 +18,14 @@ from opentelemetry.sdk.metrics._internal.metric_reader_storage import ( MetricReaderStorage, ) -from opentelemetry.sdk.metrics._internal.point import MetricsData +from opentelemetry.sdk.metrics._internal.point import ( + MetricsData, + ResourceMetrics, + ScopeMetrics, +) +from opentelemetry.sdk.resources import Resource + +_logger = getLogger(__name__) class MeasurementConsumer(ABC): @@ -133,8 +141,83 @@ def collect( result = self._reader_storages[metric_reader].collect() + if producer_scope_metrics := self._collect_from_producers( + metric_reader, deadline_ns + ): + result = self._merge_producer_metrics( + result, producer_scope_metrics, self._sdk_config.resource + ) + return result + @staticmethod + def _collect_from_producers( + metric_reader: "opentelemetry.sdk.metrics.export.MetricReader", + deadline_ns: float, + ) -> list[ScopeMetrics]: + """Collect ScopeMetrics from the reader's MetricProducers. + + Called while holding ``self._lock`` so ``produce()`` calls are + serialized. A producer that fails or times out is isolated so + metrics collected by the SDK are not dropped. + """ + producer_scope_metrics: list[ScopeMetrics] = [] + # pylint: disable-next=protected-access + for producer in metric_reader._metric_producers: + remaining_millis = (deadline_ns - time_ns()) / 1e6 + if remaining_millis <= 0: + _logger.warning( + "Timed out collecting from metric producers, " + "skipping remaining producers." + ) + break + + try: + scopes = list( + producer.produce(timeout_millis=remaining_millis) + ) + # pylint: disable-next=broad-except + except Exception: + _logger.exception( + "Metric producer %s failed to produce metrics, skipping.", + producer, + ) + continue + + producer_scope_metrics.extend(scopes) + + return producer_scope_metrics + + @staticmethod + def _merge_producer_metrics( + result: MetricsData | None, + producer_scope_metrics: list[ScopeMetrics], + resource: Resource, + ) -> MetricsData: + if result is not None and result.resource_metrics: + sdk_resource_metrics = result.resource_metrics[0] + merged = ResourceMetrics( + resource=sdk_resource_metrics.resource, + scope_metrics=[ + *sdk_resource_metrics.scope_metrics, + *producer_scope_metrics, + ], + schema_url=sdk_resource_metrics.schema_url, + ) + return MetricsData( + resource_metrics=[merged, *result.resource_metrics[1:]] + ) + + return MetricsData( + resource_metrics=[ + ResourceMetrics( + resource=resource, + scope_metrics=producer_scope_metrics, + schema_url=resource.schema_url, + ) + ] + ) + def add_metric_reader( self, metric_reader: "opentelemetry.sdk.metrics.MetricReader" ) -> None: diff --git a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/export/__init__.py b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/export/__init__.py index 56034c0e35c..7128cbd7155 100644 --- a/opentelemetry-sdk/src/opentelemetry/sdk/metrics/export/__init__.py +++ b/opentelemetry-sdk/src/opentelemetry/sdk/metrics/export/__init__.py @@ -10,6 +10,7 @@ InMemoryMetricReader, MetricExporter, MetricExportResult, + MetricProducer, MetricReader, PeriodicExportingMetricReader, ) @@ -39,6 +40,7 @@ "InMemoryMetricReader", "MetricExporter", "MetricExportResult", + "MetricProducer", "MetricReader", "PeriodicExportingMetricReader", "DataPointT", diff --git a/opentelemetry-sdk/tests/metrics/test_measurement_consumer.py b/opentelemetry-sdk/tests/metrics/test_measurement_consumer.py index d09368c75dc..4da890590b3 100644 --- a/opentelemetry-sdk/tests/metrics/test_measurement_consumer.py +++ b/opentelemetry-sdk/tests/metrics/test_measurement_consumer.py @@ -30,7 +30,7 @@ def test_parent(self, _): def test_creates_metric_reader_storages(self, MockMetricReaderStorage): """It should create one MetricReaderStorage per metric reader passed in the SdkConfiguration""" - reader_mocks = [Mock() for _ in range(5)] + reader_mocks = [Mock(_metric_producers=[]) for _ in range(5)] SynchronousMeasurementConsumer( SdkConfiguration( exemplar_filter=Mock(), @@ -44,7 +44,7 @@ def test_creates_metric_reader_storages(self, MockMetricReaderStorage): def test_measurements_passed_to_each_reader_storage( self, MockMetricReaderStorage ): - reader_mocks = [Mock() for _ in range(5)] + reader_mocks = [Mock(_metric_producers=[]) for _ in range(5)] reader_storage_mocks = [Mock() for _ in range(5)] MockMetricReaderStorage.side_effect = reader_storage_mocks @@ -66,7 +66,7 @@ def test_measurements_passed_to_each_reader_storage( def test_collect_passed_to_reader_stage(self, MockMetricReaderStorage): """Its collect() method should defer to the underlying MetricReaderStorage""" - reader_mocks = [Mock() for _ in range(5)] + reader_mocks = [Mock(_metric_producers=[]) for _ in range(5)] reader_storage_mocks = [Mock() for _ in range(5)] MockMetricReaderStorage.side_effect = reader_storage_mocks @@ -86,7 +86,7 @@ def test_collect_passed_to_reader_stage(self, MockMetricReaderStorage): def test_collect_calls_async_instruments(self, MockMetricReaderStorage): """Its collect() method should invoke async instruments and pass measurements to the corresponding metric reader storage""" - reader_mock = Mock() + reader_mock = Mock(_metric_producers=[]) reader_storage_mock = Mock() MockMetricReaderStorage.return_value = reader_storage_mock consumer = SynchronousMeasurementConsumer( @@ -117,7 +117,7 @@ def test_collect_calls_async_instruments(self, MockMetricReaderStorage): self.assertFalse(reader_storage_mock.consume_measurement.call_args[1]) def test_collect_timeout(self, MockMetricReaderStorage): - reader_mock = Mock() + reader_mock = Mock(_metric_producers=[]) reader_storage_mock = Mock() MockMetricReaderStorage.return_value = reader_storage_mock consumer = SynchronousMeasurementConsumer( @@ -151,7 +151,7 @@ def sleep_1(*args, **kwargs): def test_collect_deadline( self, mock_time_ns, mock_callback_options, MockMetricReaderStorage ): - reader_mock = Mock() + reader_mock = Mock(_metric_producers=[]) reader_storage_mock = Mock() MockMetricReaderStorage.return_value = reader_storage_mock consumer = SynchronousMeasurementConsumer( diff --git a/opentelemetry-sdk/tests/metrics/test_metric_producer.py b/opentelemetry-sdk/tests/metrics/test_metric_producer.py new file mode 100644 index 00000000000..5d74bce533a --- /dev/null +++ b/opentelemetry-sdk/tests/metrics/test_metric_producer.py @@ -0,0 +1,270 @@ +# Copyright The OpenTelemetry Authors +# SPDX-License-Identifier: Apache-2.0 + +# pylint: disable=protected-access + +from __future__ import annotations + +from collections.abc import Iterable +from unittest import TestCase +from unittest.mock import patch + +from opentelemetry.sdk.metrics import Counter, MeterProvider +from opentelemetry.sdk.metrics.export import ( + AggregationTemporality, + InMemoryMetricReader, + Metric, + MetricProducer, + NumberDataPoint, + ScopeMetrics, + Sum, +) +from opentelemetry.sdk.resources import Resource +from opentelemetry.sdk.util.instrumentation import InstrumentationScope + +_TIME_NS = "opentelemetry.sdk.metrics._internal.measurement_consumer.time_ns" + + +def _make_scope_metrics(name: str, value: int = 7) -> ScopeMetrics: + return ScopeMetrics( + scope=InstrumentationScope(name=name), + metrics=[ + Metric( + name="produced.metric", + description="", + unit="", + data=Sum( + data_points=[ + NumberDataPoint( + attributes={}, + start_time_unix_nano=0, + time_unix_nano=1, + value=value, + ) + ], + aggregation_temporality=AggregationTemporality.CUMULATIVE, + is_monotonic=True, + ), + ) + ], + schema_url="", + ) + + +def _get_scope_names(metrics_data) -> set[str]: + names = set() + for resource_metrics in metrics_data.resource_metrics: + for scope_metrics in resource_metrics.scope_metrics: + names.add(scope_metrics.scope.name) + return names + + +def _find_scope(metrics_data, scope_name: str) -> ScopeMetrics | None: + for resource_metrics in metrics_data.resource_metrics: + for scope_metrics in resource_metrics.scope_metrics: + if scope_metrics.scope.name == scope_name: + return scope_metrics + return None + + +def _point_value(scope_metrics: ScopeMetrics, metric_name: str): + for metric in scope_metrics.metrics: + if metric.name == metric_name: + return metric.data.data_points[0].value + raise AssertionError(f"metric {metric_name!r} not found") + + +class _FakeProducer(MetricProducer): + def __init__( + self, scope_name: str = "fake-producer", value: int = 7 + ) -> None: + self.produce_calls = 0 + self._scope_name = scope_name + self._value = value + + def produce( + self, timeout_millis: float = 10_000 + ) -> Iterable[ScopeMetrics]: + self.produce_calls += 1 + return [_make_scope_metrics(self._scope_name, self._value)] + + +class _FailingProducer(MetricProducer): + def produce( + self, timeout_millis: float = 10_000 + ) -> Iterable[ScopeMetrics]: + raise RuntimeError("produce failed") + + +class _ClockAdvancingProducer(MetricProducer): + def __init__(self, clock: list[int], advance_ns: int) -> None: + self._clock = clock + self._advance_ns = advance_ns + self.produce_calls = 0 + + def produce( + self, timeout_millis: float = 10_000 + ) -> Iterable[ScopeMetrics]: + self.produce_calls += 1 + self._clock[0] += self._advance_ns + return [_make_scope_metrics("slow-producer")] + + +class TestMetricProducer(TestCase): + def test_producer_metrics_merged_with_sdk_metrics(self): + resource = Resource.create({"service.name": "test"}) + producer = _FakeProducer(value=42) + reader = InMemoryMetricReader(metric_producers=[producer]) + provider = MeterProvider(metric_readers=[reader], resource=resource) + + meter = provider.get_meter("sdk-scope") + counter: Counter = meter.create_counter("sdk.counter") + counter.add(1) + + metrics_data = reader.get_metrics_data() + + self.assertEqual(producer.produce_calls, 1) + + self.assertEqual(len(metrics_data.resource_metrics), 1) + resource_metrics = metrics_data.resource_metrics[0] + self.assertIs(resource_metrics.resource, resource) + self.assertEqual( + resource_metrics.resource.attributes["service.name"], "test" + ) + self.assertEqual( + _get_scope_names(metrics_data), {"sdk-scope", "fake-producer"} + ) + + self.assertEqual( + _point_value( + _find_scope(metrics_data, "sdk-scope"), "sdk.counter" + ), + 1, + ) + self.assertEqual( + _point_value( + _find_scope(metrics_data, "fake-producer"), "produced.metric" + ), + 42, + ) + + def test_producer_only_no_sdk_metrics(self): + resource = Resource.create({"service.name": "test"}) + producer = _FakeProducer(value=13) + reader = InMemoryMetricReader(metric_producers=[producer]) + MeterProvider(metric_readers=[reader], resource=resource) + + metrics_data = reader.get_metrics_data() + + self.assertEqual(len(metrics_data.resource_metrics), 1) + resource_metrics = metrics_data.resource_metrics[0] + self.assertIs(resource_metrics.resource, resource) + self.assertEqual( + resource_metrics.resource.attributes["service.name"], "test" + ) + self.assertEqual(_get_scope_names(metrics_data), {"fake-producer"}) + self.assertEqual( + _point_value( + _find_scope(metrics_data, "fake-producer"), "produced.metric" + ), + 13, + ) + + def test_multiple_producers(self): + resource = Resource.create({}) + producers = [ + _FakeProducer("p1", value=1), + _FakeProducer("p2", value=2), + ] + reader = InMemoryMetricReader(metric_producers=producers) + MeterProvider(metric_readers=[reader], resource=resource) + + metrics_data = reader.get_metrics_data() + + self.assertEqual(_get_scope_names(metrics_data), {"p1", "p2"}) + self.assertEqual( + _point_value(_find_scope(metrics_data, "p1"), "produced.metric"), 1 + ) + self.assertEqual( + _point_value(_find_scope(metrics_data, "p2"), "produced.metric"), 2 + ) + + def test_no_producers(self): + reader = InMemoryMetricReader() + provider = MeterProvider(metric_readers=[reader]) + meter = provider.get_meter("sdk-scope") + meter.create_counter("sdk.counter").add(1) + + metrics_data = reader.get_metrics_data() + + self.assertEqual(_get_scope_names(metrics_data), {"sdk-scope"}) + + def test_failing_producer_is_isolated(self): + resource = Resource.create({}) + good = _FakeProducer("good", value=99) + reader = InMemoryMetricReader( + metric_producers=[_FailingProducer(), good] + ) + provider = MeterProvider(metric_readers=[reader], resource=resource) + provider.get_meter("sdk-scope").create_counter("sdk.counter").add(1) + + with self.assertLogs(level="WARNING") as cm: + metrics_data = reader.get_metrics_data() + + self.assertEqual(_get_scope_names(metrics_data), {"sdk-scope", "good"}) + self.assertEqual(good.produce_calls, 1) + self.assertEqual( + _point_value( + _find_scope(metrics_data, "sdk-scope"), "sdk.counter" + ), + 1, + ) + self.assertEqual( + _point_value(_find_scope(metrics_data, "good"), "produced.metric"), + 99, + ) + self.assertTrue( + any( + "failed to produce metrics" in message for message in cm.output + ) + ) + + # pylint: disable-next=no-self-use + def test_producer_receives_remaining_timeout_budget(self): + producer = _FakeProducer() + reader = InMemoryMetricReader(metric_producers=[producer]) + MeterProvider(metric_readers=[reader]) + + # Freeze the clock so the producer receives exactly the full budget. + with patch.object( + producer, "produce", wraps=producer.produce + ) as produce_mock: + with patch(_TIME_NS, lambda: 0): + reader.collect(timeout_millis=5_000) + + produce_mock.assert_called_once_with(timeout_millis=5_000) + + def test_timeout_softly_skips_remaining_producers(self): + clock = [0] + slow = _ClockAdvancingProducer(clock, advance_ns=int(20 * 1e6)) + never = _FakeProducer("never") + reader = InMemoryMetricReader(metric_producers=[slow, never]) + MeterProvider(metric_readers=[reader]) + + with patch(_TIME_NS, lambda: clock[0]): + with self.assertLogs(level="WARNING") as cm: + reader.collect(timeout_millis=10) + + metrics_data = reader._metrics_data + self.assertEqual(slow.produce_calls, 1) + self.assertEqual(never.produce_calls, 0) + self.assertIn("slow-producer", _get_scope_names(metrics_data)) + self.assertEqual( + _point_value( + _find_scope(metrics_data, "slow-producer"), "produced.metric" + ), + 7, + ) + self.assertTrue( + any("Timed out collecting" in message for message in cm.output) + )