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