LCOV - code coverage report
Current view: top level - corosio/native/detail/epoll - epoll_scheduler.hpp (source / functions) Coverage Total Hit Missed
Test: coverage_remapped.info Lines: 99.3 % 153 152 1
Test Date: 2026-09-09 02:31:18 Functions: 100.0 % 12 12

           TLA  Line data    Source code
       1                 : //
       2                 : // Copyright (c) 2026 Steve Gerbino
       3                 : // Copyright (c) 2026 Michael Vandeberg
       4                 : //
       5                 : // Distributed under the Boost Software License, Version 1.0. (See accompanying
       6                 : // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
       7                 : //
       8                 : // Official repository: https://github.com/cppalliance/corosio
       9                 : //
      10                 : 
      11                 : #ifndef BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP
      12                 : #define BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP
      13                 : 
      14                 : #include <boost/corosio/detail/platform.hpp>
      15                 : 
      16                 : #if BOOST_COROSIO_HAS_EPOLL
      17                 : 
      18                 : #include <boost/corosio/detail/config.hpp>
      19                 : #include <boost/capy/ex/execution_context.hpp>
      20                 : 
      21                 : #include <boost/corosio/native/detail/reactor/reactor_scheduler.hpp>
      22                 : #include <boost/corosio/native/detail/reactor/reactor_signal_pipe.hpp>
      23                 : 
      24                 : #include <boost/corosio/native/detail/epoll/epoll_traits.hpp>
      25                 : #include <boost/corosio/detail/timer_service.hpp>
      26                 : #include <boost/corosio/native/detail/make_err.hpp>
      27                 : #include <boost/corosio/native/detail/posix/posix_resolver_service.hpp>
      28                 : #include <boost/corosio/native/detail/posix/posix_signal_service.hpp>
      29                 : #include <boost/corosio/native/detail/posix/posix_stream_file_service.hpp>
      30                 : #include <boost/corosio/native/detail/posix/posix_random_access_file_service.hpp>
      31                 : 
      32                 : #include <boost/corosio/detail/except.hpp>
      33                 : 
      34                 : #include <atomic>
      35                 : #include <chrono>
      36                 : #include <cstdint>
      37                 : #include <mutex>
      38                 : #include <vector>
      39                 : 
      40                 : #include <errno.h>
      41                 : #include <sys/epoll.h>
      42                 : #include <sys/eventfd.h>
      43                 : #include <sys/timerfd.h>
      44                 : #include <unistd.h>
      45                 : 
      46                 : namespace boost::corosio::detail {
      47                 : 
      48                 : /** Linux scheduler using epoll for I/O multiplexing.
      49                 : 
      50                 :     This scheduler implements the scheduler interface using Linux epoll
      51                 :     for efficient I/O event notification. It uses a single reactor model
      52                 :     where one thread runs epoll_wait while other threads
      53                 :     wait on a condition variable for handler work. This design provides:
      54                 : 
      55                 :     - Handler parallelism: N posted handlers can execute on N threads
      56                 :     - No thundering herd: condition_variable wakes exactly one thread
      57                 :     - IOCP parity: Behavior matches Windows I/O completion port semantics
      58                 : 
      59                 :     When threads call run(), they first try to execute queued handlers.
      60                 :     If the queue is empty and no reactor is running, one thread becomes
      61                 :     the reactor and runs epoll_wait. Other threads wait on a condition
      62                 :     variable until handlers are available.
      63                 : 
      64                 :     @par Thread Safety
      65                 :     All public member functions are thread-safe.
      66                 : */
      67                 : class BOOST_COROSIO_DECL epoll_scheduler final : public reactor_scheduler
      68                 : {
      69                 : public:
      70                 :     /** Construct the scheduler.
      71                 : 
      72                 :         Creates an epoll instance, eventfd for reactor interruption,
      73                 :         and timerfd for kernel-managed timer expiry.
      74                 : 
      75                 :         @param ctx Reference to the owning execution_context.
      76                 :         @param concurrency_hint Hint for expected thread count (unused).
      77                 :     */
      78                 :     epoll_scheduler(capy::execution_context& ctx, int concurrency_hint = -1);
      79                 : 
      80                 :     /// Destroy the scheduler.
      81                 :     ~epoll_scheduler() override;
      82                 : 
      83                 :     epoll_scheduler(epoll_scheduler const&)            = delete;
      84                 :     epoll_scheduler& operator=(epoll_scheduler const&) = delete;
      85                 : 
      86                 :     /// Shut down the scheduler, draining pending operations.
      87                 :     void shutdown() override;
      88                 : 
      89                 :     /// Apply runtime configuration, resizing the event buffer.
      90                 :     void configure_reactor(
      91                 :         unsigned max_events,
      92                 :         unsigned budget_init,
      93                 :         unsigned budget_max,
      94                 :         unsigned unassisted) override;
      95                 : 
      96                 :     /** Return the epoll file descriptor.
      97                 : 
      98                 :         Used by socket services to register file descriptors
      99                 :         for I/O event notification.
     100                 : 
     101                 :         @return The epoll file descriptor.
     102                 :     */
     103                 :     int epoll_fd() const noexcept
     104                 :     {
     105                 :         return epoll_fd_;
     106                 :     }
     107                 : 
     108                 :     /** Register a descriptor for persistent monitoring.
     109                 : 
     110                 :         The fd is registered once and stays registered until explicitly
     111                 :         deregistered. Events are dispatched via reactor_descriptor_state which
     112                 :         tracks pending read/write/connect operations.
     113                 : 
     114                 :         @param fd The file descriptor to register.
     115                 :         @param desc Pointer to descriptor data (stored in epoll_event.data.ptr).
     116                 : 
     117                 :         @return The error if registration fails, otherwise a default
     118                 :         constructed error code.
     119                 :     */
     120                 :     std::error_code
     121                 :     register_descriptor(int fd, reactor_descriptor_state* desc) const;
     122                 : 
     123                 :     /** Deregister a persistently registered descriptor.
     124                 : 
     125                 :         @param fd The file descriptor to deregister.
     126                 :     */
     127                 :     void deregister_descriptor(int fd) const;
     128                 : 
     129                 :     /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp).
     130 HIT          69 :     [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
     131                 :     {
     132              69 :         return register_descriptor(read_fd, signal_pipe_reader_.arm());
     133                 :     }
     134                 : 
     135                 : private:
     136                 :     void run_task(lock_type& lock, context_type& ctx, long timeout_us) override;
     137                 :     void interrupt_reactor() const override;
     138                 :     void update_timerfd() const;
     139                 : 
     140                 :     int epoll_fd_;
     141                 :     int event_fd_;
     142                 :     int timer_fd_;
     143                 : 
     144                 :     // Watches the global signal self-pipe's read end (armed lazily by
     145                 :     // register_signal_reader on the first signal registration).
     146                 :     reactor_signal_pipe_reader signal_pipe_reader_;
     147                 : 
     148                 :     // Edge-triggered eventfd state
     149                 :     mutable std::atomic<bool> eventfd_armed_{false};
     150                 : 
     151                 :     // Set when the earliest timer changes; flushed before epoll_wait
     152                 :     mutable std::atomic<bool> timerfd_stale_{false};
     153                 : 
     154                 :     // Event buffer sized from max_events_per_poll_ (set at construction,
     155                 :     // resized by configure_reactor via io_context_options).
     156                 :     std::vector<epoll_event> event_buffer_;
     157                 : };
     158                 : 
     159            1232 : inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
     160            1232 :     : epoll_fd_(-1)
     161            1232 :     , event_fd_(-1)
     162            1232 :     , timer_fd_(-1)
     163            2464 :     , event_buffer_(max_events_per_poll_)
     164                 : {
     165            1232 :     epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
     166            1232 :     if (epoll_fd_ < 0)
     167               1 :         detail::throw_system_error(make_err(errno), "epoll_create1");
     168                 : 
     169            1231 :     event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
     170            1231 :     if (event_fd_ < 0)
     171                 :     {
     172               1 :         int errn = errno;
     173               1 :         ::close(epoll_fd_);
     174               1 :         detail::throw_system_error(make_err(errn), "eventfd");
     175                 :     }
     176                 : 
     177            1230 :     timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
     178            1230 :     if (timer_fd_ < 0)
     179                 :     {
     180               1 :         int errn = errno;
     181               1 :         ::close(event_fd_);
     182               1 :         ::close(epoll_fd_);
     183               1 :         detail::throw_system_error(make_err(errn), "timerfd_create");
     184                 :     }
     185                 : 
     186            1229 :     epoll_event ev{};
     187            1229 :     ev.events   = EPOLLIN | EPOLLET;
     188            1229 :     ev.data.ptr = nullptr;
     189            1229 :     if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0)
     190                 :     {
     191               1 :         int errn = errno;
     192               1 :         ::close(timer_fd_);
     193               1 :         ::close(event_fd_);
     194               1 :         ::close(epoll_fd_);
     195               1 :         detail::throw_system_error(make_err(errn), "epoll_ctl");
     196                 :     }
     197                 : 
     198            1228 :     epoll_event timer_ev{};
     199            1228 :     timer_ev.events   = EPOLLIN | EPOLLERR;
     200            1228 :     timer_ev.data.ptr = &timer_fd_;
     201            1228 :     if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0)
     202                 :     {
     203               1 :         int errn = errno;
     204               1 :         ::close(timer_fd_);
     205               1 :         ::close(event_fd_);
     206               1 :         ::close(epoll_fd_);
     207               1 :         detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)");
     208                 :     }
     209                 : 
     210            1227 :     timer_svc_ = &get_timer_service(ctx, *this);
     211            1227 :     timer_svc_->set_on_earliest_changed(
     212            5586 :         timer_service::callback(this, [](void* p) {
     213            4359 :             auto* self = static_cast<epoll_scheduler*>(p);
     214            4359 :             self->timerfd_stale_.store(true, std::memory_order_release);
     215            4359 :             self->interrupt_reactor();
     216            4359 :         }));
     217                 : 
     218            1227 :     get_resolver_service(ctx, *this);
     219            1227 :     get_signal_service(ctx, *this);
     220            1227 :     get_stream_file_service(ctx, *this);
     221            1227 :     get_random_access_file_service(ctx, *this);
     222                 : 
     223            1227 :     completed_ops_.push(&task_op_);
     224            1242 : }
     225                 : 
     226            2454 : inline epoll_scheduler::~epoll_scheduler()
     227                 : {
     228            1227 :     if (timer_fd_ >= 0)
     229            1227 :         ::close(timer_fd_);
     230            1227 :     if (event_fd_ >= 0)
     231            1227 :         ::close(event_fd_);
     232            1227 :     if (epoll_fd_ >= 0)
     233            1227 :         ::close(epoll_fd_);
     234            2454 : }
     235                 : 
     236                 : inline void
     237            1227 : epoll_scheduler::shutdown()
     238                 : {
     239            1227 :     shutdown_drain();
     240                 : 
     241            1227 :     if (event_fd_ >= 0)
     242            1227 :         interrupt_reactor();
     243            1227 : }
     244                 : 
     245                 : inline void
     246              23 : epoll_scheduler::configure_reactor(
     247                 :     unsigned max_events,
     248                 :     unsigned budget_init,
     249                 :     unsigned budget_max,
     250                 :     unsigned unassisted)
     251                 : {
     252              23 :     reactor_scheduler::configure_reactor(
     253                 :         max_events, budget_init, budget_max, unassisted);
     254              21 :     event_buffer_.resize(max_events_per_poll_);
     255              21 : }
     256                 : 
     257                 : inline std::error_code
     258            5748 : epoll_scheduler::register_descriptor(
     259                 :     int fd, reactor_descriptor_state* desc) const
     260                 : {
     261            5748 :     epoll_event ev{};
     262            5748 :     ev.events   = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
     263            5748 :     ev.data.ptr = desc;
     264                 : 
     265            5748 :     if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
     266               7 :         return make_err(errno);
     267                 : 
     268            5741 :     desc->registered_events = ev.events;
     269            5741 :     desc->fd                = fd;
     270            5741 :     desc->scheduler_        = this;
     271            5741 :     desc->mutex.set_enabled(reactor_io_locking_);
     272            5741 :     desc->ready_events_.store(0, std::memory_order_relaxed);
     273                 : 
     274            5741 :     conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
     275            5741 :     desc->impl_ref_.reset();
     276            5741 :     desc->read_ready  = false;
     277            5741 :     desc->write_ready = false;
     278            5741 :     return {};
     279            5741 : }
     280                 : 
     281                 : inline void
     282            5673 : epoll_scheduler::deregister_descriptor(int fd) const
     283                 : {
     284            5673 :     ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
     285            5673 : }
     286                 : 
     287                 : inline void
     288            6764 : epoll_scheduler::interrupt_reactor() const
     289                 : {
     290            6764 :     bool expected = false;
     291            6764 :     if (eventfd_armed_.compare_exchange_strong(
     292                 :             expected, true, std::memory_order_release,
     293                 :             std::memory_order_relaxed))
     294                 :     {
     295            5222 :         std::uint64_t val = 1;
     296            5222 :         if (::write(event_fd_, &val, sizeof(val)) < 0)
     297                 :         {
     298                 :             // The flag is what coalesces later interrupts into a byte
     299                 :             // already in the eventfd; a write that failed put no byte
     300                 :             // there, so leaving it armed would swallow every interrupt
     301                 :             // that follows. Disarming keeps the cost to the interrupts
     302                 :             // already in flight -- the next one arms and writes again,
     303                 :             // instead of every one after this coalescing into a byte
     304                 :             // that does not exist.
     305               2 :             eventfd_armed_.store(false, std::memory_order_release);
     306                 :         }
     307                 :     }
     308            6764 : }
     309                 : 
     310                 : inline void
     311           11143 : epoll_scheduler::update_timerfd() const
     312                 : {
     313           11143 :     auto nearest = timer_svc_->nearest_expiry();
     314                 : 
     315           11143 :     itimerspec ts{};
     316           11143 :     int flags = 0;
     317                 : 
     318           11143 :     if (nearest == timer_service::time_point::max())
     319                 :     {
     320                 :         // No timers — disarm by setting to 0 (relative)
     321                 :     }
     322                 :     else
     323                 :     {
     324            9948 :         auto now = std::chrono::steady_clock::now();
     325            9948 :         if (nearest <= now)
     326                 :         {
     327                 :             // Use 1ns instead of 0 — zero disarms the timerfd
     328            1167 :             ts.it_value.tv_nsec = 1;
     329                 :         }
     330                 :         else
     331                 :         {
     332            8781 :             auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
     333            8781 :                             nearest - now)
     334            8781 :                             .count();
     335            8781 :             ts.it_value.tv_sec  = nsec / 1000000000;
     336            8781 :             ts.it_value.tv_nsec = nsec % 1000000000;
     337            8781 :             if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
     338 MIS           0 :                 ts.it_value.tv_nsec = 1;
     339                 :         }
     340                 :     }
     341                 : 
     342 HIT       11143 :     if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
     343               1 :         detail::throw_system_error(make_err(errno), "timerfd_settime");
     344           11142 : }
     345                 : 
     346                 : inline void
     347           39736 : epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
     348                 : {
     349                 :     int timeout_ms;
     350           39736 :     if (task_interrupted_)
     351           28575 :         timeout_ms = 0;
     352           11161 :     else if (timeout_us < 0)
     353           11146 :         timeout_ms = -1;
     354                 :     else
     355              15 :         timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
     356                 : 
     357           39736 :     if (lock.owns_lock())
     358           11163 :         lock.unlock();
     359                 : 
     360           39736 :     task_cleanup on_exit{this, &lock, ctx};
     361                 : 
     362                 :     // Flush deferred timerfd programming before blocking
     363           39736 :     if (timerfd_stale_.exchange(false, std::memory_order_acquire))
     364            3637 :         update_timerfd();
     365                 : 
     366           39735 :     int nfds = ::epoll_wait(
     367           39735 :         epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()),
     368                 :         timeout_ms);
     369                 : 
     370           39735 :     if (nfds < 0 && errno != EINTR)
     371               1 :         detail::throw_system_error(make_err(errno), "epoll_wait");
     372                 : 
     373           39734 :     bool check_timers = false;
     374           39734 :     ready_queue local_ops;
     375                 : 
     376           86129 :     for (int i = 0; i < nfds; ++i)
     377                 :     {
     378           46395 :         if (event_buffer_[i].data.ptr == nullptr)
     379                 :         {
     380                 :             std::uint64_t val;
     381                 :             // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
     382            3993 :             [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
     383            3993 :             eventfd_armed_.store(false, std::memory_order_relaxed);
     384            3993 :             continue;
     385            3993 :         }
     386                 : 
     387           42402 :         if (event_buffer_[i].data.ptr == &timer_fd_)
     388                 :         {
     389                 :             std::uint64_t expirations;
     390                 :             // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
     391                 :             [[maybe_unused]] auto r =
     392            7506 :                 ::read(timer_fd_, &expirations, sizeof(expirations));
     393            7506 :             check_timers = true;
     394            7506 :             continue;
     395            7506 :         }
     396                 : 
     397                 :         auto* desc =
     398           34896 :             static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
     399           34896 :         desc->add_ready_events(event_buffer_[i].events);
     400                 : 
     401           34896 :         bool expected = false;
     402           34896 :         if (desc->is_enqueued_.compare_exchange_strong(
     403                 :                 expected, true, std::memory_order_release,
     404                 :                 std::memory_order_relaxed))
     405                 :         {
     406           34896 :             local_ops.push(desc);
     407                 :         }
     408                 :     }
     409                 : 
     410           39734 :     if (check_timers)
     411                 :     {
     412            7506 :         timer_svc_->process_expired();
     413            7506 :         update_timerfd();
     414                 :     }
     415                 : 
     416           39734 :     lock.lock();
     417                 : 
     418           39734 :     completed_ops_.splice(local_ops);
     419           39736 : }
     420                 : 
     421                 : } // namespace boost::corosio::detail
     422                 : 
     423                 : #endif // BOOST_COROSIO_HAS_EPOLL
     424                 : 
     425                 : #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP
        

Generated by: LCOV version 2.3