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