This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new e0b5cf2f853 branch-4.1: [fix](cloud) Drain bvar timer callbacks before
destruction (#68522) (#68523)
e0b5cf2f853 is described below
commit e0b5cf2f853af3e02f17cf2c036926eb4e52c883
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]