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

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