99.34% Lines (151/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   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   /** Deregister a persistently registered descriptor. 119   /** Deregister a persistently registered descriptor.
120   120  
121   @param fd The file descriptor to deregister. 121   @param fd The file descriptor to deregister.
122   */ 122   */
123   void deregister_descriptor(int fd) const; 123   void deregister_descriptor(int fd) const;
124   124  
125   /// 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 126   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
127   { 127   {
HITCBC 128   76 return register_descriptor(read_fd, signal_pipe_reader_.arm()); 128   76 return register_descriptor(read_fd, signal_pipe_reader_.arm());
129   } 129   }
130   130  
131   private: 131   private:
132   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;
133   void interrupt_reactor() const override; 133   void interrupt_reactor() const override;
134   void update_timerfd() const; 134   void update_timerfd() const;
135   135  
136   int epoll_fd_; 136   int epoll_fd_;
137   int event_fd_; 137   int event_fd_;
138   int timer_fd_; 138   int timer_fd_;
139   139  
140   // 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
141   // register_signal_reader on the first signal registration). 141   // register_signal_reader on the first signal registration).
142   reactor_signal_pipe_reader signal_pipe_reader_; 142   reactor_signal_pipe_reader signal_pipe_reader_;
143   143  
144   // Edge-triggered eventfd state 144   // Edge-triggered eventfd state
145   mutable std::atomic<bool> eventfd_armed_{false}; 145   mutable std::atomic<bool> eventfd_armed_{false};
146   146  
147   // Set when the earliest timer changes; flushed before epoll_wait 147   // Set when the earliest timer changes; flushed before epoll_wait
148   mutable std::atomic<bool> timerfd_stale_{false}; 148   mutable std::atomic<bool> timerfd_stale_{false};
149   149  
150   // Event buffer sized from max_events_per_poll_ (set at construction, 150   // Event buffer sized from max_events_per_poll_ (set at construction,
151   // resized by configure_reactor via io_context_options). 151   // resized by configure_reactor via io_context_options).
152   std::vector<epoll_event> event_buffer_; 152   std::vector<epoll_event> event_buffer_;
153   }; 153   };
154   154  
HITCBC 155   1305 inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int) 155   1339 inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
HITCBC 156   1305 : epoll_fd_(-1) 156   1339 : epoll_fd_(-1)
HITCBC 157   1305 , event_fd_(-1) 157   1339 , event_fd_(-1)
HITCBC 158   1305 , timer_fd_(-1) 158   1339 , timer_fd_(-1)
HITCBC 159   2610 , event_buffer_(max_events_per_poll_) 159   2678 , event_buffer_(max_events_per_poll_)
160   { 160   {
HITCBC 161   1305 epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC); 161   1339 epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
HITCBC 162   1305 if (epoll_fd_ < 0) 162   1339 if (epoll_fd_ < 0)
HITCBC 163   1 detail::throw_system_error(make_err(errno), "epoll_create1"); 163   1 detail::throw_system_error(make_err(errno), "epoll_create1");
164   164  
HITCBC 165   1304 event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC); 165   1338 event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
HITCBC 166   1304 if (event_fd_ < 0) 166   1338 if (event_fd_ < 0)
167   { 167   {
HITCBC 168   1 int errn = errno; 168   1 int errn = errno;
HITCBC 169   1 ::close(epoll_fd_); 169   1 ::close(epoll_fd_);
HITCBC 170   1 detail::throw_system_error(make_err(errn), "eventfd"); 170   1 detail::throw_system_error(make_err(errn), "eventfd");
171   } 171   }
172   172  
HITCBC 173   1303 timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC); 173   1337 timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
HITCBC 174   1303 if (timer_fd_ < 0) 174   1337 if (timer_fd_ < 0)
175   { 175   {
HITCBC 176   1 int errn = errno; 176   1 int errn = errno;
HITCBC 177   1 ::close(event_fd_); 177   1 ::close(event_fd_);
HITCBC 178   1 ::close(epoll_fd_); 178   1 ::close(epoll_fd_);
HITCBC 179   1 detail::throw_system_error(make_err(errn), "timerfd_create"); 179   1 detail::throw_system_error(make_err(errn), "timerfd_create");
180   } 180   }
181   181  
HITCBC 182   1302 epoll_event ev{}; 182   1336 epoll_event ev{};
HITCBC 183   1302 ev.events = EPOLLIN | EPOLLET; 183   1336 ev.events = EPOLLIN | EPOLLET;
HITCBC 184   1302 ev.data.ptr = nullptr; 184   1336 ev.data.ptr = nullptr;
HITCBC 185   1302 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0) 185   1336 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0)
186   { 186   {
HITCBC 187   1 int errn = errno; 187   1 int errn = errno;
HITCBC 188   1 ::close(timer_fd_); 188   1 ::close(timer_fd_);
HITCBC 189   1 ::close(event_fd_); 189   1 ::close(event_fd_);
HITCBC 190   1 ::close(epoll_fd_); 190   1 ::close(epoll_fd_);
HITCBC 191   1 detail::throw_system_error(make_err(errn), "epoll_ctl"); 191   1 detail::throw_system_error(make_err(errn), "epoll_ctl");
192   } 192   }
193   193  
HITCBC 194   1301 epoll_event timer_ev{}; 194   1335 epoll_event timer_ev{};
HITCBC 195   1301 timer_ev.events = EPOLLIN | EPOLLERR; 195   1335 timer_ev.events = EPOLLIN | EPOLLERR;
HITCBC 196   1301 timer_ev.data.ptr = &timer_fd_; 196   1335 timer_ev.data.ptr = &timer_fd_;
HITCBC 197   1301 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0) 197   1335 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0)
198   { 198   {
HITCBC 199   1 int errn = errno; 199   1 int errn = errno;
HITCBC 200   1 ::close(timer_fd_); 200   1 ::close(timer_fd_);
HITCBC 201   1 ::close(event_fd_); 201   1 ::close(event_fd_);
HITCBC 202   1 ::close(epoll_fd_); 202   1 ::close(epoll_fd_);
HITCBC 203   1 detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)"); 203   1 detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)");
204   } 204   }
205   205  
HITCBC 206   1300 timer_svc_ = &get_timer_service(ctx, *this); 206   1334 timer_svc_ = &get_timer_service(ctx, *this);
HITCBC 207   1300 timer_svc_->set_on_earliest_changed( 207   1334 timer_svc_->set_on_earliest_changed(
HITCBC 208   5667 timer_service::callback(this, [](void* p) { 208   5516 timer_service::callback(this, [](void* p) {
HITCBC 209   4367 auto* self = static_cast<epoll_scheduler*>(p); 209   4182 auto* self = static_cast<epoll_scheduler*>(p);
HITCBC 210   4367 self->timerfd_stale_.store(true, std::memory_order_release); 210   4182 self->timerfd_stale_.store(true, std::memory_order_release);
HITCBC 211   4367 self->interrupt_reactor(); 211   4182 self->interrupt_reactor();
HITCBC 212   4367 })); 212   4182 }));
213   213  
HITCBC 214   1300 completed_ops_.push(&task_op_); 214   1334 completed_ops_.push(&task_op_);
HITCBC 215   1315 } 215   1349 }
216   216  
HITCBC 217   2600 inline epoll_scheduler::~epoll_scheduler() 217   2668 inline epoll_scheduler::~epoll_scheduler()
218   { 218   {
HITCBC 219   1300 if (timer_fd_ >= 0) 219   1334 if (timer_fd_ >= 0)
HITCBC 220   1300 ::close(timer_fd_); 220   1334 ::close(timer_fd_);
HITCBC 221   1300 if (event_fd_ >= 0) 221   1334 if (event_fd_ >= 0)
HITCBC 222   1300 ::close(event_fd_); 222   1334 ::close(event_fd_);
HITCBC 223   1300 if (epoll_fd_ >= 0) 223   1334 if (epoll_fd_ >= 0)
HITCBC 224   1300 ::close(epoll_fd_); 224   1334 ::close(epoll_fd_);
HITCBC 225   2600 } 225   2668 }
226   226  
227   inline void 227   inline void
HITCBC 228   1300 epoll_scheduler::shutdown() 228   1334 epoll_scheduler::shutdown()
229   { 229   {
HITCBC 230   1300 shutdown_drain(); 230   1334 shutdown_drain();
231   231  
HITCBC 232   1300 if (event_fd_ >= 0) 232   1334 if (event_fd_ >= 0)
HITCBC 233   1300 interrupt_reactor(); 233   1334 interrupt_reactor();
HITCBC 234   1300 } 234   1334 }
235   235  
236   inline void 236   inline void
HITCBC 237   27 epoll_scheduler::configure_reactor( 237   28 epoll_scheduler::configure_reactor(
238   unsigned max_events, 238   unsigned max_events,
239   unsigned budget_init, 239   unsigned budget_init,
240   unsigned budget_max, 240   unsigned budget_max,
241   unsigned unassisted) 241   unsigned unassisted)
242   { 242   {
HITCBC 243   27 reactor_scheduler::configure_reactor( 243   28 reactor_scheduler::configure_reactor(
244   max_events, budget_init, budget_max, unassisted); 244   max_events, budget_init, budget_max, unassisted);
HITCBC 245   25 event_buffer_.resize(max_events_per_poll_); 245   26 event_buffer_.resize(max_events_per_poll_);
HITCBC 246   25 } 246   26 }
247   247  
248   inline std::error_code 248   inline std::error_code
HITCBC 249   5780 epoll_scheduler::register_descriptor( 249   5792 epoll_scheduler::register_descriptor(
250   int fd, reactor_descriptor_state* desc) const 250   int fd, reactor_descriptor_state* desc) const
251   { 251   {
HITCBC 252   5780 epoll_event ev{}; 252   5792 epoll_event ev{};
HITCBC 253   5780 ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP; 253   5792 ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
HITCBC 254   5780 ev.data.ptr = desc; 254   5792 ev.data.ptr = desc;
255   255  
HITCBC 256   5780 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0) 256   5792 if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
HITCBC 257   7 return make_err(errno); 257   7 return make_err(errno);
258   258  
HITCBC 259   5773 desc->registered_events = ev.events; 259   5785 desc->registered_events = ev.events;
HITCBC 260   5773 desc->fd = fd; 260   5785 desc->fd = fd;
HITCBC 261   5773 desc->scheduler_ = this; 261   5785 desc->scheduler_ = this;
HITCBC 262   5773 desc->mutex.set_enabled(reactor_io_locking_); 262   5785 desc->mutex.set_enabled(reactor_io_locking_);
HITCBC 263   5773 desc->ready_events_.store(0, std::memory_order_relaxed); 263   5785 desc->ready_events_.store(0, std::memory_order_relaxed);
264   264  
HITCBC 265   5773 conditionally_enabled_mutex::scoped_lock lock(desc->mutex); 265   5785 conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
HITCBC 266   5773 desc->impl_ref_.reset(); 266   5785 desc->impl_ref_.reset();
HITCBC 267   5773 desc->read_ready = false; 267   5785 desc->read_ready = false;
HITCBC 268   5773 desc->write_ready = false; 268   5785 desc->write_ready = false;
HITCBC 269   5773 return {}; 269   5785 return {};
HITCBC 270   5773 } 270   5785 }
271   271  
272   inline void 272   inline void
HITCBC 273   5698 epoll_scheduler::deregister_descriptor(int fd) const 273   5710 epoll_scheduler::deregister_descriptor(int fd) const
274   { 274   {
HITCBC 275   5698 ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr); 275   5710 ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
HITCBC 276   5698 } 276   5710 }
277   277  
278   inline void 278   inline void
HITCBC 279   7812 epoll_scheduler::interrupt_reactor() const 279   7551 epoll_scheduler::interrupt_reactor() const
280   { 280   {
HITCBC 281   7812 bool expected = false; 281   7551 bool expected = false;
HITCBC 282   7812 if (eventfd_armed_.compare_exchange_strong( 282   7551 if (eventfd_armed_.compare_exchange_strong(
283   expected, true, std::memory_order_release, 283   expected, true, std::memory_order_release,
284   std::memory_order_relaxed)) 284   std::memory_order_relaxed))
285   { 285   {
HITCBC 286   6223 std::uint64_t val = 1; 286   6134 std::uint64_t val = 1;
HITCBC 287   6223 if (::write(event_fd_, &val, sizeof(val)) < 0) 287   6134 if (::write(event_fd_, &val, sizeof(val)) < 0)
288   { 288   {
289   // The flag is what coalesces later interrupts into a byte 289   // The flag is what coalesces later interrupts into a byte
290   // already in the eventfd; a write that failed put no byte 290   // already in the eventfd; a write that failed put no byte
291   // there, so leaving it armed would swallow every interrupt 291   // there, so leaving it armed would swallow every interrupt
292   // that follows. Disarming keeps the cost to the interrupts 292   // that follows. Disarming keeps the cost to the interrupts
293   // already in flight -- the next one arms and writes again, 293   // already in flight -- the next one arms and writes again,
294   // instead of every one after this coalescing into a byte 294   // instead of every one after this coalescing into a byte
295   // that does not exist. 295   // that does not exist.
HITCBC 296   2 eventfd_armed_.store(false, std::memory_order_release); 296   2 eventfd_armed_.store(false, std::memory_order_release);
297   } 297   }
298   } 298   }
HITCBC 299   7812 } 299   7551 }
300   300  
301   inline void 301   inline void
HITCBC 302   10904 epoll_scheduler::update_timerfd() const 302   10885 epoll_scheduler::update_timerfd() const
303   { 303   {
HITCBC 304   10904 auto nearest = timer_svc_->nearest_expiry(); 304   10885 auto nearest = timer_svc_->nearest_expiry();
305   305  
HITCBC 306   10904 itimerspec ts{}; 306   10885 itimerspec ts{};
HITCBC 307   10904 int flags = 0; 307   10885 int flags = 0;
308   308  
HITCBC 309   10904 if (nearest == timer_service::time_point::max()) 309   10885 if (nearest == timer_service::time_point::max())
310   { 310   {
311   // No timers — disarm by setting to 0 (relative) 311   // No timers — disarm by setting to 0 (relative)
312   } 312   }
313   else 313   else
314   { 314   {
HITCBC 315   9726 auto now = std::chrono::steady_clock::now(); 315   9697 auto now = std::chrono::steady_clock::now();
HITCBC 316   9726 if (nearest <= now) 316   9697 if (nearest <= now)
317   { 317   {
318   // Use 1ns instead of 0 — zero disarms the timerfd 318   // Use 1ns instead of 0 — zero disarms the timerfd
HITCBC 319   1246 ts.it_value.tv_nsec = 1; 319   1258 ts.it_value.tv_nsec = 1;
320   } 320   }
321   else 321   else
322   { 322   {
HITCBC 323   8480 auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>( 323   8439 auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
HITCBC 324   8480 nearest - now) 324   8439 nearest - now)
HITCBC 325   8480 .count(); 325   8439 .count();
HITCBC 326   8480 ts.it_value.tv_sec = nsec / 1000000000; 326   8439 ts.it_value.tv_sec = nsec / 1000000000;
HITCBC 327   8480 ts.it_value.tv_nsec = nsec % 1000000000; 327   8439 ts.it_value.tv_nsec = nsec % 1000000000;
HITCBC 328   8480 if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0) 328   8439 if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
MISUBC 329   ✗ ts.it_value.tv_nsec = 1; 329   ✗ ts.it_value.tv_nsec = 1;
330   } 330   }
331   } 331   }
332   332  
HITCBC 333   10904 if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0) 333   10885 if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
HITCBC 334   1 detail::throw_system_error(make_err(errno), "timerfd_settime"); 334   1 detail::throw_system_error(make_err(errno), "timerfd_settime");
HITCBC 335   10903 } 335   10884 }
336   336  
337   inline void 337   inline void
HITCBC 338   41348 epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us) 338   38343 epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
339   { 339   {
340   int timeout_ms; 340   int timeout_ms;
HITCBC 341   41348 if (task_interrupted_) 341   38343 if (task_interrupted_)
HITCBC 342   29917 timeout_ms = 0; 342   26939 timeout_ms = 0;
HITCBC 343   11431 else if (timeout_us < 0) 343   11404 else if (timeout_us < 0)
HITCBC 344   11086 timeout_ms = -1; 344   11126 timeout_ms = -1;
345   else 345   else
HITCBC 346   345 timeout_ms = static_cast<int>((timeout_us + 999) / 1000); 346   278 timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
347   347  
HITCBC 348   41348 if (lock.owns_lock()) 348   38343 if (lock.owns_lock())
HITCBC 349   11433 lock.unlock(); 349   11406 lock.unlock();
350   350  
HITCBC 351   41348 task_cleanup on_exit{this, &lock, ctx}; 351   38343 task_cleanup on_exit{this, &lock, ctx};
352   352  
353   // Flush deferred timerfd programming before blocking 353   // Flush deferred timerfd programming before blocking
HITCBC 354   41348 if (timerfd_stale_.exchange(false, std::memory_order_acquire)) 354   38343 if (timerfd_stale_.exchange(false, std::memory_order_acquire))
HITCBC 355   3668 update_timerfd(); 355   3677 update_timerfd();
356   356  
HITCBC 357   41347 int nfds = ::epoll_wait( 357   38342 int nfds = ::epoll_wait(
HITCBC 358   41347 epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()), 358   38342 epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()),
359   timeout_ms); 359   timeout_ms);
360   360  
HITCBC 361   41347 if (nfds < 0 && errno != EINTR) 361   38342 if (nfds < 0 && errno != EINTR)
HITCBC 362   1 detail::throw_system_error(make_err(errno), "epoll_wait"); 362   1 detail::throw_system_error(make_err(errno), "epoll_wait");
363   363  
HITCBC 364   41346 bool check_timers = false; 364   38341 bool check_timers = false;
HITCBC 365   41346 ready_queue local_ops; 365   38341 ready_queue local_ops;
366   366  
HITCBC 367   88839 for (int i = 0; i < nfds; ++i) 367   84441 for (int i = 0; i < nfds; ++i)
368   { 368   {
HITCBC 369   47493 if (event_buffer_[i].data.ptr == nullptr) 369   46100 if (event_buffer_[i].data.ptr == nullptr)
370   { 370   {
371   std::uint64_t val; 371   std::uint64_t val;
372   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection) 372   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
HITCBC 373   4921 [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val)); 373   4798 [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
HITCBC 374   4921 eventfd_armed_.store(false, std::memory_order_relaxed); 374   4798 eventfd_armed_.store(false, std::memory_order_relaxed);
HITCBC 375   4921 continue; 375   4798 continue;
HITCBC 376   4921 } 376   4798 }
377   377  
HITCBC 378   42572 if (event_buffer_[i].data.ptr == &timer_fd_) 378   41302 if (event_buffer_[i].data.ptr == &timer_fd_)
379   { 379   {
380   std::uint64_t expirations; 380   std::uint64_t expirations;
381   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection) 381   // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
382   [[maybe_unused]] auto r = 382   [[maybe_unused]] auto r =
HITCBC 383   7236 ::read(timer_fd_, &expirations, sizeof(expirations)); 383   7208 ::read(timer_fd_, &expirations, sizeof(expirations));
HITCBC 384   7236 check_timers = true; 384   7208 check_timers = true;
HITCBC 385   7236 continue; 385   7208 continue;
HITCBC 386   7236 } 386   7208 }
387   387  
388   auto* desc = 388   auto* desc =
HITCBC 389   35336 static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr); 389   34094 static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
ECB 390 - 35336 desc->add_ready_events(event_buffer_[i].events); 390 +
  391 + // A pipe or tty whose peer closed reports EPOLLHUP on its own --
  392 + // no EPOLLIN, no EPOLLERR -- and EPOLLHUP maps to no
  393 + // reactor_event_* bit, so invoke_deferred_io() would take no
  394 + // branch and, the registration being edge-triggered, never get
  395 + // another chance.
  396 + //
  397 + // Sockets are unaffected because they never report EPOLLHUP
  398 + // alone. tcp_poll() and unix_poll() raise it only once
  399 + // sk_shutdown is SHUTDOWN_MASK (or the state is TCP_CLOSE), and
  400 + // both also report the socket readable and writable there --
  401 + // tcp_poll() takes an explicit `else mask |= EPOLLOUT` branch
  402 + // once SEND_SHUTDOWN is set, because a send on a shut-down
  403 + // socket fails fast rather than blocking. Measured on every
  404 + // state that produces EPOLLHUP -- peer close plus local
  405 + // SHUT_WR/SHUT_RDWR, RST, RST with the send buffer full, and
  406 + // the AF_UNIX equivalents -- the mask is always IN|OUT|HUP
  407 + // (0x15), or IN|OUT|ERR|HUP (0x1d) for a reset. The forced bits
  408 + // are therefore already set on every socket path.
