From eeec4b26ea8dc518641cc4a6bf769bccd5cd4f30 Mon Sep 17 00:00:00 2001 From: Deniz Ariyan Date: Wed, 9 Sep 2026 18:25:39 +0200 Subject: [PATCH 1/6] [SDK] Fix MetricReader Shutdown invoking OnShutDown() multiple times --- CHANGELOG.md | 3 + .../opentelemetry/sdk/metrics/metric_reader.h | 5 ++ sdk/src/metrics/metric_reader.cc | 8 +-- sdk/test/metrics/metric_reader_test.cc | 64 +++++++++++++++++++ 4 files changed, 76 insertions(+), 4 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 59f7d6529b..fea7621377 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -15,6 +15,9 @@ Increment the: ## [Unreleased] +* [SDK] Fix `MetricReader::Shutdown()` invoking `OnShutDown()` multiple times. + [#4536](https://github.com/open-telemetry/opentelemetry-cpp/issues/4536) + * [DOC] Fix and clarify the `StartSpanOptions` documentation [#4526](https://github.com/open-telemetry/opentelemetry-cpp/pull/4526) diff --git a/sdk/include/opentelemetry/sdk/metrics/metric_reader.h b/sdk/include/opentelemetry/sdk/metrics/metric_reader.h index 4813861314..5a77a99740 100644 --- a/sdk/include/opentelemetry/sdk/metrics/metric_reader.h +++ b/sdk/include/opentelemetry/sdk/metrics/metric_reader.h @@ -71,6 +71,11 @@ class MetricReader /** * Shutdown the metric reader. + * + * Idempotent. Only the first call performs the shutdown, later calls log a warning and + * return true. + * + * @return the result of OnShutDown() for the first call, true for any subsequent call. */ bool Shutdown(std::chrono::microseconds timeout = (std::chrono::microseconds::max)()) noexcept; diff --git a/sdk/src/metrics/metric_reader.cc b/sdk/src/metrics/metric_reader.cc index c9312fc183..acc219d776 100644 --- a/sdk/src/metrics/metric_reader.cc +++ b/sdk/src/metrics/metric_reader.cc @@ -48,13 +48,13 @@ bool MetricReader::Collect( bool MetricReader::Shutdown(std::chrono::microseconds timeout) noexcept { - bool status = true; - if (IsShutdown()) + bool expected = false; + if (!shutdown_.compare_exchange_strong(expected, true, std::memory_order_acq_rel)) { OTEL_INTERNAL_LOG_WARN("MetricReader::Shutdown - Cannot invoke shutdown twice!"); + return true; } - - shutdown_.store(true, std::memory_order_release); + bool status = true; if (!OnShutDown(timeout)) { diff --git a/sdk/test/metrics/metric_reader_test.cc b/sdk/test/metrics/metric_reader_test.cc index 0860fa1ef0..502f9f7cdb 100644 --- a/sdk/test/metrics/metric_reader_test.cc +++ b/sdk/test/metrics/metric_reader_test.cc @@ -2,9 +2,12 @@ // SPDX-License-Identifier: Apache-2.0 #include +#include #include #include +#include #include +#include #include "common.h" #include "opentelemetry/sdk/instrumentationscope/instrumentation_scope.h" @@ -147,3 +150,64 @@ TEST(MetricReaderTest, CardinalityLimitsExplicitSdkDefaultIsHonoured) EXPECT_EQ(metric_reader->GetCardinalityLimit(InstrumentType::kHistogram), 1500); EXPECT_EQ(metric_reader->GetCardinalityLimit(InstrumentType::kUpDownCounter), 1500); } + +namespace +{ + +class CountingMetricReader : public MetricReader +{ +public: + AggregationTemporality GetAggregationTemporality(InstrumentType) const noexcept override + { + return AggregationTemporality::kCumulative; + } + + int shutdown_count{0}; + int force_flush_count{0}; + +private: + bool OnForceFlush(std::chrono::microseconds) noexcept override + { + ++force_flush_count; + return true; + } + + bool OnShutDown(std::chrono::microseconds) noexcept override + { + ++shutdown_count; + return true; + } +}; + +} // namespace + +TEST(MetricReaderTest, ShutdownIsInvokedOnce) +{ + CountingMetricReader reader; + + EXPECT_TRUE(reader.Shutdown()); + EXPECT_TRUE(reader.IsShutdown()); + EXPECT_EQ(reader.shutdown_count, 1); + + EXPECT_TRUE(reader.Shutdown()); + EXPECT_TRUE(reader.Shutdown()); + EXPECT_EQ(reader.shutdown_count, 1); +} + +TEST(MetricReaderTest, ConcurrentShutdownIsInvokedOnce) +{ + CountingMetricReader reader; + + std::vector threads; + threads.reserve(8); + for (int i = 0; i < 8; i++) + { + threads.emplace_back([&reader]() { reader.Shutdown(); }); + } + for (auto &thread : threads) + { + thread.join(); + } + + EXPECT_EQ(reader.shutdown_count, 1); +} From 35c0e765007408883aebe9dee2d3a2cfe655d017 Mon Sep 17 00:00:00 2001 From: Deniz Ariyan Date: Wed, 9 Sep 2026 18:38:57 +0200 Subject: [PATCH 2/6] [SDK] Fix MetricReader invoking OnForceFlush on shutdown instance --- CHANGELOG.md | 4 ++++ .../opentelemetry/sdk/metrics/metric_reader.h | 3 +++ sdk/src/metrics/metric_reader.cc | 7 +++++-- sdk/test/metrics/metric_reader_test.cc | 13 +++++++++++++ 4 files changed, 25 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index fea7621377..d3242cdcf2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -18,6 +18,10 @@ Increment the: * [SDK] Fix `MetricReader::Shutdown()` invoking `OnShutDown()` multiple times. [#4536](https://github.com/open-telemetry/opentelemetry-cpp/issues/4536) +* [SDK] Fix `MetricReader::ForceFlush()` invoking `OnForceFlush()` on a + shutdown reader. + [#4548](https://github.com/open-telemetry/opentelemetry-cpp/pull/4548) + * [DOC] Fix and clarify the `StartSpanOptions` documentation [#4526](https://github.com/open-telemetry/opentelemetry-cpp/pull/4526) diff --git a/sdk/include/opentelemetry/sdk/metrics/metric_reader.h b/sdk/include/opentelemetry/sdk/metrics/metric_reader.h index 5a77a99740..17ab4fff0f 100644 --- a/sdk/include/opentelemetry/sdk/metrics/metric_reader.h +++ b/sdk/include/opentelemetry/sdk/metrics/metric_reader.h @@ -81,6 +81,9 @@ class MetricReader /** * Force flush the metric read by the reader. + * + * @return false without invoking OnForceFlush() if the reader is already shut down, otherwise + * the result of OnForceFlush(). */ bool ForceFlush(std::chrono::microseconds timeout = (std::chrono::microseconds::max)()) noexcept; diff --git a/sdk/src/metrics/metric_reader.cc b/sdk/src/metrics/metric_reader.cc index acc219d776..9e29199d87 100644 --- a/sdk/src/metrics/metric_reader.cc +++ b/sdk/src/metrics/metric_reader.cc @@ -67,11 +67,14 @@ bool MetricReader::Shutdown(std::chrono::microseconds timeout) noexcept /** Flush metric read by this reader **/ bool MetricReader::ForceFlush(std::chrono::microseconds timeout) noexcept { - bool status = true; if (IsShutdown()) { - OTEL_INTERNAL_LOG_WARN("MetricReader::Shutdown Cannot invoke Force flush on shutdown reader!"); + OTEL_INTERNAL_LOG_WARN( + "MetricReader::ForceFlush Cannot invoke Force flush on shutdown reader!"); + return false; } + + bool status = true; if (!OnForceFlush(timeout)) { status = false; diff --git a/sdk/test/metrics/metric_reader_test.cc b/sdk/test/metrics/metric_reader_test.cc index 502f9f7cdb..5f75b3bd26 100644 --- a/sdk/test/metrics/metric_reader_test.cc +++ b/sdk/test/metrics/metric_reader_test.cc @@ -211,3 +211,16 @@ TEST(MetricReaderTest, ConcurrentShutdownIsInvokedOnce) EXPECT_EQ(reader.shutdown_count, 1); } + +TEST(MetricReaderTest, ForceFlushAfterShutdownIsNoOp) +{ + CountingMetricReader reader; + + EXPECT_TRUE(reader.ForceFlush()); + EXPECT_EQ(reader.force_flush_count, 1); + + EXPECT_TRUE(reader.Shutdown()); + + EXPECT_FALSE(reader.ForceFlush()); + EXPECT_EQ(reader.force_flush_count, 1); +} From 6c1f4f846a9e5edb3b0b55b3f7acf47a9b23cdaf Mon Sep 17 00:00:00 2001 From: Deniz Ariyan Date: Wed, 9 Sep 2026 18:50:41 +0200 Subject: [PATCH 3/6] [SDK] Fix race in test --- sdk/test/metrics/metric_reader_test.cc | 25 ++++++++++++++++++------- 1 file changed, 18 insertions(+), 7 deletions(-) diff --git a/sdk/test/metrics/metric_reader_test.cc b/sdk/test/metrics/metric_reader_test.cc index 5f75b3bd26..96362ef213 100644 --- a/sdk/test/metrics/metric_reader_test.cc +++ b/sdk/test/metrics/metric_reader_test.cc @@ -2,6 +2,7 @@ // SPDX-License-Identifier: Apache-2.0 #include +#include #include #include #include @@ -10,6 +11,8 @@ #include #include "common.h" +#include "opentelemetry/nostd/shared_ptr.h" +#include "opentelemetry/sdk/common/global_log_handler.h" #include "opentelemetry/sdk/instrumentationscope/instrumentation_scope.h" #include "opentelemetry/sdk/metrics/cardinality_limits.h" #include "opentelemetry/sdk/metrics/export/metric_producer.h" @@ -162,8 +165,8 @@ class CountingMetricReader : public MetricReader return AggregationTemporality::kCumulative; } - int shutdown_count{0}; - int force_flush_count{0}; + std::atomic shutdown_count{0}; + std::atomic force_flush_count{0}; private: bool OnForceFlush(std::chrono::microseconds) noexcept override @@ -187,17 +190,23 @@ TEST(MetricReaderTest, ShutdownIsInvokedOnce) EXPECT_TRUE(reader.Shutdown()); EXPECT_TRUE(reader.IsShutdown()); - EXPECT_EQ(reader.shutdown_count, 1); + EXPECT_EQ(reader.shutdown_count.load(), 1); EXPECT_TRUE(reader.Shutdown()); EXPECT_TRUE(reader.Shutdown()); - EXPECT_EQ(reader.shutdown_count, 1); + EXPECT_EQ(reader.shutdown_count.load(), 1); } TEST(MetricReaderTest, ConcurrentShutdownIsInvokedOnce) { + namespace internal_log = opentelemetry::sdk::common::internal_log; CountingMetricReader reader; + // default logger is not thread-safe + auto previous_handler = internal_log::GlobalLogHandler::GetLogHandler(); + internal_log::GlobalLogHandler::SetLogHandler( + nostd::shared_ptr(new internal_log::NoopLogHandler())); + std::vector threads; threads.reserve(8); for (int i = 0; i < 8; i++) @@ -209,7 +218,9 @@ TEST(MetricReaderTest, ConcurrentShutdownIsInvokedOnce) thread.join(); } - EXPECT_EQ(reader.shutdown_count, 1); + internal_log::GlobalLogHandler::SetLogHandler(previous_handler); + + EXPECT_EQ(reader.shutdown_count.load(), 1); } TEST(MetricReaderTest, ForceFlushAfterShutdownIsNoOp) @@ -217,10 +228,10 @@ TEST(MetricReaderTest, ForceFlushAfterShutdownIsNoOp) CountingMetricReader reader; EXPECT_TRUE(reader.ForceFlush()); - EXPECT_EQ(reader.force_flush_count, 1); + EXPECT_EQ(reader.force_flush_count.load(), 1); EXPECT_TRUE(reader.Shutdown()); EXPECT_FALSE(reader.ForceFlush()); - EXPECT_EQ(reader.force_flush_count, 1); + EXPECT_EQ(reader.force_flush_count.load(), 1); } From eee6e32e57b5db199b1bec325cde29332101f89c Mon Sep 17 00:00:00 2001 From: Deniz Ariyan Date: Wed, 9 Sep 2026 20:14:32 +0200 Subject: [PATCH 4/6] [SDK] Fix IWYU warnings --- sdk/test/metrics/metric_reader_test.cc | 1 - 1 file changed, 1 deletion(-) diff --git a/sdk/test/metrics/metric_reader_test.cc b/sdk/test/metrics/metric_reader_test.cc index 96362ef213..7162de20a3 100644 --- a/sdk/test/metrics/metric_reader_test.cc +++ b/sdk/test/metrics/metric_reader_test.cc @@ -4,7 +4,6 @@ #include #include #include -#include #include #include #include From e0aeaf49b04573344edfad9183e64afdcdaf7811 Mon Sep 17 00:00:00 2001 From: Deniz Ariyan Date: Fri, 11 Sep 2026 18:05:23 +0200 Subject: [PATCH 5/6] [SDK] Fix concurrent MetricReader Shutdown returning early --- CHANGELOG.md | 1 + .../opentelemetry/sdk/metrics/metric_reader.h | 6 +- sdk/src/metrics/metric_reader.cc | 8 +- sdk/test/metrics/metric_reader_test.cc | 104 +++++++++++++++++- 4 files changed, 112 insertions(+), 7 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index b26e6f7033..85f40930c7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,6 +16,7 @@ Increment the: ## [Unreleased] * [SDK] Fix `MetricReader::Shutdown()` invoking `OnShutDown()` multiple times. + Concurrent calls now block until the first call's shutdown has completed. [#4536](https://github.com/open-telemetry/opentelemetry-cpp/issues/4536) * [SDK] Fix `MetricReader::ForceFlush()` invoking `OnForceFlush()` on a diff --git a/sdk/include/opentelemetry/sdk/metrics/metric_reader.h b/sdk/include/opentelemetry/sdk/metrics/metric_reader.h index 17ab4fff0f..64bf62e245 100644 --- a/sdk/include/opentelemetry/sdk/metrics/metric_reader.h +++ b/sdk/include/opentelemetry/sdk/metrics/metric_reader.h @@ -6,6 +6,7 @@ #include #include #include +#include #include "opentelemetry/nostd/function_ref.h" #include "opentelemetry/sdk/metrics/cardinality_limits.h" @@ -72,8 +73,8 @@ class MetricReader /** * Shutdown the metric reader. * - * Idempotent. Only the first call performs the shutdown, later calls log a warning and - * return true. + * Idempotent and a completion barrier: only the first call runs OnShutDown(), and a + * concurrent call blocks until the first call has finished before returning. * * @return the result of OnShutDown() for the first call, true for any subsequent call. */ @@ -104,6 +105,7 @@ class MetricReader protected: private: MetricProducer *metric_producer_{nullptr}; + std::mutex shutdown_m_; std::atomic shutdown_{false}; CardinalityLimits cardinality_limits_; }; diff --git a/sdk/src/metrics/metric_reader.cc b/sdk/src/metrics/metric_reader.cc index 9e29199d87..3220aa61f8 100644 --- a/sdk/src/metrics/metric_reader.cc +++ b/sdk/src/metrics/metric_reader.cc @@ -2,6 +2,7 @@ // SPDX-License-Identifier: Apache-2.0 #include "opentelemetry/sdk/metrics/metric_reader.h" +#include #include "opentelemetry/sdk/common/global_log_handler.h" #include "opentelemetry/sdk/metrics/cardinality_limits.h" #include "opentelemetry/sdk/metrics/export/metric_producer.h" @@ -48,10 +49,11 @@ bool MetricReader::Collect( bool MetricReader::Shutdown(std::chrono::microseconds timeout) noexcept { - bool expected = false; - if (!shutdown_.compare_exchange_strong(expected, true, std::memory_order_acq_rel)) + // Serialize so concurrent calls block until the first call's shutdown has completed. + std::lock_guard shutdown_guard{shutdown_m_}; + if (shutdown_.exchange(true, std::memory_order_release)) { - OTEL_INTERNAL_LOG_WARN("MetricReader::Shutdown - Cannot invoke shutdown twice!"); + OTEL_INTERNAL_LOG_WARN("MetricReader::Shutdown - Already shutdown!"); return true; } bool status = true; diff --git a/sdk/test/metrics/metric_reader_test.cc b/sdk/test/metrics/metric_reader_test.cc index 7162de20a3..c648ab661a 100644 --- a/sdk/test/metrics/metric_reader_test.cc +++ b/sdk/test/metrics/metric_reader_test.cc @@ -4,6 +4,7 @@ #include #include #include +#include #include #include #include @@ -159,6 +160,9 @@ namespace class CountingMetricReader : public MetricReader { public: + // hook_result is what both OnForceFlush() and OnShutDown() report back. + explicit CountingMetricReader(bool hook_result = true) : hook_result_(hook_result) {} + AggregationTemporality GetAggregationTemporality(InstrumentType) const noexcept override { return AggregationTemporality::kCumulative; @@ -171,14 +175,16 @@ class CountingMetricReader : public MetricReader bool OnForceFlush(std::chrono::microseconds) noexcept override { ++force_flush_count; - return true; + return hook_result_; } bool OnShutDown(std::chrono::microseconds) noexcept override { ++shutdown_count; - return true; + return hook_result_; } + + const bool hook_result_; }; } // namespace @@ -234,3 +240,97 @@ TEST(MetricReaderTest, ForceFlushAfterShutdownIsNoOp) EXPECT_FALSE(reader.ForceFlush()); EXPECT_EQ(reader.force_flush_count.load(), 1); } + +TEST(MetricReaderTest, FailedShutdownIsReportedAndNotRetried) +{ + CountingMetricReader reader{/* hook_result= */ false}; + + EXPECT_FALSE(reader.ForceFlush()); + EXPECT_EQ(reader.force_flush_count.load(), 1); + + // The first call reports the hook's failure. + EXPECT_FALSE(reader.Shutdown()); + EXPECT_TRUE(reader.IsShutdown()); + EXPECT_EQ(reader.shutdown_count.load(), 1); + + // A later call is a no-op that succeeds, without re-entering the failed hook. + EXPECT_TRUE(reader.Shutdown()); + EXPECT_EQ(reader.shutdown_count.load(), 1); + + // Flush stays rejected even after a failed shutdown, without re-entering the hook. + EXPECT_FALSE(reader.ForceFlush()); + EXPECT_EQ(reader.force_flush_count.load(), 1); +} + +namespace +{ + +// Parks inside OnShutDown() until released, to observe what a concurrent caller sees. +class BlockingMetricReader : public MetricReader +{ +public: + AggregationTemporality GetAggregationTemporality(InstrumentType) const noexcept override + { + return AggregationTemporality::kCumulative; + } + + std::promise entered_shutdown; + std::promise release_shutdown; + std::atomic shutdown_finished{false}; + +private: + bool OnForceFlush(std::chrono::microseconds) noexcept override { return true; } + + bool OnShutDown(std::chrono::microseconds) noexcept override + { + entered_shutdown.set_value(); + release_shutdown.get_future().wait(); + shutdown_finished.store(true, std::memory_order_release); + return true; + } +}; + +} // namespace + +TEST(MetricReaderTest, ConcurrentShutdownWaitsForCleanupToComplete) +{ + namespace internal_log = opentelemetry::sdk::common::internal_log; + BlockingMetricReader reader; + + // default logger is not thread-safe + auto previous_handler = internal_log::GlobalLogHandler::GetLogHandler(); + internal_log::GlobalLogHandler::SetLogHandler( + nostd::shared_ptr(new internal_log::NoopLogHandler())); + + auto entered = reader.entered_shutdown.get_future(); + std::thread first([&reader]() { EXPECT_TRUE(reader.Shutdown()); }); + + // The first caller owns the shutdown and is now parked inside OnShutDown(). + entered.wait(); + EXPECT_TRUE(reader.IsShutdown()); + EXPECT_FALSE(reader.shutdown_finished.load(std::memory_order_acquire)); + + std::atomic second_returned{false}; + std::promise second_started; + auto started = second_started.get_future(); + std::thread second([&]() { + second_started.set_value(); + // Block until first caller releases the shutdown, then return true. + EXPECT_TRUE(reader.Shutdown()); + EXPECT_TRUE(reader.shutdown_finished.load(std::memory_order_acquire)); + second_returned.store(true, std::memory_order_release); + }); + + started.wait(); + // Arbitrary sleep to ensure second is still blocked and not just we were too fast to check. + std::this_thread::sleep_for(std::chrono::milliseconds(50)); + EXPECT_FALSE(second_returned.load(std::memory_order_acquire)); + + reader.release_shutdown.set_value(); + second.join(); + first.join(); + + internal_log::GlobalLogHandler::SetLogHandler(previous_handler); + + EXPECT_TRUE(second_returned.load(std::memory_order_acquire)); +} From 5bf688a55ea2dbf6c3c3ed2ebaa8fc3b11029935 Mon Sep 17 00:00:00 2001 From: Deniz Ariyan Date: Sat, 12 Sep 2026 16:45:23 +0200 Subject: [PATCH 6/6] Fix CHANGELOG entry --- CHANGELOG.md | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index b7b191f467..1189ddac9e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -15,15 +15,15 @@ Increment the: ## [Unreleased] -* [SDK] Fix `MetricReader::ForceFlush()` invoking `OnForceFlush()` on a - shutdown reader. - [#4548](https://github.com/open-telemetry/opentelemetry-cpp/pull/4548) - * [BUG] Install a curl seek callback so an OTLP/HTTP export body can be rewound when libcurl restarts an upload, instead of failing with `CURLE_SEND_FAIL_REWIND` and dropping the batch ([#4549](https://github.com/open-telemetry/opentelemetry-cpp/issues/4549)) - + +* [SDK] Fix `MetricReader::ForceFlush()` invoking `OnForceFlush()` on a + shutdown reader. + [#4548](https://github.com/open-telemetry/opentelemetry-cpp/pull/4548) + * [SDK] Fix `MetricReader::Shutdown()` invoking `OnShutDown()` multiple times. Concurrent calls now block until the first call's shutdown has completed. [#4536](https://github.com/open-telemetry/opentelemetry-cpp/issues/4536)