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

86.8% Lines (145/167) 100.0% List of 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 #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 41x void register_signal_reader(int read_fd) override
131 {
132 41x register_descriptor(read_fd, signal_pipe_reader_.arm());
133 41x }
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 585x inline select_scheduler::select_scheduler(capy::execution_context& ctx, int)
155 585x : pipe_fds_{-1, -1}
156 585x , max_fd_(-1)
157 {
158 585x if (::pipe(pipe_fds_) < 0)
159 detail::throw_system_error(make_err(errno), "pipe");
160
161 1755x for (int i = 0; i < 2; ++i)
162 {
163 1170x int flags = ::fcntl(pipe_fds_[i], F_GETFL, 0);
164 1170x if (flags == -1)
165 {
166 int errn = errno;
167 ::close(pipe_fds_[0]);
168 ::close(pipe_fds_[1]);
169 detail::throw_system_error(make_err(errn), "fcntl F_GETFL");
170 }
171 1170x if (::fcntl(pipe_fds_[i], F_SETFL, flags | O_NONBLOCK) == -1)
172 {
173 int errn = errno;
174 ::close(pipe_fds_[0]);
175 ::close(pipe_fds_[1]);
176 detail::throw_system_error(make_err(errn), "fcntl F_SETFL");
177 }
178 1170x if (::fcntl(pipe_fds_[i], F_SETFD, FD_CLOEXEC) == -1)
179 {
180 int errn = errno;
181 ::close(pipe_fds_[0]);
182 ::close(pipe_fds_[1]);
183 detail::throw_system_error(make_err(errn), "fcntl F_SETFD");
184 }
185 }
186
187 585x timer_svc_ = &get_timer_service(ctx, *this);
188 585x timer_svc_->set_on_earliest_changed(
189 4423x timer_service::callback(this, [](void* p) {
190 3838x static_cast<select_scheduler*>(p)->interrupt_reactor();
191 3838x }));
192
193 585x get_resolver_service(ctx, *this);
194 585x get_signal_service(ctx, *this);
195 585x get_stream_file_service(ctx, *this);
196 585x get_random_access_file_service(ctx, *this);
197
198 585x completed_ops_.push(&task_op_);
199 585x }
200
201 1170x inline select_scheduler::~select_scheduler()
202 {
203 585x if (pipe_fds_[0] >= 0)
204 585x ::close(pipe_fds_[0]);
205 585x if (pipe_fds_[1] >= 0)
206 585x ::close(pipe_fds_[1]);
207 1170x }
208
209 inline void
210 585x select_scheduler::shutdown()
211 {
212 585x shutdown_drain();
213
214 585x if (pipe_fds_[1] >= 0)
215 585x interrupt_reactor();
216 585x }
217
218 inline void
219 6630x select_scheduler::register_descriptor(
220 int fd, reactor_descriptor_state* desc) const
221 {
222 6630x if (fd < 0 || fd >= FD_SETSIZE)
223 detail::throw_system_error(make_err(EINVAL), "select: fd out of range");
224
225 6630x desc->registered_events = reactor_event_read | reactor_event_write;
226 6630x desc->fd = fd;
227 6630x desc->scheduler_ = this;
228 6630x desc->mutex.set_enabled(reactor_io_locking_);
229 6630x desc->ready_events_.store(0, std::memory_order_relaxed);
230
231 {
232 6630x conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
233 6630x desc->impl_ref_.reset();
234 6630x desc->read_ready = false;
235 6630x desc->write_ready = false;
236 6630x }
237
238 {
239 6630x mutex_type::scoped_lock lock(mutex_);
240 6630x registered_descs_[fd] = desc;
241 6630x if (fd > max_fd_)
242 6625x max_fd_ = fd;
243 6630x }
244
245 6630x interrupt_reactor();
246 6630x }
247
248 inline void
249 6589x select_scheduler::deregister_descriptor(int fd) const
250 {
251 6589x mutex_type::scoped_lock lock(mutex_);
252
253 6589x auto it = registered_descs_.find(fd);
254 6589x if (it == registered_descs_.end())
255 return;
256
257 6589x registered_descs_.erase(it);
258
259 6589x if (fd == max_fd_)
260 {
261 6469x max_fd_ = pipe_fds_[0];
262 12685x for (auto& [registered_fd, state] : registered_descs_)
263 {
264 6216x if (registered_fd > max_fd_)
265 6169x max_fd_ = registered_fd;
266 }
267 }
268 6589x }
269
270 inline void
271 3198x select_scheduler::notify_reactor() const
272 {
273 3198x interrupt_reactor();
274 3198x }
275
276 inline void
277 14698x select_scheduler::interrupt_reactor() const
278 {
279 14698x char byte = 1;
280 14698x [[maybe_unused]] auto r = ::write(pipe_fds_[1], &byte, 1);
281 14698x }
282
283 inline long
284 407093x select_scheduler::calculate_timeout(long requested_timeout_us) const
285 {
286 407093x if (requested_timeout_us == 0)
287 return 0;
288
289 407093x auto nearest = timer_svc_->nearest_expiry();
290 407093x if (nearest == timer_service::time_point::max())
291 554x return requested_timeout_us;
292
293 406539x auto now = std::chrono::steady_clock::now();
294 406539x if (nearest <= now)
295 664x return 0;
296
297 auto timer_timeout_us =
298 405875x std::chrono::duration_cast<std::chrono::microseconds>(nearest - now)
299 405875x .count();
300
301 405875x constexpr auto long_max =
302 static_cast<long long>((std::numeric_limits<long>::max)());
303 auto capped_timer_us =
304 405875x (std::min)((std::max)(static_cast<long long>(timer_timeout_us),
305 405875x static_cast<long long>(0)),
306 405875x long_max);
307
308 405875x if (requested_timeout_us < 0)
309 405875x return static_cast<long>(capped_timer_us);
310
311 return static_cast<long>(
312 (std::min)(static_cast<long long>(requested_timeout_us),
313 capped_timer_us));
314 }
315
316 inline void
317 434839x select_scheduler::run_task(
318 lock_type& lock, context_type* ctx, long timeout_us)
319 {
320 long effective_timeout_us =
321 434839x 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 434839x int snapshot_count = 0;
336
337 1136675x for (auto& [fd, desc] : registered_descs_)
338 {
339 701836x if (snapshot_count < FD_SETSIZE)
340 {
341 701836x conditionally_enabled_mutex::scoped_lock desc_lock(desc->mutex);
342 701836x snapshot[snapshot_count].fd = fd;
343 701836x snapshot[snapshot_count].desc = desc;
344 701836x snapshot[snapshot_count].needs_write =
345 701836x (desc->write_op || desc->connect_op);
346 701836x ++snapshot_count;
347 701836x }
348 }
349
350 434839x if (lock.owns_lock())
351 407094x lock.unlock();
352
353 434839x task_cleanup on_exit{this, &lock, ctx};
354
355 fd_set read_fds, write_fds, except_fds;
356 7392263x FD_ZERO(&read_fds);
357 7392263x FD_ZERO(&write_fds);
358 7392263x FD_ZERO(&except_fds);
359
360 434839x FD_SET(pipe_fds_[0], &read_fds);
361 434839x int nfds = pipe_fds_[0];
362
363 1136675x for (int i = 0; i < snapshot_count; ++i)
364 {
365 701836x int fd = snapshot[i].fd;
366 701836x FD_SET(fd, &read_fds);
367 701836x if (snapshot[i].needs_write)
368 20256x FD_SET(fd, &write_fds);
369 701836x FD_SET(fd, &except_fds);
370 701836x if (fd > nfds)
371 434525x nfds = fd;
372 }
373
374 struct timeval tv;
375 434839x struct timeval* tv_ptr = nullptr;
376 434839x if (effective_timeout_us >= 0)
377 {
378 434289x tv.tv_sec = effective_timeout_us / 1000000;
379 434289x tv.tv_usec = effective_timeout_us % 1000000;
380 434289x tv_ptr = &tv;
381 }
382
383 434839x 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 434839x if (ready < 0)
389 {
390 if (errno == EINTR || errno == EBADF)
391 return;
392 detail::throw_system_error(make_err(errno), "select");
393 }
394
395 // Process timers outside the lock
396 434839x timer_svc_->process_expired();
397
398 434839x ready_queue local_ops;
399
400 434839x if (ready > 0)
401 {
402 414445x if (FD_ISSET(pipe_fds_[0], &read_fds))
403 {
404 char buf[256];
405 13740x while (::read(pipe_fds_[0], buf, sizeof(buf)) > 0)
406 {
407 }
408 }
409
410 1057045x for (int i = 0; i < snapshot_count; ++i)
411 {
412 642600x int fd = snapshot[i].fd;
413 642600x reactor_descriptor_state* desc = snapshot[i].desc;
414
415 642600x std::uint32_t flags = 0;
416 642600x if (FD_ISSET(fd, &read_fds))
417 410926x flags |= reactor_event_read;
418 642600x if (FD_ISSET(fd, &write_fds))
419 3195x flags |= reactor_event_write;
420 642600x if (FD_ISSET(fd, &except_fds))
421 flags |= reactor_event_error;
422
423 642600x if (flags == 0)
424 228490x continue;
425
426 414110x desc->add_ready_events(flags);
427
428 414110x bool expected = false;
429 414110x if (desc->is_enqueued_.compare_exchange_strong(
430 expected, true, std::memory_order_release,
431 std::memory_order_relaxed))
432 {
433 414110x local_ops.push(desc);
434 }
435 }
436 }
437
438 434839x lock.lock();
439
440 434839x completed_ops_.splice(local_ops);
441 434839x }
442
443 } // namespace boost::corosio::detail
444
445 #endif // BOOST_COROSIO_HAS_SELECT
446
447 #endif // BOOST_COROSIO_NATIVE_DETAIL_SELECT_SELECT_SCHEDULER_HPP
448