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 [[nodiscard]] std::error_code
131 69x register_signal_reader(int read_fd) override
132 {
133 69x return register_descriptor(read_fd, signal_pipe_reader_.arm());
134 }
135
136 private:
137 void
138 run_task(lock_type& lock, context_type& ctx,
139 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 1226x inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
163 1226x : epoll_fd_(-1)
164 1226x , event_fd_(-1)
165 1226x , timer_fd_(-1)
166 2452x , event_buffer_(max_events_per_poll_)
167 {
168 1226x epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
169 1226x if (epoll_fd_ < 0)
170 1x detail::throw_system_error(make_err(errno), "epoll_create1");
171
172 1225x event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
173 1225x 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 1224x timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
181 1224x 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 1223x epoll_event ev{};
190 1223x ev.events = EPOLLIN | EPOLLET;
191 1223x ev.data.ptr = nullptr;
192 1223x 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 1222x epoll_event timer_ev{};
202 1222x timer_ev.events = EPOLLIN | EPOLLERR;
203 1222x timer_ev.data.ptr = &timer_fd_;
204 1222x 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 1221x timer_svc_ = &get_timer_service(ctx, *this);
214 1221x timer_svc_->set_on_earliest_changed(
215 7293x timer_service::callback(this, [](void* p) {
216 6072x auto* self = static_cast<epoll_scheduler*>(p);
217 6072x self->timerfd_stale_.store(true, std::memory_order_release);
218 6072x self->interrupt_reactor();
219 6072x }));
220
221 1221x get_resolver_service(ctx, *this);
222 1221x get_signal_service(ctx, *this);
223 1221x get_stream_file_service(ctx, *this);
224 1221x get_random_access_file_service(ctx, *this);
225
226 1221x completed_ops_.push(&task_op_);
227 1236x }
228
229 2442x inline epoll_scheduler::~epoll_scheduler()
230 {
231 1221x if (timer_fd_ >= 0)
232 1221x ::close(timer_fd_);
233 1221x if (event_fd_ >= 0)
234 1221x ::close(event_fd_);
235 1221x if (epoll_fd_ >= 0)
236 1221x ::close(epoll_fd_);
237 2442x }
238
239 inline void
240 1221x epoll_scheduler::shutdown()
241 {
242 1221x shutdown_drain();
243
244 1221x if (event_fd_ >= 0)
245 1221x interrupt_reactor();
246 1221x }
247
248 inline void
249 23x epoll_scheduler::configure_reactor(
250 unsigned max_events,
251 unsigned budget_init,
252 unsigned budget_max,
253 unsigned unassisted)
254 {
255 23x reactor_scheduler::configure_reactor(
256 max_events, budget_init, budget_max, unassisted);
257 21x event_buffer_.resize(max_events_per_poll_);
258 21x }
259
260 inline std::error_code
261 9487x epoll_scheduler::register_descriptor(int fd, reactor_descriptor_state* desc) const
262 {
263 9487x epoll_event ev{};
264 9487x ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
265 9487x ev.data.ptr = desc;
266
267 9487x if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
268 7x return make_err(errno);
269
270 9480x desc->registered_events = ev.events;
271 9480x desc->fd = fd;
272 9480x desc->scheduler_ = this;
273 9480x desc->mutex.set_enabled(reactor_io_locking_);
274 9480x desc->ready_events_.store(0, std::memory_order_relaxed);
275
276 9480x conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
277 9480x desc->impl_ref_.reset();
278 9480x desc->read_ready = false;
279 9480x desc->write_ready = false;
280 9480x return {};
281 9480x }
282
283 inline void
284 9412x epoll_scheduler::deregister_descriptor(int fd) const
285 {
286 9412x ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
287 9412x }
288
289 inline void
290 8464x epoll_scheduler::interrupt_reactor() const
291 {
292 8464x bool expected = false;
293 8464x if (eventfd_armed_.compare_exchange_strong(
294 expected, true, std::memory_order_release,
295 std::memory_order_relaxed))
296 {
297 7034x std::uint64_t val = 1;
298 7034x if (::write(event_fd_, &val, sizeof(val)) < 0)
299 {
300 // The flag is what coalesces later interrupts into a byte
301 // already in the eventfd; a write that failed put no byte
302 // there, so leaving it armed would swallow every interrupt
303 // that follows. Disarming keeps the cost to the interrupts
304 // already in flight -- the next one arms and writes again,
305 // instead of every one after this coalescing into a byte
306 // that does not exist.
307 2x eventfd_armed_.store(false, std::memory_order_release);
308 }
309 }
310 8464x }
311
312 inline void
313 14962x epoll_scheduler::update_timerfd() const
314 {
315 14962x auto nearest = timer_svc_->nearest_expiry();
316
317 14962x itimerspec ts{};
318 14962x int flags = 0;
319
320 14962x if (nearest == timer_service::time_point::max())
321 {
322 // No timers — disarm by setting to 0 (relative)
323 }
324 else
325 {
326 13819x auto now = std::chrono::steady_clock::now();
327 13819x if (nearest <= now)
328 {
329 // Use 1ns instead of 0 — zero disarms the timerfd
330 1246x ts.it_value.tv_nsec = 1;
331 }
332 else
333 {
334 12573x auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
335 12573x nearest - now)
336 12573x .count();
337 12573x ts.it_value.tv_sec = nsec / 1000000000;
338 12573x ts.it_value.tv_nsec = nsec % 1000000000;
339 12573x if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
340 ts.it_value.tv_nsec = 1;
341 }
342 }
343
344 14962x if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
345 1x detail::throw_system_error(make_err(errno), "timerfd_settime");
346 14961x }
347
348 inline void
349 47003x epoll_scheduler::run_task(
350 lock_type& lock, context_type& ctx, long timeout_us)
351 {
352 int timeout_ms;
353 47003x if (task_interrupted_)
354 30232x timeout_ms = 0;
355 16771x else if (timeout_us < 0)
356 16756x timeout_ms = -1;
357 else
358 15x timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
359
360 47003x if (lock.owns_lock())
361 16773x lock.unlock();
362
363 47003x task_cleanup on_exit{this, &lock, ctx};
364
365 // Flush deferred timerfd programming before blocking
366 47003x if (timerfd_stale_.exchange(false, std::memory_order_acquire))
367 5457x update_timerfd();
368
369 47002x int nfds = ::epoll_wait(
370 epoll_fd_, event_buffer_.data(),
371 47002x static_cast<int>(event_buffer_.size()), timeout_ms);
372
373 47002x if (nfds < 0 && errno != EINTR)
374 1x detail::throw_system_error(make_err(errno), "epoll_wait");
375
376 47001x bool check_timers = false;
377 47001x ready_queue local_ops;
378
379 100534x for (int i = 0; i < nfds; ++i)
380 {
381 53533x if (event_buffer_[i].data.ptr == nullptr)
382 {
383 std::uint64_t val;
384 // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
385 5811x [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
386 5811x eventfd_armed_.store(false, std::memory_order_relaxed);
387 5811x continue;
388 5811x }
389
390 47722x if (event_buffer_[i].data.ptr == &timer_fd_)
391 {
392 std::uint64_t expirations;
393 // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
394 [[maybe_unused]] auto r =
395 9505x ::read(timer_fd_, &expirations, sizeof(expirations));
396 9505x check_timers = true;
397 9505x continue;
398 9505x }
399
400 auto* desc =
401 38217x static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
402 38217x desc->add_ready_events(event_buffer_[i].events);
403
404 38217x bool expected = false;
405 38217x if (desc->is_enqueued_.compare_exchange_strong(
406 expected, true, std::memory_order_release,
407 std::memory_order_relaxed))
408 {
409 38217x local_ops.push(desc);
410 }
411 }
412
413 47001x if (check_timers)
414 {
415 9505x timer_svc_->process_expired();
416 9505x update_timerfd();
417 }
418
419 47001x lock.lock();
420
421 47001x completed_ops_.splice(local_ops);
422 47003x }
423
424 } // namespace boost::corosio::detail
425
426 #endif // BOOST_COROSIO_HAS_EPOLL
427
428 #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP
429