include/boost/corosio/native/detail/select/select_scheduler.hpp

98.8% Lines (168/0/170) 100.0% List of functions (11/0/11)
select_scheduler.hpp
f(x) Functions (11)
Line TLA Hits 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 55x register_signal_reader(int read_fd) override
137 {
138 55x 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 881x inline select_scheduler::select_scheduler(capy::execution_context& ctx, int)
161 881x : pipe_fds_{-1, -1}
162 881x , max_fd_(-1)
163 {
164 881x if (::pipe(pipe_fds_) < 0)
165 1x detail::throw_system_error(make_err(errno), "pipe");
166
167 2631x for (int i = 0; i < 2; ++i)
168 {
169 1757x int flags = ::fcntl(pipe_fds_[i], F_GETFL, 0);
170 1757x if (flags == -1)
171 {
172 2x int errn = errno;
173 2x ::close(pipe_fds_[0]);
174 2x ::close(pipe_fds_[1]);
175 2x detail::throw_system_error(make_err(errn), "fcntl F_GETFL");
176 }
177 1755x if (::fcntl(pipe_fds_[i], F_SETFL, flags | O_NONBLOCK) == -1)
178 {
179 2x int errn = errno;
180 2x ::close(pipe_fds_[0]);
181 2x ::close(pipe_fds_[1]);
182 2x detail::throw_system_error(make_err(errn), "fcntl F_SETFL");
183 }
184 1753x if (::fcntl(pipe_fds_[i], F_SETFD, FD_CLOEXEC) == -1)
185 {
186 2x int errn = errno;
187 2x ::close(pipe_fds_[0]);
188 2x ::close(pipe_fds_[1]);
189 2x detail::throw_system_error(make_err(errn), "fcntl F_SETFD");
190 }
191 }
192
193 874x timer_svc_ = &get_timer_service(ctx, *this);
194 874x timer_svc_->set_on_earliest_changed(
195 3714x timer_service::callback(this, [](void* p) {
196 2840x static_cast<select_scheduler*>(p)->interrupt_reactor();
197 2840x }));
198
199 874x get_resolver_service(ctx, *this);
200 874x get_signal_service(ctx, *this);
201 874x get_stream_file_service(ctx, *this);
202 874x get_random_access_file_service(ctx, *this);
203
204 874x completed_ops_.push(&task_op_);
205 895x }
206
207 1748x inline select_scheduler::~select_scheduler()
208 {
209 874x if (pipe_fds_[0] >= 0)
210 874x ::close(pipe_fds_[0]);
211 874x if (pipe_fds_[1] >= 0)
212 874x ::close(pipe_fds_[1]);
213 1748x }
214
215 inline void
216 874x select_scheduler::shutdown()
217 {
218 874x shutdown_drain();
219
220 874x if (pipe_fds_[1] >= 0)
221 874x interrupt_reactor();
222 874x }
223
224 inline std::error_code
225 5234x select_scheduler::register_descriptor(
226 int fd, reactor_descriptor_state* desc) const
227 {
228 5234x if (fd < 0 || fd >= FD_SETSIZE)
229 1x return make_err(EMFILE);
230
231 5233x desc->registered_events = reactor_event_read | reactor_event_write;
232 5233x desc->fd = fd;
233 5233x desc->scheduler_ = this;
234 5233x desc->mutex.set_enabled(reactor_io_locking_);
235 5233x desc->ready_events_.store(0, std::memory_order_relaxed);
236
237 {
238 5233x conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
239 5233x desc->impl_ref_.reset();
240 5233x desc->read_ready = false;
241 5233x desc->write_ready = false;
242 5233x }
243
244 {
245 5233x mutex_type::scoped_lock lock(mutex_);
246 try
247 {
248 5233x registered_descs_[fd] = desc;
249 }
250 1x catch (std::bad_alloc const&)
251 {
252 1x return make_err(ENOMEM);
253 1x }
254 5232x if (fd > max_fd_)
255 5178x max_fd_ = fd;
256 5233x }
257
258 5232x interrupt_reactor();
259 5232x return {};
260 }
261
262 inline void
263 5178x select_scheduler::deregister_descriptor(int fd) const
264 {
265 5178x mutex_type::scoped_lock lock(mutex_);
266
267 5178x auto it = registered_descs_.find(fd);
268 5178x if (it == registered_descs_.end())
269 return;
270
271 5178x registered_descs_.erase(it);
272
273 5178x if (fd == max_fd_)
274 {
275 4841x max_fd_ = pipe_fds_[0];
276 9283x for (auto& [registered_fd, state] : registered_descs_)
277 {
278 4442x if (registered_fd > max_fd_)
279 4349x max_fd_ = registered_fd;
280 }
281 }
282 5178x }
283
284 inline void
285 2351x select_scheduler::notify_reactor() const
286 {
287 2351x interrupt_reactor();
288 2351x }
289
290 inline void
291 12090x select_scheduler::interrupt_reactor() const
292 {
293 12090x char byte = 1;
294 12090x [[maybe_unused]] auto r = ::write(pipe_fds_[1], &byte, 1);
295 12090x }
296
297 inline long
298 331843x select_scheduler::calculate_timeout(long requested_timeout_us) const
299 {
300 331843x if (requested_timeout_us == 0)
301 return 0;
302
303 331843x auto nearest = timer_svc_->nearest_expiry();
304 331843x if (nearest == timer_service::time_point::max())
305 733x return requested_timeout_us;
306
307 331110x auto now = std::chrono::steady_clock::now();
308 331110x if (nearest <= now)
309 519x return 0;
310
311 auto timer_timeout_us =
312 330591x std::chrono::duration_cast<std::chrono::microseconds>(nearest - now)
313 330591x .count();
314
315 330591x constexpr auto long_max =
316 static_cast<long long>((std::numeric_limits<long>::max)());
317 auto capped_timer_us =
318 330591x (std::min)((std::max)(static_cast<long long>(timer_timeout_us),
319 330591x static_cast<long long>(0)),
320 330591x long_max);
321
322 330591x if (requested_timeout_us < 0)
323 330589x return static_cast<long>(capped_timer_us);
324
325 return static_cast<long>(
326 2x (std::min)(static_cast<long long>(requested_timeout_us),
327 2x capped_timer_us));
328 }
329
330 inline void
331 356670x select_scheduler::run_task(
332 lock_type& lock, context_type& ctx, long timeout_us)
333 {
334 long effective_timeout_us =
335 356670x 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 356670x int snapshot_count = 0;
352
353 912560x for (auto& [fd, desc] : registered_descs_)
354 {
355 555890x if (snapshot_count < FD_SETSIZE)
356 {
357 555890x conditionally_enabled_mutex::scoped_lock desc_lock(desc->mutex);
358 555890x snapshot[snapshot_count].fd = fd;
359 555890x snapshot[snapshot_count].desc = desc;
360 555890x snapshot[snapshot_count].needs_write =
361 1098709x (desc->write_op || desc->connect_op ||
362 542819x desc->wait_write_op);
363 555890x ++snapshot_count;
364 555890x }
365 }
366
367 356670x if (lock.owns_lock())
368 331844x lock.unlock();
369
370 356670x task_cleanup on_exit{this, &lock, ctx};
371
372 fd_set read_fds, write_fds, except_fds;
373 6063390x FD_ZERO(&read_fds);
374 6063390x FD_ZERO(&write_fds);
375 6063390x FD_ZERO(&except_fds);
376
377 356670x FD_SET(pipe_fds_[0], &read_fds);
378 356670x int nfds = pipe_fds_[0];
379
380 912560x for (int i = 0; i < snapshot_count; ++i)
381 {
382 555890x int fd = snapshot[i].fd;
383 555890x FD_SET(fd, &read_fds);
384 555890x if (snapshot[i].needs_write)
385 13077x FD_SET(fd, &write_fds);
386 555890x FD_SET(fd, &except_fds);
387 555890x if (fd > nfds)
388 356273x nfds = fd;
389 }
390
391 struct timeval tv;
392 356670x struct timeval* tv_ptr = nullptr;
393 356670x if (effective_timeout_us >= 0)
394 {
395 355953x tv.tv_sec = effective_timeout_us / 1000000;
396 355953x tv.tv_usec = effective_timeout_us % 1000000;
397 355953x tv_ptr = &tv;
398 }
399
400 356670x 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 356670x if (ready < 0)
409 {
410 3x if (errno != EINTR && errno != EBADF)
411 1x detail::throw_system_error(make_err(errno), "select");
412 2x ready = 0;
413 }
414
415 // Process timers outside the lock
416 356669x timer_svc_->process_expired();
417
418 356669x ready_queue local_ops;
419
420 356669x if (ready > 0)
421 {
422 339472x if (FD_ISSET(pipe_fds_[0], &read_fds))
423 {
424 char buf[256];
425 10808x while (::read(pipe_fds_[0], buf, sizeof(buf)) > 0)
426 {
427 }
428 }
429
430 843531x for (int i = 0; i < snapshot_count; ++i)
431 {
432 504059x int fd = snapshot[i].fd;
433 504059x reactor_descriptor_state* desc = snapshot[i].desc;
434
435 504059x std::uint32_t flags = 0;
436 504059x if (FD_ISSET(fd, &read_fds))
437 338797x flags |= reactor_event_read;
438 504059x if (FD_ISSET(fd, &write_fds))
439 2343x flags |= reactor_event_write;
440 504059x if (FD_ISSET(fd, &except_fds))
441 16x flags |= reactor_event_error;
442
443 504059x if (flags == 0)
444 162928x continue;
445
446 341131x desc->add_ready_events(flags);
447
448 341131x bool expected = false;
449 341131x if (desc->is_enqueued_.compare_exchange_strong(
450 expected, true, std::memory_order_release,
451 std::memory_order_relaxed))
452 {
453 341131x local_ops.push(desc);
454 }
455 }
456 }
457
458 356669x lock.lock();
459
460 356669x completed_ops_.splice(local_ops);
461 356670x }
462
463 } // namespace boost::corosio::detail
464
465 #endif // BOOST_COROSIO_HAS_SELECT
466
467 #endif // BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP
468