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 : /// No-op: write readiness is watched from registration on.
120 : std::error_code
121 HIT 2618 : ensure_write_registered(int, reactor_descriptor_state*) const noexcept
122 : {
123 2618 : return {};
124 : }
125 :
126 : /** Deregister a persistently registered descriptor.
127 :
128 : @param fd The file descriptor to deregister.
129 : */
130 : void deregister_descriptor(int fd) const;
131 :
132 : /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp).
133 76 : [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
134 : {
135 76 : return register_descriptor(read_fd, signal_pipe_reader_.arm());
136 : }
137 :
138 : private:
139 : void run_task(lock_type& lock, context_type& ctx, 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 1577 : inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
163 1577 : : epoll_fd_(-1)
164 1577 : , event_fd_(-1)
165 1577 : , timer_fd_(-1)
166 3154 : , event_buffer_(max_events_per_poll_)
167 : {
168 1577 : epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
169 1577 : if (epoll_fd_ < 0)
170 1 : detail::throw_system_error(make_err(errno), "epoll_create1");
171 :
172 1576 : event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
173 1576 : if (event_fd_ < 0)
174 : {
175 1 : int errn = errno;
176 1 : ::close(epoll_fd_);
177 1 : detail::throw_system_error(make_err(errn), "eventfd");
178 : }
179 :
180 1575 : timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
181 1575 : if (timer_fd_ < 0)
182 : {
183 1 : int errn = errno;
184 1 : ::close(event_fd_);
185 1 : ::close(epoll_fd_);
186 1 : detail::throw_system_error(make_err(errn), "timerfd_create");
187 : }
188 :
189 1574 : epoll_event ev{};
190 1574 : ev.events = EPOLLIN | EPOLLET;
191 1574 : ev.data.ptr = nullptr;
192 1574 : if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0)
193 : {
194 1 : int errn = errno;
195 1 : ::close(timer_fd_);
196 1 : ::close(event_fd_);
197 1 : ::close(epoll_fd_);
198 1 : detail::throw_system_error(make_err(errn), "epoll_ctl");
199 : }
200 :
201 1573 : epoll_event timer_ev{};
202 1573 : timer_ev.events = EPOLLIN | EPOLLERR;
203 1573 : timer_ev.data.ptr = &timer_fd_;
204 1573 : if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0)
205 : {
206 1 : int errn = errno;
207 1 : ::close(timer_fd_);
208 1 : ::close(event_fd_);
209 1 : ::close(epoll_fd_);
210 1 : detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)");
211 : }
212 :
213 1572 : timer_svc_ = &get_timer_service(ctx, *this);
214 1572 : timer_svc_->set_on_earliest_changed(
215 5979 : timer_service::callback(this, [](void* p) {
216 4407 : auto* self = static_cast<epoll_scheduler*>(p);
217 4407 : self->timerfd_stale_.store(true, std::memory_order_release);
218 4407 : self->interrupt_reactor();
219 4407 : }));
220 :
221 1572 : completed_ops_.push(&task_op_);
222 1587 : }
223 :
224 3144 : inline epoll_scheduler::~epoll_scheduler()
225 : {
226 1572 : if (timer_fd_ >= 0)
227 1572 : ::close(timer_fd_);
228 1572 : if (event_fd_ >= 0)
229 1572 : ::close(event_fd_);
230 1572 : if (epoll_fd_ >= 0)
231 1572 : ::close(epoll_fd_);
232 3144 : }
233 :
234 : inline void
235 1572 : epoll_scheduler::shutdown()
236 : {
237 1572 : shutdown_drain();
238 :
239 1572 : if (event_fd_ >= 0)
240 1572 : interrupt_reactor();
241 1572 : }
242 :
243 : inline void
244 28 : epoll_scheduler::configure_reactor(
245 : unsigned max_events,
246 : unsigned budget_init,
247 : unsigned budget_max,
248 : unsigned unassisted)
249 : {
250 28 : reactor_scheduler::configure_reactor(
251 : max_events, budget_init, budget_max, unassisted);
252 26 : event_buffer_.resize(max_events_per_poll_);
253 26 : }
254 :
255 : inline std::error_code
256 5923 : epoll_scheduler::register_descriptor(
257 : int fd, reactor_descriptor_state* desc) const
258 : {
259 5923 : epoll_event ev{};
260 5923 : ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
261 5923 : ev.data.ptr = desc;
262 :
263 5923 : bool unpollable = false;
264 5923 : if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
265 : {
266 : // EPERM: a file type epoll cannot watch. Its I/O does not
267 : // block, so adopt it unwatched, as asio does.
268 11 : if (errno != EPERM)
269 7 : return make_err(errno);
270 4 : unpollable = true;
271 : }
272 :
273 5916 : desc->registered_events = unpollable ? 0 : ev.events;
274 5916 : desc->unpollable = unpollable;
275 5916 : desc->fd = fd;
276 5916 : desc->scheduler_ = this;
277 5916 : desc->mutex.set_enabled(reactor_io_locking_);
278 5916 : desc->ready_events_.store(0, std::memory_order_relaxed);
279 :
280 5916 : conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
281 5916 : desc->impl_ref_.reset();
282 5916 : desc->read_ready = false;
283 5916 : desc->write_ready = false;
284 5916 : return {};
285 5916 : }
286 :
287 : inline void
288 5837 : epoll_scheduler::deregister_descriptor(int fd) const
289 : {
290 5837 : ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
291 5837 : }
292 :
293 : inline void
294 8498 : epoll_scheduler::interrupt_reactor() const
295 : {
296 8498 : bool expected = false;
297 8498 : if (eventfd_armed_.compare_exchange_strong(
298 : expected, true, std::memory_order_release,
299 : std::memory_order_relaxed))
300 : {
301 6645 : std::uint64_t val = 1;
302 6645 : if (::write(event_fd_, &val, sizeof(val)) < 0)
303 : {
304 : // The flag is what coalesces later interrupts into a byte
305 : // already in the eventfd; a write that failed put no byte
306 : // there, so leaving it armed would swallow every interrupt
307 : // that follows. Disarming keeps the cost to the interrupts
308 : // already in flight -- the next one arms and writes again,
309 : // instead of every one after this coalescing into a byte
310 : // that does not exist.
311 2 : eventfd_armed_.store(false, std::memory_order_release);
312 : }
313 : }
314 8498 : }
315 :
316 : inline void
317 10780 : epoll_scheduler::update_timerfd() const
318 : {
319 10780 : auto nearest = timer_svc_->nearest_expiry();
320 :
321 10780 : itimerspec ts{};
322 10780 : int flags = 0;
323 :
324 10780 : if (nearest == timer_service::time_point::max())
325 : {
326 : // No timers — disarm by setting to 0 (relative)
327 : }
328 : else
329 : {
330 9671 : auto now = std::chrono::steady_clock::now();
331 9671 : if (nearest <= now)
332 : {
333 : // Use 1ns instead of 0 — zero disarms the timerfd
334 1464 : ts.it_value.tv_nsec = 1;
335 : }
336 : else
337 : {
338 8207 : auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
339 8207 : nearest - now)
340 8207 : .count();
341 8207 : ts.it_value.tv_sec = nsec / 1000000000;
342 8207 : ts.it_value.tv_nsec = nsec % 1000000000;
343 8207 : if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
344 MIS 0 : ts.it_value.tv_nsec = 1;
345 : }
346 : }
347 :
348 HIT 10780 : if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
349 1 : detail::throw_system_error(make_err(errno), "timerfd_settime");
350 10779 : }
351 :
352 : inline void
353 41807 : epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
354 : {
355 : int timeout_ms;
356 41807 : if (task_interrupted_)
357 30227 : timeout_ms = 0;
358 11580 : else if (timeout_us < 0)
359 11232 : timeout_ms = -1;
360 : else
361 348 : timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
362 :
363 41807 : if (lock.owns_lock())
364 11582 : lock.unlock();
365 :
366 41807 : task_cleanup on_exit{this, &lock, ctx};
367 :
368 : // Flush deferred timerfd programming before blocking
369 41807 : if (timerfd_stale_.exchange(false, std::memory_order_acquire))
370 3696 : update_timerfd();
371 :
372 41806 : int nfds = ::epoll_wait(
373 41806 : epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()),
374 : timeout_ms);
375 :
376 41806 : if (nfds < 0 && errno != EINTR)
377 1 : detail::throw_system_error(make_err(errno), "epoll_wait");
378 :
379 41805 : bool check_timers = false;
380 41805 : ready_queue local_ops;
381 :
382 89644 : for (int i = 0; i < nfds; ++i)
383 : {
384 47839 : if (event_buffer_[i].data.ptr == nullptr)
385 : {
386 : std::uint64_t val;
387 : // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
388 5071 : [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
389 5071 : eventfd_armed_.store(false, std::memory_order_relaxed);
390 5071 : continue;
391 5071 : }
392 :
393 42768 : if (event_buffer_[i].data.ptr == &timer_fd_)
394 : {
395 : std::uint64_t expirations;
396 : // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
397 : [[maybe_unused]] auto r =
398 7084 : ::read(timer_fd_, &expirations, sizeof(expirations));
399 7084 : check_timers = true;
400 7084 : continue;
401 7084 : }
402 :
403 : auto* desc =
404 35684 : static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
405 :
406 : // EPOLLHUP maps to no reactor_event_* bit, so a HUP the widening
407 : // below did not translate would take no branch in
408 : // invoke_deferred_io() and, the registration being
409 : // edge-triggered, never get another chance.
410 : //
411 : // HUP without OUT means a non-socket: a pipe or tty whose peer
412 : // closed, with or without data still unread (HUP alone, or
413 : // IN|HUP). Treat it as readable, writable and faulted, the way
414 : // poll() and io_uring report it, so a parked error wait
415 : // completes too. Sockets carry at least OUT with HUP (a fresh
416 : // unconnected stream socket reports HUP|OUT), so they take the
417 : // plain widening, where forcing IN costs at most one spurious
418 : // EAGAIN.
419 35684 : std::uint32_t ev = event_buffer_[i].events;
420 35684 : if ((ev & (EPOLLHUP | EPOLLOUT)) == EPOLLHUP)
421 8 : ev |= EPOLLIN | EPOLLOUT | EPOLLERR;
422 35676 : else if (ev & EPOLLHUP)
423 235 : ev |= EPOLLIN | EPOLLOUT;
424 35684 : desc->add_ready_events(ev);
425 :
426 35684 : bool expected = false;
427 35684 : if (desc->is_enqueued_.compare_exchange_strong(
428 : expected, true, std::memory_order_release,
429 : std::memory_order_relaxed))
430 : {
431 35684 : local_ops.push(desc);
432 : }
433 : }
434 :
435 41805 : if (check_timers)
436 : {
437 7084 : timer_svc_->process_expired();
438 7084 : update_timerfd();
439 : }
440 :
441 41805 : lock.lock();
442 :
443 41805 : completed_ops_.splice(local_ops);
444 41807 : }
445 :
446 : } // namespace boost::corosio::detail
447 :
448 : #endif // BOOST_COROSIO_HAS_EPOLL
449 :
450 : #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP
|