100.00% Lines (140/140) 100.00% Functions (29/29)
TLA Baseline Branch
Line Hits Code Line Hits Code
1   // 1   //
2   // Copyright (c) 2025 Vinnie Falco (vinnie.falco@gmail.com) 2   // Copyright (c) 2025 Vinnie Falco (vinnie.falco@gmail.com)
3   // Copyright (c) 2026 Michael Vandeberg 3   // Copyright (c) 2026 Michael Vandeberg
4   // 4   //
5   // Distributed under the Boost Software License, Version 1.0. (See accompanying 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) 6   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7   // 7   //
8   // Official repository: https://github.com/boostorg/capy 8   // Official repository: https://github.com/boostorg/capy
9   // 9   //
10   10  
11   #include <boost/capy/ex/thread_pool.hpp> 11   #include <boost/capy/ex/thread_pool.hpp>
12   #include <boost/capy/continuation.hpp> 12   #include <boost/capy/continuation.hpp>
13   #include <boost/capy/detail/thread_local_ptr.hpp> 13   #include <boost/capy/detail/thread_local_ptr.hpp>
14   #include <boost/capy/ex/frame_allocator.hpp> 14   #include <boost/capy/ex/frame_allocator.hpp>
15   #include <boost/capy/test/thread_name.hpp> 15   #include <boost/capy/test/thread_name.hpp>
16   #include <algorithm> 16   #include <algorithm>
17   #include <atomic> 17   #include <atomic>
18   #include <condition_variable> 18   #include <condition_variable>
19   #include <cstdio> 19   #include <cstdio>
20   #include <mutex> 20   #include <mutex>
21   #include <thread> 21   #include <thread>
22   #include <vector> 22   #include <vector>
23   23  
24   /* 24   /*
25   Thread pool implementation using a shared work queue. 25   Thread pool implementation using a shared work queue.
26   26  
27   Work items are continuations linked via their intrusive next pointer, 27   Work items are continuations linked via their intrusive next pointer,
28   stored in a single queue protected by a mutex. No per-post heap 28   stored in a single queue protected by a mutex. No per-post heap
29   allocation: the continuation is owned by the caller and linked 29   allocation: the continuation is owned by the caller and linked
30   directly. Worker threads wait on a condition_variable until work 30   directly. Worker threads wait on a condition_variable until work
31   is available or stop is requested. 31   is available or stop is requested.
32   32  
33   Threads are started lazily on first post() via std::call_once to avoid 33   Threads are started lazily on first post() via std::call_once to avoid
34   spawning threads for pools that are constructed but never used. Each 34   spawning threads for pools that are constructed but never used. Each
35   thread is named with a configurable prefix plus index for debugger 35   thread is named with a configurable prefix plus index for debugger
36   visibility. 36   visibility.
37   37  
38   Work tracking: on_work_started/on_work_finished maintain the atomic 38   Work tracking: on_work_started/on_work_finished maintain the atomic
39   outstanding_work_ counter. on_work_started is lock-free; the worker 39   outstanding_work_ counter. on_work_started is lock-free; the worker
40   that drives the count to zero takes mutex_ and re-reads the count 40   that drives the count to zero takes mutex_ and re-reads the count
41   before deciding to stop, so the count and the stop decision stay 41   before deciding to stop, so the count and the stop decision stay
42   consistent even if work is started in between. join() blocks until 42   consistent even if work is started in between. join() blocks until
43   this counter reaches zero, then signals workers to stop and joins 43   this counter reaches zero, then signals workers to stop and joins
44   threads. 44   threads.
45   45  
46   Two shutdown paths: 46   Two shutdown paths:
47   - join(): waits for outstanding work to drain, then stops workers. 47   - join(): waits for outstanding work to drain, then stops workers.
48   - stop(): immediately signals workers to exit; queued work is abandoned. 48   - stop(): immediately signals workers to exit; queued work is abandoned.
49   - Destructor: stop() then join() (abandon + wait for threads). 49   - Destructor: stop() then join() (abandon + wait for threads).
50   */ 50   */
51   51  
52   namespace boost { 52   namespace boost {
53   namespace capy { 53   namespace capy {
54   54  
55   //------------------------------------------------------------------------------ 55   //------------------------------------------------------------------------------
56   56  
57   class thread_pool::impl 57   class thread_pool::impl
58   { 58   {
59   // Identifies the pool owning the current worker thread, or 59   // Identifies the pool owning the current worker thread, or
60   // nullptr if the calling thread is not a pool worker. Checked 60   // nullptr if the calling thread is not a pool worker. Checked
61   // by dispatch() to decide between symmetric transfer (inline 61   // by dispatch() to decide between symmetric transfer (inline
62   // resume) and post. 62   // resume) and post.
63   static inline detail::thread_local_ptr<impl const> current_; 63   static inline detail::thread_local_ptr<impl const> current_;
64   64  
65   // Intrusive queue of continuations: the next link is stored in 65   // Intrusive queue of continuations: the next link is stored in
66   // continuation::reserved (typed continuation* round-tripped through 66   // continuation::reserved (typed continuation* round-tripped through
67   // void*). No per-post allocation: the continuation is owned by the caller. 67   // void*). No per-post allocation: the continuation is owned by the caller.
68   continuation* head_ = nullptr; 68   continuation* head_ = nullptr;
69   continuation* tail_ = nullptr; 69   continuation* tail_ = nullptr;
70   70  
HITCBC 71   19680 void push(continuation* c) noexcept 71   19978 void push(continuation* c) noexcept
72   { 72   {
HITCBC 73   19680 c->reserved = nullptr; 73   19978 c->reserved = nullptr;
HITCBC 74   19680 if(tail_) 74   19978 if(tail_)
HITCBC 75   1824 tail_->reserved = c; 75   2615 tail_->reserved = c;
76   else 76   else
HITCBC 77   17856 head_ = c; 77   17363 head_ = c;
HITCBC 78   19680 tail_ = c; 78   19978 tail_ = c;
HITCBC 79   19680 } 79   19978 }
80   80  
HITCBC 81   20087 continuation* pop() noexcept 81   20380 continuation* pop() noexcept
82   { 82   {
HITCBC 83   20087 if(!head_) 83   20380 if(!head_)
HITCBC 84   407 return nullptr; 84   402 return nullptr;
HITCBC 85   19680 continuation* c = head_; 85   19978 continuation* c = head_;
HITCBC 86   19680 head_ = static_cast<continuation*>(head_->reserved); 86   19978 head_ = static_cast<continuation*>(head_->reserved);
HITCBC 87   19680 if(!head_) 87   19978 if(!head_)
HITCBC 88   17856 tail_ = nullptr; 88   17363 tail_ = nullptr;
HITCBC 89   19680 return c; 89   19978 return c;
90   } 90   }
91   91  
HITCBC 92   38048 bool empty() const noexcept 92   37574 bool empty() const noexcept
93   { 93   {
HITCBC 94   38048 return head_ == nullptr; 94   37574 return head_ == nullptr;
95   } 95   }
96   96  
97   std::mutex mutex_; 97   std::mutex mutex_;
98   std::condition_variable work_cv_; 98   std::condition_variable work_cv_;
99   std::condition_variable done_cv_; 99   std::condition_variable done_cv_;
100   std::vector<std::thread> threads_; 100   std::vector<std::thread> threads_;
101   std::atomic<std::size_t> outstanding_work_{0}; 101   std::atomic<std::size_t> outstanding_work_{0};
102   bool stop_{false}; 102   bool stop_{false};
103   bool joined_{false}; 103   bool joined_{false};
104   std::size_t num_threads_; 104   std::size_t num_threads_;
105   char thread_name_prefix_[13]{}; // 12 chars max + null terminator 105   char thread_name_prefix_[13]{}; // 12 chars max + null terminator
106   std::once_flag start_flag_; 106   std::once_flag start_flag_;
107   107  
108   public: 108   public:
HITCBC 109   407 ~impl() = default; 109   402 ~impl() = default;
110   110  
111   bool 111   bool
HITCBC 112   695 running_in_this_thread() const noexcept 112   690 running_in_this_thread() const noexcept
113   { 113   {
HITCBC 114   695 return current_.get() == this; 114   690 return current_.get() == this;
115   } 115   }
116   116  
117   // Destroy abandoned coroutine frames. Must be called 117   // Destroy abandoned coroutine frames. Must be called
118   // before execution_context::shutdown()/destroy() so 118   // before execution_context::shutdown()/destroy() so
119   // that suspended-frame destructors touching services 119   // that suspended-frame destructors touching services
120   // (e.g. cancelling registrations) run while those 120   // (e.g. cancelling registrations) run while those
121   // services are still valid. 121   // services are still valid.
122   void 122   void
HITCBC 123   407 drain_abandoned() noexcept 123   402 drain_abandoned() noexcept
124   { 124   {
HITCBC 125   579 while(auto* c = pop()) 125   579 while(auto* c = pop())
126   { 126   {
HITCBC 127   172 auto h = c->h; 127   177 auto h = c->h;
HITCBC 128   172 if(h && h != std::noop_coroutine()) 128   177 if(h && h != std::noop_coroutine())
HITCBC 129   121 h.destroy(); 129   126 h.destroy();
HITCBC 130   172 } 130   177 }
HITCBC 131   407 } 131   402 }
132   132  
HITCBC 133   407 impl(std::size_t num_threads, std::string_view thread_name_prefix) 133   402 impl(std::size_t num_threads, std::string_view thread_name_prefix)
HITCBC 134   407 : num_threads_(num_threads) 134   402 : num_threads_(num_threads)
135   { 135   {
HITCBC 136   407 if(num_threads_ == 0) 136   402 if(num_threads_ == 0)
HITCBC 137   8 num_threads_ = std::max( 137   8 num_threads_ = std::max(
HITCBC 138   4 std::thread::hardware_concurrency(), 1u); 138   4 std::thread::hardware_concurrency(), 1u);
139   139  
140   // Truncate prefix to 12 chars, leaving room for up to 3-digit index. 140   // Truncate prefix to 12 chars, leaving room for up to 3-digit index.
HITCBC 141   407 auto n = thread_name_prefix.copy(thread_name_prefix_, 12); 141   402 auto n = thread_name_prefix.copy(thread_name_prefix_, 12);
HITCBC 142   407 thread_name_prefix_[n] = '\0'; 142   402 thread_name_prefix_[n] = '\0';
HITCBC 143   407 } 143   402 }
144   144  
145   void 145   void
HITCBC 146   19680 post(continuation& c) 146   19978 post(continuation& c)
147   { 147   {
HITCBC 148   19680 ensure_started(); 148   19978 ensure_started();
149   { 149   {
HITCBC 150   19680 std::lock_guard<std::mutex> lock(mutex_); 150   19978 std::lock_guard<std::mutex> lock(mutex_);
HITCBC 151   19680 push(&c); 151   19978 push(&c);
152   // Under the lock so the pool cannot drain, join, and 152   // Under the lock so the pool cannot drain, join, and
153   // destroy the condition variable mid-signal. 153   // destroy the condition variable mid-signal.
HITCBC 154   19680 work_cv_.notify_one(); 154   19978 work_cv_.notify_one();
HITCBC 155   19680 } 155   19978 }
HITCBC 156   19680 } 156   19978 }
157   157  
158   void 158   void
HITCBC 159   694 on_work_started() noexcept 159   689 on_work_started() noexcept
160   { 160   {
HITCBC 161   694 outstanding_work_.fetch_add(1, std::memory_order_acq_rel); 161   689 outstanding_work_.fetch_add(1, std::memory_order_acq_rel);
HITCBC 162   694 } 162   689 }
163   163  
164   void 164   void
HITCBC 165   694 on_work_finished() noexcept 165   689 on_work_finished() noexcept
166   { 166   {
HITCBC 167   694 if(outstanding_work_.fetch_sub( 167   689 if(outstanding_work_.fetch_sub(
HITCBC 168   694 1, std::memory_order_acq_rel) == 1) 168   689 1, std::memory_order_acq_rel) == 1)
169   { 169   {
170   // fetch_sub's result can be stale: a concurrent 170   // fetch_sub's result can be stale: a concurrent
171   // on_work_started() may raise the count before we take the 171   // on_work_started() may raise the count before we take the
172   // lock, so re-read it here rather than trust the decrement. 172   // lock, so re-read it here rather than trust the decrement.
HITCBC 173   331 std::lock_guard<std::mutex> lock(mutex_); 173   327 std::lock_guard<std::mutex> lock(mutex_);
HITCBC 174   331 if(outstanding_work_.load( 174   327 if(outstanding_work_.load(
HITCBC 175   331 std::memory_order_acquire) == 0 && joined_ && !stop_) 175   327 std::memory_order_acquire) == 0 && joined_ && !stop_)
176   { 176   {
HITCBC 177   180 stop_ = true; 177   175 stop_ = true;
HITCBC 178   180 done_cv_.notify_all(); 178   175 done_cv_.notify_all();
HITCBC 179   180 work_cv_.notify_all(); 179   175 work_cv_.notify_all();
180   } 180   }
HITCBC 181   331 } 181   327 }
HITCBC 182   694 } 182   689 }
183   183  
184   void 184   void
HITCBC 185   680 join() noexcept 185   670 join() noexcept
186   { 186   {
187   { 187   {
HITCBC 188   680 std::unique_lock<std::mutex> lock(mutex_); 188   670 std::unique_lock<std::mutex> lock(mutex_);
HITCBC 189   680 if(joined_) 189   670 if(joined_)
HITCBC 190   273 return; 190   268 return;
HITCBC 191   407 joined_ = true; 191   402 joined_ = true;
192   192  
HITCBC 193   407 if(outstanding_work_.load( 193   402 if(outstanding_work_.load(
HITCBC 194   407 std::memory_order_acquire) == 0) 194   402 std::memory_order_acquire) == 0)
195   { 195   {
HITCBC 196   174 stop_ = true; 196   173 stop_ = true;
HITCBC 197   174 work_cv_.notify_all(); 197   173 work_cv_.notify_all();
198   } 198   }
199   else 199   else
200   { 200   {
HITCBC 201   233 done_cv_.wait(lock, [this]{ 201   229 done_cv_.wait(lock, [this]{
HITCBC 202   414 return stop_; 202   405 return stop_;
203   }); 203   });
204   } 204   }
HITCBC 205   680 } 205   670 }
206   206  
HITCBC 207   868 for(auto& t : threads_) 207   858 for(auto& t : threads_)
HITCBC 208   461 if(t.joinable()) 208   456 if(t.joinable())
HITCBC 209   461 t.join(); 209   456 t.join();
210   } 210   }
211   211  
212   void 212   void
HITCBC 213   409 stop() noexcept 213   404 stop() noexcept
214   { 214   {
215   { 215   {
HITCBC 216   409 std::lock_guard<std::mutex> lock(mutex_); 216   404 std::lock_guard<std::mutex> lock(mutex_);
HITCBC 217   409 stop_ = true; 217   404 stop_ = true;
HITCBC 218   409 } 218   404 }
HITCBC 219   409 work_cv_.notify_all(); 219   404 work_cv_.notify_all();
HITCBC 220   409 done_cv_.notify_all(); 220   404 done_cv_.notify_all();
HITCBC 221   409 } 221   404 }
222   222  
223   private: 223   private:
224   void 224   void
HITCBC 225   19680 ensure_started() 225   19978 ensure_started()
226   { 226   {
HITCBC 227   19680 std::call_once(start_flag_, [this]{ 227   19978 std::call_once(start_flag_, [this]{
HITCBC 228   357 threads_.reserve(num_threads_); 228   352 threads_.reserve(num_threads_);
HITCBC 229   818 for(std::size_t i = 0; i < num_threads_; ++i) 229   808 for(std::size_t i = 0; i < num_threads_; ++i)
HITCBC 230   922 threads_.emplace_back([this, i]{ run(i); }); 230   912 threads_.emplace_back([this, i]{ run(i); });
HITCBC 231   357 }); 231   352 });
HITCBC 232   19680 } 232   19978 }
233   233  
234   void 234   void
HITCBC 235   461 run(std::size_t index) 235   456 run(std::size_t index)
236   { 236   {
237   // Build name; set_current_thread_name truncates to platform limits. 237   // Build name; set_current_thread_name truncates to platform limits.
238   char name[16]; 238   char name[16];
HITCBC 239   461 std::snprintf(name, sizeof(name), "%s%zu", thread_name_prefix_, index); 239   456 std::snprintf(name, sizeof(name), "%s%zu", thread_name_prefix_, index);
HITCBC 240   461 set_current_thread_name(name); 240   456 set_current_thread_name(name);
241   241  
242   // Mark this thread as a worker of this pool so dispatch() 242   // Mark this thread as a worker of this pool so dispatch()
243   // can symmetric-transfer when called from within pool work. 243   // can symmetric-transfer when called from within pool work.
244   struct scoped_pool 244   struct scoped_pool
245   { 245   {
HITCBC 246   461 scoped_pool(impl const* p) noexcept { current_.set(p); } 246   456 scoped_pool(impl const* p) noexcept { current_.set(p); }
HITCBC 247   461 ~scoped_pool() noexcept { current_.set(nullptr); } 247   456 ~scoped_pool() noexcept { current_.set(nullptr); }
HITCBC 248   461 } guard(this); 248   456 } guard(this);
249   249  
250   for(;;) 250   for(;;)
251   { 251   {
HITCBC 252   19969 continuation* c = nullptr; 252   20257 continuation* c = nullptr;
253   { 253   {
HITCBC 254   19969 std::unique_lock<std::mutex> lock(mutex_); 254   20257 std::unique_lock<std::mutex> lock(mutex_);
HITCBC 255   19969 work_cv_.wait(lock, [this]{ 255   20257 work_cv_.wait(lock, [this]{
HITCBC 256   56497 return !empty() || 256   55261 return !empty() ||
HITCBC 257   56497 stop_; 257   55261 stop_;
258   }); 258   });
HITCBC 259   19969 if(stop_) 259   20257 if(stop_)
HITCBC 260   922 return; 260   912 return;
HITCBC 261   19508 c = pop(); 261   19801 c = pop();
HITCBC 262   19969 } 262   20257 }
HITCBC 263   19508 if(c) 263   19801 if(c)
HITCBC 264   19508 safe_resume(c->h); 264   19801 safe_resume(c->h);
HITCBC 265   19508 } 265   19801 }
HITCBC 266   461 } 266   456 }
267   }; 267   };
268   268  
269   //------------------------------------------------------------------------------ 269   //------------------------------------------------------------------------------
270   270  
HITCBC 271   407 thread_pool:: 271   402 thread_pool::
272   ~thread_pool() 272   ~thread_pool()
273   { 273   {
HITCBC 274   407 impl_->stop(); 274   402 impl_->stop();
HITCBC 275   407 impl_->join(); 275   402 impl_->join();
HITCBC 276   407 impl_->drain_abandoned(); 276   402 impl_->drain_abandoned();
HITCBC 277   407 shutdown(); 277   402 shutdown();
HITCBC 278   407 destroy(); 278   402 destroy();
HITCBC 279   407 delete impl_; 279   402 delete impl_;
HITCBC 280   407 } 280   402 }
281   281  
HITCBC 282   407 thread_pool:: 282   402 thread_pool::
HITCBC 283   407 thread_pool(std::size_t num_threads, std::string_view thread_name_prefix) 283   402 thread_pool(std::size_t num_threads, std::string_view thread_name_prefix)
HITCBC 284   407 : impl_(new impl(num_threads, thread_name_prefix)) 284   402 : impl_(new impl(num_threads, thread_name_prefix))
285   { 285   {
HITCBC 286   407 this->set_frame_allocator(std::allocator<void>{}); 286   402 this->set_frame_allocator(std::allocator<void>{});
HITCBC 287   407 } 287   402 }
288   288  
289   void 289   void
HITCBC 290   273 thread_pool:: 290   268 thread_pool::
291   join() noexcept 291   join() noexcept
292   { 292   {
HITCBC 293   273 impl_->join(); 293   268 impl_->join();
HITCBC 294   273 } 294   268 }
295   295  
296   void 296   void
HITCBC 297   2 thread_pool:: 297   2 thread_pool::
298   stop() noexcept 298   stop() noexcept
299   { 299   {
HITCBC 300   2 impl_->stop(); 300   2 impl_->stop();
HITCBC 301   2 } 301   2 }
302   302  
303   //------------------------------------------------------------------------------ 303   //------------------------------------------------------------------------------
304   304  
305   thread_pool::executor_type 305   thread_pool::executor_type
HITCBC 306   11819 thread_pool:: 306   11814 thread_pool::
307   get_executor() const noexcept 307   get_executor() const noexcept
308   { 308   {
HITCBC 309   11819 return executor_type( 309   11814 return executor_type(
HITCBC 310   11819 const_cast<thread_pool&>(*this)); 310   11814 const_cast<thread_pool&>(*this));
311   } 311   }
312   312  
313   void 313   void
HITCBC 314   694 thread_pool::executor_type:: 314   689 thread_pool::executor_type::
315   on_work_started() const noexcept 315   on_work_started() const noexcept
316   { 316   {
HITCBC 317   694 pool_->impl_->on_work_started(); 317   689 pool_->impl_->on_work_started();
HITCBC 318   694 } 318   689 }
319   319  
320   void 320   void
HITCBC 321   694 thread_pool::executor_type:: 321   689 thread_pool::executor_type::
322   on_work_finished() const noexcept 322   on_work_finished() const noexcept
323   { 323   {
HITCBC 324   694 pool_->impl_->on_work_finished(); 324   689 pool_->impl_->on_work_finished();
HITCBC 325   694 } 325   689 }
326   326  
327   void 327   void
HITCBC 328   18992 thread_pool::executor_type:: 328   19295 thread_pool::executor_type::
329   post(continuation& c) const 329   post(continuation& c) const
330   { 330   {
HITCBC 331   18992 pool_->impl_->post(c); 331   19295 pool_->impl_->post(c);
HITCBC 332   18992 } 332   19295 }
333   333  
334   std::coroutine_handle<> 334   std::coroutine_handle<>
HITCBC 335   695 thread_pool::executor_type:: 335   690 thread_pool::executor_type::
336   dispatch(continuation& c) const 336   dispatch(continuation& c) const
337   { 337   {
HITCBC 338   695 if(pool_->impl_->running_in_this_thread()) 338   690 if(pool_->impl_->running_in_this_thread())
HITCBC 339   7 return c.h; 339   7 return c.h;
HITCBC 340   688 pool_->impl_->post(c); 340   683 pool_->impl_->post(c);
HITCBC 341   688 return std::noop_coroutine(); 341   683 return std::noop_coroutine();
342   } 342   }
343   343  
344   } // capy 344   } // capy
345   } // boost 345   } // boost