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 4356x ensure_write_registered(int, reactor_descriptor_state*) const noexcept
122 {
123 4356x 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 1952x inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
163 1952x : epoll_fd_(-1)
164 1952x , event_fd_(-1)
165 1952x , timer_fd_(-1)
166 3904x , event_buffer_(max_events_per_poll_)
167 {
168 1952x epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
169 1952x if (epoll_fd_ < 0)
170 1x detail::throw_system_error(make_err(errno), "epoll_create1");
171
172 1951x event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
173 1951x 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 1950x timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
181 1950x 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 1949x epoll_event ev{};
190 1949x ev.events = EPOLLIN | EPOLLET;
191 1949x ev.data.ptr = nullptr;
192 1949x 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 1948x epoll_event timer_ev{};
202 1948x timer_ev.events = EPOLLIN | EPOLLERR;
203 1948x timer_ev.data.ptr = &timer_fd_;
204 1948x 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 1947x timer_svc_ = &get_timer_service(ctx, *this);
214 1947x timer_svc_->set_on_earliest_changed(
215 8143x timer_service::callback(this, [](void* p) {
216 6196x auto* self = static_cast<epoll_scheduler*>(p);
217 6196x self->timerfd_stale_.store(true, std::memory_order_release);
218 6196x self->interrupt_reactor();
219 6196x }));
220
221 1947x completed_ops_.push(&task_op_);
222 1962x }
223
224 3894x inline epoll_scheduler::~epoll_scheduler()
225 {
226 1947x if (timer_fd_ >= 0)
227 1947x ::close(timer_fd_);
228 1947x if (event_fd_ >= 0)
229 1947x ::close(event_fd_);
230 1947x if (epoll_fd_ >= 0)
231 1947x ::close(epoll_fd_);
232 3894x }
233
234 inline void
235 1947x epoll_scheduler::shutdown()
236 {
237 1947x shutdown_drain();
238
239 1947x if (event_fd_ >= 0)
240 1947x interrupt_reactor();
241 1947x }
242
243 inline void
244 30x epoll_scheduler::configure_reactor(
245 unsigned max_events,
246 unsigned budget_init,
247 unsigned budget_max,
248 unsigned unassisted)
249 {
250 30x reactor_scheduler::configure_reactor(
251 max_events, budget_init, budget_max, unassisted);
252 28x event_buffer_.resize(max_events_per_poll_);
253 28x }
254
255 inline std::error_code
256 10683x epoll_scheduler::register_descriptor(
257 int fd, reactor_descriptor_state* desc) const
258 {
259 10683x epoll_event ev{};
260 10683x ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
261 10683x ev.data.ptr = desc;
262
263 10683x bool unpollable = false;
264 10683x 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 10676x desc->registered_events = unpollable ? 0 : ev.events;
274 10676x desc->unpollable = unpollable;
275 10676x desc->fd = fd;
276 10676x desc->scheduler_ = this;
277 10676x desc->mutex.set_enabled(reactor_io_locking_);
278 10676x desc->ready_events_.store(0, std::memory_order_relaxed);
279
280 10676x conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
281 10676x desc->impl_ref_.reset();
282 10676x desc->read_ready = false;
283 10676x desc->write_ready = false;
284 10676x return {};
285 10676x }
286
287 inline void
288 10597x epoll_scheduler::deregister_descriptor(int fd) const
289 {
290 10597x ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
291 10597x }
292
293 inline void
294 13656x epoll_scheduler::interrupt_reactor() const
295 {
296 13656x bool expected = false;
297 13656x if (eventfd_armed_.compare_exchange_strong(
298 expected, true, std::memory_order_release,
299 std::memory_order_relaxed))
300 {
301 11406x std::uint64_t val = 1;
302 11406x 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 13656x }
315
316 inline void
317 11267x epoll_scheduler::update_timerfd() const
318 {
319 11267x auto nearest = timer_svc_->nearest_expiry();
320
321 11267x itimerspec ts{};
322 11267x int flags = 0;
323
324 11267x if (nearest == timer_service::time_point::max())
325 {
326 // No timers — disarm by setting to 0 (relative)
327 }
328 else
329 {
330 10051x auto now = std::chrono::steady_clock::now();
331 10051x if (nearest <= now)
332 {
333 // Use 1ns instead of 0 — zero disarms the timerfd
334 283x ts.it_value.tv_nsec = 1;
335 }
336 else
337 {
338 9768x auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
339 9768x nearest - now)
340 9768x .count();
341 9768x ts.it_value.tv_sec = nsec / 1000000000;
342 9768x ts.it_value.tv_nsec = nsec % 1000000000;
343 9768x if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
344 ✗ ts.it_value.tv_nsec = 1;
345 }
346 }
347
348 11267x if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
349 1x detail::throw_system_error(make_err(errno), "timerfd_settime");
350 11266x }
351
352 inline void
353 54656x epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
354 {
355 int timeout_ms;
356 54656x if (task_interrupted_)
357 37896x timeout_ms = 0;
358 16760x else if (timeout_us < 0)
359 16082x timeout_ms = -1;
360 else
361 678x timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
362
363 54656x if (lock.owns_lock())
364 16762x lock.unlock();
365
366 54656x task_cleanup on_exit{this, &lock, ctx};
367
368 // Flush deferred timerfd programming before blocking
369 54656x if (timerfd_stale_.exchange(false, std::memory_order_acquire))
370 5442x update_timerfd();
371
372 54655x int nfds = ::epoll_wait(
373 54655x epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()),
374 timeout_ms);
375
376 54655x if (nfds < 0 && errno != EINTR)
377 1x detail::throw_system_error(make_err(errno), "epoll_wait");
378
379 54654x bool check_timers = false;
380 54654x ready_queue local_ops;
381
382 117204x for (int i = 0; i < nfds; ++i)
383 {
384 62550x if (event_buffer_[i].data.ptr == nullptr)
385 {
386 std::uint64_t val;
387 // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
388 9457x [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
389 9457x eventfd_armed_.store(false, std::memory_order_relaxed);
390 9457x continue;
391 9457x }
392
393 53093x 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 5825x ::read(timer_fd_, &expirations, sizeof(expirations));
399 5825x check_timers = true;
400 5825x continue;
401 5825x }
402
403 auto* desc =
404 47268x 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 47268x std::uint32_t ev = event_buffer_[i].events;
420 47268x if ((ev & (EPOLLHUP | EPOLLOUT)) == EPOLLHUP)
421 8x ev |= EPOLLIN | EPOLLOUT | EPOLLERR;
422 47260x else if (ev & EPOLLHUP)
423 1522x ev |= EPOLLIN | EPOLLOUT;
424 47268x desc->add_ready_events(ev);
425
426 47268x bool expected = false;
427 47268x if (desc->is_enqueued_.compare_exchange_strong(
428 expected, true, std::memory_order_release,
429 std::memory_order_relaxed))
430 {
431 47268x local_ops.push(desc);
432 }
433 }
434
435 54654x if (check_timers)
436 {
437 5825x timer_svc_->process_expired();
438 5825x update_timerfd();
439 }
440
441 54654x lock.lock();
442
443 54654x completed_ops_.splice(local_ops);
444 54656x }
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