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