This is an automated email from the ASF dual-hosted git repository.

lxy-9602 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-cpp.git


The following commit(s) were added to refs/heads/main by this push:
     new 86449758 fix(executor): support destruction from worker threads (#233)
86449758 is described below

commit 864497581488637094990366aeb5d21a05c7c1f3
Author: Mr Dk. <[email protected]>
AuthorDate: Fri Aug 28 16:47:52 2026 +0800

    fix(executor): support destruction from worker threads (#233)
---
 .../common/executor/default_executor_test.cpp      | 51 ++++++++++++++++++
 src/paimon/common/executor/executor.cpp            | 63 +++++++++++-----------
 src/paimon/fs/s3/s3_file_system_test.cpp           | 50 +++++++++++++++++
 3 files changed, 133 insertions(+), 31 deletions(-)

diff --git a/src/paimon/common/executor/default_executor_test.cpp 
b/src/paimon/common/executor/default_executor_test.cpp
index 91f2c94c..07e78681 100644
--- a/src/paimon/common/executor/default_executor_test.cpp
+++ b/src/paimon/common/executor/default_executor_test.cpp
@@ -125,6 +125,57 @@ TEST(DefaultExecutorTest, 
TestAddTaskAfterShutdownNowIgnored) {
     ASSERT_EQ(executed_count.load(), 0);
 }
 
+TEST(DefaultExecutorTest, TestConcurrentShutdownNow) {
+    constexpr int32_t kShutdownThreadCount = 2;
+    constexpr int32_t kAttempts = 50;
+    for (int32_t attempt = 0; attempt < kAttempts; ++attempt) {
+        ASSERT_OK_AND_ASSIGN(auto executor, 
CreateDefaultExecutor(/*thread_count=*/4));
+        std::atomic<int32_t> ready_shutdown_count = 0;
+        std::promise<void> start_signal;
+        std::shared_future<void> start_future = 
start_signal.get_future().share();
+        std::vector<std::thread> shutdown_threads;
+        shutdown_threads.reserve(kShutdownThreadCount);
+
+        for (int32_t thread_index = 0; thread_index < kShutdownThreadCount; 
++thread_index) {
+            shutdown_threads.emplace_back([&]() {
+                ++ready_shutdown_count;
+                start_future.wait();
+                executor->ShutdownNow();
+            });
+        }
+        for (int32_t retry = 0; retry < 100 && ready_shutdown_count.load() < 
kShutdownThreadCount;
+             ++retry) {
+            std::this_thread::sleep_for(std::chrono::milliseconds(10));
+        }
+        const int32_t ready_count_before_start = ready_shutdown_count.load();
+        start_signal.set_value();
+        for (std::thread& shutdown_thread : shutdown_threads) {
+            shutdown_thread.join();
+        }
+        ASSERT_EQ(kShutdownThreadCount, ready_count_before_start);
+    }
+}
+
+TEST(DefaultExecutorTest, TestDestroyFromWorkerThread) {
+    std::unique_ptr<Executor> created = CreateDefaultExecutor();
+    std::shared_ptr<Executor> executor(std::move(created));
+    std::shared_ptr<Executor> task_executor = executor;
+    auto release = std::make_shared<std::promise<void>>();
+    std::shared_future<void> release_future = release->get_future().share();
+    auto destroyed = std::make_shared<std::promise<void>>();
+    std::future<void> future = destroyed->get_future();
+
+    executor->Add([executor = std::move(task_executor), release_future, 
destroyed]() mutable {
+        release_future.wait();
+        executor.reset();
+        destroyed->set_value();
+    });
+
+    executor.reset();
+    release->set_value();
+    ASSERT_EQ(std::future_status::ready, 
future.wait_for(std::chrono::seconds(5)));
+}
+
 TEST(DefaultExecutorTest, TestAddTaskFromMultipleThreads) {
     ASSERT_OK_AND_ASSIGN(auto executor, 
CreateDefaultExecutor(/*thread_count=*/4));
 
diff --git a/src/paimon/common/executor/executor.cpp 
b/src/paimon/common/executor/executor.cpp
index cd3f699d..45253fe1 100644
--- a/src/paimon/common/executor/executor.cpp
+++ b/src/paimon/common/executor/executor.cpp
@@ -40,23 +40,26 @@ class DefaultExecutor : public Executor {
     uint32_t GetThreadNum() const override;
 
  private:
-    void WorkerThread();
+    struct State {
+        std::queue<std::function<void()>> tasks;
+        std::mutex mutex;
+        std::condition_variable condition;
+        bool stop = false;
+    };
+
+    static void WorkerThread(std::shared_ptr<State> state);
     void ShutdownInternal(bool wait_for_pending_tasks);
 
  private:
     uint32_t thread_count_;
     std::vector<std::thread> workers_;
-    std::queue<std::function<void()>> tasks_;
-    std::mutex queue_mutex_;
-    std::condition_variable condition_;
-    bool stop_ = false;
-    int32_t active_tasks_ = 0;
+    std::shared_ptr<State> state_ = std::make_shared<State>();
 };
 
 DefaultExecutor::DefaultExecutor(uint32_t thread_count) : 
thread_count_(thread_count) {
     assert(thread_count > 0);
     for (uint32_t i = 0; i < thread_count_; ++i) {
-        workers_.emplace_back(&DefaultExecutor::WorkerThread, this);
+        workers_.emplace_back(&DefaultExecutor::WorkerThread, state_);
     }
 }
 
@@ -66,21 +69,25 @@ uint32_t DefaultExecutor::GetThreadNum() const {
 
 void DefaultExecutor::ShutdownInternal(bool wait_for_pending_tasks) {
     {
-        std::unique_lock<std::mutex> lock(queue_mutex_);
-        if (stop_) {
+        std::unique_lock<std::mutex> lock(state_->mutex);
+        if (state_->stop) {
             return;
         }
-        stop_ = true;
+        state_->stop = true;
         if (!wait_for_pending_tasks) {
             // Discard all pending tasks immediately.
             std::queue<std::function<void()>> empty;
-            tasks_.swap(empty);
+            state_->tasks.swap(empty);
         }
-        condition_.notify_all();
+        state_->condition.notify_all();
     }
     for (std::thread& worker : workers_) {
         if (worker.joinable()) {
-            worker.join();
+            if (worker.get_id() == std::this_thread::get_id()) {
+                worker.detach();
+            } else {
+                worker.join();
+            }
         }
     }
 }
@@ -100,38 +107,32 @@ void DefaultExecutor::Add(std::function<void()> func) {
         return;
     }
     {
-        std::unique_lock<std::mutex> lock(queue_mutex_);
-        if (stop_) {
+        std::unique_lock<std::mutex> lock(state_->mutex);
+        if (state_->stop) {
             return;
         }
-        tasks_.emplace(std::move(func));
+        state_->tasks.emplace(std::move(func));
     }
-    condition_.notify_one();
+    state_->condition.notify_one();
 }
 
-void DefaultExecutor::WorkerThread() {
+void DefaultExecutor::WorkerThread(std::shared_ptr<State> state) {
     while (true) {
         std::function<void()> task;
         {
-            std::unique_lock<std::mutex> lock(queue_mutex_);
-            condition_.wait(lock, [this] { return stop_ || !tasks_.empty(); });
-            if (stop_ && tasks_.empty() && active_tasks_ == 0) {
-                condition_.notify_all();
+            std::unique_lock<std::mutex> lock(state->mutex);
+            state->condition.wait(lock, [&state] { return state->stop || 
!state->tasks.empty(); });
+            if (state->stop && state->tasks.empty()) {
+                state->condition.notify_all();
                 return;
             }
-            if (!tasks_.empty()) {
-                task = std::move(tasks_.front());
-                tasks_.pop();
-                ++active_tasks_;
+            if (!state->tasks.empty()) {
+                task = std::move(state->tasks.front());
+                state->tasks.pop();
             }
         }
         if (task) {
             task();
-            std::unique_lock<std::mutex> lock(queue_mutex_);
-            --active_tasks_;
-            if (tasks_.empty() && active_tasks_ == 0) {
-                condition_.notify_all();
-            }
         }
     }
 }
diff --git a/src/paimon/fs/s3/s3_file_system_test.cpp 
b/src/paimon/fs/s3/s3_file_system_test.cpp
index 09b1bb4b..063a7dcb 100644
--- a/src/paimon/fs/s3/s3_file_system_test.cpp
+++ b/src/paimon/fs/s3/s3_file_system_test.cpp
@@ -21,11 +21,17 @@
 
 #include <gtest/gtest.h>
 
+#include <chrono>
+#include <condition_variable>
 #include <cstdlib>
 #include <cstring>
 #include <filesystem>
 #include <fstream>
+#include <functional>
+#include <future>
+#include <mutex>
 #include <optional>
+#include <thread>
 #include <utility>
 #include <vector>
 
@@ -40,6 +46,9 @@ class MockHttpClient : public HttpClient {
  public:
     Result<HttpResponse> Execute(const HttpRequest& request,
                                  const HttpBodyConsumer& consumer) const 
override {
+        if (before_execute_) {
+            before_execute_();
+        }
         request_ = request;
         HttpResponse response;
         response.status_code = status_code_;
@@ -55,6 +64,7 @@ class MockHttpClient : public HttpClient {
     int32_t status_code_ = 200;
     HttpHeaders response_headers_;
     std::string body_;
+    std::function<void()> before_execute_;
 };
 
 class ScopedEnvironmentVariable {
@@ -435,6 +445,46 @@ TEST(S3ObjectStoreClientTest, TestRangeAndListObjects) {
     ASSERT_NE(http->request_.url.find("continuation-token=old%20token"), 
std::string::npos);
 }
 
+TEST(S3ObjectStoreClientTest, TestGetObjectRangeAsyncClientLifetime) {
+    auto http = std::make_shared<MockHttpClient>();
+    http->body_ = "data";
+    std::mutex mutex;
+    std::condition_variable condition;
+    bool request_started = false;
+    bool release_request = false;
+    http->before_execute_ = [&] {
+        std::unique_lock<std::mutex> lock(mutex);
+        request_started = true;
+        condition.notify_one();
+        condition.wait(lock, [&] { return release_request; });
+    };
+    ASSERT_OK_AND_ASSIGN(std::shared_ptr<ObjectStoreClient> client,
+                         MakeS3ObjectStoreClient(StaticOptions(), http));
+    char buffer[4];
+    auto promise = std::make_shared<std::promise<Status>>();
+    std::future<Status> future = promise->get_future();
+
+    client->GetObjectRangeAsync({"bucket", "key"}, 0, 4, buffer, 
[promise](Status status) {
+        promise->set_value(std::move(status));
+    });
+    {
+        std::unique_lock<std::mutex> lock(mutex);
+        ASSERT_TRUE(
+            condition.wait_for(lock, std::chrono::seconds(5), [&] { return 
request_started; }));
+    }
+    std::thread destruction_thread([client = std::move(client)]() mutable { 
client.reset(); });
+    {
+        std::lock_guard<std::mutex> lock(mutex);
+        release_request = true;
+    }
+    condition.notify_one();
+
+    destruction_thread.join();
+    ASSERT_EQ(std::future_status::ready, 
future.wait_for(std::chrono::seconds(5)));
+    ASSERT_OK(future.get());
+    ASSERT_EQ("data", std::string(buffer, sizeof(buffer)));
+}
+
 TEST(S3ObjectStoreClientTest, TestUrlEncodedListObjects) {
     auto http = std::make_shared<MockHttpClient>();
     ASSERT_OK_AND_ASSIGN(std::shared_ptr<ObjectStoreClient> client,

Reply via email to