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

99.3% Lines (152 / 153) 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 #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 76x [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
131 {
132 76x 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 1303x inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
160 1303x : epoll_fd_(-1)
161 1303x , event_fd_(-1)
162 1303x , timer_fd_(-1)
163 2606x , event_buffer_(max_events_per_poll_)
164 {
165 1303x epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
166 1303x if (epoll_fd_ < 0)
167 1x detail::throw_system_error(make_err(errno), "epoll_create1");
168
169 1302x event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
170 1302x 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 1301x timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
178 1301x 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 1300x epoll_event ev{};
187 1300x ev.events = EPOLLIN | EPOLLET;
188 1300x ev.data.ptr = nullptr;
189 1300x 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 1299x epoll_event timer_ev{};
199 1299x timer_ev.events = EPOLLIN | EPOLLERR;
200 1299x timer_ev.data.ptr = &timer_fd_;
201 1299x 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 1298x timer_svc_ = &get_timer_service(ctx, *this);
211 1298x timer_svc_->set_on_earliest_changed(
212 5472x timer_service::callback(this, [](void* p) {
213 4174x auto* self = static_cast<epoll_scheduler*>(p);
214 4174x self->timerfd_stale_.store(true, std::memory_order_release);
215 4174x self->interrupt_reactor();
216 4174x }));
217
218 1298x get_resolver_service(ctx, *this);
219 1298x get_signal_service(ctx, *this);
220 1298x get_stream_file_service(ctx, *this);
221 1298x get_random_access_file_service(ctx, *this);
222
223 1298x completed_ops_.push(&task_op_);
224 1313x }
225
226 2596x inline epoll_scheduler::~epoll_scheduler()
227 {
228 1298x if (timer_fd_ >= 0)
229 1298x ::close(timer_fd_);
230 1298x if (event_fd_ >= 0)
231 1298x ::close(event_fd_);
232 1298x if (epoll_fd_ >= 0)
233 1298x ::close(epoll_fd_);
234 2596x }
235
236 inline void
237 1298x epoll_scheduler::shutdown()
238 {
239 1298x shutdown_drain();
240
241 1298x if (event_fd_ >= 0)
242 1298x interrupt_reactor();
243 1298x }
244
245 inline void
246 27x epoll_scheduler::configure_reactor(
247 unsigned max_events,
248 unsigned budget_init,
249 unsigned budget_max,
250 unsigned unassisted)
251 {
252 27x reactor_scheduler::configure_reactor(
253 max_events, budget_init, budget_max, unassisted);
254 25x event_buffer_.resize(max_events_per_poll_);
255 25x }
256
257 inline std::error_code
258 5549x epoll_scheduler::register_descriptor(
259 int fd, reactor_descriptor_state* desc) const
260 {
261 5549x epoll_event ev{};
262 5549x ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
263 5549x ev.data.ptr = desc;
264
265 5549x if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
266 7x return make_err(errno);
267
268 5542x desc->registered_events = ev.events;
269 5542x desc->fd = fd;
270 5542x desc->scheduler_ = this;
271 5542x desc->mutex.set_enabled(reactor_io_locking_);
272 5542x desc->ready_events_.store(0, std::memory_order_relaxed);
273
274 5542x conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
275 5542x desc->impl_ref_.reset();
276 5542x desc->read_ready = false;
277 5542x desc->write_ready = false;
278 5542x return {};
279 5542x }
280
281 inline void
282 5467x epoll_scheduler::deregister_descriptor(int fd) const
283 {
284 5467x ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
285 5467x }
286
287 inline void
288 7568x epoll_scheduler::interrupt_reactor() const
289 {
290 7568x bool expected = false;
291 7568x if (eventfd_armed_.compare_exchange_strong(
292 expected, true, std::memory_order_release,
293 std::memory_order_relaxed))
294 {
295 6058x std::uint64_t val = 1;
296 6058x 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 7568x }
309
310 inline void
311 10598x epoll_scheduler::update_timerfd() const
312 {
313 10598x auto nearest = timer_svc_->nearest_expiry();
314
315 10598x itimerspec ts{};
316 10598x int flags = 0;
317
318 10598x if (nearest == timer_service::time_point::max())
319 {
320 // No timers — disarm by setting to 0 (relative)
321 }
322 else
323 {
324 9421x auto now = std::chrono::steady_clock::now();
325 9421x if (nearest <= now)
326 {
327 // Use 1ns instead of 0 — zero disarms the timerfd
328 1128x ts.it_value.tv_nsec = 1;
329 }
330 else
331 {
332 8293x auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
333 8293x nearest - now)
334 8293x .count();
335 8293x ts.it_value.tv_sec = nsec / 1000000000;
336 8293x ts.it_value.tv_nsec = nsec % 1000000000;
337 8293x if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
338 ✗ ts.it_value.tv_nsec = 1;
339 }
340 }
341
342 10598x if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
343 1x detail::throw_system_error(make_err(errno), "timerfd_settime");
344 10597x }
345
346 inline void
347 40977x epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
348 {
349 int timeout_ms;
350 40977x if (task_interrupted_)
351 29995x timeout_ms = 0;
352 10982x else if (timeout_us < 0)
353 10679x timeout_ms = -1;
354 else
355 303x timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
356
357 40977x if (lock.owns_lock())
358 10984x lock.unlock();
359
360 40977x task_cleanup on_exit{this, &lock, ctx};
361
362 // Flush deferred timerfd programming before blocking
363 40977x if (timerfd_stale_.exchange(false, std::memory_order_acquire))
364 3555x update_timerfd();
365
366 40976x int nfds = ::epoll_wait(
367 40976x epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()),
368 timeout_ms);
369
370 40976x if (nfds < 0 && errno != EINTR)
371 1x detail::throw_system_error(make_err(errno), "epoll_wait");
372
373 40975x bool check_timers = false;
374 40975x ready_queue local_ops;
375
376 90006x for (int i = 0; i < nfds; ++i)
377 {
378 49031x if (event_buffer_[i].data.ptr == nullptr)
379 {
380 std::uint64_t val;
381 // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
382 4758x [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
383 4758x eventfd_armed_.store(false, std::memory_order_relaxed);
384 4758x continue;
385 4758x }
386
387 44273x 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 7043x ::read(timer_fd_, &expirations, sizeof(expirations));
393 7043x check_timers = true;
394 7043x continue;
395 7043x }
396
397 auto* desc =
398 37230x static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
399 37230x desc->add_ready_events(event_buffer_[i].events);
400
401 37230x bool expected = false;
402 37230x if (desc->is_enqueued_.compare_exchange_strong(
403 expected, true, std::memory_order_release,
404 std::memory_order_relaxed))
405 {
406 37230x local_ops.push(desc);
407 }
408 }
409
410 40975x if (check_timers)
411 {
412 7043x timer_svc_->process_expired();
413 7043x update_timerfd();
414 }
415
416 40975x lock.lock();
417
418 40975x completed_ops_.splice(local_ops);
419 40977x }
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