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 21048x 100.0% 100.0% boost::capy::thread_pool::impl::pop() :81 21321x 100.0% 100.0% boost::capy::thread_pool::impl::empty() const :92 41318x 100.0% 100.0% boost::capy::thread_pool::impl::~impl() :109 273x 100.0% 100.0% boost::capy::thread_pool::impl::running_in_this_thread() const :112 456x 100.0% 100.0% boost::capy::thread_pool::impl::drain_abandoned() :123 273x 100.0% 100.0% boost::capy::thread_pool::impl::impl(unsigned long, std::basic_string_view<char, std::char_traits<char> >) :133 273x 100.0% 72.0% boost::capy::thread_pool::impl::post(boost::capy::continuation&) :146 21048x 100.0% 100.0% boost::capy::thread_pool::impl::on_work_started() :157 456x 100.0% 100.0% boost::capy::thread_pool::impl::on_work_finished() :163 456x 100.0% 81.0% boost::capy::thread_pool::impl::join() :183 416x 100.0% 85.0% boost::capy::thread_pool::impl::join()::{lambda()#1}::operator()() const :199 202x 100.0% 100.0% boost::capy::thread_pool::impl::stop() :211 275x 100.0% 100.0% boost::capy::thread_pool::impl::ensure_started() :223 21048x 100.0% 100.0% boost::capy::thread_pool::impl::ensure_started()::{lambda()#1}::operator()() const :225 226x 100.0% 100.0% boost::capy::thread_pool::impl::ensure_started()::{lambda()#1}::operator()() const::{lambda()#1}::operator()() const :228 307x 100.0% 100.0% boost::capy::thread_pool::impl::run(unsigned long) :233 307x 100.0% 78.0% boost::capy::thread_pool::impl::run(unsigned long)::scoped_pool::scoped_pool(boost::capy::thread_pool::impl const*) :244 307x 100.0% 100.0% boost::capy::thread_pool::impl::run(unsigned long)::scoped_pool::~scoped_pool() :245 307x 100.0% 100.0% boost::capy::thread_pool::impl::run(unsigned long)::{lambda()#1}::operator()() const :253 41318x 100.0% 100.0% boost::capy::thread_pool::~thread_pool() :269 273x 100.0% 100.0% boost::capy::thread_pool::thread_pool(unsigned long, std::basic_string_view<char, std::char_traits<char> >) :280 273x 100.0% 55.0% boost::capy::thread_pool::join() :288 143x 100.0% 100.0% boost::capy::thread_pool::stop() :295 2x 100.0% 100.0% boost::capy::thread_pool::get_executor() const :304 11679x 100.0% 100.0% boost::capy::thread_pool::executor_type::on_work_started() const :312 456x 100.0% 100.0% boost::capy::thread_pool::executor_type::on_work_finished() const :319 456x 100.0% 100.0% boost::capy::thread_pool::executor_type::post(boost::capy::continuation&) const :326 20597x 100.0% 100.0% boost::capy::thread_pool::executor_type::dispatch(boost::capy::continuation&) const :333 456x 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 21048x void push(continuation* c) noexcept
72 {
73 21048x c->reserved = nullptr;
74 21048x if(tail_)
75 2711x tail_->reserved = c;
76 else
77 18337x head_ = c;
78 21048x tail_ = c;
79 21048x }
80
81 21321x continuation* pop() noexcept
82 {
83 21321x if(!head_)
84 273x return nullptr;
85 21048x continuation* c = head_;
86 21048x head_ = static_cast<continuation*>(head_->reserved);
87 21048x if(!head_)
88 18337x tail_ = nullptr;
89 21048x return c;
90 }
91
92 41318x bool empty() const noexcept
93 {
94 41318x 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 273x ~impl() = default;
110
111 bool
112 456x running_in_this_thread() const noexcept
113 {
114 456x 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 273x drain_abandoned() noexcept
124 {
125 476x while(auto* c = pop())
126 {
127 203x auto h = c->h;
128 203x if(h && h != std::noop_coroutine())
129 152x h.destroy();
130 203x }
131 273x }
132
133 273x impl(std::size_t num_threads, std::string_view thread_name_prefix)
134 273x : num_threads_(num_threads)
135 {
136 273x 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 273x auto n = thread_name_prefix.copy(thread_name_prefix_, 12);
142 273x thread_name_prefix_[n] = '\0';
143 273x }
144
145 void
146 21048x post(continuation& c)
147 {
148 21048x ensure_started();
149 {
150 21048x std::lock_guard<std::mutex> lock(mutex_);
151 21048x push(&c);
152 21048x }
153 21048x work_cv_.notify_one();
154 21048x }
155
156 void
157 456x on_work_started() noexcept
158 {
159 456x outstanding_work_.fetch_add(1, std::memory_order_acq_rel);
160 456x }
161
162 void
163 456x on_work_finished() noexcept
164 {
165 456x if(outstanding_work_.fetch_sub(
166 456x 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 204x std::lock_guard<std::mutex> lock(mutex_);
172 204x if(outstanding_work_.load(
173 204x std::memory_order_acquire) == 0 && joined_ && !stop_)
174 {
175 73x stop_ = true;
176 73x done_cv_.notify_all();
177 73x work_cv_.notify_all();
178 }
179 204x }
180 456x }
181
182 void
183 416x join() noexcept
184 {
185 {
186 416x std::unique_lock<std::mutex> lock(mutex_);
187 416x if(joined_)
188 143x return;
189 273x joined_ = true;
190
191 273x if(outstanding_work_.load(
192 273x std::memory_order_acquire) == 0)
193 {
194 145x stop_ = true;
195 145x work_cv_.notify_all();
196 }
197 else
198 {
199 128x done_cv_.wait(lock, [this]{
200 202x return stop_;
201 });
202 }
203 416x }
204
205 580x for(auto& t : threads_)
206 307x if(t.joinable())
207 307x t.join();
208 }
209
210 void
211 275x stop() noexcept
212 {
213 {
214 275x std::lock_guard<std::mutex> lock(mutex_);
215 275x stop_ = true;
216 275x }
217 275x work_cv_.notify_all();
218 275x done_cv_.notify_all();
219 275x }
220
221 private:
222 void
223 21048x ensure_started()
224 {
225 21048x std::call_once(start_flag_, [this]{
226 226x threads_.reserve(num_threads_);
227 533x for(std::size_t i = 0; i < num_threads_; ++i)
228 614x threads_.emplace_back([this, i]{ run(i); });
229 226x });
230 21048x }
231
232 void
233 307x run(std::size_t index)
234 {
235 // Build name; set_current_thread_name truncates to platform limits.
236 char name[16];
237 307x std::snprintf(name, sizeof(name), "%s%zu", thread_name_prefix_, index);
238 307x 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 307x scoped_pool(impl const* p) noexcept { current_.set(p); }
245 307x ~scoped_pool() noexcept { current_.set(nullptr); }
246 307x } guard(this);
247
248 for(;;)
249 {
250 21152x continuation* c = nullptr;
251 {
252 21152x std::unique_lock<std::mutex> lock(mutex_);
253 21152x work_cv_.wait(lock, [this]{
254 61692x return !empty() ||
255 61692x stop_;
256 });
257 21152x if(stop_)
258 614x return;
259 20845x c = pop();
260 21152x }
261 20845x if(c)
262 20845x safe_resume(c->h);
263 20845x }
264 307x }
265 };
266
267 //------------------------------------------------------------------------------
268
269 273x thread_pool::
270 ~thread_pool()
271 {
272 273x impl_->stop();
273 273x impl_->join();
274 273x impl_->drain_abandoned();
275 273x shutdown();
276 273x destroy();
277 273x delete impl_;
278 273x }
279
280 273x thread_pool::
281 273x thread_pool(std::size_t num_threads, std::string_view thread_name_prefix)
282 273x : impl_(new impl(num_threads, thread_name_prefix))
283 {
284 273x this->set_frame_allocator(std::allocator<void>{});
285 273x }
286
287 void
288 143x thread_pool::
289 join() noexcept
290 {
291 143x impl_->join();
292 143x }
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 11679x thread_pool::
305 get_executor() const noexcept
306 {
307 11679x return executor_type(
308 11679x const_cast<thread_pool&>(*this));
309 }
310
311 void
312 456x thread_pool::executor_type::
313 on_work_started() const noexcept
314 {
315 456x pool_->impl_->on_work_started();
316 456x }
317
318 void
319 456x thread_pool::executor_type::
320 on_work_finished() const noexcept
321 {
322 456x pool_->impl_->on_work_finished();
323 456x }
324
325 void
326 20597x thread_pool::executor_type::
327 post(continuation& c) const
328 {
329 20597x pool_->impl_->post(c);
330 20597x }
331
332 std::coroutine_handle<>
333 456x thread_pool::executor_type::
334 dispatch(continuation& c) const
335 {
336 456x if(pool_->impl_->running_in_this_thread())
337 5x return c.h;
338 451x pool_->impl_->post(c);
339 451x return std::noop_coroutine();
340 }
341
342 } // capy
343 } // boost
344