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

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