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