86.84% Lines (66/76) 100.00% Functions (11/11)
TLA Baseline Branch
Line Hits Code Line Hits Code
1   // 1   //
2   // Copyright (c) 2026 Steve Gerbino 2   // Copyright (c) 2026 Steve Gerbino
3   // 3   //
4   // Distributed under the Boost Software License, Version 1.0. (See accompanying 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) 5   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
6   // 6   //
7   // Official repository: https://github.com/cppalliance/corosio 7   // Official repository: https://github.com/cppalliance/corosio
8   // 8   //
9   9  
10   #ifndef BOOST_COROSIO_DETAIL_THREAD_POOL_HPP 10   #ifndef BOOST_COROSIO_DETAIL_THREAD_POOL_HPP
11   #define BOOST_COROSIO_DETAIL_THREAD_POOL_HPP 11   #define BOOST_COROSIO_DETAIL_THREAD_POOL_HPP
12   12  
13   #include <boost/corosio/detail/config.hpp> 13   #include <boost/corosio/detail/config.hpp>
14   #include <boost/corosio/detail/intrusive.hpp> 14   #include <boost/corosio/detail/intrusive.hpp>
  15 + #include <boost/capy/error.hpp>
15   #include <boost/capy/ex/execution_context.hpp> 16   #include <boost/capy/ex/execution_context.hpp>
16   #include <boost/capy/test/thread_name.hpp> 17   #include <boost/capy/test/thread_name.hpp>
17   18  
  19 + #include <atomic>
18   #include <condition_variable> 20   #include <condition_variable>
19   #include <cstdio> 21   #include <cstdio>
20   #include <mutex> 22   #include <mutex>
21   #include <stdexcept> 23   #include <stdexcept>
  24 + #include <system_error>
22   #include <thread> 25   #include <thread>
23   #include <vector> 26   #include <vector>
24   27  
25   namespace boost::corosio::detail { 28   namespace boost::corosio::detail {
26   29  
27   /** Base class for thread pool work items. 30   /** Base class for thread pool work items.
28   31  
29   Derive from this to create work that can be posted to a 32   Derive from this to create work that can be posted to a
30   @ref thread_pool. Uses static function pointer dispatch, 33   @ref thread_pool. Uses static function pointer dispatch,
31   consistent with the IOCP `op` pattern. 34   consistent with the IOCP `op` pattern.
32   35  
33   @par Example 36   @par Example
34   @code 37   @code
35   struct my_work : pool_work_item 38   struct my_work : pool_work_item
36   { 39   {
37   int* result; 40   int* result;
38   static void execute( pool_work_item* w ) noexcept 41   static void execute( pool_work_item* w ) noexcept
39   { 42   {
40   auto* self = static_cast<my_work*>( w ); 43   auto* self = static_cast<my_work*>( w );
41   *self->result = 42; 44   *self->result = 42;
42   } 45   }
43   }; 46   };
44   47  
45   my_work w; 48   my_work w;
46   w.func_ = &my_work::execute; 49   w.func_ = &my_work::execute;
47   w.result = &r; 50   w.result = &r;
48 - pool.post( &w ); 51 + auto ec = pool.post( &w );
49   @endcode 52   @endcode
50   */ 53   */
51   struct pool_work_item : intrusive_queue<pool_work_item>::node 54   struct pool_work_item : intrusive_queue<pool_work_item>::node
52   { 55   {
53   /// Static dispatch function signature. 56   /// Static dispatch function signature.
54   using func_type = void (*)(pool_work_item*) noexcept; 57   using func_type = void (*)(pool_work_item*) noexcept;
55   58  
56   /// Completion handler invoked by the worker thread. 59   /// Completion handler invoked by the worker thread.
57   func_type func_ = nullptr; 60   func_type func_ = nullptr;
58   }; 61   };
59   62  
60   /** Shared thread pool for dispatching blocking operations. 63   /** Shared thread pool for dispatching blocking operations.
61   64  
62   Provides a fixed pool of reusable worker threads for operations 65   Provides a fixed pool of reusable worker threads for operations
63   that cannot be integrated with async I/O (e.g. blocking DNS 66   that cannot be integrated with async I/O (e.g. blocking DNS
64   calls). Registered as an `execution_context::service` so it 67   calls). Registered as an `execution_context::service` so it
65   is a singleton per io_context. 68   is a singleton per io_context.
66   69  
67 - Threads are created eagerly in the constructor. The default 70 + The service is created with its context, but the workers start on
68 - thread count is 1. 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.
69   74  
70   @par Thread Safety 75   @par Thread Safety
71   All public member functions are thread-safe. 76   All public member functions are thread-safe.
72   77  
73   @par Shutdown 78   @par Shutdown
74   Sets a shutdown flag, notifies all threads, and joins them. 79   Sets a shutdown flag, notifies all threads, and joins them.
75   In-flight blocking calls complete naturally before the thread 80   In-flight blocking calls complete naturally before the thread
76   exits. 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).
77   */ 94   */
78 - class thread_pool final : public capy::execution_context::service 95 + class BOOST_COROSIO_SYMBOL_VISIBLE thread_pool final
  96 + : public capy::execution_context::service
79   { 97   {
80   std::mutex mutex_; 98   std::mutex mutex_;
81   std::condition_variable cv_; 99   std::condition_variable cv_;
82   intrusive_queue<pool_work_item> work_queue_; 100   intrusive_queue<pool_work_item> work_queue_;
83   std::vector<std::thread> threads_; 101   std::vector<std::thread> threads_;
  102 + unsigned num_threads_;
84   bool shutdown_ = false; 103   bool shutdown_ = false;
85   104  
86   void worker_loop(unsigned index); 105   void worker_loop(unsigned index);
  106 + std::error_code start_workers() noexcept;
87   107  
88   public: 108   public:
89   using key_type = thread_pool; 109   using key_type = thread_pool;
90   110  
91   /** Construct the thread pool service. 111   /** Construct the thread pool service.
92   112  
93 - Eagerly creates all worker threads. 113 + Records the worker count. The workers themselves start on the
  114 + first `post()`.
94   115  
95   @par Exception Safety 116   @par Exception Safety
96 - Strong guarantee. If thread creation fails, all 117 + Strong guarantee.
97 - already-created threads are shut down and joined  
98 - before the exception propagates.  
99   118  
100   @param ctx Reference to the owning execution_context. 119   @param ctx Reference to the owning execution_context.
101   @param num_threads Number of worker threads. Must be 120   @param num_threads Number of worker threads. Must be
102   at least 1. 121   at least 1.
103   122  
104   @throws std::logic_error If `num_threads` is 0. 123   @throws std::logic_error If `num_threads` is 0.
105   */ 124   */
HITCBC 106   1605 explicit thread_pool( 125   1792 explicit thread_pool(
107   [[maybe_unused]] capy::execution_context& ctx, 126   [[maybe_unused]] capy::execution_context& ctx,
108   unsigned num_threads = 1) 127   unsigned num_threads = 1)
HITGNC   128 + 1792 : num_threads_(num_threads)
ECB 109   1605 { 129   {
HITCBC 110   1605 if (!num_threads) 130   1792 if (!num_threads)
DCB 111 - 1 threads_.reserve(num_threads);  
DCB 112 - 1604 try  
113 - {  
114 - for (unsigned i = 0; i < num_threads; ++i)  
DCB 115 - 3217 threads_.emplace_back([this, i] { worker_loop(i + 1); });  
DCB 116 - 3226 }  
117 - catch (...)  
DUB 118 - {  
119 - shutdown();  
DUB 120 - throw;  
DUB 121 - }  
HITGBC 122   throw std::logic_error("thread_pool requires at least 1 thread"); 131   1 throw std::logic_error("thread_pool requires at least 1 thread");
HITCBC 123   1607 } 132   1794 }
124   133  
ECB 125 - 3207 ~thread_pool() override = default; 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 + */
HITGNC   142 + 3581 ~thread_pool() override
HITGNC   143 + 1791 {
HITGNC   144 + 1791 if (!threads_.empty())
MISUNC   145 + shutdown();
HITGNC   146 + 3581 }
126   147  
127   thread_pool(thread_pool const&) = delete; 148   thread_pool(thread_pool const&) = delete;
128   thread_pool& operator=(thread_pool const&) = delete; 149   thread_pool& operator=(thread_pool const&) = delete;
129   150  
130   /** Enqueue a work item for execution on the thread pool. 151   /** Enqueue a work item for execution on the thread pool.
131   152  
132 - Zero-allocation: the caller owns the work item's storage. 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.
133   162  
134   @param w The work item to execute. Must remain valid until 163   @param w The work item to execute. Must remain valid until
135   its `func_` has been called. 164   its `func_` has been called.
136   165  
137 - @return `true` if the item was enqueued, `false` if the 166 + @return An empty code if the item was enqueued;
138 - pool has already shut down. 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.
139   */ 170   */
140 - bool post(pool_work_item* w) noexcept; 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 + */
HITGNC   181 + 5 unsigned worker_count() noexcept
  182 + {
HITGNC   183 + 5 std::lock_guard<std::mutex> lock(mutex_);
HITGNC   184 + 5 return static_cast<unsigned>(threads_.size());
HITGNC   185 + 5 }
141   186  
142   /** Shut down the thread pool. 187   /** Shut down the thread pool.
143   188  
144   Signals all threads to exit after draining any 189   Signals all threads to exit after draining any
145   remaining queued work, then joins them. 190   remaining queued work, then joins them.
146   */ 191   */
147   void shutdown() override; 192   void shutdown() override;
148   }; 193   };
149   194  
150   inline void 195   inline void
HITCBC 151   1613 thread_pool::worker_loop(unsigned index) 196   136 thread_pool::worker_loop(unsigned index)
152   { 197   {
153   // Name format chosen to fit Linux's 15-char pthread limit: 198   // Name format chosen to fit Linux's 15-char pthread limit:
154   // "tpool-svc-" (10) + up to 4 digit index leaves "tpool-svc-9999". 199   // "tpool-svc-" (10) + up to 4 digit index leaves "tpool-svc-9999".
155   char name[16]; 200   char name[16];
HITCBC 156   1613 std::snprintf(name, sizeof(name), "tpool-svc-%u", index); 201   136 std::snprintf(name, sizeof(name), "tpool-svc-%u", index);
HITCBC 157   1613 capy::set_current_thread_name(name); 202   136 capy::set_current_thread_name(name);
158   203  
159   for (;;) 204   for (;;)
160   { 205   {
161   pool_work_item* w; 206   pool_work_item* w;
162   { 207   {
HITCBC 163   2005 std::unique_lock<std::mutex> lock(mutex_); 208   591 std::unique_lock<std::mutex> lock(mutex_);
HITCBC 164   2005 cv_.wait( 209   591 cv_.wait(
HITCBC 165   2668 lock, [this] { return shutdown_ || !work_queue_.empty(); }); 210   757 lock, [this] { return shutdown_ || !work_queue_.empty(); });
166   211  
HITCBC 167   2005 w = work_queue_.pop(); 212   591 w = work_queue_.pop();
HITCBC 168   2005 if (!w) 213   591 if (!w)
169   { 214   {
HITCBC 170   1613 if (shutdown_) 215   136 if (shutdown_)
HITCBC 171   3226 return; 216   272 return;
MISUBC 172   continue; 217   continue;
173   } 218   }
HITCBC 174   2005 } 219   591 }
HITCBC 175   392 w->func_(w); 220   455 w->func_(w);
HITCBC 176   392 } 221   455 }
177   } 222   }
178   223  
179 - inline bool 224 + // Called with mutex_ held, so the workers are started once however
  225 + // many threads race the first post.
  226 + inline std::error_code
HITGNC   227 + 455 thread_pool::start_workers() noexcept
  228 + {
HITGNC   229 + 455 if (!threads_.empty())
HITGNC   230 + 322 return {};
HITGNC   231 + 133 std::error_code ec;
  232 + try
  233 + {
HITGNC   234 + 133 threads_.reserve(num_threads_);
HITGNC   235 + 269 for (unsigned i = 0; i < num_threads_; ++i)
HITGNC   236 + 272 threads_.emplace_back([this, i] { worker_loop(i + 1); });
  237 + }
MISUNC   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.
MISUNC   244 + ec = e.code();
MISUNC   245 + }
MISUNC   246 + catch (...)
  247 + {
MISUNC   248 + ec = std::make_error_code(std::errc::resource_unavailable_try_again);
MISUNC   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.
HITGNC   255 + 133 if (!threads_.empty())
HITGNC   256 + 133 return {};
MISUNC   257 + return ec;
  258 + }
  259 +
  260 + inline std::error_code
HITCBC 180   393 thread_pool::post(pool_work_item* w) noexcept 261   456 thread_pool::post(pool_work_item* w) noexcept
181   { 262   {
182   { 263   {
HITCBC 183   393 std::lock_guard<std::mutex> lock(mutex_); 264   456 std::lock_guard<std::mutex> lock(mutex_);
HITCBC 184   393 if (shutdown_) 265   456 if (shutdown_)
HITCBC 185 - 1 return false; 266 + 1 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.
HITGNC   270 + 455 if (auto ec = start_workers())
MISUNC   271 + return ec;
HITCBC 186   392 work_queue_.push(w); 272   455 work_queue_.push(w);
HITCBC 187   393 } 273   456 }
HITCBC 188   392 cv_.notify_one(); 274   455 cv_.notify_one();
HITCBC 189 - 392 return true; 275 + 455 return {};
190   } 276   }
191   277  
192   inline void 278   inline void
HITCBC 193   1608 thread_pool::shutdown() 279   1796 thread_pool::shutdown()
194   { 280   {
195   { 281   {
HITCBC 196   1608 std::lock_guard<std::mutex> lock(mutex_); 282   1796 std::lock_guard<std::mutex> lock(mutex_);
HITCBC 197   1608 shutdown_ = true; 283   1796 shutdown_ = true;
HITCBC 198   1608 } 284   1796 }
HITCBC 199   1608 cv_.notify_all(); 285   1796 cv_.notify_all();
200   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.
HITCBC 201   3221 for (auto& t : threads_) 291   1932 for (auto& t : threads_)
202   { 292   {
HITCBC 203   1613 if (t.joinable()) 293   136 if (t.joinable())
HITCBC 204   1613 t.join(); 294   136 t.join();
205   } 295   }
HITCBC 206   1608 threads_.clear(); 296   1796 threads_.clear();
207   297  
208   { 298   {
HITCBC 209   1608 std::lock_guard<std::mutex> lock(mutex_); 299   1796 std::lock_guard<std::mutex> lock(mutex_);
HITCBC 210   1608 while (work_queue_.pop()) 300   1796 while (work_queue_.pop())
211   ; 301   ;
HITCBC 212   1608 } 302   1796 }
HITCBC 213   1608 } 303   1796 }
  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 + */
HITGNC   338 + 5370 explicit thread_pool_ref(capy::execution_context& ctx) noexcept
HITGNC   339 + 5370 : ctx_(ctx)
  340 + {
HITGNC   341 + 5370 }
  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 + */
HITGNC   364 + 433 thread_pool& get()
  365 + {
HITGNC   366 + 433 auto* p = pool_.load(std::memory_order_acquire);
HITGNC   367 + 433 if (!p)
  368 + {
HITGNC   369 + 129 p = &ctx_.use_service<thread_pool>();
HITGNC   370 + 129 pool_.store(p, std::memory_order_release);
  371 + }
HITGNC   372 + 433 return *p;
  373 + }
  374 + };
214   375  
215   } // namespace boost::corosio::detail 376   } // namespace boost::corosio::detail
216   377  
217   #endif // BOOST_COROSIO_DETAIL_THREAD_POOL_HPP 378   #endif // BOOST_COROSIO_DETAIL_THREAD_POOL_HPP