src/ex/thread_pool.cpp

100.0% Lines (140/140) 100.0% List of functions (29/29)
thread_pool.cpp
f(x) Functions (29)
Function Calls Lines Blocks
boost::capy::thread_pool::impl::push(boost::capy::continuation*) :71 20629x 100.0% 100.0% boost::capy::thread_pool::impl::pop() :81 20900x 100.0% 100.0% boost::capy::thread_pool::impl::empty() const :92 40423x 100.0% 100.0% boost::capy::thread_pool::impl::~impl() :109 271x 100.0% 100.0% boost::capy::thread_pool::impl::running_in_this_thread() const :112 454x 100.0% 100.0% boost::capy::thread_pool::impl::drain_abandoned() :123 271x 100.0% 100.0% boost::capy::thread_pool::impl::impl(unsigned long, std::basic_string_view<char, std::char_traits<char> >) :133 271x 100.0% 72.0% boost::capy::thread_pool::impl::post(boost::capy::continuation&) :146 20629x 100.0% 100.0% boost::capy::thread_pool::impl::on_work_started() :157 454x 100.0% 100.0% boost::capy::thread_pool::impl::on_work_finished() :163 454x 100.0% 81.0% boost::capy::thread_pool::impl::join() :183 412x 100.0% 85.0% boost::capy::thread_pool::impl::join()::{lambda()#1}::operator()() const :199 203x 100.0% 100.0% boost::capy::thread_pool::impl::stop() :211 273x 100.0% 100.0% boost::capy::thread_pool::impl::ensure_started() :223 20629x 100.0% 100.0% boost::capy::thread_pool::impl::ensure_started()::{lambda()#1}::operator()() const :225 224x 100.0% 100.0% boost::capy::thread_pool::impl::ensure_started()::{lambda()#1}::operator()() const::{lambda()#1}::operator()() const :228 305x 100.0% 100.0% boost::capy::thread_pool::impl::run(unsigned long) :233 305x 100.0% 78.0% boost::capy::thread_pool::impl::run(unsigned long)::scoped_pool::scoped_pool(boost::capy::thread_pool::impl const*) :244 305x 100.0% 100.0% boost::capy::thread_pool::impl::run(unsigned long)::scoped_pool::~scoped_pool() :245 305x 100.0% 100.0% boost::capy::thread_pool::impl::run(unsigned long)::{lambda()#1}::operator()() const :253 40423x 100.0% 100.0% boost::capy::thread_pool::~thread_pool() :269 271x 100.0% 100.0% boost::capy::thread_pool::thread_pool(unsigned long, std::basic_string_view<char, std::char_traits<char> >) :280 271x 100.0% 55.0% boost::capy::thread_pool::join() :288 141x 100.0% 100.0% boost::capy::thread_pool::stop() :295 2x 100.0% 100.0% boost::capy::thread_pool::get_executor() const :304 11677x 100.0% 100.0% boost::capy::thread_pool::executor_type::on_work_started() const :312 454x 100.0% 100.0% boost::capy::thread_pool::executor_type::on_work_finished() const :319 454x 100.0% 100.0% boost::capy::thread_pool::executor_type::post(boost::capy::continuation&) const :326 20180x 100.0% 100.0% boost::capy::thread_pool::executor_type::dispatch(boost::capy::continuation&) const :333 454x 100.0% 100.0%
Line TLA Hits Source Code
1 //
2 // Copyright (c) 2025 Vinnie Falco (vinnie.falco@gmail.com)
3 // Copyright (c) 2026 Michael Vandeberg
4 //
5 // Distributed under the Boost Software License, Version 1.0. (See accompanying
6 // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7 //
8 // Official repository: https://github.com/boostorg/capy
9 //
10
11 #include <boost/capy/ex/thread_pool.hpp>
12 #include <boost/capy/continuation.hpp>
13 #include <boost/capy/detail/thread_local_ptr.hpp>
14 #include <boost/capy/ex/frame_allocator.hpp>
15 #include <boost/capy/test/thread_name.hpp>
16 #include <algorithm>
17 #include <atomic>
18 #include <condition_variable>
19 #include <cstdio>
20 #include <mutex>
21 #include <thread>
22 #include <vector>
23
24 /*
25 Thread pool implementation using a shared work queue.
26
27 Work items are continuations linked via their intrusive next pointer,
28 stored in a single queue protected by a mutex. No per-post heap
29 allocation: the continuation is owned by the caller and linked
30 directly. Worker threads wait on a condition_variable until work
31 is available or stop is requested.
32
33 Threads are started lazily on first post() via std::call_once to avoid
34 spawning threads for pools that are constructed but never used. Each
35 thread is named with a configurable prefix plus index for debugger
36 visibility.
37
38 Work tracking: on_work_started/on_work_finished maintain the atomic
39 outstanding_work_ counter. on_work_started is lock-free; the worker
40 that drives the count to zero takes mutex_ and re-reads the count
41 before deciding to stop, so the count and the stop decision stay
42 consistent even if work is started in between. join() blocks until
43 this counter reaches zero, then signals workers to stop and joins
44 threads.
45
46 Two shutdown paths:
47 - join(): waits for outstanding work to drain, then stops workers.
48 - stop(): immediately signals workers to exit; queued work is abandoned.
49 - Destructor: stop() then join() (abandon + wait for threads).
50 */
51
52 namespace boost {
53 namespace capy {
54
55 //------------------------------------------------------------------------------
56
57 class thread_pool::impl
58 {
59 // Identifies the pool owning the current worker thread, or
60 // nullptr if the calling thread is not a pool worker. Checked
61 // by dispatch() to decide between symmetric transfer (inline
62 // resume) and post.
63 static inline detail::thread_local_ptr<impl const> current_;
64
65 // Intrusive queue of continuations: the next link is stored in
66 // continuation::reserved (typed continuation* round-tripped through
67 // void*). No per-post allocation: the continuation is owned by the caller.
68 continuation* head_ = nullptr;
69 continuation* tail_ = nullptr;
70
71 20629x void push(continuation* c) noexcept
72 {
73 20629x c->reserved = nullptr;
74 20629x if(tail_)
75 2242x tail_->reserved = c;
76 else
77 18387x head_ = c;
78 20629x tail_ = c;
79 20629x }
80
81 20900x continuation* pop() noexcept
82 {
83 20900x if(!head_)
84 271x return nullptr;
85 20629x continuation* c = head_;
86 20629x head_ = static_cast<continuation*>(head_->reserved);
87 20629x if(!head_)
88 18387x tail_ = nullptr;
89 20629x return c;
90 }
91
92 40423x bool empty() const noexcept
93 {
94 40423x return head_ == nullptr;
95 }
96
97 std::mutex mutex_;
98 std::condition_variable work_cv_;
99 std::condition_variable done_cv_;
100 std::vector<std::thread> threads_;
101 std::atomic<std::size_t> outstanding_work_{0};
102 bool stop_{false};
103 bool joined_{false};
104 std::size_t num_threads_;
105 char thread_name_prefix_[13]{}; // 12 chars max + null terminator
106 std::once_flag start_flag_;
107
108 public:
109 271x ~impl() = default;
110
111 bool
112 454x running_in_this_thread() const noexcept
113 {
114 454x return current_.get() == this;
115 }
116
117 // Destroy abandoned coroutine frames. Must be called
118 // before execution_context::shutdown()/destroy() so
119 // that suspended-frame destructors touching services
120 // (e.g. cancelling registrations) run while those
121 // services are still valid.
122 void
123 271x drain_abandoned() noexcept
124 {
125 477x while(auto* c = pop())
126 {
127 206x auto h = c->h;
128 206x if(h && h != std::noop_coroutine())
129 155x h.destroy();
130 206x }
131 271x }
132
133 271x impl(std::size_t num_threads, std::string_view thread_name_prefix)
134 271x : num_threads_(num_threads)
135 {
136 271x if(num_threads_ == 0)
137 4x num_threads_ = std::max(
138 2x std::thread::hardware_concurrency(), 1u);
139
140 // Truncate prefix to 12 chars, leaving room for up to 3-digit index.
141 271x auto n = thread_name_prefix.copy(thread_name_prefix_, 12);
142 271x thread_name_prefix_[n] = '\0';
143 271x }
144
145 void
146 20629x post(continuation& c)
147 {
148 20629x ensure_started();
149 {
150 20629x std::lock_guard<std::mutex> lock(mutex_);
151 20629x push(&c);
152 20629x }
153 20629x work_cv_.notify_one();
154 20629x }
155
156 void
157 454x on_work_started() noexcept
158 {
159 454x outstanding_work_.fetch_add(1, std::memory_order_acq_rel);
160 454x }
161
162 void
163 454x on_work_finished() noexcept
164 {
165 454x if(outstanding_work_.fetch_sub(
166 454x 1, std::memory_order_acq_rel) == 1)
167 {
168 // fetch_sub's result can be stale: a concurrent
169 // on_work_started() may raise the count before we take the
170 // lock, so re-read it here rather than trust the decrement.
171 202x std::lock_guard<std::mutex> lock(mutex_);
172 202x if(outstanding_work_.load(
173 202x std::memory_order_acquire) == 0 && joined_ && !stop_)
174 {
175 74x stop_ = true;
176 74x done_cv_.notify_all();
177 74x work_cv_.notify_all();
178 }
179 202x }
180 454x }
181
182 void
183 412x join() noexcept
184 {
185 {
186 412x std::unique_lock<std::mutex> lock(mutex_);
187 412x if(joined_)
188 141x return;
189 271x joined_ = true;
190
191 271x if(outstanding_work_.load(
192 271x std::memory_order_acquire) == 0)
193 {
194 143x stop_ = true;
195 143x work_cv_.notify_all();
196 }
197 else
198 {
199 128x done_cv_.wait(lock, [this]{
200 203x return stop_;
201 });
202 }
203 412x }
204
205 576x for(auto& t : threads_)
206 305x if(t.joinable())
207 305x t.join();
208 }
209
210 void
211 273x stop() noexcept
212 {
213 {
214 273x std::lock_guard<std::mutex> lock(mutex_);
215 273x stop_ = true;
216 273x }
217 273x work_cv_.notify_all();
218 273x done_cv_.notify_all();
219 273x }
220
221 private:
222 void
223 20629x ensure_started()
224 {
225 20629x std::call_once(start_flag_, [this]{
226 224x threads_.reserve(num_threads_);
227 529x for(std::size_t i = 0; i < num_threads_; ++i)
228 610x threads_.emplace_back([this, i]{ run(i); });
229 224x });
230 20629x }
231
232 void
233 305x run(std::size_t index)
234 {
235 // Build name; set_current_thread_name truncates to platform limits.
236 char name[16];
237 305x std::snprintf(name, sizeof(name), "%s%zu", thread_name_prefix_, index);
238 305x set_current_thread_name(name);
239
240 // Mark this thread as a worker of this pool so dispatch()
241 // can symmetric-transfer when called from within pool work.
242 struct scoped_pool
243 {
244 305x scoped_pool(impl const* p) noexcept { current_.set(p); }
245 305x ~scoped_pool() noexcept { current_.set(nullptr); }
246 305x } guard(this);
247
248 for(;;)
249 {
250 20728x continuation* c = nullptr;
251 {
252 20728x std::unique_lock<std::mutex> lock(mutex_);
253 20728x work_cv_.wait(lock, [this]{
254 60328x return !empty() ||
255 60328x stop_;
256 });
257 20728x if(stop_)
258 610x return;
259 20423x c = pop();
260 20728x }
261 20423x if(c)
262 20423x safe_resume(c->h);
263 20423x }
264 305x }
265 };
266
267 //------------------------------------------------------------------------------
268
269 271x thread_pool::
270 ~thread_pool()
271 {
272 271x impl_->stop();
273 271x impl_->join();
274 271x impl_->drain_abandoned();
275 271x shutdown();
276 271x destroy();
277 271x delete impl_;
278 271x }
279
280 271x thread_pool::
281 271x thread_pool(std::size_t num_threads, std::string_view thread_name_prefix)
282 271x : impl_(new impl(num_threads, thread_name_prefix))
283 {
284 271x this->set_frame_allocator(std::allocator<void>{});
285 271x }
286
287 void
288 141x thread_pool::
289 join() noexcept
290 {
291 141x impl_->join();
292 141x }
293
294 void
295 2x thread_pool::
296 stop() noexcept
297 {
298 2x impl_->stop();
299 2x }
300
301 //------------------------------------------------------------------------------
302
303 thread_pool::executor_type
304 11677x thread_pool::
305 get_executor() const noexcept
306 {
307 11677x return executor_type(
308 11677x const_cast<thread_pool&>(*this));
309 }
310
311 void
312 454x thread_pool::executor_type::
313 on_work_started() const noexcept
314 {
315 454x pool_->impl_->on_work_started();
316 454x }
317
318 void
319 454x thread_pool::executor_type::
320 on_work_finished() const noexcept
321 {
322 454x pool_->impl_->on_work_finished();
323 454x }
324
325 void
326 20180x thread_pool::executor_type::
327 post(continuation& c) const
328 {
329 20180x pool_->impl_->post(c);
330 20180x }
331
332 std::coroutine_handle<>
333 454x thread_pool::executor_type::
334 dispatch(continuation& c) const
335 {
336 454x if(pool_->impl_->running_in_this_thread())
337 5x return c.h;
338 449x pool_->impl_->post(c);
339 449x return std::noop_coroutine();
340 }
341
342 } // capy
343 } // boost
344