TLA Line data 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 HIT 20272 : void push(continuation* c) noexcept
72 : {
73 20272 : c->reserved = nullptr;
74 20272 : if(tail_)
75 6108 : tail_->reserved = c;
76 : else
77 14164 : head_ = c;
78 20272 : tail_ = c;
79 20272 : }
80 :
81 20631 : continuation* pop() noexcept
82 : {
83 20631 : if(!head_)
84 359 : return nullptr;
85 20272 : continuation* c = head_;
86 20272 : head_ = static_cast<continuation*>(head_->reserved);
87 20272 : if(!head_)
88 14164 : tail_ = nullptr;
89 20272 : return c;
90 : }
91 :
92 34728 : bool empty() const noexcept
93 : {
94 34728 : 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 359 : ~impl() = default;
110 :
111 : bool
112 647 : running_in_this_thread() const noexcept
113 : {
114 647 : 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 359 : drain_abandoned() noexcept
124 : {
125 702 : while(auto* c = pop())
126 : {
127 343 : auto h = c->h;
128 343 : if(h && h != std::noop_coroutine())
129 212 : h.destroy();
130 343 : }
131 359 : }
132 :
133 359 : impl(std::size_t num_threads, std::string_view thread_name_prefix)
134 359 : : num_threads_(num_threads)
135 : {
136 359 : if(num_threads_ == 0)
137 8 : num_threads_ = std::max(
138 4 : std::thread::hardware_concurrency(), 1u);
139 :
140 : // Truncate prefix to 12 chars, leaving room for up to 3-digit index.
141 359 : auto n = thread_name_prefix.copy(thread_name_prefix_, 12);
142 359 : thread_name_prefix_[n] = '\0';
143 359 : }
144 :
145 : void
146 20272 : post(continuation& c)
147 : {
148 20272 : ensure_started();
149 : {
150 20272 : std::lock_guard<std::mutex> lock(mutex_);
151 20272 : push(&c);
152 : // Under the lock so the pool cannot drain, join, and
153 : // destroy the condition variable mid-signal.
154 20272 : work_cv_.notify_one();
155 20272 : }
156 20272 : }
157 :
158 : void
159 646 : on_work_started() noexcept
160 : {
161 646 : outstanding_work_.fetch_add(1, std::memory_order_acq_rel);
162 646 : }
163 :
164 : void
165 646 : on_work_finished() noexcept
166 : {
167 646 : if(outstanding_work_.fetch_sub(
168 646 : 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 283 : std::lock_guard<std::mutex> lock(mutex_);
174 283 : if(outstanding_work_.load(
175 283 : std::memory_order_acquire) == 0 && joined_ && !stop_)
176 : {
177 111 : stop_ = true;
178 111 : done_cv_.notify_all();
179 111 : work_cv_.notify_all();
180 : }
181 283 : }
182 646 : }
183 :
184 : void
185 584 : join() noexcept
186 : {
187 : {
188 584 : std::unique_lock<std::mutex> lock(mutex_);
189 584 : if(joined_)
190 225 : return;
191 359 : joined_ = true;
192 :
193 359 : if(outstanding_work_.load(
194 359 : std::memory_order_acquire) == 0)
195 : {
196 196 : stop_ = true;
197 196 : work_cv_.notify_all();
198 : }
199 : else
200 : {
201 163 : done_cv_.wait(lock, [this]{
202 275 : return stop_;
203 : });
204 : }
205 584 : }
206 :
207 772 : for(auto& t : threads_)
208 413 : if(t.joinable())
209 413 : t.join();
210 : }
211 :
212 : void
213 361 : stop() noexcept
214 : {
215 : {
216 361 : std::lock_guard<std::mutex> lock(mutex_);
217 361 : stop_ = true;
218 361 : }
219 361 : work_cv_.notify_all();
220 361 : done_cv_.notify_all();
221 361 : }
222 :
223 : private:
224 : void
225 20272 : ensure_started()
226 : {
227 20272 : std::call_once(start_flag_, [this]{
228 309 : threads_.reserve(num_threads_);
229 722 : for(std::size_t i = 0; i < num_threads_; ++i)
230 826 : threads_.emplace_back([this, i]{ run(i); });
231 309 : });
232 20272 : }
233 :
234 : void
235 413 : run(std::size_t index)
236 : {
237 : // Build name; set_current_thread_name truncates to platform limits.
238 : char name[16];
239 413 : std::snprintf(name, sizeof(name), "%s%zu", thread_name_prefix_, index);
240 413 : 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 413 : scoped_pool(impl const* p) noexcept { current_.set(p); }
247 413 : ~scoped_pool() noexcept { current_.set(nullptr); }
248 413 : } guard(this);
249 :
250 : for(;;)
251 : {
252 20342 : continuation* c = nullptr;
253 : {
254 20342 : std::unique_lock<std::mutex> lock(mutex_);
255 20342 : work_cv_.wait(lock, [this]{
256 49425 : return !empty() ||
257 49425 : stop_;
258 : });
259 20342 : if(stop_)
260 826 : return;
261 19929 : c = pop();
262 20342 : }
263 19929 : if(c)
264 19929 : safe_resume(c->h);
265 19929 : }
266 413 : }
267 : };
268 :
269 : //------------------------------------------------------------------------------
270 :
271 359 : thread_pool::
272 : ~thread_pool()
273 : {
274 359 : impl_->stop();
275 359 : impl_->join();
276 359 : impl_->drain_abandoned();
277 359 : shutdown();
278 359 : destroy();
279 359 : delete impl_;
280 359 : }
281 :
282 359 : thread_pool::
283 359 : thread_pool(std::size_t num_threads, std::string_view thread_name_prefix)
284 359 : : impl_(new impl(num_threads, thread_name_prefix))
285 : {
286 359 : this->set_frame_allocator(std::allocator<void>{});
287 359 : }
288 :
289 : void
290 225 : thread_pool::
291 : join() noexcept
292 : {
293 225 : impl_->join();
294 225 : }
295 :
296 : void
297 2 : thread_pool::
298 : stop() noexcept
299 : {
300 2 : impl_->stop();
301 2 : }
302 :
303 : //------------------------------------------------------------------------------
304 :
305 : thread_pool::executor_type
306 11771 : thread_pool::
307 : get_executor() const noexcept
308 : {
309 11771 : return executor_type(
310 11771 : const_cast<thread_pool&>(*this));
311 : }
312 :
313 : void
314 646 : thread_pool::executor_type::
315 : on_work_started() const noexcept
316 : {
317 646 : pool_->impl_->on_work_started();
318 646 : }
319 :
320 : void
321 646 : thread_pool::executor_type::
322 : on_work_finished() const noexcept
323 : {
324 646 : pool_->impl_->on_work_finished();
325 646 : }
326 :
327 : void
328 19632 : thread_pool::executor_type::
329 : post(continuation& c) const
330 : {
331 19632 : pool_->impl_->post(c);
332 19632 : }
333 :
334 : std::coroutine_handle<>
335 647 : thread_pool::executor_type::
336 : dispatch(continuation& c) const
337 : {
338 647 : if(pool_->impl_->running_in_this_thread())
339 7 : return c.h;
340 640 : pool_->impl_->post(c);
341 640 : return std::noop_coroutine();
342 : }
343 :
344 : } // capy
345 : } // boost
|