86.18% Lines (131/152) 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 - 51 void register_signal_reader(int read_fd) override 130 + [[nodiscard]] std::error_code
HITGNC   131 + 51 register_signal_reader(int read_fd) override
131   { 132   {
HITCBC 132 - 51 if (auto ec = register_descriptor(read_fd, signal_pipe_reader_.arm())) 133 + 51 return register_descriptor(read_fd, signal_pipe_reader_.arm());
DUB 133 - detail::throw_system_error(ec, "epoll_ctl (register)");  
ECB 134   51 } 134   }
135   135  
136   private: 136   private:
137   void 137   void
138   run_task(lock_type& lock, context_type* ctx, 138   run_task(lock_type& lock, context_type* ctx,
139   long timeout_us) override; 139   long timeout_us) override;
140   void interrupt_reactor() const override; 140   void interrupt_reactor() const override;
141   void update_timerfd() const; 141   void update_timerfd() const;
142   142  
143   int epoll_fd_; 143   int epoll_fd_;
144   int event_fd_; 144   int event_fd_;
145   int timer_fd_; 145   int timer_fd_;
146   146  
147   // 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
148   // register_signal_reader on the first signal registration). 148   // register_signal_reader on the first signal registration).
149   reactor_signal_pipe_reader signal_pipe_reader_; 149   reactor_signal_pipe_reader signal_pipe_reader_;
150   150  
151   // Edge-triggered eventfd state 151   // Edge-triggered eventfd state
152   mutable std::atomic<bool> eventfd_armed_{false}; 152   mutable std::atomic<bool> eventfd_armed_{false};
153   153  
154   // Set when the earliest timer changes; flushed before epoll_wait 154   // Set when the earliest timer changes; flushed before epoll_wait
155   mutable std::atomic<bool> timerfd_stale_{false}; 155   mutable std::atomic<bool> timerfd_stale_{false};
156   156  
157   // Event buffer sized from max_events_per_poll_ (set at construction, 157   // Event buffer sized from max_events_per_poll_ (set at construction,
158   // resized by configure_reactor via io_context_options). 158   // resized by configure_reactor via io_context_options).
159   std::vector<epoll_event> event_buffer_; 159   std::vector<epoll_event> event_buffer_;
160   }; 160   };
161   161  
HITCBC 162   867 inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int) 162   902 inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
HITCBC 163   867 : epoll_fd_(-1) 163   902 : epoll_fd_(-1)
HITCBC 164   867 , event_fd_(-1) 164   902 , event_fd_(-1)
HITCBC 165   867 , timer_fd_(-1) 165   902 , timer_fd_(-1)
HITCBC 166   1734 , event_buffer_(max_events_per_poll_) 166   1804 , event_buffer_(max_events_per_poll_)
167   { 167   {
HITCBC 168   867 epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC); 168   902 epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
HITCBC 169   867 if (epoll_fd_ < 0) 169   902 if (epoll_fd_ < 0)
MISUBC 170   detail::throw_system_error(make_err(errno), "epoll_create1"); 170   detail::throw_system_error(make_err(errno), "epoll_create1");
171   171  
HITCBC 172   867 event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC); 172   902 event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
HITCBC 173   867 if (event_fd_ < 0) 173   902 if (event_fd_ < 0)
174   { 174   {
MISUBC 175   int errn = errno; 175   int errn = errno;
MISUBC 176   ::close(epoll_fd_); 176   ::close(epoll_fd_);
MISUBC 177   detail::throw_system_error(make_err(errn), "eventfd"); 177   detail::throw_system_error(make_err(errn), "eventfd");
178   } 178   }
179   179  
HITCBC 180   867 timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC); 180   902 timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
HITCBC 181   867 if (timer_fd_ < 0) 181   902 if (timer_fd_ < 0)
182   { 182   {
MISUBC 183   int errn = errno; 183   int errn = errno;
MISUBC 184   ::close(event_fd_); 184   ::close(event_fd_);
MISUBC 185   ::close(epoll_fd_); 185   ::close(epoll_fd_);
MISUBC 186   detail::throw_system_error(make_err(errn), "timerfd_create"); 186   detail::throw_system_error(make_err(errn), "timerfd_create");
187   } 187   }
188   188  
HITCBC 189   867 epoll_event ev{}; 189   902 epoll_event ev{};
HITCBC 190   867 ev.events = EPOLLIN | EPOLLET; 190   902 ev.events = EPOLLIN | EPOLLET;
HITCBC 191   867 ev.data.ptr = nullptr; 191   902 ev.data.ptr = nullptr;
HITCBC 192   867 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0) 192   902 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0)
193   { 193   {
MISUBC 194   int errn = errno; 194   int errn = errno;
MISUBC 195   ::close(timer_fd_); 195   ::close(timer_fd_);
MISUBC 196   ::close(event_fd_); 196   ::close(event_fd_);
MISUBC 197   ::close(epoll_fd_); 197   ::close(epoll_fd_);
MISUBC 198   detail::throw_system_error(make_err(errn), "epoll_ctl"); 198   detail::throw_system_error(make_err(errn), "epoll_ctl");
199   } 199   }
200   200  
HITCBC 201   867 epoll_event timer_ev{}; 201   902 epoll_event timer_ev{};
HITCBC 202   867 timer_ev.events = EPOLLIN | EPOLLERR; 202   902 timer_ev.events = EPOLLIN | EPOLLERR;
HITCBC 203   867 timer_ev.data.ptr = &timer_fd_; 203   902 timer_ev.data.ptr = &timer_fd_;
HITCBC 204   867 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0) 204   902 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0)
205   { 205   {
MISUBC 206   int errn = errno; 206   int errn = errno;
MISUBC 207   ::close(timer_fd_); 207   ::close(timer_fd_);
MISUBC 208   ::close(event_fd_); 208   ::close(event_fd_);
MISUBC 209   ::close(epoll_fd_); 209   ::close(epoll_fd_);
MISUBC 210   detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)"); 210   detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)");
211   } 211   }
212   212  
HITCBC 213   867 timer_svc_ = &get_timer_service(ctx, *this); 213   902 timer_svc_ = &get_timer_service(ctx, *this);
HITCBC 214   867 timer_svc_->set_on_earliest_changed( 214   902 timer_svc_->set_on_earliest_changed(
HITCBC 215   5055 timer_service::callback(this, [](void* p) { 215   5467 timer_service::callback(this, [](void* p) {
HITCBC 216   4188 auto* self = static_cast<epoll_scheduler*>(p); 216   4565 auto* self = static_cast<epoll_scheduler*>(p);
HITCBC 217   4188 self->timerfd_stale_.store(true, std::memory_order_release); 217   4565 self->timerfd_stale_.store(true, std::memory_order_release);
HITCBC 218   4188 self->interrupt_reactor(); 218   4565 self->interrupt_reactor();
HITCBC 219   4188 })); 219   4565 }));
220   220  
HITCBC 221   867 get_resolver_service(ctx, *this); 221   902 get_resolver_service(ctx, *this);
HITCBC 222   867 get_signal_service(ctx, *this); 222   902 get_signal_service(ctx, *this);
HITCBC 223   867 get_stream_file_service(ctx, *this); 223   902 get_stream_file_service(ctx, *this);
HITCBC 224   867 get_random_access_file_service(ctx, *this); 224   902 get_random_access_file_service(ctx, *this);
225   225  
HITCBC 226   867 completed_ops_.push(&task_op_); 226   902 completed_ops_.push(&task_op_);
HITCBC 227   867 } 227   902 }
228   228  
HITCBC 229   1734 inline epoll_scheduler::~epoll_scheduler() 229   1804 inline epoll_scheduler::~epoll_scheduler()
230   { 230   {
HITCBC 231   867 if (timer_fd_ >= 0) 231   902 if (timer_fd_ >= 0)
HITCBC 232   867 ::close(timer_fd_); 232   902 ::close(timer_fd_);
HITCBC 233   867 if (event_fd_ >= 0) 233   902 if (event_fd_ >= 0)
HITCBC 234   867 ::close(event_fd_); 234   902 ::close(event_fd_);
HITCBC 235   867 if (epoll_fd_ >= 0) 235   902 if (epoll_fd_ >= 0)
HITCBC 236   867 ::close(epoll_fd_); 236   902 ::close(epoll_fd_);
HITCBC 237   1734 } 237   1804 }
238   238  
239   inline void 239   inline void
HITCBC 240   867 epoll_scheduler::shutdown() 240   902 epoll_scheduler::shutdown()
241   { 241   {
HITCBC 242   867 shutdown_drain(); 242   902 shutdown_drain();
243   243  
HITCBC 244   867 if (event_fd_ >= 0) 244   902 if (event_fd_ >= 0)
HITCBC 245   867 interrupt_reactor(); 245   902 interrupt_reactor();
HITCBC 246   867 } 246   902 }
247   247  
248   inline void 248   inline void
HITCBC 249   19 epoll_scheduler::configure_reactor( 249   19 epoll_scheduler::configure_reactor(
250   unsigned max_events, 250   unsigned max_events,
251   unsigned budget_init, 251   unsigned budget_init,
252   unsigned budget_max, 252   unsigned budget_max,
253   unsigned unassisted) 253   unsigned unassisted)
254   { 254   {
HITCBC 255   19 reactor_scheduler::configure_reactor( 255   19 reactor_scheduler::configure_reactor(
256   max_events, budget_init, budget_max, unassisted); 256   max_events, budget_init, budget_max, unassisted);
HITCBC 257   18 event_buffer_.resize(max_events_per_poll_); 257   18 event_buffer_.resize(max_events_per_poll_);
HITCBC 258   18 } 258   18 }
259   259  
260   inline std::error_code 260   inline std::error_code
HITCBC 261   7475 epoll_scheduler::register_descriptor(int fd, reactor_descriptor_state* desc) const 261   8085 epoll_scheduler::register_descriptor(int fd, reactor_descriptor_state* desc) const
262   { 262   {
HITCBC 263   7475 epoll_event ev{}; 263   8085 epoll_event ev{};
HITCBC 264   7475 ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP; 264   8085 ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
HITCBC 265   7475 ev.data.ptr = desc; 265   8085 ev.data.ptr = desc;
266   266  
HITCBC 267   7475 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0) 267   8085 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
HITCBC 268   1 return make_err(errno); 268   1 return make_err(errno);
269   269  
HITCBC 270   7474 desc->registered_events = ev.events; 270   8084 desc->registered_events = ev.events;
HITCBC 271   7474 desc->fd = fd; 271   8084 desc->fd = fd;
HITCBC 272   7474 desc->scheduler_ = this; 272   8084 desc->scheduler_ = this;
HITCBC 273   7474 desc->mutex.set_enabled(reactor_io_locking_); 273   8084 desc->mutex.set_enabled(reactor_io_locking_);
HITCBC 274   7474 desc->ready_events_.store(0, std::memory_order_relaxed); 274   8084 desc->ready_events_.store(0, std::memory_order_relaxed);
275   275  
HITCBC 276   7474 conditionally_enabled_mutex::scoped_lock lock(desc->mutex); 276   8084 conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
HITCBC 277   7474 desc->impl_ref_.reset(); 277   8084 desc->impl_ref_.reset();
HITCBC 278   7474 desc->read_ready = false; 278   8084 desc->read_ready = false;
HITCBC 279   7474 desc->write_ready = false; 279   8084 desc->write_ready = false;
HITCBC 280   7474 return {}; 280   8084 return {};
HITCBC 281   7474 } 281   8084 }
282   282  
283   inline void 283   inline void
HITCBC 284   7423 epoll_scheduler::deregister_descriptor(int fd) const 284   8033 epoll_scheduler::deregister_descriptor(int fd) const
285   { 285   {
HITCBC 286   7423 ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr); 286   8033 ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
HITCBC 287   7423 } 287   8033 }
288   288  
289   inline void 289   inline void
HITCBC 290   5841 epoll_scheduler::interrupt_reactor() const 290   6296 epoll_scheduler::interrupt_reactor() const
291   { 291   {
HITCBC 292   5841 bool expected = false; 292   6296 bool expected = false;
HITCBC 293   5841 if (eventfd_armed_.compare_exchange_strong( 293   6296 if (eventfd_armed_.compare_exchange_strong(
294   expected, true, std::memory_order_release, 294   expected, true, std::memory_order_release,
295   std::memory_order_relaxed)) 295   std::memory_order_relaxed))
296   { 296   {
HITCBC 297   4703 std::uint64_t val = 1; 297   5058 std::uint64_t val = 1;
HITCBC 298   4703 [[maybe_unused]] auto r = ::write(event_fd_, &val, sizeof(val)); 298   5058 [[maybe_unused]] auto r = ::write(event_fd_, &val, sizeof(val));
299   } 299   }
HITCBC 300   5841 } 300   6296 }
301   301  
302   inline void 302   inline void
HITCBC 303   7192 epoll_scheduler::update_timerfd() const 303   7779 epoll_scheduler::update_timerfd() const
304   { 304   {
HITCBC 305   7192 auto nearest = timer_svc_->nearest_expiry(); 305   7779 auto nearest = timer_svc_->nearest_expiry();
306   306  
HITCBC 307   7192 itimerspec ts{}; 307   7779 itimerspec ts{};
HITCBC 308   7192 int flags = 0; 308   7779 int flags = 0;
309   309  
HITCBC 310   7192 if (nearest == timer_service::time_point::max()) 310   7779 if (nearest == timer_service::time_point::max())
311   { 311   {
312   // No timers — disarm by setting to 0 (relative) 312   // No timers — disarm by setting to 0 (relative)
313   } 313   }
314   else 314   else
315   { 315   {
HITCBC 316   7015 auto now = std::chrono::steady_clock::now(); 316   7602 auto now = std::chrono::steady_clock::now();
HITCBC 317   7015 if (nearest <= now) 317   7602 if (nearest <= now)
318   { 318   {
319   // Use 1ns instead of 0 — zero disarms the timerfd 319   // Use 1ns instead of 0 — zero disarms the timerfd
HITCBC 320   483 ts.it_value.tv_nsec = 1; 320   896 ts.it_value.tv_nsec = 1;
321   } 321   }
322   else 322   else
323   { 323   {
HITCBC 324   6532 auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>( 324   6706 auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
HITCBC 325   6532 nearest - now) 325   6706 nearest - now)
HITCBC 326   6532 .count(); 326   6706 .count();
HITCBC 327   6532 ts.it_value.tv_sec = nsec / 1000000000; 327   6706 ts.it_value.tv_sec = nsec / 1000000000;
HITCBC 328   6532 ts.it_value.tv_nsec = nsec % 1000000000; 328   6706 ts.it_value.tv_nsec = nsec % 1000000000;
HITCBC 329   6532 if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0) 329   6706 if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
MISUBC 330   ts.it_value.tv_nsec = 1; 330   ts.it_value.tv_nsec = 1;
331   } 331   }
332   } 332   }
333   333  
HITCBC 334   7192 if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0) 334   7779 if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
MISUBC 335   detail::throw_system_error(make_err(errno), "timerfd_settime"); 335   detail::throw_system_error(make_err(errno), "timerfd_settime");
HITCBC 336   7192 } 336   7779 }
337   337  
338   inline void 338   inline void
HITCBC 339   33211 epoll_scheduler::run_task( 339   49383 epoll_scheduler::run_task(
340   lock_type& lock, context_type* ctx, long timeout_us) 340   lock_type& lock, context_type* ctx, long timeout_us)
341   { 341   {
342   int timeout_ms; 342   int timeout_ms;
HITCBC 343   33211 if (task_interrupted_) 343   49383 if (task_interrupted_)
HITCBC 344   24425 timeout_ms = 0; 344   38744 timeout_ms = 0;
HITCBC 345   8786 else if (timeout_us < 0) 345   10639 else if (timeout_us < 0)
HITCBC 346   8782 timeout_ms = -1; 346   10636 timeout_ms = -1;
347   else 347   else
HITCBC 348   4 timeout_ms = static_cast<int>((timeout_us + 999) / 1000); 348   3 timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
349   349  
HITCBC 350   33211 if (lock.owns_lock()) 350   49383 if (lock.owns_lock())
HITCBC 351   8788 lock.unlock(); 351   10641 lock.unlock();
352   352  
HITCBC 353   33211 task_cleanup on_exit{this, &lock, ctx}; 353   49383 task_cleanup on_exit{this, &lock, ctx};
354   354  
355   // Flush deferred timerfd programming before blocking 355   // Flush deferred timerfd programming before blocking
HITCBC 356   33211 if (timerfd_stale_.exchange(false, std::memory_order_acquire)) 356   49383 if (timerfd_stale_.exchange(false, std::memory_order_acquire))
HITCBC 357   3603 update_timerfd(); 357   3897 update_timerfd();
358   358  
HITCBC 359   33211 int nfds = ::epoll_wait( 359   49383 int nfds = ::epoll_wait(
360   epoll_fd_, event_buffer_.data(), 360   epoll_fd_, event_buffer_.data(),
HITCBC 361   33211 static_cast<int>(event_buffer_.size()), timeout_ms); 361   49383 static_cast<int>(event_buffer_.size()), timeout_ms);
362   362  
HITCBC 363   33211 if (nfds < 0 && errno != EINTR) 363   49383 if (nfds < 0 && errno != EINTR)
MISUBC 364   detail::throw_system_error(make_err(errno), "epoll_wait"); 364   detail::throw_system_error(make_err(errno), "epoll_wait");
365   365  
HITCBC 366   33211 bool check_timers = false; 366   49383 bool check_timers = false;
HITCBC 367   33211 ready_queue local_ops; 367   49383 ready_queue local_ops;
368   368  
HITCBC 369   79244 for (int i = 0; i < nfds; ++i) 369   113995 for (int i = 0; i < nfds; ++i)
370   { 370   {
HITCBC 371   46033 if (event_buffer_[i].data.ptr == nullptr) 371   64612 if (event_buffer_[i].data.ptr == nullptr)
372   { 372   {
373   std::uint64_t val; 373   std::uint64_t val;
374   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection) 374   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
HITCBC 375   3836 [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val)); 375   4156 [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
HITCBC 376   3836 eventfd_armed_.store(false, std::memory_order_relaxed); 376   4156 eventfd_armed_.store(false, std::memory_order_relaxed);
HITCBC 377   3836 continue; 377   4156 continue;
HITCBC 378   3836 } 378   4156 }
379   379  
HITCBC 380   42197 if (event_buffer_[i].data.ptr == &timer_fd_) 380   60456 if (event_buffer_[i].data.ptr == &timer_fd_)
381   { 381   {
382   std::uint64_t expirations; 382   std::uint64_t expirations;
383   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection) 383   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
384   [[maybe_unused]] auto r = 384   [[maybe_unused]] auto r =
HITCBC 385   3589 ::read(timer_fd_, &expirations, sizeof(expirations)); 385   3882 ::read(timer_fd_, &expirations, sizeof(expirations));
HITCBC 386   3589 check_timers = true; 386   3882 check_timers = true;
HITCBC 387   3589 continue; 387   3882 continue;
HITCBC 388   3589 } 388   3882 }
389   389  
390   auto* desc = 390   auto* desc =
HITCBC 391   38608 static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr); 391   56574 static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
HITCBC 392   38608 desc->add_ready_events(event_buffer_[i].events); 392   56574 desc->add_ready_events(event_buffer_[i].events);
393   393  
HITCBC 394   38608 bool expected = false; 394   56574 bool expected = false;
HITCBC 395   38608 if (desc->is_enqueued_.compare_exchange_strong( 395   56574 if (desc->is_enqueued_.compare_exchange_strong(
396   expected, true, std::memory_order_release, 396   expected, true, std::memory_order_release,
397   std::memory_order_relaxed)) 397   std::memory_order_relaxed))
398   { 398   {
HITCBC 399   38608 local_ops.push(desc); 399   56574 local_ops.push(desc);
400   } 400   }
401   } 401   }
402   402  
HITCBC 403   33211 if (check_timers) 403   49383 if (check_timers)
404   { 404   {
HITCBC 405   3589 timer_svc_->process_expired(); 405   3882 timer_svc_->process_expired();
HITCBC 406   3589 update_timerfd(); 406   3882 update_timerfd();
407   } 407   }
408   408  
HITCBC 409   33211 lock.lock(); 409   49383 lock.lock();
410   410  
HITCBC 411   33211 completed_ops_.splice(local_ops); 411   49383 completed_ops_.splice(local_ops);
HITCBC 412   33211 } 412   49383 }
413   413  
414   } // namespace boost::corosio::detail 414   } // namespace boost::corosio::detail
415   415  
416   #endif // BOOST_COROSIO_HAS_EPOLL 416   #endif // BOOST_COROSIO_HAS_EPOLL
417   417  
418   #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP 418   #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP