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

99.4% Lines (159 / 160) 100.0% Functions (12 / 12)
epoll_scheduler.hpp
f(x) Functions (12)
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 /// No-op: write readiness is watched from registration on.
120 std::error_code
121 2618x ensure_write_registered(int, reactor_descriptor_state*) const noexcept
122 {
123 2618x return {};
124 }
125
126 /** Deregister a persistently registered descriptor.
127
128 @param fd The file descriptor to deregister.
129 */
130 void deregister_descriptor(int fd) const;
131
132 /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp).
133 76x [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
134 {
135 76x return register_descriptor(read_fd, signal_pipe_reader_.arm());
136 }
137
138 private:
139 void run_task(lock_type& lock, context_type& ctx, long timeout_us) override;
140 void interrupt_reactor() const override;
141 void update_timerfd() const;
142
143 int epoll_fd_;
144 int event_fd_;
145 int timer_fd_;
146
147 // Watches the global signal self-pipe's read end (armed lazily by
148 // register_signal_reader on the first signal registration).
149 reactor_signal_pipe_reader signal_pipe_reader_;
150
151 // Edge-triggered eventfd state
152 mutable std::atomic<bool> eventfd_armed_{false};
153
154 // Set when the earliest timer changes; flushed before epoll_wait
155 mutable std::atomic<bool> timerfd_stale_{false};
156
157 // Event buffer sized from max_events_per_poll_ (set at construction,
158 // resized by configure_reactor via io_context_options).
159 std::vector<epoll_event> event_buffer_;
160 };
161
162 1577x inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
163 1577x : epoll_fd_(-1)
164 1577x , event_fd_(-1)
165 1577x , timer_fd_(-1)
166 3154x , event_buffer_(max_events_per_poll_)
167 {
168 1577x epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
169 1577x if (epoll_fd_ < 0)
170 1x detail::throw_system_error(make_err(errno), "epoll_create1");
171
172 1576x event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
173 1576x if (event_fd_ < 0)
174 {
175 1x int errn = errno;
176 1x ::close(epoll_fd_);
177 1x detail::throw_system_error(make_err(errn), "eventfd");
178 }
179
180 1575x timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
181 1575x if (timer_fd_ < 0)
182 {
183 1x int errn = errno;
184 1x ::close(event_fd_);
185 1x ::close(epoll_fd_);
186 1x detail::throw_system_error(make_err(errn), "timerfd_create");
187 }
188
189 1574x epoll_event ev{};
190 1574x ev.events = EPOLLIN | EPOLLET;
191 1574x ev.data.ptr = nullptr;
192 1574x if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0)
193 {
194 1x int errn = errno;
195 1x ::close(timer_fd_);
196 1x ::close(event_fd_);
197 1x ::close(epoll_fd_);
198 1x detail::throw_system_error(make_err(errn), "epoll_ctl");
199 }
200
201 1573x epoll_event timer_ev{};
202 1573x timer_ev.events = EPOLLIN | EPOLLERR;
203 1573x timer_ev.data.ptr = &timer_fd_;
204 1573x if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0)
205 {
206 1x int errn = errno;
207 1x ::close(timer_fd_);
208 1x ::close(event_fd_);
209 1x ::close(epoll_fd_);
210 1x detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)");
211 }
212
213 1572x timer_svc_ = &get_timer_service(ctx, *this);
214 1572x timer_svc_->set_on_earliest_changed(
215 5979x timer_service::callback(this, [](void* p) {
216 4407x auto* self = static_cast<epoll_scheduler*>(p);
217 4407x self->timerfd_stale_.store(true, std::memory_order_release);
218 4407x self->interrupt_reactor();
219 4407x }));
220
221 1572x completed_ops_.push(&task_op_);
222 1587x }
223
224 3144x inline epoll_scheduler::~epoll_scheduler()
225 {
226 1572x if (timer_fd_ >= 0)
227 1572x ::close(timer_fd_);
228 1572x if (event_fd_ >= 0)
229 1572x ::close(event_fd_);
230 1572x if (epoll_fd_ >= 0)
231 1572x ::close(epoll_fd_);
232 3144x }
233
234 inline void
235 1572x epoll_scheduler::shutdown()
236 {
237 1572x shutdown_drain();
238
239 1572x if (event_fd_ >= 0)
240 1572x interrupt_reactor();
241 1572x }
242
243 inline void
244 28x epoll_scheduler::configure_reactor(
245 unsigned max_events,
246 unsigned budget_init,
247 unsigned budget_max,
248 unsigned unassisted)
249 {
250 28x reactor_scheduler::configure_reactor(
251 max_events, budget_init, budget_max, unassisted);
252 26x event_buffer_.resize(max_events_per_poll_);
253 26x }
254
255 inline std::error_code
256 5923x epoll_scheduler::register_descriptor(
257 int fd, reactor_descriptor_state* desc) const
258 {
259 5923x epoll_event ev{};
260 5923x ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
261 5923x ev.data.ptr = desc;
262
263 5923x bool unpollable = false;
264 5923x if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
265 {
266 // EPERM: a file type epoll cannot watch. Its I/O does not
267 // block, so adopt it unwatched, as asio does.
268 11x if (errno != EPERM)
269 7x return make_err(errno);
270 4x unpollable = true;
271 }
272
273 5916x desc->registered_events = unpollable ? 0 : ev.events;
274 5916x desc->unpollable = unpollable;
275 5916x desc->fd = fd;
276 5916x desc->scheduler_ = this;
277 5916x desc->mutex.set_enabled(reactor_io_locking_);
278 5916x desc->ready_events_.store(0, std::memory_order_relaxed);
279
280 5916x conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
281 5916x desc->impl_ref_.reset();
282 5916x desc->read_ready = false;
283 5916x desc->write_ready = false;
284 5916x return {};
285 5916x }
286
287 inline void
288 5837x epoll_scheduler::deregister_descriptor(int fd) const
289 {
290 5837x ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
291 5837x }
292
293 inline void
294 8498x epoll_scheduler::interrupt_reactor() const
295 {
296 8498x bool expected = false;
297 8498x if (eventfd_armed_.compare_exchange_strong(
298 expected, true, std::memory_order_release,
299 std::memory_order_relaxed))
300 {
301 6645x std::uint64_t val = 1;
302 6645x if (::write(event_fd_, &val, sizeof(val)) < 0)
303 {
304 // The flag is what coalesces later interrupts into a byte
305 // already in the eventfd; a write that failed put no byte
306 // there, so leaving it armed would swallow every interrupt
307 // that follows. Disarming keeps the cost to the interrupts
308 // already in flight -- the next one arms and writes again,
309 // instead of every one after this coalescing into a byte
310 // that does not exist.
311 2x eventfd_armed_.store(false, std::memory_order_release);
312 }
313 }
314 8498x }
315
316 inline void
317 10780x epoll_scheduler::update_timerfd() const
318 {
319 10780x auto nearest = timer_svc_->nearest_expiry();
320
321 10780x itimerspec ts{};
322 10780x int flags = 0;
323
324 10780x if (nearest == timer_service::time_point::max())
325 {
326 // No timers — disarm by setting to 0 (relative)
327 }
328 else
329 {
330 9671x auto now = std::chrono::steady_clock::now();
331 9671x if (nearest <= now)
332 {
333 // Use 1ns instead of 0 — zero disarms the timerfd
334 1464x ts.it_value.tv_nsec = 1;
335 }
336 else
337 {
338 8207x auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
339 8207x nearest - now)
340 8207x .count();
341 8207x ts.it_value.tv_sec = nsec / 1000000000;
342 8207x ts.it_value.tv_nsec = nsec % 1000000000;
343 8207x if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
344 ✗ ts.it_value.tv_nsec = 1;
345 }
346 }
347
348 10780x if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
349 1x detail::throw_system_error(make_err(errno), "timerfd_settime");
350 10779x }
351
352 inline void
353 41807x epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
354 {
355 int timeout_ms;
356 41807x if (task_interrupted_)
357 30227x timeout_ms = 0;
358 11580x else if (timeout_us < 0)
359 11232x timeout_ms = -1;
360 else
361 348x timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
362
363 41807x if (lock.owns_lock())
364 11582x lock.unlock();
365
366 41807x task_cleanup on_exit{this, &lock, ctx};
367
368 // Flush deferred timerfd programming before blocking
369 41807x if (timerfd_stale_.exchange(false, std::memory_order_acquire))
370 3696x update_timerfd();
371
372 41806x int nfds = ::epoll_wait(
373 41806x epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()),
374 timeout_ms);
375
376 41806x if (nfds < 0 && errno != EINTR)
377 1x detail::throw_system_error(make_err(errno), "epoll_wait");
378
379 41805x bool check_timers = false;
380 41805x ready_queue local_ops;
381
382 89644x for (int i = 0; i < nfds; ++i)
383 {
384 47839x if (event_buffer_[i].data.ptr == nullptr)
385 {
386 std::uint64_t val;
387 // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
388 5071x [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
389 5071x eventfd_armed_.store(false, std::memory_order_relaxed);
390 5071x continue;
391 5071x }
392
393 42768x if (event_buffer_[i].data.ptr == &timer_fd_)
394 {
395 std::uint64_t expirations;
396 // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
397 [[maybe_unused]] auto r =
398 7084x ::read(timer_fd_, &expirations, sizeof(expirations));
399 7084x check_timers = true;
400 7084x continue;
401 7084x }
402
403 auto* desc =
404 35684x static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
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.
419 35684x std::uint32_t ev = event_buffer_[i].events;
420 35684x if ((ev & (EPOLLHUP | EPOLLOUT)) == EPOLLHUP)
421 8x ev |= EPOLLIN | EPOLLOUT | EPOLLERR;
422 35676x else if (ev & EPOLLHUP)
423 235x ev |= EPOLLIN | EPOLLOUT;
424 35684x desc->add_ready_events(ev);
425
426 35684x bool expected = false;
427 35684x if (desc->is_enqueued_.compare_exchange_strong(
428 expected, true, std::memory_order_release,
429 std::memory_order_relaxed))
430 {
431 35684x local_ops.push(desc);
432 }
433 }
434
435 41805x if (check_timers)
436 {
437 7084x timer_svc_->process_expired();
438 7084x update_timerfd();
439 }
440
441 41805x lock.lock();
442
443 41805x completed_ops_.splice(local_ops);
444 41807x }
445
446 } // namespace boost::corosio::detail
447
448 #endif // BOOST_COROSIO_HAS_EPOLL
449
450 #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP
451