src/ex/thread_pool.cpp

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