99.35% Lines (152/153) 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_EPOLL_EPOLL_SCHEDULER_HPP 11   #ifndef BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP
12   #define BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP 12   #define BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_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_EPOLL 16   #if BOOST_COROSIO_HAS_EPOLL
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/epoll/epoll_traits.hpp> 24   #include <boost/corosio/native/detail/epoll/epoll_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 <atomic> 34   #include <atomic>
35   #include <chrono> 35   #include <chrono>
36   #include <cstdint> 36   #include <cstdint>
37   #include <mutex> 37   #include <mutex>
38   #include <vector> 38   #include <vector>
39   39  
40   #include <errno.h> 40   #include <errno.h>
41   #include <sys/epoll.h> 41   #include <sys/epoll.h>
42   #include <sys/eventfd.h> 42   #include <sys/eventfd.h>
43   #include <sys/timerfd.h> 43   #include <sys/timerfd.h>
44   #include <unistd.h> 44   #include <unistd.h>
45   45  
46   namespace boost::corosio::detail { 46   namespace boost::corosio::detail {
47   47  
48   /** Linux scheduler using epoll for I/O multiplexing. 48   /** Linux scheduler using epoll for I/O multiplexing.
49   49  
50   This scheduler implements the scheduler interface using Linux epoll 50   This scheduler implements the scheduler interface using Linux epoll
51   for efficient I/O event notification. It uses a single reactor model 51   for efficient I/O event notification. It uses a single reactor model
52   where one thread runs epoll_wait while other threads 52   where one thread runs epoll_wait while other threads
53   wait on a condition variable for handler work. This design provides: 53   wait on a condition variable for handler work. This design provides:
54   54  
55   - Handler parallelism: N posted handlers can execute on N threads 55   - Handler parallelism: N posted handlers can execute on N threads
56   - No thundering herd: condition_variable wakes exactly one thread 56   - No thundering herd: condition_variable wakes exactly one thread
57   - IOCP parity: Behavior matches Windows I/O completion port semantics 57   - IOCP parity: Behavior matches Windows I/O completion port semantics
58   58  
59   When threads call run(), they first try to execute queued handlers. 59   When threads call run(), they first try to execute queued handlers.
60   If the queue is empty and no reactor is running, one thread becomes 60   If the queue is empty and no reactor is running, one thread becomes
61   the reactor and runs epoll_wait. Other threads wait on a condition 61   the reactor and runs epoll_wait. Other threads wait on a condition
62   variable until handlers are available. 62   variable until handlers are available.
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 epoll_scheduler final : public reactor_scheduler 67   class BOOST_COROSIO_DECL epoll_scheduler final : public reactor_scheduler
68   { 68   {
69   public: 69   public:
70   /** Construct the scheduler. 70   /** Construct the scheduler.
71   71  
72   Creates an epoll instance, eventfd for reactor interruption, 72   Creates an epoll instance, eventfd for reactor interruption,
73   and timerfd for kernel-managed timer expiry. 73   and timerfd for kernel-managed timer expiry.
74   74  
75   @param ctx Reference to the owning execution_context. 75   @param ctx Reference to the owning execution_context.
76   @param concurrency_hint Hint for expected thread count (unused). 76   @param concurrency_hint Hint for expected thread count (unused).
77   */ 77   */
78   epoll_scheduler(capy::execution_context& ctx, int concurrency_hint = -1); 78   epoll_scheduler(capy::execution_context& ctx, int concurrency_hint = -1);
79   79  
80   /// Destroy the scheduler. 80   /// Destroy the scheduler.
81   ~epoll_scheduler() override; 81   ~epoll_scheduler() override;
82   82  
83   epoll_scheduler(epoll_scheduler const&) = delete; 83   epoll_scheduler(epoll_scheduler const&) = delete;
84   epoll_scheduler& operator=(epoll_scheduler const&) = delete; 84   epoll_scheduler& operator=(epoll_scheduler const&) = delete;
85   85  
86   /// Shut down the scheduler, draining pending operations. 86   /// Shut down the scheduler, draining pending operations.
87   void shutdown() override; 87   void shutdown() override;
88   88  
89   /// Apply runtime configuration, resizing the event buffer. 89   /// Apply runtime configuration, resizing the event buffer.
90   void configure_reactor( 90   void configure_reactor(
91   unsigned max_events, 91   unsigned max_events,
92   unsigned budget_init, 92   unsigned budget_init,
93   unsigned budget_max, 93   unsigned budget_max,
94   unsigned unassisted) override; 94   unsigned unassisted) override;
95   95  
96   /** Return the epoll file descriptor. 96   /** Return the epoll file descriptor.
97   97  
98   Used by socket services to register file descriptors 98   Used by socket services to register file descriptors
99   for I/O event notification. 99   for I/O event notification.
100   100  
101   @return The epoll file descriptor. 101   @return The epoll file descriptor.
102   */ 102   */
103   int epoll_fd() const noexcept 103   int epoll_fd() const noexcept
104   { 104   {
105   return epoll_fd_; 105   return epoll_fd_;
106   } 106   }
107   107  
108   /** Register a descriptor for persistent monitoring. 108   /** Register a descriptor for persistent monitoring.
109   109  
110   The fd is registered once and stays registered until explicitly 110   The fd is registered once and stays registered until explicitly
111   deregistered. Events are dispatched via reactor_descriptor_state which 111   deregistered. Events are dispatched via reactor_descriptor_state which
112   tracks pending read/write/connect operations. 112   tracks pending read/write/connect operations.
113   113  
114   @param fd The file descriptor to register. 114   @param fd The file descriptor to register.
115   @param desc Pointer to descriptor data (stored in epoll_event.data.ptr). 115   @param desc Pointer to descriptor data (stored in epoll_event.data.ptr).
116   116  
117   @return The error if registration fails, otherwise a default 117   @return The error if registration fails, otherwise a default
118   constructed error code. 118   constructed error code.
119   */ 119   */
120   std::error_code 120   std::error_code
121   register_descriptor(int fd, reactor_descriptor_state* desc) const; 121   register_descriptor(int fd, reactor_descriptor_state* desc) const;
122   122  
123   /** Deregister a persistently registered descriptor. 123   /** Deregister a persistently registered descriptor.
124   124  
125   @param fd The file descriptor to deregister. 125   @param fd The file descriptor to deregister.
126   */ 126   */
127   void deregister_descriptor(int fd) const; 127   void deregister_descriptor(int fd) const;
128   128  
129   /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp). 129   /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp).
HITGIC 130 - [[nodiscard]] std::error_code 130 + 69 [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
DCB 131 - 69 register_signal_reader(int read_fd) override  
132   { 131   {
HITCBC 133   69 return register_descriptor(read_fd, signal_pipe_reader_.arm()); 132   69 return register_descriptor(read_fd, signal_pipe_reader_.arm());
134   } 133   }
135   134  
136   private: 135   private:
137 - void 136 + void run_task(lock_type& lock, context_type& ctx, long timeout_us) override;
138 - run_task(lock_type& lock, context_type& ctx,  
139 - long timeout_us) override;  
140   void interrupt_reactor() const override; 137   void interrupt_reactor() const override;
141   void update_timerfd() const; 138   void update_timerfd() const;
142   139  
143   int epoll_fd_; 140   int epoll_fd_;
144   int event_fd_; 141   int event_fd_;
145   int timer_fd_; 142   int timer_fd_;
146   143  
147   // Watches the global signal self-pipe's read end (armed lazily by 144   // Watches the global signal self-pipe's read end (armed lazily by
148   // register_signal_reader on the first signal registration). 145   // register_signal_reader on the first signal registration).
149   reactor_signal_pipe_reader signal_pipe_reader_; 146   reactor_signal_pipe_reader signal_pipe_reader_;
150   147  
151   // Edge-triggered eventfd state 148   // Edge-triggered eventfd state
152   mutable std::atomic<bool> eventfd_armed_{false}; 149   mutable std::atomic<bool> eventfd_armed_{false};
153   150  
154   // Set when the earliest timer changes; flushed before epoll_wait 151   // Set when the earliest timer changes; flushed before epoll_wait
155   mutable std::atomic<bool> timerfd_stale_{false}; 152   mutable std::atomic<bool> timerfd_stale_{false};
156   153  
157   // Event buffer sized from max_events_per_poll_ (set at construction, 154   // Event buffer sized from max_events_per_poll_ (set at construction,
158   // resized by configure_reactor via io_context_options). 155   // resized by configure_reactor via io_context_options).
159   std::vector<epoll_event> event_buffer_; 156   std::vector<epoll_event> event_buffer_;
160   }; 157   };
161   158  
HITCBC 162   1232 inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int) 159   1232 inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
HITCBC 163   1232 : epoll_fd_(-1) 160   1232 : epoll_fd_(-1)
HITCBC 164   1232 , event_fd_(-1) 161   1232 , event_fd_(-1)
HITCBC 165   1232 , timer_fd_(-1) 162   1232 , timer_fd_(-1)
HITCBC 166   2464 , event_buffer_(max_events_per_poll_) 163   2464 , event_buffer_(max_events_per_poll_)
167   { 164   {
HITCBC 168   1232 epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC); 165   1232 epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
HITCBC 169   1232 if (epoll_fd_ < 0) 166   1232 if (epoll_fd_ < 0)
HITCBC 170   1 detail::throw_system_error(make_err(errno), "epoll_create1"); 167   1 detail::throw_system_error(make_err(errno), "epoll_create1");
171   168  
HITCBC 172   1231 event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC); 169   1231 event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
HITCBC 173   1231 if (event_fd_ < 0) 170   1231 if (event_fd_ < 0)
174   { 171   {
HITCBC 175   1 int errn = errno; 172   1 int errn = errno;
HITCBC 176   1 ::close(epoll_fd_); 173   1 ::close(epoll_fd_);
HITCBC 177   1 detail::throw_system_error(make_err(errn), "eventfd"); 174   1 detail::throw_system_error(make_err(errn), "eventfd");
178   } 175   }
179   176  
HITCBC 180   1230 timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC); 177   1230 timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
HITCBC 181   1230 if (timer_fd_ < 0) 178   1230 if (timer_fd_ < 0)
182   { 179   {
HITCBC 183   1 int errn = errno; 180   1 int errn = errno;
HITCBC 184   1 ::close(event_fd_); 181   1 ::close(event_fd_);
HITCBC 185   1 ::close(epoll_fd_); 182   1 ::close(epoll_fd_);
HITCBC 186   1 detail::throw_system_error(make_err(errn), "timerfd_create"); 183   1 detail::throw_system_error(make_err(errn), "timerfd_create");
187   } 184   }
188   185  
HITCBC 189   1229 epoll_event ev{}; 186   1229 epoll_event ev{};
HITCBC 190   1229 ev.events = EPOLLIN | EPOLLET; 187   1229 ev.events = EPOLLIN | EPOLLET;
HITCBC 191   1229 ev.data.ptr = nullptr; 188   1229 ev.data.ptr = nullptr;
HITCBC 192   1229 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0) 189   1229 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0)
193   { 190   {
HITCBC 194   1 int errn = errno; 191   1 int errn = errno;
HITCBC 195   1 ::close(timer_fd_); 192   1 ::close(timer_fd_);
HITCBC 196   1 ::close(event_fd_); 193   1 ::close(event_fd_);
HITCBC 197   1 ::close(epoll_fd_); 194   1 ::close(epoll_fd_);
HITCBC 198   1 detail::throw_system_error(make_err(errn), "epoll_ctl"); 195   1 detail::throw_system_error(make_err(errn), "epoll_ctl");
199   } 196   }
200   197  
HITCBC 201   1228 epoll_event timer_ev{}; 198   1228 epoll_event timer_ev{};
HITCBC 202   1228 timer_ev.events = EPOLLIN | EPOLLERR; 199   1228 timer_ev.events = EPOLLIN | EPOLLERR;
HITCBC 203   1228 timer_ev.data.ptr = &timer_fd_; 200   1228 timer_ev.data.ptr = &timer_fd_;
HITCBC 204   1228 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0) 201   1228 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0)
205   { 202   {
HITCBC 206   1 int errn = errno; 203   1 int errn = errno;
HITCBC 207   1 ::close(timer_fd_); 204   1 ::close(timer_fd_);
HITCBC 208   1 ::close(event_fd_); 205   1 ::close(event_fd_);
HITCBC 209   1 ::close(epoll_fd_); 206   1 ::close(epoll_fd_);
HITCBC 210   1 detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)"); 207   1 detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)");
211   } 208   }
212   209  
HITCBC 213   1227 timer_svc_ = &get_timer_service(ctx, *this); 210   1227 timer_svc_ = &get_timer_service(ctx, *this);
HITCBC 214   1227 timer_svc_->set_on_earliest_changed( 211   1227 timer_svc_->set_on_earliest_changed(
HITCBC 215   5537 timer_service::callback(this, [](void* p) { 212   5586 timer_service::callback(this, [](void* p) {
HITCBC 216   4310 auto* self = static_cast<epoll_scheduler*>(p); 213   4359 auto* self = static_cast<epoll_scheduler*>(p);
HITCBC 217   4310 self->timerfd_stale_.store(true, std::memory_order_release); 214   4359 self->timerfd_stale_.store(true, std::memory_order_release);
HITCBC 218   4310 self->interrupt_reactor(); 215   4359 self->interrupt_reactor();
HITCBC 219   4310 })); 216   4359 }));
220   217  
HITCBC 221   1227 get_resolver_service(ctx, *this); 218   1227 get_resolver_service(ctx, *this);
HITCBC 222   1227 get_signal_service(ctx, *this); 219   1227 get_signal_service(ctx, *this);
HITCBC 223   1227 get_stream_file_service(ctx, *this); 220   1227 get_stream_file_service(ctx, *this);
HITCBC 224   1227 get_random_access_file_service(ctx, *this); 221   1227 get_random_access_file_service(ctx, *this);
225   222  
HITCBC 226   1227 completed_ops_.push(&task_op_); 223   1227 completed_ops_.push(&task_op_);
HITCBC 227   1242 } 224   1242 }
228   225  
HITCBC 229   2454 inline epoll_scheduler::~epoll_scheduler() 226   2454 inline epoll_scheduler::~epoll_scheduler()
230   { 227   {
HITCBC 231   1227 if (timer_fd_ >= 0) 228   1227 if (timer_fd_ >= 0)
HITCBC 232   1227 ::close(timer_fd_); 229   1227 ::close(timer_fd_);
HITCBC 233   1227 if (event_fd_ >= 0) 230   1227 if (event_fd_ >= 0)
HITCBC 234   1227 ::close(event_fd_); 231   1227 ::close(event_fd_);
HITCBC 235   1227 if (epoll_fd_ >= 0) 232   1227 if (epoll_fd_ >= 0)
HITCBC 236   1227 ::close(epoll_fd_); 233   1227 ::close(epoll_fd_);
HITCBC 237   2454 } 234   2454 }
238   235  
239   inline void 236   inline void
HITCBC 240   1227 epoll_scheduler::shutdown() 237   1227 epoll_scheduler::shutdown()
241   { 238   {
HITCBC 242   1227 shutdown_drain(); 239   1227 shutdown_drain();
243   240  
HITCBC 244   1227 if (event_fd_ >= 0) 241   1227 if (event_fd_ >= 0)
HITCBC 245   1227 interrupt_reactor(); 242   1227 interrupt_reactor();
HITCBC 246   1227 } 243   1227 }
247   244  
248   inline void 245   inline void
HITCBC 249   23 epoll_scheduler::configure_reactor( 246   23 epoll_scheduler::configure_reactor(
250   unsigned max_events, 247   unsigned max_events,
251   unsigned budget_init, 248   unsigned budget_init,
252   unsigned budget_max, 249   unsigned budget_max,
253   unsigned unassisted) 250   unsigned unassisted)
254   { 251   {
HITCBC 255   23 reactor_scheduler::configure_reactor( 252   23 reactor_scheduler::configure_reactor(
256   max_events, budget_init, budget_max, unassisted); 253   max_events, budget_init, budget_max, unassisted);
HITCBC 257   21 event_buffer_.resize(max_events_per_poll_); 254   21 event_buffer_.resize(max_events_per_poll_);
HITCBC 258   21 } 255   21 }
259   256  
260   inline std::error_code 257   inline std::error_code
HITCBC 261 - 5733 epoll_scheduler::register_descriptor(int fd, reactor_descriptor_state* desc) const 258 + 5748 epoll_scheduler::register_descriptor(
  259 + int fd, reactor_descriptor_state* desc) const
