LCOV - code coverage report
Current view: top level - corosio/native/detail/reactor - reactor_scheduler.hpp (source / functions) Coverage Total Hit Missed
Test: coverage_remapped.info Lines: 96.4 % 329 317 12
Test Date: 2026-09-25 21:36:08 Functions: 100.0 % 43 43

           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_NATIVE_DETAIL_REACTOR_REACTOR_SCHEDULER_HPP
      11                 : #define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_SCHEDULER_HPP
      12                 : 
      13                 : #include <boost/corosio/detail/config.hpp>
      14                 : #include <boost/capy/ex/execution_context.hpp>
      15                 : 
      16                 : #include <boost/corosio/detail/ready_queue.hpp>
      17                 : #include <boost/corosio/detail/scheduler.hpp>
      18                 : #include <boost/corosio/detail/scheduler_op.hpp>
      19                 : #include <boost/corosio/detail/thread_local_ptr.hpp>
      20                 : 
      21                 : #include <atomic>
      22                 : #include <chrono>
      23                 : #include <coroutine>
      24                 : #include <cstddef>
      25                 : #include <cstdint>
      26                 : #include <limits>
      27                 : #include <memory>
      28                 : #include <stdexcept>
      29                 : 
      30                 : #include <boost/corosio/detail/conditionally_enabled_mutex.hpp>
      31                 : #include <boost/corosio/detail/conditionally_enabled_event.hpp>
      32                 : 
      33                 : namespace boost::corosio::detail {
      34                 : 
      35                 : // Forward declarations
      36                 : class reactor_scheduler;
      37                 : class timer_service;
      38                 : 
      39                 : /** Per-thread state for a reactor scheduler.
      40                 : 
      41                 :     Each thread running a scheduler's event loop has one of these
      42                 :     on a thread-local stack. It holds a private work queue and
      43                 :     inline completion budget for speculative I/O fast paths.
      44                 : */
      45                 : struct BOOST_COROSIO_SYMBOL_VISIBLE reactor_scheduler_context
      46                 : {
      47                 :     /// Scheduler this context belongs to.
      48                 :     reactor_scheduler const* key;
      49                 : 
      50                 :     /// Next context frame on this thread's stack.
      51                 :     reactor_scheduler_context* next;
      52                 : 
      53                 :     /// Private work queue for reduced contention.
      54                 :     ready_queue private_queue;
      55                 : 
      56                 :     /// Unflushed work count for the private queue.
      57                 :     std::int64_t private_outstanding_work;
      58                 : 
      59                 :     /// Remaining inline completions allowed this cycle.
      60                 :     int inline_budget;
      61                 : 
      62                 :     /// Maximum inline budget (adaptive, 2-16).
      63                 :     int inline_budget_max;
      64                 : 
      65                 :     /// True if no other thread absorbed queued work last cycle.
      66                 :     bool unassisted;
      67                 : 
      68                 :     /// Construct a context frame linked to @a n.
      69                 :     reactor_scheduler_context(
      70                 :         reactor_scheduler const* k, reactor_scheduler_context* n);
      71                 : };
      72                 : 
      73                 : /// Thread-local context stack for reactor schedulers.
      74                 : inline thread_local_ptr<reactor_scheduler_context> reactor_context_stack;
      75                 : 
      76                 : /// Find the context frame for a scheduler on this thread.
      77                 : inline reactor_scheduler_context*
      78 HIT      878613 : reactor_find_context(reactor_scheduler const* self) noexcept
      79                 : {
      80          878613 :     for (auto* c = reactor_context_stack.get(); c != nullptr; c = c->next)
      81                 :     {
      82          853720 :         if (c->key == self)
      83          853720 :             return c;
      84                 :     }
      85           24893 :     return nullptr;
      86                 : }
      87                 : 
      88                 : /** Non-template base for reactor-backed scheduler implementations.
      89                 : 
      90                 :     Provides the complete threading model shared by epoll, kqueue,
      91                 :     and select schedulers: signal state machine, inline completion
      92                 :     budget, work counting, run/poll methods, and the do_one event
      93                 :     loop.
      94                 : 
      95                 :     Derived classes provide platform-specific hooks by overriding:
      96                 :     - `run_task(lock, ctx)` to run the reactor poll
      97                 :     - `interrupt_reactor()` to wake a blocked reactor
      98                 : 
      99                 :     De-templated from the original CRTP design to eliminate
     100                 :     duplicate instantiations when multiple backends are compiled
     101                 :     into the same binary. Virtual dispatch for run_task (called
     102                 :     once per reactor cycle, before a blocking syscall) has
     103                 :     negligible overhead.
     104                 : 
     105                 :     @par Thread Safety
     106                 :     All public member functions are thread-safe.
     107                 : */
     108                 : class reactor_scheduler
     109                 :     : public scheduler
     110                 :     , public capy::execution_context::service
     111                 : {
     112                 : public:
     113                 :     using key_type     = scheduler;
     114                 :     using context_type = reactor_scheduler_context;
     115                 :     using mutex_type   = conditionally_enabled_mutex;
     116                 :     using lock_type    = mutex_type::scoped_lock;
     117                 :     using event_type   = conditionally_enabled_event;
     118                 : 
     119                 :     /// Post a coroutine for deferred execution.
     120                 :     void post(std::coroutine_handle<> h) const override;
     121                 : 
     122                 :     /// Post a scheduler operation for deferred execution.
     123                 :     void post(scheduler_op* h) const override;
     124                 : 
     125                 :     /// Post a continuation for deferred execution.
     126                 :     void post(capy::continuation&) const override;
     127                 : 
     128                 :     /// Return true if called from a thread running this scheduler.
     129                 :     bool running_in_this_thread() const noexcept override;
     130                 : 
     131                 :     /// Request the scheduler to stop dispatching handlers.
     132                 :     void stop() override;
     133                 : 
     134                 :     /// Return true if the scheduler has been stopped.
     135                 :     bool stopped() const noexcept override;
     136                 : 
     137                 :     /// Reset the stopped state so `run()` can resume.
     138                 :     void restart() override;
     139                 : 
     140                 :     /// Run the event loop until no work remains.
     141                 :     std::size_t run() override;
     142                 : 
     143                 :     /// Run until one handler completes or no work remains.
     144                 :     std::size_t run_one() override;
     145                 : 
     146                 :     /// Run until one handler completes or @a usec elapses.
     147                 :     std::size_t wait_one(long usec) override;
     148                 : 
     149                 :     /// Run ready handlers without blocking.
     150                 :     std::size_t poll() override;
     151                 : 
     152                 :     /// Run at most one ready handler without blocking.
     153                 :     std::size_t poll_one() override;
     154                 : 
     155                 :     /// Increment the outstanding work count.
     156                 :     void work_started() noexcept override;
     157                 : 
     158                 :     /// Decrement the outstanding work count, stopping on zero.
     159                 :     void work_finished() noexcept override;
     160                 : 
     161                 :     /** Reset the thread's inline completion budget.
     162                 : 
     163                 :         Called at the start of each posted completion handler to
     164                 :         grant a fresh budget for speculative inline completions.
     165                 :     */
     166                 :     void reset_inline_budget() const noexcept;
     167                 : 
     168                 :     /** Consume one unit of inline budget if available.
     169                 : 
     170                 :         @return True if budget was available and consumed.
     171                 :     */
     172                 :     bool try_consume_inline_budget() const noexcept;
     173                 : 
     174                 :     /** Offset a forthcoming work_finished from work_cleanup.
     175                 : 
     176                 :         Called by descriptor_state when all I/O returned EAGAIN and
     177                 :         no handler will be executed. Must be called from a scheduler
     178                 :         thread.
     179                 :     */
     180                 :     void compensating_work_started() const noexcept;
     181                 : 
     182                 :     /** Post completed operations for deferred invocation.
     183                 : 
     184                 :         If called from a thread running this scheduler, operations
     185                 :         go to the thread's private queue (fast path). Otherwise,
     186                 :         operations are added to the global queue under mutex and a
     187                 :         waiter is signaled.
     188                 : 
     189                 :         @pre work_started() must have been called for each operation.
     190                 : 
     191                 :         @param ops Queue of operations to post.
     192                 :     */
     193                 :     void post_deferred_completions(ready_queue& ops) const;
     194                 : 
     195                 :     /** Apply runtime configuration to the scheduler.
     196                 : 
     197                 :         Called by `io_context` after construction. Values that do
     198                 :         not apply to this backend are silently ignored.
     199                 : 
     200                 :         @param max_events  Event buffer size for epoll/kqueue.
     201                 :         @param budget_init Starting inline completion budget.
     202                 :         @param budget_max  Hard ceiling on adaptive budget ramp-up.
     203                 :         @param unassisted  Budget when single-threaded.
     204                 :     */
     205                 :     virtual void configure_reactor(
     206                 :         unsigned max_events,
     207                 :         unsigned budget_init,
     208                 :         unsigned budget_max,
     209                 :         unsigned unassisted);
     210                 : 
     211                 :     /// Return the configured initial inline budget.
     212            5451 :     unsigned inline_budget_initial() const noexcept
     213                 :     {
     214            5451 :         return inline_budget_initial_;
     215                 :     }
     216                 : 
     217                 :     /// Return true when scheduler locking is disabled (fully-lockless tier).
     218             502 :     bool scheduler_locking_disabled() const noexcept override
     219                 :     {
     220             502 :         return scheduler_locking_disabled_;
     221                 :     }
     222                 : 
     223            2241 :     void configure_threading(threading_config cfg) noexcept override
     224                 :     {
     225            2241 :         scheduler_locking_disabled_ = !cfg.scheduler_locking;
     226                 :         // reactor_io_locking takes effect at descriptor registration (see the
     227                 :         // register_descriptor overrides), not here.
     228            2241 :         reactor_io_locking_ = cfg.reactor_io_locking;
     229            2241 :         one_thread_         = cfg.one_thread;
     230            2241 :         mutex_.set_enabled(cfg.scheduler_locking);
     231            2241 :         cond_.set_enabled(cfg.scheduler_locking);
     232            2241 :     }
     233                 : 
     234                 : protected:
     235                 :     timer_service* timer_svc_        = nullptr;
     236                 :     bool scheduler_locking_disabled_ = false;
     237                 :     bool reactor_io_locking_         = true;
     238                 :     bool one_thread_                 = false;
     239                 : 
     240            2253 :     reactor_scheduler() = default;
     241                 : 
     242                 :     /** Drain completed_ops during shutdown.
     243                 : 
     244                 :         Pops all operations from the global queue and destroys them,
     245                 :         skipping the task sentinel. Signals all waiting threads.
     246                 :         Derived classes call this from their shutdown() override
     247                 :         before performing platform-specific cleanup.
     248                 :     */
     249                 :     void shutdown_drain();
     250                 : 
     251                 :     /// RAII guard that re-inserts the task sentinel after `run_task`.
     252                 :     struct task_cleanup
     253                 :     {
     254                 :         reactor_scheduler const* sched;
     255                 :         lock_type* lock;
     256                 :         context_type& ctx;
     257                 :         ~task_cleanup();
     258                 :     };
     259                 : 
     260                 :     mutable mutex_type mutex_{true};
     261                 :     mutable event_type cond_{true};
     262                 :     mutable ready_queue completed_ops_;
     263                 :     mutable std::atomic<std::int64_t> outstanding_work_{0};
     264                 :     std::atomic<bool> stopped_{false};
     265                 :     mutable std::atomic<bool> task_running_{false};
     266                 :     mutable bool task_interrupted_ = false;
     267                 : 
     268                 :     // Runtime-configurable reactor tuning parameters.
     269                 :     // Defaults match the library's built-in values.
     270                 :     unsigned max_events_per_poll_   = 128;
     271                 :     unsigned inline_budget_initial_ = 2;
     272                 :     unsigned inline_budget_max_     = 16;
     273                 :     unsigned unassisted_budget_     = 4;
     274                 : 
     275                 :     /// Bit 0 of `state_`: set when the condvar should be signaled.
     276                 :     static constexpr std::size_t signaled_bit = 1;
     277                 : 
     278                 :     /// Increment per waiting thread in `state_`.
     279                 :     static constexpr std::size_t waiter_increment = 2;
     280                 :     mutable std::size_t state_                    = 0;
     281                 : 
     282                 :     /// Sentinel op that triggers a reactor poll when dequeued.
     283                 :     struct task_op final : scheduler_op
     284                 :     {
     285                 :         // LCOV_EXCL_START: the sentinel is intercepted by pointer
     286                 :         // identity; its virtuals exist for vtable completeness.
     287                 :         void operator()() override {}
     288                 :         void destroy() override {}
     289                 :         // LCOV_EXCL_STOP
     290                 :     };
     291                 :     task_op task_op_;
     292                 : 
     293                 :     /** Run the platform-specific reactor poll.
     294                 : 
     295                 :         @par Postconditions
     296                 :         `lock` is owned on return, however the poll ended. An
     297                 :         implementation that unlocks around the blocking call owes the
     298                 :         caller a matching re-acquire on every path out, including the
     299                 :         errors it retries rather than reports.
     300                 :     */
     301                 :     virtual void
     302                 :     run_task(lock_type& lock, context_type& ctx, long timeout_us) = 0;
     303                 : 
     304                 :     /// Wake a blocked reactor (e.g. write to eventfd or pipe).
     305                 :     virtual void interrupt_reactor() const = 0;
     306                 : 
     307                 : private:
     308                 :     struct work_cleanup
     309                 :     {
     310                 :         reactor_scheduler* sched;
     311                 :         lock_type* lock;
     312                 :         context_type& ctx;
     313                 :         ~work_cleanup();
     314                 :     };
     315                 : 
     316                 :     std::size_t do_one(lock_type& lock, long timeout_us, context_type& ctx);
     317                 : 
     318                 :     void signal_all(lock_type& lock) const;
     319                 :     bool maybe_unlock_and_signal_one(lock_type& lock) const;
     320                 :     bool unlock_and_signal_one(lock_type& lock) const;
     321                 :     void clear_signal() const;
     322                 :     void wait_for_signal(lock_type& lock) const;
     323                 :     void wait_for_signal_for(lock_type& lock, long timeout_us) const;
     324                 :     void wake_one_thread_and_unlock(lock_type& lock) const;
     325                 : };
     326                 : 
     327                 : /** RAII guard that pushes/pops a scheduler context frame.
     328                 : 
     329                 :     On construction, pushes a new context frame onto the
     330                 :     thread-local stack. On destruction, drains any remaining
     331                 :     private queue items to the global queue and pops the frame.
     332                 : */
     333                 : struct reactor_thread_context_guard
     334                 : {
     335                 :     /// The context frame managed by this guard.
     336                 :     reactor_scheduler_context frame_;
     337                 : 
     338                 :     /// Construct the guard, pushing a frame for @a sched.
     339            5451 :     explicit reactor_thread_context_guard(
     340                 :         reactor_scheduler const* sched) noexcept
     341            5451 :         : frame_(sched, reactor_context_stack.get())
     342                 :     {
     343            5451 :         reactor_context_stack.set(&frame_);
     344            5451 :     }
     345                 : 
     346                 :     /** Destroy the guard, popping the frame.
     347                 : 
     348                 :         The private queue is empty here by invariant: work_cleanup and
     349                 :         task_cleanup splice it to the global queue after every handler
     350                 :         and every reactor pass.
     351                 :     */
     352            5451 :     ~reactor_thread_context_guard() noexcept
     353                 :     {
     354            5451 :         reactor_context_stack.set(frame_.next);
     355            5451 :     }
     356                 : };
     357                 : 
     358                 : // ---- Inline implementations ------------------------------------------------
     359                 : 
     360            5451 : inline reactor_scheduler_context::reactor_scheduler_context(
     361            5451 :     reactor_scheduler const* k, reactor_scheduler_context* n)
     362            5451 :     : key(k)
     363            5451 :     , next(n)
     364            5451 :     , private_outstanding_work(0)
     365            5451 :     , inline_budget(0)
     366            5451 :     , inline_budget_max(static_cast<int>(k->inline_budget_initial()))
     367            5451 :     , unassisted(false)
     368                 : {
     369            5451 : }
     370                 : 
     371                 : inline void
     372              44 : reactor_scheduler::configure_reactor(
     373                 :     unsigned max_events,
     374                 :     unsigned budget_init,
     375                 :     unsigned budget_max,
     376                 :     unsigned unassisted)
     377                 : {
     378              86 :     if (max_events < 1 ||
     379              42 :         max_events > static_cast<unsigned>(std::numeric_limits<int>::max()))
     380               2 :         throw std::out_of_range("max_events_per_poll must be in [1, INT_MAX]");
     381              42 :     if (budget_max > static_cast<unsigned>(std::numeric_limits<int>::max()))
     382               2 :         throw std::out_of_range("inline_budget_max must be in [0, INT_MAX]");
     383                 : 
     384                 :     // Clamp initial and unassisted to budget_max.
     385              40 :     if (budget_init > budget_max)
     386              16 :         budget_init = budget_max;
     387              40 :     if (unassisted > budget_max)
     388              16 :         unassisted = budget_max;
     389                 : 
     390              40 :     max_events_per_poll_   = max_events;
     391              40 :     inline_budget_initial_ = budget_init;
     392              40 :     inline_budget_max_     = budget_max;
     393              40 :     unassisted_budget_     = unassisted;
     394              40 : }
     395                 : 
     396                 : inline void
     397           86044 : reactor_scheduler::reset_inline_budget() const noexcept
     398                 : {
     399                 :     // When budget is disabled (max==0), all paths below would no-op
     400                 :     // (inline_budget stays 0). Skip the TLS lookup entirely.
     401           86044 :     if (inline_budget_max_ == 0)
     402              56 :         return;
     403           85988 :     if (auto* ctx = reactor_find_context(this))
     404                 :     {
     405                 :         // Cap when no other thread absorbed queued work
     406           85988 :         if (ctx->unassisted)
     407                 :         {
     408           85988 :             ctx->inline_budget_max = static_cast<int>(unassisted_budget_);
     409           85988 :             ctx->inline_budget     = static_cast<int>(unassisted_budget_);
     410           85988 :             return;
     411                 :         }
     412                 :         // Ramp up when previous cycle fully consumed budget.
     413                 :         // max(1, ...) ensures the doubling escapes zero.
     414 MIS           0 :         if (ctx->inline_budget == 0)
     415               0 :             ctx->inline_budget_max =
     416               0 :                 (std::min)((std::max)(1, ctx->inline_budget_max) * 2,
     417               0 :                            static_cast<int>(inline_budget_max_));
     418               0 :         else if (ctx->inline_budget < ctx->inline_budget_max)
     419               0 :             ctx->inline_budget_max = static_cast<int>(inline_budget_initial_);
     420               0 :         ctx->inline_budget = ctx->inline_budget_max;
     421                 :     }
     422                 : }
     423                 : 
     424                 : inline bool
     425 HIT      395233 : reactor_scheduler::try_consume_inline_budget() const noexcept
     426                 : {
     427          395233 :     if (inline_budget_max_ == 0)
     428              40 :         return false;
     429          395193 :     if (auto* ctx = reactor_find_context(this))
     430                 :     {
     431          395193 :         if (ctx->inline_budget > 0)
     432                 :         {
     433          315977 :             --ctx->inline_budget;
     434          315977 :             return true;
     435                 :         }
     436                 :     }
     437           79216 :     return false;
     438                 : }
     439                 : 
     440                 : inline void
     441            3756 : reactor_scheduler::post(std::coroutine_handle<> h) const
     442                 : {
     443                 :     struct post_handler final : scheduler_op
     444                 :     {
     445                 :         std::coroutine_handle<> h_;
     446                 : 
     447            3756 :         explicit post_handler(std::coroutine_handle<> h) : h_(h) {}
     448            7512 :         ~post_handler() override = default;
     449                 : 
     450            3744 :         void operator()() override
     451                 :         {
     452            3744 :             auto saved = h_;
     453            3744 :             delete this;
     454            3744 :             saved.resume();
     455            3744 :         }
     456                 : 
     457              12 :         void destroy() override
     458                 :         {
     459              12 :             auto saved = h_;
     460              12 :             delete this;
     461              12 :             saved.destroy();
     462              12 :         }
     463                 :     };
     464                 : 
     465            3756 :     auto ph = std::make_unique<post_handler>(h);
     466                 : 
     467            3756 :     if (auto* ctx = reactor_find_context(this))
     468                 :     {
     469              96 :         ++ctx->private_outstanding_work;
     470              96 :         ctx->private_queue.push(ph.release());
     471              96 :         return;
     472                 :     }
     473                 : 
     474            3660 :     outstanding_work_.fetch_add(1, std::memory_order_relaxed);
     475                 : 
     476            3660 :     lock_type lock(mutex_);
     477            3660 :     completed_ops_.push(ph.release());
     478            3660 :     wake_one_thread_and_unlock(lock);
     479            3756 : }
     480                 : 
     481                 : inline void
     482           97697 : reactor_scheduler::post(scheduler_op* h) const
     483                 : {
     484           97697 :     if (auto* ctx = reactor_find_context(this))
     485                 :     {
     486           96698 :         ++ctx->private_outstanding_work;
     487           96698 :         ctx->private_queue.push(h);
     488           96698 :         return;
     489                 :     }
     490                 : 
     491             999 :     outstanding_work_.fetch_add(1, std::memory_order_relaxed);
     492                 : 
     493             999 :     lock_type lock(mutex_);
     494             999 :     completed_ops_.push(h);
     495             999 :     wake_one_thread_and_unlock(lock);
     496             999 : }
     497                 : 
     498                 : inline void
     499           24257 : reactor_scheduler::post(capy::continuation& c) const
     500                 : {
     501           24257 :     if (auto* ctx = reactor_find_context(this))
     502                 :     {
     503           13877 :         ++ctx->private_outstanding_work;
     504           13877 :         ctx->private_queue.push(c);
     505           13877 :         return;
     506                 :     }
     507                 : 
     508           10380 :     outstanding_work_.fetch_add(1, std::memory_order_relaxed);
     509                 : 
     510           10380 :     lock_type lock(mutex_);
     511           10380 :     completed_ops_.push(c);
     512           10380 :     wake_one_thread_and_unlock(lock);
     513           10380 : }
     514                 : 
     515                 : inline bool
     516           10717 : reactor_scheduler::running_in_this_thread() const noexcept
     517                 : {
     518           10717 :     return reactor_find_context(this) != nullptr;
     519                 : }
     520                 : 
     521                 : inline void
     522            3733 : reactor_scheduler::stop()
     523                 : {
     524            3733 :     lock_type lock(mutex_);
     525            3733 :     if (!stopped_.load(std::memory_order_acquire))
     526                 :     {
     527            2865 :         stopped_.store(true, std::memory_order_release);
     528            2865 :         signal_all(lock);
     529            2865 :         interrupt_reactor();
     530                 :     }
     531            3733 : }
     532                 : 
     533                 : inline bool
     534            2506 : reactor_scheduler::stopped() const noexcept
     535                 : {
     536            2506 :     return stopped_.load(std::memory_order_acquire);
     537                 : }
     538                 : 
     539                 : inline void
     540            1439 : reactor_scheduler::restart()
     541                 : {
     542            1439 :     stopped_.store(false, std::memory_order_release);
     543            1439 : }
     544                 : 
     545                 : inline std::size_t
     546            2102 : reactor_scheduler::run()
     547                 : {
     548            4204 :     if (outstanding_work_.load(std::memory_order_acquire) == 0)
     549                 :     {
     550             101 :         stop();
     551             101 :         return 0;
     552                 :     }
     553                 : 
     554            2001 :     reactor_thread_context_guard ctx(this);
     555            2001 :     lock_type lock(mutex_);
     556                 : 
     557            2001 :     std::size_t n = 0;
     558                 :     for (;;)
     559                 :     {
     560          393098 :         if (!do_one(lock, -1, ctx.frame_))
     561            1998 :             break;
     562          391097 :         if (n != (std::numeric_limits<std::size_t>::max)())
     563          391097 :             ++n;
     564          391097 :         if (!lock.owns_lock())
     565          295953 :             lock.lock();
     566                 :     }
     567            1998 :     return n;
     568            2004 : }
     569                 : 
     570                 : inline std::size_t
     571             112 : reactor_scheduler::run_one()
     572                 : {
     573             224 :     if (outstanding_work_.load(std::memory_order_acquire) == 0)
     574                 :     {
     575               3 :         stop();
     576               3 :         return 0;
     577                 :     }
     578                 : 
     579             109 :     reactor_thread_context_guard ctx(this);
     580             109 :     lock_type lock(mutex_);
     581             109 :     return do_one(lock, -1, ctx.frame_);
     582             109 : }
     583                 : 
     584                 : inline std::size_t
     585            4093 : reactor_scheduler::wait_one(long usec)
     586                 : {
     587            8186 :     if (outstanding_work_.load(std::memory_order_acquire) == 0)
     588                 :     {
     589             792 :         stop();
     590             792 :         return 0;
     591                 :     }
     592                 : 
     593            3301 :     reactor_thread_context_guard ctx(this);
     594            3301 :     lock_type lock(mutex_);
     595            3301 :     return do_one(lock, usec, ctx.frame_);
     596            3301 : }
     597                 : 
     598                 : inline std::size_t
     599              49 : reactor_scheduler::poll()
     600                 : {
     601              98 :     if (outstanding_work_.load(std::memory_order_acquire) == 0)
     602                 :     {
     603              15 :         stop();
     604              15 :         return 0;
     605                 :     }
     606                 : 
     607              34 :     reactor_thread_context_guard ctx(this);
     608              34 :     lock_type lock(mutex_);
     609                 : 
     610              34 :     std::size_t n = 0;
     611                 :     for (;;)
     612                 :     {
     613              75 :         if (!do_one(lock, 0, ctx.frame_))
     614              34 :             break;
     615              41 :         if (n != (std::numeric_limits<std::size_t>::max)())
     616              41 :             ++n;
     617              41 :         if (!lock.owns_lock())
     618              41 :             lock.lock();
     619                 :     }
     620              34 :     return n;
     621              34 : }
     622                 : 
     623                 : inline std::size_t
     624              11 : reactor_scheduler::poll_one()
     625                 : {
     626              22 :     if (outstanding_work_.load(std::memory_order_acquire) == 0)
     627                 :     {
     628               5 :         stop();
     629               5 :         return 0;
     630                 :     }
     631                 : 
     632               6 :     reactor_thread_context_guard ctx(this);
     633               6 :     lock_type lock(mutex_);
     634               6 :     return do_one(lock, 0, ctx.frame_);
     635               6 : }
     636                 : 
     637                 : inline void
     638           31259 : reactor_scheduler::work_started() noexcept
     639                 : {
     640           31259 :     outstanding_work_.fetch_add(1, std::memory_order_relaxed);
     641           31259 : }
     642                 : 
     643                 : inline void
     644           61730 : reactor_scheduler::work_finished() noexcept
     645                 : {
     646          123460 :     if (outstanding_work_.fetch_sub(1, std::memory_order_acq_rel) == 1)
     647            2802 :         stop();
     648           61730 : }
     649                 : 
     650                 : inline void
     651          261003 : reactor_scheduler::compensating_work_started() const noexcept
     652                 : {
     653          261003 :     auto* ctx = reactor_find_context(this);
     654          261003 :     if (ctx)
     655          261003 :         ++ctx->private_outstanding_work;
     656          261003 : }
     657                 : 
     658                 : inline void
     659            6263 : reactor_scheduler::post_deferred_completions(ready_queue& ops) const
     660                 : {
     661            6263 :     if (ops.empty())
     662            6263 :         return;
     663                 : 
     664               2 :     if (auto* ctx = reactor_find_context(this))
     665                 :     {
     666               2 :         ctx->private_queue.splice(ops);
     667               2 :         return;
     668                 :     }
     669                 : 
     670 MIS           0 :     lock_type lock(mutex_);
     671               0 :     completed_ops_.splice(ops);
     672               0 :     wake_one_thread_and_unlock(lock);
     673               0 : }
     674                 : 
     675                 : inline void
     676 HIT        2241 : reactor_scheduler::shutdown_drain()
     677                 : {
     678            2241 :     lock_type lock(mutex_);
     679                 : 
     680            4862 :     while (auto e = completed_ops_.pop())
     681                 :     {
     682            2621 :         if (ready_is_continuation(e))
     683                 :         {
     684               8 :             lock.unlock();
     685               8 :             if (auto h = ready_as_cont(e)->h)
     686               8 :                 h.destroy();
     687               8 :             lock.lock();
     688                 :         }
     689                 :         else
     690                 :         {
     691            2613 :             auto* op = ready_as_op(e);
     692            2613 :             if (op == &task_op_)
     693            2238 :                 continue;
     694             375 :             lock.unlock();
     695             375 :             op->destroy();
     696             375 :             lock.lock();
     697                 :         }
     698            2621 :     }
     699                 : 
     700            2241 :     signal_all(lock);
     701            2241 : }
     702                 : 
     703                 : inline void
     704            5106 : reactor_scheduler::signal_all(lock_type&) const
     705                 : {
     706            5106 :     state_ |= signaled_bit;
     707            5106 :     cond_.notify_all();
     708            5106 : }
     709                 : 
     710                 : inline bool
     711           15039 : reactor_scheduler::maybe_unlock_and_signal_one(lock_type& lock) const
     712                 : {
     713           15039 :     state_ |= signaled_bit;
     714           15039 :     if (state_ > signaled_bit)
     715                 :     {
     716              31 :         lock.unlock();
     717              31 :         cond_.notify_one();
     718              31 :         return true;
     719                 :     }
     720           15008 :     return false;
     721                 : }
     722                 : 
     723                 : inline bool
     724          445819 : reactor_scheduler::unlock_and_signal_one(lock_type& lock) const
     725                 : {
     726          445819 :     state_ |= signaled_bit;
     727          445819 :     bool have_waiters = state_ > signaled_bit;
     728          445819 :     lock.unlock();
     729          445819 :     if (have_waiters)
     730               6 :         cond_.notify_one();
     731          445819 :     return have_waiters;
     732                 : }
     733                 : 
     734                 : inline void
     735              50 : reactor_scheduler::clear_signal() const
     736                 : {
     737              50 :     state_ &= ~signaled_bit;
     738              50 : }
     739                 : 
     740                 : inline void
     741               6 : reactor_scheduler::wait_for_signal(lock_type& lock) const
     742                 : {
     743              14 :     while ((state_ & signaled_bit) == 0)
     744                 :     {
     745               8 :         state_ += waiter_increment;
     746               8 :         cond_.wait(lock);
     747               8 :         state_ -= waiter_increment;
     748                 :     }
     749               6 : }
     750                 : 
     751                 : inline void
     752              44 : reactor_scheduler::wait_for_signal_for(lock_type& lock, long timeout_us) const
     753                 : {
     754              44 :     if ((state_ & signaled_bit) == 0)
     755                 :     {
     756              44 :         state_ += waiter_increment;
     757              44 :         cond_.wait_for(lock, std::chrono::microseconds(timeout_us));
     758              44 :         state_ -= waiter_increment;
     759                 :     }
     760              44 : }
     761                 : 
     762                 : inline void
     763           15039 : reactor_scheduler::wake_one_thread_and_unlock(lock_type& lock) const
     764                 : {
     765           15039 :     if (maybe_unlock_and_signal_one(lock))
     766              31 :         return;
     767                 : 
     768           15008 :     if (task_running_.load(std::memory_order_relaxed) && !task_interrupted_)
     769                 :     {
     770             891 :         task_interrupted_ = true;
     771             891 :         lock.unlock();
     772             891 :         interrupt_reactor();
     773                 :     }
     774                 :     else
     775                 :     {
     776           14117 :         lock.unlock();
     777                 :     }
     778                 : }
     779                 : 
     780          392894 : inline reactor_scheduler::work_cleanup::~work_cleanup()
     781                 : {
     782          392894 :     std::int64_t produced = ctx.private_outstanding_work;
     783          392894 :     if (produced > 1)
     784             345 :         sched->outstanding_work_.fetch_add(
     785                 :             produced - 1, std::memory_order_relaxed);
     786          392549 :     else if (produced < 1)
     787           36804 :         sched->work_finished();
     788          392894 :     ctx.private_outstanding_work = 0;
     789                 : 
     790          392894 :     if (!ctx.private_queue.empty())
     791                 :     {
     792           95406 :         lock->lock();
     793           95406 :         sched->completed_ops_.splice(ctx.private_queue);
     794                 :     }
     795          392894 : }
     796                 : 
     797          284131 : inline reactor_scheduler::task_cleanup::~task_cleanup()
     798                 : {
     799          284131 :     if (ctx.private_outstanding_work > 0)
     800                 :     {
     801            7849 :         sched->outstanding_work_.fetch_add(
     802            7849 :             ctx.private_outstanding_work, std::memory_order_relaxed);
     803            7849 :         ctx.private_outstanding_work = 0;
     804                 :     }
     805                 : 
     806          284131 :     if (!ctx.private_queue.empty())
     807                 :     {
     808            7849 :         if (!lock->owns_lock())
     809 MIS           0 :             lock->lock();
     810 HIT        7849 :         sched->completed_ops_.splice(ctx.private_queue);
     811                 :     }
     812          284131 : }
     813                 : 
     814                 : inline std::size_t
     815          396589 : reactor_scheduler::do_one(lock_type& lock, long timeout_us, context_type& ctx)
     816                 : {
     817                 :     for (;;)
     818                 :     {
     819          679109 :         if (stopped_.load(std::memory_order_acquire))
     820            2002 :             return 0;
     821                 : 
     822          677107 :         std::uintptr_t e = completed_ops_.pop();
     823          677107 :         scheduler_op* op = ready_is_continuation(e) ? nullptr : ready_as_op(e);
     824                 : 
     825                 :         // Handle reactor sentinel — time to poll for I/O
     826          677107 :         if (op == &task_op_)
     827                 :         {
     828          284162 :             bool more_handlers = !completed_ops_.empty();
     829                 : 
     830          515335 :             if (!more_handlers &&
     831          462346 :                 (outstanding_work_.load(std::memory_order_acquire) == 0 ||
     832                 :                  timeout_us == 0))
     833                 :             {
     834              31 :                 completed_ops_.push(&task_op_);
     835              31 :                 return 0;
     836                 :             }
     837                 : 
     838          284131 :             long task_timeout_us = more_handlers ? 0 : timeout_us;
     839          284131 :             task_interrupted_    = task_timeout_us == 0;
     840          284131 :             task_running_.store(true, std::memory_order_release);
     841                 : 
     842                 :             // Wake a peer to take the pending handlers while this thread
     843                 :             // polls the reactor; skipped when one_thread_ (no peer exists).
     844          284131 :             if (more_handlers && !one_thread_)
     845           52982 :                 unlock_and_signal_one(lock);
     846                 : 
     847                 :             try
     848                 :             {
     849          284131 :                 run_task(lock, ctx, task_timeout_us);
     850                 :             }
     851               3 :             catch (...)
     852                 :             {
     853               3 :                 task_running_.store(false, std::memory_order_relaxed);
     854               3 :                 throw;
     855               3 :             }
     856                 : 
     857          284128 :             task_running_.store(false, std::memory_order_relaxed);
     858          284128 :             completed_ops_.push(&task_op_);
     859          284128 :             if (timeout_us > 0)
     860            1658 :                 return 0;
     861          282470 :             continue;
     862          282470 :         }
     863                 : 
     864                 :         // Handle ready entry (op or continuation)
     865          392945 :         if (e != 0)
     866                 :         {
     867          392894 :             bool more = !completed_ops_.empty();
     868                 : 
     869          392894 :             if (more && !one_thread_)
     870                 :             {
     871                 :                 // Wake a peer for the remaining work; unassisted if none
     872                 :                 // was parked to take it.
     873          392837 :                 ctx.unassisted = !unlock_and_signal_one(lock);
     874                 :             }
     875                 :             else
     876                 :             {
     877                 :                 // No peer to wake (one_thread_, or nothing more queued).
     878              57 :                 ctx.unassisted = more;
     879              57 :                 lock.unlock();
     880                 :             }
     881                 : 
     882          392894 :             [[maybe_unused]] work_cleanup on_exit{this, &lock, ctx};
     883                 : 
     884          392894 :             if (ready_is_continuation(e))
     885           24249 :                 ready_as_cont(e)->h.resume();
     886                 :             else
     887          368645 :                 (*op)();
     888          392894 :             return 1;
     889          392894 :         }
     890                 : 
     891             102 :         if (outstanding_work_.load(std::memory_order_acquire) == 0 ||
     892                 :             timeout_us == 0)
     893               1 :             return 0;
     894                 : 
     895              50 :         clear_signal();
     896              50 :         if (timeout_us < 0)
     897               6 :             wait_for_signal(lock);
     898                 :         else
     899              44 :             wait_for_signal_for(lock, timeout_us);
     900          282520 :     }
     901                 : }
     902                 : 
     903                 : } // namespace boost::corosio::detail
     904                 : 
     905                 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_SCHEDULER_HPP
        

Generated by: LCOV version 2.3