include/boost/corosio/native/detail/reactor/reactor_descriptor_state.hpp

85.1% Lines (114/134) 100.0% List of functions (6/6) 75.0% Branches (66/88)
reactor_descriptor_state.hpp
f(x) Functions (6)
Line Branch 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_REACTOR_REACTOR_DESCRIPTOR_STATE_HPP
12 #define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_STATE_HPP
13
14 #include <boost/corosio/native/detail/reactor/reactor_events.hpp>
15 #include <boost/corosio/native/detail/reactor/reactor_op_base.hpp>
16 #include <boost/corosio/native/detail/reactor/reactor_scheduler.hpp>
17 #include <boost/corosio/detail/ready_queue.hpp>
18
19 #include <boost/corosio/detail/conditionally_enabled_mutex.hpp>
20
21 #include <atomic>
22 #include <cstdint>
23 #include <memory>
24
25 #include <errno.h>
26 #include <sys/socket.h>
27
28 namespace boost::corosio::detail {
29
30 /** Per-descriptor state shared across reactor backends.
31
32 Tracks pending operations for a file descriptor. The fd is registered
33 once with the reactor and stays registered until closed. Uses deferred
34 I/O: the reactor sets ready_events atomically, then enqueues this state.
35 When popped by the scheduler, invoke_deferred_io() performs I/O under
36 the mutex and queues completed ops.
37
38 Non-template: uses reactor_op_base pointers so the scheduler and
39 descriptor_state code exist as a single copy in the binary regardless
40 of how many backends are compiled in.
41
42 @par Thread Safety
43 The mutex protects operation pointers and ready flags. ready_events_
44 and is_enqueued_ are atomic for lock-free reactor access.
45 */
46 31446x struct reactor_descriptor_state : scheduler_op
47 {
48 /// Protects operation pointers and ready/cancel flags.
49 /// Becomes a no-op in single-threaded mode.
50 31446x conditionally_enabled_mutex mutex{true};
51
52 /// Pending read operation (guarded by `mutex`).
53 31446x reactor_op_base* read_op = nullptr;
54
55 /// Pending write operation (guarded by `mutex`).
56 31446x reactor_op_base* write_op = nullptr;
57
58 /// Pending connect operation (guarded by `mutex`).
59 31446x reactor_op_base* connect_op = nullptr;
60
61 /// Pending wait-for-read operation (guarded by `mutex`).
62 31446x reactor_op_base* wait_read_op = nullptr;
63
64 /// Pending wait-for-write operation (guarded by `mutex`).
65 31446x reactor_op_base* wait_write_op = nullptr;
66
67 /// Pending wait-for-error operation (guarded by `mutex`).
68 31446x reactor_op_base* wait_error_op = nullptr;
69
70 /// True if a read edge event arrived before an op was registered.
71 31446x bool read_ready = false;
72
73 /// True if a write edge event arrived before an op was registered.
74 31446x bool write_ready = false;
75
76 /// Event mask set during registration (no mutex needed).
77 31446x std::uint32_t registered_events = 0;
78
79 /// The reactor refused to watch this fd (e.g. /dev/null on epoll);
80 /// its I/O never blocks, and an op that would park must not.
81 31446x bool unpollable = false;
82
83 /// File descriptor this state tracks.
84 31446x int fd = -1;
85
86 /// Accumulated ready events (set by reactor, read by scheduler).
87 31446x std::atomic<std::uint32_t> ready_events_{0};
88
89 /// True while this state is queued in the scheduler's completed_ops.
90 31446x std::atomic<bool> is_enqueued_{false};
91
92 /// Owning scheduler for posting completions.
93 31446x reactor_scheduler const* scheduler_ = nullptr;
94
95 /// Prevents impl destruction while queued in the scheduler.
96 std::shared_ptr<void> impl_ref_;
97
98 /// Add ready events atomically.
99 /// Release pairs with the consumer's acquire exchange on
100 /// ready_events_ so the consumer sees all flags. On x86 (TSO)
101 /// this compiles to the same LOCK OR as relaxed.
102 124669x void add_ready_events(std::uint32_t ev) noexcept
103 {
104 124669x ready_events_.fetch_or(ev, std::memory_order_release);
105 124669x }
106
107 /// Invoke deferred I/O and dispatch completions.
108 124460x void operator()() override
109 {
110 124460x invoke_deferred_io();
111 124460x }
112
113 /// Destroy without invoking.
114 /// Called during scheduler::shutdown() drain. Clear impl_ref_ to break
115 /// the self-referential cycle set by close_socket().
116 134x void destroy() override
117 {
118 134x impl_ref_.reset();
119 134x }
120
121 /** Perform deferred I/O and queue completions.
122
123 Performs I/O under the mutex and queues completed ops. EAGAIN
124 ops stay parked in their slot for re-delivery on the next
125 edge event.
126 */
127 void invoke_deferred_io();
128 };
129
130 inline void
131 124460x reactor_descriptor_state::invoke_deferred_io()
132 {
133 124460x std::shared_ptr<void> prevent_impl_destruction;
134 124460x ready_queue local_ops;
135
136 {
137
1/2
✓ Branch 0 taken 124460 times.
✗ Branch 1 not taken.
124460x conditionally_enabled_mutex::scoped_lock lock(mutex);
138
139 // Must clear is_enqueued_ and move impl_ref_ under the same
140 // lock that processes I/O. close_socket() checks is_enqueued_
141 // under this mutex — without atomicity between the flag store
142 // and the ref move, close_socket() could see is_enqueued_==false,
143 // skip setting impl_ref_, and destroy the impl under us.
144 124460x prevent_impl_destruction = std::move(impl_ref_);
145 124460x is_enqueued_.store(false, std::memory_order_release);
146
147 124460x std::uint32_t ev = ready_events_.exchange(0, std::memory_order_acquire);
148
1/2
✓ Branch 0 taken 124460 times.
✗ Branch 1 not taken.
124460x if (ev == 0)
149 {
150 // Mutex unlocks here; compensate for work_cleanup's decrement
151 ✗ scheduler_->compensating_work_started();
152 ✗ return;
153 }
154
155 124460x int err = 0;
156
2/2
✓ Branch 0 taken 124415 times.
✓ Branch 1 taken 45 times.
124460x if (ev & reactor_event_error)
157 {
158 // Force the read/write dispatch below to run: an
159 // edge-triggered EPOLLERR can arrive alone, and without this
160 // a parked op never calls perform_io() and, the edge being
161 // one-shot, never gets another chance -- a permanent hang.
162 // Every parked op then re-runs its own syscall or probe.
163 //
164 // Assumes at least one parked op's own syscall makes
165 // non-EAGAIN progress; if every op re-parks with EAGAIN this
166 // sticky error is never redelivered and they hang. No such
167 // case is known -- a future descriptor type that hits one
168 // should be handled here.
169 45x ev |= reactor_event_read | reactor_event_write;
170
171 // SO_ERROR clears on read, so take it only for a parked op
172 // that reports it. A readiness wait reports readiness and
173 // leaves the error for the next read or write to name, as
174 // asio does; reading it here with nothing to report it to
175 // would turn a reset into a clean EOF.
176 45x bool const reports_error =
177
6/6
✓ Branch 0 taken 32 times.
✓ Branch 1 taken 13 times.
✓ Branch 2 taken 31 times.
✓ Branch 3 taken 1 time.
✓ Branch 4 taken 18 times.
✓ Branch 5 taken 13 times.
45x read_op || write_op || connect_op || wait_error_op;
178 45x socklen_t len = sizeof(err);
179
4/4
✓ Branch 0 taken 10 times.
✓ Branch 1 taken 35 times.
✓ Branch 2 taken 4 times.
✓ Branch 3 taken 31 times.
80x if (reports_error &&
180
1/2
✓ Branch 0 taken 35 times.
✗ Branch 1 not taken.
35x ::getsockopt(fd, SOL_SOCKET, SO_ERROR, &err, &len) < 0)
181 {
182 // Non-socket fd (pipe, chardev, ...): no SO_ERROR, so
183 // let the op's own syscall name the real failure.
184
3/6
✓ Branch 0 taken 4 times.
✗ Branch 1 not taken.
✗ Branch 2 not taken.
✓ Branch 3 taken 4 times.
✓ Branch 4 taken 4 times.
✗ Branch 5 not taken.
4x err = (errno == ENOTSOCK) ? 0 : errno;
185 4x }
186 // select raises its exceptional set for out-of-band/urgent
187 // data as well as for genuine faults; on a healthy socket the
188 // probe then reads SO_ERROR == 0. Faulting a pending read or
189 // write on that is wrong, so an I/O operation completes only
190 // on a real (non-zero) error. wait(error) still names a code
191 // below.
192 45x }
193
194
2/2
✓ Branch 0 taken 10144 times.
✓ Branch 1 taken 114316 times.
124460x if (ev & reactor_event_read)
195 {
196
2/2
✓ Branch 0 taken 7579 times.
✓ Branch 1 taken 106737 times.
114316x if (read_op)
197 {
198 106737x auto* rd = read_op;
199
2/2
✓ Branch 0 taken 12 times.
✓ Branch 1 taken 106725 times.
106737x if (err)
200 12x rd->complete(err, 0);
201 else
202 106725x rd->perform_io();
203
204
3/4
✓ Branch 0 taken 106171 times.
✓ Branch 1 taken 566 times.
✗ Branch 2 not taken.
✓ Branch 3 taken 106171 times.
106737x if (rd->errn == EAGAIN || rd->errn == EWOULDBLOCK)
205 {
206 566x rd->errn = 0;
207 566x }
208 else
209 {
210 106171x read_op = nullptr;
211 106171x local_ops.push(rd);
212 }
213 106737x }
214 else
215 {
216 7579x read_ready = true;
217 }
218
219 // The event does not prove the socket is still readable: a
220 // parked read op above may have drained it, or a speculative
221 // read consumed the data before this dispatch ran. The wait
222 // op's perform_io() re-probes and reports EAGAIN to stay
223 // parked.
224
2/2
✓ Branch 0 taken 114280 times.
✓ Branch 1 taken 36 times.
114316x if (wait_read_op)
225 {
226 36x auto* wo = wait_read_op;
227 36x wo->perform_io();
228
229
2/4
✓ Branch 0 taken 36 times.
✗ Branch 1 not taken.
✗ Branch 2 not taken.
✓ Branch 3 taken 36 times.
36x if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
230 {
231 ✗ wo->errn = 0;
232 ✗ }
233 else
234 {
235 36x wait_read_op = nullptr;
236 36x local_ops.push(wo);
237 }
238 36x }
239 114316x }
240
2/2
✓ Branch 0 taken 114221 times.
✓ Branch 1 taken 10239 times.
124460x if (ev & reactor_event_write)
241 {
242
2/2
✓ Branch 0 taken 2103 times.
✓ Branch 1 taken 8136 times.
10239x bool had_write_op = (connect_op || write_op);
243 // A writable event on a socket still in SYN_SENT (e.g. the
244 // spurious pre-connect readiness of a fresh socket) must
245 // not complete the connect; perform_io() reports EAGAIN
246 // until a peer is actually established.
247
2/2
✓ Branch 0 taken 2103 times.
✓ Branch 1 taken 8136 times.
10239x if (connect_op)
248 {
249 8136x auto* cn = connect_op;
250
2/2
✓ Branch 0 taken 18 times.
✓ Branch 1 taken 8118 times.
8136x if (err)
251 18x cn->complete(err, 0);
252 else
253 8118x cn->perform_io();
254
255
2/4
✓ Branch 0 taken 8136 times.
✗ Branch 1 not taken.
✗ Branch 2 not taken.
✓ Branch 3 taken 8136 times.
8136x if (cn->errn == EAGAIN || cn->errn == EWOULDBLOCK)
256 {
257 ✗ cn->errn = 0;
258 ✗ }
259 else
260 {
261 8136x connect_op = nullptr;
262 8136x local_ops.push(cn);
263 }
264 8136x }
265
2/2
✓ Branch 0 taken 9226 times.
✓ Branch 1 taken 1013 times.
10239x if (write_op)
266 {
267 1013x auto* wr = write_op;
268
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 1013 times.
1013x if (err)
269 ✗ wr->complete(err, 0);
270 else
271 1013x wr->perform_io();
272
273
2/4
✓ Branch 0 taken 1013 times.
✗ Branch 1 not taken.
✗ Branch 2 not taken.
✓ Branch 3 taken 1013 times.
1013x if (wr->errn == EAGAIN || wr->errn == EWOULDBLOCK)
274 {
275 ✗ wr->errn = 0;
276 ✗ }
277 else
278 {
279 1013x write_op = nullptr;
280 1013x local_ops.push(wr);
281 }
282 1013x }
283
2/2
✓ Branch 0 taken 1090 times.
✓ Branch 1 taken 9149 times.
10239x if (!had_write_op)
284 1090x write_ready = true;
285
286 // Same re-probe discipline as the wait-for-read dispatch.
287
2/2
✓ Branch 0 taken 10230 times.
✓ Branch 1 taken 9 times.
10239x if (wait_write_op)
288 {
289 9x auto* wo = wait_write_op;
290 9x wo->perform_io();
291
292
2/4
✓ Branch 0 taken 9 times.
✗ Branch 1 not taken.
✗ Branch 2 not taken.
✓ Branch 3 taken 9 times.
9x if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
293 {
294 ✗ wo->errn = 0;
295 ✗ }
296 else
297 {
298 9x wait_write_op = nullptr;
299 9x local_ops.push(wo);
300 }
301 9x }
302 10239x }
303 // Complete a parked wait-for-error on any error condition.
304
2/2
✓ Branch 0 taken 124415 times.
✓ Branch 1 taken 45 times.
124460x if (ev & reactor_event_error)
305 {
306
2/2
✓ Branch 0 taken 42 times.
✓ Branch 1 taken 3 times.
45x if (wait_error_op)
307 {
308 // wait(error) fired on the exceptional condition; name a
309 // code even when the kernel exposed none (e.g. urgent
310 // data leaves SO_ERROR == 0).
311
1/2
✓ Branch 0 taken 3 times.
✗ Branch 1 not taken.
3x int const werr = err ? err : EIO;
312 3x wait_error_op->complete(werr, 0);
313 3x local_ops.push(std::exchange(wait_error_op, nullptr));
314 3x }
315 45x }
316
2/2
✓ Branch 0 taken 124427 times.
✓ Branch 1 taken 33 times.
124460x if (err)
317 {
318
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 33 times.
33x if (read_op)
319 {
320 ✗ read_op->complete(err, 0);
321 ✗ local_ops.push(std::exchange(read_op, nullptr));
322 ✗ }
323
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 33 times.
33x if (write_op)
324 {
325 ✗ write_op->complete(err, 0);
326 ✗ local_ops.push(std::exchange(write_op, nullptr));
327 ✗ }
328
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 33 times.
33x if (connect_op)
329 {
330 ✗ connect_op->complete(err, 0);
331 ✗ local_ops.push(std::exchange(connect_op, nullptr));
332 ✗ }
333 33x }
334 124460x }
335
336 // Execute first handler inline — the scheduler's work_cleanup
337 // accounts for this as the "consumed" work item. local_ops holds
338 // only ops, so the popped entry decodes directly.
339 124460x scheduler_op* first = ready_as_op(local_ops.pop());
340
2/2
✓ Branch 0 taken 115366 times.
✓ Branch 1 taken 9094 times.
124460x if (first)
341 {
342
1/2
✓ Branch 0 taken 115366 times.
✗ Branch 1 not taken.
115366x scheduler_->post_deferred_completions(local_ops);
343
1/2
✓ Branch 0 taken 115366 times.
✗ Branch 1 not taken.
115366x (*first)();
344 115366x }
345 else
346 {
347 9094x scheduler_->compensating_work_started();
348 }
349 124460x }
350
351 } // namespace boost::corosio::detail
352
353 #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_STATE_HPP
354