This is an automated email from the ASF dual-hosted git repository.
hello-stephen pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 6472420a19a [fix](cloud) Drain bvar timer callbacks before destruction
(#68522)
6472420a19a is described below
commit 6472420a19aaa72add975c40a8da42805e8ab3bd
Author: Gavin Chou <[email protected]>
AuthorDate: Thu Oct 8 19:38:43 2026 +0800
[fix](cloud) Drain bvar timer callbacks before destruction (#68522)
Problem Summary:
Cloud ASAN UT could finish every gtest case and then crash during
process teardown in
MBvarLatencyRecorderWithStatus::ScheduledLatencyUpdater::update while
accessing a bvar MultiDimension map. One observed failure is TeamCity
build 1055210:
http://43.132.222.7:8111/buildConfiguration/Doris_DorisCloudUt_CloudUt/1055210
The investigation showed that this is a lifecycle race in the Doris
wrapper rather than an FDB or brpc data-path failure:
1. The test binary reached global test teardown before the crash, and
the failing thread was the bthread TimerThread.
2. ScheduledLatencyUpdater keeps a raw pointer to its parent and
periodically reads the parent recorder map, mutexes, and status bvars.
3. With the implicit parent destructor, those dependencies begin
destruction before recorder_ releases its updater objects.
4. In brpc 1.4.0, bthread_timer_del returns 1 when a callback is already
running and does not join it. The old stop path ignored that return
value.
The fix serializes timer start, reschedule, and stop; tracks whether a
callback is scheduled or running; waits for an active callback to
finish; and explicitly drains every updater at the beginning of parent
destruction. The parent map lock is released before waiting so a
callback can finish its dimension lookup without deadlocking. This
guarantees that callbacks no longer access the parent before its members
are destroyed.
A deterministic regression test blocks an update callback, destroys the
parent concurrently, verifies that destruction waits, then releases the
callback.
---
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 19d8c2c01ed..44e2c66d1ee 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]