src/ex/thread_pool.cpp

100.0% Lines (140/0/140) 100.0% List of functions (29/0/29)
thread_pool.cpp
f(x) Functions (29)
Function Calls Lines Blocks
boost::capy::thread_pool::impl::push(boost::capy::continuation*) :71 20136x 100.0% 100.0% boost::capy::thread_pool::impl::pop() :81 20540x 100.0% 100.0% boost::capy::thread_pool::impl::empty() const :92 39383x 100.0% 100.0% boost::capy::thread_pool::impl::~impl() :109 404x 100.0% 100.0% boost::capy::thread_pool::impl::running_in_this_thread() const :112 692x 100.0% 100.0% boost::capy::thread_pool::impl::drain_abandoned() :123 404x 100.0% 100.0% boost::capy::thread_pool::impl::impl(unsigned long, std::basic_string_view<char, std::char_traits<char> >) :133 404x 100.0% 72.0% boost::capy::thread_pool::impl::post(boost::capy::continuation&) :146 20136x 100.0% 100.0% boost::capy::thread_pool::impl::on_work_started() :159 691x 100.0% 100.0% boost::capy::thread_pool::impl::on_work_finished() :165 691x 100.0% 81.0% boost::capy::thread_pool::impl::join() :185 674x 100.0% 85.0% boost::capy::thread_pool::impl::join()::{lambda()#1}::operator()() const :201 425x 100.0% 100.0% boost::capy::thread_pool::impl::stop() :213 406x 100.0% 100.0% boost::capy::thread_pool::impl::ensure_started() :225 20136x 100.0% 100.0% boost::capy::thread_pool::impl::ensure_started()::{lambda()#1}::operator()() const :227 354x 100.0% 100.0% boost::capy::thread_pool::impl::ensure_started()::{lambda()#1}::operator()() const::{lambda()#1}::operator()() const :230 458x 100.0% 100.0% boost::capy::thread_pool::impl::run(unsigned long) :235 458x 100.0% 78.0% boost::capy::thread_pool::impl::run(unsigned long)::scoped_pool::scoped_pool(boost::capy::thread_pool::impl const*) :246 458x 100.0% 100.0% boost::capy::thread_pool::impl::run(unsigned long)::scoped_pool::~scoped_pool() :247 458x 100.0% 100.0% boost::capy::thread_pool::impl::run(unsigned long)::{lambda()#1}::operator()() const :255 39383x 100.0% 100.0% boost::capy::thread_pool::~thread_pool() :271 404x 100.0% 100.0% boost::capy::thread_pool::thread_pool(unsigned long, std::basic_string_view<char, std::char_traits<char> >) :282 404x 100.0% 55.0% boost::capy::thread_pool::join() :290 270x 100.0% 100.0% boost::capy::thread_pool::stop() :297 2x 100.0% 100.0% boost::capy::thread_pool::get_executor() const :306 11816x 100.0% 100.0% boost::capy::thread_pool::executor_type::on_work_started() const :314 691x 100.0% 100.0% boost::capy::thread_pool::executor_type::on_work_finished() const :321 691x 100.0% 100.0% boost::capy::thread_pool::executor_type::post(boost::capy::continuation&) const :328 19451x 100.0% 100.0% boost::capy::thread_pool::executor_type::dispatch(boost::capy::continuation&) const :335 692x 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 20136x void push(continuation* c) noexcept
72 {
73 20136x c->reserved = nullptr;
74 20136x if(tail_)
75 1528x tail_->reserved = c;
76 else
77 18608x head_ = c;
78 20136x tail_ = c;
79 20136x }
80
81 20540x continuation* pop() noexcept
82 {
83 20540x if(!head_)
84 404x return nullptr;
85 20136x continuation* c = head_;
86 20136x head_ = static_cast<continuation*>(head_->reserved);
87 20136x if(!head_)
88 18608x tail_ = nullptr;
89 20136x return c;
90 }
91
92 39383x bool empty() const noexcept
93 {
94 39383x 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 404x ~impl() = default;
110
111 bool
112 692x running_in_this_thread() const noexcept
113 {
114 692x 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 404x drain_abandoned() noexcept
124 {
125 614x while(auto* c = pop())
126 {
127 210x auto h = c->h;
128 210x if(h && h != std::noop_coroutine())
129 159x h.destroy();
130 210x }
131 404x }
132
133 404x impl(std::size_t num_threads, std::string_view thread_name_prefix)
134 404x : num_threads_(num_threads)
135 {
136 404x if(num_threads_ == 0)
137 8x num_threads_ = std::max(
138 4x std::thread::hardware_concurrency(), 1u);
139
140 // Truncate prefix to 12 chars, leaving room for up to 3-digit index.
141 404x auto n = thread_name_prefix.copy(thread_name_prefix_, 12);
142 404x thread_name_prefix_[n] = '\0';
143 404x }
144
145 void
146 20136x post(continuation& c)
147 {
148 20136x ensure_started();
149 {
150 20136x std::lock_guard<std::mutex> lock(mutex_);
151 20136x push(&c);
152 // Under the lock so the pool cannot drain, join, and
153 // destroy the condition variable mid-signal.
154 20136x work_cv_.notify_one();
155 20136x }
156 20136x }
157
158 void
159 691x on_work_started() noexcept
160 {
161 691x outstanding_work_.fetch_add(1, std::memory_order_acq_rel);
162 691x }
163
164 void
165 691x on_work_finished() noexcept
166 {
167 691x if(outstanding_work_.fetch_sub(
168 691x 1, std::memory_order_acq_rel) == 1)
169 {
170 // fetch_sub's result can be stale: a concurrent
171 // on_work_started() may raise the count before we take the
172 // lock, so re-read it here rather than trust the decrement.
173 328x std::lock_guard<std::mutex> lock(mutex_);
174 328x if(outstanding_work_.load(
175 328x std::memory_order_acquire) == 0 && joined_ && !stop_)
176 {
177 186x stop_ = true;
178 186x done_cv_.notify_all();
179 186x work_cv_.notify_all();
180 }
181 328x }
182 691x }
183
184 void
185 674x join() noexcept
186 {
187 {
188 674x std::unique_lock<std::mutex> lock(mutex_);
189 674x if(joined_)
190 270x return;
191 404x joined_ = true;
192
193 404x if(outstanding_work_.load(
194 404x std::memory_order_acquire) == 0)
195 {
196 166x stop_ = true;
197 166x work_cv_.notify_all();
198 }
199 else
200 {
201 238x done_cv_.wait(lock, [this]{
202 425x return stop_;
203 });
204 }
205 674x }
206
207 862x for(auto& t : threads_)
208 458x if(t.joinable())
209 458x t.join();
210 }
211
212 void
213 406x stop() noexcept
214 {
215 {
216 406x std::lock_guard<std::mutex> lock(mutex_);
217 406x stop_ = true;
218 406x }
219 406x work_cv_.notify_all();
220 406x done_cv_.notify_all();
221 406x }
222
223 private:
224 void
225 20136x ensure_started()
226 {
227 20136x std::call_once(start_flag_, [this]{
228 354x threads_.reserve(num_threads_);
229 812x for(std::size_t i = 0; i < num_threads_; ++i)
230 916x threads_.emplace_back([this, i]{ run(i); });
231 354x });
232 20136x }
233
234 void
235 458x run(std::size_t index)
236 {
237 // Build name; set_current_thread_name truncates to platform limits.
238 char name[16];
239 458x std::snprintf(name, sizeof(name), "%s%zu", thread_name_prefix_, index);
240 458x set_current_thread_name(name);
241
242 // Mark this thread as a worker of this pool so dispatch()
243 // can symmetric-transfer when called from within pool work.
244 struct scoped_pool
245 {
246 458x scoped_pool(impl const* p) noexcept { current_.set(p); }
247 458x ~scoped_pool() noexcept { current_.set(nullptr); }
248 458x } guard(this);
249
250 for(;;)
251 {
252 20384x continuation* c = nullptr;
253 {
254 20384x std::unique_lock<std::mutex> lock(mutex_);
255 20384x work_cv_.wait(lock, [this]{
256 58749x return !empty() ||
257 58749x stop_;
258 });
259 20384x if(stop_)
260 916x return;
261 19926x c = pop();
262 20384x }
263 19926x if(c)
264 19926x safe_resume(c->h);
265 19926x }
266 458x }
267 };
268
269 //------------------------------------------------------------------------------
270
271 404x thread_pool::
272 ~thread_pool()
273 {
274 404x impl_->stop();
275 404x impl_->join();
276 404x impl_->drain_abandoned();
277 404x shutdown();
278 404x destroy();
279 404x delete impl_;
280 404x }
281
282 404x thread_pool::
283 404x thread_pool(std::size_t num_threads, std::string_view thread_name_prefix)
284 404x : impl_(new impl(num_threads, thread_name_prefix))
285 {
286 404x this->set_frame_allocator(std::allocator<void>{});
287 404x }
288
289 void
290 270x thread_pool::
291 join() noexcept
292 {
293 270x impl_->join();
294 270x }
295
296 void
297 2x thread_pool::
298 stop() noexcept
299 {
300 2x impl_->stop();
301 2x }
302
303 //------------------------------------------------------------------------------
304
305 thread_pool::executor_type
306 11816x thread_pool::
307 get_executor() const noexcept
308 {
309 11816x return executor_type(
310 11816x const_cast<thread_pool&>(*this));
311 }
312
313 void
314 691x thread_pool::executor_type::
315 on_work_started() const noexcept
316 {
317 691x pool_->impl_->on_work_started();
318 691x }
319
320 void
321 691x thread_pool::executor_type::
322 on_work_finished() const noexcept
323 {
324 691x pool_->impl_->on_work_finished();
325 691x }
326
327 void
328 19451x thread_pool::executor_type::
329 post(continuation& c) const
330 {
331 19451x pool_->impl_->post(c);
332 19451x }
333
334 std::coroutine_handle<>
335 692x thread_pool::executor_type::
336 dispatch(continuation& c) const
337 {
338 692x if(pool_->impl_->running_in_this_thread())
339 7x return c.h;
340 685x pool_->impl_->post(c);
341 685x return std::noop_coroutine();
342 }
343
344 } // capy
345 } // boost
346