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