From fabd66021eebefaa67b24a9a5125020ecdd42a46 Mon Sep 17 00:00:00 2001 From: Patrick Summerer Date: Wed, 26 Aug 2026 10:17:09 +0200 Subject: [PATCH 1/5] [METRICS] Fix stale async attribute sets in cumulative exports (#4108) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Async instruments (ObservableCounter, ObservableGauge, ObservableUpDownCounter) under cumulative temporality were emitting attribute sets indefinitely after the callback stopped reporting them, violating the OTel spec requirement: "The implementation SHOULD NOT produce aggregated metric data for a previously-observed attribute set which is not observed during a successful callback." Root cause: `TemporalMetricStorage::buildMetrics()` unconditionally carried every entry from `last_reported_metrics_` into the output even when it was absent from the current delta. Fix: - Add `is_async_` flag (default false) to `TemporalMetricStorage`. The cumulative merge now skips entries not present in the current delta for async instruments, while sync instruments retain the existing carry-forward behaviour. - Pass `is_async = true` when constructing `TemporalMetricStorage` from `AsyncMetricStorage`. - Do NOT prune `cumulative_hash_map_` in `AsyncMetricStorage::Collect()` so that the absolute-value baseline is preserved across absent cycles. This ensures correct delta computation (new - last_seen, not the full new value) when an attribute set reappears after a gap — consistent with the approach taken by opentelemetry-dotnet#6883. Tests added in async_metric_storage_test.cc: - StaleAttributeSetDroppedInCumulativeExport: verifies that an attribute set absent from the callback is not emitted in subsequent cumulative exports. - AttributeReappearanceAfterGapDeltaTemporality: verifies that an attribute set reappearing after an absent cycle emits only the increment since last observed (delta = 1, not 11), confirming the baseline is correctly preserved. Fixes #4108 Co-authored-by: pranitaurlam <227409059+pranitaurlam@users.noreply.github.com> --- .../sdk/metrics/state/async_metric_storage.h | 8 +- .../metrics/state/temporal_metric_storage.h | 4 +- .../metrics/state/temporal_metric_storage.cc | 11 +- sdk/test/metrics/async_metric_storage_test.cc | 165 ++++++++++++++++++ 4 files changed, 183 insertions(+), 5 deletions(-) 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 674863428b..dc3aa8f85f 100644 --- a/sdk/include/opentelemetry/sdk/metrics/state/async_metric_storage.h +++ b/sdk/include/opentelemetry/sdk/metrics/state/async_metric_storage.h @@ -53,7 +53,7 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora exemplar_filter_type_(exemplar_filter_type), exemplar_reservoir_(std::move(exemplar_reservoir)), #endif - temporal_metric_storage_(instrument_descriptor, aggregation_type, aggregation_config) + temporal_metric_storage_(instrument_descriptor, aggregation_type, aggregation_config, true) {} template @@ -134,6 +134,12 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora delta_metrics = std::move(delta_hash_map_); delta_hash_map_ = std::make_unique(aggregation_config_->cardinality_limit_); + // cumulative_hash_map_ is intentionally NOT pruned here. + // It preserves the last-seen absolute value for every attribute set so that + // delta computation in Record() remains correct if an attribute set reappears + // after being absent for one or more collection cycles. + // Stale entries are suppressed at export time by the is_async_ guard in + // TemporalMetricStorage::buildMetrics() instead. } auto status = diff --git a/sdk/include/opentelemetry/sdk/metrics/state/temporal_metric_storage.h b/sdk/include/opentelemetry/sdk/metrics/state/temporal_metric_storage.h index d86093c376..c0f8c77bd5 100644 --- a/sdk/include/opentelemetry/sdk/metrics/state/temporal_metric_storage.h +++ b/sdk/include/opentelemetry/sdk/metrics/state/temporal_metric_storage.h @@ -34,7 +34,8 @@ class TemporalMetricStorage public: TemporalMetricStorage(InstrumentDescriptor instrument_descriptor, AggregationType aggregation_type, - const AggregationConfig *aggregation_config); + const AggregationConfig *aggregation_config, + bool is_async = false); bool buildMetrics(CollectorHandle *collector, nostd::span> collectors, @@ -64,6 +65,7 @@ class TemporalMetricStorage // See https://github.com/open-telemetry/opentelemetry-specification (logs/metrics // SDK specs) and issue #4062. const opentelemetry::common::SystemTimestamp instrument_creation_ts_; + bool is_async_ = false; }; } // namespace metrics } // namespace sdk diff --git a/sdk/src/metrics/state/temporal_metric_storage.cc b/sdk/src/metrics/state/temporal_metric_storage.cc index 85a0a4b3c0..60e522421f 100644 --- a/sdk/src/metrics/state/temporal_metric_storage.cc +++ b/sdk/src/metrics/state/temporal_metric_storage.cc @@ -32,11 +32,13 @@ namespace metrics TemporalMetricStorage::TemporalMetricStorage(InstrumentDescriptor instrument_descriptor, AggregationType aggregation_type, - const AggregationConfig *aggregation_config) + const AggregationConfig *aggregation_config, + bool is_async) : instrument_descriptor_(std::move(instrument_descriptor)), aggregation_type_(aggregation_type), aggregation_config_(aggregation_config), - instrument_creation_ts_(std::chrono::system_clock::now()) + instrument_creation_ts_(std::chrono::system_clock::now()), + is_async_(is_async) {} bool TemporalMetricStorage::buildMetrics(CollectorHandle *collector, @@ -159,8 +161,11 @@ bool TemporalMetricStorage::buildMetrics(CollectorHandle *collector, { merged_metrics->Set(attributes, agg->Merge(aggregation)); } - else + else if (!is_async_) { + // For sync instruments, carry forward attribute sets not observed this cycle. + // For async instruments, drop them per the spec: the SDK SHOULD NOT produce + // aggregated metric data for attribute sets not observed in the current callback. auto def_agg = DefaultAggregation::CreateAggregation( aggregation_type_, instrument_descriptor_, aggregation_config_); merged_metrics->Set(attributes, def_agg->Merge(aggregation)); diff --git a/sdk/test/metrics/async_metric_storage_test.cc b/sdk/test/metrics/async_metric_storage_test.cc index 54eec782ec..6e7ce2275a 100644 --- a/sdk/test/metrics/async_metric_storage_test.cc +++ b/sdk/test/metrics/async_metric_storage_test.cc @@ -322,4 +322,169 @@ INSTANTIATE_TEST_SUITE_P(WritableMetricStorageTestObservableGaugeFixtureLong, ::testing::Values(AggregationTemporality::kCumulative, AggregationTemporality::kDelta)); +// Regression test for https://github.com/open-telemetry/opentelemetry-cpp/issues/4108 +// +// Async instruments under cumulative temporality must NOT carry forward attribute sets that were +// not reported by the callback in the current collection cycle. +TEST(AsyncMetricStorageRegressionTest, StaleAttributeSetDroppedInCumulativeExport) +{ + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", InstrumentType::kObservableCounter, + InstrumentValueType::kLong}; + + auto sdk_start_ts = std::chrono::system_clock::now(); + // Some computation here + auto collection_ts = sdk_start_ts + std::chrono::seconds(5); + + std::shared_ptr collector( + new MockCollectorHandle(AggregationTemporality::kCumulative)); + std::vector> collectors; + collectors.push_back(collector); + + opentelemetry::sdk::metrics::AsyncMetricStorage storage( + instr_desc, AggregationType::kSum, +#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW + ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), +#endif + nullptr); + + // Collection 1: both GET and PUT reported. + std::unordered_map measurements1 = { + {{{"RequestType", "GET"}}, 10}, {{{"RequestType", "PUT"}}, 5}}; + storage.RecordLong(measurements1, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + int get_count = 0; + int put_count = 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 &key = opentelemetry::nostd::get( + data_attr.attributes.find("RequestType")->second); + if (key == "GET") + get_count++; + else if (key == "PUT") + put_count++; + } + return true; + }); + EXPECT_EQ(get_count, 1); + EXPECT_EQ(put_count, 1); + + // Collection 2: only GET reported – PUT is dropped by callback. + std::unordered_map measurements2 = { + {{{"RequestType", "GET"}}, 20}}; + storage.RecordLong(measurements2, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + get_count = 0; + put_count = 0; + int64_t get_value = 0; + storage.Collect(collector.get(), collectors, sdk_start_ts, + collection_ts + std::chrono::seconds(5), [&](const MetricData &metric_data) { + for (const auto &data_attr : metric_data.point_data_attr_) + { + const auto &key = opentelemetry::nostd::get( + data_attr.attributes.find("RequestType")->second); + if (key == "GET") + { + get_count++; + get_value = opentelemetry::nostd::get( + opentelemetry::nostd::get(data_attr.point_data).value_); + } + else if (key == "PUT") + { + put_count++; + } + } + return true; + }); + + // PUT must not appear – it was absent from the callback this cycle. + EXPECT_EQ(put_count, 0) << "Stale PUT attribute set must be dropped from cumulative export"; + EXPECT_EQ(get_count, 1); + EXPECT_EQ(get_value, 20); +} + +// Regression test for https://github.com/open-telemetry/opentelemetry-cpp/issues/4108 +// +// Under delta temporality an attribute set that disappears for one collection cycle and then +// reappears must emit only the increment since the last observed value, not the full new absolute +// value. The cumulative baseline (cumulative_hash_map_) is preserved across absent cycles so that +// the delta computation in Record() remains correct. +TEST(AsyncMetricStorageRegressionTest, AttributeReappearanceAfterGapDeltaTemporality) +{ + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", InstrumentType::kObservableCounter, + InstrumentValueType::kLong}; + + auto sdk_start_ts = std::chrono::system_clock::now(); + // Some computation here + auto collection_ts = sdk_start_ts + std::chrono::seconds(5); + + std::shared_ptr collector( + new MockCollectorHandle(AggregationTemporality::kDelta)); + std::vector> collectors; + collectors.push_back(collector); + + opentelemetry::sdk::metrics::AsyncMetricStorage storage( + instr_desc, AggregationType::kSum, +#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW + ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), +#endif + nullptr); + + // Collection 1: A=10 → delta should be 10. + std::unordered_map measurements1 = { + {{{"attr", "A"}}, 10}}; + storage.RecordLong(measurements1, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + int64_t delta_value = -1; + storage.Collect(collector.get(), collectors, sdk_start_ts, collection_ts, + [&](const MetricData &metric_data) { + for (const auto &data_attr : metric_data.point_data_attr_) + { + delta_value = opentelemetry::nostd::get( + opentelemetry::nostd::get(data_attr.point_data).value_); + } + return true; + }); + EXPECT_EQ(delta_value, 10); + + // Collection 2: attribute A absent – nothing recorded, nothing emitted. + std::unordered_map measurements2; + storage.RecordLong(measurements2, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + int attr_count = 0; + storage.Collect(collector.get(), collectors, sdk_start_ts, + collection_ts + std::chrono::seconds(5), [&](const MetricData &metric_data) { + attr_count += static_cast(metric_data.point_data_attr_.size()); + return true; + }); + EXPECT_EQ(attr_count, 0) << "No data points expected when attribute set is absent"; + + // Collection 3: A reappears with absolute value 11. + // The cumulative baseline (10) was preserved across the absent cycle, so + // delta = 11 - 10 = 1 — the correct increment since the attribute was last seen. + std::unordered_map measurements3 = { + {{{"attr", "A"}}, 11}}; + storage.RecordLong(measurements3, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + delta_value = -1; + storage.Collect(collector.get(), collectors, sdk_start_ts, + collection_ts + std::chrono::seconds(10), [&](const MetricData &metric_data) { + for (const auto &data_attr : metric_data.point_data_attr_) + { + delta_value = opentelemetry::nostd::get( + opentelemetry::nostd::get(data_attr.point_data).value_); + } + return true; + }); + // Baseline was preserved: delta = new_value - last_seen_value = 11 - 10 = 1. + EXPECT_EQ(delta_value, 1) + << "After a gap, reappearing attribute must emit the increment since last seen"; +} + } // namespace From a3f8710f8d39e5c79faf519ab053427973c27401 Mon Sep 17 00:00:00 2001 From: Patrick Summerer Date: Thu, 10 Sep 2026 10:06:17 +0200 Subject: [PATCH 2/5] [METRICS] Refine fix of stale async attribute sets in cumulative exports (#4108) --- .../metrics/state/temporal_metric_storage.cc | 38 +++-- sdk/test/metrics/async_metric_storage_test.cc | 154 ++++++++++++++++++ 2 files changed, 180 insertions(+), 12 deletions(-) diff --git a/sdk/src/metrics/state/temporal_metric_storage.cc b/sdk/src/metrics/state/temporal_metric_storage.cc index 7bd87afe2c..0e32a9b056 100644 --- a/sdk/src/metrics/state/temporal_metric_storage.cc +++ b/sdk/src/metrics/state/temporal_metric_storage.cc @@ -160,11 +160,11 @@ bool TemporalMetricStorage::buildMetrics(CollectorHandle *collector, { merged_metrics->Set(attributes, agg->Merge(aggregation)); } - else if (!is_async_) + else { - // For sync instruments, carry forward attribute sets not observed this cycle. - // For async instruments, drop them per the spec: the SDK SHOULD NOT produce - // aggregated metric data for attribute sets not observed in the current callback. + // Always carry forward cumulative state so the baseline survives gaps. + // For async instruments, stale suppression is applied only when building + // the exported MetricData below. auto def_agg = DefaultAggregation::CreateAggregation( aggregation_type_, instrument_descriptor_, aggregation_config_); merged_metrics->Set(attributes, def_agg->Merge(aggregation)); @@ -197,14 +197,28 @@ bool TemporalMetricStorage::buildMetrics(CollectorHandle *collector, metric_data.aggregation_temporality = aggregation_temporarily; metric_data.start_ts = last_collection_ts; metric_data.end_ts = collection_ts; - result_to_export->GetAllEntries( - [&metric_data](const MetricAttributes &attributes, Aggregation &aggregation) { - PointDataAttributes point_data_attr; - point_data_attr.point_data = aggregation.ToPoint(); - point_data_attr.attributes = attributes; - metric_data.point_data_attr_.emplace_back(std::move(point_data_attr)); - return true; - }); + result_to_export->GetAllEntries([&metric_data, &delta_metrics, this]( + const MetricAttributes &attributes, + Aggregation &aggregation) { + if (is_async_ && metric_data.aggregation_temporality == AggregationTemporality::kCumulative && + !delta_metrics->Has(attributes)) + { + // Async cumulative exports must omit attribute sets that were not observed + // in the current callback cycle, while keeping the internal cumulative state. + return true; + } + + PointDataAttributes point_data_attr; + point_data_attr.point_data = aggregation.ToPoint(); + point_data_attr.attributes = attributes; + metric_data.point_data_attr_.emplace_back(std::move(point_data_attr)); + return true; + }); + + if (metric_data.point_data_attr_.empty()) + { + return true; + } return callback(metric_data); } diff --git a/sdk/test/metrics/async_metric_storage_test.cc b/sdk/test/metrics/async_metric_storage_test.cc index 6e7ce2275a..931427174d 100644 --- a/sdk/test/metrics/async_metric_storage_test.cc +++ b/sdk/test/metrics/async_metric_storage_test.cc @@ -487,4 +487,158 @@ TEST(AsyncMetricStorageRegressionTest, AttributeReappearanceAfterGapDeltaTempora << "After a gap, reappearing attribute must emit the increment since last seen"; } +// Regression test for the cumulative reappearance scenario. +// +// For async cumulative exports, stale attribute sets must be suppressed when absent, +// but the internal cumulative baseline must still be preserved. If A=10, then absent, +// then A=30, the reappearance must export 30 (full cumulative), not 20 (increment only). +TEST(AsyncMetricStorageRegressionTest, AttributeReappearanceAfterGapCumulativeTemporality) +{ + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", InstrumentType::kObservableCounter, + InstrumentValueType::kLong}; + + auto sdk_start_ts = std::chrono::system_clock::now(); + auto collection_ts = sdk_start_ts + std::chrono::seconds(5); + + std::shared_ptr collector( + new MockCollectorHandle(AggregationTemporality::kCumulative)); + std::vector> collectors; + collectors.push_back(collector); + opentelemetry::sdk::metrics::AsyncMetricStorage storage( + instr_desc, AggregationType::kSum, +#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW + ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), +#endif + nullptr); + + // Collection 1: A=10 -> cumulative export should be 10. + std::unordered_map measurements1 = { + {{{"attr", "A"}}, 10}}; + storage.RecordLong(measurements1, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + int64_t cumulative_value = -1; + storage.Collect(collector.get(), collectors, sdk_start_ts, collection_ts, + [&](const MetricData &metric_data) { + for (const auto &data_attr : metric_data.point_data_attr_) + { + cumulative_value = opentelemetry::nostd::get( + opentelemetry::nostd::get(data_attr.point_data).value_); + } + return true; + }); + EXPECT_EQ(cumulative_value, 10); + + // Collection 2: A absent -> no export for A. + std::unordered_map measurements2; + storage.RecordLong(measurements2, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + int attr_count = 0; + storage.Collect(collector.get(), collectors, sdk_start_ts, + collection_ts + std::chrono::seconds(5), [&](const MetricData &metric_data) { + attr_count += static_cast(metric_data.point_data_attr_.size()); + return true; + }); + EXPECT_EQ(attr_count, 0) << "No data points expected when attribute set is absent"; + + // Collection 3: A=30 -> cumulative export must be 30 (not 20). + std::unordered_map measurements3 = { + {{{"attr", "A"}}, 30}}; + storage.RecordLong(measurements3, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + cumulative_value = -1; + storage.Collect(collector.get(), collectors, sdk_start_ts, + collection_ts + std::chrono::seconds(10), [&](const MetricData &metric_data) { + for (const auto &data_attr : metric_data.point_data_attr_) + { + cumulative_value = opentelemetry::nostd::get( + opentelemetry::nostd::get(data_attr.point_data).value_); + } + return true; + }); + EXPECT_EQ(cumulative_value, 30) + << "Cumulative reappearance must preserve baseline and export full cumulative value"; +} + +// Regression test for the multi-collector delta slow path. +// +// With more than one collector, delta temporality no longer short-circuits on the single-collector +// fast path it goes through the slow path instead. This exercises +// the delta branch that sets start_ts to the previous collection's timestamp, and verifies both +// collectors independently receive the increment since their own last collection. +TEST(AsyncMetricStorageRegressionTest, MultiCollectorDeltaTemporality) +{ + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", InstrumentType::kObservableCounter, + InstrumentValueType::kLong}; + + auto sdk_start_ts = std::chrono::system_clock::now(); + auto collection_ts1 = sdk_start_ts + std::chrono::seconds(5); + auto collection_ts2 = sdk_start_ts + std::chrono::seconds(10); + + std::shared_ptr collector1( + new MockCollectorHandle(AggregationTemporality::kDelta)); + std::shared_ptr collector2( + new MockCollectorHandle(AggregationTemporality::kDelta)); + std::vector> collectors; + collectors.push_back(collector1); + collectors.push_back(collector2); + + opentelemetry::sdk::metrics::AsyncMetricStorage storage( + instr_desc, AggregationType::kSum, +#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW + ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), +#endif + nullptr); + + auto collect_value = [&](CollectorHandle *collector, + opentelemetry::common::SystemTimestamp collection_ts, int64_t &value_out, + opentelemetry::common::SystemTimestamp &start_ts_out) { + storage.Collect(collector, collectors, sdk_start_ts, collection_ts, + [&](const MetricData &metric_data) { + start_ts_out = metric_data.start_ts; + for (const auto &data_attr : metric_data.point_data_attr_) + { + value_out = opentelemetry::nostd::get( + opentelemetry::nostd::get(data_attr.point_data).value_); + } + return true; + }); + }; + + // Cycle 1: A=10 observed once, both collectors drain the same delta -> each sees 10. + std::unordered_map measurements1 = { + {{{"attr", "A"}}, 10}}; + storage.RecordLong(measurements1, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + int64_t c1_value = -1; + int64_t c2_value = -1; + opentelemetry::common::SystemTimestamp c1_start; + opentelemetry::common::SystemTimestamp c2_start; + collect_value(collector1.get(), collection_ts1, c1_value, c1_start); + collect_value(collector2.get(), collection_ts1, c2_value, c2_start); + EXPECT_EQ(c1_value, 10); + EXPECT_EQ(c2_value, 10); + + // Cycle 2: A=30 observed once -> delta since last seen is 20 for each collector. + std::unordered_map measurements2 = { + {{{"attr", "A"}}, 30}}; + storage.RecordLong(measurements2, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + + c1_value = -1; + c2_value = -1; + collect_value(collector1.get(), collection_ts2, c1_value, c1_start); + collect_value(collector2.get(), collection_ts2, c2_value, c2_start); + EXPECT_EQ(c1_value, 20) << "Delta since previous collection must be 30 - 10 = 20"; + EXPECT_EQ(c2_value, 20) << "Second collector must independently receive the same increment"; + // The slow-path delta branch sets start_ts to the previous collection's timestamp. + EXPECT_EQ(c1_start, opentelemetry::common::SystemTimestamp(collection_ts1)) + << "Delta start_ts must continue from the previous collection"; + EXPECT_EQ(c2_start, opentelemetry::common::SystemTimestamp(collection_ts1)) + << "Delta start_ts must continue from the previous collection"; +} + } // namespace From 35630e2fe2ab31c1c41dd3f98a014d3e9e1f6f2b Mon Sep 17 00:00:00 2001 From: Patrick Summerer Date: Fri, 11 Sep 2026 10:43:22 +0200 Subject: [PATCH 3/5] [METRICS] Fix multi-collector cadence gap in async cumulative stale-attribute suppression The is_async_ cumulative export guard tested the shared delta_metrics map, which is only populated for whichever collector drains it first in a given cycle. Every other collector saw an empty map and had its already-observed attribute sets incorrectly suppressed for that cycle. Replace the check with a per-collector "observed this cycle" set, captured from each collector's own merged unreported deltas before the cumulative baseline is merged in. This preserves the #4108 stale-drop behavior while making it independent of collection order across collectors. Adds regression coverage for: - stale attribute suppression with two collectors - a collector lagging behind another by multiple cycles - start_ts/end_ts correctness for the multi-collector cumulative and delta paths --- .../metrics/state/temporal_metric_storage.cc | 51 ++- sdk/test/metrics/async_metric_storage_test.cc | 346 ++++++++++++------ 2 files changed, 273 insertions(+), 124 deletions(-) diff --git a/sdk/src/metrics/state/temporal_metric_storage.cc b/sdk/src/metrics/state/temporal_metric_storage.cc index 0e32a9b056..d57908d694 100644 --- a/sdk/src/metrics/state/temporal_metric_storage.cc +++ b/sdk/src/metrics/state/temporal_metric_storage.cc @@ -6,6 +6,7 @@ #include #include #include +#include #include #include @@ -137,6 +138,23 @@ bool TemporalMetricStorage::buildMetrics(CollectorHandle *collector, return true; }); } + + // For async cumulative exports, capture the attribute sets actually observed for THIS collector + // during THIS cycle (the freshly merged deltas), before the cumulative baseline is merged in. + // delta_metrics only carries the deltas for whichever collector drains the shared map first, so + // it cannot be used as the "observed this cycle" signal in multi-collector setups. + const bool async_cumulative = + is_async_ && aggregation_temporarily == AggregationTemporality::kCumulative; + std::unordered_set observed_this_cycle; + if (async_cumulative) + { + merged_metrics->GetAllEntries( + [&observed_this_cycle](const MetricAttributes &attributes, Aggregation &) { + observed_this_cycle.insert(attributes); + return true; + }); + } + // Get the last reported metrics for the `collector` from `last reported metrics` stash // - If the aggregation_temporarily for the collector is cumulative // - Merge the last reported metrics with unreported metrics (which is in merged_metrics), @@ -197,23 +215,22 @@ bool TemporalMetricStorage::buildMetrics(CollectorHandle *collector, metric_data.aggregation_temporality = aggregation_temporarily; metric_data.start_ts = last_collection_ts; metric_data.end_ts = collection_ts; - result_to_export->GetAllEntries([&metric_data, &delta_metrics, this]( - const MetricAttributes &attributes, - Aggregation &aggregation) { - if (is_async_ && metric_data.aggregation_temporality == AggregationTemporality::kCumulative && - !delta_metrics->Has(attributes)) - { - // Async cumulative exports must omit attribute sets that were not observed - // in the current callback cycle, while keeping the internal cumulative state. - return true; - } - - PointDataAttributes point_data_attr; - point_data_attr.point_data = aggregation.ToPoint(); - point_data_attr.attributes = attributes; - metric_data.point_data_attr_.emplace_back(std::move(point_data_attr)); - return true; - }); + result_to_export->GetAllEntries( + [&metric_data, &observed_this_cycle, async_cumulative](const MetricAttributes &attributes, + Aggregation &aggregation) { + if (async_cumulative && observed_this_cycle.find(attributes) == observed_this_cycle.end()) + { + // Async cumulative exports must omit attribute sets that were not observed + // in the current callback cycle, while keeping the internal cumulative state. + return true; + } + + PointDataAttributes point_data_attr; + point_data_attr.point_data = aggregation.ToPoint(); + point_data_attr.attributes = attributes; + metric_data.point_data_attr_.emplace_back(std::move(point_data_attr)); + return true; + }); if (metric_data.point_data_attr_.empty()) { diff --git a/sdk/test/metrics/async_metric_storage_test.cc b/sdk/test/metrics/async_metric_storage_test.cc index 931427174d..e3dcc9cf17 100644 --- a/sdk/test/metrics/async_metric_storage_test.cc +++ b/sdk/test/metrics/async_metric_storage_test.cc @@ -52,6 +52,10 @@ class WritableMetricStorageTestObservableGaugeFixture : public ::testing::TestWithParam {}; +class AsyncMetricStorageStaleAttributeFixture + : public ::testing::TestWithParam +{}; + TEST_P(AsyncWritableMetricStorageTestFixture, TestAggregation) { AggregationTemporality temporality = GetParam(); @@ -322,12 +326,12 @@ INSTANTIATE_TEST_SUITE_P(WritableMetricStorageTestObservableGaugeFixtureLong, ::testing::Values(AggregationTemporality::kCumulative, AggregationTemporality::kDelta)); -// Regression test for https://github.com/open-telemetry/opentelemetry-cpp/issues/4108 -// -// Async instruments under cumulative temporality must NOT carry forward attribute sets that were -// not reported by the callback in the current collection cycle. -TEST(AsyncMetricStorageRegressionTest, StaleAttributeSetDroppedInCumulativeExport) +// Async instruments must NOT carry forward attribute sets that were not reported by the callback in +// the current collection cycle. +TEST_P(AsyncMetricStorageStaleAttributeFixture, StaleAttributeSetDropped) { + const AggregationTemporality temporality = GetParam(); + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", InstrumentType::kObservableCounter, InstrumentValueType::kLong}; @@ -335,8 +339,7 @@ TEST(AsyncMetricStorageRegressionTest, StaleAttributeSetDroppedInCumulativeExpor // Some computation here auto collection_ts = sdk_start_ts + std::chrono::seconds(5); - std::shared_ptr collector( - new MockCollectorHandle(AggregationTemporality::kCumulative)); + std::shared_ptr collector(new MockCollectorHandle(temporality)); std::vector> collectors; collectors.push_back(collector); @@ -401,19 +404,20 @@ TEST(AsyncMetricStorageRegressionTest, StaleAttributeSetDroppedInCumulativeExpor }); // PUT must not appear – it was absent from the callback this cycle. - EXPECT_EQ(put_count, 0) << "Stale PUT attribute set must be dropped from cumulative export"; + EXPECT_EQ(put_count, 0) << "Stale PUT attribute set must be dropped"; EXPECT_EQ(get_count, 1); - EXPECT_EQ(get_value, 20); + // Cumulative exports the absolute value (20); delta exports the increment since last seen (10). + const int64_t expected_get_value = (temporality == AggregationTemporality::kCumulative) ? 20 : 10; + EXPECT_EQ(get_value, expected_get_value); } -// Regression test for https://github.com/open-telemetry/opentelemetry-cpp/issues/4108 -// -// Under delta temporality an attribute set that disappears for one collection cycle and then -// reappears must emit only the increment since the last observed value, not the full new absolute -// value. The cumulative baseline (cumulative_hash_map_) is preserved across absent cycles so that -// the delta computation in Record() remains correct. -TEST(AsyncMetricStorageRegressionTest, AttributeReappearanceAfterGapDeltaTemporality) +// An attribute set that disappears for one collection cycle and then reappears must preserve the +// cumulative baseline across the absent cycle. On reappearance cumulative exports the full absolute +// value, while delta exports only the increment since the attribute was last seen. +TEST_P(AsyncMetricStorageStaleAttributeFixture, AttributeReappearanceAfterGap) { + const AggregationTemporality temporality = GetParam(); + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", InstrumentType::kObservableCounter, InstrumentValueType::kLong}; @@ -421,8 +425,7 @@ TEST(AsyncMetricStorageRegressionTest, AttributeReappearanceAfterGapDeltaTempora // Some computation here auto collection_ts = sdk_start_ts + std::chrono::seconds(5); - std::shared_ptr collector( - new MockCollectorHandle(AggregationTemporality::kDelta)); + std::shared_ptr collector(new MockCollectorHandle(temporality)); std::vector> collectors; collectors.push_back(collector); @@ -433,23 +436,23 @@ TEST(AsyncMetricStorageRegressionTest, AttributeReappearanceAfterGapDeltaTempora #endif nullptr); - // Collection 1: A=10 → delta should be 10. + // Collection 1: A=10 -> both temporalities export 10. std::unordered_map measurements1 = { {{{"attr", "A"}}, 10}}; storage.RecordLong(measurements1, opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); - int64_t delta_value = -1; + int64_t value = -1; storage.Collect(collector.get(), collectors, sdk_start_ts, collection_ts, [&](const MetricData &metric_data) { for (const auto &data_attr : metric_data.point_data_attr_) { - delta_value = opentelemetry::nostd::get( + value = opentelemetry::nostd::get( opentelemetry::nostd::get(data_attr.point_data).value_); } return true; }); - EXPECT_EQ(delta_value, 10); + EXPECT_EQ(value, 10); // Collection 2: attribute A absent – nothing recorded, nothing emitted. std::unordered_map measurements2; @@ -464,46 +467,48 @@ TEST(AsyncMetricStorageRegressionTest, AttributeReappearanceAfterGapDeltaTempora }); EXPECT_EQ(attr_count, 0) << "No data points expected when attribute set is absent"; - // Collection 3: A reappears with absolute value 11. - // The cumulative baseline (10) was preserved across the absent cycle, so - // delta = 11 - 10 = 1 — the correct increment since the attribute was last seen. + // Collection 3: A reappears with absolute value 30. std::unordered_map measurements3 = { - {{{"attr", "A"}}, 11}}; + {{{"attr", "A"}}, 30}}; storage.RecordLong(measurements3, opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); - delta_value = -1; + value = -1; storage.Collect(collector.get(), collectors, sdk_start_ts, collection_ts + std::chrono::seconds(10), [&](const MetricData &metric_data) { for (const auto &data_attr : metric_data.point_data_attr_) { - delta_value = opentelemetry::nostd::get( + value = opentelemetry::nostd::get( opentelemetry::nostd::get(data_attr.point_data).value_); } return true; }); - // Baseline was preserved: delta = new_value - last_seen_value = 11 - 10 = 1. - EXPECT_EQ(delta_value, 1) - << "After a gap, reappearing attribute must emit the increment since last seen"; + // Baseline (10) preserved across the gap: cumulative = 30, delta = 30 - 10 = 20. + const int64_t expected_value = (temporality == AggregationTemporality::kCumulative) ? 30 : 20; + EXPECT_EQ(value, expected_value) + << "After a gap, reappearing attribute must preserve the baseline across the absent cycle"; } -// Regression test for the cumulative reappearance scenario. -// -// For async cumulative exports, stale attribute sets must be suppressed when absent, -// but the internal cumulative baseline must still be preserved. If A=10, then absent, -// then A=30, the reappearance must export 30 (full cumulative), not 20 (increment only). -TEST(AsyncMetricStorageRegressionTest, AttributeReappearanceAfterGapCumulativeTemporality) +// Stale suppression must apply to every collector, not just the one that drains the shared delta +// first. With two collectors, an attribute dropped by the callback must disappear from both +// collectors' exports while the still-reported attribute keeps its correct value for both. +TEST_P(AsyncMetricStorageStaleAttributeFixture, StaleAttributeSetDroppedMultiCollector) { + const 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 = sdk_start_ts + std::chrono::seconds(5); + auto sdk_start_ts = std::chrono::system_clock::now(); + auto collection_ts1 = sdk_start_ts + std::chrono::seconds(5); + auto collection_ts2 = sdk_start_ts + std::chrono::seconds(10); - std::shared_ptr collector( - new MockCollectorHandle(AggregationTemporality::kCumulative)); + std::shared_ptr collector1(new MockCollectorHandle(temporality)); + std::shared_ptr collector2(new MockCollectorHandle(temporality)); std::vector> collectors; - collectors.push_back(collector); + collectors.push_back(collector1); + collectors.push_back(collector2); + opentelemetry::sdk::metrics::AsyncMetricStorage storage( instr_desc, AggregationType::kSum, #ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW @@ -511,65 +516,159 @@ TEST(AsyncMetricStorageRegressionTest, AttributeReappearanceAfterGapCumulativeTe #endif nullptr); - // Collection 1: A=10 -> cumulative export should be 10. + auto collect_request_types = [&](CollectorHandle *collector, + opentelemetry::common::SystemTimestamp collection_ts, + int &get_count, int &put_count, int64_t &get_value) { + get_count = 0; + put_count = 0; + get_value = 0; + storage.Collect(collector, collectors, sdk_start_ts, collection_ts, + [&](const MetricData &metric_data) { + for (const auto &data_attr : metric_data.point_data_attr_) + { + const auto &key = opentelemetry::nostd::get( + data_attr.attributes.find("RequestType")->second); + if (key == "GET") + { + get_count++; + get_value = opentelemetry::nostd::get( + opentelemetry::nostd::get(data_attr.point_data).value_); + } + else if (key == "PUT") + { + put_count++; + } + } + return true; + }); + }; + + // Collection 1: both GET and PUT reported – both collectors see both attribute sets. std::unordered_map measurements1 = { - {{{"attr", "A"}}, 10}}; + {{{"RequestType", "GET"}}, 10}, {{{"RequestType", "PUT"}}, 5}}; storage.RecordLong(measurements1, opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); - int64_t cumulative_value = -1; - storage.Collect(collector.get(), collectors, sdk_start_ts, collection_ts, - [&](const MetricData &metric_data) { - for (const auto &data_attr : metric_data.point_data_attr_) - { - cumulative_value = opentelemetry::nostd::get( - opentelemetry::nostd::get(data_attr.point_data).value_); - } - return true; - }); - EXPECT_EQ(cumulative_value, 10); - - // Collection 2: A absent -> no export for A. - std::unordered_map measurements2; + int c1_get = 0, c1_put = 0, c2_get = 0, c2_put = 0; + int64_t c1_get_value = 0, c2_get_value = 0; + collect_request_types(collector1.get(), collection_ts1, c1_get, c1_put, c1_get_value); + collect_request_types(collector2.get(), collection_ts1, c2_get, c2_put, c2_get_value); + EXPECT_EQ(c1_get, 1); + EXPECT_EQ(c1_put, 1); + EXPECT_EQ(c2_get, 1); + EXPECT_EQ(c2_put, 1) << "Second collector must observe PUT even though the first drained the " + "shared delta"; + + // Collection 2: only GET reported – PUT is dropped by the callback. + std::unordered_map measurements2 = { + {{{"RequestType", "GET"}}, 20}}; storage.RecordLong(measurements2, opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); - int attr_count = 0; - storage.Collect(collector.get(), collectors, sdk_start_ts, - collection_ts + std::chrono::seconds(5), [&](const MetricData &metric_data) { - attr_count += static_cast(metric_data.point_data_attr_.size()); - return true; - }); - EXPECT_EQ(attr_count, 0) << "No data points expected when attribute set is absent"; + collect_request_types(collector1.get(), collection_ts2, c1_get, c1_put, c1_get_value); + collect_request_types(collector2.get(), collection_ts2, c2_get, c2_put, c2_get_value); - // Collection 3: A=30 -> cumulative export must be 30 (not 20). - std::unordered_map measurements3 = { - {{{"attr", "A"}}, 30}}; - storage.RecordLong(measurements3, + // Cumulative exports the absolute value (20); delta exports the increment since last seen (10). + const int64_t expected_get_value = (temporality == AggregationTemporality::kCumulative) ? 20 : 10; + + EXPECT_EQ(c1_put, 0) << "Stale PUT must be dropped for the first collector"; + EXPECT_EQ(c2_put, 0) << "Stale PUT must be dropped for the second collector too"; + EXPECT_EQ(c1_get, 1); + EXPECT_EQ(c2_get, 1); + EXPECT_EQ(c1_get_value, expected_get_value); + EXPECT_EQ(c2_get_value, expected_get_value) + << "Second collector must independently receive the still-reported GET value"; +} + +// A collector that skips a collection cycle must not lose values observed during its interval. +// Deltas drained by another collector are stashed for every collector, so when the lagging +// collector finally collects it receives the accumulated value — even though delta_metrics is empty +// at that moment (already drained by the other collector). This is the exact case the per-collector +// observed set fixes: suppression is driven by the collector's own unreported deltas, not by +// delta_metrics. +TEST_P(AsyncMetricStorageStaleAttributeFixture, MultiCollectorLaggingCollector) +{ + const AggregationTemporality temporality = GetParam(); + + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", InstrumentType::kObservableCounter, + InstrumentValueType::kLong}; + + auto sdk_start_ts = std::chrono::system_clock::now(); + auto collection_ts1 = sdk_start_ts + std::chrono::seconds(5); + auto collection_ts2 = sdk_start_ts + std::chrono::seconds(10); + + std::shared_ptr collector1(new MockCollectorHandle(temporality)); + std::shared_ptr collector2(new MockCollectorHandle(temporality)); + std::vector> collectors; + collectors.push_back(collector1); + collectors.push_back(collector2); + + opentelemetry::sdk::metrics::AsyncMetricStorage storage( + instr_desc, AggregationType::kSum, +#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW + ExemplarFilterType::kAlwaysOff, ExemplarReservoir::GetNoExemplarReservoir(), +#endif + nullptr); + + auto collect = [&](CollectorHandle *collector, + opentelemetry::common::SystemTimestamp collection_ts, int &count, + int64_t &value_out) { + count = 0; + value_out = 0; + storage.Collect(collector, collectors, sdk_start_ts, collection_ts, + [&](const MetricData &metric_data) { + for (const auto &data_attr : metric_data.point_data_attr_) + { + count++; + value_out = opentelemetry::nostd::get( + opentelemetry::nostd::get(data_attr.point_data).value_); + } + return true; + }); + }; + + int c1_count = 0, c2_count = 0; + int64_t c1_value = 0, c2_value = 0; + + // Cycle 1: A=10 observed. Only collector1 collects; collector2 lags behind. + std::unordered_map measurements1 = { + {{{"attr", "A"}}, 10}}; + storage.RecordLong(measurements1, opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + collect(collector1.get(), collection_ts1, c1_count, c1_value); + EXPECT_EQ(c1_count, 1); + EXPECT_EQ(c1_value, 10); - cumulative_value = -1; - storage.Collect(collector.get(), collectors, sdk_start_ts, - collection_ts + std::chrono::seconds(10), [&](const MetricData &metric_data) { - for (const auto &data_attr : metric_data.point_data_attr_) - { - cumulative_value = opentelemetry::nostd::get( - opentelemetry::nostd::get(data_attr.point_data).value_); - } - return true; - }); - EXPECT_EQ(cumulative_value, 30) - << "Cumulative reappearance must preserve baseline and export full cumulative value"; + // Cycle 2: A=20 observed. Only collector1 collects again; collector2 still lags. + std::unordered_map measurements2 = { + {{{"attr", "A"}}, 20}}; + storage.RecordLong(measurements2, + opentelemetry::common::SystemTimestamp(std::chrono::system_clock::now())); + collect(collector1.get(), collection_ts2, c1_count, c1_value); + EXPECT_EQ(c1_count, 1); + // Cumulative: absolute 20; delta: increment since previous collection (20 - 10 = 10). + EXPECT_EQ(c1_value, (temporality == AggregationTemporality::kCumulative) ? 20 : 10); + + // collector2 finally collects. delta_metrics is empty now (collector1 already drained it), but + // collector2's stash accumulated both cycles -> it must still receive the value. + collect(collector2.get(), collection_ts2, c2_count, c2_value); + EXPECT_EQ(c2_count, 1) << "Lagging collector must not drop a value observed during its interval"; + // Both temporalities yield 20 here: cumulative absolute is 20, and delta since collector2's first + // (never) collection is the full accumulated 10 + 10 = 20. + EXPECT_EQ(c2_value, 20) + << "Lagging collector must receive the accumulated value from the cycles it missed"; } -// Regression test for the multi-collector delta slow path. -// // With more than one collector, delta temporality no longer short-circuits on the single-collector -// fast path it goes through the slow path instead. This exercises -// the delta branch that sets start_ts to the previous collection's timestamp, and verifies both -// collectors independently receive the increment since their own last collection. -TEST(AsyncMetricStorageRegressionTest, MultiCollectorDeltaTemporality) +// fast path — it goes through the slow path instead. This exercises the delta branch that sets +// start_ts to the previous collection's timestamp, and verifies both collectors independently +// receive the correct value. Under cumulative every collector must observe the absolute value even +// though only the first collector of a cycle drains the shared delta (the per-collector observed +// set drives stale suppression, not delta_metrics), and start_ts stays at the SDK start. +TEST_P(AsyncMetricStorageStaleAttributeFixture, MultiCollector) { + const AggregationTemporality temporality = GetParam(); + InstrumentDescriptor instr_desc = {"name", "desc", "1unit", InstrumentType::kObservableCounter, InstrumentValueType::kLong}; @@ -577,10 +676,8 @@ TEST(AsyncMetricStorageRegressionTest, MultiCollectorDeltaTemporality) auto collection_ts1 = sdk_start_ts + std::chrono::seconds(5); auto collection_ts2 = sdk_start_ts + std::chrono::seconds(10); - std::shared_ptr collector1( - new MockCollectorHandle(AggregationTemporality::kDelta)); - std::shared_ptr collector2( - new MockCollectorHandle(AggregationTemporality::kDelta)); + std::shared_ptr collector1(new MockCollectorHandle(temporality)); + std::shared_ptr collector2(new MockCollectorHandle(temporality)); std::vector> collectors; collectors.push_back(collector1); collectors.push_back(collector2); @@ -594,10 +691,12 @@ TEST(AsyncMetricStorageRegressionTest, MultiCollectorDeltaTemporality) auto collect_value = [&](CollectorHandle *collector, opentelemetry::common::SystemTimestamp collection_ts, int64_t &value_out, - opentelemetry::common::SystemTimestamp &start_ts_out) { + opentelemetry::common::SystemTimestamp &start_ts_out, + opentelemetry::common::SystemTimestamp &end_ts_out) { storage.Collect(collector, collectors, sdk_start_ts, collection_ts, [&](const MetricData &metric_data) { start_ts_out = metric_data.start_ts; + end_ts_out = metric_data.end_ts; for (const auto &data_attr : metric_data.point_data_attr_) { value_out = opentelemetry::nostd::get( @@ -607,7 +706,7 @@ TEST(AsyncMetricStorageRegressionTest, MultiCollectorDeltaTemporality) }); }; - // Cycle 1: A=10 observed once, both collectors drain the same delta -> each sees 10. + // Cycle 1: A=10 observed once, both collectors drain the same value -> each sees 10. std::unordered_map measurements1 = { {{{"attr", "A"}}, 10}}; storage.RecordLong(measurements1, @@ -617,12 +716,19 @@ TEST(AsyncMetricStorageRegressionTest, MultiCollectorDeltaTemporality) int64_t c2_value = -1; opentelemetry::common::SystemTimestamp c1_start; opentelemetry::common::SystemTimestamp c2_start; - collect_value(collector1.get(), collection_ts1, c1_value, c1_start); - collect_value(collector2.get(), collection_ts1, c2_value, c2_start); + opentelemetry::common::SystemTimestamp c1_end; + opentelemetry::common::SystemTimestamp c2_end; + collect_value(collector1.get(), collection_ts1, c1_value, c1_start, c1_end); + collect_value(collector2.get(), collection_ts1, c2_value, c2_start, c2_end); EXPECT_EQ(c1_value, 10); - EXPECT_EQ(c2_value, 10); - - // Cycle 2: A=30 observed once -> delta since last seen is 20 for each collector. + EXPECT_EQ(c2_value, 10) << "Second collector must observe the value even though the first " + "collector drained the shared delta"; + EXPECT_EQ(c1_end, opentelemetry::common::SystemTimestamp(collection_ts1)) + << "end_ts must equal the collection timestamp"; + EXPECT_EQ(c2_end, opentelemetry::common::SystemTimestamp(collection_ts1)) + << "end_ts must equal the collection timestamp"; + + // Cycle 2: A=30 observed once. std::unordered_map measurements2 = { {{{"attr", "A"}}, 30}}; storage.RecordLong(measurements2, @@ -630,15 +736,41 @@ TEST(AsyncMetricStorageRegressionTest, MultiCollectorDeltaTemporality) c1_value = -1; c2_value = -1; - collect_value(collector1.get(), collection_ts2, c1_value, c1_start); - collect_value(collector2.get(), collection_ts2, c2_value, c2_start); - EXPECT_EQ(c1_value, 20) << "Delta since previous collection must be 30 - 10 = 20"; - EXPECT_EQ(c2_value, 20) << "Second collector must independently receive the same increment"; - // The slow-path delta branch sets start_ts to the previous collection's timestamp. - EXPECT_EQ(c1_start, opentelemetry::common::SystemTimestamp(collection_ts1)) - << "Delta start_ts must continue from the previous collection"; - EXPECT_EQ(c2_start, opentelemetry::common::SystemTimestamp(collection_ts1)) - << "Delta start_ts must continue from the previous collection"; + collect_value(collector1.get(), collection_ts2, c1_value, c1_start, c1_end); + collect_value(collector2.get(), collection_ts2, c2_value, c2_start, c2_end); + + // end_ts always advances to the current collection timestamp, for both temporalities. + EXPECT_EQ(c1_end, opentelemetry::common::SystemTimestamp(collection_ts2)) + << "end_ts must equal the current collection timestamp"; + EXPECT_EQ(c2_end, opentelemetry::common::SystemTimestamp(collection_ts2)) + << "end_ts must equal the current collection timestamp"; + + if (temporality == AggregationTemporality::kCumulative) + { + // Cumulative reports the absolute value and keeps start_ts at the SDK start. + EXPECT_EQ(c1_value, 30); + EXPECT_EQ(c2_value, 30) << "Second collector must independently receive the absolute value"; + EXPECT_EQ(c1_start, opentelemetry::common::SystemTimestamp(sdk_start_ts)) + << "Cumulative start_ts must stay at the SDK start"; + EXPECT_EQ(c2_start, opentelemetry::common::SystemTimestamp(sdk_start_ts)) + << "Cumulative start_ts must stay at the SDK start"; + } + else + { + // Delta reports the increment since each collector's previous collection. + EXPECT_EQ(c1_value, 20) << "Delta since previous collection must be 30 - 10 = 20"; + EXPECT_EQ(c2_value, 20) << "Second collector must independently receive the same increment"; + // The slow-path delta branch sets start_ts to the previous collection's timestamp. + EXPECT_EQ(c1_start, opentelemetry::common::SystemTimestamp(collection_ts1)) + << "Delta start_ts must continue from the previous collection"; + EXPECT_EQ(c2_start, opentelemetry::common::SystemTimestamp(collection_ts1)) + << "Delta start_ts must continue from the previous collection"; + } } +INSTANTIATE_TEST_SUITE_P(AsyncMetricStorageRegression, + AsyncMetricStorageStaleAttributeFixture, + ::testing::Values(AggregationTemporality::kCumulative, + AggregationTemporality::kDelta)); + } // namespace From 420e4336518163bbb6a6fdde9b4da3aba82776bb Mon Sep 17 00:00:00 2001 From: Patrick Summerer Date: Fri, 11 Sep 2026 10:43:22 +0200 Subject: [PATCH 4/5] [METRICS] Fix multi-collector cadence gap in async cumulative stale-attribute suppression The is_async_ cumulative export guard tested the shared delta_metrics map, which is only populated for whichever collector drains it first in a given cycle. Every other collector saw an empty map and had its already-observed attribute sets incorrectly suppressed for that cycle. Replace the check with a per-collector "observed this cycle" set, captured from each collector's own merged unreported deltas before the cumulative baseline is merged in. This preserves the #4108 stale-drop behavior while making it independent of collection order across collectors. Adds regression coverage for: - stale attribute suppression with two collectors - a collector lagging behind another by multiple cycles - start_ts/end_ts correctness for the multi-collector cumulative and delta paths --- .../opentelemetry/sdk/metrics/state/async_metric_storage.h | 6 ------ 1 file changed, 6 deletions(-) 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 c2fc160610..8578599fe2 100644 --- a/sdk/include/opentelemetry/sdk/metrics/state/async_metric_storage.h +++ b/sdk/include/opentelemetry/sdk/metrics/state/async_metric_storage.h @@ -134,12 +134,6 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora delta_metrics = std::move(delta_hash_map_); delta_hash_map_ = std::make_unique(aggregation_config_->cardinality_limit_); - // cumulative_hash_map_ is intentionally NOT pruned here. - // It preserves the last-seen absolute value for every attribute set so that - // delta computation in Record() remains correct if an attribute set reappears - // after being absent for one or more collection cycles. - // Stale entries are suppressed at export time by the is_async_ guard in - // TemporalMetricStorage::buildMetrics() instead. } auto status = From 038608ef6bcb96968ba37a40e79a1ae744e9e775 Mon Sep 17 00:00:00 2001 From: Patrick Summerer Date: Thu, 17 Sep 2026 07:57:03 +0200 Subject: [PATCH 5/5] [METRICS] Introduce instrument async decision as util function --- .../opentelemetry/sdk/metrics/instruments.h | 23 +++++++++++++++++++ .../sdk/metrics/state/async_metric_storage.h | 2 +- .../metrics/state/temporal_metric_storage.h | 4 +--- .../metrics/state/temporal_metric_storage.cc | 9 ++++---- 4 files changed, 29 insertions(+), 9 deletions(-) diff --git a/sdk/include/opentelemetry/sdk/metrics/instruments.h b/sdk/include/opentelemetry/sdk/metrics/instruments.h index 938d87fd8a..3051952cb8 100644 --- a/sdk/include/opentelemetry/sdk/metrics/instruments.h +++ b/sdk/include/opentelemetry/sdk/metrics/instruments.h @@ -155,6 +155,29 @@ struct InstrumentDescriptorUtil return "Unknown"; } } + + static bool IsInstrumentTypeAsync(InstrumentType type) noexcept + { + switch (type) + { + case InstrumentType::kCounter: + return false; + case InstrumentType::kUpDownCounter: + return false; + case InstrumentType::kHistogram: + return false; + case InstrumentType::kObservableCounter: + return true; + case InstrumentType::kObservableUpDownCounter: + return true; + case InstrumentType::kObservableGauge: + return true; + case InstrumentType::kGauge: + return false; + default: + return false; + } + } }; struct InstrumentEqualNameCaseInsensitive 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 8578599fe2..b1c31cf324 100644 --- a/sdk/include/opentelemetry/sdk/metrics/state/async_metric_storage.h +++ b/sdk/include/opentelemetry/sdk/metrics/state/async_metric_storage.h @@ -53,7 +53,7 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora exemplar_filter_type_(exemplar_filter_type), exemplar_reservoir_(std::move(exemplar_reservoir)), #endif - temporal_metric_storage_(instrument_descriptor, aggregation_type, aggregation_config, true) + temporal_metric_storage_(instrument_descriptor, aggregation_type, aggregation_config) {} template diff --git a/sdk/include/opentelemetry/sdk/metrics/state/temporal_metric_storage.h b/sdk/include/opentelemetry/sdk/metrics/state/temporal_metric_storage.h index a752c7430e..50fa1a2240 100644 --- a/sdk/include/opentelemetry/sdk/metrics/state/temporal_metric_storage.h +++ b/sdk/include/opentelemetry/sdk/metrics/state/temporal_metric_storage.h @@ -34,8 +34,7 @@ class TemporalMetricStorage public: TemporalMetricStorage(InstrumentDescriptor instrument_descriptor, AggregationType aggregation_type, - const AggregationConfig *aggregation_config, - bool is_async = false); + const AggregationConfig *aggregation_config); bool buildMetrics(CollectorHandle *collector, nostd::span> collectors, @@ -65,7 +64,6 @@ class TemporalMetricStorage // See https://github.com/open-telemetry/opentelemetry-specification (logs/metrics // SDK specs) and issue #4062. const opentelemetry::common::SystemTimestamp instrument_creation_ts_; - bool is_async_ = false; }; } // namespace metrics } // namespace sdk diff --git a/sdk/src/metrics/state/temporal_metric_storage.cc b/sdk/src/metrics/state/temporal_metric_storage.cc index d57908d694..2187811e0f 100644 --- a/sdk/src/metrics/state/temporal_metric_storage.cc +++ b/sdk/src/metrics/state/temporal_metric_storage.cc @@ -32,13 +32,11 @@ namespace metrics TemporalMetricStorage::TemporalMetricStorage(InstrumentDescriptor instrument_descriptor, AggregationType aggregation_type, - const AggregationConfig *aggregation_config, - bool is_async) + const AggregationConfig *aggregation_config) : instrument_descriptor_(std::move(instrument_descriptor)), aggregation_type_(aggregation_type), aggregation_config_(aggregation_config), - instrument_creation_ts_(std::chrono::system_clock::now()), - is_async_(is_async) + instrument_creation_ts_(std::chrono::system_clock::now()) {} bool TemporalMetricStorage::buildMetrics(CollectorHandle *collector, @@ -144,7 +142,8 @@ bool TemporalMetricStorage::buildMetrics(CollectorHandle *collector, // delta_metrics only carries the deltas for whichever collector drains the shared map first, so // it cannot be used as the "observed this cycle" signal in multi-collector setups. const bool async_cumulative = - is_async_ && aggregation_temporarily == AggregationTemporality::kCumulative; + InstrumentDescriptorUtil::IsInstrumentTypeAsync(instrument_descriptor_.type_) && + aggregation_temporarily == AggregationTemporality::kCumulative; std::unordered_set observed_this_cycle; if (async_cumulative) {