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

99.3% Lines (148 / 149) 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 1305x inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
156 1305x : epoll_fd_(-1)
157 1305x , event_fd_(-1)
158 1305x , timer_fd_(-1)
159 2610x , event_buffer_(max_events_per_poll_)
160 {
161 1305x epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
162 1305x if (epoll_fd_ < 0)
163 1x detail::throw_system_error(make_err(errno), "epoll_create1");
164
165 1304x event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
166 1304x 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 1303x timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
174 1303x 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 1302x epoll_event ev{};
183 1302x ev.events = EPOLLIN | EPOLLET;
184 1302x ev.data.ptr = nullptr;
185 1302x 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 1301x epoll_event timer_ev{};
195 1301x timer_ev.events = EPOLLIN | EPOLLERR;
196 1301x timer_ev.data.ptr = &timer_fd_;
197 1301x 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 1300x timer_svc_ = &get_timer_service(ctx, *this);
207 1300x timer_svc_->set_on_earliest_changed(
208 5470x timer_service::callback(this, [](void* p) {
209 4170x auto* self = static_cast<epoll_scheduler*>(p);
210 4170x self->timerfd_stale_.store(true, std::memory_order_release);
211 4170x self->interrupt_reactor();
212 4170x }));
213
214 1300x completed_ops_.push(&task_op_);
215 1315x }
216
217 2600x inline epoll_scheduler::~epoll_scheduler()
218 {
219 1300x if (timer_fd_ >= 0)
220 1300x ::close(timer_fd_);
221 1300x if (event_fd_ >= 0)
222 1300x ::close(event_fd_);
223 1300x if (epoll_fd_ >= 0)
224 1300x ::close(epoll_fd_);
225 2600x }
226
227 inline void
228 1300x epoll_scheduler::shutdown()
229 {
230 1300x shutdown_drain();
231
232 1300x if (event_fd_ >= 0)
233 1300x interrupt_reactor();
234 1300x }
235
236 inline void
237 27x epoll_scheduler::configure_reactor(
238 unsigned max_events,
239 unsigned budget_init,
240 unsigned budget_max,
241 unsigned unassisted)
242 {
243 27x reactor_scheduler::configure_reactor(
244 max_events, budget_init, budget_max, unassisted);
245 25x event_buffer_.resize(max_events_per_poll_);
246 25x }
247
248 inline std::error_code
249 5883x epoll_scheduler::register_descriptor(
250 int fd, reactor_descriptor_state* desc) const
251 {
252 5883x epoll_event ev{};
253 5883x ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
254 5883x ev.data.ptr = desc;
255
256 5883x if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
257 7x return make_err(errno);
258
259 5876x desc->registered_events = ev.events;
260 5876x desc->fd = fd;
261 5876x desc->scheduler_ = this;
262 5876x desc->mutex.set_enabled(reactor_io_locking_);
263 5876x desc->ready_events_.store(0, std::memory_order_relaxed);
264
265 5876x conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
266 5876x desc->impl_ref_.reset();
267 5876x desc->read_ready = false;
268 5876x desc->write_ready = false;
269 5876x return {};
270 5876x }
271
272 inline void
273 5801x epoll_scheduler::deregister_descriptor(int fd) const
274 {
275 5801x ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
276 5801x }
277
278 inline void
279 7504x epoll_scheduler::interrupt_reactor() const
280 {
281 7504x bool expected = false;
282 7504x 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 7504x }
300
301 inline void
302 11268x epoll_scheduler::update_timerfd() const
303 {
304 11268x auto nearest = timer_svc_->nearest_expiry();
305
306 11268x itimerspec ts{};
307 11268x int flags = 0;
308
309 11268x if (nearest == timer_service::time_point::max())
310 {
311 // No timers — disarm by setting to 0 (relative)
312 }
313 else
314 {
315 10080x auto now = std::chrono::steady_clock::now();
316 10080x if (nearest <= now)
317 {
318 // Use 1ns instead of 0 — zero disarms the timerfd
319 1300x ts.it_value.tv_nsec = 1;
320 }
321 else
322 {
323 8780x auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
324 8780x nearest - now)
325 8780x .count();
326 8780x ts.it_value.tv_sec = nsec / 1000000000;
327 8780x ts.it_value.tv_nsec = nsec % 1000000000;
328 8780x if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
329 ✗ ts.it_value.tv_nsec = 1;
330 }
331 }
332
333 11268x if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
334 1x detail::throw_system_error(make_err(errno), "timerfd_settime");
335 11267x }
336
337 inline void
338 40730x epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
339 {
340 int timeout_ms;
341 40730x if (task_interrupted_)
342 29071x timeout_ms = 0;
343 11659x else if (timeout_us < 0)
344 11342x timeout_ms = -1;
345 else
346 317x timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
347
348 40730x if (lock.owns_lock())
349 11661x lock.unlock();
350
351 40730x task_cleanup on_exit{this, &lock, ctx};
352
353 // Flush deferred timerfd programming before blocking
354 40730x if (timerfd_stale_.exchange(false, std::memory_order_acquire))
355 3687x update_timerfd();
356
357 40729x int nfds = ::epoll_wait(
358 40729x epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()),
359 timeout_ms);
360
361 40729x if (nfds < 0 && errno != EINTR)
362 1x detail::throw_system_error(make_err(errno), "epoll_wait");
363
364 40728x bool check_timers = false;
365 40728x ready_queue local_ops;
366
367 89348x for (int i = 0; i < nfds; ++i)
368 {
369 48620x if (event_buffer_[i].data.ptr == nullptr)
370 {
371 std::uint64_t val;
372 // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
373 4832x [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
374 4832x eventfd_armed_.store(false, std::memory_order_relaxed);
375 4832x continue;
376 4832x }
377
378 43788x 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 7581x ::read(timer_fd_, &expirations, sizeof(expirations));
384 7581x check_timers = true;
385 7581x continue;
386 7581x }
387
388 auto* desc =
389 36207x static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
390 36207x desc->add_ready_events(event_buffer_[i].events);
391
392 36207x bool expected = false;
393 36207x if (desc->is_enqueued_.compare_exchange_strong(
394 expected, true, std::memory_order_release,
395 std::memory_order_relaxed))
396 {
397 36207x local_ops.push(desc);
398 }
399 }
400
401 40728x if (check_timers)
402 {
403 7581x timer_svc_->process_expired();
404 7581x update_timerfd();
405 }
406
407 40728x lock.lock();
408
409 40728x completed_ops_.splice(local_ops);
410 40730x }
411
412 } // namespace boost::corosio::detail
413
414 #endif // BOOST_COROSIO_HAS_EPOLL
415
416 #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP
417