85.29% Lines (145/170) 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   #include <boost/corosio/native/detail/posix/posix_resolver_service.hpp> 27   #include <boost/corosio/native/detail/posix/posix_resolver_service.hpp>
28   #include <boost/corosio/native/detail/posix/posix_signal_service.hpp> 28   #include <boost/corosio/native/detail/posix/posix_signal_service.hpp>
29   #include <boost/corosio/native/detail/posix/posix_stream_file_service.hpp> 29   #include <boost/corosio/native/detail/posix/posix_stream_file_service.hpp>
30   #include <boost/corosio/native/detail/posix/posix_random_access_file_service.hpp> 30   #include <boost/corosio/native/detail/posix/posix_random_access_file_service.hpp>
31   31  
32   #include <boost/corosio/detail/except.hpp> 32   #include <boost/corosio/detail/except.hpp>
33   33  
34   #include <sys/select.h> 34   #include <sys/select.h>
35   #include <unistd.h> 35   #include <unistd.h>
36   #include <errno.h> 36   #include <errno.h>
37   #include <fcntl.h> 37   #include <fcntl.h>
38   38  
39   #include <atomic> 39   #include <atomic>
40   #include <chrono> 40   #include <chrono>
41   #include <cstdint> 41   #include <cstdint>
42   #include <limits> 42   #include <limits>
43   #include <mutex> 43   #include <mutex>
44   #include <new> 44   #include <new>
45   #include <unordered_map> 45   #include <unordered_map>
46   46  
47   namespace boost::corosio::detail { 47   namespace boost::corosio::detail {
48   48  
49   struct select_op; 49   struct select_op;
50   50  
51   /** POSIX scheduler using select() for I/O multiplexing. 51   /** POSIX scheduler using select() for I/O multiplexing.
52   52  
53   This scheduler implements the scheduler interface using the POSIX select() 53   This scheduler implements the scheduler interface using the POSIX select()
54   call for I/O event notification. It inherits the shared reactor threading 54   call for I/O event notification. It inherits the shared reactor threading
55   model from reactor_scheduler: signal state machine, inline completion 55   model from reactor_scheduler: signal state machine, inline completion
56   budget, work counting, and the do_one event loop. 56   budget, work counting, and the do_one event loop.
57   57  
58   The design mirrors epoll_scheduler for behavioral consistency: 58   The design mirrors epoll_scheduler for behavioral consistency:
59   - Same single-reactor thread coordination model 59   - Same single-reactor thread coordination model
60   - Same deferred I/O pattern (reactor marks ready; workers do I/O) 60   - Same deferred I/O pattern (reactor marks ready; workers do I/O)
61   - Same timer integration pattern 61   - Same timer integration pattern
62   62  
63   Known Limitations: 63   Known Limitations:
64   - FD_SETSIZE (~1024) limits maximum concurrent connections 64   - FD_SETSIZE (~1024) limits maximum concurrent connections
65   - O(n) scanning: rebuilds fd_sets each iteration 65   - O(n) scanning: rebuilds fd_sets each iteration
66   - Level-triggered only (no edge-triggered mode) 66   - Level-triggered only (no edge-triggered mode)
67   67  
68   @par Thread Safety 68   @par Thread Safety
69   All public member functions are thread-safe. 69   All public member functions are thread-safe.
70   */ 70   */
71   class BOOST_COROSIO_DECL select_scheduler final : public reactor_scheduler 71   class BOOST_COROSIO_DECL select_scheduler final : public reactor_scheduler
72   { 72   {
73   public: 73   public:
74   /** Construct the scheduler. 74   /** Construct the scheduler.
75   75  
76   Creates a self-pipe for reactor interruption. 76   Creates a self-pipe for reactor interruption.
77   77  
78   @param ctx Reference to the owning execution_context. 78   @param ctx Reference to the owning execution_context.
79   @param concurrency_hint Hint for expected thread count (unused). 79   @param concurrency_hint Hint for expected thread count (unused).
80   */ 80   */
81   select_scheduler(capy::execution_context& ctx, int concurrency_hint = -1); 81   select_scheduler(capy::execution_context& ctx, int concurrency_hint = -1);
82   82  
83   /// Destroy the scheduler. 83   /// Destroy the scheduler.
84   ~select_scheduler() override; 84   ~select_scheduler() override;
85   85  
86   select_scheduler(select_scheduler const&) = delete; 86   select_scheduler(select_scheduler const&) = delete;
87   select_scheduler& operator=(select_scheduler const&) = delete; 87   select_scheduler& operator=(select_scheduler const&) = delete;
88   88  
89   /// Shut down the scheduler, draining pending operations. 89   /// Shut down the scheduler, draining pending operations.
90   void shutdown() override; 90   void shutdown() override;
91   91  
92   /** Return the maximum file descriptor value supported. 92   /** Return the maximum file descriptor value supported.
93   93  
94   Returns FD_SETSIZE - 1, the maximum fd value that can be 94   Returns FD_SETSIZE - 1, the maximum fd value that can be
95   monitored by select(). Operations with fd >= FD_SETSIZE 95   monitored by select(). Operations with fd >= FD_SETSIZE
96   will fail with EINVAL. 96   will fail with EINVAL.
97   97  
98   @return The maximum supported file descriptor value. 98   @return The maximum supported file descriptor value.
99   */ 99   */
100   static constexpr int max_fd() noexcept 100   static constexpr int max_fd() noexcept
101   { 101   {
102   return FD_SETSIZE - 1; 102   return FD_SETSIZE - 1;
103   } 103   }
104   104  
105   /** Register a descriptor for persistent monitoring. 105   /** Register a descriptor for persistent monitoring.
106   106  
107   The fd is added to the registered_descs_ map and will be 107   The fd is added to the registered_descs_ map and will be
108   included in subsequent select() calls. The reactor is 108   included in subsequent select() calls. The reactor is
109   interrupted so a blocked select() rebuilds its fd_sets. 109   interrupted so a blocked select() rebuilds its fd_sets.
110   110  
111   @param fd The file descriptor to register. 111   @param fd The file descriptor to register.
112   @param desc Pointer to descriptor state for this fd. 112   @param desc Pointer to descriptor state for this fd.
113   113  
114   @return The error if the fd cannot be tracked, otherwise a 114   @return The error if the fd cannot be tracked, otherwise a
115   default constructed error code. 115   default constructed error code.
116   */ 116   */
117   std::error_code 117   std::error_code
118   register_descriptor(int fd, reactor_descriptor_state* desc) const; 118   register_descriptor(int fd, reactor_descriptor_state* desc) const;
119   119  
120   /** Deregister a persistently registered descriptor. 120   /** Deregister a persistently registered descriptor.
121   121  
122   @param fd The file descriptor to deregister. 122   @param fd The file descriptor to deregister.
123   */ 123   */
124   void deregister_descriptor(int fd) const; 124   void deregister_descriptor(int fd) const;
125   125  
126   /** Interrupt the reactor so it rebuilds its fd_sets. 126   /** Interrupt the reactor so it rebuilds its fd_sets.
127   127  
128   Called when a write, connect, or write-wait op is registered 128   Called when a write, connect, or write-wait op is registered
129   after the reactor's snapshot was taken. Without this, 129   after the reactor's snapshot was taken. Without this,
130   select() may block not watching for writability on the fd. 130   select() may block not watching for writability on the fd.
131   */ 131   */
132   void notify_reactor() const; 132   void notify_reactor() const;
133   133  
134   /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp). 134   /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp).
ECB 135 - 41 void register_signal_reader(int read_fd) override 135 + [[nodiscard]] std::error_code
HITGNC   136 + 41 register_signal_reader(int read_fd) override
136   { 137   {
HITCBC 137 - 41 if (auto ec = register_descriptor(read_fd, signal_pipe_reader_.arm())) 138 + 41 return register_descriptor(read_fd, signal_pipe_reader_.arm());
DUB 138 - detail::throw_system_error(ec, "select: register");  
ECB 139   41 } 139   }
140   140  
141   private: 141   private:
142   void 142   void
143   run_task(lock_type& lock, context_type* ctx, 143   run_task(lock_type& lock, context_type* ctx,
144   long timeout_us) override; 144   long timeout_us) override;
145   void interrupt_reactor() const override; 145   void interrupt_reactor() const override;
146   long calculate_timeout(long requested_timeout_us) const; 146   long calculate_timeout(long requested_timeout_us) const;
147   147  
148   // 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
149   // register_signal_reader on the first signal registration). 149   // register_signal_reader on the first signal registration).
150   reactor_signal_pipe_reader signal_pipe_reader_; 150   reactor_signal_pipe_reader signal_pipe_reader_;
151   151  
152   // Self-pipe for interrupting select() 152   // Self-pipe for interrupting select()
153   int pipe_fds_[2]; // [0]=read, [1]=write 153   int pipe_fds_[2]; // [0]=read, [1]=write
154   154  
155   // Per-fd tracking for fd_set building 155   // Per-fd tracking for fd_set building
156   mutable std::unordered_map<int, reactor_descriptor_state*> registered_descs_; 156   mutable std::unordered_map<int, reactor_descriptor_state*> registered_descs_;
157   mutable int max_fd_ = -1; 157   mutable int max_fd_ = -1;
158   }; 158   };
159   159  
HITCBC 160   666 inline select_scheduler::select_scheduler(capy::execution_context& ctx, int) 160   701 inline select_scheduler::select_scheduler(capy::execution_context& ctx, int)
HITCBC 161   666 : pipe_fds_{-1, -1} 161   701 : pipe_fds_{-1, -1}
HITCBC 162   666 , max_fd_(-1) 162   701 , max_fd_(-1)
163   { 163   {
HITCBC 164   666 if (::pipe(pipe_fds_) < 0) 164   701 if (::pipe(pipe_fds_) < 0)
MISUBC 165   detail::throw_system_error(make_err(errno), "pipe"); 165   detail::throw_system_error(make_err(errno), "pipe");
166   166  
HITCBC 167   1998 for (int i = 0; i < 2; ++i) 167   2103 for (int i = 0; i < 2; ++i)
168   { 168   {
HITCBC 169   1332 int flags = ::fcntl(pipe_fds_[i], F_GETFL, 0); 169   1402 int flags = ::fcntl(pipe_fds_[i], F_GETFL, 0);
HITCBC 170   1332 if (flags == -1) 170   1402 if (flags == -1)
171   { 171   {
MISUBC 172   int errn = errno; 172   int errn = errno;
MISUBC 173   ::close(pipe_fds_[0]); 173   ::close(pipe_fds_[0]);
MISUBC 174   ::close(pipe_fds_[1]); 174   ::close(pipe_fds_[1]);
MISUBC 175   detail::throw_system_error(make_err(errn), "fcntl F_GETFL"); 175   detail::throw_system_error(make_err(errn), "fcntl F_GETFL");
176   } 176   }
HITCBC 177   1332 if (::fcntl(pipe_fds_[i], F_SETFL, flags | O_NONBLOCK) == -1) 177   1402 if (::fcntl(pipe_fds_[i], F_SETFL, flags | O_NONBLOCK) == -1)
178   { 178   {
MISUBC 179   int errn = errno; 179   int errn = errno;
MISUBC 180   ::close(pipe_fds_[0]); 180   ::close(pipe_fds_[0]);
MISUBC 181   ::close(pipe_fds_[1]); 181   ::close(pipe_fds_[1]);
MISUBC 182   detail::throw_system_error(make_err(errn), "fcntl F_SETFL"); 182   detail::throw_system_error(make_err(errn), "fcntl F_SETFL");
183   } 183   }
HITCBC 184   1332 if (::fcntl(pipe_fds_[i], F_SETFD, FD_CLOEXEC) == -1) 184   1402 if (::fcntl(pipe_fds_[i], F_SETFD, FD_CLOEXEC) == -1)
185   { 185   {
MISUBC 186   int errn = errno; 186   int errn = errno;
MISUBC 187   ::close(pipe_fds_[0]); 187   ::close(pipe_fds_[0]);
MISUBC 188   ::close(pipe_fds_[1]); 188   ::close(pipe_fds_[1]);
MISUBC 189   detail::throw_system_error(make_err(errn), "fcntl F_SETFD"); 189   detail::throw_system_error(make_err(errn), "fcntl F_SETFD");
190   } 190   }
191   } 191   }
192   192  
HITCBC 193   666 timer_svc_ = &get_timer_service(ctx, *this); 193   701 timer_svc_ = &get_timer_service(ctx, *this);
HITCBC 194   666 timer_svc_->set_on_earliest_changed( 194   701 timer_svc_->set_on_earliest_changed(
HITCBC 195   3805 timer_service::callback(this, [](void* p) { 195   4882 timer_service::callback(this, [](void* p) {
HITCBC 196   3139 static_cast<select_scheduler*>(p)->interrupt_reactor(); 196   4181 static_cast<select_scheduler*>(p)->interrupt_reactor();
HITCBC 197   3139 })); 197   4181 }));
198   198  
HITCBC 199   666 get_resolver_service(ctx, *this); 199   701 get_resolver_service(ctx, *this);
HITCBC 200   666 get_signal_service(ctx, *this); 200   701 get_signal_service(ctx, *this);
HITCBC 201   666 get_stream_file_service(ctx, *this); 201   701 get_stream_file_service(ctx, *this);
HITCBC 202   666 get_random_access_file_service(ctx, *this); 202   701 get_random_access_file_service(ctx, *this);
203   203  
HITCBC 204   666 completed_ops_.push(&task_op_); 204   701 completed_ops_.push(&task_op_);
HITCBC 205   666 } 205   701 }
206   206  
HITCBC 207   1332 inline select_scheduler::~select_scheduler() 207   1402 inline select_scheduler::~select_scheduler()
208   { 208   {
HITCBC 209   666 if (pipe_fds_[0] >= 0) 209   701 if (pipe_fds_[0] >= 0)
HITCBC 210   666 ::close(pipe_fds_[0]); 210   701 ::close(pipe_fds_[0]);
HITCBC 211   666 if (pipe_fds_[1] >= 0) 211   701 if (pipe_fds_[1] >= 0)
HITCBC 212   666 ::close(pipe_fds_[1]); 212   701 ::close(pipe_fds_[1]);
HITCBC 213   1332 } 213   1402 }
214   214  
215   inline void 215   inline void
HITCBC 216   666 select_scheduler::shutdown() 216   701 select_scheduler::shutdown()
217   { 217   {
HITCBC 218   666 shutdown_drain(); 218   701 shutdown_drain();
219   219  
HITCBC 220   666 if (pipe_fds_[1] >= 0) 220   701 if (pipe_fds_[1] >= 0)
HITCBC 221   666 interrupt_reactor(); 221   701 interrupt_reactor();
HITCBC 222   666 } 222   701 }
223   223  
224   inline std::error_code 224   inline std::error_code
HITCBC 225   5143 select_scheduler::register_descriptor( 225   7118 select_scheduler::register_descriptor(
226   int fd, reactor_descriptor_state* desc) const 226   int fd, reactor_descriptor_state* desc) const
227   { 227   {
HITCBC 228   5143 if (fd < 0 || fd >= FD_SETSIZE) 228   7118 if (fd < 0 || fd >= FD_SETSIZE)
MISUBC 229   return make_err(EMFILE); 229   return make_err(EMFILE);
230   230  
HITCBC 231   5143 desc->registered_events = reactor_event_read | reactor_event_write; 231   7118 desc->registered_events = reactor_event_read | reactor_event_write;
HITCBC 232   5143 desc->fd = fd; 232   7118 desc->fd = fd;
HITCBC 233   5143 desc->scheduler_ = this; 233   7118 desc->scheduler_ = this;
HITCBC 234   5143 desc->mutex.set_enabled(reactor_io_locking_); 234   7118 desc->mutex.set_enabled(reactor_io_locking_);
HITCBC 235   5143 desc->ready_events_.store(0, std::memory_order_relaxed); 235   7118 desc->ready_events_.store(0, std::memory_order_relaxed);
236   236  
237   { 237   {
HITCBC 238   5143 conditionally_enabled_mutex::scoped_lock lock(desc->mutex); 238   7118 conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
HITCBC 239   5143 desc->impl_ref_.reset(); 239   7118 desc->impl_ref_.reset();
HITCBC 240   5143 desc->read_ready = false; 240   7118 desc->read_ready = false;
HITCBC 241   5143 desc->write_ready = false; 241   7118 desc->write_ready = false;
HITCBC 242   5143 } 242   7118 }
243   243  
244   { 244   {
HITCBC 245   5143 mutex_type::scoped_lock lock(mutex_); 245   7118 mutex_type::scoped_lock lock(mutex_);
246   try 246   try
247   { 247   {
HITCBC 248   5143 registered_descs_[fd] = desc; 248   7118 registered_descs_[fd] = desc;
249   } 249   }
MISUBC 250   catch (std::bad_alloc const&) 250   catch (std::bad_alloc const&)
251   { 251   {
MISUBC 252   return make_err(ENOMEM); 252   return make_err(ENOMEM);
MISUBC 253   } 253   }
HITCBC 254   5143 if (fd > max_fd_) 254   7118 if (fd > max_fd_)
HITCBC 255   5133 max_fd_ = fd; 255   7108 max_fd_ = fd;
HITCBC 256   5143 } 256   7118 }
257   257  
HITCBC 258   5143 interrupt_reactor(); 258   7118 interrupt_reactor();
HITCBC 259   5143 return {}; 259   7118 return {};
260   } 260   }
261   261  
262   inline void 262   inline void
HITCBC 263   5102 select_scheduler::deregister_descriptor(int fd) const 263   7077 select_scheduler::deregister_descriptor(int fd) const
264   { 264   {
HITCBC 265   5102 mutex_type::scoped_lock lock(mutex_); 265   7077 mutex_type::scoped_lock lock(mutex_);
266   266  
HITCBC 267   5102 auto it = registered_descs_.find(fd); 267   7077 auto it = registered_descs_.find(fd);
HITCBC 268   5102 if (it == registered_descs_.end()) 268   7077 if (it == registered_descs_.end())
MISUBC 269   return; 269   return;
270   270  
HITCBC 271   5102 registered_descs_.erase(it); 271   7077 registered_descs_.erase(it);
272   272  
HITCBC 273   5102 if (fd == max_fd_) 273   7077 if (fd == max_fd_)
274   { 274   {
HITCBC 275   4946 max_fd_ = pipe_fds_[0]; 275   6921 max_fd_ = pipe_fds_[0];
HITCBC 276   9593 for (auto& [registered_fd, state] : registered_descs_) 276   13533 for (auto& [registered_fd, state] : registered_descs_)
277   { 277   {
HITCBC 278   4647 if (registered_fd > max_fd_) 278   6612 if (registered_fd > max_fd_)
HITCBC 279   4592 max_fd_ = registered_fd; 279   6557 max_fd_ = registered_fd;
280   } 280   }
281   } 281   }
HITCBC 282   5102 } 282   7077 }
283   283  
284   inline void 284   inline void
HITCBC 285   2409 select_scheduler::notify_reactor() const 285   3392 select_scheduler::notify_reactor() const
286   { 286   {
HITCBC 287   2409 interrupt_reactor(); 287   3392 interrupt_reactor();
HITCBC 288   2409 } 288   3392 }
289   289  
290   inline void 290   inline void
HITCBC 291   11903 select_scheduler::interrupt_reactor() const 291   15971 select_scheduler::interrupt_reactor() const
292   { 292   {
HITCBC 293   11903 char byte = 1; 293   15971 char byte = 1;
HITCBC 294   11903 [[maybe_unused]] auto r = ::write(pipe_fds_[1], &byte, 1); 294   15971 [[maybe_unused]] auto r = ::write(pipe_fds_[1], &byte, 1);
HITCBC 295   11903 } 295   15971 }
296   296  
297   inline long 297   inline long
HITCBC 298   275799 select_scheduler::calculate_timeout(long requested_timeout_us) const 298   432368 select_scheduler::calculate_timeout(long requested_timeout_us) const
299   { 299   {
HITCBC 300   275799 if (requested_timeout_us == 0) 300   432368 if (requested_timeout_us == 0)
MISUBC 301   return 0; 301   return 0;
302   302  
HITCBC 303   275799 auto nearest = timer_svc_->nearest_expiry(); 303   432368 auto nearest = timer_svc_->nearest_expiry();
HITCBC 304   275799 if (nearest == timer_service::time_point::max()) 304   432368 if (nearest == timer_service::time_point::max())
HITCBC 305   594 return requested_timeout_us; 305   607 return requested_timeout_us;
306   306  
HITCBC 307   275205 auto now = std::chrono::steady_clock::now(); 307   431761 auto now = std::chrono::steady_clock::now();
HITCBC 308   275205 if (nearest <= now) 308   431761 if (nearest <= now)
HITCBC 309   1286 return 0; 309   637 return 0;
310   310  
311   auto timer_timeout_us = 311   auto timer_timeout_us =
HITCBC 312   273919 std::chrono::duration_cast<std::chrono::microseconds>(nearest - now) 312   431124 std::chrono::duration_cast<std::chrono::microseconds>(nearest - now)
HITCBC 313   273919 .count(); 313   431124 .count();
314   314  
HITCBC 315   273919 constexpr auto long_max = 315   431124 constexpr auto long_max =
316   static_cast<long long>((std::numeric_limits<long>::max)()); 316   static_cast<long long>((std::numeric_limits<long>::max)());
317   auto capped_timer_us = 317   auto capped_timer_us =
HITCBC 318   273919 (std::min)((std::max)(static_cast<long long>(timer_timeout_us), 318   431124 (std::min)((std::max)(static_cast<long long>(timer_timeout_us),
HITCBC 319   273919 static_cast<long long>(0)), 319   431124 static_cast<long long>(0)),
HITCBC 320   273919 long_max); 320   431124 long_max);
321   321  
HITCBC 322   273919 if (requested_timeout_us < 0) 322   431124 if (requested_timeout_us < 0)
HITCBC 323   273919 return static_cast<long>(capped_timer_us); 323   431124 return static_cast<long>(capped_timer_us);
324   324  
325   return static_cast<long>( 325   return static_cast<long>(
MISUBC 326   (std::min)(static_cast<long long>(requested_timeout_us), 326   (std::min)(static_cast<long long>(requested_timeout_us),
MISUBC 327   capped_timer_us)); 327   capped_timer_us));
328   } 328   }
329   329  
330   inline void 330   inline void
HITCBC 331   298386 select_scheduler::run_task( 331   469285 select_scheduler::run_task(
332   lock_type& lock, context_type* ctx, long timeout_us) 332   lock_type& lock, context_type* ctx, long timeout_us)
333   { 333   {
334   long effective_timeout_us = 334   long effective_timeout_us =
HITCBC 335   298386 task_interrupted_ ? 0 : calculate_timeout(timeout_us); 335   469285 task_interrupted_ ? 0 : calculate_timeout(timeout_us);
336   336  
337   // Snapshot registered descriptors while holding lock. 337   // Snapshot registered descriptors while holding lock.
338   // Record which fds need write monitoring to avoid a hot loop: 338   // Record which fds need write monitoring to avoid a hot loop:
339   // select is level-triggered so writable sockets (nearly always 339   // select is level-triggered so writable sockets (nearly always
340   // writable) would cause select() to return immediately every 340   // writable) would cause select() to return immediately every
341   // iteration if unconditionally added to write_fds. Membership 341   // iteration if unconditionally added to write_fds. Membership
342   // stays opt-in: a parked write wait opts in the same way a 342   // stays opt-in: a parked write wait opts in the same way a
343   // parked write or connect op does. 343   // parked write or connect op does.
344   struct fd_entry 344   struct fd_entry
345   { 345   {
346   int fd; 346   int fd;
347   reactor_descriptor_state* desc; 347   reactor_descriptor_state* desc;
348   bool needs_write; 348   bool needs_write;
349   }; 349   };
350   fd_entry snapshot[FD_SETSIZE]; 350   fd_entry snapshot[FD_SETSIZE];
HITCBC 351   298386 int snapshot_count = 0; 351   469285 int snapshot_count = 0;
352   352  
HITCBC 353   723504 for (auto& [fd, desc] : registered_descs_) 353   1178656 for (auto& [fd, desc] : registered_descs_)
354   { 354   {
HITCBC 355   425118 if (snapshot_count < FD_SETSIZE) 355   709371 if (snapshot_count < FD_SETSIZE)
356   { 356   {
HITCBC 357   425118 conditionally_enabled_mutex::scoped_lock desc_lock(desc->mutex); 357   709371 conditionally_enabled_mutex::scoped_lock desc_lock(desc->mutex);
HITCBC 358   425118 snapshot[snapshot_count].fd = fd; 358   709371 snapshot[snapshot_count].fd = fd;
HITCBC 359   425118 snapshot[snapshot_count].desc = desc; 359   709371 snapshot[snapshot_count].desc = desc;
HITCBC 360   425118 snapshot[snapshot_count].needs_write = 360   709371 snapshot[snapshot_count].needs_write =
HITCBC 361   837005 (desc->write_op || desc->connect_op || 361   1404630 (desc->write_op || desc->connect_op ||
HITCBC 362   411887 desc->wait_write_op); 362   695259 desc->wait_write_op);
HITCBC 363   425118 ++snapshot_count; 363   709371 ++snapshot_count;
HITCBC 364   425118 } 364   709371 }
365   } 365   }
366   366  
HITCBC 367   298386 if (lock.owns_lock()) 367   469285 if (lock.owns_lock())
HITCBC 368   275800 lock.unlock(); 368   432369 lock.unlock();
369   369  
HITCBC 370   298386 task_cleanup on_exit{this, &lock, ctx}; 370   469285 task_cleanup on_exit{this, &lock, ctx};
371   371  
372   fd_set read_fds, write_fds, except_fds; 372   fd_set read_fds, write_fds, except_fds;
HITCBC 373   5072562 FD_ZERO(&read_fds); 373   7977845 FD_ZERO(&read_fds);
HITCBC 374   5072562 FD_ZERO(&write_fds); 374   7977845 FD_ZERO(&write_fds);
HITCBC 375   5072562 FD_ZERO(&except_fds); 375   7977845 FD_ZERO(&except_fds);
376   376  
HITCBC 377   298386 FD_SET(pipe_fds_[0], &read_fds); 377   469285 FD_SET(pipe_fds_[0], &read_fds);
HITCBC 378   298386 int nfds = pipe_fds_[0]; 378   469285 int nfds = pipe_fds_[0];
379   379  
HITCBC 380   723504 for (int i = 0; i < snapshot_count; ++i) 380   1178656 for (int i = 0; i < snapshot_count; ++i)
381   { 381   {
HITCBC 382   425118 int fd = snapshot[i].fd; 382   709371 int fd = snapshot[i].fd;
HITCBC 383   425118 FD_SET(fd, &read_fds); 383   709371 FD_SET(fd, &read_fds);
HITCBC 384   425118 if (snapshot[i].needs_write) 384   709371 if (snapshot[i].needs_write)
HITCBC 385   13235 FD_SET(fd, &write_fds); 385   14116 FD_SET(fd, &write_fds);
HITCBC 386   425118 FD_SET(fd, &except_fds); 386   709371 FD_SET(fd, &except_fds);
HITCBC 387   425118 if (fd > nfds) 387   709371 if (fd > nfds)
HITCBC 388   298022 nfds = fd; 388   468863 nfds = fd;
389   } 389   }
390   390  
391   struct timeval tv; 391   struct timeval tv;
HITCBC 392   298386 struct timeval* tv_ptr = nullptr; 392   469285 struct timeval* tv_ptr = nullptr;
HITCBC 393   298386 if (effective_timeout_us >= 0) 393   469285 if (effective_timeout_us >= 0)
394   { 394   {
HITCBC 395   297796 tv.tv_sec = effective_timeout_us / 1000000; 395   468681 tv.tv_sec = effective_timeout_us / 1000000;
HITCBC 396   297796 tv.tv_usec = effective_timeout_us % 1000000; 396   468681 tv.tv_usec = effective_timeout_us % 1000000;
HITCBC 397   297796 tv_ptr = &tv; 397   468681 tv_ptr = &tv;
398   } 398   }
399   399  
HITCBC 400   298386 int ready = ::select(nfds + 1, &read_fds, &write_fds, &except_fds, tv_ptr); 400   469285 int ready = ::select(nfds + 1, &read_fds, &write_fds, &except_fds, tv_ptr);
401   401  
402   // EINTR: signal interrupted select(), just retry. 402   // EINTR: signal interrupted select(), just retry.
403   // EBADF: an fd was closed between snapshot and select(); retry 403   // EBADF: an fd was closed between snapshot and select(); retry
404   // with a fresh snapshot from registered_descs_. 404   // with a fresh snapshot from registered_descs_.
HITCBC 405   298386 if (ready < 0) 405   469285 if (ready < 0)
406   { 406   {
MISUBC 407   if (errno == EINTR || errno == EBADF) 407   if (errno == EINTR || errno == EBADF)
MISUBC 408   return; 408   return;
MISUBC 409   detail::throw_system_error(make_err(errno), "select"); 409   detail::throw_system_error(make_err(errno), "select");
410   } 410   }
411   411  
412   // Process timers outside the lock 412   // Process timers outside the lock
HITCBC 413   298386 timer_svc_->process_expired(); 413   469285 timer_svc_->process_expired();
414   414  
HITCBC 415   298386 ready_queue local_ops; 415   469285 ready_queue local_ops;
416   416  
HITCBC 417   298386 if (ready > 0) 417   469285 if (ready > 0)
418   { 418   {
HITCBC 419   283601 if (FD_ISSET(pipe_fds_[0], &read_fds)) 419   443837 if (FD_ISSET(pipe_fds_[0], &read_fds))
420   { 420   {
421   char buf[256]; 421   char buf[256];
HITCBC 422   10740 while (::read(pipe_fds_[0], buf, sizeof(buf)) > 0) 422   14726 while (::read(pipe_fds_[0], buf, sizeof(buf)) > 0)
423   { 423   {
424   } 424   }
425   } 425   }
426   426  
HITCBC 427   666256 for (int i = 0; i < snapshot_count; ++i) 427   1078914 for (int i = 0; i < snapshot_count; ++i)
428   { 428   {
HITCBC 429   382655 int fd = snapshot[i].fd; 429   635077 int fd = snapshot[i].fd;
HITCBC 430   382655 reactor_descriptor_state* desc = snapshot[i].desc; 430   635077 reactor_descriptor_state* desc = snapshot[i].desc;
431   431  
HITCBC 432   382655 std::uint32_t flags = 0; 432   635077 std::uint32_t flags = 0;
HITCBC 433   382655 if (FD_ISSET(fd, &read_fds)) 433   635077 if (FD_ISSET(fd, &read_fds))
HITCBC 434   280828 flags |= reactor_event_read; 434   440061 flags |= reactor_event_read;
HITCBC 435   382655 if (FD_ISSET(fd, &write_fds)) 435   635077 if (FD_ISSET(fd, &write_fds))
HITCBC 436   2404 flags |= reactor_event_write; 436   3387 flags |= reactor_event_write;
HITCBC 437   382655 if (FD_ISSET(fd, &except_fds)) 437   635077 if (FD_ISSET(fd, &except_fds))
MISUBC 438   flags |= reactor_event_error; 438   flags |= reactor_event_error;
439   439  
HITCBC 440   382655 if (flags == 0) 440   635077 if (flags == 0)
HITCBC 441   99433 continue; 441   191640 continue;
442   442  
HITCBC 443   283222 desc->add_ready_events(flags); 443   443437 desc->add_ready_events(flags);
444   444  
HITCBC 445   283222 bool expected = false; 445   443437 bool expected = false;
HITCBC 446   283222 if (desc->is_enqueued_.compare_exchange_strong( 446   443437 if (desc->is_enqueued_.compare_exchange_strong(
447   expected, true, std::memory_order_release, 447   expected, true, std::memory_order_release,
448   std::memory_order_relaxed)) 448   std::memory_order_relaxed))
449   { 449   {
HITCBC 450   283222 local_ops.push(desc); 450   443437 local_ops.push(desc);
451   } 451   }
452   } 452   }
453   } 453   }
454   454  
HITCBC 455   298386 lock.lock(); 455   469285 lock.lock();
456   456  
HITCBC 457   298386 completed_ops_.splice(local_ops); 457   469285 completed_ops_.splice(local_ops);
HITCBC 458   298386 } 458   469285 }
459   459  
460   } // namespace boost::corosio::detail 460   } // namespace boost::corosio::detail
461   461  
462   #endif // BOOST_COROSIO_HAS_SELECT 462   #endif // BOOST_COROSIO_HAS_SELECT
463   463  
464   #endif // BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP 464   #endif // BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP