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/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   /** Deregister a persistently registered descriptor. 116   /** Deregister a persistently registered descriptor.
117   117  
118   @param fd The file descriptor to deregister. 118   @param fd The file descriptor to deregister.
119   */ 119   */
120   void deregister_descriptor(int fd) const; 120   void deregister_descriptor(int fd) const;
121   121  
122   /** Interrupt the reactor so it rebuilds its fd_sets. 122   /** Interrupt the reactor so it rebuilds its fd_sets.
123   123  
124   Called when a write, connect, or write-wait op is registered 124   Called when a write, connect, or write-wait op is registered
125   after the reactor's snapshot was taken. Without this, 125   after the reactor's snapshot was taken. Without this,
126   select() may block not watching for writability on the fd. 126   select() may block not watching for writability on the fd.
127   */ 127   */
128   void notify_reactor() const; 128   void notify_reactor() const;
129   129  
130   /// 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 131   61 [[nodiscard]] std::error_code register_signal_reader(int read_fd) override 131   62 [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
132   { 132   {
HITCBC 133   61 return register_descriptor(read_fd, signal_pipe_reader_.arm()); 133   62 return register_descriptor(read_fd, signal_pipe_reader_.arm());
134   } 134   }
135   135  
136   private: 136   private:
137   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;
138   void interrupt_reactor() const override; 138   void interrupt_reactor() const override;
139   long calculate_timeout(long requested_timeout_us) const; 139   long calculate_timeout(long requested_timeout_us) const;
140   140  
141   // 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
142   // register_signal_reader on the first signal registration). 142   // register_signal_reader on the first signal registration).
143   reactor_signal_pipe_reader signal_pipe_reader_; 143   reactor_signal_pipe_reader signal_pipe_reader_;
144   144  
145   // Self-pipe for interrupting select() 145   // Self-pipe for interrupting select()
146   int pipe_fds_[2]; // [0]=read, [1]=write 146   int pipe_fds_[2]; // [0]=read, [1]=write
147   147  
148   // Per-fd tracking for fd_set building 148   // Per-fd tracking for fd_set building
149   mutable std::unordered_map<int, reactor_descriptor_state*> 149   mutable std::unordered_map<int, reactor_descriptor_state*>
150   registered_descs_; 150   registered_descs_;
151   mutable int max_fd_ = -1; 151   mutable int max_fd_ = -1;
152   }; 152   };
153   153  
HITCBC 154   950 inline select_scheduler::select_scheduler(capy::execution_context& ctx, int) 154   974 inline select_scheduler::select_scheduler(capy::execution_context& ctx, int)
HITCBC 155   950 : pipe_fds_{-1, -1} 155   974 : pipe_fds_{-1, -1}
HITCBC 156   950 , max_fd_(-1) 156   974 , max_fd_(-1)
157   { 157   {
HITCBC 158   950 if (::pipe(pipe_fds_) < 0) 158   974 if (::pipe(pipe_fds_) < 0)
HITCBC 159   1 detail::throw_system_error(make_err(errno), "pipe"); 159   1 detail::throw_system_error(make_err(errno), "pipe");
160   160  
HITCBC 161   2838 for (int i = 0; i < 2; ++i) 161   2910 for (int i = 0; i < 2; ++i)
162   { 162   {
HITCBC 163   1895 int flags = ::fcntl(pipe_fds_[i], F_GETFL, 0); 163   1943 int flags = ::fcntl(pipe_fds_[i], F_GETFL, 0);
HITCBC 164   1895 if (flags == -1) 164   1943 if (flags == -1)
165   { 165   {
HITCBC 166   2 int errn = errno; 166   2 int errn = errno;
HITCBC 167   2 ::close(pipe_fds_[0]); 167   2 ::close(pipe_fds_[0]);
HITCBC 168   2 ::close(pipe_fds_[1]); 168   2 ::close(pipe_fds_[1]);
HITCBC 169   2 detail::throw_system_error(make_err(errn), "fcntl F_GETFL"); 169   2 detail::throw_system_error(make_err(errn), "fcntl F_GETFL");
170   } 170   }
HITCBC 171   1893 if (::fcntl(pipe_fds_[i], F_SETFL, flags | O_NONBLOCK) == -1) 171   1941 if (::fcntl(pipe_fds_[i], F_SETFL, flags | O_NONBLOCK) == -1)
172   { 172   {
HITCBC 173   2 int errn = errno; 173   2 int errn = errno;
HITCBC 174   2 ::close(pipe_fds_[0]); 174   2 ::close(pipe_fds_[0]);
HITCBC 175   2 ::close(pipe_fds_[1]); 175   2 ::close(pipe_fds_[1]);
HITCBC 176   2 detail::throw_system_error(make_err(errn), "fcntl F_SETFL"); 176   2 detail::throw_system_error(make_err(errn), "fcntl F_SETFL");
177   } 177   }
HITCBC 178   1891 if (::fcntl(pipe_fds_[i], F_SETFD, FD_CLOEXEC) == -1) 178   1939 if (::fcntl(pipe_fds_[i], F_SETFD, FD_CLOEXEC) == -1)
179   { 179   {
HITCBC 180   2 int errn = errno; 180   2 int errn = errno;
HITCBC 181   2 ::close(pipe_fds_[0]); 181   2 ::close(pipe_fds_[0]);
HITCBC 182   2 ::close(pipe_fds_[1]); 182   2 ::close(pipe_fds_[1]);
HITCBC 183   2 detail::throw_system_error(make_err(errn), "fcntl F_SETFD"); 183   2 detail::throw_system_error(make_err(errn), "fcntl F_SETFD");
184   } 184   }
185   } 185   }
186   186  
HITCBC 187   943 timer_svc_ = &get_timer_service(ctx, *this); 187   967 timer_svc_ = &get_timer_service(ctx, *this);
HITCBC 188   943 timer_svc_->set_on_earliest_changed( 188   967 timer_svc_->set_on_earliest_changed(
HITCBC 189   3778 timer_service::callback(this, [](void* p) { 189   3751 timer_service::callback(this, [](void* p) {
HITCBC 190   2835 static_cast<select_scheduler*>(p)->interrupt_reactor(); 190   2784 static_cast<select_scheduler*>(p)->interrupt_reactor();
HITCBC 191   2835 })); 191   2784 }));
192   192  
HITCBC 193   943 completed_ops_.push(&task_op_); 193   967 completed_ops_.push(&task_op_);
HITCBC 194   964 } 194   988 }
195   195  
HITCBC 196   1886 inline select_scheduler::~select_scheduler() 196   1934 inline select_scheduler::~select_scheduler()
197   { 197   {
HITCBC 198   943 if (pipe_fds_[0] >= 0) 198   967 if (pipe_fds_[0] >= 0)
HITCBC 199   943 ::close(pipe_fds_[0]); 199   967 ::close(pipe_fds_[0]);
HITCBC 200   943 if (pipe_fds_[1] >= 0) 200   967 if (pipe_fds_[1] >= 0)
HITCBC 201   943 ::close(pipe_fds_[1]); 201   967 ::close(pipe_fds_[1]);
HITCBC 202   1886 } 202   1934 }
203   203  
204   inline void 204   inline void
HITCBC 205   943 select_scheduler::shutdown() 205   967 select_scheduler::shutdown()
206   { 206   {
HITCBC 207   943 shutdown_drain(); 207   967 shutdown_drain();
208   208  
HITCBC 209   943 if (pipe_fds_[1] >= 0) 209   967 if (pipe_fds_[1] >= 0)
HITCBC 210   943 interrupt_reactor(); 210   967 interrupt_reactor();
HITCBC 211   943 } 211   967 }
212   212  
213   inline std::error_code 213   inline std::error_code
HITCBC 214   5060 select_scheduler::register_descriptor( 214   5146 select_scheduler::register_descriptor(
215   int fd, reactor_descriptor_state* desc) const 215   int fd, reactor_descriptor_state* desc) const
216   { 216   {
HITCBC 217   5060 if (fd < 0 || fd >= FD_SETSIZE) 217   5146 if (fd < 0 || fd >= FD_SETSIZE)
HITCBC 218   1 return make_err(EMFILE); 218   1 return make_err(EMFILE);
219   219  
HITCBC 220   5059 desc->registered_events = reactor_event_read | reactor_event_write; 220   5145 desc->registered_events = reactor_event_read | reactor_event_write;
HITCBC 221   5059 desc->fd = fd; 221   5145 desc->fd = fd;
HITCBC 222   5059 desc->scheduler_ = this; 222   5145 desc->scheduler_ = this;
HITCBC 223   5059 desc->mutex.set_enabled(reactor_io_locking_); 223   5145 desc->mutex.set_enabled(reactor_io_locking_);
HITCBC 224   5059 desc->ready_events_.store(0, std::memory_order_relaxed); 224   5145 desc->ready_events_.store(0, std::memory_order_relaxed);
225   225  
226   { 226   {
HITCBC 227   5059 conditionally_enabled_mutex::scoped_lock lock(desc->mutex); 227   5145 conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
HITCBC 228 - 5059 desc->impl_ref_.reset(); 228 + 5145 desc->object_ref_.reset();
HITCBC 229   5059 desc->read_ready = false; 229   5145 desc->read_ready = false;
HITCBC 230   5059 desc->write_ready = false; 230   5145 desc->write_ready = false;
HITCBC 231   5059 } 231   5145 }
232   232  
233   { 233   {
HITCBC 234   5059 mutex_type::scoped_lock lock(mutex_); 234   5145 mutex_type::scoped_lock lock(mutex_);
235   try 235   try
236   { 236   {
HITCBC 237   5059 registered_descs_[fd] = desc; 237   5145 registered_descs_[fd] = desc;
238   } 238   }
HITCBC 239   1 catch (std::bad_alloc const&) 239   1 catch (std::bad_alloc const&)
240   { 240   {
HITCBC 241   1 return make_err(ENOMEM); 241   1 return make_err(ENOMEM);
HITCBC 242   1 } 242   1 }
HITCBC 243   5058 if (fd > max_fd_) 243   5144 if (fd > max_fd_)
HITCBC 244   5008 max_fd_ = fd; 244   5081 max_fd_ = fd;
HITCBC 245   5059 } 245   5145 }
246   246  
HITCBC 247   5058 interrupt_reactor(); 247   5144 interrupt_reactor();
HITCBC 248   5058 return {}; 248   5144 return {};
249   } 249   }
250   250  
251   inline void 251   inline void
HITCBC 252   4998 select_scheduler::deregister_descriptor(int fd) const 252   5083 select_scheduler::deregister_descriptor(int fd) const
253   { 253   {
HITCBC 254   4998 mutex_type::scoped_lock lock(mutex_); 254   5083 mutex_type::scoped_lock lock(mutex_);
255   255  
HITCBC 256   4998 auto it = registered_descs_.find(fd); 256   5083 auto it = registered_descs_.find(fd);
HITCBC 257   4998 if (it == registered_descs_.end()) 257   5083 if (it == registered_descs_.end())
MISUBC 258   ✗ return; 258   ✗ return;
259   259  
HITCBC 260   4998 registered_descs_.erase(it); 260   5083 registered_descs_.erase(it);
261   261  
HITCBC 262   4998 if (fd == max_fd_) 262   5083 if (fd == max_fd_)
263   { 263   {
HITCBC 264   4663 max_fd_ = pipe_fds_[0]; 264   4687 max_fd_ = pipe_fds_[0];
HITCBC 265   8920 for (auto& [registered_fd, state] : registered_descs_) 265   8967 for (auto& [registered_fd, state] : registered_descs_)
266   { 266   {
HITCBC 267   4257 if (registered_fd > max_fd_) 267   4280 if (registered_fd > max_fd_)
HITCBC 268   4164 max_fd_ = registered_fd; 268   4170 max_fd_ = registered_fd;
269   } 269   }
270   } 270   }
HITCBC 271   4998 } 271   5083 }
272   272  
273   inline void 273   inline void
HITCBC 274   2258 select_scheduler::notify_reactor() const 274   2282 select_scheduler::notify_reactor() const
275   { 275   {
HITCBC 276   2258 interrupt_reactor(); 276   2282 interrupt_reactor();
HITCBC 277   2258 } 277   2282 }
278   278  
279   inline void 279   inline void
HITCBC 280   12863 select_scheduler::interrupt_reactor() const 280   13111 select_scheduler::interrupt_reactor() const
281   { 281   {
HITCBC 282   12863 char byte = 1; 282   13111 char byte = 1;
HITCBC 283   12863 [[maybe_unused]] auto r = ::write(pipe_fds_[1], &byte, 1); 283   13111 [[maybe_unused]] auto r = ::write(pipe_fds_[1], &byte, 1);
HITCBC 284   12863 } 284   13111 }
285   285  
286   inline long 286   inline long
HITCBC 287   270283 select_scheduler::calculate_timeout(long requested_timeout_us) const 287   253504 select_scheduler::calculate_timeout(long requested_timeout_us) const
288   { 288   {
HITCBC 289   270283 if (requested_timeout_us == 0) 289   253504 if (requested_timeout_us == 0)
290   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
291   291  
HITCBC 292   270283 auto nearest = timer_svc_->nearest_expiry(); 292   253504 auto nearest = timer_svc_->nearest_expiry();
HITCBC 293   270283 if (nearest == timer_service::time_point::max()) 293   253504 if (nearest == timer_service::time_point::max())
HITCBC 294   1207 return requested_timeout_us; 294   1336 return requested_timeout_us;
295   295  
HITCBC 296   269076 auto now = std::chrono::steady_clock::now(); 296   252168 auto now = std::chrono::steady_clock::now();
HITCBC 297   269076 if (nearest <= now) 297   252168 if (nearest <= now)
HITCBC 298   687 return 0; 298   106 return 0;
299   299  
300   auto timer_timeout_us = 300   auto timer_timeout_us =
HITCBC 301   268389 std::chrono::duration_cast<std::chrono::microseconds>(nearest - now) 301   252062 std::chrono::duration_cast<std::chrono::microseconds>(nearest - now)
HITCBC 302   268389 .count(); 302   252062 .count();
303   303  
HITCBC 304   268389 constexpr auto long_max = 304   252062 constexpr auto long_max =
305   static_cast<long long>((std::numeric_limits<long>::max)()); 305   static_cast<long long>((std::numeric_limits<long>::max)());
306   auto capped_timer_us = 306   auto capped_timer_us =
HITCBC 307   268389 (std::min)((std::max)(static_cast<long long>(timer_timeout_us), 307   252062 (std::min)((std::max)(static_cast<long long>(timer_timeout_us),
HITCBC 308   268389 static_cast<long long>(0)), 308   252062 static_cast<long long>(0)),
HITCBC 309   268389 long_max); 309   252062 long_max);
310   310  
HITCBC 311   268389 if (requested_timeout_us < 0) 311   252062 if (requested_timeout_us < 0)
HITCBC 312   268387 return static_cast<long>(capped_timer_us); 312   252060 return static_cast<long>(capped_timer_us);
313   313  
314   return static_cast<long>( 314   return static_cast<long>(
HITCBC 315   2 (std::min)(static_cast<long long>(requested_timeout_us), 315   2 (std::min)(static_cast<long long>(requested_timeout_us),
HITCBC 316   2 capped_timer_us)); 316   2 capped_timer_us));
317   } 317   }
318   318  
319   inline void 319   inline void
HITCBC 320   295281 select_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us) 320   279931 select_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
321   { 321   {
322   long effective_timeout_us = 322   long effective_timeout_us =
HITCBC 323   295281 task_interrupted_ ? 0 : calculate_timeout(timeout_us); 323   279931 task_interrupted_ ? 0 : calculate_timeout(timeout_us);
324   324  
325   // Snapshot registered descriptors while holding lock. 325   // Snapshot registered descriptors while holding lock.
326   // Record which fds need write monitoring to avoid a hot loop: 326   // Record which fds need write monitoring to avoid a hot loop:
327   // select is level-triggered so writable sockets (nearly always 327   // select is level-triggered so writable sockets (nearly always
328   // writable) would cause select() to return immediately every 328   // writable) would cause select() to return immediately every
329   // iteration if unconditionally added to write_fds. Membership 329   // iteration if unconditionally added to write_fds. Membership
330   // 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
331   // parked write or connect op does. 331   // parked write or connect op does.
332   struct fd_entry 332   struct fd_entry
333   { 333   {
334   int fd; 334   int fd;
335   reactor_descriptor_state* desc; 335   reactor_descriptor_state* desc;
336   bool needs_write; 336   bool needs_write;
337   }; 337   };
338   fd_entry snapshot[FD_SETSIZE]; 338   fd_entry snapshot[FD_SETSIZE];
HITCBC 339   295281 int snapshot_count = 0; 339   279931 int snapshot_count = 0;
340   340  
HITCBC 341   774219 for (auto& [fd, desc] : registered_descs_) 341   750294 for (auto& [fd, desc] : registered_descs_)
342   { 342   {
HITCBC 343   478938 if (snapshot_count < FD_SETSIZE) 343   470363 if (snapshot_count < FD_SETSIZE)
344   { 344   {
HITCBC 345   478938 conditionally_enabled_mutex::scoped_lock desc_lock(desc->mutex); 345   470363 conditionally_enabled_mutex::scoped_lock desc_lock(desc->mutex);
HITCBC 346   478938 snapshot[snapshot_count].fd = fd; 346   470363 snapshot[snapshot_count].fd = fd;
HITCBC 347   478938 snapshot[snapshot_count].desc = desc; 347   470363 snapshot[snapshot_count].desc = desc;
HITCBC 348   478938 snapshot[snapshot_count].needs_write = 348   470363 snapshot[snapshot_count].needs_write =
HITCBC 349   478938 (desc->write_op || desc->connect_op || desc->wait_write_op); 349   470363 (desc->write_op || desc->connect_op || desc->wait_write_op);
HITCBC 350   478938 ++snapshot_count; 350   470363 ++snapshot_count;
HITCBC 351   478938 } 351   470363 }
352   } 352   }
353   353  
HITCBC 354   295281 if (lock.owns_lock()) 354   279931 if (lock.owns_lock())
HITCBC 355   270284 lock.unlock(); 355   253505 lock.unlock();
356   356  
HITCBC 357   295281 task_cleanup on_exit{this, &lock, ctx}; 357   279931 task_cleanup on_exit{this, &lock, ctx};
358   358  
359   fd_set read_fds, write_fds, except_fds; 359   fd_set read_fds, write_fds, except_fds;
HITCBC 360   5019777 FD_ZERO(&read_fds); 360   4758827 FD_ZERO(&read_fds);
HITCBC 361   5019777 FD_ZERO(&write_fds); 361   4758827 FD_ZERO(&write_fds);
HITCBC 362   5019777 FD_ZERO(&except_fds); 362   4758827 FD_ZERO(&except_fds);
363   363  
HITCBC 364   295281 FD_SET(pipe_fds_[0], &read_fds); 364   279931 FD_SET(pipe_fds_[0], &read_fds);
HITCBC 365   295281 int nfds = pipe_fds_[0]; 365   279931 int nfds = pipe_fds_[0];
366   366  
HITCBC 367   774219 for (int i = 0; i < snapshot_count; ++i) 367   750294 for (int i = 0; i < snapshot_count; ++i)
368   { 368   {
HITCBC 369   478938 int fd = snapshot[i].fd; 369   470363 int fd = snapshot[i].fd;
HITCBC 370   478938 FD_SET(fd, &read_fds); 370   470363 FD_SET(fd, &read_fds);
HITCBC 371   478938 if (snapshot[i].needs_write) 371   470363 if (snapshot[i].needs_write)
HITCBC 372   13408 FD_SET(fd, &write_fds); 372   8845 FD_SET(fd, &write_fds);
HITCBC 373   478938 FD_SET(fd, &except_fds); 373   470363 FD_SET(fd, &except_fds);
HITCBC 374   478938 if (fd > nfds) 374   470363 if (fd > nfds)
HITCBC 375   294636 nfds = fd; 375   279220 nfds = fd;
376   } 376   }
377   377  
378   struct timeval tv; 378   struct timeval tv;
HITCBC 379   295281 struct timeval* tv_ptr = nullptr; 379   279931 struct timeval* tv_ptr = nullptr;
HITCBC 380   295281 if (effective_timeout_us >= 0) 380   279931 if (effective_timeout_us >= 0)
381   { 381   {
HITCBC 382   294453 tv.tv_sec = effective_timeout_us / 1000000; 382   279035 tv.tv_sec = effective_timeout_us / 1000000;
HITCBC 383   294453 tv.tv_usec = effective_timeout_us % 1000000; 383   279035 tv.tv_usec = effective_timeout_us % 1000000;
HITCBC 384   294453 tv_ptr = &tv; 384   279035 tv_ptr = &tv;
385   } 385   }
386   386  
HITCBC 387   295281 int ready = ::select(nfds + 1, &read_fds, &write_fds, &except_fds, tv_ptr); 387   279931 int ready = ::select(nfds + 1, &read_fds, &write_fds, &except_fds, tv_ptr);
388   388  
389   // EINTR: signal interrupted select(), just retry. 389   // EINTR: signal interrupted select(), just retry.
390   // EBADF: an fd was closed between snapshot and select(); retry 390   // EBADF: an fd was closed between snapshot and select(); retry
391   // with a fresh snapshot from registered_descs_. 391   // with a fresh snapshot from registered_descs_.
392   // Both fall through with no ready descriptors rather than 392   // Both fall through with no ready descriptors rather than
393   // returning: the caller handed this function an owned lock that 393   // returning: the caller handed this function an owned lock that
394   // only the epilogue below re-acquires. 394   // only the epilogue below re-acquires.
HITCBC 395   295281 if (ready < 0) 395   279931 if (ready < 0)
396   { 396   {
HITCBC 397   3 if (errno != EINTR && errno != EBADF) 397   3 if (errno != EINTR && errno != EBADF)
HITCBC 398   1 detail::throw_system_error(make_err(errno), "select"); 398   1 detail::throw_system_error(make_err(errno), "select");
HITCBC 399   2 ready = 0; 399   2 ready = 0;
400   } 400   }
401   401  
402   // Process timers outside the lock 402   // Process timers outside the lock
HITCBC 403   295280 timer_svc_->process_expired(); 403   279930 timer_svc_->process_expired();
404   404  
HITCBC 405   295280 ready_queue local_ops; 405   279930 ready_queue local_ops;
406   406  
HITCBC 407   295280 if (ready > 0) 407   279930 if (ready > 0)
408   { 408   {
HITCBC 409   278637 if (FD_ISSET(pipe_fds_[0], &read_fds)) 409   262074 if (FD_ISSET(pipe_fds_[0], &read_fds))
410   { 410   {
411   char buf[256]; 411   char buf[256];
HITCBC 412   12282 while (::read(pipe_fds_[0], buf, sizeof(buf)) > 0) 412   12670 while (::read(pipe_fds_[0], buf, sizeof(buf)) > 0)
413   { 413   {
414   } 414   }
415   } 415   }
416   416  
HITCBC 417   709488 for (int i = 0; i < snapshot_count; ++i) 417   681491 for (int i = 0; i < snapshot_count; ++i)
418   { 418   {
HITCBC 419   430851 int fd = snapshot[i].fd; 419   419417 int fd = snapshot[i].fd;
HITCBC 420   430851 reactor_descriptor_state* desc = snapshot[i].desc; 420   419417 reactor_descriptor_state* desc = snapshot[i].desc;
421   421  
HITCBC 422   430851 std::uint32_t flags = 0; 422   419417 std::uint32_t flags = 0;
HITCBC 423   430851 if (FD_ISSET(fd, &read_fds)) 423   419417 if (FD_ISSET(fd, &read_fds))
HITCBC 424   277266 flags |= reactor_event_read; 424   260574 flags |= reactor_event_read;
HITCBC 425   430851 if (FD_ISSET(fd, &write_fds)) 425   419417 if (FD_ISSET(fd, &write_fds))
HITCBC 426   2250 flags |= reactor_event_write; 426   2274 flags |= reactor_event_write;
HITCBC 427   430851 if (FD_ISSET(fd, &except_fds)) 427   419417 if (FD_ISSET(fd, &except_fds))
HITCBC 428   16 flags |= reactor_event_error; 428   16 flags |= reactor_event_error;
429   429  
HITCBC 430   430851 if (flags == 0) 430   419417 if (flags == 0)
HITCBC 431   151343 continue; 431   156578 continue;
432   432  
HITCBC 433   279508 desc->add_ready_events(flags); 433   262839 desc->add_ready_events(flags);
434   434  
HITCBC 435   279508 bool expected = false; 435   262839 bool expected = false;
HITCBC 436   279508 if (desc->is_enqueued_.compare_exchange_strong( 436   262839 if (desc->is_enqueued_.compare_exchange_strong(
437   expected, true, std::memory_order_release, 437   expected, true, std::memory_order_release,
438   std::memory_order_relaxed)) 438   std::memory_order_relaxed))
439   { 439   {
HITCBC 440   279508 local_ops.push(desc); 440   262839 local_ops.push(desc);
441   } 441   }
442   } 442   }
443   } 443   }
444   444  
HITCBC 445   295280 lock.lock(); 445   279930 lock.lock();
446   446  
HITCBC 447   295280 completed_ops_.splice(local_ops); 447   279930 completed_ops_.splice(local_ops);
HITCBC 448   295281 } 448   279931 }
449   449  
450   } // namespace boost::corosio::detail 450   } // namespace boost::corosio::detail
451   451  
452   #endif // BOOST_COROSIO_HAS_SELECT 452   #endif // BOOST_COROSIO_HAS_SELECT
453   453  
454   #endif // BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP 454   #endif // BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP