99.39% Lines (163/164) 100.00% Functions (11/11)
TLA Baseline Branch
Line Hits Code Line Hits Code
1   // 1   //
2   // Copyright (c) 2026 Steve Gerbino 2   // Copyright (c) 2026 Steve Gerbino
3   // Copyright (c) 2026 Michael Vandeberg 3   // Copyright (c) 2026 Michael Vandeberg
4   // 4   //
5   // Distributed under the Boost Software License, Version 1.0. (See accompanying 5   // Distributed under the Boost Software License, Version 1.0. (See accompanying
6   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt) 6   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7   // 7   //
8   // Official repository: https://github.com/cppalliance/corosio 8   // Official repository: https://github.com/cppalliance/corosio
9   // 9   //
10   10  
11   #ifndef BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP 11   #ifndef BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP
12   #define BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP 12   #define BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP
13   13  
14   #include <boost/corosio/detail/platform.hpp> 14   #include <boost/corosio/detail/platform.hpp>
15   15  
16   #if BOOST_COROSIO_HAS_SELECT 16   #if BOOST_COROSIO_HAS_SELECT
17   17  
18   #include <boost/corosio/detail/config.hpp> 18   #include <boost/corosio/detail/config.hpp>
19   #include <boost/capy/ex/execution_context.hpp> 19   #include <boost/capy/ex/execution_context.hpp>
20   20  
21   #include <boost/corosio/native/detail/reactor/reactor_scheduler.hpp> 21   #include <boost/corosio/native/detail/reactor/reactor_scheduler.hpp>
22   #include <boost/corosio/native/detail/reactor/reactor_signal_pipe.hpp> 22   #include <boost/corosio/native/detail/reactor/reactor_signal_pipe.hpp>
23   23  
24   #include <boost/corosio/native/detail/select/select_traits.hpp> 24   #include <boost/corosio/native/detail/select/select_traits.hpp>
25   #include <boost/corosio/detail/timer_service.hpp> 25   #include <boost/corosio/detail/timer_service.hpp>
26 - #include <boost/corosio/native/detail/posix/posix_resolver_service.hpp>  
27 - #include <boost/corosio/native/detail/posix/posix_signal_service.hpp>  
28 - #include <boost/corosio/native/detail/posix/posix_stream_file_service.hpp>  
29 - #include <boost/corosio/native/detail/posix/posix_random_access_file_service.hpp>  
30   #include <boost/corosio/native/detail/make_err.hpp> 26   #include <boost/corosio/native/detail/make_err.hpp>
31   27  
32   #include <boost/corosio/detail/except.hpp> 28   #include <boost/corosio/detail/except.hpp>
33   29  
34   #include <sys/select.h> 30   #include <sys/select.h>
35   #include <unistd.h> 31   #include <unistd.h>
36   #include <errno.h> 32   #include <errno.h>
37   #include <fcntl.h> 33   #include <fcntl.h>
38   34  
39   #include <atomic> 35   #include <atomic>
40   #include <chrono> 36   #include <chrono>
41   #include <cstdint> 37   #include <cstdint>
42   #include <limits> 38   #include <limits>
43   #include <mutex> 39   #include <mutex>
44   #include <new> 40   #include <new>
45   #include <unordered_map> 41   #include <unordered_map>
46   42  
47   namespace boost::corosio::detail { 43   namespace boost::corosio::detail {
48   44  
49   struct select_op; 45   struct select_op;
50   46  
51   /** POSIX scheduler using select() for I/O multiplexing. 47   /** POSIX scheduler using select() for I/O multiplexing.
52   48  
53   This scheduler implements the scheduler interface using the POSIX select() 49   This scheduler implements the scheduler interface using the POSIX select()
54   call for I/O event notification. It inherits the shared reactor threading 50   call for I/O event notification. It inherits the shared reactor threading
55   model from reactor_scheduler: signal state machine, inline completion 51   model from reactor_scheduler: signal state machine, inline completion
56   budget, work counting, and the do_one event loop. 52   budget, work counting, and the do_one event loop.
57   53  
58   The design mirrors epoll_scheduler for behavioral consistency: 54   The design mirrors epoll_scheduler for behavioral consistency:
59   - Same single-reactor thread coordination model 55   - Same single-reactor thread coordination model
60   - Same deferred I/O pattern (reactor marks ready; workers do I/O) 56   - Same deferred I/O pattern (reactor marks ready; workers do I/O)
61   - Same timer integration pattern 57   - Same timer integration pattern
62   58  
63   Known Limitations: 59   Known Limitations:
64   - FD_SETSIZE (~1024) limits maximum concurrent connections 60   - FD_SETSIZE (~1024) limits maximum concurrent connections
65   - O(n) scanning: rebuilds fd_sets each iteration 61   - O(n) scanning: rebuilds fd_sets each iteration
66   - Level-triggered only (no edge-triggered mode) 62   - Level-triggered only (no edge-triggered mode)
67   63  
68   @par Thread Safety 64   @par Thread Safety
69   All public member functions are thread-safe. 65   All public member functions are thread-safe.
70   */ 66   */
71   class BOOST_COROSIO_DECL select_scheduler final : public reactor_scheduler 67   class BOOST_COROSIO_DECL select_scheduler final : public reactor_scheduler
72   { 68   {
73   public: 69   public:
74   /** Construct the scheduler. 70   /** Construct the scheduler.
75   71  
76   Creates a self-pipe for reactor interruption. 72   Creates a self-pipe for reactor interruption.
77   73  
78   @param ctx Reference to the owning execution_context. 74   @param ctx Reference to the owning execution_context.
79   @param concurrency_hint Hint for expected thread count (unused). 75   @param concurrency_hint Hint for expected thread count (unused).
80   */ 76   */
81   select_scheduler(capy::execution_context& ctx, int concurrency_hint = -1); 77   select_scheduler(capy::execution_context& ctx, int concurrency_hint = -1);
82   78  
83   /// Destroy the scheduler. 79   /// Destroy the scheduler.
84   ~select_scheduler() override; 80   ~select_scheduler() override;
85   81  
86   select_scheduler(select_scheduler const&) = delete; 82   select_scheduler(select_scheduler const&) = delete;
87   select_scheduler& operator=(select_scheduler const&) = delete; 83   select_scheduler& operator=(select_scheduler const&) = delete;
88   84  
89   /// Shut down the scheduler, draining pending operations. 85   /// Shut down the scheduler, draining pending operations.
90   void shutdown() override; 86   void shutdown() override;
91   87  
92   /** Return the maximum file descriptor value supported. 88   /** Return the maximum file descriptor value supported.
93   89  
94   Returns FD_SETSIZE - 1, the maximum fd value that can be 90   Returns FD_SETSIZE - 1, the maximum fd value that can be
95   monitored by select(). Operations with fd >= FD_SETSIZE 91   monitored by select(). Operations with fd >= FD_SETSIZE
96   will fail with EINVAL. 92   will fail with EINVAL.
97   93  
98   @return The maximum supported file descriptor value. 94   @return The maximum supported file descriptor value.
99   */ 95   */
100   static constexpr int max_fd() noexcept 96   static constexpr int max_fd() noexcept
101   { 97   {
102   return FD_SETSIZE - 1; 98   return FD_SETSIZE - 1;
103   } 99   }
104   100  
105   /** Register a descriptor for persistent monitoring. 101   /** Register a descriptor for persistent monitoring.
106   102  
107   The fd is added to the registered_descs_ map and will be 103   The fd is added to the registered_descs_ map and will be
108   included in subsequent select() calls. The reactor is 104   included in subsequent select() calls. The reactor is
109   interrupted so a blocked select() rebuilds its fd_sets. 105   interrupted so a blocked select() rebuilds its fd_sets.
110   106  
111   @param fd The file descriptor to register. 107   @param fd The file descriptor to register.
112   @param desc Pointer to descriptor state for this fd. 108   @param desc Pointer to descriptor state for this fd.
113   109  
114   @return The error if the fd cannot be tracked, otherwise a 110   @return The error if the fd cannot be tracked, otherwise a
115   default constructed error code. 111   default constructed error code.
116   */ 112   */
117   std::error_code 113   std::error_code
118   register_descriptor(int fd, reactor_descriptor_state* desc) const; 114   register_descriptor(int fd, reactor_descriptor_state* desc) const;
119   115  
120   /** Deregister a persistently registered descriptor. 116   /** Deregister a persistently registered descriptor.
121   117  
122   @param fd The file descriptor to deregister. 118   @param fd The file descriptor to deregister.
123   */ 119   */
124   void deregister_descriptor(int fd) const; 120   void deregister_descriptor(int fd) const;
125   121  
126   /** Interrupt the reactor so it rebuilds its fd_sets. 122   /** Interrupt the reactor so it rebuilds its fd_sets.
127   123  
128   Called when a write, connect, or write-wait op is registered 124   Called when a write, connect, or write-wait op is registered
129   after the reactor's snapshot was taken. Without this, 125   after the reactor's snapshot was taken. Without this,
130   select() may block not watching for writability on the fd. 126   select() may block not watching for writability on the fd.
131   */ 127   */
132   void notify_reactor() const; 128   void notify_reactor() const;
133   129  
134   /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp). 130   /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp).
HITCBC 135   61 [[nodiscard]] std::error_code register_signal_reader(int read_fd) override 131   61 [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
136   { 132   {
HITCBC 137   61 return register_descriptor(read_fd, signal_pipe_reader_.arm()); 133   61 return register_descriptor(read_fd, signal_pipe_reader_.arm());
138   } 134   }
139   135  
140   private: 136   private:
141   void run_task(lock_type& lock, context_type& ctx, long timeout_us) override; 137   void run_task(lock_type& lock, context_type& ctx, long timeout_us) override;
142   void interrupt_reactor() const override; 138   void interrupt_reactor() const override;
143   long calculate_timeout(long requested_timeout_us) const; 139   long calculate_timeout(long requested_timeout_us) const;
144   140  
145   // Watches the global signal self-pipe's read end (armed lazily by 141   // Watches the global signal self-pipe's read end (armed lazily by
146   // register_signal_reader on the first signal registration). 142   // register_signal_reader on the first signal registration).
147   reactor_signal_pipe_reader signal_pipe_reader_; 143   reactor_signal_pipe_reader signal_pipe_reader_;
148   144  
149   // Self-pipe for interrupting select() 145   // Self-pipe for interrupting select()
150   int pipe_fds_[2]; // [0]=read, [1]=write 146   int pipe_fds_[2]; // [0]=read, [1]=write
151   147  
152   // Per-fd tracking for fd_set building 148   // Per-fd tracking for fd_set building
153   mutable std::unordered_map<int, reactor_descriptor_state*> 149   mutable std::unordered_map<int, reactor_descriptor_state*>
154   registered_descs_; 150   registered_descs_;
155   mutable int max_fd_ = -1; 151   mutable int max_fd_ = -1;
156   }; 152   };
157   153  
HITCBC 158   950 inline select_scheduler::select_scheduler(capy::execution_context& ctx, int) 154   950 inline select_scheduler::select_scheduler(capy::execution_context& ctx, int)
HITCBC 159   950 : pipe_fds_{-1, -1} 155   950 : pipe_fds_{-1, -1}
HITCBC 160   950 , max_fd_(-1) 156   950 , max_fd_(-1)
161   { 157   {
HITCBC 162   950 if (::pipe(pipe_fds_) < 0) 158   950 if (::pipe(pipe_fds_) < 0)
HITCBC 163   1 detail::throw_system_error(make_err(errno), "pipe"); 159   1 detail::throw_system_error(make_err(errno), "pipe");
164   160  
HITCBC 165   2838 for (int i = 0; i < 2; ++i) 161   2838 for (int i = 0; i < 2; ++i)
166   { 162   {
HITCBC 167   1895 int flags = ::fcntl(pipe_fds_[i], F_GETFL, 0); 163   1895 int flags = ::fcntl(pipe_fds_[i], F_GETFL, 0);
HITCBC 168   1895 if (flags == -1) 164   1895 if (flags == -1)
169   { 165   {
HITCBC 170   2 int errn = errno; 166   2 int errn = errno;
HITCBC 171   2 ::close(pipe_fds_[0]); 167   2 ::close(pipe_fds_[0]);
HITCBC 172   2 ::close(pipe_fds_[1]); 168   2 ::close(pipe_fds_[1]);
HITCBC 173   2 detail::throw_system_error(make_err(errn), "fcntl F_GETFL"); 169   2 detail::throw_system_error(make_err(errn), "fcntl F_GETFL");
174   } 170   }
HITCBC 175   1893 if (::fcntl(pipe_fds_[i], F_SETFL, flags | O_NONBLOCK) == -1) 171   1893 if (::fcntl(pipe_fds_[i], F_SETFL, flags | O_NONBLOCK) == -1)
176   { 172   {
HITCBC 177   2 int errn = errno; 173   2 int errn = errno;
HITCBC 178   2 ::close(pipe_fds_[0]); 174   2 ::close(pipe_fds_[0]);
HITCBC 179   2 ::close(pipe_fds_[1]); 175   2 ::close(pipe_fds_[1]);
HITCBC 180   2 detail::throw_system_error(make_err(errn), "fcntl F_SETFL"); 176   2 detail::throw_system_error(make_err(errn), "fcntl F_SETFL");
181   } 177   }
HITCBC 182   1891 if (::fcntl(pipe_fds_[i], F_SETFD, FD_CLOEXEC) == -1) 178   1891 if (::fcntl(pipe_fds_[i], F_SETFD, FD_CLOEXEC) == -1)
183   { 179   {
HITCBC 184   2 int errn = errno; 180   2 int errn = errno;
HITCBC 185   2 ::close(pipe_fds_[0]); 181   2 ::close(pipe_fds_[0]);
HITCBC 186   2 ::close(pipe_fds_[1]); 182   2 ::close(pipe_fds_[1]);
HITCBC 187   2 detail::throw_system_error(make_err(errn), "fcntl F_SETFD"); 183   2 detail::throw_system_error(make_err(errn), "fcntl F_SETFD");
188   } 184   }
189   } 185   }
190   186  
HITCBC 191   943 timer_svc_ = &get_timer_service(ctx, *this); 187   943 timer_svc_ = &get_timer_service(ctx, *this);
HITCBC 192   943 timer_svc_->set_on_earliest_changed( 188   943 timer_svc_->set_on_earliest_changed(
HITCBC 193   3758 timer_service::callback(this, [](void* p) { 189   3704 timer_service::callback(this, [](void* p) {
HITCBC 194   2815 static_cast<select_scheduler*>(p)->interrupt_reactor(); 190   2761 static_cast<select_scheduler*>(p)->interrupt_reactor();
DCB 195 - 2815  
196 - get_resolver_service(ctx, *this);  
DCB 197 - 943 get_signal_service(ctx, *this);  
DCB 198 - 943 get_stream_file_service(ctx, *this);  
DCB 199 - 943 get_random_access_file_service(ctx, *this);  
HITCBC 200   943 })); 191   2761 }));
201   192  
HITCBC 202   943 completed_ops_.push(&task_op_); 193   943 completed_ops_.push(&task_op_);
HITCBC 203   964 } 194   964 }
204   195  
HITCBC 205   1886 inline select_scheduler::~select_scheduler() 196   1886 inline select_scheduler::~select_scheduler()
206   { 197   {
HITCBC 207   943 if (pipe_fds_[0] >= 0) 198   943 if (pipe_fds_[0] >= 0)
HITCBC 208   943 ::close(pipe_fds_[0]); 199   943 ::close(pipe_fds_[0]);
HITCBC 209   943 if (pipe_fds_[1] >= 0) 200   943 if (pipe_fds_[1] >= 0)
HITCBC 210   943 ::close(pipe_fds_[1]); 201   943 ::close(pipe_fds_[1]);
HITCBC 211   1886 } 202   1886 }
212   203  
213   inline void 204   inline void
HITCBC 214   943 select_scheduler::shutdown() 205   943 select_scheduler::shutdown()
215   { 206   {
HITCBC 216   943 shutdown_drain(); 207   943 shutdown_drain();
217   208  
HITCBC 218   943 if (pipe_fds_[1] >= 0) 209   943 if (pipe_fds_[1] >= 0)
HITCBC 219   943 interrupt_reactor(); 210   943 interrupt_reactor();
HITCBC 220   943 } 211   943 }
221   212  
222   inline std::error_code 213   inline std::error_code
HITCBC 223   4906 select_scheduler::register_descriptor( 214   4764 select_scheduler::register_descriptor(
224   int fd, reactor_descriptor_state* desc) const 215   int fd, reactor_descriptor_state* desc) const
225   { 216   {
HITCBC 226   4906 if (fd < 0 || fd >= FD_SETSIZE) 217   4764 if (fd < 0 || fd >= FD_SETSIZE)
HITCBC 227   1 return make_err(EMFILE); 218   1 return make_err(EMFILE);
228   219  
HITCBC 229   4905 desc->registered_events = reactor_event_read | reactor_event_write; 220   4763 desc->registered_events = reactor_event_read | reactor_event_write;
HITCBC 230   4905 desc->fd = fd; 221   4763 desc->fd = fd;
HITCBC 231   4905 desc->scheduler_ = this; 222   4763 desc->scheduler_ = this;
HITCBC 232   4905 desc->mutex.set_enabled(reactor_io_locking_); 223   4763 desc->mutex.set_enabled(reactor_io_locking_);
HITCBC 233   4905 desc->ready_events_.store(0, std::memory_order_relaxed); 224   4763 desc->ready_events_.store(0, std::memory_order_relaxed);
234   225  
235   { 226   {
HITCBC 236   4905 conditionally_enabled_mutex::scoped_lock lock(desc->mutex); 227   4763 conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
HITCBC 237   4905 desc->impl_ref_.reset(); 228   4763 desc->impl_ref_.reset();
HITCBC 238   4905 desc->read_ready = false; 229   4763 desc->read_ready = false;
HITCBC 239   4905 desc->write_ready = false; 230   4763 desc->write_ready = false;
HITCBC 240   4905 } 231   4763 }
241   232  
242   { 233   {
HITCBC 243   4905 mutex_type::scoped_lock lock(mutex_); 234   4763 mutex_type::scoped_lock lock(mutex_);
244   try 235   try
245   { 236   {
HITCBC 246   4905 registered_descs_[fd] = desc; 237   4763 registered_descs_[fd] = desc;
247   } 238   }
HITCBC 248   1 catch (std::bad_alloc const&) 239   1 catch (std::bad_alloc const&)
249   { 240   {
HITCBC 250   1 return make_err(ENOMEM); 241   1 return make_err(ENOMEM);
HITCBC 251   1 } 242   1 }
HITCBC 252   4904 if (fd > max_fd_) 243   4762 if (fd > max_fd_)
HITCBC 253   4850 max_fd_ = fd; 244   4712 max_fd_ = fd;
HITCBC 254   4905 } 245   4763 }
255   246  
HITCBC 256   4904 interrupt_reactor(); 247   4762 interrupt_reactor();
HITCBC 257   4904 return {}; 248   4762 return {};
258   } 249   }
259   250  
260   inline void 251   inline void
HITCBC 261   4844 select_scheduler::deregister_descriptor(int fd) const 252   4702 select_scheduler::deregister_descriptor(int fd) const
262   { 253   {
HITCBC 263   4844 mutex_type::scoped_lock lock(mutex_); 254   4702 mutex_type::scoped_lock lock(mutex_);
264   255  
HITCBC 265   4844 auto it = registered_descs_.find(fd); 256   4702 auto it = registered_descs_.find(fd);
HITCBC 266   4844 if (it == registered_descs_.end()) 257   4702 if (it == registered_descs_.end())
MISUBC 267   ✗ return; 258   ✗ return;
268   259  
HITCBC 269   4844 registered_descs_.erase(it); 260   4702 registered_descs_.erase(it);
270   261  
HITCBC 271   4844 if (fd == max_fd_) 262   4702 if (fd == max_fd_)
272   { 263   {
HITCBC 273   4497 max_fd_ = pipe_fds_[0]; 264   4367 max_fd_ = pipe_fds_[0];
HITCBC 274   8584 for (auto& [registered_fd, state] : registered_descs_) 265   8328 for (auto& [registered_fd, state] : registered_descs_)
275   { 266   {
HITCBC 276   4087 if (registered_fd > max_fd_) 267   3961 if (registered_fd > max_fd_)
HITCBC 277   3994 max_fd_ = registered_fd; 268   3868 max_fd_ = registered_fd;
278   } 269   }
279   } 270   }
HITCBC 280   4844 } 271   4702 }
281   272  
282   inline void 273   inline void
HITCBC 283   2177 select_scheduler::notify_reactor() const 274   2110 select_scheduler::notify_reactor() const
284   { 275   {
HITCBC 285   2177 interrupt_reactor(); 276   2110 interrupt_reactor();
HITCBC 286   2177 } 277   2110 }
287   278  
288   inline void 279   inline void
HITCBC 289   12452 select_scheduler::interrupt_reactor() const 280   12255 select_scheduler::interrupt_reactor() const
290   { 281   {
HITCBC 291   12452 char byte = 1; 282   12255 char byte = 1;
HITCBC 292   12452 [[maybe_unused]] auto r = ::write(pipe_fds_[1], &byte, 1); 283   12255 [[maybe_unused]] auto r = ::write(pipe_fds_[1], &byte, 1);
HITCBC 293   12452 } 284   12255 }
294   285  
295   inline long 286   inline long
HITCBC 296   257933 select_scheduler::calculate_timeout(long requested_timeout_us) const 287   275633 select_scheduler::calculate_timeout(long requested_timeout_us) const
297   { 288   {
HITCBC 298   257933 if (requested_timeout_us == 0) 289   275633 if (requested_timeout_us == 0)
299   return 0; // LCOV_EXCL_LINE run_task passes 0 via task_interrupted_, never through this argument 290   return 0; // LCOV_EXCL_LINE run_task passes 0 via task_interrupted_, never through this argument
300   291  
HITCBC 301   257933 auto nearest = timer_svc_->nearest_expiry(); 292   275633 auto nearest = timer_svc_->nearest_expiry();
HITCBC 302   257933 if (nearest == timer_service::time_point::max()) 293   275633 if (nearest == timer_service::time_point::max())
HITCBC 303   1034 return requested_timeout_us; 294   1128 return requested_timeout_us;
304   295  
HITCBC 305   256899 auto now = std::chrono::steady_clock::now(); 296   274505 auto now = std::chrono::steady_clock::now();
HITCBC 306   256899 if (nearest <= now) 297   274505 if (nearest <= now)
HITCBC 307   684 return 0; 298   475 return 0;
308   299  
309   auto timer_timeout_us = 300   auto timer_timeout_us =
HITCBC 310   256215 std::chrono::duration_cast<std::chrono::microseconds>(nearest - now) 301   274030 std::chrono::duration_cast<std::chrono::microseconds>(nearest - now)
HITCBC 311   256215 .count(); 302   274030 .count();
312   303  
HITCBC 313   256215 constexpr auto long_max = 304   274030 constexpr auto long_max =
314   static_cast<long long>((std::numeric_limits<long>::max)()); 305   static_cast<long long>((std::numeric_limits<long>::max)());
315   auto capped_timer_us = 306   auto capped_timer_us =
HITCBC 316   256215 (std::min)((std::max)(static_cast<long long>(timer_timeout_us), 307   274030 (std::min)((std::max)(static_cast<long long>(timer_timeout_us),
HITCBC 317   256215 static_cast<long long>(0)), 308   274030 static_cast<long long>(0)),
HITCBC 318   256215 long_max); 309   274030 long_max);
319   310  
HITCBC 320   256215 if (requested_timeout_us < 0) 311   274030 if (requested_timeout_us < 0)
HITCBC 321   256213 return static_cast<long>(capped_timer_us); 312   274028 return static_cast<long>(capped_timer_us);
322   313  
323   return static_cast<long>( 314   return static_cast<long>(
HITCBC 324   2 (std::min)(static_cast<long long>(requested_timeout_us), 315   2 (std::min)(static_cast<long long>(requested_timeout_us),
HITCBC 325   2 capped_timer_us)); 316   2 capped_timer_us));
326   } 317   }
327   318  
328   inline void 319   inline void
HITCBC 329   282466 select_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us) 320   297590 select_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
330   { 321   {
331   long effective_timeout_us = 322   long effective_timeout_us =
HITCBC 332   282466 task_interrupted_ ? 0 : calculate_timeout(timeout_us); 323   297590 task_interrupted_ ? 0 : calculate_timeout(timeout_us);
333   324  
334   // Snapshot registered descriptors while holding lock. 325   // Snapshot registered descriptors while holding lock.
335   // Record which fds need write monitoring to avoid a hot loop: 326   // Record which fds need write monitoring to avoid a hot loop:
336   // select is level-triggered so writable sockets (nearly always 327   // select is level-triggered so writable sockets (nearly always
337   // writable) would cause select() to return immediately every 328   // writable) would cause select() to return immediately every
338   // iteration if unconditionally added to write_fds. Membership 329   // iteration if unconditionally added to write_fds. Membership
339   // stays opt-in: a parked write wait opts in the same way a 330   // stays opt-in: a parked write wait opts in the same way a
340   // parked write or connect op does. 331   // parked write or connect op does.
341   struct fd_entry 332   struct fd_entry
342   { 333   {
343   int fd; 334   int fd;
344   reactor_descriptor_state* desc; 335   reactor_descriptor_state* desc;
345   bool needs_write; 336   bool needs_write;
346   }; 337   };
347   fd_entry snapshot[FD_SETSIZE]; 338   fd_entry snapshot[FD_SETSIZE];
HITCBC 348   282466 int snapshot_count = 0; 339   297590 int snapshot_count = 0;
349   340  
HITCBC 350   747613 for (auto& [fd, desc] : registered_descs_) 341   793019 for (auto& [fd, desc] : registered_descs_)
351   { 342   {
HITCBC 352   465147 if (snapshot_count < FD_SETSIZE) 343   495429 if (snapshot_count < FD_SETSIZE)
353   { 344   {
HITCBC 354   465147 conditionally_enabled_mutex::scoped_lock desc_lock(desc->mutex); 345   495429 conditionally_enabled_mutex::scoped_lock desc_lock(desc->mutex);
HITCBC 355   465147 snapshot[snapshot_count].fd = fd; 346   495429 snapshot[snapshot_count].fd = fd;
HITCBC 356   465147 snapshot[snapshot_count].desc = desc; 347   495429 snapshot[snapshot_count].desc = desc;
HITCBC 357   465147 snapshot[snapshot_count].needs_write = 348   495429 snapshot[snapshot_count].needs_write =
HITCBC 358   465147 (desc->write_op || desc->connect_op || desc->wait_write_op); 349   495429 (desc->write_op || desc->connect_op || desc->wait_write_op);
HITCBC 359   465147 ++snapshot_count; 350   495429 ++snapshot_count;
HITCBC 360   465147 } 351   495429 }
361   } 352   }
362   353  
HITCBC 363   282466 if (lock.owns_lock()) 354   297590 if (lock.owns_lock())
HITCBC 364   257934 lock.unlock(); 355   275634 lock.unlock();
365   356  
HITCBC 366   282466 task_cleanup on_exit{this, &lock, ctx}; 357   297590 task_cleanup on_exit{this, &lock, ctx};
367   358  
368   fd_set read_fds, write_fds, except_fds; 359   fd_set read_fds, write_fds, except_fds;
HITCBC 369   4801922 FD_ZERO(&read_fds); 360   5059030 FD_ZERO(&read_fds);
HITCBC 370   4801922 FD_ZERO(&write_fds); 361   5059030 FD_ZERO(&write_fds);
HITCBC 371   4801922 FD_ZERO(&except_fds); 362   5059030 FD_ZERO(&except_fds);
372   363  
HITCBC 373   282466 FD_SET(pipe_fds_[0], &read_fds); 364   297590 FD_SET(pipe_fds_[0], &read_fds);
HITCBC 374   282466 int nfds = pipe_fds_[0]; 365   297590 int nfds = pipe_fds_[0];
375   366  
HITCBC 376   747613 for (int i = 0; i < snapshot_count; ++i) 367   793019 for (int i = 0; i < snapshot_count; ++i)
377   { 368   {
HITCBC 378   465147 int fd = snapshot[i].fd; 369   495429 int fd = snapshot[i].fd;
HITCBC 379   465147 FD_SET(fd, &read_fds); 370   495429 FD_SET(fd, &read_fds);
HITCBC 380   465147 if (snapshot[i].needs_write) 371   495429 if (snapshot[i].needs_write)
HITCBC 381   13127 FD_SET(fd, &write_fds); 372   12839 FD_SET(fd, &write_fds);
HITCBC 382   465147 FD_SET(fd, &except_fds); 373   495429 FD_SET(fd, &except_fds);
HITCBC 383   465147 if (fd > nfds) 374   495429 if (fd > nfds)
HITCBC 384   281874 nfds = fd; 375   296999 nfds = fd;
385   } 376   }
386   377  
387   struct timeval tv; 378   struct timeval tv;
HITCBC 388   282466 struct timeval* tv_ptr = nullptr; 379   297590 struct timeval* tv_ptr = nullptr;
HITCBC 389   282466 if (effective_timeout_us >= 0) 380   297590 if (effective_timeout_us >= 0)
390   { 381   {
HITCBC 391   281692 tv.tv_sec = effective_timeout_us / 1000000; 382   296831 tv.tv_sec = effective_timeout_us / 1000000;
HITCBC 392   281692 tv.tv_usec = effective_timeout_us % 1000000; 383   296831 tv.tv_usec = effective_timeout_us % 1000000;
HITCBC 393   281692 tv_ptr = &tv; 384   296831 tv_ptr = &tv;
394   } 385   }
395   386  
HITCBC 396   282466 int ready = ::select(nfds + 1, &read_fds, &write_fds, &except_fds, tv_ptr); 387   297590 int ready = ::select(nfds + 1, &read_fds, &write_fds, &except_fds, tv_ptr);
397   388  
398   // EINTR: signal interrupted select(), just retry. 389   // EINTR: signal interrupted select(), just retry.
399   // EBADF: an fd was closed between snapshot and select(); retry 390   // EBADF: an fd was closed between snapshot and select(); retry
400   // with a fresh snapshot from registered_descs_. 391   // with a fresh snapshot from registered_descs_.
401   // Both fall through with no ready descriptors rather than 392   // Both fall through with no ready descriptors rather than
402   // returning: the caller handed this function an owned lock that 393   // returning: the caller handed this function an owned lock that
403   // only the epilogue below re-acquires. 394   // only the epilogue below re-acquires.
HITCBC 404   282466 if (ready < 0) 395   297590 if (ready < 0)
405   { 396   {
HITCBC 406   3 if (errno != EINTR && errno != EBADF) 397   3 if (errno != EINTR && errno != EBADF)
HITCBC 407   1 detail::throw_system_error(make_err(errno), "select"); 398   1 detail::throw_system_error(make_err(errno), "select");
HITCBC 408   2 ready = 0; 399   2 ready = 0;
409   } 400   }
410   401  
411   // Process timers outside the lock 402   // Process timers outside the lock
HITCBC 412   282465 timer_svc_->process_expired(); 403   297589 timer_svc_->process_expired();
413   404  
HITCBC 414   282465 ready_queue local_ops; 405   297589 ready_queue local_ops;
415   406  
HITCBC 416   282465 if (ready > 0) 407   297589 if (ready > 0)
417   { 408   {
HITCBC 418   265887 if (FD_ISSET(pipe_fds_[0], &read_fds)) 409   282510 if (FD_ISSET(pipe_fds_[0], &read_fds))
419   { 410   {
420   char buf[256]; 411   char buf[256];
HITCBC 421   11642 while (::read(pipe_fds_[0], buf, sizeof(buf)) > 0) 412   11514 while (::read(pipe_fds_[0], buf, sizeof(buf)) > 0)
422   { 413   {
423   } 414   }
424   } 415   }
425   416  
HITCBC 426   681388 for (int i = 0; i < snapshot_count; ++i) 417   733214 for (int i = 0; i < snapshot_count; ++i)
427   { 418   {
HITCBC 428   415501 int fd = snapshot[i].fd; 419   450704 int fd = snapshot[i].fd;
HITCBC 429   415501 reactor_descriptor_state* desc = snapshot[i].desc; 420   450704 reactor_descriptor_state* desc = snapshot[i].desc;
430   421  
HITCBC 431   415501 std::uint32_t flags = 0; 422   450704 std::uint32_t flags = 0;
HITCBC 432   415501 if (FD_ISSET(fd, &read_fds)) 423   450704 if (FD_ISSET(fd, &read_fds))
HITCBC 433   264795 flags |= reactor_event_read; 424   281273 flags |= reactor_event_read;
HITCBC 434   415501 if (FD_ISSET(fd, &write_fds)) 425   450704 if (FD_ISSET(fd, &write_fds))
HITCBC 435   2169 flags |= reactor_event_write; 426   2102 flags |= reactor_event_write;
HITCBC 436   415501 if (FD_ISSET(fd, &except_fds)) 427   450704 if (FD_ISSET(fd, &except_fds))
HITCBC 437   16 flags |= reactor_event_error; 428   16 flags |= reactor_event_error;
438   429  
HITCBC 439   415501 if (flags == 0) 430   450704 if (flags == 0)
HITCBC 440   148546 continue; 431   167338 continue;
441   432  
HITCBC 442   266955 desc->add_ready_events(flags); 433   283366 desc->add_ready_events(flags);
443   434  
HITCBC 444   266955 bool expected = false; 435   283366 bool expected = false;
HITCBC 445   266955 if (desc->is_enqueued_.compare_exchange_strong( 436   283366 if (desc->is_enqueued_.compare_exchange_strong(
446   expected, true, std::memory_order_release, 437   expected, true, std::memory_order_release,
447   std::memory_order_relaxed)) 438   std::memory_order_relaxed))
448   { 439   {
HITCBC 449   266955 local_ops.push(desc); 440   283366 local_ops.push(desc);
450   } 441   }
451   } 442   }
452   } 443   }
453   444  
HITCBC 454   282465 lock.lock(); 445   297589 lock.lock();
455   446  
HITCBC 456   282465 completed_ops_.splice(local_ops); 447   297589 completed_ops_.splice(local_ops);
HITCBC 457   282466 } 448   297590 }
458   449  
459   } // namespace boost::corosio::detail 450   } // namespace boost::corosio::detail
460   451  
461   #endif // BOOST_COROSIO_HAS_SELECT 452   #endif // BOOST_COROSIO_HAS_SELECT
462   453  
463   #endif // BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP 454   #endif // BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP