This is an automated email from the ASF dual-hosted git repository. yiguolei pushed a commit to branch branch-4.2 in repository https://gitbox.apache.org/repos/asf/doris.git
commit 7f117e69f799b8254f20d66b3c00f719a986873f Author: Gavin Chou <[email protected]> AuthorDate: Mon Sep 28 09:35:43 2026 +0800 branch-4.1: [fix](cloud) Drain bvar timer callbacks before destruction (#68522) (#68523) ### What problem does this PR solve? Issue Number: N/A Related PR: #68522 Problem Summary: This backports #68522 to branch-4.1. Cloud ASAN unit tests can finish every gtest case and then crash during process teardown in bvar::MultiDimension from a bthread TimerThread callback. MBvarLatencyRecorderWithStatus owns recurring updater callbacks through recorder_, but implicit reverse-order member destruction tears down the status bvars and mutexes before recorder_. A callback can therefore keep using the parent after those dependencies start being destroyed. The old stop path also ignored that bthread_timer_del returns 1 when a callback is already running and does not join it. The fix serializes each updater lifecycle, prevents rescheduling after stop, waits for an active callback to finish, and explicitly drains every updater in the parent destructor while all callback dependencies are still alive. A deterministic synchronization-point regression test holds a callback in flight while another thread destroys the parent and verifies that destruction waits. ### Verification - local_to_dev clean ASAN build on branch-4.1 - bvars_test: 2/2 passed, including DestructorWaitsForRunningUpdate - meta_service_test: 223/223 passed; global teardown and process exit completed normally without ASAN errors - clang-format 16 dry run passed - clang-tidy with the ASAN compilation database completed with no user-code diagnostics ### Release note None ### Check List (For Author) - Test - [x] Regression test - [x] Unit Test - Behavior changed: - [x] No. - Does this need documentation? - [x] No. --- cloud/src/common/bvars.h | 188 ++++++++++++++++++++++++++++++---------------- cloud/test/bvars_test.cpp | 50 +++++++++++- 2 files changed, 173 insertions(+), 65 deletions(-) diff --git a/cloud/src/common/bvars.h b/cloud/src/common/bvars.h index 7a02be09ae0..bbd1f10adc5 100644 --- a/cloud/src/common/bvars.h +++ b/cloud/src/common/bvars.h @@ -19,6 +19,7 @@ #include <aws/core/external/cjson/cJSON.h> #include <bthread/bthread.h> +#include <bthread/condition_variable.h> #include <bthread/mutex.h> #include <bthread/unstable.h> #include <bvar/bvar.h> @@ -30,7 +31,6 @@ #include <cpp/sync_point.h> #include <gmock/gmock-actions.h> -#include <atomic> #include <cstdint> #include <initializer_list> #include <map> @@ -39,6 +39,7 @@ #include <string> #include <type_traits> #include <utility> +#include <vector> #include "common/logging.h" @@ -345,17 +346,15 @@ private: * @return true if the timer was successfully started, false otherwise */ bool start() { - if (!_started.load()) { - { - std::lock_guard<bthread::Mutex> l(init_mutex_); - if (!_started.load()) { - if (!schedule()) { - return false; - } - _started.store(true); - } - return true; - } + std::lock_guard<bthread::Mutex> l(lifecycle_mutex_); + if (_started) { + return true; + } + + _started = true; + if (!schedule_locked()) { + _started = false; + return false; } return true; } @@ -365,12 +364,14 @@ private: * Scheduling a one-time task. * This is useful if you want to reset the timer interval. */ - bool schedule() { + // lifecycle_mutex_ must be held so stop() cannot race with replacing _timer. + bool schedule_locked() { if (bthread_timer_add(&_timer, butil::seconds_from_now(_interval_s), update, this) != 0) { LOG(WARNING) << "Failed to add bthread timer for ScheduledLatencyUpdater"; return false; } + _callback_pending = true; return true; } @@ -384,8 +385,8 @@ private: */ static void update(void* arg) { auto* latency_updater = static_cast<ScheduledLatencyUpdater*>(arg); - if (!latency_updater || !latency_updater->_started) { - LOG(WARNING) << "Invalid ScheduledLatencyUpdater in timer callback"; + CHECK(latency_updater != nullptr); + if (!latency_updater->begin_callback()) { return; } @@ -395,53 +396,48 @@ private: auto* parent = static_cast<MBvarLatencyRecorderWithStatus*>(latency_updater->_arg); if (!parent) { LOG(WARNING) << "Invalid parent container in timer callback"; - return; - } - - std::list<std::string> current_dim_list; - { - std::lock_guard<bthread::Mutex> l(parent->recorder_mutex_); - for (const auto& it : parent->recorder_) { - if (it.second.get() == latency_updater) { - current_dim_list = it.first; - break; + } else { + std::list<std::string> current_dim_list; + { + std::lock_guard<bthread::Mutex> l(parent->recorder_mutex_); + for (const auto& it : parent->recorder_) { + if (it.second.get() == latency_updater) { + current_dim_list = it.first; + break; + } } } - } - - if (current_dim_list.empty()) { - LOG(WARNING) << "Could not find dimension for ScheduledLatencyUpdater"; - return; - } - - { - std::lock_guard<bthread::Mutex> l(parent->timer_mutex_); - - bvar::Status<int64_t>* max_status = parent->max_status_.get_stats(current_dim_list); - bvar::Status<int64_t>* avg_status = parent->avg_status_.get_stats(current_dim_list); - bvar::Status<int64_t>* count_status = - parent->count_status_.get_stats(current_dim_list); - - VLOG_DEBUG << "Updating latency recorder status for dimension, " - << "max_latency: " << latency_updater->max_latency() - << ", avg_latency: " << latency_updater->latency(); - TEST_SYNC_POINT("mBvarLatencyRecorderWithStatus::update"); - if (max_status) { - max_status->set_value(latency_updater->max_latency()); - } - if (avg_status) { - avg_status->set_value(latency_updater->latency()); - } - if (count_status) { - count_status->set_value(latency_updater->count()); + if (current_dim_list.empty()) { + LOG(WARNING) << "Could not find dimension for ScheduledLatencyUpdater"; + } else { + std::lock_guard<bthread::Mutex> l(parent->timer_mutex_); + + bvar::Status<int64_t>* max_status = + parent->max_status_.get_stats(current_dim_list); + bvar::Status<int64_t>* avg_status = + parent->avg_status_.get_stats(current_dim_list); + bvar::Status<int64_t>* count_status = + parent->count_status_.get_stats(current_dim_list); + + VLOG_DEBUG << "Updating latency recorder status for dimension, " + << "max_latency: " << latency_updater->max_latency() + << ", avg_latency: " << latency_updater->latency(); + TEST_SYNC_POINT("mBvarLatencyRecorderWithStatus::update"); + + if (max_status) { + max_status->set_value(latency_updater->max_latency()); + } + if (avg_status) { + avg_status->set_value(latency_updater->latency()); + } + if (count_status) { + count_status->set_value(latency_updater->count()); + } } } - if (latency_updater->_started && !latency_updater->schedule()) { - LOG(WARNING) << "Failed to reschedule timer for ScheduledLatencyUpdater"; - latency_updater->_started = false; - } + latency_updater->finish_callback(); } /** @@ -451,18 +447,64 @@ private: * any pending callbacks from accessing potentially freed resources. */ void stop() { - if (_started.load()) { - bthread_timer_del(_timer); - _started = false; + std::unique_lock<bthread::Mutex> l(lifecycle_mutex_); + if (!_started && !_callback_pending) { + return; + } + + // Prevent a running callback from scheduling the next timer before trying to + // cancel the current one. bthread_timer_del() returns 1 when the callback is + // already running; in that case the updater and its parent must stay alive until + // finish_callback() signals that the callback no longer accesses either object. + _started = false; + if (!_callback_pending) { + return; + } + + const int timer_state = bthread_timer_del(_timer); + if (timer_state == 0) { + _callback_pending = false; + return; + } + + CHECK_EQ(1, timer_state); + TEST_SYNC_POINT("mBvarLatencyRecorderWithStatus::stop"); + while (_callback_pending) { + callback_finished_.wait(l); } } private: - int _interval_s; // Timer interval in seconds - void* _arg; // Argument to pass to the callback - bthread_timer_t _timer; // The bthread timer handle - std::atomic_bool _started {false}; // Whether the timer has been started - bthread::Mutex init_mutex_; // Mutex for timer_map_ + bool begin_callback() { + std::lock_guard<bthread::Mutex> l(lifecycle_mutex_); + if (_started) { + return true; + } + + // stop() may observe the timer as running before this function acquires the + // lifecycle mutex. Acknowledge the canceled callback so stop() can finish. + _callback_pending = false; + callback_finished_.notify_all(); + return false; + } + + void finish_callback() { + std::lock_guard<bthread::Mutex> l(lifecycle_mutex_); + _callback_pending = false; + if (_started && !schedule_locked()) { + LOG(WARNING) << "Failed to reschedule timer for ScheduledLatencyUpdater"; + _started = false; + } + callback_finished_.notify_all(); + } + + int _interval_s; // Timer interval in seconds + void* _arg; // Argument to pass to the callback + bthread_timer_t _timer; // The bthread timer handle + bool _started = false; // Whether callbacks should keep running + bool _callback_pending = false; // A timer callback is scheduled or running + bthread::Mutex lifecycle_mutex_; + bthread::ConditionVariable callback_finished_; }; public: @@ -483,6 +525,24 @@ public: const std::initializer_list<std::string>& dim_names) : MBvarLatencyRecorderWithStatus(prefix + "_" + metric_name, dim_names) {} + ~MBvarLatencyRecorderWithStatus() { + // Members are destroyed in reverse declaration order, so recorder_ (which owns the + // timer callbacks) would otherwise be destroyed after the status bvars and mutexes + // used by those callbacks. Stop and drain every callback while all parent members are + // still alive. + std::vector<std::shared_ptr<ScheduledLatencyUpdater>> latency_updaters; + { + std::lock_guard<bthread::Mutex> l(recorder_mutex_); + latency_updaters.reserve(recorder_.size()); + for (const auto& entry : recorder_) { + latency_updaters.push_back(entry.second); + } + } + for (const auto& latency_updater : latency_updaters) { + latency_updater->stop(); + } + } + /** * @brief Record a latency value * diff --git a/cloud/test/bvars_test.cpp b/cloud/test/bvars_test.cpp index 95062465a4e..d82820823f8 100644 --- a/cloud/test/bvars_test.cpp +++ b/cloud/test/bvars_test.cpp @@ -19,6 +19,8 @@ #include <brpc/server.h> #include <bthread/bthread.h> +#include <bthread/countdown_event.h> +#include <butil/time.h> #include <glog/logging.h> #include <gtest/gtest.h> @@ -49,11 +51,17 @@ int main(int argc, char** argv) { class BvarsTest : public ::testing::Test { public: void SetUp() override { + auto* sp = SyncPoint::get_instance(); + sp->disable_processing(); + sp->clear_all_call_backs(); if (server.Start("0.0.0.0:0", &options) == -1) { perror("Start brpc server"); } } void TearDown() override { + auto* sp = SyncPoint::get_instance(); + sp->disable_processing(); + sp->clear_all_call_backs(); server.Stop(0); server.Join(); } @@ -126,4 +134,44 @@ TEST(BvarsTest, MultiThreadRecordMetrics) { ASSERT_GT(update_count.load(), 200); } -} // namespace doris::cloud \ No newline at end of file +TEST(BvarsTest, DestructorWaitsForRunningUpdate) { + constexpr int callback_timeout_s = 5; + int interval_s = 1; + bthread::CountdownEvent update_started; + bthread::CountdownEvent allow_update_to_finish; + bthread::CountdownEvent stop_waiting; + + auto* sp = SyncPoint::get_instance(); + sp->set_call_back("mBvarLatencyRecorderWithStatus::put", [&interval_s](auto&& args) { + auto* interval = try_any_cast<int*>(args[0]); + *interval = interval_s; + }); + sp->set_call_back("mBvarLatencyRecorderWithStatus::update", [&](auto&&) { + update_started.signal(); + allow_update_to_finish.wait(); + }); + sp->set_call_back("mBvarLatencyRecorderWithStatus::stop", + [&](auto&&) { stop_waiting.signal(); }); + sp->enable_processing(); + + auto recorder = std::make_unique<MBvarLatencyRecorderWithStatus<60>>( + "destructor_wait_test", std::initializer_list<std::string> {"instance_id"}); + recorder->put({"instance"}, 1); + const int update_result = + update_started.timed_wait(butil::seconds_from_now(callback_timeout_s)); + if (update_result != 0) { + // Keep cleanup non-blocking even when the timer thread did not invoke the callback. + allow_update_to_finish.signal(); + EXPECT_EQ(0, update_result); + return; + } + + std::thread destroyer([&] { recorder.reset(); }); + const int stop_result = stop_waiting.timed_wait(butil::seconds_from_now(callback_timeout_s)); + allow_update_to_finish.signal(); + destroyer.join(); + + EXPECT_EQ(0, stop_result); +} + +} // namespace doris::cloud --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
