99.38% Lines (159/160) 100.00% Functions (12/12)
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   27  
28   #include <boost/corosio/detail/except.hpp> 28   #include <boost/corosio/detail/except.hpp>
29   29  
30   #include <atomic> 30   #include <atomic>
31   #include <chrono> 31   #include <chrono>
32   #include <cstdint> 32   #include <cstdint>
33   #include <mutex> 33   #include <mutex>
34   #include <vector> 34   #include <vector>
35   35  
36   #include <errno.h> 36   #include <errno.h>
37   #include <sys/epoll.h> 37   #include <sys/epoll.h>
38   #include <sys/eventfd.h> 38   #include <sys/eventfd.h>
39   #include <sys/timerfd.h> 39   #include <sys/timerfd.h>
40   #include <unistd.h> 40   #include <unistd.h>
41   41  
42   namespace boost::corosio::detail { 42   namespace boost::corosio::detail {
43   43  
44   /** Linux scheduler using epoll for I/O multiplexing. 44   /** Linux scheduler using epoll for I/O multiplexing.
45   45  
46   This scheduler implements the scheduler interface using Linux epoll 46   This scheduler implements the scheduler interface using Linux epoll
47   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
48   where one thread runs epoll_wait while other threads 48   where one thread runs epoll_wait while other threads
49   wait on a condition variable for handler work. This design provides: 49   wait on a condition variable for handler work. This design provides:
50   50  
51   - Handler parallelism: N posted handlers can execute on N threads 51   - Handler parallelism: N posted handlers can execute on N threads
52   - No thundering herd: condition_variable wakes exactly one thread 52   - No thundering herd: condition_variable wakes exactly one thread
53   - IOCP parity: Behavior matches Windows I/O completion port semantics 53   - IOCP parity: Behavior matches Windows I/O completion port semantics
54   54  
55   When threads call run(), they first try to execute queued handlers. 55   When threads call run(), they first try to execute queued handlers.
56   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
57   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
58   variable until handlers are available. 58   variable until handlers are available.
59   59  
60   @par Thread Safety 60   @par Thread Safety
61   All public member functions are thread-safe. 61   All public member functions are thread-safe.
62   */ 62   */
63   class BOOST_COROSIO_DECL epoll_scheduler final : public reactor_scheduler 63   class BOOST_COROSIO_DECL epoll_scheduler final : public reactor_scheduler
64   { 64   {
65   public: 65   public:
66   /** Construct the scheduler. 66   /** Construct the scheduler.
67   67  
68   Creates an epoll instance, eventfd for reactor interruption, 68   Creates an epoll instance, eventfd for reactor interruption,
69   and timerfd for kernel-managed timer expiry. 69   and timerfd for kernel-managed timer expiry.
70   70  
71   @param ctx Reference to the owning execution_context. 71   @param ctx Reference to the owning execution_context.
72   @param concurrency_hint Hint for expected thread count (unused). 72   @param concurrency_hint Hint for expected thread count (unused).
73   */ 73   */
74   epoll_scheduler(capy::execution_context& ctx, int concurrency_hint = -1); 74   epoll_scheduler(capy::execution_context& ctx, int concurrency_hint = -1);
75   75  
76   /// Destroy the scheduler. 76   /// Destroy the scheduler.
77   ~epoll_scheduler() override; 77   ~epoll_scheduler() override;
78   78  
79   epoll_scheduler(epoll_scheduler const&) = delete; 79   epoll_scheduler(epoll_scheduler const&) = delete;
80   epoll_scheduler& operator=(epoll_scheduler const&) = delete; 80   epoll_scheduler& operator=(epoll_scheduler const&) = delete;
81   81  
82   /// Shut down the scheduler, draining pending operations. 82   /// Shut down the scheduler, draining pending operations.
83   void shutdown() override; 83   void shutdown() override;
84   84  
85   /// Apply runtime configuration, resizing the event buffer. 85   /// Apply runtime configuration, resizing the event buffer.
86   void configure_reactor( 86   void configure_reactor(
87   unsigned max_events, 87   unsigned max_events,
88   unsigned budget_init, 88   unsigned budget_init,
89   unsigned budget_max, 89   unsigned budget_max,
90   unsigned unassisted) override; 90   unsigned unassisted) override;
91   91  
92   /** Return the epoll file descriptor. 92   /** Return the epoll file descriptor.
93   93  
94   Used by socket services to register file descriptors 94   Used by socket services to register file descriptors
95   for I/O event notification. 95   for I/O event notification.
96   96  
97   @return The epoll file descriptor. 97   @return The epoll file descriptor.
98   */ 98   */
99   int epoll_fd() const noexcept 99   int epoll_fd() const noexcept
100   { 100   {
101   return epoll_fd_; 101   return epoll_fd_;
102   } 102   }
103   103  
104   /** Register a descriptor for persistent monitoring. 104   /** Register a descriptor for persistent monitoring.
105   105  
106   The fd is registered once and stays registered until explicitly 106   The fd is registered once and stays registered until explicitly
107   deregistered. Events are dispatched via reactor_descriptor_state which 107   deregistered. Events are dispatched via reactor_descriptor_state which
108   tracks pending read/write/connect operations. 108   tracks pending read/write/connect operations.
109   109  
110   @param fd The file descriptor to register. 110   @param fd The file descriptor to register.
111   @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).
112   112  
113   @return The error if registration fails, otherwise a default 113   @return The error if registration fails, otherwise a default
114   constructed error code. 114   constructed error code.
115   */ 115   */
116   std::error_code 116   std::error_code
117   register_descriptor(int fd, reactor_descriptor_state* desc) const; 117   register_descriptor(int fd, reactor_descriptor_state* desc) const;
118   118  
  119 + /// No-op: write readiness is watched from registration on.
  120 + std::error_code
HITGNC   121 + 2618 ensure_write_registered(int, reactor_descriptor_state*) const noexcept
  122 + {
HITGNC   123 + 2618 return {};
  124 + }
  125 +
119   /** Deregister a persistently registered descriptor. 126   /** Deregister a persistently registered descriptor.
120   127  
121   @param fd The file descriptor to deregister. 128   @param fd The file descriptor to deregister.
122   */ 129   */
123   void deregister_descriptor(int fd) const; 130   void deregister_descriptor(int fd) const;
124   131  
125   /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp). 132   /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp).
HITCBC 126   76 [[nodiscard]] std::error_code register_signal_reader(int read_fd) override 133   76 [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
127   { 134   {
HITCBC 128   76 return register_descriptor(read_fd, signal_pipe_reader_.arm()); 135   76 return register_descriptor(read_fd, signal_pipe_reader_.arm());
129   } 136   }
130   137  
131   private: 138   private:
132   void run_task(lock_type& lock, context_type& ctx, long timeout_us) override; 139   void run_task(lock_type& lock, context_type& ctx, long timeout_us) override;
133   void interrupt_reactor() const override; 140   void interrupt_reactor() const override;
134   void update_timerfd() const; 141   void update_timerfd() const;
135   142  
136   int epoll_fd_; 143   int epoll_fd_;
137   int event_fd_; 144   int event_fd_;
138   int timer_fd_; 145   int timer_fd_;
139   146  
140   // 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
141   // register_signal_reader on the first signal registration). 148   // register_signal_reader on the first signal registration).
142   reactor_signal_pipe_reader signal_pipe_reader_; 149   reactor_signal_pipe_reader signal_pipe_reader_;
143   150  
144   // Edge-triggered eventfd state 151   // Edge-triggered eventfd state
145   mutable std::atomic<bool> eventfd_armed_{false}; 152   mutable std::atomic<bool> eventfd_armed_{false};
146   153  
147   // Set when the earliest timer changes; flushed before epoll_wait 154   // Set when the earliest timer changes; flushed before epoll_wait
148   mutable std::atomic<bool> timerfd_stale_{false}; 155   mutable std::atomic<bool> timerfd_stale_{false};
149   156  
150   // Event buffer sized from max_events_per_poll_ (set at construction, 157   // Event buffer sized from max_events_per_poll_ (set at construction,
151   // resized by configure_reactor via io_context_options). 158   // resized by configure_reactor via io_context_options).
152   std::vector<epoll_event> event_buffer_; 159   std::vector<epoll_event> event_buffer_;
153   }; 160   };
154   161  
HITCBC 155   1305 inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int) 162   1577 inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
HITCBC 156   1305 : epoll_fd_(-1) 163   1577 : epoll_fd_(-1)
HITCBC 157   1305 , event_fd_(-1) 164   1577 , event_fd_(-1)
HITCBC 158   1305 , timer_fd_(-1) 165   1577 , timer_fd_(-1)
HITCBC 159   2610 , event_buffer_(max_events_per_poll_) 166   3154 , event_buffer_(max_events_per_poll_)
160   { 167   {
HITCBC 161   1305 epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC); 168   1577 epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
HITCBC 162   1305 if (epoll_fd_ < 0) 169   1577 if (epoll_fd_ < 0)
HITCBC 163   1 detail::throw_system_error(make_err(errno), "epoll_create1"); 170   1 detail::throw_system_error(make_err(errno), "epoll_create1");
164   171  
HITCBC 165   1304 event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC); 172   1576 event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
HITCBC 166   1304 if (event_fd_ < 0) 173   1576 if (event_fd_ < 0)
167   { 174   {
HITCBC 168   1 int errn = errno; 175   1 int errn = errno;
HITCBC 169   1 ::close(epoll_fd_); 176   1 ::close(epoll_fd_);
HITCBC 170   1 detail::throw_system_error(make_err(errn), "eventfd"); 177   1 detail::throw_system_error(make_err(errn), "eventfd");
171   } 178   }
172   179  
HITCBC 173   1303 timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC); 180   1575 timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
HITCBC 174   1303 if (timer_fd_ < 0) 181   1575 if (timer_fd_ < 0)
175   { 182   {
HITCBC 176   1 int errn = errno; 183   1 int errn = errno;
HITCBC 177   1 ::close(event_fd_); 184   1 ::close(event_fd_);
HITCBC 178   1 ::close(epoll_fd_); 185   1 ::close(epoll_fd_);
HITCBC 179   1 detail::throw_system_error(make_err(errn), "timerfd_create"); 186   1 detail::throw_system_error(make_err(errn), "timerfd_create");
180   } 187   }
181   188  
HITCBC 182   1302 epoll_event ev{}; 189   1574 epoll_event ev{};
HITCBC 183   1302 ev.events = EPOLLIN | EPOLLET; 190   1574 ev.events = EPOLLIN | EPOLLET;
HITCBC 184   1302 ev.data.ptr = nullptr; 191   1574 ev.data.ptr = nullptr;
HITCBC 185   1302 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0) 192   1574 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0)
186   { 193   {
HITCBC 187   1 int errn = errno; 194   1 int errn = errno;
HITCBC 188   1 ::close(timer_fd_); 195   1 ::close(timer_fd_);
HITCBC 189   1 ::close(event_fd_); 196   1 ::close(event_fd_);
HITCBC 190   1 ::close(epoll_fd_); 197   1 ::close(epoll_fd_);
HITCBC 191   1 detail::throw_system_error(make_err(errn), "epoll_ctl"); 198   1 detail::throw_system_error(make_err(errn), "epoll_ctl");
192   } 199   }
193   200  
HITCBC 194   1301 epoll_event timer_ev{}; 201   1573 epoll_event timer_ev{};
HITCBC 195   1301 timer_ev.events = EPOLLIN | EPOLLERR; 202   1573 timer_ev.events = EPOLLIN | EPOLLERR;
HITCBC 196   1301 timer_ev.data.ptr = &timer_fd_; 203   1573 timer_ev.data.ptr = &timer_fd_;
HITCBC 197   1301 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0) 204   1573 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0)
198   { 205   {
HITCBC 199   1 int errn = errno; 206   1 int errn = errno;
HITCBC 200   1 ::close(timer_fd_); 207   1 ::close(timer_fd_);
HITCBC 201   1 ::close(event_fd_); 208   1 ::close(event_fd_);
HITCBC 202   1 ::close(epoll_fd_); 209   1 ::close(epoll_fd_);
HITCBC 203   1 detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)"); 210   1 detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)");
204   } 211   }
205   212  
HITCBC 206   1300 timer_svc_ = &get_timer_service(ctx, *this); 213   1572 timer_svc_ = &get_timer_service(ctx, *this);
HITCBC 207   1300 timer_svc_->set_on_earliest_changed( 214   1572 timer_svc_->set_on_earliest_changed(
HITCBC 208   5580 timer_service::callback(this, [](void* p) { 215   5979 timer_service::callback(this, [](void* p) {
HITCBC 209   4280 auto* self = static_cast<epoll_scheduler*>(p); 216   4407 auto* self = static_cast<epoll_scheduler*>(p);
HITCBC 210   4280 self->timerfd_stale_.store(true, std::memory_order_release); 217   4407 self->timerfd_stale_.store(true, std::memory_order_release);
HITCBC 211   4280 self->interrupt_reactor(); 218   4407 self->interrupt_reactor();
HITCBC 212   4280 })); 219   4407 }));
213   220  
HITCBC 214   1300 completed_ops_.push(&task_op_); 221   1572 completed_ops_.push(&task_op_);
HITCBC 215   1315 } 222   1587 }
216   223  
HITCBC 217   2600 inline epoll_scheduler::~epoll_scheduler() 224   3144 inline epoll_scheduler::~epoll_scheduler()
218   { 225   {
HITCBC 219   1300 if (timer_fd_ >= 0) 226   1572 if (timer_fd_ >= 0)
HITCBC 220   1300 ::close(timer_fd_); 227   1572 ::close(timer_fd_);
HITCBC 221   1300 if (event_fd_ >= 0) 228   1572 if (event_fd_ >= 0)
HITCBC 222   1300 ::close(event_fd_); 229   1572 ::close(event_fd_);
HITCBC 223   1300 if (epoll_fd_ >= 0) 230   1572 if (epoll_fd_ >= 0)
HITCBC 224   1300 ::close(epoll_fd_); 231   1572 ::close(epoll_fd_);
HITCBC 225   2600 } 232   3144 }
226   233  
227   inline void 234   inline void
HITCBC 228   1300 epoll_scheduler::shutdown() 235   1572 epoll_scheduler::shutdown()
229   { 236   {
HITCBC 230   1300 shutdown_drain(); 237   1572 shutdown_drain();
231   238  
HITCBC 232   1300 if (event_fd_ >= 0) 239   1572 if (event_fd_ >= 0)
HITCBC 233   1300 interrupt_reactor(); 240   1572 interrupt_reactor();
HITCBC 234   1300 } 241   1572 }
235   242  
236   inline void 243   inline void
HITCBC 237   27 epoll_scheduler::configure_reactor( 244   28 epoll_scheduler::configure_reactor(
238   unsigned max_events, 245   unsigned max_events,
239   unsigned budget_init, 246   unsigned budget_init,
240   unsigned budget_max, 247   unsigned budget_max,
241   unsigned unassisted) 248   unsigned unassisted)
242   { 249   {
HITCBC 243   27 reactor_scheduler::configure_reactor( 250   28 reactor_scheduler::configure_reactor(
244   max_events, budget_init, budget_max, unassisted); 251   max_events, budget_init, budget_max, unassisted);
HITCBC 245   25 event_buffer_.resize(max_events_per_poll_); 252   26 event_buffer_.resize(max_events_per_poll_);
HITCBC 246   25 } 253   26 }
247   254  
248   inline std::error_code 255   inline std::error_code
HITCBC 249   5928 epoll_scheduler::register_descriptor( 256   5923 epoll_scheduler::register_descriptor(
250   int fd, reactor_descriptor_state* desc) const 257   int fd, reactor_descriptor_state* desc) const
251   { 258   {
HITCBC 252   5928 epoll_event ev{}; 259   5923 epoll_event ev{};
HITCBC 253   5928 ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP; 260   5923 ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
HITCBC 254   5928 ev.data.ptr = desc; 261   5923 ev.data.ptr = desc;
255   262  
HITGNC   263 + 5923 bool unpollable = false;
HITCBC 256   5928 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0) 264   5923 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
ECB 257 - 7 return make_err(errno); 265 + {
  266 + // EPERM: a file type epoll cannot watch. Its I/O does not
  267 + // block, so adopt it unwatched, as asio does.
HITGNC   268 + 11 if (errno != EPERM)
HITGNC   269 + 7 return make_err(errno);
HITGNC   270 + 4 unpollable = true;
  271 + }
258   272  
HITCBC 259 - 5921 desc->registered_events = ev.events; 273 + 5916 desc->registered_events = unpollable ? 0 : ev.events;
HITGNC   274 + 5916 desc->unpollable = unpollable;
HITCBC 260   5921 desc->fd = fd; 275   5916 desc->fd = fd;
HITCBC 261   5921 desc->scheduler_ = this; 276   5916 desc->scheduler_ = this;
HITCBC 262   5921 desc->mutex.set_enabled(reactor_io_locking_); 277   5916 desc->mutex.set_enabled(reactor_io_locking_);
HITCBC 263   5921 desc->ready_events_.store(0, std::memory_order_relaxed); 278   5916 desc->ready_events_.store(0, std::memory_order_relaxed);
264   279  
HITCBC 265   5921 conditionally_enabled_mutex::scoped_lock lock(desc->mutex); 280   5916 conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
HITCBC 266   5921 desc->impl_ref_.reset(); 281   5916 desc->impl_ref_.reset();
HITCBC 267   5921 desc->read_ready = false; 282   5916 desc->read_ready = false;
HITCBC 268   5921 desc->write_ready = false; 283   5916 desc->write_ready = false;
HITCBC 269   5921 return {}; 284   5916 return {};
HITCBC 270   5921 } 285   5916 }
271   286  
272   inline void 287   inline void
HITCBC 273   5846 epoll_scheduler::deregister_descriptor(int fd) const 288   5837 epoll_scheduler::deregister_descriptor(int fd) const
274   { 289   {
HITCBC 275   5846 ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr); 290   5837 ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
HITCBC 276   5846 } 291   5837 }
277   292  
278   inline void 293   inline void
HITCBC 279   7689 epoll_scheduler::interrupt_reactor() const 294   8498 epoll_scheduler::interrupt_reactor() const
280   { 295   {
HITCBC 281   7689 bool expected = false; 296   8498 bool expected = false;
HITCBC 282   7689 if (eventfd_armed_.compare_exchange_strong( 297   8498 if (eventfd_armed_.compare_exchange_strong(
283   expected, true, std::memory_order_release, 298   expected, true, std::memory_order_release,
284   std::memory_order_relaxed)) 299   std::memory_order_relaxed))
285   { 300   {
HITCBC 286   6218 std::uint64_t val = 1; 301   6645 std::uint64_t val = 1;
HITCBC 287   6218 if (::write(event_fd_, &val, sizeof(val)) < 0) 302   6645 if (::write(event_fd_, &val, sizeof(val)) < 0)
288   { 303   {
289   // The flag is what coalesces later interrupts into a byte 304   // The flag is what coalesces later interrupts into a byte
290   // already in the eventfd; a write that failed put no byte 305   // already in the eventfd; a write that failed put no byte
291   // there, so leaving it armed would swallow every interrupt 306   // there, so leaving it armed would swallow every interrupt
292   // that follows. Disarming keeps the cost to the interrupts 307   // that follows. Disarming keeps the cost to the interrupts
293   // already in flight -- the next one arms and writes again, 308   // already in flight -- the next one arms and writes again,
294   // instead of every one after this coalescing into a byte 309   // instead of every one after this coalescing into a byte
295   // that does not exist. 310   // that does not exist.
HITCBC 296   2 eventfd_armed_.store(false, std::memory_order_release); 311   2 eventfd_armed_.store(false, std::memory_order_release);
297   } 312   }
298   } 313   }
HITCBC 299   7689 } 314   8498 }
300   315  
301   inline void 316   inline void
HITCBC 302   11081 epoll_scheduler::update_timerfd() const 317   10780 epoll_scheduler::update_timerfd() const
303   { 318   {
HITCBC 304   11081 auto nearest = timer_svc_->nearest_expiry(); 319   10780 auto nearest = timer_svc_->nearest_expiry();
305   320  
HITCBC 306   11081 itimerspec ts{}; 321   10780 itimerspec ts{};
HITCBC 307   11081 int flags = 0; 322   10780 int flags = 0;
308   323  
HITCBC 309   11081 if (nearest == timer_service::time_point::max()) 324   10780 if (nearest == timer_service::time_point::max())
310   { 325   {
311   // No timers — disarm by setting to 0 (relative) 326   // No timers — disarm by setting to 0 (relative)
312   } 327   }
313   else 328   else
314   { 329   {
HITCBC 315   9894 auto now = std::chrono::steady_clock::now(); 330   9671 auto now = std::chrono::steady_clock::now();
HITCBC 316   9894 if (nearest <= now) 331   9671 if (nearest <= now)
317   { 332   {
318   // Use 1ns instead of 0 — zero disarms the timerfd 333   // Use 1ns instead of 0 — zero disarms the timerfd
HITCBC 319   1139 ts.it_value.tv_nsec = 1; 334   1464 ts.it_value.tv_nsec = 1;
320   } 335   }
321   else 336   else
322   { 337   {
HITCBC 323   8755 auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>( 338   8207 auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
HITCBC 324   8755 nearest - now) 339   8207 nearest - now)
HITCBC 325   8755 .count(); 340   8207 .count();
HITCBC 326   8755 ts.it_value.tv_sec = nsec / 1000000000; 341   8207 ts.it_value.tv_sec = nsec / 1000000000;
HITCBC 327   8755 ts.it_value.tv_nsec = nsec % 1000000000; 342   8207 ts.it_value.tv_nsec = nsec % 1000000000;
HITCBC 328   8755 if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0) 343   8207 if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
MISUBC 329   ✗ ts.it_value.tv_nsec = 1; 344   ✗ ts.it_value.tv_nsec = 1;
330   } 345   }
331   } 346   }
332   347  
HITCBC 333   11081 if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0) 348   10780 if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
HITCBC 334   1 detail::throw_system_error(make_err(errno), "timerfd_settime"); 349   1 detail::throw_system_error(make_err(errno), "timerfd_settime");
HITCBC 335   11080 } 350   10779 }
336   351  
337   inline void 352   inline void
HITCBC 338   42354 epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us) 353   41807 epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
339   { 354   {
340   int timeout_ms; 355   int timeout_ms;
HITCBC 341   42354 if (task_interrupted_) 356   41807 if (task_interrupted_)
HITCBC 342   31171 timeout_ms = 0; 357   30227 timeout_ms = 0;
HITCBC 343   11183 else if (timeout_us < 0) 358   11580 else if (timeout_us < 0)
HITCBC 344   10743 timeout_ms = -1; 359   11232 timeout_ms = -1;
345   else 360   else
HITCBC 346   440 timeout_ms = static_cast<int>((timeout_us + 999) / 1000); 361   348 timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
347   362  
HITCBC 348   42354 if (lock.owns_lock()) 363   41807 if (lock.owns_lock())
HITCBC 349   11185 lock.unlock(); 364   11582 lock.unlock();
350   365  
HITCBC 351   42354 task_cleanup on_exit{this, &lock, ctx}; 366   41807 task_cleanup on_exit{this, &lock, ctx};
352   367  
353   // Flush deferred timerfd programming before blocking 368   // Flush deferred timerfd programming before blocking
HITCBC 354   42354 if (timerfd_stale_.exchange(false, std::memory_order_acquire)) 369   41807 if (timerfd_stale_.exchange(false, std::memory_order_acquire))
HITCBC 355   3697 update_timerfd(); 370   3696 update_timerfd();
356   371  
HITCBC 357   42353 int nfds = ::epoll_wait( 372   41806 int nfds = ::epoll_wait(
HITCBC 358   42353 epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()), 373   41806 epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()),
359   timeout_ms); 374   timeout_ms);
360   375  
HITCBC 361   42353 if (nfds < 0 && errno != EINTR) 376   41806 if (nfds < 0 && errno != EINTR)
HITCBC 362   1 detail::throw_system_error(make_err(errno), "epoll_wait"); 377   1 detail::throw_system_error(make_err(errno), "epoll_wait");
363   378  
HITCBC 364   42352 bool check_timers = false; 379   41805 bool check_timers = false;
HITCBC 365   42352 ready_queue local_ops; 380   41805 ready_queue local_ops;
366   381  
HITCBC 367   92030 for (int i = 0; i < nfds; ++i) 382   89644 for (int i = 0; i < nfds; ++i)
368   { 383   {
HITCBC 369   49678 if (event_buffer_[i].data.ptr == nullptr) 384   47839 if (event_buffer_[i].data.ptr == nullptr)
370   { 385   {
371   std::uint64_t val; 386   std::uint64_t val;
372   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection) 387   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
HITCBC 373   4916 [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val)); 388   5071 [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
HITCBC 374   4916 eventfd_armed_.store(false, std::memory_order_relaxed); 389   5071 eventfd_armed_.store(false, std::memory_order_relaxed);
HITCBC 375   4916 continue; 390   5071 continue;
HITCBC 376   4916 } 391   5071 }
377   392  
HITCBC 378   44762 if (event_buffer_[i].data.ptr == &timer_fd_) 393   42768 if (event_buffer_[i].data.ptr == &timer_fd_)
379   { 394   {
380   std::uint64_t expirations; 395   std::uint64_t expirations;
381   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection) 396   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
382   [[maybe_unused]] auto r = 397   [[maybe_unused]] auto r =
HITCBC 383   7384 ::read(timer_fd_, &expirations, sizeof(expirations)); 398   7084 ::read(timer_fd_, &expirations, sizeof(expirations));
HITCBC 384   7384 check_timers = true; 399   7084 check_timers = true;
HITCBC 385   7384 continue; 400   7084 continue;
HITCBC 386   7384 } 401   7084 }
387   402  
388   auto* desc = 403   auto* desc =
HITCBC 389   37378 static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr); 404   35684 static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
ECB 390 - 37378 desc->add_ready_events(event_buffer_[i].events); 405 +
  406 + // EPOLLHUP maps to no reactor_event_* bit, so a HUP the widening
  407 + // below did not translate would take no branch in
  408 + // invoke_deferred_io() and, the registration being
  409 + // edge-triggered, never get another chance.
  410 + //
  411 + // HUP without OUT means a non-socket: a pipe or tty whose peer
  412 + // closed, with or without data still unread (HUP alone, or
  413 + // IN|HUP). Treat it as readable, writable and faulted, the way
  414 + // poll() and io_uring report it, so a parked error wait
  415 + // completes too. Sockets carry at least OUT with HUP (a fresh
  416 + // unconnected stream socket reports HUP|OUT), so they take the
  417 + // plain widening, where forcing IN costs at most one spurious
  418 + // EAGAIN.
