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