TLA Line data 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 HIT 76 : [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
127 : {
128 76 : 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 1339 : inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
156 1339 : : epoll_fd_(-1)
157 1339 : , event_fd_(-1)
158 1339 : , timer_fd_(-1)
159 2678 : , event_buffer_(max_events_per_poll_)
160 : {
161 1339 : epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
162 1339 : if (epoll_fd_ < 0)
163 1 : detail::throw_system_error(make_err(errno), "epoll_create1");
164 :
165 1338 : event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
166 1338 : if (event_fd_ < 0)
167 : {
168 1 : int errn = errno;
169 1 : ::close(epoll_fd_);
170 1 : detail::throw_system_error(make_err(errn), "eventfd");
171 : }
172 :
173 1337 : timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
174 1337 : if (timer_fd_ < 0)
175 : {
176 1 : int errn = errno;
177 1 : ::close(event_fd_);
178 1 : ::close(epoll_fd_);
179 1 : detail::throw_system_error(make_err(errn), "timerfd_create");
180 : }
181 :
182 1336 : epoll_event ev{};
183 1336 : ev.events = EPOLLIN | EPOLLET;
184 1336 : ev.data.ptr = nullptr;
185 1336 : if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0)
186 : {
187 1 : int errn = errno;
188 1 : ::close(timer_fd_);
189 1 : ::close(event_fd_);
190 1 : ::close(epoll_fd_);
191 1 : detail::throw_system_error(make_err(errn), "epoll_ctl");
192 : }
193 :
194 1335 : epoll_event timer_ev{};
195 1335 : timer_ev.events = EPOLLIN | EPOLLERR;
196 1335 : timer_ev.data.ptr = &timer_fd_;
197 1335 : if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0)
198 : {
199 1 : int errn = errno;
200 1 : ::close(timer_fd_);
201 1 : ::close(event_fd_);
202 1 : ::close(epoll_fd_);
203 1 : detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)");
204 : }
205 :
206 1334 : timer_svc_ = &get_timer_service(ctx, *this);
207 1334 : timer_svc_->set_on_earliest_changed(
208 5516 : timer_service::callback(this, [](void* p) {
209 4182 : auto* self = static_cast<epoll_scheduler*>(p);
210 4182 : self->timerfd_stale_.store(true, std::memory_order_release);
211 4182 : self->interrupt_reactor();
212 4182 : }));
213 :
214 1334 : completed_ops_.push(&task_op_);
215 1349 : }
216 :
217 2668 : inline epoll_scheduler::~epoll_scheduler()
218 : {
219 1334 : if (timer_fd_ >= 0)
220 1334 : ::close(timer_fd_);
221 1334 : if (event_fd_ >= 0)
222 1334 : ::close(event_fd_);
223 1334 : if (epoll_fd_ >= 0)
224 1334 : ::close(epoll_fd_);
225 2668 : }
226 :
227 : inline void
228 1334 : epoll_scheduler::shutdown()
229 : {
230 1334 : shutdown_drain();
231 :
232 1334 : if (event_fd_ >= 0)
233 1334 : interrupt_reactor();
234 1334 : }
235 :
236 : inline void
237 28 : epoll_scheduler::configure_reactor(
238 : unsigned max_events,
239 : unsigned budget_init,
240 : unsigned budget_max,
241 : unsigned unassisted)
242 : {
243 28 : reactor_scheduler::configure_reactor(
244 : max_events, budget_init, budget_max, unassisted);
245 26 : event_buffer_.resize(max_events_per_poll_);
246 26 : }
247 :
248 : inline std::error_code
249 5792 : epoll_scheduler::register_descriptor(
250 : int fd, reactor_descriptor_state* desc) const
251 : {
252 5792 : epoll_event ev{};
253 5792 : ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
254 5792 : ev.data.ptr = desc;
255 :
256 5792 : if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
257 7 : return make_err(errno);
258 :
259 5785 : desc->registered_events = ev.events;
260 5785 : desc->fd = fd;
261 5785 : desc->scheduler_ = this;
262 5785 : desc->mutex.set_enabled(reactor_io_locking_);
263 5785 : desc->ready_events_.store(0, std::memory_order_relaxed);
264 :
265 5785 : conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
266 5785 : desc->impl_ref_.reset();
267 5785 : desc->read_ready = false;
268 5785 : desc->write_ready = false;
269 5785 : return {};
270 5785 : }
271 :
272 : inline void
273 5710 : epoll_scheduler::deregister_descriptor(int fd) const
274 : {
275 5710 : ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
276 5710 : }
277 :
278 : inline void
279 7551 : epoll_scheduler::interrupt_reactor() const
280 : {
281 7551 : bool expected = false;
282 7551 : if (eventfd_armed_.compare_exchange_strong(
283 : expected, true, std::memory_order_release,
284 : std::memory_order_relaxed))
285 : {
286 6134 : std::uint64_t val = 1;
287 6134 : 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 2 : eventfd_armed_.store(false, std::memory_order_release);
297 : }
298 : }
299 7551 : }
300 :
301 : inline void
302 10885 : epoll_scheduler::update_timerfd() const
303 : {
304 10885 : auto nearest = timer_svc_->nearest_expiry();
305 :
306 10885 : itimerspec ts{};
307 10885 : int flags = 0;
308 :
309 10885 : if (nearest == timer_service::time_point::max())
310 : {
311 : // No timers — disarm by setting to 0 (relative)
312 : }
313 : else
314 : {
315 9697 : auto now = std::chrono::steady_clock::now();
316 9697 : if (nearest <= now)
317 : {
318 : // Use 1ns instead of 0 — zero disarms the timerfd
319 1258 : ts.it_value.tv_nsec = 1;
320 : }
321 : else
322 : {
323 8439 : auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
324 8439 : nearest - now)
325 8439 : .count();
326 8439 : ts.it_value.tv_sec = nsec / 1000000000;
327 8439 : ts.it_value.tv_nsec = nsec % 1000000000;
328 8439 : if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
329 MIS 0 : ts.it_value.tv_nsec = 1;
330 : }
331 : }
332 :
333 HIT 10885 : if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
334 1 : detail::throw_system_error(make_err(errno), "timerfd_settime");
335 10884 : }
336 :
337 : inline void
338 38343 : epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
339 : {
340 : int timeout_ms;
341 38343 : if (task_interrupted_)
342 26939 : timeout_ms = 0;
343 11404 : else if (timeout_us < 0)
344 11126 : timeout_ms = -1;
345 : else
346 278 : timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
347 :
348 38343 : if (lock.owns_lock())
349 11406 : lock.unlock();
350 :
351 38343 : task_cleanup on_exit{this, &lock, ctx};
352 :
353 : // Flush deferred timerfd programming before blocking
354 38343 : if (timerfd_stale_.exchange(false, std::memory_order_acquire))
355 3677 : update_timerfd();
356 :
357 38342 : int nfds = ::epoll_wait(
358 38342 : epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()),
359 : timeout_ms);
360 :
361 38342 : if (nfds < 0 && errno != EINTR)
362 1 : detail::throw_system_error(make_err(errno), "epoll_wait");
363 :
364 38341 : bool check_timers = false;
365 38341 : ready_queue local_ops;
366 :
367 84441 : for (int i = 0; i < nfds; ++i)
368 : {
369 46100 : if (event_buffer_[i].data.ptr == nullptr)
370 : {
371 : std::uint64_t val;
372 : // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
373 4798 : [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
374 4798 : eventfd_armed_.store(false, std::memory_order_relaxed);
375 4798 : continue;
376 4798 : }
377 :
378 41302 : 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 7208 : ::read(timer_fd_, &expirations, sizeof(expirations));
384 7208 : check_timers = true;
385 7208 : continue;
386 7208 : }
387 :
388 : auto* desc =
389 34094 : static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
390 :
391 : // A pipe or tty whose peer closed reports EPOLLHUP on its own --
392 : // no EPOLLIN, no EPOLLERR -- and EPOLLHUP maps to no
393 : // reactor_event_* bit, so invoke_deferred_io() would take no
394 : // branch and, the registration being edge-triggered, never get
395 : // another chance.
396 : //
397 : // Sockets are unaffected because they never report EPOLLHUP
398 : // alone. tcp_poll() and unix_poll() raise it only once
399 : // sk_shutdown is SHUTDOWN_MASK (or the state is TCP_CLOSE), and
400 : // both also report the socket readable and writable there --
401 : // tcp_poll() takes an explicit `else mask |= EPOLLOUT` branch
402 : // once SEND_SHUTDOWN is set, because a send on a shut-down
403 : // socket fails fast rather than blocking. Measured on every
404 : // state that produces EPOLLHUP -- peer close plus local
405 : // SHUT_WR/SHUT_RDWR, RST, RST with the send buffer full, and
406 : // the AF_UNIX equivalents -- the mask is always IN|OUT|HUP
407 : // (0x15), or IN|OUT|ERR|HUP (0x1d) for a reset. The forced bits
408 : // are therefore already set on every socket path.
409 34094 : std::uint32_t ev = event_buffer_[i].events;
410 34094 : if (ev & EPOLLHUP)
411 229 : ev |= EPOLLIN | EPOLLOUT;
412 34094 : desc->add_ready_events(ev);
413 :
414 34094 : bool expected = false;
415 34094 : if (desc->is_enqueued_.compare_exchange_strong(
416 : expected, true, std::memory_order_release,
417 : std::memory_order_relaxed))
418 : {
419 34094 : local_ops.push(desc);
420 : }
421 : }
422 :
423 38341 : if (check_timers)
424 : {
425 7208 : timer_svc_->process_expired();
426 7208 : update_timerfd();
427 : }
428 :
429 38341 : lock.lock();
430 :
431 38341 : completed_ops_.splice(local_ops);
432 38343 : }
433 :
434 : } // namespace boost::corosio::detail
435 :
436 : #endif // BOOST_COROSIO_HAS_EPOLL
437 :
438 : #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP
|