include/boost/corosio/detail/thread_pool.hpp

98.7% Lines (75/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 2109x explicit thread_pool(
126 [[maybe_unused]] capy::execution_context& ctx,
127 unsigned num_threads = 1)
128 2109x : num_threads_(num_threads)
129 {
130 2109x if (!num_threads)
131 1x throw std::logic_error("thread_pool requires at least 1 thread");
132 2111x }
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 4214x ~thread_pool() override
143 2108x {
144 2108x if (!threads_.empty())
145 1x shutdown();
146 4214x }
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 6x unsigned worker_count() noexcept
182 {
183 6x std::lock_guard<std::mutex> lock(mutex_);
184 6x return static_cast<unsigned>(threads_.size());
185 6x }
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 185x 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 185x std::snprintf(name, sizeof(name), "tpool-svc-%u", index);
202 185x capy::set_current_thread_name(name);
203
204 for (;;)
205 {
206 pool_work_item* w;
207 {
208 697x std::unique_lock<std::mutex> lock(mutex_);
209 697x cv_.wait(
210 936x lock, [this] { return shutdown_ || !work_queue_.empty(); });
211
212 697x w = work_queue_.pop();
213 697x if (!w)
214 {
215 185x if (shutdown_)
216 370x return;
217 continue;
218 }
219 697x }
220 512x w->func_(w);
221 512x }
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 516x thread_pool::start_workers() noexcept
228 {
229 516x if (!threads_.empty())
230 330x return {};
231 186x std::error_code ec;
232 try
233 {
234 186x threads_.reserve(num_threads_);
235 370x for (unsigned i = 0; i < num_threads_; ++i)
236 373x 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 186x if (!threads_.empty())
256 182x return {};
257 4x return ec;
258 }
259
260 inline std::error_code
261 527x thread_pool::post(pool_work_item* w) noexcept
262 {
263 {
264 527x std::lock_guard<std::mutex> lock(mutex_);
265 527x 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 516x if (auto ec = start_workers())
271 4x return ec;
272 512x work_queue_.push(w);
273 527x }
274 512x cv_.notify_one();
275 512x return {};
276 }
277
278 inline void
279 2118x thread_pool::shutdown()
280 {
281 {
282 2118x std::lock_guard<std::mutex> lock(mutex_);
283 2118x shutdown_ = true;
284 2118x }
285 2118x 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 2303x for (auto& t : threads_)
292 {
293 185x if (t.joinable())
294 185x t.join();
295 }
296 2118x threads_.clear();
297
298 {
299 2118x std::lock_guard<std::mutex> lock(mutex_);
300 2118x while (work_queue_.pop())
301 ;
302 2118x }
303 2118x }
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 6318x explicit thread_pool_ref(capy::execution_context& ctx) noexcept
339 6318x : ctx_(ctx)
340 {
341 6318x }
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 503x thread_pool& get()
365 {
366 503x auto* p = pool_.load(std::memory_order_acquire);
367 503x if (!p)
368 {
369 186x p = &ctx_.use_service<thread_pool>();
370 186x pool_.store(p, std::memory_order_release);
371 }
372 503x return *p;
373 }
374 };
375
376 } // namespace boost::corosio::detail
377
378 #endif // BOOST_COROSIO_DETAIL_THREAD_POOL_HPP
379