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

99.4% Lines (163 / 164, 1 excl) 100.0% Functions (11 / 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
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 /** Deregister a persistently registered descriptor.
117
118 @param fd The file descriptor to deregister.
119 */
120 void deregister_descriptor(int fd) const;
121
122 /** Interrupt the reactor so it rebuilds its fd_sets.
123
124 Called when a write, connect, or write-wait op is registered
125 after the reactor's snapshot was taken. Without this,
126 select() may block not watching for writability on the fd.
127 */
128 void notify_reactor() const;
129
130 /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp).
131 61x [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
132 {
133 61x return register_descriptor(read_fd, signal_pipe_reader_.arm());
134 }
135
136 private:
137 void run_task(lock_type& lock, context_type& ctx, long timeout_us) override;
138 void interrupt_reactor() const override;
139 long calculate_timeout(long requested_timeout_us) const;
140
141 // Watches the global signal self-pipe's read end (armed lazily by
142 // register_signal_reader on the first signal registration).
143 reactor_signal_pipe_reader signal_pipe_reader_;
144
145 // Self-pipe for interrupting select()
146 int pipe_fds_[2]; // [0]=read, [1]=write
147
148 // Per-fd tracking for fd_set building
149 mutable std::unordered_map<int, reactor_descriptor_state*>
150 registered_descs_;
151 mutable int max_fd_ = -1;
152 };
153
154 950x inline select_scheduler::select_scheduler(capy::execution_context& ctx, int)
155 950x : pipe_fds_{-1, -1}
156 950x , max_fd_(-1)
157 {
158 950x if (::pipe(pipe_fds_) < 0)
159 1x detail::throw_system_error(make_err(errno), "pipe");
160
161 2838x for (int i = 0; i < 2; ++i)
162 {
163 1895x int flags = ::fcntl(pipe_fds_[i], F_GETFL, 0);
164 1895x if (flags == -1)
165 {
166 2x int errn = errno;
167 2x ::close(pipe_fds_[0]);
168 2x ::close(pipe_fds_[1]);
169 2x detail::throw_system_error(make_err(errn), "fcntl F_GETFL");
170 }
171 1893x if (::fcntl(pipe_fds_[i], F_SETFL, flags | O_NONBLOCK) == -1)
172 {
173 2x int errn = errno;
174 2x ::close(pipe_fds_[0]);
175 2x ::close(pipe_fds_[1]);
176 2x detail::throw_system_error(make_err(errn), "fcntl F_SETFL");
177 }
178 1891x if (::fcntl(pipe_fds_[i], F_SETFD, FD_CLOEXEC) == -1)
179 {
180 2x int errn = errno;
181 2x ::close(pipe_fds_[0]);
182 2x ::close(pipe_fds_[1]);
183 2x detail::throw_system_error(make_err(errn), "fcntl F_SETFD");
184 }
185 }
186
187 943x timer_svc_ = &get_timer_service(ctx, *this);
188 943x timer_svc_->set_on_earliest_changed(
189 3212x timer_service::callback(this, [](void* p) {
190 2269x static_cast<select_scheduler*>(p)->interrupt_reactor();
191 2269x }));
192
193 943x completed_ops_.push(&task_op_);
194 964x }
195
196 1886x inline select_scheduler::~select_scheduler()
197 {
198 943x if (pipe_fds_[0] >= 0)
199 943x ::close(pipe_fds_[0]);
200 943x if (pipe_fds_[1] >= 0)
201 943x ::close(pipe_fds_[1]);
202 1886x }
203
204 inline void
205 943x select_scheduler::shutdown()
206 {
207 943x shutdown_drain();
208
209 943x if (pipe_fds_[1] >= 0)
210 943x interrupt_reactor();
211 943x }
212
213 inline std::error_code
214 4162x select_scheduler::register_descriptor(
215 int fd, reactor_descriptor_state* desc) const
216 {
217 4162x if (fd < 0 || fd >= FD_SETSIZE)
218 1x return make_err(EMFILE);
219
220 4161x desc->registered_events = reactor_event_read | reactor_event_write;
221 4161x desc->fd = fd;
222 4161x desc->scheduler_ = this;
223 4161x desc->mutex.set_enabled(reactor_io_locking_);
224 4161x desc->ready_events_.store(0, std::memory_order_relaxed);
225
226 {
227 4161x conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
228 4161x desc->impl_ref_.reset();
229 4161x desc->read_ready = false;
230 4161x desc->write_ready = false;
231 4161x }
232
233 {
234 4161x mutex_type::scoped_lock lock(mutex_);
235 try
236 {
237 4161x registered_descs_[fd] = desc;
238 }
239 1x catch (std::bad_alloc const&)
240 {
241 1x return make_err(ENOMEM);
242 1x }
243 4160x if (fd > max_fd_)
244 4110x max_fd_ = fd;
245 4161x }
246
247 4160x interrupt_reactor();
248 4160x return {};
249 }
250
251 inline void
252 4100x select_scheduler::deregister_descriptor(int fd) const
253 {
254 4100x mutex_type::scoped_lock lock(mutex_);
255
256 4100x auto it = registered_descs_.find(fd);
257 4100x if (it == registered_descs_.end())
258 ✗ return;
259
260 4100x registered_descs_.erase(it);
261
262 4100x if (fd == max_fd_)
263 {
264 3765x max_fd_ = pipe_fds_[0];
265 7124x for (auto& [registered_fd, state] : registered_descs_)
266 {
267 3359x if (registered_fd > max_fd_)
268 3266x max_fd_ = registered_fd;
269 }
270 }
271 4100x }
272
273 inline void
274 1809x select_scheduler::notify_reactor() const
275 {
276 1809x interrupt_reactor();
277 1809x }
278
279 inline void
280 10841x select_scheduler::interrupt_reactor() const
281 {
282 10841x char byte = 1;
283 10841x [[maybe_unused]] auto r = ::write(pipe_fds_[1], &byte, 1);
284 10841x }
285
286 inline long
287 132022x select_scheduler::calculate_timeout(long requested_timeout_us) const
288 {
289 132022x if (requested_timeout_us == 0)
290 − return 0; // LCOV_EXCL_LINE run_task passes 0 via task_interrupted_, never through this argument
291
292 132022x auto nearest = timer_svc_->nearest_expiry();
293 132022x if (nearest == timer_service::time_point::max())
294 1091x return requested_timeout_us;
295
296 130931x auto now = std::chrono::steady_clock::now();
297 130931x if (nearest <= now)
298 484x return 0;
299
300 auto timer_timeout_us =
301 130447x std::chrono::duration_cast<std::chrono::microseconds>(nearest - now)
302 130447x .count();
303
304 130447x constexpr auto long_max =
305 static_cast<long long>((std::numeric_limits<long>::max)());
306 auto capped_timer_us =
307 130447x (std::min)((std::max)(static_cast<long long>(timer_timeout_us),
308 130447x static_cast<long long>(0)),
309 130447x long_max);
310
311 130447x if (requested_timeout_us < 0)
312 130445x return static_cast<long>(capped_timer_us);
313
314 return static_cast<long>(
315 2x (std::min)(static_cast<long long>(requested_timeout_us),
316 2x capped_timer_us));
317 }
318
319 inline void
320 142811x select_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
321 {
322 long effective_timeout_us =
323 142811x task_interrupted_ ? 0 : calculate_timeout(timeout_us);
324
325 // Snapshot registered descriptors while holding lock.
326 // Record which fds need write monitoring to avoid a hot loop:
327 // select is level-triggered so writable sockets (nearly always
328 // writable) would cause select() to return immediately every
329 // iteration if unconditionally added to write_fds. Membership
330 // stays opt-in: a parked write wait opts in the same way a
331 // parked write or connect op does.
332 struct fd_entry
333 {
334 int fd;
335 reactor_descriptor_state* desc;
336 bool needs_write;
337 };
338 fd_entry snapshot[FD_SETSIZE];
339 142811x int snapshot_count = 0;
340
341 356310x for (auto& [fd, desc] : registered_descs_)
342 {
343 213499x if (snapshot_count < FD_SETSIZE)
344 {
345 213499x conditionally_enabled_mutex::scoped_lock desc_lock(desc->mutex);
346 213499x snapshot[snapshot_count].fd = fd;
347 213499x snapshot[snapshot_count].desc = desc;
348 213499x snapshot[snapshot_count].needs_write =
349 213499x (desc->write_op || desc->connect_op || desc->wait_write_op);
350 213499x ++snapshot_count;
351 213499x }
352 }
353
354 142811x if (lock.owns_lock())
355 132023x lock.unlock();
356
357 142811x task_cleanup on_exit{this, &lock, ctx};
358
359 fd_set read_fds, write_fds, except_fds;
360 2427787x FD_ZERO(&read_fds);
361 2427787x FD_ZERO(&write_fds);
362 2427787x FD_ZERO(&except_fds);
363
364 142811x FD_SET(pipe_fds_[0], &read_fds);
365 142811x int nfds = pipe_fds_[0];
366
367 356310x for (int i = 0; i < snapshot_count; ++i)
368 {
369 213499x int fd = snapshot[i].fd;
370 213499x FD_SET(fd, &read_fds);
371 213499x if (snapshot[i].needs_write)
372 5204x FD_SET(fd, &write_fds);
373 213499x FD_SET(fd, &except_fds);
374 213499x if (fd > nfds)
375 142125x nfds = fd;
376 }
377
378 struct timeval tv;
379 142811x struct timeval* tv_ptr = nullptr;
380 142811x if (effective_timeout_us >= 0)
381 {
382 141994x tv.tv_sec = effective_timeout_us / 1000000;
383 141994x tv.tv_usec = effective_timeout_us % 1000000;
384 141994x tv_ptr = &tv;
385 }
386
387 142811x int ready = ::select(nfds + 1, &read_fds, &write_fds, &except_fds, tv_ptr);
388
389 // EINTR: signal interrupted select(), just retry.
390 // EBADF: an fd was closed between snapshot and select(); retry
391 // with a fresh snapshot from registered_descs_.
392 // Both fall through with no ready descriptors rather than
393 // returning: the caller handed this function an owned lock that
394 // only the epilogue below re-acquires.
395 142811x if (ready < 0)
396 {
397 3x if (errno != EINTR && errno != EBADF)
398 1x detail::throw_system_error(make_err(errno), "select");
399 2x ready = 0;
400 }
401
402 // Process timers outside the lock
403 142810x timer_svc_->process_expired();
404
405 142810x ready_queue local_ops;
406
407 142810x if (ready > 0)
408 {
409 135167x if (FD_ISSET(pipe_fds_[0], &read_fds))
410 {
411 char buf[256];
412 10116x while (::read(pipe_fds_[0], buf, sizeof(buf)) > 0)
413 {
414 }
415 }
416
417 333947x for (int i = 0; i < snapshot_count; ++i)
418 {
419 198780x int fd = snapshot[i].fd;
420 198780x reactor_descriptor_state* desc = snapshot[i].desc;
421
422 198780x std::uint32_t flags = 0;
423 198780x if (FD_ISSET(fd, &read_fds))
424 132628x flags |= reactor_event_read;
425 198780x if (FD_ISSET(fd, &write_fds))
426 1801x flags |= reactor_event_write;
427 198780x if (FD_ISSET(fd, &except_fds))
428 16x flags |= reactor_event_error;
429
430 198780x if (flags == 0)
431 64360x continue;
432
433 134420x desc->add_ready_events(flags);
434
435 134420x bool expected = false;
436 134420x if (desc->is_enqueued_.compare_exchange_strong(
437 expected, true, std::memory_order_release,
438 std::memory_order_relaxed))
439 {
440 134420x local_ops.push(desc);
441 }
442 }
443 }
444
445 142810x lock.lock();
446
447 142810x completed_ops_.splice(local_ops);
448 142811x }
449
450 } // namespace boost::corosio::detail
451
452 #endif // BOOST_COROSIO_HAS_SELECT
453
454 #endif // BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP
455