include/boost/corosio/native/detail/epoll/epoll_scheduler.hpp

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