include/boost/corosio/detail/thread_pool.hpp

86.8% Lines (66/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 1792x explicit thread_pool(
126 [[maybe_unused]] capy::execution_context& ctx,
127 unsigned num_threads = 1)
128 1792x : num_threads_(num_threads)
129 {
130 1792x if (!num_threads)
131 1x throw std::logic_error("thread_pool requires at least 1 thread");
132 1794x }
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 3581x ~thread_pool() override
143 1791x {
144 1791x if (!threads_.empty())
145 shutdown();
146 3581x }
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 136x 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 136x std::snprintf(name, sizeof(name), "tpool-svc-%u", index);
202 136x capy::set_current_thread_name(name);
203
204 for (;;)
205 {
206 pool_work_item* w;
207 {
208 591x std::unique_lock<std::mutex> lock(mutex_);
209 591x cv_.wait(
210 757x lock, [this] { return shutdown_ || !work_queue_.empty(); });
211
212 591x w = work_queue_.pop();
213 591x if (!w)
214 {
215 136x if (shutdown_)
216 272x return;
217 continue;
218 }
219 591x }
220 455x w->func_(w);
221 455x }
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 455x thread_pool::start_workers() noexcept
228 {
229 455x if (!threads_.empty())
230 322x return {};
231 133x std::error_code ec;
232 try
233 {
234 133x threads_.reserve(num_threads_);
235 269x for (unsigned i = 0; i < num_threads_; ++i)
236 272x threads_.emplace_back([this, i] { worker_loop(i + 1); });
237 }
238 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 ec = e.code();
245 }
246 catch (...)
247 {
248 ec = std::make_error_code(std::errc::resource_unavailable_try_again);
249 }
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 133x if (!threads_.empty())
256 133x return {};
257 return ec;
258 }
259
260 inline std::error_code
261 456x thread_pool::post(pool_work_item* w) noexcept
262 {
263 {
264 456x std::lock_guard<std::mutex> lock(mutex_);
265 456x if (shutdown_)
266 1x 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 455x if (auto ec = start_workers())
271 return ec;
272 455x work_queue_.push(w);
273 456x }
274 455x cv_.notify_one();
275 455x return {};
276 }
277
278 inline void
279 1796x thread_pool::shutdown()
280 {
281 {
282 1796x std::lock_guard<std::mutex> lock(mutex_);
283 1796x shutdown_ = true;
284 1796x }
285 1796x 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 1932x for (auto& t : threads_)
292 {
293 136x if (t.joinable())
294 136x t.join();
295 }
296 1796x threads_.clear();
297
298 {
299 1796x std::lock_guard<std::mutex> lock(mutex_);
300 1796x while (work_queue_.pop())
301 ;
302 1796x }
303 1796x }
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 5370x explicit thread_pool_ref(capy::execution_context& ctx) noexcept
339 5370x : ctx_(ctx)
340 {
341 5370x }
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 433x thread_pool& get()
365 {
366 433x auto* p = pool_.load(std::memory_order_acquire);
367 433x if (!p)
368 {
369 129x p = &ctx_.use_service<thread_pool>();
370 129x pool_.store(p, std::memory_order_release);
371 }
372 433x return *p;
373 }
374 };
375
376 } // namespace boost::corosio::detail
377
378 #endif // BOOST_COROSIO_DETAIL_THREAD_POOL_HPP
379