include/boost/corosio/detail/thread_pool.hpp

97.4% Lines (74/0/76) 100.0% List of functions (11/0/11)
thread_pool.hpp
f(x) Functions (11)
Line TLA Hits Source Code
1 //
2 // Copyright (c) 2026 Steve Gerbino
3 //
4 // Distributed under the Boost Software License, Version 1.0. (See accompanying
5 // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
6 //
7 // Official repository: https://github.com/cppalliance/corosio
8 //
9
10 #ifndef BOOST_COROSIO_DETAIL_THREAD_POOL_HPP
11 #define BOOST_COROSIO_DETAIL_THREAD_POOL_HPP
12
13 #include <boost/corosio/detail/config.hpp>
14 #include <boost/corosio/detail/intrusive.hpp>
15 #include <boost/capy/error.hpp>
16 #include <boost/capy/ex/execution_context.hpp>
17 #include <boost/capy/test/thread_name.hpp>
18
19 #include <atomic>
20 #include <condition_variable>
21 #include <cstdio>
22 #include <mutex>
23 #include <stdexcept>
24 #include <system_error>
25 #include <thread>
26 #include <vector>
27
28 namespace boost::corosio::detail {
29
30 /** Base class for thread pool work items.
31
32 Derive from this to create work that can be posted to a
33 @ref thread_pool. Uses static function pointer dispatch,
34 consistent with the IOCP `op` pattern.
35
36 @par Example
37 @code
38 struct my_work : pool_work_item
39 {
40 int* result;
41 static void execute( pool_work_item* w ) noexcept
42 {
43 auto* self = static_cast<my_work*>( w );
44 *self->result = 42;
45 }
46 };
47
48 my_work w;
49 w.func_ = &my_work::execute;
50 w.result = &r;
51 auto ec = pool.post( &w );
52 @endcode
53 */
54 struct pool_work_item : intrusive_queue<pool_work_item>::node
55 {
56 /// Static dispatch function signature.
57 using func_type = void (*)(pool_work_item*) noexcept;
58
59 /// Completion handler invoked by the worker thread.
60 func_type func_ = nullptr;
61 };
62
63 /** Shared thread pool for dispatching blocking operations.
64
65 Provides a fixed pool of reusable worker threads for operations
66 that cannot be integrated with async I/O (e.g. blocking DNS
67 calls). Registered as an `execution_context::service` so it
68 is a singleton per io_context.
69
70 The service is created with its context, but the workers start on
71 the first `post()`: a context that never opens a file and never
72 resolves a name never pays for a thread. The default thread count
73 is 1.
74
75 @par Thread Safety
76 All public member functions are thread-safe.
77
78 @par Shutdown
79 Sets a shutdown flag, notifies all threads, and joins them.
80 In-flight blocking calls complete naturally before the thread
81 exits.
82
83 @note Create this service after the scheduler its work items post
84 completions to. Services shut down newest first, so a pool created
85 earlier joins its workers only after the scheduler has drained its
86 completion queue, and the completion the last worker posts is then
87 neither run nor destroyed.
88
89 @note The type is symbol-visible because services are keyed by type
90 identity: with RTTI, hidden behind a shared library boundary, a
91 module that asks for the pool would look up, and create, one of its
92 own (the no-RTTI key is a template static whose visibility follows
93 the template it is instantiated from).
94 */
95 class BOOST_COROSIO_SYMBOL_VISIBLE thread_pool final
96 : public capy::execution_context::service
97 {
98 std::mutex mutex_;
99 std::condition_variable cv_;
100 intrusive_queue<pool_work_item> work_queue_;
101 std::vector<std::thread> threads_;
102 unsigned num_threads_;
103 bool shutdown_ = false;
104
105 void worker_loop(unsigned index);
106 std::error_code start_workers() noexcept;
107
108 public:
109 using key_type = thread_pool;
110
111 /** Construct the thread pool service.
112
113 Records the worker count. The workers themselves start on the
114 first `post()`.
115
116 @par Exception Safety
117 Strong guarantee.
118
119 @param ctx Reference to the owning execution_context.
120 @param num_threads Number of worker threads. Must be
121 at least 1.
122
123 @throws std::logic_error If `num_threads` is 0.
124 */
125 2097x explicit thread_pool(
126 [[maybe_unused]] capy::execution_context& ctx,
127 unsigned num_threads = 1)
128 2097x : num_threads_(num_threads)
129 {
130 2097x if (!num_threads)
131 1x throw std::logic_error("thread_pool requires at least 1 thread");
132 2099x }
133
134 /** Destroy the pool, joining any worker `shutdown()` never reached.
135
136 The context's shutdown walk is the normal path; this only
137 catches a pool created after that walk, whose `shutdown()` is
138 therefore never called and whose joinable threads would
139 otherwise terminate the process. A pool that was never posted
140 to holds no thread and needs neither.
141 */
142 4191x ~thread_pool() override
143 2096x {
144 2096x if (!threads_.empty())
145 shutdown();
146 4191x }
147
148 thread_pool(thread_pool const&) = delete;
149 thread_pool& operator=(thread_pool const&) = delete;
150
151 /** Enqueue a work item for execution on the thread pool.
152
153 The first item posted starts the workers. Zero-allocation:
154 the caller owns the work item's storage.
155
156 A refusal answers with the code the caller reports for the
157 operation it was starting, so that a system that will not give
158 the pool a thread is not mistaken for a cancellation.
159
160 @par Thread Safety
161 Safe. Racing first posts start the workers once.
162
163 @param w The work item to execute. Must remain valid until
164 its `func_` has been called.
165
166 @return An empty code if the item was enqueued;
167 `capy::error::canceled` if the pool has already shut
168 down; otherwise the code of the thread the system
169 refused, which left the pool with no worker at all.
170 */
171 [[nodiscard]] std::error_code post(pool_work_item* w) noexcept;
172
173 /** Return the number of workers the pool has started.
174
175 Zero until the first `post()`, and zero again once
176 `shutdown()` has joined them.
177
178 @par Thread Safety
179 Safe.
180 */
181 5x unsigned worker_count() noexcept
182 {
183 5x std::lock_guard<std::mutex> lock(mutex_);
184 5x return static_cast<unsigned>(threads_.size());
185 5x }
186
187 /** Shut down the thread pool.
188
189 Signals all threads to exit after draining any
190 remaining queued work, then joins them.
191 */
192 void shutdown() override;
193 };
194
195 inline void
196 180x thread_pool::worker_loop(unsigned index)
197 {
198 // Name format chosen to fit Linux's 15-char pthread limit:
199 // "tpool-svc-" (10) + up to 4 digit index leaves "tpool-svc-9999".
200 char name[16];
201 180x std::snprintf(name, sizeof(name), "tpool-svc-%u", index);
202 180x capy::set_current_thread_name(name);
203
204 for (;;)
205 {
206 pool_work_item* w;
207 {
208 685x std::unique_lock<std::mutex> lock(mutex_);
209 685x cv_.wait(
210 897x lock, [this] { return shutdown_ || !work_queue_.empty(); });
211
212 685x w = work_queue_.pop();
213 685x if (!w)
214 {
215 180x if (shutdown_)
216 360x return;
217 continue;
218 }
219 685x }
220 505x w->func_(w);
221 505x }
222 }
223
224 // Called with mutex_ held, so the workers are started once however
225 // many threads race the first post.
226 inline std::error_code
227 509x thread_pool::start_workers() noexcept
228 {
229 509x if (!threads_.empty())
230 328x return {};
231 181x std::error_code ec;
232 try
233 {
234 181x threads_.reserve(num_threads_);
235 360x for (unsigned i = 0; i < num_threads_; ++i)
236 363x threads_.emplace_back([this, i] { worker_loop(i + 1); });
237 }
238 4x catch (std::system_error const& e)
239 {
240 // The refusal is carried out, not swallowed: a thread the
241 // system will not give is a real error and the operation that
242 // asked for it says so, rather than reporting the cancellation
243 // that belongs to a stop token.
244 2x ec = e.code();
245 2x }
246 2x catch (...)
247 {
248 2x ec = std::make_error_code(std::errc::resource_unavailable_try_again);
249 2x }
250 // A pool short of workers still runs everything posted to it, only
251 // less of it at once, so a partial start is a start. What it does
252 // not do is come back for the rest: the size is a tuning knob, and
253 // topping it up would put a thread creation on the initiator's
254 // path for every operation after a refusal.
255 181x if (!threads_.empty())
256 177x return {};
257 4x return ec;
258 }
259
260 inline std::error_code
261 520x thread_pool::post(pool_work_item* w) noexcept
262 {
263 {
264 520x std::lock_guard<std::mutex> lock(mutex_);
265 520x if (shutdown_)
266 11x return capy::error::canceled;
267 // The system can refuse a thread, and an initiator has no way
268 // to throw; a refused post is the failure the callers already
269 // report through the operation they were starting.
270 509x if (auto ec = start_workers())
271 4x return ec;
272 505x work_queue_.push(w);
273 520x }
274 505x cv_.notify_one();
275 505x return {};
276 }
277
278 inline void
279 2106x thread_pool::shutdown()
280 {
281 {
282 2106x std::lock_guard<std::mutex> lock(mutex_);
283 2106x shutdown_ = true;
284 2106x }
285 2106x cv_.notify_all();
286
287 // Unlocked, though a post may add to threads_: the flag above is
288 // published under the same mutex, so a post that has not taken it
289 // yet will find it set and start nothing, and one already inside
290 // released the mutex before this thread acquired it.
291 2286x for (auto& t : threads_)
292 {
293 180x if (t.joinable())
294 180x t.join();
295 }
296 2106x threads_.clear();
297
298 {
299 2106x std::lock_guard<std::mutex> lock(mutex_);
300 2106x while (work_queue_.pop())
301 ;
302 2106x }
303 2106x }
304
305 /** A reference to the context's shared thread pool, bound on first use.
306
307 Services that hand blocking work to the pool hold one of these
308 instead of a reference bound at construction. They are constructed
309 from the scheduler's constructor, where the pool they created would
310 be older than the scheduler and would join too late; binding on
311 first use puts the pool after it instead.
312
313 The owning `io_context` creates the pool service during
314 construction, so by the time any operation can run the binding only
315 ever finds it. That is what keeps `get()` from constructing
316 anything on an initiator's thread, and so from throwing where an
317 initiator may not: the throwing spelling exists for a scheduler
318 driven without an `io_context`. What the service defers is its
319 workers, and those are started by `post()`, which reports a refusal
320 rather than throwing it.
321
322 @par Thread Safety
323 Distinct objects: Safe.
324 Shared objects: Safe.
325
326 @see thread_pool
327 */
328 class thread_pool_ref
329 {
330 capy::execution_context& ctx_;
331 std::atomic<thread_pool*> pool_{nullptr};
332
333 public:
334 /** Construct a reference into the given context.
335
336 @param ctx The context whose pool is used.
337 */
338 6285x explicit thread_pool_ref(capy::execution_context& ctx) noexcept
339 6285x : ctx_(ctx)
340 {
341 6285x }
342
343 thread_pool_ref(thread_pool_ref const&) = delete;
344 thread_pool_ref& operator=(thread_pool_ref const&) = delete;
345
346 /** Return the pool, creating it if this is the first use.
347
348 @par Preconditions
349 For the throwing clauses below to be unreachable, the owning
350 context must already hold the pool service. Every `io_context`
351 constructor installs it — what waits for a first post is the
352 service's workers, not the service — so the creating branch is
353 reached only by a scheduler driven without one.
354
355 @par Exception Safety
356 Strong guarantee.
357
358 @throws std::bad_alloc If the service cannot be allocated.
359
360 @throws std::logic_error If the pool is asked for zero threads.
361
362 @return The context's shared thread pool.
363 */
364 497x thread_pool& get()
365 {
366 497x auto* p = pool_.load(std::memory_order_acquire);
367 497x if (!p)
368 {
369 182x p = &ctx_.use_service<thread_pool>();
370 182x pool_.store(p, std::memory_order_release);
371 }
372 497x return *p;
373 }
374 };
375
376 } // namespace boost::corosio::detail
377
378 #endif // BOOST_COROSIO_DETAIL_THREAD_POOL_HPP
379