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

82.7% Lines (86/0/104) 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 500011x void add_ready_events(std::uint32_t ev) noexcept
116 {
117 500011x ready_events_.fetch_or(ev, std::memory_order_release);
118 500011x }
119
120 /// Invoke deferred I/O and dispatch completions.
121 499788x void operator()() override
122 {
123 499788x invoke_deferred_io();
124 499788x }
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 223x void destroy() override
130 {
131 223x impl_ref_.reset();
132 223x }
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 499788x reactor_descriptor_state::invoke_deferred_io()
145 {
146 499788x std::shared_ptr<void> prevent_impl_destruction;
147 499788x ready_queue local_ops;
148
149 {
150 499788x 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 499788x prevent_impl_destruction = std::move(impl_ref_);
158 499788x is_enqueued_.store(false, std::memory_order_release);
159
160 499788x std::uint32_t ev = ready_events_.exchange(0, std::memory_order_acquire);
161 499788x 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 499784x int err = 0;
169 499784x if (ev & reactor_event_error)
170 {
171 13x socklen_t len = sizeof(err);
172 13x if (::getsockopt(fd, SOL_SOCKET, SO_ERROR, &err, &len) < 0)
173 err = errno;
174 13x if (err == 0)
175 1x err = EIO;
176 }
177
178 499784x if (ev & reactor_event_read)
179 {
180 458450x if (read_op)
181 {
182 7853x auto* rd = read_op;
183 7853x if (err)
184 2x rd->complete(err, 0);
185 else
186 7851x rd->perform_io();
187
188 7853x if (rd->errn == EAGAIN || rd->errn == EWOULDBLOCK)
189 {
190 355x rd->errn = 0;
191 }
192 else
193 {
194 7498x read_op = nullptr;
195 7498x local_ops.push(rd);
196 }
197 }
198 else
199 {
200 450597x read_ready = true;
201 }
202
203 // The event does not prove the socket is still readable: a
204 // parked read op above may have drained it, or a speculative
205 // read consumed the data before this dispatch ran. The wait
206 // op's perform_io() re-probes and reports EAGAIN to stay
207 // parked.
208 458450x if (wait_read_op)
209 {
210 22x auto* wo = wait_read_op;
211 22x if (err)
212 wo->complete(err, 0);
213 else
214 22x wo->perform_io();
215
216 22x if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
217 {
218 4x wo->errn = 0;
219 }
220 else
221 {
222 18x wait_read_op = nullptr;
223 18x local_ops.push(wo);
224 }
225 }
226 }
227 499784x if (ev & reactor_event_write)
228 {
229 55919x bool had_write_op = (connect_op || write_op);
230 // A writable event on a socket still in SYN_SENT (e.g. the
231 // spurious pre-connect readiness of a fresh socket) must
232 // not complete the connect; perform_io() reports EAGAIN
233 // until a peer is actually established.
234 55919x if (connect_op)
235 {
236 7014x auto* cn = connect_op;
237 7014x if (err)
238 8x cn->complete(err, 0);
239 else
240 7006x cn->perform_io();
241
242 7014x if (cn->errn == EAGAIN || cn->errn == EWOULDBLOCK)
243 {
244 cn->errn = 0;
245 }
246 else
247 {
248 7014x connect_op = nullptr;
249 7014x local_ops.push(cn);
250 }
251 }
252 55919x if (write_op)
253 {
254 190x auto* wr = write_op;
255 190x if (err)
256 wr->complete(err, 0);
257 else
258 190x wr->perform_io();
259
260 190x if (wr->errn == EAGAIN || wr->errn == EWOULDBLOCK)
261 {
262 1x wr->errn = 0;
263 }
264 else
265 {
266 189x write_op = nullptr;
267 189x local_ops.push(wr);
268 }
269 }
270 55919x if (!had_write_op)
271 48715x write_ready = true;
272
273 // Same re-probe discipline as the wait-for-read dispatch.
274 55919x if (wait_write_op)
275 {
276 4x auto* wo = wait_write_op;
277 4x if (err)
278 wo->complete(err, 0);
279 else
280 4x wo->perform_io();
281
282 4x if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
283 {
284 wo->errn = 0;
285 }
286 else
287 {
288 4x wait_write_op = nullptr;
289 4x local_ops.push(wo);
290 }
291 }
292 }
293 // Complete a parked wait-for-error on any error condition.
294 499784x if ((ev & reactor_event_error) || err)
295 {
296 13x if (wait_error_op)
297 {
298 wait_error_op->complete(err, 0);
299 local_ops.push(std::exchange(wait_error_op, nullptr));
300 }
301 }
302 499784x if (err)
303 {
304 13x if (read_op)
305 {
306 read_op->complete(err, 0);
307 local_ops.push(std::exchange(read_op, nullptr));
308 }
309 13x if (write_op)
310 {
311 write_op->complete(err, 0);
312 local_ops.push(std::exchange(write_op, nullptr));
313 }
314 13x if (connect_op)
315 {
316 connect_op->complete(err, 0);
317 local_ops.push(std::exchange(connect_op, nullptr));
318 }
319 13x if (wait_read_op)
320 {
321 wait_read_op->complete(err, 0);
322 local_ops.push(std::exchange(wait_read_op, nullptr));
323 }
324 13x if (wait_write_op)
325 {
326 wait_write_op->complete(err, 0);
327 local_ops.push(std::exchange(wait_write_op, nullptr));
328 }
329 }
330 499788x }
331
332 // Execute first handler inline — the scheduler's work_cleanup
333 // accounts for this as the "consumed" work item. local_ops holds
334 // only ops, so the popped entry decodes directly.
335 499784x scheduler_op* first = ready_as_op(local_ops.pop());
336 499784x if (first)
337 {
338 14723x scheduler_->post_deferred_completions(local_ops);
339 14723x (*first)();
340 }
341 else
342 {
343 485061x scheduler_->compensating_work_started();
344 }
345 499788x }
346
347 } // namespace boost::corosio::detail
348
349 #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_STATE_HPP
350