HITGNC   419 + 35684 std::uint32_t ev = event_buffer_[i].events;
HITGNC   420 + 35684 if ((ev & (EPOLLHUP | EPOLLOUT)) == EPOLLHUP)
HITGNC   421 + 8 ev |= EPOLLIN | EPOLLOUT | EPOLLERR;
HITGNC   422 + 35676 else if (ev & EPOLLHUP)
HITGNC   423 + 235 ev |= EPOLLIN | EPOLLOUT;
HITGNC   424 + 35684 desc->add_ready_events(ev);
391   425  
HITCBC 392   37378 bool expected = false; 426   35684 bool expected = false;
HITCBC 393   37378 if (desc->is_enqueued_.compare_exchange_strong( 427   35684 if (desc->is_enqueued_.compare_exchange_strong(
394   expected, true, std::memory_order_release, 428   expected, true, std::memory_order_release,
395   std::memory_order_relaxed)) 429   std::memory_order_relaxed))
396   { 430   {
HITCBC 397   37378 local_ops.push(desc); 431   35684 local_ops.push(desc);
398   } 432   }
399   } 433   }
400   434  
HITCBC 401   42352 if (check_timers) 435   41805 if (check_timers)
402   { 436   {
HITCBC 403   7384 timer_svc_->process_expired(); 437   7084 timer_svc_->process_expired();
HITCBC 404   7384 update_timerfd(); 438   7084 update_timerfd();
405   } 439   }
406   440  
HITCBC 407   42352 lock.lock(); 441   41805 lock.lock();
408   442  
HITCBC 409   42352 completed_ops_.splice(local_ops); 443   41805 completed_ops_.splice(local_ops);
HITCBC 410   42354 } 444   41807 }
411   445  
412   } // namespace boost::corosio::detail 446   } // namespace boost::corosio::detail
413   447  
414   #endif // BOOST_COROSIO_HAS_EPOLL 448   #endif // BOOST_COROSIO_HAS_EPOLL
415   449  
416   #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP 450   #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP