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).
ECB 130 - 69 [[nodiscard]] std::error_code register_signal_reader(int read_fd) override 130 + [[nodiscard]] std::error_code
HITGNC   131 + 69 register_signal_reader(int read_fd) override
131   { 132   {
HITCBC 132   69 return register_descriptor(read_fd, signal_pipe_reader_.arm()); 133   69 return register_descriptor(read_fd, signal_pipe_reader_.arm());
133   } 134   }
134   135  
135   private: 136   private:
136 - void run_task(lock_type& lock, context_type& ctx, long timeout_us) override; 137 + void
  138 + run_task(lock_type& lock, context_type& ctx,
  139 + long timeout_us) override;
137   void interrupt_reactor() const override; 140   void interrupt_reactor() const override;
138   void update_timerfd() const; 141   void update_timerfd() const;
139   142  
140   int epoll_fd_; 143   int epoll_fd_;
141   int event_fd_; 144   int event_fd_;
142   int timer_fd_; 145   int timer_fd_;
143   146  
144   // Watches the global signal self-pipe's read end (armed lazily by 147   // Watches the global signal self-pipe's read end (armed lazily by
145   // register_signal_reader on the first signal registration). 148   // register_signal_reader on the first signal registration).
146   reactor_signal_pipe_reader signal_pipe_reader_; 149   reactor_signal_pipe_reader signal_pipe_reader_;
147   150  
148   // Edge-triggered eventfd state 151   // Edge-triggered eventfd state
149   mutable std::atomic<bool> eventfd_armed_{false}; 152   mutable std::atomic<bool> eventfd_armed_{false};
150   153  
151   // Set when the earliest timer changes; flushed before epoll_wait 154   // Set when the earliest timer changes; flushed before epoll_wait
152   mutable std::atomic<bool> timerfd_stale_{false}; 155   mutable std::atomic<bool> timerfd_stale_{false};
153   156  
154   // Event buffer sized from max_events_per_poll_ (set at construction, 157   // Event buffer sized from max_events_per_poll_ (set at construction,
155   // resized by configure_reactor via io_context_options). 158   // resized by configure_reactor via io_context_options).
156   std::vector<epoll_event> event_buffer_; 159   std::vector<epoll_event> event_buffer_;
157   }; 160   };
158   161  
HITCBC 159   1232 inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int) 162   1232 inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
HITCBC 160   1232 : epoll_fd_(-1) 163   1232 : epoll_fd_(-1)
HITCBC 161   1232 , event_fd_(-1) 164   1232 , event_fd_(-1)
HITCBC 162   1232 , timer_fd_(-1) 165   1232 , timer_fd_(-1)
HITCBC 163   2464 , event_buffer_(max_events_per_poll_) 166   2464 , event_buffer_(max_events_per_poll_)
164   { 167   {
HITCBC 165   1232 epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC); 168   1232 epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
HITCBC 166   1232 if (epoll_fd_ < 0) 169   1232 if (epoll_fd_ < 0)
HITCBC 167   1 detail::throw_system_error(make_err(errno), "epoll_create1"); 170   1 detail::throw_system_error(make_err(errno), "epoll_create1");
168   171  
HITCBC 169   1231 event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC); 172   1231 event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
HITCBC 170   1231 if (event_fd_ < 0) 173   1231 if (event_fd_ < 0)
171   { 174   {
HITCBC 172   1 int errn = errno; 175   1 int errn = errno;
HITCBC 173   1 ::close(epoll_fd_); 176   1 ::close(epoll_fd_);
HITCBC 174   1 detail::throw_system_error(make_err(errn), "eventfd"); 177   1 detail::throw_system_error(make_err(errn), "eventfd");
175   } 178   }
176   179  
HITCBC 177   1230 timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC); 180   1230 timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
HITCBC 178   1230 if (timer_fd_ < 0) 181   1230 if (timer_fd_ < 0)
179   { 182   {
HITCBC 180   1 int errn = errno; 183   1 int errn = errno;
HITCBC 181   1 ::close(event_fd_); 184   1 ::close(event_fd_);
HITCBC 182   1 ::close(epoll_fd_); 185   1 ::close(epoll_fd_);
HITCBC 183   1 detail::throw_system_error(make_err(errn), "timerfd_create"); 186   1 detail::throw_system_error(make_err(errn), "timerfd_create");
184   } 187   }
185   188  
HITCBC 186   1229 epoll_event ev{}; 189   1229 epoll_event ev{};
HITCBC 187   1229 ev.events = EPOLLIN | EPOLLET; 190   1229 ev.events = EPOLLIN | EPOLLET;
HITCBC 188   1229 ev.data.ptr = nullptr; 191   1229 ev.data.ptr = nullptr;
HITCBC 189   1229 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0) 192   1229 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0)
190   { 193   {
HITCBC 191   1 int errn = errno; 194   1 int errn = errno;
HITCBC 192   1 ::close(timer_fd_); 195   1 ::close(timer_fd_);
HITCBC 193   1 ::close(event_fd_); 196   1 ::close(event_fd_);
HITCBC 194   1 ::close(epoll_fd_); 197   1 ::close(epoll_fd_);
HITCBC 195   1 detail::throw_system_error(make_err(errn), "epoll_ctl"); 198   1 detail::throw_system_error(make_err(errn), "epoll_ctl");
196   } 199   }
197   200  
HITCBC 198   1228 epoll_event timer_ev{}; 201   1228 epoll_event timer_ev{};
HITCBC 199   1228 timer_ev.events = EPOLLIN | EPOLLERR; 202   1228 timer_ev.events = EPOLLIN | EPOLLERR;
HITCBC 200   1228 timer_ev.data.ptr = &timer_fd_; 203   1228 timer_ev.data.ptr = &timer_fd_;
HITCBC 201   1228 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0) 204   1228 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0)
202   { 205   {
HITCBC 203   1 int errn = errno; 206   1 int errn = errno;
HITCBC 204   1 ::close(timer_fd_); 207   1 ::close(timer_fd_);
HITCBC 205   1 ::close(event_fd_); 208   1 ::close(event_fd_);
HITCBC 206   1 ::close(epoll_fd_); 209   1 ::close(epoll_fd_);
HITCBC 207   1 detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)"); 210   1 detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)");
208   } 211   }
209   212  
HITCBC 210   1227 timer_svc_ = &get_timer_service(ctx, *this); 213   1227 timer_svc_ = &get_timer_service(ctx, *this);
HITCBC 211   1227 timer_svc_->set_on_earliest_changed( 214   1227 timer_svc_->set_on_earliest_changed(
HITCBC 212   5695 timer_service::callback(this, [](void* p) { 215   5427 timer_service::callback(this, [](void* p) {
HITCBC 213   4468 auto* self = static_cast<epoll_scheduler*>(p); 216   4200 auto* self = static_cast<epoll_scheduler*>(p);
HITCBC 214   4468 self->timerfd_stale_.store(true, std::memory_order_release); 217   4200 self->timerfd_stale_.store(true, std::memory_order_release);
HITCBC 215   4468 self->interrupt_reactor(); 218   4200 self->interrupt_reactor();
HITCBC 216   4468 })); 219   4200 }));
217   220  
HITCBC 218   1227 get_resolver_service(ctx, *this); 221   1227 get_resolver_service(ctx, *this);
HITCBC 219   1227 get_signal_service(ctx, *this); 222   1227 get_signal_service(ctx, *this);
HITCBC 220   1227 get_stream_file_service(ctx, *this); 223   1227 get_stream_file_service(ctx, *this);
HITCBC 221   1227 get_random_access_file_service(ctx, *this); 224   1227 get_random_access_file_service(ctx, *this);
222   225  
HITCBC 223   1227 completed_ops_.push(&task_op_); 226   1227 completed_ops_.push(&task_op_);
HITCBC 224   1242 } 227   1242 }
225   228  
HITCBC 226   2454 inline epoll_scheduler::~epoll_scheduler() 229   2454 inline epoll_scheduler::~epoll_scheduler()
227   { 230   {
HITCBC 228   1227 if (timer_fd_ >= 0) 231   1227 if (timer_fd_ >= 0)
HITCBC 229   1227 ::close(timer_fd_); 232   1227 ::close(timer_fd_);
HITCBC 230   1227 if (event_fd_ >= 0) 233   1227 if (event_fd_ >= 0)
HITCBC 231   1227 ::close(event_fd_); 234   1227 ::close(event_fd_);
HITCBC 232   1227 if (epoll_fd_ >= 0) 235   1227 if (epoll_fd_ >= 0)
HITCBC 233   1227 ::close(epoll_fd_); 236   1227 ::close(epoll_fd_);
HITCBC 234   2454 } 237   2454 }
235   238  
236   inline void 239   inline void
HITCBC 237   1227 epoll_scheduler::shutdown() 240   1227 epoll_scheduler::shutdown()
238   { 241   {
HITCBC 239   1227 shutdown_drain(); 242   1227 shutdown_drain();
240   243  
HITCBC 241   1227 if (event_fd_ >= 0) 244   1227 if (event_fd_ >= 0)
HITCBC 242   1227 interrupt_reactor(); 245   1227 interrupt_reactor();
HITCBC 243   1227 } 246   1227 }
244   247  
245   inline void 248   inline void
HITCBC 246   23 epoll_scheduler::configure_reactor( 249   23 epoll_scheduler::configure_reactor(
247   unsigned max_events, 250   unsigned max_events,
248   unsigned budget_init, 251   unsigned budget_init,
249   unsigned budget_max, 252   unsigned budget_max,
250   unsigned unassisted) 253   unsigned unassisted)
251   { 254   {
HITCBC 252   23 reactor_scheduler::configure_reactor( 255   23 reactor_scheduler::configure_reactor(
253   max_events, budget_init, budget_max, unassisted); 256   max_events, budget_init, budget_max, unassisted);
HITCBC 254   21 event_buffer_.resize(max_events_per_poll_); 257   21 event_buffer_.resize(max_events_per_poll_);
HITCBC 255   21 } 258   21 }
256   259  
257   inline std::error_code 260   inline std::error_code
HITCBC 258 - 5900 epoll_scheduler::register_descriptor( 261 + 5582 epoll_scheduler::register_descriptor(int fd, reactor_descriptor_state* desc) const
259 - int fd, reactor_descriptor_state* desc) const  
260   { 262   {
HITCBC 261   5900 epoll_event ev{}; 263   5582 epoll_event ev{};
HITCBC 262   5900 ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP; 264   5582 ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
HITCBC 263   5900 ev.data.ptr = desc; 265   5582 ev.data.ptr = desc;
264   266  
HITCBC 265   5900 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0) 267   5582 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
HITCBC 266   7 return make_err(errno); 268   7 return make_err(errno);
267   269  
HITCBC 268   5893 desc->registered_events = ev.events; 270   5575 desc->registered_events = ev.events;
HITCBC 269   5893 desc->fd = fd; 271   5575 desc->fd = fd;
HITCBC 270   5893 desc->scheduler_ = this; 272   5575 desc->scheduler_ = this;
HITCBC 271   5893 desc->mutex.set_enabled(reactor_io_locking_); 273   5575 desc->mutex.set_enabled(reactor_io_locking_);
HITCBC 272   5893 desc->ready_events_.store(0, std::memory_order_relaxed); 274   5575 desc->ready_events_.store(0, std::memory_order_relaxed);
273   275  
HITCBC 274   5893 conditionally_enabled_mutex::scoped_lock lock(desc->mutex); 276   5575 conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
HITCBC 275   5893 desc->impl_ref_.reset(); 277   5575 desc->impl_ref_.reset();
HITCBC 276   5893 desc->read_ready = false; 278   5575 desc->read_ready = false;
HITCBC 277   5893 desc->write_ready = false; 279   5575 desc->write_ready = false;
HITCBC 278   5893 return {}; 280   5575 return {};
HITCBC 279   5893 } 281   5575 }
280   282  
281   inline void 283   inline void
HITCBC 282   5825 epoll_scheduler::deregister_descriptor(int fd) const 284   5507 epoll_scheduler::deregister_descriptor(int fd) const
283   { 285   {
HITCBC 284   5825 ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr); 286   5507 ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
HITCBC 285   5825 } 287   5507 }
286   288  
287   inline void 289   inline void
HITCBC 288   6880 epoll_scheduler::interrupt_reactor() const 290   6613 epoll_scheduler::interrupt_reactor() const
289   { 291   {
HITCBC 290   6880 bool expected = false; 292   6613 bool expected = false;
HITCBC 291   6880 if (eventfd_armed_.compare_exchange_strong( 293   6613 if (eventfd_armed_.compare_exchange_strong(
292   expected, true, std::memory_order_release, 294   expected, true, std::memory_order_release,
293   std::memory_order_relaxed)) 295   std::memory_order_relaxed))
294   { 296   {
HITCBC 295   5308 std::uint64_t val = 1; 297   5199 std::uint64_t val = 1;
HITCBC 296   5308 if (::write(event_fd_, &val, sizeof(val)) < 0) 298   5199 if (::write(event_fd_, &val, sizeof(val)) < 0)
297   { 299   {
298   // The flag is what coalesces later interrupts into a byte 300   // The flag is what coalesces later interrupts into a byte
299   // already in the eventfd; a write that failed put no byte 301   // already in the eventfd; a write that failed put no byte
300   // there, so leaving it armed would swallow every interrupt 302   // there, so leaving it armed would swallow every interrupt
301   // that follows. Disarming keeps the cost to the interrupts 303   // that follows. Disarming keeps the cost to the interrupts
302   // already in flight -- the next one arms and writes again, 304   // already in flight -- the next one arms and writes again,
303   // instead of every one after this coalescing into a byte 305   // instead of every one after this coalescing into a byte
304   // that does not exist. 306   // that does not exist.
HITCBC 305   2 eventfd_armed_.store(false, std::memory_order_release); 307   2 eventfd_armed_.store(false, std::memory_order_release);
306   } 308   }
307   } 309   }
HITCBC 308   6880 } 310   6613 }
309   311  
310   inline void 312   inline void
HITCBC 311   11142 epoll_scheduler::update_timerfd() const 313   10850 epoll_scheduler::update_timerfd() const
312   { 314   {
HITCBC 313   11142 auto nearest = timer_svc_->nearest_expiry(); 315   10850 auto nearest = timer_svc_->nearest_expiry();
314   316  
HITCBC 315   11142 itimerspec ts{}; 317   10850 itimerspec ts{};
HITCBC 316   11142 int flags = 0; 318   10850 int flags = 0;
317   319  
HITCBC 318   11142 if (nearest == timer_service::time_point::max()) 320   10850 if (nearest == timer_service::time_point::max())
319   { 321   {
320   // No timers — disarm by setting to 0 (relative) 322   // No timers — disarm by setting to 0 (relative)
321   } 323   }
322   else 324   else
323   { 325   {
HITCBC 324   10000 auto now = std::chrono::steady_clock::now(); 326   9664 auto now = std::chrono::steady_clock::now();
HITCBC 325   10000 if (nearest <= now) 327   9664 if (nearest <= now)
326   { 328   {
327   // Use 1ns instead of 0 — zero disarms the timerfd 329   // Use 1ns instead of 0 — zero disarms the timerfd
HITCBC 328   1030 ts.it_value.tv_nsec = 1; 330   1043 ts.it_value.tv_nsec = 1;
329   } 331   }
330   else 332   else
331   { 333   {
HITCBC 332   8970 auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>( 334   8621 auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
HITCBC 333   8970 nearest - now) 335   8621 nearest - now)
HITCBC 334   8970 .count(); 336   8621 .count();
HITCBC 335   8970 ts.it_value.tv_sec = nsec / 1000000000; 337   8621 ts.it_value.tv_sec = nsec / 1000000000;
HITCBC 336   8970 ts.it_value.tv_nsec = nsec % 1000000000; 338   8621 ts.it_value.tv_nsec = nsec % 1000000000;
HITCBC 337   8970 if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0) 339   8621 if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
MISUBC 338   ts.it_value.tv_nsec = 1; 340   ts.it_value.tv_nsec = 1;
339   } 341   }
340   } 342   }
341   343  
HITCBC 342   11142 if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0) 344   10850 if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
HITCBC 343   1 detail::throw_system_error(make_err(errno), "timerfd_settime"); 345   1 detail::throw_system_error(make_err(errno), "timerfd_settime");
HITCBC 344   11141 } 346   10849 }
345   347  
346   inline void 348   inline void
HITCBC 347 - 39004 epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us) 349 + 39127 epoll_scheduler::run_task(
  350 + lock_type& lock, context_type& ctx, long timeout_us)
348   { 351   {
349   int timeout_ms; 352   int timeout_ms;
HITCBC 350   39004 if (task_interrupted_) 353   39127 if (task_interrupted_)
HITCBC 351   28380 timeout_ms = 0; 354   28098 timeout_ms = 0;
HITCBC 352   10624 else if (timeout_us < 0) 355   11029 else if (timeout_us < 0)
HITCBC 353   10608 timeout_ms = -1; 356   11013 timeout_ms = -1;
354   else 357   else
HITCBC 355   16 timeout_ms = static_cast<int>((timeout_us + 999) / 1000); 358   16 timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
356   359  
HITCBC 357   39004 if (lock.owns_lock()) 360   39127 if (lock.owns_lock())
HITCBC 358   10626 lock.unlock(); 361   11031 lock.unlock();
359   362  
HITCBC 360   39004 task_cleanup on_exit{this, &lock, ctx}; 363   39127 task_cleanup on_exit{this, &lock, ctx};
361   364  
362   // Flush deferred timerfd programming before blocking 365   // Flush deferred timerfd programming before blocking
HITCBC 363   39004 if (timerfd_stale_.exchange(false, std::memory_order_acquire)) 366   39127 if (timerfd_stale_.exchange(false, std::memory_order_acquire))
HITCBC 364   3716 update_timerfd(); 367   3605 update_timerfd();
365   368  
HITCBC 366   39003 int nfds = ::epoll_wait( 369   39126 int nfds = ::epoll_wait(
ECB 367 - 39003 epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()), 370 + epoll_fd_, event_buffer_.data(),
HITGIC 368 - timeout_ms); 371 + 39126 static_cast<int>(event_buffer_.size()), timeout_ms);
369   372  
HITCBC 370   39003 if (nfds < 0 && errno != EINTR) 373   39126 if (nfds < 0 && errno != EINTR)
HITCBC 371   1 detail::throw_system_error(make_err(errno), "epoll_wait"); 374   1 detail::throw_system_error(make_err(errno), "epoll_wait");
372   375  
HITCBC 373   39002 bool check_timers = false; 376   39125 bool check_timers = false;
HITCBC 374   39002 ready_queue local_ops; 377   39125 ready_queue local_ops;
375   378  
HITCBC 376   86934 for (int i = 0; i < nfds; ++i) 379   86230 for (int i = 0; i < nfds; ++i)
377   { 380   {
HITCBC 378   47932 if (event_buffer_[i].data.ptr == nullptr) 381   47105 if (event_buffer_[i].data.ptr == nullptr)
379   { 382   {
380   std::uint64_t val; 383   std::uint64_t val;
381   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection) 384   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
HITCBC 382   4079 [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val)); 385   3970 [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
HITCBC 383   4079 eventfd_armed_.store(false, std::memory_order_relaxed); 386   3970 eventfd_armed_.store(false, std::memory_order_relaxed);
HITCBC 384   4079 continue; 387   3970 continue;
HITCBC 385   4079 } 388   3970 }
386   389  
HITCBC 387   43853 if (event_buffer_[i].data.ptr == &timer_fd_) 390   43135 if (event_buffer_[i].data.ptr == &timer_fd_)
388   { 391   {
389   std::uint64_t expirations; 392   std::uint64_t expirations;
390   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection) 393   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
391   [[maybe_unused]] auto r = 394   [[maybe_unused]] auto r =
HITCBC 392   7426 ::read(timer_fd_, &expirations, sizeof(expirations)); 395   7245 ::read(timer_fd_, &expirations, sizeof(expirations));
HITCBC 393   7426 check_timers = true; 396   7245 check_timers = true;
HITCBC 394   7426 continue; 397   7245 continue;
HITCBC 395   7426 } 398   7245 }
396   399  
397   auto* desc = 400   auto* desc =
HITCBC 398   36427 static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr); 401   35890 static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
HITCBC 399   36427 desc->add_ready_events(event_buffer_[i].events); 402   35890 desc->add_ready_events(event_buffer_[i].events);
400   403  
HITCBC 401   36427 bool expected = false; 404   35890 bool expected = false;
HITCBC 402   36427 if (desc->is_enqueued_.compare_exchange_strong( 405   35890 if (desc->is_enqueued_.compare_exchange_strong(
403   expected, true, std::memory_order_release, 406   expected, true, std::memory_order_release,
404   std::memory_order_relaxed)) 407   std::memory_order_relaxed))
405   { 408   {
HITCBC 406   36427 local_ops.push(desc); 409   35890 local_ops.push(desc);
407   } 410   }
408   } 411   }
409   412  
HITCBC 410   39002 if (check_timers) 413   39125 if (check_timers)
411   { 414   {
HITCBC 412   7426 timer_svc_->process_expired(); 415   7245 timer_svc_->process_expired();
HITCBC 413   7426 update_timerfd(); 416   7245 update_timerfd();
414   } 417   }
415   418  
HITCBC 416   39002 lock.lock(); 419   39125 lock.lock();
417   420  
HITCBC 418   39002 completed_ops_.splice(local_ops); 421   39125 completed_ops_.splice(local_ops);
HITCBC 419   39004 } 422   39127 }
420   423  
421   } // namespace boost::corosio::detail 424   } // namespace boost::corosio::detail
422   425  
423   #endif // BOOST_COROSIO_HAS_EPOLL 426   #endif // BOOST_COROSIO_HAS_EPOLL
424   427  
425   #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP 428   #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP