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