LCOV - code coverage report
Current view: top level - corosio/detail - thread_pool.hpp (source / functions) Coverage Total Hit Missed
Test: coverage_remapped.info Lines: 86.8 % 76 66 10
Test Date: 2026-09-02 21:27:06 Functions: 100.0 % 12 12

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

Generated by: LCOV version 2.3