LCOV - code coverage report
Current view: top level - corosio/native/detail/select - select_scheduler.hpp (source / functions) Coverage Total Hit Missed
Test: coverage_remapped.info Lines: 99.4 % 174 173 1
Test Date: 2026-10-08 17:58:26 Functions: 100.0 % 13 13

           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_SELECT_SELECT_SCHEDULER_HPP
      12                 : #define BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP
      13                 : 
      14                 : #include <boost/corosio/detail/platform.hpp>
      15                 : 
      16                 : #if BOOST_COROSIO_HAS_SELECT
      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/select/select_traits.hpp>
      25                 : #include <boost/corosio/detail/timer_service.hpp>
      26                 : #include <boost/corosio/native/detail/make_err.hpp>
      27                 : 
      28                 : #include <boost/corosio/detail/except.hpp>
      29                 : 
      30                 : #include <sys/select.h>
      31                 : #include <unistd.h>
      32                 : #include <errno.h>
      33                 : #include <fcntl.h>
      34                 : 
      35                 : #include <atomic>
      36                 : #include <chrono>
      37                 : #include <cstdint>
      38                 : #include <limits>
      39                 : #include <mutex>
      40                 : #include <new>
      41                 : #include <unordered_map>
      42                 : 
      43                 : namespace boost::corosio::detail {
      44                 : 
      45                 : struct select_op;
      46                 : 
      47                 : /** POSIX scheduler using select() for I/O multiplexing.
      48                 : 
      49                 :     This scheduler implements the scheduler interface using the POSIX select()
      50                 :     call for I/O event notification. It inherits the shared reactor threading
      51                 :     model from reactor_scheduler: signal state machine, inline completion
      52                 :     budget, work counting, and the do_one event loop.
      53                 : 
      54                 :     The design mirrors epoll_scheduler for behavioral consistency:
      55                 :     - Same single-reactor thread coordination model
      56                 :     - Same deferred I/O pattern (reactor marks ready; workers do I/O)
      57                 :     - Same timer integration pattern
      58                 : 
      59                 :     Known Limitations:
      60                 :     - FD_SETSIZE (~1024) limits maximum concurrent connections
      61                 :     - O(n) scanning: rebuilds fd_sets each iteration
      62                 :     - Level-triggered only (no edge-triggered mode)
      63                 : 
      64                 :     @par Thread Safety
      65                 :     All public member functions are thread-safe.
      66                 : */
      67                 : class BOOST_COROSIO_DECL select_scheduler final : public reactor_scheduler
      68                 : {
      69                 : public:
      70                 :     /** Construct the scheduler.
      71                 : 
      72                 :         Creates a self-pipe for reactor interruption.
      73                 : 
      74                 :         @param ctx Reference to the owning execution_context.
      75                 :         @param concurrency_hint Hint for expected thread count (unused).
      76                 :     */
      77                 :     select_scheduler(capy::execution_context& ctx, int concurrency_hint = -1);
      78                 : 
      79                 :     /// Destroy the scheduler.
      80                 :     ~select_scheduler() override;
      81                 : 
      82                 :     select_scheduler(select_scheduler const&)            = delete;
      83                 :     select_scheduler& operator=(select_scheduler const&) = delete;
      84                 : 
      85                 :     /// Shut down the scheduler, draining pending operations.
      86                 :     void shutdown() override;
      87                 : 
      88                 :     /** Return the maximum file descriptor value supported.
      89                 : 
      90                 :         Returns FD_SETSIZE - 1, the maximum fd value that can be
      91                 :         monitored by select(). Operations with fd >= FD_SETSIZE
      92                 :         will fail with EINVAL.
      93                 : 
      94                 :         @return The maximum supported file descriptor value.
      95                 :     */
      96                 :     static constexpr int max_fd() noexcept
      97                 :     {
      98                 :         return FD_SETSIZE - 1;
      99                 :     }
     100                 : 
     101                 :     /** Register a descriptor for persistent monitoring.
     102                 : 
     103                 :         The fd is added to the registered_descs_ map and will be
     104                 :         included in subsequent select() calls. The reactor is
     105                 :         interrupted so a blocked select() rebuilds its fd_sets.
     106                 : 
     107                 :         @param fd The file descriptor to register.
     108                 :         @param desc Pointer to descriptor state for this fd.
     109                 : 
     110                 :         @return The error if the fd cannot be tracked, otherwise a
     111                 :         default constructed error code.
     112                 :     */
     113                 :     std::error_code
     114                 :     register_descriptor(int fd, reactor_descriptor_state* desc) const;
     115                 : 
     116                 :     /// No-op: write readiness is watched from registration on.
     117                 :     std::error_code
     118 HIT        2240 :     ensure_write_registered(int, reactor_descriptor_state*) const noexcept
     119                 :     {
     120            2240 :         return {};
     121                 :     }
     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                 :     /** Interrupt the reactor so it rebuilds its fd_sets.
     130                 : 
     131                 :         Called when a write, connect, or write-wait op is registered
     132                 :         after the reactor's snapshot was taken. Without this,
     133                 :         select() may block not watching for writability on the fd.
     134                 :     */
     135                 :     void notify_reactor() const;
     136                 : 
     137                 :     /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp).
     138              61 :     [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
     139                 :     {
     140              61 :         return register_descriptor(read_fd, signal_pipe_reader_.arm());
     141                 :     }
     142                 : 
     143                 : private:
     144                 :     void run_task(lock_type& lock, context_type& ctx, long timeout_us) override;
     145                 :     void interrupt_reactor() const override;
     146                 :     long calculate_timeout(long requested_timeout_us) const;
     147                 : 
     148                 :     // Watches the global signal self-pipe's read end (armed lazily by
     149                 :     // register_signal_reader on the first signal registration).
     150                 :     reactor_signal_pipe_reader signal_pipe_reader_;
     151                 : 
     152                 :     // Self-pipe for interrupting select()
     153                 :     int pipe_fds_[2]; // [0]=read, [1]=write
     154                 : 
     155                 :     // Per-fd tracking for fd_set building
     156                 :     mutable std::unordered_map<int, reactor_descriptor_state*>
     157                 :         registered_descs_;
     158                 :     mutable int max_fd_ = -1;
     159                 : };
     160                 : 
     161            1220 : inline select_scheduler::select_scheduler(capy::execution_context& ctx, int)
     162            1220 :     : pipe_fds_{-1, -1}
     163            1220 :     , max_fd_(-1)
     164                 : {
     165            1220 :     if (::pipe(pipe_fds_) < 0)
     166               1 :         detail::throw_system_error(make_err(errno), "pipe");
     167                 : 
     168            3648 :     for (int i = 0; i < 2; ++i)
     169                 :     {
     170            2435 :         int flags = ::fcntl(pipe_fds_[i], F_GETFL, 0);
     171            2435 :         if (flags == -1)
     172                 :         {
     173               2 :             int errn = errno;
     174               2 :             ::close(pipe_fds_[0]);
     175               2 :             ::close(pipe_fds_[1]);
     176               2 :             detail::throw_system_error(make_err(errn), "fcntl F_GETFL");
     177                 :         }
     178            2433 :         if (::fcntl(pipe_fds_[i], F_SETFL, flags | O_NONBLOCK) == -1)
     179                 :         {
     180               2 :             int errn = errno;
     181               2 :             ::close(pipe_fds_[0]);
     182               2 :             ::close(pipe_fds_[1]);
     183               2 :             detail::throw_system_error(make_err(errn), "fcntl F_SETFL");
     184                 :         }
     185            2431 :         if (::fcntl(pipe_fds_[i], F_SETFD, FD_CLOEXEC) == -1)
     186                 :         {
     187               2 :             int errn = errno;
     188               2 :             ::close(pipe_fds_[0]);
     189               2 :             ::close(pipe_fds_[1]);
     190               2 :             detail::throw_system_error(make_err(errn), "fcntl F_SETFD");
     191                 :         }
     192                 :     }
     193                 : 
     194            1213 :     timer_svc_ = &get_timer_service(ctx, *this);
     195            1213 :     timer_svc_->set_on_earliest_changed(
     196            4081 :         timer_service::callback(this, [](void* p) {
     197            2868 :             static_cast<select_scheduler*>(p)->interrupt_reactor();
     198            2868 :         }));
     199                 : 
     200            1213 :     completed_ops_.push(&task_op_);
     201            1234 : }
     202                 : 
     203            2426 : inline select_scheduler::~select_scheduler()
     204                 : {
     205            1213 :     if (pipe_fds_[0] >= 0)
     206            1213 :         ::close(pipe_fds_[0]);
     207            1213 :     if (pipe_fds_[1] >= 0)
     208            1213 :         ::close(pipe_fds_[1]);
     209            2426 : }
     210                 : 
     211                 : inline void
     212            1213 : select_scheduler::shutdown()
     213                 : {
     214            1213 :     shutdown_drain();
     215                 : 
     216            1213 :     if (pipe_fds_[1] >= 0)
     217            1213 :         interrupt_reactor();
     218            1213 : }
     219                 : 
     220                 : inline std::error_code
     221            5064 : select_scheduler::register_descriptor(
     222                 :     int fd, reactor_descriptor_state* desc) const
     223                 : {
     224            5064 :     if (fd < 0 || fd >= FD_SETSIZE)
     225               1 :         return make_err(EMFILE);
     226                 : 
     227            5063 :     desc->registered_events = reactor_event_read | reactor_event_write;
     228            5063 :     desc->unpollable        = false; // the state is reused across adoptions
     229            5063 :     desc->fd                = fd;
     230            5063 :     desc->scheduler_        = this;
     231            5063 :     desc->mutex.set_enabled(reactor_io_locking_);
     232            5063 :     desc->ready_events_.store(0, std::memory_order_relaxed);
     233                 : 
     234                 :     {
     235            5063 :         conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
     236            5063 :         desc->impl_ref_.reset();
     237            5063 :         desc->read_ready  = false;
     238            5063 :         desc->write_ready = false;
     239            5063 :     }
     240                 : 
     241                 :     {
     242            5063 :         mutex_type::scoped_lock lock(mutex_);
     243                 :         try
     244                 :         {
     245            5063 :             registered_descs_[fd] = desc;
     246                 :         }
     247               1 :         catch (std::bad_alloc const&)
     248                 :         {
     249               1 :             return make_err(ENOMEM);
     250               1 :         }
     251            5062 :         if (fd > max_fd_)
     252            5013 :             max_fd_ = fd;
     253            5063 :     }
     254                 : 
     255            5062 :     interrupt_reactor();
     256            5062 :     return {};
     257                 : }
     258                 : 
     259                 : inline void
     260            5002 : select_scheduler::deregister_descriptor(int fd) const
     261                 : {
     262            5002 :     mutex_type::scoped_lock lock(mutex_);
     263                 : 
     264            5002 :     auto it = registered_descs_.find(fd);
     265            5002 :     if (it == registered_descs_.end())
     266 MIS           0 :         return;
     267                 : 
     268 HIT        5002 :     registered_descs_.erase(it);
     269                 : 
     270            5002 :     if (fd == max_fd_)
     271                 :     {
     272            4654 :         max_fd_ = pipe_fds_[0];
     273            8856 :         for (auto& [registered_fd, state] : registered_descs_)
     274                 :         {
     275            4202 :             if (registered_fd > max_fd_)
     276            4104 :                 max_fd_ = registered_fd;
     277                 :         }
     278                 :     }
     279            5002 : }
     280                 : 
     281                 : inline void
     282            4919 : select_scheduler::notify_reactor() const
     283                 : {
     284            4919 :     interrupt_reactor();
     285            4919 : }
     286                 : 
     287                 : inline void
     288           16224 : select_scheduler::interrupt_reactor() const
     289                 : {
     290           16224 :     char byte               = 1;
     291           16224 :     [[maybe_unused]] auto r = ::write(pipe_fds_[1], &byte, 1);
     292           16224 : }
     293                 : 
     294                 : inline long
     295            7390 : select_scheduler::calculate_timeout(long requested_timeout_us) const
     296                 : {
     297            7390 :     if (requested_timeout_us == 0)
     298                 :         return 0; // LCOV_EXCL_LINE run_task passes 0 via task_interrupted_, never through this argument
     299                 : 
     300            7390 :     auto nearest = timer_svc_->nearest_expiry();
     301            7390 :     if (nearest == timer_service::time_point::max())
     302            1465 :         return requested_timeout_us;
     303                 : 
     304            5925 :     auto now = std::chrono::steady_clock::now();
     305            5925 :     if (nearest <= now)
     306             577 :         return 0;
     307                 : 
     308                 :     auto timer_timeout_us =
     309            5348 :         std::chrono::duration_cast<std::chrono::microseconds>(nearest - now)
     310            5348 :             .count();
     311                 : 
     312            5348 :     constexpr auto long_max =
     313                 :         static_cast<long long>((std::numeric_limits<long>::max)());
     314                 :     auto capped_timer_us =
     315            5348 :         (std::min)((std::max)(static_cast<long long>(timer_timeout_us),
     316            5348 :                               static_cast<long long>(0)),
     317            5348 :                    long_max);
     318                 : 
     319            5348 :     if (requested_timeout_us < 0)
     320            5346 :         return static_cast<long>(capped_timer_us);
     321                 : 
     322                 :     return static_cast<long>(
     323               2 :         (std::min)(static_cast<long long>(requested_timeout_us),
     324               2 :                    capped_timer_us));
     325                 : }
     326                 : 
     327                 : inline void
     328           33270 : select_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
     329                 : {
     330                 :     long effective_timeout_us =
     331           33270 :         task_interrupted_ ? 0 : calculate_timeout(timeout_us);
     332                 : 
     333                 :     // Snapshot registered descriptors while holding lock.
     334                 :     // Record which directions each fd needs monitored to avoid a hot
     335                 :     // loop: select is level-triggered, so a writable socket (nearly
     336                 :     // always writable) or an always-readable fd (/dev/zero, a pipe at
     337                 :     // EOF) would return select() immediately every iteration if
     338                 :     // unconditionally added. Membership in both sets is opt-in: a
     339                 :     // parked op or wait in a direction opts that direction in. The
     340                 :     // exceptional set is opt-in too, for any parked op or wait:
     341                 :     // Darwin reports a character device (/dev/zero, /dev/null) as
     342                 :     // exceptional on every call.
     343                 :     struct fd_entry
     344                 :     {
     345                 :         int fd;
     346                 :         reactor_descriptor_state* desc;
     347                 :         std::uint32_t want;
     348                 :     };
     349                 :     fd_entry snapshot[FD_SETSIZE];
     350           33270 :     int snapshot_count = 0;
     351                 : 
     352          110528 :     for (auto& [fd, desc] : registered_descs_)
     353                 :     {
     354           77258 :         if (snapshot_count < FD_SETSIZE)
     355                 :         {
     356           77258 :             conditionally_enabled_mutex::scoped_lock desc_lock(desc->mutex);
     357           77258 :             snapshot[snapshot_count].fd   = fd;
     358           77258 :             snapshot[snapshot_count].desc = desc;
     359           77258 :             snapshot[snapshot_count].want =
     360           77258 :                 ((desc->read_op || desc->wait_read_op) ? reactor_event_read
     361           77258 :                                                        : 0) |
     362           76878 :                 ((desc->write_op || desc->connect_op || desc->wait_write_op)
     363          154136 :                      ? reactor_event_write
     364           77258 :                      : 0) |
     365           77258 :                 (desc->wait_error_op ? reactor_event_error : 0);
     366           77258 :             ++snapshot_count;
     367           77258 :         }
     368                 :     }
     369                 : 
     370           33270 :     if (lock.owns_lock())
     371            7391 :         lock.unlock();
     372                 : 
     373           33270 :     task_cleanup on_exit{this, &lock, ctx};
     374                 : 
     375                 :     fd_set read_fds, write_fds, except_fds;
     376          565590 :     FD_ZERO(&read_fds);
     377          565590 :     FD_ZERO(&write_fds);
     378          565590 :     FD_ZERO(&except_fds);
     379                 : 
     380           33270 :     FD_SET(pipe_fds_[0], &read_fds);
     381           33270 :     int nfds = pipe_fds_[0];
     382                 : 
     383          110528 :     for (int i = 0; i < snapshot_count; ++i)
     384                 :     {
     385           77258 :         int fd = snapshot[i].fd;
     386           77258 :         if (snapshot[i].want & reactor_event_read)
     387            9470 :             FD_SET(fd, &read_fds);
     388           77258 :         if (snapshot[i].want & reactor_event_write)
     389            2522 :             FD_SET(fd, &write_fds);
     390           77258 :         if (snapshot[i].want != 0)
     391           12282 :             FD_SET(fd, &except_fds);
     392           77258 :         if (fd > nfds)
     393           32216 :             nfds = fd;
     394                 :     }
     395                 : 
     396                 :     struct timeval tv;
     397           33270 :     struct timeval* tv_ptr = nullptr;
     398           33270 :     if (effective_timeout_us >= 0)
     399                 :     {
     400           32165 :         tv.tv_sec  = effective_timeout_us / 1000000;
     401           32165 :         tv.tv_usec = effective_timeout_us % 1000000;
     402           32165 :         tv_ptr     = &tv;
     403                 :     }
     404                 : 
     405           33270 :     int ready = ::select(nfds + 1, &read_fds, &write_fds, &except_fds, tv_ptr);
     406                 : 
     407                 :     // EINTR: signal interrupted select(), just retry.
     408                 :     // EBADF: an fd was closed between snapshot and select(); retry
     409                 :     // with a fresh snapshot from registered_descs_.
     410                 :     // Both fall through with no ready descriptors rather than
     411                 :     // returning: the caller handed this function an owned lock that
     412                 :     // only the epilogue below re-acquires.
     413           33270 :     if (ready < 0)
     414                 :     {
     415               3 :         if (errno != EINTR && errno != EBADF)
     416               1 :             detail::throw_system_error(make_err(errno), "select");
     417               2 :         ready = 0;
     418                 :     }
     419                 : 
     420                 :     // Process timers outside the lock
     421           33269 :     timer_svc_->process_expired();
     422                 : 
     423           33269 :     ready_queue local_ops;
     424                 : 
     425           33269 :     if (ready > 0)
     426                 :     {
     427            6966 :         if (FD_ISSET(pipe_fds_[0], &read_fds))
     428                 :         {
     429                 :             char buf[256];
     430           13400 :             while (::read(pipe_fds_[0], buf, sizeof(buf)) > 0)
     431                 :             {
     432                 :             }
     433                 :         }
     434                 : 
     435           17066 :         for (int i = 0; i < snapshot_count; ++i)
     436                 :         {
     437           10100 :             int fd                         = snapshot[i].fd;
     438           10100 :             reactor_descriptor_state* desc = snapshot[i].desc;
     439                 : 
     440           10100 :             std::uint32_t flags = 0;
     441           10100 :             if (FD_ISSET(fd, &read_fds))
     442            2545 :                 flags |= reactor_event_read;
     443           10100 :             if (FD_ISSET(fd, &write_fds))
     444            2228 :                 flags |= reactor_event_write;
     445           10100 :             if (FD_ISSET(fd, &except_fds))
     446               6 :                 flags |= reactor_event_error;
     447                 : 
     448           10100 :             if (flags == 0)
     449            5325 :                 continue;
     450                 : 
     451            4775 :             desc->add_ready_events(flags);
     452                 : 
     453            4775 :             bool expected = false;
     454            4775 :             if (desc->is_enqueued_.compare_exchange_strong(
     455                 :                     expected, true, std::memory_order_release,
     456                 :                     std::memory_order_relaxed))
     457                 :             {
     458            4775 :                 local_ops.push(desc);
     459                 :             }
     460                 :         }
     461                 :     }
     462                 : 
     463           33269 :     lock.lock();
     464                 : 
     465           33269 :     completed_ops_.splice(local_ops);
     466           33270 : }
     467                 : 
     468                 : } // namespace boost::corosio::detail
     469                 : 
     470                 : #endif // BOOST_COROSIO_HAS_SELECT
     471                 : 
     472                 : #endif // BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP
        

Generated by: LCOV version 2.3