HITGNC   409 + 34094 std::uint32_t ev = event_buffer_[i].events;
HITGNC   410 + 34094 if (ev & EPOLLHUP)
HITGNC   411 + 229 ev |= EPOLLIN | EPOLLOUT;
HITGNC   412 + 34094 desc->add_ready_events(ev);
391   413  
HITCBC 392   35336 bool expected = false; 414   34094 bool expected = false;
HITCBC 393   35336 if (desc->is_enqueued_.compare_exchange_strong( 415   34094 if (desc->is_enqueued_.compare_exchange_strong(
394   expected, true, std::memory_order_release, 416   expected, true, std::memory_order_release,
395   std::memory_order_relaxed)) 417   std::memory_order_relaxed))
396   { 418   {
HITCBC 397   35336 local_ops.push(desc); 419   34094 local_ops.push(desc);
398   } 420   }
399   } 421   }
400   422  
HITCBC 401   41346 if (check_timers) 423   38341 if (check_timers)
402   { 424   {
HITCBC 403   7236 timer_svc_->process_expired(); 425   7208 timer_svc_->process_expired();
HITCBC 404   7236 update_timerfd(); 426   7208 update_timerfd();
405   } 427   }
406   428  
HITCBC 407   41346 lock.lock(); 429   38341 lock.lock();
408   430  
HITCBC 409   41346 completed_ops_.splice(local_ops); 431   38341 completed_ops_.splice(local_ops);
HITCBC 410   41348 } 432   38343 }
411   433  
412   } // namespace boost::corosio::detail 434   } // namespace boost::corosio::detail
413   435  
414   #endif // BOOST_COROSIO_HAS_EPOLL 436   #endif // BOOST_COROSIO_HAS_EPOLL
415   437  
416   #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP 438   #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP