Skip to content

Commit c8ae8cb

Browse files
mrdrivingduckcodex
andcommitted
fix(executor): support destruction from worker threads
Decouple worker state from the executor object so that a task can safely destroy its final executor owner from a worker thread. Cover the lifecycle with executor and S3 asynchronous range-read tests. Co-authored-by: GPT-5.6 Terra <codex@users.noreply.github.com>
1 parent 05d6497 commit c8ae8cb

3 files changed

Lines changed: 101 additions & 33 deletions

File tree

src/paimon/common/executor/default_executor_test.cpp

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -125,6 +125,26 @@ TEST(DefaultExecutorTest, TestAddTaskAfterShutdownNowIgnored) {
125125
ASSERT_EQ(executed_count.load(), 0);
126126
}
127127

128+
TEST(DefaultExecutorTest, TestDestroyFromWorkerThread) {
129+
std::unique_ptr<Executor> created = CreateDefaultExecutor();
130+
std::shared_ptr<Executor> executor(std::move(created));
131+
std::shared_ptr<Executor> task_executor = executor;
132+
auto release = std::make_shared<std::promise<void>>();
133+
std::shared_future<void> release_future = release->get_future().share();
134+
auto destroyed = std::make_shared<std::promise<void>>();
135+
std::future<void> future = destroyed->get_future();
136+
137+
executor->Add([executor = std::move(task_executor), release_future, destroyed]() mutable {
138+
release_future.wait();
139+
executor.reset();
140+
destroyed->set_value();
141+
});
142+
143+
executor.reset();
144+
release->set_value();
145+
ASSERT_EQ(std::future_status::ready, future.wait_for(std::chrono::seconds(5)));
146+
}
147+
128148
TEST(DefaultExecutorTest, TestAddTaskFromMultipleThreads) {
129149
ASSERT_OK_AND_ASSIGN(auto executor, CreateDefaultExecutor(/*thread_count=*/4));
130150

src/paimon/common/executor/executor.cpp

Lines changed: 31 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -40,23 +40,26 @@ class DefaultExecutor : public Executor {
4040
uint32_t GetThreadNum() const override;
4141

4242
private:
43-
void WorkerThread();
43+
struct State {
44+
std::queue<std::function<void()>> tasks;
45+
std::mutex mutex;
46+
std::condition_variable condition;
47+
bool stop = false;
48+
};
49+
50+
static void WorkerThread(std::shared_ptr<State> state);
4451
void ShutdownInternal(bool wait_for_pending_tasks);
4552

4653
private:
4754
uint32_t thread_count_;
4855
std::vector<std::thread> workers_;
49-
std::queue<std::function<void()>> tasks_;
50-
std::mutex queue_mutex_;
51-
std::condition_variable condition_;
52-
bool stop_ = false;
53-
int32_t active_tasks_ = 0;
56+
std::shared_ptr<State> state_ = std::make_shared<State>();
5457
};
5558

5659
DefaultExecutor::DefaultExecutor(uint32_t thread_count) : thread_count_(thread_count) {
5760
assert(thread_count > 0);
5861
for (uint32_t i = 0; i < thread_count_; ++i) {
59-
workers_.emplace_back(&DefaultExecutor::WorkerThread, this);
62+
workers_.emplace_back(&DefaultExecutor::WorkerThread, state_);
6063
}
6164
}
6265

@@ -66,21 +69,22 @@ uint32_t DefaultExecutor::GetThreadNum() const {
6669

6770
void DefaultExecutor::ShutdownInternal(bool wait_for_pending_tasks) {
6871
{
69-
std::unique_lock<std::mutex> lock(queue_mutex_);
70-
if (stop_) {
71-
return;
72-
}
73-
stop_ = true;
72+
std::unique_lock<std::mutex> lock(state_->mutex);
73+
state_->stop = true;
7474
if (!wait_for_pending_tasks) {
7575
// Discard all pending tasks immediately.
7676
std::queue<std::function<void()>> empty;
77-
tasks_.swap(empty);
77+
state_->tasks.swap(empty);
7878
}
79-
condition_.notify_all();
79+
state_->condition.notify_all();
8080
}
8181
for (std::thread& worker : workers_) {
8282
if (worker.joinable()) {
83-
worker.join();
83+
if (worker.get_id() == std::this_thread::get_id()) {
84+
worker.detach();
85+
} else {
86+
worker.join();
87+
}
8488
}
8589
}
8690
}
@@ -100,38 +104,32 @@ void DefaultExecutor::Add(std::function<void()> func) {
100104
return;
101105
}
102106
{
103-
std::unique_lock<std::mutex> lock(queue_mutex_);
104-
if (stop_) {
107+
std::unique_lock<std::mutex> lock(state_->mutex);
108+
if (state_->stop) {
105109
return;
106110
}
107-
tasks_.emplace(std::move(func));
111+
state_->tasks.emplace(std::move(func));
108112
}
109-
condition_.notify_one();
113+
state_->condition.notify_one();
110114
}
111115

112-
void DefaultExecutor::WorkerThread() {
116+
void DefaultExecutor::WorkerThread(std::shared_ptr<State> state) {
113117
while (true) {
114118
std::function<void()> task;
115119
{
116-
std::unique_lock<std::mutex> lock(queue_mutex_);
117-
condition_.wait(lock, [this] { return stop_ || !tasks_.empty(); });
118-
if (stop_ && tasks_.empty() && active_tasks_ == 0) {
119-
condition_.notify_all();
120+
std::unique_lock<std::mutex> lock(state->mutex);
121+
state->condition.wait(lock, [&state] { return state->stop || !state->tasks.empty(); });
122+
if (state->stop && state->tasks.empty()) {
123+
state->condition.notify_all();
120124
return;
121125
}
122-
if (!tasks_.empty()) {
123-
task = std::move(tasks_.front());
124-
tasks_.pop();
125-
++active_tasks_;
126+
if (!state->tasks.empty()) {
127+
task = std::move(state->tasks.front());
128+
state->tasks.pop();
126129
}
127130
}
128131
if (task) {
129132
task();
130-
std::unique_lock<std::mutex> lock(queue_mutex_);
131-
--active_tasks_;
132-
if (tasks_.empty() && active_tasks_ == 0) {
133-
condition_.notify_all();
134-
}
135133
}
136134
}
137135
}

src/paimon/fs/s3/s3_file_system_test.cpp

Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,11 +21,17 @@
2121

2222
#include <gtest/gtest.h>
2323

24+
#include <chrono>
25+
#include <condition_variable>
2426
#include <cstdlib>
2527
#include <cstring>
2628
#include <filesystem>
2729
#include <fstream>
30+
#include <functional>
31+
#include <future>
32+
#include <mutex>
2833
#include <optional>
34+
#include <thread>
2935
#include <utility>
3036
#include <vector>
3137

@@ -40,6 +46,9 @@ class MockHttpClient : public HttpClient {
4046
public:
4147
Result<HttpResponse> Execute(const HttpRequest& request,
4248
const HttpBodyConsumer& consumer) const override {
49+
if (before_execute_) {
50+
before_execute_();
51+
}
4352
request_ = request;
4453
HttpResponse response;
4554
response.status_code = status_code_;
@@ -55,6 +64,7 @@ class MockHttpClient : public HttpClient {
5564
int32_t status_code_ = 200;
5665
HttpHeaders response_headers_;
5766
std::string body_;
67+
std::function<void()> before_execute_;
5868
};
5969

6070
class ScopedEnvironmentVariable {
@@ -435,6 +445,46 @@ TEST(S3ObjectStoreClientTest, TestRangeAndListObjects) {
435445
ASSERT_NE(http->request_.url.find("continuation-token=old%20token"), std::string::npos);
436446
}
437447

448+
TEST(S3ObjectStoreClientTest, TestGetObjectRangeAsyncClientLifetime) {
449+
auto http = std::make_shared<MockHttpClient>();
450+
http->body_ = "data";
451+
std::mutex mutex;
452+
std::condition_variable condition;
453+
bool request_started = false;
454+
bool release_request = false;
455+
http->before_execute_ = [&] {
456+
std::unique_lock<std::mutex> lock(mutex);
457+
request_started = true;
458+
condition.notify_one();
459+
condition.wait(lock, [&] { return release_request; });
460+
};
461+
ASSERT_OK_AND_ASSIGN(std::shared_ptr<ObjectStoreClient> client,
462+
MakeS3ObjectStoreClient(StaticOptions(), http));
463+
char buffer[4];
464+
auto promise = std::make_shared<std::promise<Status>>();
465+
std::future<Status> future = promise->get_future();
466+
467+
client->GetObjectRangeAsync({"bucket", "key"}, 0, 4, buffer, [promise](Status status) {
468+
promise->set_value(std::move(status));
469+
});
470+
{
471+
std::unique_lock<std::mutex> lock(mutex);
472+
ASSERT_TRUE(
473+
condition.wait_for(lock, std::chrono::seconds(5), [&] { return request_started; }));
474+
}
475+
std::thread destruction_thread([client = std::move(client)]() mutable { client.reset(); });
476+
{
477+
std::lock_guard<std::mutex> lock(mutex);
478+
release_request = true;
479+
}
480+
condition.notify_one();
481+
482+
destruction_thread.join();
483+
ASSERT_EQ(std::future_status::ready, future.wait_for(std::chrono::seconds(5)));
484+
ASSERT_OK(future.get());
485+
ASSERT_EQ("data", std::string(buffer, sizeof(buffer)));
486+
}
487+
438488
TEST(S3ObjectStoreClientTest, TestUrlEncodedListObjects) {
439489
auto http = std::make_shared<MockHttpClient>();
440490
ASSERT_OK_AND_ASSIGN(std::shared_ptr<ObjectStoreClient> client,

0 commit comments

Comments
 (0)