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 : #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 HIT 69 : [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
131 : {
132 69 : 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 1232 : inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
160 1232 : : epoll_fd_(-1)
161 1232 : , event_fd_(-1)
162 1232 : , timer_fd_(-1)
163 2464 : , event_buffer_(max_events_per_poll_)
164 : {
165 1232 : epoll_fd_ = ::epoll_create1(EPOLL_CLOEXEC);
166 1232 : if (epoll_fd_ < 0)
167 1 : detail::throw_system_error(make_err(errno), "epoll_create1");
168 :
169 1231 : event_fd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC);
170 1231 : if (event_fd_ < 0)
171 : {
172 1 : int errn = errno;
173 1 : ::close(epoll_fd_);
174 1 : detail::throw_system_error(make_err(errn), "eventfd");
175 : }
176 :
177 1230 : timer_fd_ = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
178 1230 : if (timer_fd_ < 0)
179 : {
180 1 : int errn = errno;
181 1 : ::close(event_fd_);
182 1 : ::close(epoll_fd_);
183 1 : detail::throw_system_error(make_err(errn), "timerfd_create");
184 : }
185 :
186 1229 : epoll_event ev{};
187 1229 : ev.events = EPOLLIN | EPOLLET;
188 1229 : ev.data.ptr = nullptr;
189 1229 : if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, event_fd_, &ev) < 0)
190 : {
191 1 : int errn = errno;
192 1 : ::close(timer_fd_);
193 1 : ::close(event_fd_);
194 1 : ::close(epoll_fd_);
195 1 : detail::throw_system_error(make_err(errn), "epoll_ctl");
196 : }
197 :
198 1228 : epoll_event timer_ev{};
199 1228 : timer_ev.events = EPOLLIN | EPOLLERR;
200 1228 : timer_ev.data.ptr = &timer_fd_;
201 1228 : if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, timer_fd_, &timer_ev) < 0)
202 : {
203 1 : int errn = errno;
204 1 : ::close(timer_fd_);
205 1 : ::close(event_fd_);
206 1 : ::close(epoll_fd_);
207 1 : detail::throw_system_error(make_err(errn), "epoll_ctl (timerfd)");
208 : }
209 :
210 1227 : timer_svc_ = &get_timer_service(ctx, *this);
211 1227 : timer_svc_->set_on_earliest_changed(
212 5586 : timer_service::callback(this, [](void* p) {
213 4359 : auto* self = static_cast<epoll_scheduler*>(p);
214 4359 : self->timerfd_stale_.store(true, std::memory_order_release);
215 4359 : self->interrupt_reactor();
216 4359 : }));
217 :
218 1227 : get_resolver_service(ctx, *this);
219 1227 : get_signal_service(ctx, *this);
220 1227 : get_stream_file_service(ctx, *this);
221 1227 : get_random_access_file_service(ctx, *this);
222 :
223 1227 : completed_ops_.push(&task_op_);
224 1242 : }
225 :
226 2454 : inline epoll_scheduler::~epoll_scheduler()
227 : {
228 1227 : if (timer_fd_ >= 0)
229 1227 : ::close(timer_fd_);
230 1227 : if (event_fd_ >= 0)
231 1227 : ::close(event_fd_);
232 1227 : if (epoll_fd_ >= 0)
233 1227 : ::close(epoll_fd_);
234 2454 : }
235 :
236 : inline void
237 1227 : epoll_scheduler::shutdown()
238 : {
239 1227 : shutdown_drain();
240 :
241 1227 : if (event_fd_ >= 0)
242 1227 : interrupt_reactor();
243 1227 : }
244 :
245 : inline void
246 23 : epoll_scheduler::configure_reactor(
247 : unsigned max_events,
248 : unsigned budget_init,
249 : unsigned budget_max,
250 : unsigned unassisted)
251 : {
252 23 : reactor_scheduler::configure_reactor(
253 : max_events, budget_init, budget_max, unassisted);
254 21 : event_buffer_.resize(max_events_per_poll_);
255 21 : }
256 :
257 : inline std::error_code
258 5748 : epoll_scheduler::register_descriptor(
259 : int fd, reactor_descriptor_state* desc) const
260 : {
261 5748 : epoll_event ev{};
262 5748 : ev.events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLERR | EPOLLHUP;
263 5748 : ev.data.ptr = desc;
264 :
265 5748 : if (::epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, fd, &ev) < 0)
266 7 : return make_err(errno);
267 :
268 5741 : desc->registered_events = ev.events;
269 5741 : desc->fd = fd;
270 5741 : desc->scheduler_ = this;
271 5741 : desc->mutex.set_enabled(reactor_io_locking_);
272 5741 : desc->ready_events_.store(0, std::memory_order_relaxed);
273 :
274 5741 : conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
275 5741 : desc->impl_ref_.reset();
276 5741 : desc->read_ready = false;
277 5741 : desc->write_ready = false;
278 5741 : return {};
279 5741 : }
280 :
281 : inline void
282 5673 : epoll_scheduler::deregister_descriptor(int fd) const
283 : {
284 5673 : ::epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, fd, nullptr);
285 5673 : }
286 :
287 : inline void
288 6764 : epoll_scheduler::interrupt_reactor() const
289 : {
290 6764 : bool expected = false;
291 6764 : if (eventfd_armed_.compare_exchange_strong(
292 : expected, true, std::memory_order_release,
293 : std::memory_order_relaxed))
294 : {
295 5222 : std::uint64_t val = 1;
296 5222 : 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 2 : eventfd_armed_.store(false, std::memory_order_release);
306 : }
307 : }
308 6764 : }
309 :
310 : inline void
311 11143 : epoll_scheduler::update_timerfd() const
312 : {
313 11143 : auto nearest = timer_svc_->nearest_expiry();
314 :
315 11143 : itimerspec ts{};
316 11143 : int flags = 0;
317 :
318 11143 : if (nearest == timer_service::time_point::max())
319 : {
320 : // No timers — disarm by setting to 0 (relative)
321 : }
322 : else
323 : {
324 9948 : auto now = std::chrono::steady_clock::now();
325 9948 : if (nearest <= now)
326 : {
327 : // Use 1ns instead of 0 — zero disarms the timerfd
328 1167 : ts.it_value.tv_nsec = 1;
329 : }
330 : else
331 : {
332 8781 : auto nsec = std::chrono::duration_cast<std::chrono::nanoseconds>(
333 8781 : nearest - now)
334 8781 : .count();
335 8781 : ts.it_value.tv_sec = nsec / 1000000000;
336 8781 : ts.it_value.tv_nsec = nsec % 1000000000;
337 8781 : if (ts.it_value.tv_sec == 0 && ts.it_value.tv_nsec == 0)
338 MIS 0 : ts.it_value.tv_nsec = 1;
339 : }
340 : }
341 :
342 HIT 11143 : if (::timerfd_settime(timer_fd_, flags, &ts, nullptr) < 0)
343 1 : detail::throw_system_error(make_err(errno), "timerfd_settime");
344 11142 : }
345 :
346 : inline void
347 39736 : epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
348 : {
349 : int timeout_ms;
350 39736 : if (task_interrupted_)
351 28575 : timeout_ms = 0;
352 11161 : else if (timeout_us < 0)
353 11146 : timeout_ms = -1;
354 : else
355 15 : timeout_ms = static_cast<int>((timeout_us + 999) / 1000);
356 :
357 39736 : if (lock.owns_lock())
358 11163 : lock.unlock();
359 :
360 39736 : task_cleanup on_exit{this, &lock, ctx};
361 :
362 : // Flush deferred timerfd programming before blocking
363 39736 : if (timerfd_stale_.exchange(false, std::memory_order_acquire))
364 3637 : update_timerfd();
365 :
366 39735 : int nfds = ::epoll_wait(
367 39735 : epoll_fd_, event_buffer_.data(), static_cast<int>(event_buffer_.size()),
368 : timeout_ms);
369 :
370 39735 : if (nfds < 0 && errno != EINTR)
371 1 : detail::throw_system_error(make_err(errno), "epoll_wait");
372 :
373 39734 : bool check_timers = false;
374 39734 : ready_queue local_ops;
375 :
376 86129 : for (int i = 0; i < nfds; ++i)
377 : {
378 46395 : if (event_buffer_[i].data.ptr == nullptr)
379 : {
380 : std::uint64_t val;
381 : // NOLINTNEXTLINE(clang-analyzer-unix.BlockInCriticalSection)
382 3993 : [[maybe_unused]] auto r = ::read(event_fd_, &val, sizeof(val));
383 3993 : eventfd_armed_.store(false, std::memory_order_relaxed);
384 3993 : continue;
385 3993 : }
386 :
387 42402 : 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 7506 : ::read(timer_fd_, &expirations, sizeof(expirations));
393 7506 : check_timers = true;
394 7506 : continue;
395 7506 : }
396 :
397 : auto* desc =
398 34896 : static_cast<reactor_descriptor_state*>(event_buffer_[i].data.ptr);
399 34896 : desc->add_ready_events(event_buffer_[i].events);
400 :
401 34896 : bool expected = false;
402 34896 : if (desc->is_enqueued_.compare_exchange_strong(
403 : expected, true, std::memory_order_release,
404 : std::memory_order_relaxed))
405 : {
406 34896 : local_ops.push(desc);
407 : }
408 : }
409 :
410 39734 : if (check_timers)
411 : {
412 7506 : timer_svc_->process_expired();
413 7506 : update_timerfd();
414 : }
415 :
416 39734 : lock.lock();
417 :
418 39734 : completed_ops_.splice(local_ops);
419 39736 : }
420 :
421 : } // namespace boost::corosio::detail
422 :
423 : #endif // BOOST_COROSIO_HAS_EPOLL
424 :
425 : #endif // BOOST_COROSIO_NATIVE_DETAIL_EPOLL_EPOLL_SCHEDULER_HPP
|