262   { 260   {
HITCBC 263   5733 epoll_event ev{}; 261   5748 epoll_event ev{};
HITCBC 264   5733 ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP; 262   5748 ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
HITCBC 265   5733 ev.data.ptr = desc; 263   5748 ev.data.ptr = desc;
266   264  
HITCBC 267   5733 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0) 265   5748 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
HITCBC 268   7 return make_err(errno); 266   7 return make_err(errno);
269   267  
HITCBC 270   5726 desc->registered_events = ev.events; 268   5741 desc->registered_events = ev.events;
HITCBC 271   5726 desc->fd = fd; 269   5741 desc->fd = fd;
HITCBC 272   5726 desc->scheduler_ = this; 270   5741 desc->scheduler_ = this;
HITCBC 273   5726 desc->mutex.set_enabled(reactor_io_locking_); 271   5741 desc->mutex.set_enabled(reactor_io_locking_);
HITCBC 274   5726 desc->ready_events_.store(0, std::memory_order_relaxed); 272   5741 desc->ready_events_.store(0, std::memory_order_relaxed);
275   273  
HITCBC 276   5726 conditionally_enabled_mutex::scoped_lock lock(desc->mutex); 274   5741 conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
HITCBC 277   5726 desc->impl_ref_.reset(); 275   5741 desc->impl_ref_.reset();
HITCBC 278   5726 desc->read_ready = false; 276   5741 desc->read_ready = false;
HITCBC 279   5726 desc->write_ready = false; 277   5741 desc->write_ready = false;
HITCBC 280   5726 return {}; 278   5741 return {};
HITCBC 281   5726 } 279   5741 }
282   280  
283   inline void 281   inline void
HITCBC 284   5658 epoll_scheduler::deregister_descriptor(int fd) const 282   5673 epoll_scheduler::deregister_descriptor(int fd) const
285   { 283   {
HITCBC 286   5658 ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr); 284   5673 ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
HITCBC 287   5658 } 285   5673 }
288   286  
289   inline void 287   inline void
HITCBC 290   6725 epoll_scheduler::interrupt_reactor() const 288   6764 epoll_scheduler::interrupt_reactor() const
291   { 289   {
HITCBC 292   6725 bool expected = false; 290   6764 bool expected = false;
HITCBC 293   6725 if (eventfd_armed_.compare_exchange_strong( 291   6764 if (eventfd_armed_.compare_exchange_strong(
294   expected, true, std::memory_order_release, 292   expected, true, std::memory_order_release,
295   std::memory_order_relaxed)) 293   std::memory_order_relaxed))
296   { 294   {
HITCBC 297   5279 std::uint64_t val = 1; 295   5222 std::uint64_t val = 1;
HITCBC 298   5279 if (::write(event_fd_, &val, sizeof(val)) < 0) 296   5222 if (::write(event_fd_, &val, sizeof(val)) < 0)
299   { 297   {
300   // The flag is what coalesces later interrupts into a byte 298   // The flag is what coalesces later interrupts into a byte
301   // already in the eventfd; a write that failed put no byte 299   // already in the eventfd; a write that failed put no byte
302   // there, so leaving it armed would swallow every interrupt 300   // there, so leaving it armed would swallow every interrupt
303   // that follows. Disarming keeps the cost to the interrupts 301   // that follows. Disarming keeps the cost to the interrupts
304   // already in flight -- the next one arms and writes again, 302   // already in flight -- the next one arms and writes again,
305   // instead of every one after this coalescing into a byte 303   // instead of every one after this coalescing into a byte
306   // that does not exist. 304   // that does not exist.
HITCBC 307   2 eventfd_armed_.store(false, std::memory_order_release); 305   2 eventfd_armed_.store(false, std::memory_order_release);
308   } 306   }
309   } 307   }
HITCBC 310   6725 } 308   6764 }
311   309  
312   inline void 310   inline void
HITCBC 313   11371 epoll_scheduler::update_timerfd() const 311   11143 epoll_scheduler::update_timerfd() const
314   { 312   {
HITCBC 315   11371 auto nearest = timer_svc_->nearest_expiry(); 313   11143 auto nearest = timer_svc_->nearest_expiry();
316   314  
HITCBC 317   11371 itimerspec ts{}; 315   11143 itimerspec ts{};
HITCBC 318   11371 int flags = 0; 316   11143 int flags = 0;
319   317  
HITCBC 320   11371 if (nearest == timer_service::time_point::max()) 318   11143 if (nearest == timer_service::time_point::max())
321   { 319   {
322   // No timers — disarm by setting to 0 (relative) 320   // No timers — disarm by setting to 0 (relative)
323   } 321   }
324   else 322   else
325   { 323   {
HITCBC 326   10180 auto now = std::chrono::steady_clock::now(); 324   9948 auto now = std::chrono::steady_clock::now();
HITCBC 327   10180 if (nearest <= now) 325   9948 if (nearest <= now)
328   { 326   {
329   // Use 1ns instead of 0 — zero disarms the timerfd 327   // Use 1ns instead of 0 — zero disarms the timerfd
HITCBC 330   1134 ts.it_value.tv_nsec = 1; 328   1167 ts.it_value.tv_nsec = 1;
331   } 329   }
332   else 330   else
333   { 331   {
HITCBC 334   9046 auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>( 332   8781 auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
HITCBC 335   9046 nearest - now) 333   8781 nearest - now)
HITCBC 336   9046 .count(); 334   8781 .count();
HITCBC 337   9046 ts.it_value.tv_sec = nsec / 1000000000; 335   8781 ts.it_value.tv_sec = nsec / 1000000000;
HITCBC 338   9046 ts.it_value.tv_nsec = nsec % 1000000000; 336   8781 ts.it_value.tv_nsec = nsec % 1000000000;
HITCBC 339   9046 if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0) 337   8781 if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
MISUBC 340   ts.it_value.tv_nsec = 1; 338   ts.it_value.tv_nsec = 1;
341   } 339   }
342   } 340   }
343   341  
HITCBC 344   11371 if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0) 342   11143 if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
HITCBC 345   1 detail::throw_system_error(make_err(errno), "timerfd_settime"); 343   1 detail::throw_system_error(make_err(errno), "timerfd_settime");
HITCBC 346   11370 } 344   11142 }
347   345  
348   inline void 346   inline void
HITCBC 349 - 39476 epoll_scheduler::run_task( 347 + 39736 epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
350 - lock_type& lock, context_type& ctx, long timeout_us)  
351   { 348   {
352   int timeout_ms; 349   int timeout_ms;
HITCBC 353   39476 if (task_interrupted_) 350   39736 if (task_interrupted_)
HITCBC 354   28299 timeout_ms = 0; 351   28575 timeout_ms = 0;
HITCBC 355   11177 else if (timeout_us < 0) 352   11161 else if (timeout_us < 0)
HITCBC 356   11158 timeout_ms = -1; 353   11146 timeout_ms = -1;
357   else 354   else
HITCBC 358   19 timeout_ms = static_cast<int>((timeout_us + 999) / 1000); 355   15 timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
359   356  
HITCBC 360   39476 if (lock.owns_lock()) 357   39736 if (lock.owns_lock())
HITCBC 361   11179 lock.unlock(); 358   11163 lock.unlock();
362   359  
HITCBC 363   39476 task_cleanup on_exit{this, &lock, ctx}; 360   39736 task_cleanup on_exit{this, &lock, ctx};
364   361  
365   // Flush deferred timerfd programming before blocking 362   // Flush deferred timerfd programming before blocking
HITCBC 366   39476 if (timerfd_stale_.exchange(false, std::memory_order_acquire)) 363   39736 if (timerfd_stale_.exchange(false, std::memory_order_acquire))
HITCBC 367   3684 update_timerfd(); 364   3637 update_timerfd();
368   365  
HITCBC 369   39475 int nfds = ::epoll_wait( 366   39735 int nfds = ::epoll_wait(
HITGIC 370 - epoll_fd_, event_buffer_.data(), 367 + 39735 epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()),
ECB 371 - 39475 static_cast<int>(event_buffer_.size()), timeout_ms); 368 + timeout_ms);
372   369  
HITCBC 373   39475 if (nfds < 0 && errno != EINTR) 370   39735 if (nfds < 0 && errno != EINTR)
HITCBC 374   1 detail::throw_system_error(make_err(errno), "epoll_wait"); 371   1 detail::throw_system_error(make_err(errno), "epoll_wait");
375   372  
HITCBC 376   39474 bool check_timers = false; 373   39734 bool check_timers = false;
HITCBC 377   39474 ready_queue local_ops; 374   39734 ready_queue local_ops;
378   375  
HITCBC 379   87128 for (int i = 0; i < nfds; ++i) 376   86129 for (int i = 0; i < nfds; ++i)
380   { 377   {
HITCBC 381   47654 if (event_buffer_[i].data.ptr == nullptr) 378   46395 if (event_buffer_[i].data.ptr == nullptr)
382   { 379   {
383   std::uint64_t val; 380   std::uint64_t val;
384   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection) 381   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
HITCBC 385   4050 [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val)); 382   3993 [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
HITCBC 386   4050 eventfd_armed_.store(false, std::memory_order_relaxed); 383   3993 eventfd_armed_.store(false, std::memory_order_relaxed);
HITCBC 387   4050 continue; 384   3993 continue;
HITCBC 388   4050 } 385   3993 }
389   386  
HITCBC 390   43604 if (event_buffer_[i].data.ptr == &timer_fd_) 387   42402 if (event_buffer_[i].data.ptr == &timer_fd_)
391   { 388   {
392   std::uint64_t expirations; 389   std::uint64_t expirations;
393   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection) 390   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
394   [[maybe_unused]] auto r = 391   [[maybe_unused]] auto r =
HITCBC 395   7687 ::read(timer_fd_, &expirations, sizeof(expirations)); 392   7506 ::read(timer_fd_, &expirations, sizeof(expirations));
HITCBC 396   7687 check_timers = true; 393   7506 check_timers = true;
HITCBC 397   7687 continue; 394   7506 continue;
HITCBC 398   7687 } 395   7506 }
399   396  
400   auto* desc = 397   auto* desc =
HITCBC 401   35917 static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr); 398   34896 static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
HITCBC 402   35917 desc->add_ready_events(event_buffer_[i].events); 399   34896 desc->add_ready_events(event_buffer_[i].events);
403   400  
HITCBC 404   35917 bool expected = false; 401   34896 bool expected = false;
HITCBC 405   35917 if (desc->is_enqueued_.compare_exchange_strong( 402   34896 if (desc->is_enqueued_.compare_exchange_strong(
406   expected, true, std::memory_order_release, 403   expected, true, std::memory_order_release,
407   std::memory_order_relaxed)) 404   std::memory_order_relaxed))
408   { 405   {
HITCBC 409   35917 local_ops.push(desc); 406   34896 local_ops.push(desc);
410   } 407   }
411   } 408   }
412   409  
HITCBC 413   39474 if (check_timers) 410   39734 if (check_timers)
414   { 411   {
HITCBC 415   7687 timer_svc_->process_expired(); 412   7506 timer_svc_->process_expired();
HITCBC 416   7687 update_timerfd(); 413   7506 update_timerfd();
417   } 414   }
418   415  
HITCBC 419   39474 lock.lock(); 416   39734 lock.lock();
420   417  
HITCBC 421   39474 completed_ops_.splice(local_ops); 418   39734 completed_ops_.splice(local_ops);
HITCBC 422   39476 } 419   39736 }
423   420  
424   } // namespace boost::corosio::detail 421   } // namespace boost::corosio::detail
425   422  
426   #endif // BOOST_COROSIO_HAS_EPOLL 423   #endif // BOOST_COROSIO_HAS_EPOLL
427   424  
428   #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP 425   #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP