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]

Reply via email to