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_SELECT_SELECT_SCHEDULER_HPP
12 : #define BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP
13 :
14 : #include <boost/corosio/detail/platform.hpp>
15 :
16 : #if BOOST_COROSIO_HAS_SELECT
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/select/select_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 <sys/select.h>
35 : #include <unistd.h>
36 : #include <errno.h>
37 : #include <fcntl.h>
38 :
39 : #include <atomic>
40 : #include <chrono>
41 : #include <cstdint>
42 : #include <limits>
43 : #include <mutex>
44 : #include <new>
45 : #include <unordered_map>
46 :
47 : namespace boost::corosio::detail {
48 :
49 : struct select_op;
50 :
51 : /** POSIX scheduler using select() for I/O multiplexing.
52 :
53 : This scheduler implements the scheduler interface using the POSIX select()
54 : call for I/O event notification. It inherits the shared reactor threading
55 : model from reactor_scheduler: signal state machine, inline completion
56 : budget, work counting, and the do_one event loop.
57 :
58 : The design mirrors epoll_scheduler for behavioral consistency:
59 : - Same single-reactor thread coordination model
60 : - Same deferred I/O pattern (reactor marks ready; workers do I/O)
61 : - Same timer integration pattern
62 :
63 : Known Limitations:
64 : - FD_SETSIZE (~1024) limits maximum concurrent connections
65 : - O(n) scanning: rebuilds fd_sets each iteration
66 : - Level-triggered only (no edge-triggered mode)
67 :
68 : @par Thread Safety
69 : All public member functions are thread-safe.
70 : */
71 : class BOOST_COROSIO_DECL select_scheduler final : public reactor_scheduler
72 : {
73 : public:
74 : /** Construct the scheduler.
75 :
76 : Creates a self-pipe for reactor interruption.
77 :
78 : @param ctx Reference to the owning execution_context.
79 : @param concurrency_hint Hint for expected thread count (unused).
80 : */
81 : select_scheduler(capy::execution_context& ctx, int concurrency_hint = -1);
82 :
83 : /// Destroy the scheduler.
84 : ~select_scheduler() override;
85 :
86 : select_scheduler(select_scheduler const&) = delete;
87 : select_scheduler& operator=(select_scheduler const&) = delete;
88 :
89 : /// Shut down the scheduler, draining pending operations.
90 : void shutdown() override;
91 :
92 : /** Return the maximum file descriptor value supported.
93 :
94 : Returns FD_SETSIZE - 1, the maximum fd value that can be
95 : monitored by select(). Operations with fd >= FD_SETSIZE
96 : will fail with EINVAL.
97 :
98 : @return The maximum supported file descriptor value.
99 : */
100 : static constexpr int max_fd() noexcept
101 : {
102 : return FD_SETSIZE - 1;
103 : }
104 :
105 : /** Register a descriptor for persistent monitoring.
106 :
107 : The fd is added to the registered_descs_ map and will be
108 : included in subsequent select() calls. The reactor is
109 : interrupted so a blocked select() rebuilds its fd_sets.
110 :
111 : @param fd The file descriptor to register.
112 : @param desc Pointer to descriptor state for this fd.
113 :
114 : @return The error if the fd cannot be tracked, otherwise a
115 : default constructed error code.
116 : */
117 : std::error_code
118 : register_descriptor(int fd, reactor_descriptor_state* desc) const;
119 :
120 : /** Deregister a persistently registered descriptor.
121 :
122 : @param fd The file descriptor to deregister.
123 : */
124 : void deregister_descriptor(int fd) const;
125 :
126 : /** Interrupt the reactor so it rebuilds its fd_sets.
127 :
128 : Called when a write, connect, or write-wait op is registered
129 : after the reactor's snapshot was taken. Without this,
130 : select() may block not watching for writability on the fd.
131 : */
132 : void notify_reactor() const;
133 :
134 : /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp).
135 : [[nodiscard]] std::error_code
136 HIT 47 : register_signal_reader(int read_fd) override
137 : {
138 47 : return register_descriptor(read_fd, signal_pipe_reader_.arm());
139 : }
140 :
141 : private:
142 : void
143 : run_task(lock_type& lock, context_type* ctx,
144 : long timeout_us) override;
145 : void interrupt_reactor() const override;
146 : long calculate_timeout(long requested_timeout_us) const;
147 :
148 : // Watches the global signal self-pipe's read end (armed lazily by
149 : // register_signal_reader on the first signal registration).
150 : reactor_signal_pipe_reader signal_pipe_reader_;
151 :
152 : // Self-pipe for interrupting select()
153 : int pipe_fds_[2]; // [0]=read, [1]=write
154 :
155 : // Per-fd tracking for fd_set building
156 : mutable std::unordered_map<int, reactor_descriptor_state*> registered_descs_;
157 : mutable int max_fd_ = -1;
158 : };
159 :
160 788 : inline select_scheduler::select_scheduler(capy::execution_context& ctx, int)
161 788 : : pipe_fds_{-1, -1}
162 788 : , max_fd_(-1)
163 : {
164 788 : if (::pipe(pipe_fds_) < 0)
165 1 : detail::throw_system_error(make_err(errno), "pipe");
166 :
167 2352 : for (int i = 0; i < 2; ++i)
168 : {
169 1571 : int flags = ::fcntl(pipe_fds_[i], F_GETFL, 0);
170 1571 : if (flags == -1)
171 : {
172 2 : int errn = errno;
173 2 : ::close(pipe_fds_[0]);
174 2 : ::close(pipe_fds_[1]);
175 2 : detail::throw_system_error(make_err(errn), "fcntl F_GETFL");
176 : }
177 1569 : if (::fcntl(pipe_fds_[i], F_SETFL, flags | O_NONBLOCK) == -1)
178 : {
179 2 : int errn = errno;
180 2 : ::close(pipe_fds_[0]);
181 2 : ::close(pipe_fds_[1]);
182 2 : detail::throw_system_error(make_err(errn), "fcntl F_SETFL");
183 : }
184 1567 : if (::fcntl(pipe_fds_[i], F_SETFD, FD_CLOEXEC) == -1)
185 : {
186 2 : int errn = errno;
187 2 : ::close(pipe_fds_[0]);
188 2 : ::close(pipe_fds_[1]);
189 2 : detail::throw_system_error(make_err(errn), "fcntl F_SETFD");
190 : }
191 : }
192 :
193 781 : timer_svc_ = &get_timer_service(ctx, *this);
194 781 : timer_svc_->set_on_earliest_changed(
195 3604 : timer_service::callback(this, [](void* p) {
196 2823 : static_cast<select_scheduler*>(p)->interrupt_reactor();
197 2823 : }));
198 :
199 781 : get_resolver_service(ctx, *this);
200 781 : get_signal_service(ctx, *this);
201 781 : get_stream_file_service(ctx, *this);
202 781 : get_random_access_file_service(ctx, *this);
203 :
204 781 : completed_ops_.push(&task_op_);
205 802 : }
206 :
207 1562 : inline select_scheduler::~select_scheduler()
208 : {
209 781 : if (pipe_fds_[0] >= 0)
210 781 : ::close(pipe_fds_[0]);
211 781 : if (pipe_fds_[1] >= 0)
212 781 : ::close(pipe_fds_[1]);
213 1562 : }
214 :
215 : inline void
216 781 : select_scheduler::shutdown()
217 : {
218 781 : shutdown_drain();
219 :
220 781 : if (pipe_fds_[1] >= 0)
221 781 : interrupt_reactor();
222 781 : }
223 :
224 : inline std::error_code
225 4558 : select_scheduler::register_descriptor(
226 : int fd, reactor_descriptor_state* desc) const
227 : {
228 4558 : if (fd < 0 || fd >= FD_SETSIZE)
229 1 : return make_err(EMFILE);
230 :
231 4557 : desc->registered_events = reactor_event_read | reactor_event_write;
232 4557 : desc->fd = fd;
233 4557 : desc->scheduler_ = this;
234 4557 : desc->mutex.set_enabled(reactor_io_locking_);
235 4557 : desc->ready_events_.store(0, std::memory_order_relaxed);
236 :
237 : {
238 4557 : conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
239 4557 : desc->impl_ref_.reset();
240 4557 : desc->read_ready = false;
241 4557 : desc->write_ready = false;
242 4557 : }
243 :
244 : {
245 4557 : mutex_type::scoped_lock lock(mutex_);
246 : try
247 : {
248 4557 : registered_descs_[fd] = desc;
249 : }
250 MIS 0 : catch (std::bad_alloc const&)
251 : {
252 0 : return make_err(ENOMEM);
253 0 : }
254 HIT 4557 : if (fd > max_fd_)
255 4542 : max_fd_ = fd;
256 4557 : }
257 :
258 4557 : interrupt_reactor();
259 4557 : return {};
260 : }
261 :
262 : inline void
263 4511 : select_scheduler::deregister_descriptor(int fd) const
264 : {
265 4511 : mutex_type::scoped_lock lock(mutex_);
266 :
267 4511 : auto it = registered_descs_.find(fd);
268 4511 : if (it == registered_descs_.end())
269 MIS 0 : return;
270 :
271 HIT 4511 : registered_descs_.erase(it);
272 :
273 4511 : if (fd == max_fd_)
274 : {
275 4320 : max_fd_ = pipe_fds_[0];
276 8315 : for (auto& [registered_fd, state] : registered_descs_)
277 : {
278 3995 : if (registered_fd > max_fd_)
279 3908 : max_fd_ = registered_fd;
280 : }
281 : }
282 4511 : }
283 :
284 : inline void
285 2082 : select_scheduler::notify_reactor() const
286 : {
287 2082 : interrupt_reactor();
288 2082 : }
289 :
290 : inline void
291 10890 : select_scheduler::interrupt_reactor() const
292 : {
293 10890 : char byte = 1;
294 10890 : [[maybe_unused]] auto r = ::write(pipe_fds_[1], &byte, 1);
295 10890 : }
296 :
297 : inline long
298 340423 : select_scheduler::calculate_timeout(long requested_timeout_us) const
299 : {
300 340423 : if (requested_timeout_us == 0)
301 MIS 0 : return 0;
302 :
303 HIT 340423 : auto nearest = timer_svc_->nearest_expiry();
304 340423 : if (nearest == timer_service::time_point::max())
305 656 : return requested_timeout_us;
306 :
307 339767 : auto now = std::chrono::steady_clock::now();
308 339767 : if (nearest <= now)
309 557 : return 0;
310 :
311 : auto timer_timeout_us =
312 339210 : std::chrono::duration_cast<std::chrono::microseconds>(nearest - now)
313 339210 : .count();
314 :
315 339210 : constexpr auto long_max =
316 : static_cast<long long>((std::numeric_limits<long>::max)());
317 : auto capped_timer_us =
318 339210 : (std::min)((std::max)(static_cast<long long>(timer_timeout_us),
319 339210 : static_cast<long long>(0)),
320 339210 : long_max);
321 :
322 339210 : if (requested_timeout_us < 0)
323 339210 : return static_cast<long>(capped_timer_us);
324 :
325 : return static_cast<long>(
326 MIS 0 : (std::min)(static_cast<long long>(requested_timeout_us),
327 0 : capped_timer_us));
328 : }
329 :
330 : inline void
331 HIT 363119 : select_scheduler::run_task(
332 : lock_type& lock, context_type* ctx, long timeout_us)
333 : {
334 : long effective_timeout_us =
335 363119 : task_interrupted_ ? 0 : calculate_timeout(timeout_us);
336 :
337 : // Snapshot registered descriptors while holding lock.
338 : // Record which fds need write monitoring to avoid a hot loop:
339 : // select is level-triggered so writable sockets (nearly always
340 : // writable) would cause select() to return immediately every
341 : // iteration if unconditionally added to write_fds. Membership
342 : // stays opt-in: a parked write wait opts in the same way a
343 : // parked write or connect op does.
344 : struct fd_entry
345 : {
346 : int fd;
347 : reactor_descriptor_state* desc;
348 : bool needs_write;
349 : };
350 : fd_entry snapshot[FD_SETSIZE];
351 363119 : int snapshot_count = 0;
352 :
353 911495 : for (auto& [fd, desc] : registered_descs_)
354 : {
355 548376 : if (snapshot_count < FD_SETSIZE)
356 : {
357 548376 : conditionally_enabled_mutex::scoped_lock desc_lock(desc->mutex);
358 548376 : snapshot[snapshot_count].fd = fd;
359 548376 : snapshot[snapshot_count].desc = desc;
360 548376 : snapshot[snapshot_count].needs_write =
361 1084227 : (desc->write_op || desc->connect_op ||
362 535851 : desc->wait_write_op);
363 548376 : ++snapshot_count;
364 548376 : }
365 : }
366 :
367 363119 : if (lock.owns_lock())
368 340424 : lock.unlock();
369 :
370 363119 : task_cleanup on_exit{this, &lock, ctx};
371 :
372 : fd_set read_fds, write_fds, except_fds;
373 6173023 : FD_ZERO(&read_fds);
374 6173023 : FD_ZERO(&write_fds);
375 6173023 : FD_ZERO(&except_fds);
376 :
377 363119 : FD_SET(pipe_fds_[0], &read_fds);
378 363119 : int nfds = pipe_fds_[0];
379 :
380 911495 : for (int i = 0; i < snapshot_count; ++i)
381 : {
382 548376 : int fd = snapshot[i].fd;
383 548376 : FD_SET(fd, &read_fds);
384 548376 : if (snapshot[i].needs_write)
385 12531 : FD_SET(fd, &write_fds);
386 548376 : FD_SET(fd, &except_fds);
387 548376 : if (fd > nfds)
388 362663 : nfds = fd;
389 : }
390 :
391 : struct timeval tv;
392 363119 : struct timeval* tv_ptr = nullptr;
393 363119 : if (effective_timeout_us >= 0)
394 : {
395 362467 : tv.tv_sec = effective_timeout_us / 1000000;
396 362467 : tv.tv_usec = effective_timeout_us % 1000000;
397 362467 : tv_ptr = &tv;
398 : }
399 :
400 363119 : int ready = ::select(nfds + 1, &read_fds, &write_fds, &except_fds, tv_ptr);
401 :
402 : // EINTR: signal interrupted select(), just retry.
403 : // EBADF: an fd was closed between snapshot and select(); retry
404 : // with a fresh snapshot from registered_descs_.
405 : // Both fall through with no ready descriptors rather than
406 : // returning: the caller handed this function an owned lock that
407 : // only the epilogue below re-acquires.
408 363119 : if (ready < 0)
409 : {
410 3 : if (errno != EINTR && errno != EBADF)
411 1 : detail::throw_system_error(make_err(errno), "select");
412 2 : ready = 0;
413 : }
414 :
415 : // Process timers outside the lock
416 363118 : timer_svc_->process_expired();
417 :
418 363118 : ready_queue local_ops;
419 :
420 363118 : if (ready > 0)
421 : {
422 347714 : if (FD_ISSET(pipe_fds_[0], &read_fds))
423 : {
424 : char buf[256];
425 9554 : while (::read(pipe_fds_[0], buf, sizeof(buf)) > 0)
426 : {
427 : }
428 : }
429 :
430 852306 : for (int i = 0; i < snapshot_count; ++i)
431 : {
432 504592 : int fd = snapshot[i].fd;
433 504592 : reactor_descriptor_state* desc = snapshot[i].desc;
434 :
435 504592 : std::uint32_t flags = 0;
436 504592 : if (FD_ISSET(fd, &read_fds))
437 345255 : flags |= reactor_event_read;
438 504592 : if (FD_ISSET(fd, &write_fds))
439 2077 : flags |= reactor_event_write;
440 504592 : if (FD_ISSET(fd, &except_fds))
441 16 : flags |= reactor_event_error;
442 :
443 504592 : if (flags == 0)
444 157269 : continue;
445 :
446 347323 : desc->add_ready_events(flags);
447 :
448 347323 : bool expected = false;
449 347323 : if (desc->is_enqueued_.compare_exchange_strong(
450 : expected, true, std::memory_order_release,
451 : std::memory_order_relaxed))
452 : {
453 347323 : local_ops.push(desc);
454 : }
455 : }
456 : }
457 :
458 363118 : lock.lock();
459 :
460 363118 : completed_ops_.splice(local_ops);
461 363119 : }
462 :
463 : } // namespace boost::corosio::detail
464 :
465 : #endif // BOOST_COROSIO_HAS_SELECT
466 :
467 : #endif // BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP
|