diff --git a/CHANGELOG.md b/CHANGELOG.md index 4ff18bad04..83b502ad5a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -71,6 +71,15 @@ Increment the: at shutdown [#4637](https://github.com/open-telemetry/opentelemetry-cpp/pull/4637) +* [SDK] Fix spatial re-aggregation for asynchronous instruments. A view which + drops attributes is now applied to asynchronous instruments as well, and the + observations which collapse onto the same attribute set are re-aggregated + (summed up for the additive instruments) instead of the last observation + overwriting the previous ones. + Note that `AsyncMetricStorage`'s constructor now takes the view's + `AttributesProcessor`, matching `SyncMetricStorage`. + [#1724](https://github.com/open-telemetry/opentelemetry-cpp/issues/1724) + ## [1.29.0] 2026-09-13 * [RELEASE] Bump main branch to 1.29.0-dev (#4259) diff --git a/sdk/include/opentelemetry/sdk/metrics/state/async_metric_storage.h b/sdk/include/opentelemetry/sdk/metrics/state/async_metric_storage.h index b1c31cf324..72d14f378e 100644 --- a/sdk/include/opentelemetry/sdk/metrics/state/async_metric_storage.h +++ b/sdk/include/opentelemetry/sdk/metrics/state/async_metric_storage.h @@ -3,12 +3,15 @@ #pragma once +#include #include #include #include +#include #include "opentelemetry/nostd/shared_ptr.h" #include "opentelemetry/sdk/common/attributemap_hash.h" +#include "opentelemetry/sdk/metrics/aggregation/aggregation.h" #include "opentelemetry/sdk/metrics/aggregation/aggregation_config.h" #include "opentelemetry/sdk/metrics/aggregation/default_aggregation.h" @@ -37,6 +40,7 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora public: AsyncMetricStorage(const InstrumentDescriptor &instrument_descriptor, const AggregationType aggregation_type, + std::shared_ptr attributes_processor, #ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW ExemplarFilterType exemplar_filter_type, nostd::shared_ptr &&exemplar_reservoir, @@ -45,24 +49,34 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora : instrument_descriptor_(instrument_descriptor), aggregation_type_{aggregation_type}, aggregation_config_{AggregationConfig::GetOrDefault(aggregation_config)}, - cumulative_hash_map_( + attributes_processor_{std::move(attributes_processor)}, + is_monotonic_sum_{IsMonotonicSum(aggregation_type, instrument_descriptor)}, + last_observed_hash_map_( std::make_unique(aggregation_config_->cardinality_limit_)), delta_hash_map_( std::make_unique(aggregation_config_->cardinality_limit_)), + round_hash_map_( + std::make_unique(aggregation_config_->cardinality_limit_)), #ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW exemplar_filter_type_(exemplar_filter_type), exemplar_reservoir_(std::move(exemplar_reservoir)), #endif temporal_metric_storage_(instrument_descriptor, aggregation_type, aggregation_config) - {} + { + create_default_aggregation_ = [this]() -> std::unique_ptr { + return DefaultAggregation::CreateAggregation(aggregation_type_, instrument_descriptor_); + }; + } + /** + * Converts the absolute values reported by the callbacks into the deltas the temporal storage + * consumes. Runs once per callback registered on the instrument, so what it accumulates stays + * until the next Collect(). + */ template void Record(const std::unordered_map &measurements, opentelemetry::common::SystemTimestamp /* observation_time */) noexcept { - // Async counter always record monotonically increasing values, and the - // exporter/reader can request either for delta or cumulative value. - // So we convert the async counter value to delta before passing it to temporal storage. std::lock_guard guard(hashmap_lock_); #ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW const bool offer_exemplars = @@ -76,25 +90,15 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora exemplar_reservoir_->OfferMeasurement(measurement.second, measurement.first, {}); } #endif - - auto aggr = DefaultAggregation::CreateAggregation(aggregation_type_, instrument_descriptor_); - aggr->Aggregate(measurement.second); - auto prev = cumulative_hash_map_->Get(measurement.first); - if (prev) + if (is_monotonic_sum_) { - auto delta = prev->Diff(*aggr); - // store received value in cumulative map, and the diff in delta map (to pass it to temporal - // storage) - cumulative_hash_map_->Set(measurement.first, std::move(aggr)); - delta_hash_map_->Set(measurement.first, std::move(delta)); + RecordMonotonicSum(measurement.first, measurement.second); } else { - // store received value in cumulative and delta map. - cumulative_hash_map_->Set( - measurement.first, - DefaultAggregation::CloneAggregation(aggregation_type_, instrument_descriptor_, *aggr)); - delta_hash_map_->Set(measurement.first, std::move(aggr)); + round_hash_map_ + ->GetOrSetDefault(FilterAttributes(measurement.first), create_default_aggregation_) + ->Aggregate(measurement.second); } } } @@ -131,6 +135,10 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora std::shared_ptr delta_metrics = nullptr; { std::lock_guard guard(hashmap_lock_); + if (!is_monotonic_sum_) + { + BuildDeltaFromRound(); + } delta_metrics = std::move(delta_hash_map_); delta_hash_map_ = std::make_unique(aggregation_config_->cardinality_limit_); @@ -143,11 +151,140 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora } private: + /** + * Differences a monotonic sum per source series - keyed by the unfiltered attributes - before + * the view merges anything. Diffing the merged group instead would turn a series which stops + * being reported into a spurious decrease. + */ + template + void RecordMonotonicSum(const MetricAttributes &attributes, T value) + { + auto observed = create_default_aggregation_(); + observed->Aggregate(value); + + std::unique_ptr delta; + auto previous = last_observed_hash_map_->Get(attributes); + if (previous) + { + delta = previous->Diff(*observed); + } + else + { + delta = DefaultAggregation::CloneAggregation(aggregation_type_, instrument_descriptor_, + *observed); + } + last_observed_hash_map_->Set(attributes, std::move(observed)); + + AccumulateDelta(FilterAttributes(attributes), *delta); + } + + /** + * Adds one series' delta to its output point. Merging, not overwriting, is what lets the series + * a view collapses together and the observations from different callbacks all contribute. + * Over-the-limit attribute sets resolve to the same otel.metric.overflow entry, which merges + * too. + */ + void AccumulateDelta(const MetricAttributes &attributes, const Aggregation &delta) + { + auto merged = + delta_hash_map_->GetOrSetDefault(attributes, create_default_aggregation_)->Merge(delta); + delta_hash_map_->Set(attributes, std::move(merged)); + } + + /** + * Differences this round's absolute observations and starts a fresh round, for everything but + * monotonic sums. An attribute set missing from this round keeps its baseline and contributes + * no delta; suppressing the stale output is the temporal storage's job. + */ + void BuildDeltaFromRound() noexcept + { + round_hash_map_->GetAllEntries([this](const MetricAttributes &attributes, + Aggregation &aggregation) { + auto observed = DefaultAggregation::CloneAggregation(aggregation_type_, + instrument_descriptor_, aggregation); + auto previous = last_observed_hash_map_->Get(attributes); + if (previous) + { + delta_hash_map_->Set(attributes, previous->Diff(*observed)); + } + else + { + delta_hash_map_->Set(attributes, DefaultAggregation::CloneAggregation( + aggregation_type_, instrument_descriptor_, *observed)); + } + last_observed_hash_map_->Set(attributes, std::move(observed)); + return true; + }); + + round_hash_map_ = std::make_unique(aggregation_config_->cardinality_limit_); + } + + /** + * Returns a copy of the observed attributes with the attributes dropped by the view removed. + */ + MetricAttributes FilterAttributes(const MetricAttributes &attributes) const + { + MetricAttributes filtered(attributes); + if (!attributes_processor_) + { + return filtered; + } + + bool dropped = false; + for (auto iter = filtered.begin(); iter != filtered.end();) + { + if (attributes_processor_->isPresent(iter->first)) + { + ++iter; + } + else + { + iter = filtered.erase(iter); + dropped = true; + } + } + + if (dropped) + { + filtered.UpdateHash(); + } + return filtered; + } + + /** + * Whether this storage aggregates into a monotonic sum, matching what + * DefaultAggregation::CreateAggregation() would create. + */ + static bool IsMonotonicSum(AggregationType aggregation_type, + const InstrumentDescriptor &instrument_descriptor) noexcept + { + bool is_monotonic = true; + if (aggregation_type == AggregationType::kDefault) + { + const AggregationType resolved = + DefaultAggregation::GetDefaultAggregationType(instrument_descriptor.type_, is_monotonic); + return resolved == AggregationType::kSum && is_monotonic; + } + if (aggregation_type != AggregationType::kSum) + { + return false; + } + return instrument_descriptor.type_ != InstrumentType::kUpDownCounter && + instrument_descriptor.type_ != InstrumentType::kObservableUpDownCounter && + instrument_descriptor.type_ != InstrumentType::kHistogram; + } + InstrumentDescriptor instrument_descriptor_; AggregationType aggregation_type_; const AggregationConfig *aggregation_config_; - std::unique_ptr cumulative_hash_map_; + std::shared_ptr attributes_processor_; + bool is_monotonic_sum_; + std::function()> create_default_aggregation_; + + std::unique_ptr last_observed_hash_map_; std::unique_ptr delta_hash_map_; + std::unique_ptr round_hash_map_; + std::mutex hashmap_lock_; #ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW ExemplarFilterType exemplar_filter_type_; diff --git a/sdk/src/metrics/meter.cc b/sdk/src/metrics/meter.cc index 0abd75265e..12280976dd 100644 --- a/sdk/src/metrics/meter.cc +++ b/sdk/src/metrics/meter.cc @@ -615,7 +615,7 @@ std::unique_ptr Meter::RegisterAsyncMetricStorage( { WarnOnDuplicateInstrument(GetInstrumentationScope(), storage_registry_, view_instr_desc); async_storage = std::shared_ptr(new AsyncMetricStorage( - view_instr_desc, view.GetAggregationType(), + view_instr_desc, view.GetAggregationType(), view.GetAttributesProcessor(), #ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW exemplar_filter_type, GetExemplarReservoir(view.GetAggregationType(), view.GetAggregationConfig(), diff --git a/sdk/test/metrics/async_metric_storage_test.cc b/sdk/test/metrics/async_metric_storage_test.cc index 54eec782ec..a11a75a104 100644 --- a/sdk/test/metrics/async_metric_storage_test.cc +++ b/sdk/test/metrics/async_metric_storage_test.cc @@ -3,7 +3,9 @@ #include #include +#include #include +#include #include #include #include @@ -16,7 +18,9 @@ #include "opentelemetry/nostd/span.h" #include "opentelemetry/nostd/utility.h" #include "opentelemetry/nostd/variant.h" +#include "opentelemetry/sdk/common/attribute_utils.h" #include "opentelemetry/sdk/instrumentationscope/instrumentation_scope.h" +#include "opentelemetry/sdk/metrics/aggregation/aggregation_config.h" #include "opentelemetry/sdk/metrics/data/metric_data.h" #include "opentelemetry/sdk/metrics/data/point_data.h" #include "opentelemetry/sdk/metrics/export/metric_producer.h" @@ -68,7 +72,7 @@ TEST_P(AsyncWritableMetricStorageTestFixture, TestAggregation) collectors.push_back(collector); opentelemetry::sdk::metrics::AsyncMetricStorage storage( - instr_desc, AggregationType::kSum, + instr_desc, AggregationType::kSum, std::make_shared(), #ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), #endif @@ -163,7 +167,7 @@ TEST_P(WritableMetricStorageTestUpDownFixture, TestAggregation) collectors.push_back(collector); opentelemetry::sdk::metrics::AsyncMetricStorage storage( - instr_desc, AggregationType::kDefault, + instr_desc, AggregationType::kDefault, std::make_shared(), #ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), #endif @@ -259,7 +263,7 @@ TEST_P(WritableMetricStorageTestObservableGaugeFixture, TestAggregation) collectors.push_back(collector); opentelemetry::sdk::metrics::AsyncMetricStorage storage( - instr_desc, AggregationType::kLastValue, + instr_desc, AggregationType::kLastValue, std::make_shared(), #ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), #endif @@ -322,4 +326,371 @@ INSTANTIATE_TEST_SUITE_P(WritableMetricStorageTestObservableGaugeFixtureLong, ::testing::Values(AggregationTemporality::kCumulative, AggregationTemporality::kDelta)); +class WritableMetricStorageTestFilteredAttributesFixture + : public ::testing::TestWithParam +{}; + +// Observations the view collapses onto one attribute set are summed, not overwritten. +TEST_P(WritableMetricStorageTestFilteredAttributesFixture, TestAggregation) +{ + AggregationTemporality temporality = GetParam(); + + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", InstrumentType::kObservableCounter, + InstrumentValueType::kLong}; + + auto sdk_start_ts = std::chrono::system_clock::now(); + auto collection_ts = std::chrono::system_clock::now() + std::chrono::seconds(5); + + std::shared_ptr collector(new MockCollectorHandle(temporality)); + std::vector> collectors; + collectors.push_back(collector); + + FilterAttributeMap allowed_attributes; + allowed_attributes["RequestType"] = true; + std::shared_ptr attributes_processor{ + new FilteringAttributesProcessor(allowed_attributes)}; + + opentelemetry::sdk::metrics::AsyncMetricStorage storage( + instr_desc, AggregationType::kSum, attributes_processor, +#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW + ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), +#endif + nullptr); + + int64_t get_count_v1 = 20; + int64_t get_count_v2 = 10; + int64_t put_count_v1 = 5; + + std::unordered_map measurements1 = { + {{{"RequestType", "GET"}, {"version", "1"}}, get_count_v1}, + {{{"RequestType", "GET"}, {"version", "2"}}, get_count_v2}, + {{{"RequestType", "PUT"}, {"version", "1"}}, put_count_v1}}; + storage.RecordLong(measurements1, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + size_t collected_points = 0; + storage.Collect( + collector.get(), collectors, sdk_start_ts, collection_ts, [&](const MetricData &metric_data) { + for (const auto &data_attr : metric_data.point_data_attr_) + { + const auto &data = opentelemetry::nostd::get(data_attr.point_data); + EXPECT_EQ(data_attr.attributes.end(), data_attr.attributes.find("version")); + ++collected_points; + if (opentelemetry::nostd::get( + data_attr.attributes.find("RequestType")->second) == "GET") + { + EXPECT_EQ(opentelemetry::nostd::get(data.value_), get_count_v1 + get_count_v2); + } + else + { + EXPECT_EQ(opentelemetry::nostd::get(data.value_), put_count_v1); + } + } + return true; + }); + EXPECT_EQ(collected_points, 2); + + int64_t get_count_v1_2 = 50; + int64_t get_count_v2_2 = 30; + int64_t put_count_v1_2 = 8; + + std::unordered_map measurements2 = { + {{{"RequestType", "GET"}, {"version", "1"}}, get_count_v1_2}, + {{{"RequestType", "GET"}, {"version", "2"}}, get_count_v2_2}, + {{{"RequestType", "PUT"}, {"version", "1"}}, put_count_v1_2}}; + storage.RecordLong(measurements2, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + storage.Collect( + collector.get(), collectors, sdk_start_ts, collection_ts, [&](const MetricData &metric_data) { + for (const auto &data_attr : metric_data.point_data_attr_) + { + const auto &data = opentelemetry::nostd::get(data_attr.point_data); + EXPECT_EQ(data_attr.attributes.end(), data_attr.attributes.find("version")); + if (opentelemetry::nostd::get( + data_attr.attributes.find("RequestType")->second) == "GET") + { + if (temporality == AggregationTemporality::kCumulative) + { + EXPECT_EQ(opentelemetry::nostd::get(data.value_), + get_count_v1_2 + get_count_v2_2); + } + else + { + EXPECT_EQ(opentelemetry::nostd::get(data.value_), + (get_count_v1_2 + get_count_v2_2) - (get_count_v1 + get_count_v2)); + } + } + else + { + if (temporality == AggregationTemporality::kCumulative) + { + EXPECT_EQ(opentelemetry::nostd::get(data.value_), put_count_v1_2); + } + else + { + EXPECT_EQ(opentelemetry::nostd::get(data.value_), + put_count_v1_2 - put_count_v1); + } + } + } + return true; + }); +} + +INSTANTIATE_TEST_SUITE_P(WritableMetricStorageTestFilteredAttributesLong, + WritableMetricStorageTestFilteredAttributesFixture, + ::testing::Values(AggregationTemporality::kCumulative, + AggregationTemporality::kDelta)); + +// Every dimension dropped: all observations collapse onto the empty attribute set. +TEST(WritableMetricStorageTestFilteredAttributes, TestUpDownCounterAllAttributesDropped) +{ + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", + InstrumentType::kObservableUpDownCounter, + InstrumentValueType::kDouble}; + + auto sdk_start_ts = std::chrono::system_clock::now(); + auto collection_ts = std::chrono::system_clock::now() + std::chrono::seconds(5); + + std::shared_ptr collector( + new MockCollectorHandle(AggregationTemporality::kCumulative)); + std::vector> collectors; + collectors.push_back(collector); + + std::shared_ptr attributes_processor{ + new FilteringAttributesProcessor(FilterAttributeMap{})}; + + opentelemetry::sdk::metrics::AsyncMetricStorage storage( + instr_desc, AggregationType::kSum, attributes_processor, +#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW + ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), +#endif + nullptr); + + std::unordered_map measurements = { + {{{"version", "1"}}, 1.0}, {{{"version", "2"}}, 2.0}, {{{"version", "3"}}, -4.0}}; + storage.RecordDouble(measurements, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + size_t collected_points = 0; + storage.Collect( + collector.get(), collectors, sdk_start_ts, collection_ts, [&](const MetricData &metric_data) { + for (const auto &data_attr : metric_data.point_data_attr_) + { + const auto &data = opentelemetry::nostd::get(data_attr.point_data); + ++collected_points; + EXPECT_EQ(0, data_attr.attributes.size()); + EXPECT_DOUBLE_EQ(opentelemetry::nostd::get(data.value_), -1.0); + } + return true; + }); + EXPECT_EQ(collected_points, 1); +} + +namespace +{ +// Collects the emitted sum points, so assertions don't depend on iteration order. +template +std::vector> CollectSumPoints( + opentelemetry::sdk::metrics::AsyncMetricStorage &storage, + const std::shared_ptr &collector, + std::vector> &collectors) +{ + std::vector> points; + auto sdk_start_ts = std::chrono::system_clock::now(); + auto collection_ts = std::chrono::system_clock::now() + std::chrono::seconds(5); + storage.Collect( + collector.get(), collectors, sdk_start_ts, collection_ts, [&](const MetricData &metric_data) { + for (const auto &data_attr : metric_data.point_data_attr_) + { + const auto &data = opentelemetry::nostd::get(data_attr.point_data); + points.emplace_back(data_attr.attributes, opentelemetry::nostd::get(data.value_)); + } + return true; + }); + return points; +} + +std::shared_ptr AllowOnly(const std::string &key) +{ + FilterAttributeMap allowed; + allowed[key] = true; + return std::shared_ptr(new FilteringAttributesProcessor(allowed)); +} + +std::shared_ptr DropEverything() +{ + return std::shared_ptr( + new FilteringAttributesProcessor(FilterAttributeMap{})); +} +} // namespace + +// Each callback produces its own Record() call; observations the view collapses onto one point +// must combine across them rather than overwrite. +TEST(WritableMetricStorageMultiCallback, MonotonicSumCombinesAcrossCallbacks) +{ + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", InstrumentType::kObservableCounter, + InstrumentValueType::kLong}; + + std::shared_ptr collector( + new MockCollectorHandle(AggregationTemporality::kCumulative)); + std::vector> collectors{collector}; + + opentelemetry::sdk::metrics::AsyncMetricStorage storage( + instr_desc, AggregationType::kSum, DropEverything(), +#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW + ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), +#endif + nullptr); + + std::unordered_map from_first_callback = { + {{{"version", "v1"}}, 20}}; + std::unordered_map from_second_callback = { + {{{"version", "v2"}}, 10}}; + storage.RecordLong(from_first_callback, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + storage.RecordLong(from_second_callback, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + auto points = CollectSumPoints(storage, collector, collectors); + ASSERT_EQ(1, points.size()); + EXPECT_EQ(0, points[0].first.size()); + EXPECT_EQ(30, points[0].second); +} + +// Differenced per source series, so a series which stops being reported doesn't drag the merged +// group negative: 20 + 10 followed by 25 is a rise of 5, not a drop of 5. +TEST(WritableMetricStorageDisappearingSeries, MonotonicSumKeepsEarlierContribution) +{ + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", InstrumentType::kObservableCounter, + InstrumentValueType::kLong}; + + std::shared_ptr collector( + new MockCollectorHandle(AggregationTemporality::kCumulative)); + std::vector> collectors{collector}; + + opentelemetry::sdk::metrics::AsyncMetricStorage storage( + instr_desc, AggregationType::kSum, DropEverything(), +#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW + ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), +#endif + nullptr); + + std::unordered_map first_round = { + {{{"version", "v1"}}, 20}, {{{"version", "v2"}}, 10}}; + storage.RecordLong(first_round, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + auto points = CollectSumPoints(storage, collector, collectors); + ASSERT_EQ(1, points.size()); + EXPECT_EQ(30, points[0].second); + + std::unordered_map second_round = { + {{{"version", "v1"}}, 25}}; + storage.RecordLong(second_round, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + points = CollectSumPoints(storage, collector, collectors); + ASSERT_EQ(1, points.size()); + EXPECT_EQ(35, points[0].second); +} + +// Non-monotonic sums total the round first and difference once, so the reported total follows a +// series disappearing: 100 + 50 becomes 100 when the second process exits. +TEST(WritableMetricStorageDisappearingSeries, UpDownCounterFollowsTheTotal) +{ + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", + InstrumentType::kObservableUpDownCounter, + InstrumentValueType::kLong}; + + std::shared_ptr collector( + new MockCollectorHandle(AggregationTemporality::kCumulative)); + std::vector> collectors{collector}; + + opentelemetry::sdk::metrics::AsyncMetricStorage storage( + instr_desc, AggregationType::kSum, AllowOnly("host"), +#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW + ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), +#endif + nullptr); + + std::unordered_map first_process = { + {{{"host", "h1"}, {"process", "p1"}}, 100}}; + std::unordered_map second_process = { + {{{"host", "h1"}, {"process", "p2"}}, 50}}; + storage.RecordLong(first_process, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + storage.RecordLong(second_process, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + auto points = CollectSumPoints(storage, collector, collectors); + ASSERT_EQ(1, points.size()); + EXPECT_EQ(1, points[0].first.size()); + EXPECT_EQ(150, points[0].second); + + storage.RecordLong(first_process, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + points = CollectSumPoints(storage, collector, collectors); + ASSERT_EQ(1, points.size()); + EXPECT_EQ(100, points[0].second); +} + +// Attribute sets beyond the cardinality limit combine into `otel.metric.overflow`; no +// contribution may be dropped, including from a later callback in the same round. +TEST(WritableMetricStorageCardinalityLimit, OverflowCombinesContributions) +{ + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", InstrumentType::kObservableCounter, + InstrumentValueType::kLong}; + + std::shared_ptr collector( + new MockCollectorHandle(AggregationTemporality::kCumulative)); + std::vector> collectors{collector}; + + constexpr size_t kCardinalityLimit = 2; + AggregationConfig aggregation_config{kCardinalityLimit}; + + opentelemetry::sdk::metrics::AsyncMetricStorage storage( + instr_desc, AggregationType::kSum, AllowOnly("id"), +#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW + ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), +#endif + &aggregation_config); + + std::unordered_map from_first_callback = { + {{{"id", "a"}, {"version", "v1"}}, 1}, + {{{"id", "b"}, {"version", "v1"}}, 2}, + {{{"id", "c"}, {"version", "v1"}}, 3}}; + std::unordered_map from_second_callback = { + {{{"id", "d"}, {"version", "v1"}}, 4}, {{{"id", "e"}, {"version", "v1"}}, 5}}; + storage.RecordLong(from_first_callback, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + storage.RecordLong(from_second_callback, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + auto points = CollectSumPoints(storage, collector, collectors); + + ASSERT_EQ(kCardinalityLimit + 1, points.size()); + + int64_t total = 0; + bool has_overflow = false; + int64_t overflow_value = 0; + for (const auto &point : points) + { + total += point.second; + const auto overflow_it = point.first.find(kAttributesLimitOverflowKey); + if (overflow_it != point.first.end()) + { + has_overflow = true; + EXPECT_EQ(true, opentelemetry::nostd::get(overflow_it->second)); + overflow_value = point.second; + } + } + + EXPECT_TRUE(has_overflow); + EXPECT_EQ(1 + 2 + 3 + 4 + 5, total); + EXPECT_GT(overflow_value, 0); +} + } // namespace diff --git a/sdk/test/metrics/sum_aggregation_test.cc b/sdk/test/metrics/sum_aggregation_test.cc index 03119906d1..6fd1ba458c 100644 --- a/sdk/test/metrics/sum_aggregation_test.cc +++ b/sdk/test/metrics/sum_aggregation_test.cc @@ -13,7 +13,9 @@ #include "opentelemetry/common/macros.h" #include "opentelemetry/context/context.h" +#include "opentelemetry/metrics/async_instruments.h" #include "opentelemetry/metrics/meter.h" +#include "opentelemetry/metrics/observer_result.h" #include "opentelemetry/metrics/sync_instruments.h" #include "opentelemetry/nostd/function_ref.h" #include "opentelemetry/nostd/shared_ptr.h" @@ -411,6 +413,76 @@ TEST(CounterToSumFilterAttributesWithCardinalityLimit, Double) } } +namespace +{ +// Observes two different "version" dimensions for the same "attr1" value. +void ObservableCounterCallback(opentelemetry::metrics::ObserverResult observer, void * /* state */) +{ + auto observer_double = opentelemetry::nostd::get< + opentelemetry::nostd::shared_ptr>>(observer); + observer_double->Observe(1.0, + std::map{{"attr1", "val1"}, {"version", "1"}}); + observer_double->Observe(2.0, + std::map{{"attr1", "val1"}, {"version", "2"}}); +} +} // namespace + +// A view dropping "version" collapses both observations onto one point: additive, so 1 + 2, +// not the last observed value. +TEST(AsyncCounterToSumFilterAttributes, Double) +{ + MeterProvider mp; + auto m = mp.GetMeter("meter1", "version1", "schema1"); + std::string instrument_unit = "ms"; + std::string instrument_name = "observable_counter1"; + std::string instrument_desc = "observable counter metrics"; + + opentelemetry::sdk::metrics::FilterAttributeMap allowedattr; + allowedattr["attr1"] = true; + std::unique_ptr attrproc{ + new opentelemetry::sdk::metrics::FilteringAttributesProcessor(allowedattr)}; + + std::shared_ptr dummy_aggregation_config{ + new opentelemetry::sdk::metrics::AggregationConfig}; + std::unique_ptr exporter(new MockMetricExporter()); + std::shared_ptr reader{new MockMetricReader(std::move(exporter))}; + mp.AddMetricReader(reader); + + std::unique_ptr view{new View("view1", "view1_description", AggregationType::kSum, + dummy_aggregation_config, std::move(attrproc))}; + std::unique_ptr instrument_selector{ + new InstrumentSelector(InstrumentType::kObservableCounter, instrument_name, instrument_unit)}; + std::unique_ptr meter_selector{new MeterSelector("meter1", "version1", "schema1")}; + mp.AddView(std::move(instrument_selector), std::move(meter_selector), std::move(view)); + + auto c = m->CreateDoubleObservableCounter(instrument_name, instrument_desc, instrument_unit); + c->AddCallback(ObservableCounterCallback, nullptr); + + size_t collected_points = 0; + reader->Collect([&](ResourceMetrics &rm) { + for (const ScopeMetrics &smd : rm.scope_metric_data_) + { + for (const MetricData &md : smd.metric_data_) + { + EXPECT_EQ(1, md.point_data_attr_.size()); + for (const PointDataAttributes &dp : md.point_data_attr_) + { + ++collected_points; + EXPECT_EQ(3.0, opentelemetry::nostd::get( + opentelemetry::nostd::get(dp.point_data).value_)); + EXPECT_EQ(1, dp.attributes.size()); + EXPECT_NE(dp.attributes.end(), dp.attributes.find("attr1")); + EXPECT_EQ(dp.attributes.end(), dp.attributes.find("version")); + } + } + } + return true; + }); + EXPECT_EQ(1, collected_points); + + c->RemoveCallback(ObservableCounterCallback, nullptr); +} + namespace {