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

85.5% Lines (130/0/152) 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 void 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 51x void register_signal_reader(int read_fd) override
127 {
128 51x register_descriptor(read_fd, signal_pipe_reader_.arm());
129 51x }
130
131 private:
132 void
133 run_task(lock_type& lock, context_type* ctx,
134 long timeout_us) override;
135 void interrupt_reactor() const override;
136 void update_timerfd() const;
137
138 int epoll_fd_;
139 int event_fd_;
140 int timer_fd_;
141
142 // Watches the global signal self-pipe's read end (armed lazily by
143 // register_signal_reader on the first signal registration).
144 reactor_signal_pipe_reader signal_pipe_reader_;
145
146 // Edge-triggered eventfd state
147 mutable std::atomic<bool> eventfd_armed_{false};
148
149 // Set when the earliest timer changes; flushed before epoll_wait
150 mutable std::atomic<bool> timerfd_stale_{false};
151
152 // Event buffer sized from max_events_per_poll_ (set at construction,
153 // resized by configure_reactor via io_context_options).
154 std::vector<epoll_event> event_buffer_;
155 };
156
157 817x inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
158 817x : epoll_fd_(-1)
159 817x , event_fd_(-1)
160 817x , timer_fd_(-1)
161 1634x , event_buffer_(max_events_per_poll_)
162 {
163 817x epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
164 817x if (epoll_fd_ < 0)
165 detail::throw_system_error(make_err(errno), "epoll_create1");
166
167 817x event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
168 817x if (event_fd_ < 0)
169 {
170 int errn = errno;
171 ::close(epoll_fd_);
172 detail::throw_system_error(make_err(errn), "eventfd");
173 }
174
175 817x timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
176 817x if (timer_fd_ < 0)
177 {
178 int errn = errno;
179 ::close(event_fd_);
180 ::close(epoll_fd_);
181 detail::throw_system_error(make_err(errn), "timerfd_create");
182 }
183
184 817x epoll_event ev{};
185 817x ev.events = EPOLLIN | EPOLLET;
186 817x ev.data.ptr = nullptr;
187 817x if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0)
188 {
189 int errn = errno;
190 ::close(timer_fd_);
191 ::close(event_fd_);
192 ::close(epoll_fd_);
193 detail::throw_system_error(make_err(errn), "epoll_ctl");
194 }
195
196 817x epoll_event timer_ev{};
197 817x timer_ev.events = EPOLLIN | EPOLLERR;
198 817x timer_ev.data.ptr = &timer_fd_;
199 817x if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0)
200 {
201 int errn = errno;
202 ::close(timer_fd_);
203 ::close(event_fd_);
204 ::close(epoll_fd_);
205 detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)");
206 }
207
208 817x timer_svc_ = &get_timer_service(ctx, *this);
209 817x timer_svc_->set_on_earliest_changed(
210 6037x timer_service::callback(this, [](void* p) {
211 5220x auto* self = static_cast<epoll_scheduler*>(p);
212 5220x self->timerfd_stale_.store(true, std::memory_order_release);
213 5220x self->interrupt_reactor();
214 5220x }));
215
216 817x get_resolver_service(ctx, *this);
217 817x get_signal_service(ctx, *this);
218 817x get_stream_file_service(ctx, *this);
219 817x get_random_access_file_service(ctx, *this);
220
221 817x completed_ops_.push(&task_op_);
222 817x }
223
224 1634x inline epoll_scheduler::~epoll_scheduler()
225 {
226 817x if (timer_fd_ >= 0)
227 817x ::close(timer_fd_);
228 817x if (event_fd_ >= 0)
229 817x ::close(event_fd_);
230 817x if (epoll_fd_ >= 0)
231 817x ::close(epoll_fd_);
232 1634x }
233
234 inline void
235 817x epoll_scheduler::shutdown()
236 {
237 817x shutdown_drain();
238
239 817x if (event_fd_ >= 0)
240 817x interrupt_reactor();
241 817x }
242
243 inline void
244 19x epoll_scheduler::configure_reactor(
245 unsigned max_events,
246 unsigned budget_init,
247 unsigned budget_max,
248 unsigned unassisted)
249 {
250 19x reactor_scheduler::configure_reactor(
251 max_events, budget_init, budget_max, unassisted);
252 18x event_buffer_.resize(max_events_per_poll_);
253 18x }
254
255 inline void
256 9406x epoll_scheduler::register_descriptor(int fd, reactor_descriptor_state* desc) const
257 {
258 9406x epoll_event ev{};
259 9406x ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
260 9406x ev.data.ptr = desc;
261
262 9406x if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
263 detail::throw_system_error(make_err(errno), "epoll_ctl (register)");
264
265 9406x desc->registered_events = ev.events;
266 9406x desc->fd = fd;
267 9406x desc->scheduler_ = this;
268 9406x desc->mutex.set_enabled(reactor_io_locking_);
269 9406x desc->ready_events_.store(0, std::memory_order_relaxed);
270
271 9406x conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
272 9406x desc->impl_ref_.reset();
273 9406x desc->read_ready = false;
274 9406x desc->write_ready = false;
275 9406x }
276
277 inline void
278 9355x epoll_scheduler::deregister_descriptor(int fd) const
279 {
280 9355x ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
281 9355x }
282
283 inline void
284 6771x epoll_scheduler::interrupt_reactor() const
285 {
286 6771x bool expected = false;
287 6771x if (eventfd_armed_.compare_exchange_strong(
288 expected, true, std::memory_order_release,
289 std::memory_order_relaxed))
290 {
291 5661x std::uint64_t val = 1;
292 5661x [[maybe_unused]] auto r = ::write(event_fd_, &val, sizeof(val));
293 }
294 6771x }
295
296 inline void
297 9231x epoll_scheduler::update_timerfd() const
298 {
299 9231x auto nearest = timer_svc_->nearest_expiry();
300
301 9231x itimerspec ts{};
302 9231x int flags = 0;
303
304 9231x if (nearest == timer_service::time_point::max())
305 {
306 // No timers — disarm by setting to 0 (relative)
307 }
308 else
309 {
310 9060x auto now = std::chrono::steady_clock::now();
311 9060x if (nearest <= now)
312 {
313 // Use 1ns instead of 0 — zero disarms the timerfd
314 1034x ts.it_value.tv_nsec = 1;
315 }
316 else
317 {
318 8026x auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
319 8026x nearest - now)
320 8026x .count();
321 8026x ts.it_value.tv_sec = nsec / 1000000000;
322 8026x ts.it_value.tv_nsec = nsec % 1000000000;
323 8026x if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
324 ts.it_value.tv_nsec = 1;
325 }
326 }
327
328 9231x if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
329 detail::throw_system_error(make_err(errno), "timerfd_settime");
330 9231x }
331
332 inline void
333 40232x epoll_scheduler::run_task(
334 lock_type& lock, context_type* ctx, long timeout_us)
335 {
336 int timeout_ms;
337 40232x if (task_interrupted_)
338 27659x timeout_ms = 0;
339 12573x else if (timeout_us < 0)
340 12569x timeout_ms = -1;
341 else
342 4x timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
343
344 40232x if (lock.owns_lock())
345 12575x lock.unlock();
346
347 40232x task_cleanup on_exit{this, &lock, ctx};
348
349 // Flush deferred timerfd programming before blocking
350 40232x if (timerfd_stale_.exchange(false, std::memory_order_acquire))
351 4622x update_timerfd();
352
353 40232x int nfds = ::epoll_wait(
354 epoll_fd_, event_buffer_.data(),
355 40232x static_cast<int>(event_buffer_.size()), timeout_ms);
356
357 40232x if (nfds < 0 && errno != EINTR)
358 detail::throw_system_error(make_err(errno), "epoll_wait");
359
360 40232x bool check_timers = false;
361 40232x ready_queue local_ops;
362
363 90648x for (int i = 0; i < nfds; ++i)
364 {
365 50416x if (event_buffer_[i].data.ptr == nullptr)
366 {
367 std::uint64_t val;
368 // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
369 4844x [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
370 4844x eventfd_armed_.store(false, std::memory_order_relaxed);
371 4844x continue;
372 4844x }
373
374 45572x if (event_buffer_[i].data.ptr == &timer_fd_)
375 {
376 std::uint64_t expirations;
377 // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
378 [[maybe_unused]] auto r =
379 4609x ::read(timer_fd_, &expirations, sizeof(expirations));
380 4609x check_timers = true;
381 4609x continue;
382 4609x }
383
384 auto* desc =
385 40963x static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
386 40963x desc->add_ready_events(event_buffer_[i].events);
387
388 40963x bool expected = false;
389 40963x if (desc->is_enqueued_.compare_exchange_strong(
390 expected, true, std::memory_order_release,
391 std::memory_order_relaxed))
392 {
393 40963x local_ops.push(desc);
394 }
395 }
396
397 40232x if (check_timers)
398 {
399 4609x timer_svc_->process_expired();
400 4609x update_timerfd();
401 }
402
403 40232x lock.lock();
404
405 40232x completed_ops_.splice(local_ops);
406 40232x }
407
408 } // namespace boost::corosio::detail
409
410 #endif // BOOST_COROSIO_HAS_EPOLL
411
412 #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP
413