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/cppalliance/capy
9 : //
10 :
11 : #ifndef BOOST_CAPY_EX_THREAD_POOL_HPP
12 : #define BOOST_CAPY_EX_THREAD_POOL_HPP
13 :
14 : #include <boost/capy/detail/config.hpp>
15 : #include <boost/capy/continuation.hpp>
16 : #include <coroutine>
17 : #include <boost/capy/ex/execution_context.hpp>
18 : #include <cstddef>
19 : #include <string_view>
20 :
21 : namespace boost {
22 : namespace capy {
23 :
24 : /** A pool of threads for executing work concurrently.
25 :
26 : Use this when you need to run coroutines on multiple threads
27 : without the overhead of creating and destroying threads for
28 : each task. Work items are distributed across the pool using
29 : a shared queue.
30 :
31 : @par Thread Safety
32 : Distinct objects: Safe.
33 : Shared objects: Unsafe.
34 :
35 : @par Example
36 : @code
37 : thread_pool pool(4); // 4 worker threads
38 : auto ex = pool.get_executor();
39 : run_async(ex)(some_task()); // launch work; tracked so join() waits for it
40 : pool.join(); // wait for outstanding work to complete
41 : // pool destructor stops the pool, discarding any pending work
42 : @endcode
43 :
44 : @note `join()` waits only for work that holds outstanding-work
45 : counting, which `run_async` (and `make_work_guard`) provide. A bare
46 : `executor_type::post()` does not register outstanding work, so
47 : `join()` will not wait for it.
48 : */
49 : class BOOST_CAPY_DECL
50 : thread_pool
51 : : public execution_context
52 : {
53 : class impl;
54 : impl* impl_;
55 :
56 : public:
57 : class executor_type;
58 :
59 : /** Destroy the thread pool.
60 :
61 : Signals all worker threads to stop, waits for them to
62 : finish, and destroys any pending work items.
63 :
64 : @par Preconditions
65 : No thread outside this pool may post or dispatch work to it
66 : (or to a strand built on it) concurrently with, or after,
67 : destruction; doing so is undefined behavior. Submit such work
68 : through @ref run_async or @ref run and call @ref join before
69 : the pool is destroyed, so it has completed first.
70 : */
71 : ~thread_pool();
72 :
73 : /** Construct a thread pool.
74 :
75 : Creates a pool with the specified number of worker threads.
76 : If `num_threads` is zero, the number of threads is set to
77 : the hardware concurrency, or one if that cannot be determined.
78 :
79 : @param num_threads The number of worker threads, or zero
80 : for automatic selection.
81 :
82 : @param thread_name_prefix The prefix for worker thread names.
83 : Thread names appear as "{prefix}0", "{prefix}1", etc.
84 : The prefix is truncated to 12 characters. Defaults to
85 : "capy-pool-".
86 : */
87 : explicit
88 : thread_pool(
89 : std::size_t num_threads = 0,
90 : std::string_view thread_name_prefix = "capy-pool-");
91 :
92 : thread_pool(thread_pool const&) = delete;
93 : thread_pool& operator=(thread_pool const&) = delete;
94 :
95 : /** Wait for all outstanding work to complete.
96 :
97 : Releases the internal work guard, then blocks the calling
98 : thread until all outstanding work tracked by
99 : @ref executor_type::on_work_started and
100 : @ref executor_type::on_work_finished completes. After all
101 : work finishes, joins the worker threads.
102 :
103 : If @ref stop is called while `join()` is blocking, the
104 : pool stops without waiting for remaining work to
105 : complete. Worker threads finish their current item and
106 : exit; `join()` still waits for all threads to be joined
107 : before returning.
108 :
109 : This function is idempotent. The first call performs the
110 : join; subsequent calls return immediately.
111 :
112 : @par Preconditions
113 : Must not be called from a thread in this pool (undefined
114 : behavior).
115 :
116 : @par Postconditions
117 : All worker threads have been joined. The pool cannot be
118 : reused.
119 :
120 : @par Thread Safety
121 : May be called from any thread not in this pool.
122 : */
123 : void
124 : join() noexcept;
125 :
126 : /** Request all worker threads to stop.
127 :
128 : Signals all threads to exit after finishing their current
129 : work item. Queued work that has not started is abandoned.
130 : Does not wait for threads to exit.
131 :
132 : If @ref join is blocking on another thread, calling
133 : `stop()` causes it to stop waiting for outstanding
134 : work. The `join()` call still waits for worker threads
135 : to finish their current item and exit before returning.
136 : */
137 : void
138 : stop() noexcept;
139 :
140 : /** Return an executor for this thread pool.
141 :
142 : @return An executor associated with this thread pool.
143 : */
144 : executor_type
145 : get_executor() const noexcept;
146 : };
147 :
148 : /** An executor that submits work to a thread_pool.
149 :
150 : Executors are lightweight handles that can be copied and stored.
151 : All copies refer to the same underlying thread pool.
152 :
153 : @par Thread Safety
154 : Distinct objects: Safe.
155 : Shared objects: Safe.
156 : */
157 : class thread_pool::executor_type
158 : {
159 : friend class thread_pool;
160 :
161 : thread_pool* pool_ = nullptr;
162 :
163 : explicit
164 HIT 11685 : executor_type(thread_pool& pool) noexcept
165 11685 : : pool_(&pool)
166 : {
167 11685 : }
168 :
169 : public:
170 : /** Construct a default null executor.
171 :
172 : The resulting executor is not associated with any pool.
173 : `context()`, `dispatch()`, and `post()` require the
174 : executor to be associated with a pool before use.
175 : */
176 : executor_type() = default;
177 :
178 : /// Return the underlying thread pool.
179 : thread_pool&
180 11901 : context() const noexcept
181 : {
182 11901 : return *pool_;
183 : }
184 :
185 : /** Notify that work has started.
186 :
187 : Increments the outstanding work count. Must be paired
188 : with a subsequent call to @ref on_work_finished.
189 :
190 : @see on_work_finished, work_guard
191 : */
192 : BOOST_CAPY_DECL
193 : void
194 : on_work_started() const noexcept;
195 :
196 : /** Notify that work has finished.
197 :
198 : Decrements the outstanding work count. When the count
199 : reaches zero after @ref thread_pool::join has been called,
200 : the pool's worker threads are signaled to stop.
201 :
202 : @pre A preceding call to @ref on_work_started was made.
203 :
204 : @see on_work_started, work_guard
205 : */
206 : BOOST_CAPY_DECL
207 : void
208 : on_work_finished() const noexcept;
209 :
210 : /** Dispatch a continuation for execution.
211 :
212 : If the calling thread is a worker of this pool, returns
213 : `c.h` for symmetric transfer so the caller can resume the
214 : continuation inline. Otherwise, posts the continuation to
215 : the pool for execution on a worker thread and returns
216 : `std::noop_coroutine()`.
217 :
218 : @param c The continuation to execute. On the post path,
219 : must remain at a stable address until dequeued
220 : and resumed.
221 :
222 : @return `c.h` when the calling thread is a pool worker;
223 : `std::noop_coroutine()` otherwise.
224 : */
225 : BOOST_CAPY_DECL
226 : std::coroutine_handle<>
227 : dispatch(continuation& c) const;
228 :
229 : /** Post a continuation to the thread pool.
230 :
231 : The continuation will be resumed on one of the pool's
232 : worker threads. The continuation must remain at a stable
233 : address until it is dequeued and resumed.
234 :
235 : @param c The continuation to execute.
236 : */
237 : BOOST_CAPY_DECL
238 : void
239 : post(continuation& c) const;
240 :
241 : /// Return true if two executors refer to the same thread pool.
242 : bool
243 13 : operator==(executor_type const& other) const noexcept
244 : {
245 13 : return pool_ == other.pool_;
246 : }
247 : };
248 :
249 : } // capy
250 : } // boost
251 :
252 : #endif
|