TLA Line data 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 HIT 500011 : void add_ready_events(std::uint32_t ev) noexcept
116 : {
117 500011 : ready_events_.fetch_or(ev, std::memory_order_release);
118 500011 : }
119 :
120 : /// Invoke deferred I/O and dispatch completions.
121 499788 : void operator()() override
122 : {
123 499788 : invoke_deferred_io();
124 499788 : }
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 223 : void destroy() override
130 : {
131 223 : impl_ref_.reset();
132 223 : }
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 499788 : reactor_descriptor_state::invoke_deferred_io()
145 : {
146 499788 : std::shared_ptr<void> prevent_impl_destruction;
147 499788 : ready_queue local_ops;
148 :
149 : {
150 499788 : 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 499788 : prevent_impl_destruction = std::move(impl_ref_);
158 499788 : is_enqueued_.store(false, std::memory_order_release);
159 :
160 499788 : std::uint32_t ev = ready_events_.exchange(0, std::memory_order_acquire);
161 499788 : if (ev == 0)
162 : {
163 : // Mutex unlocks here; compensate for work_cleanup's decrement
164 4 : scheduler_->compensating_work_started();
165 4 : return;
166 : }
167 :
168 499784 : int err = 0;
169 499784 : if (ev & reactor_event_error)
170 : {
171 13 : socklen_t len = sizeof(err);
172 13 : if (::getsockopt(fd, SOL_SOCKET, SO_ERROR, &err, &len) < 0)
173 MIS 0 : err = errno;
174 HIT 13 : if (err == 0)
175 1 : err = EIO;
176 : }
177 :
178 499784 : if (ev & reactor_event_read)
179 : {
180 458450 : if (read_op)
181 : {
182 7853 : auto* rd = read_op;
183 7853 : if (err)
184 2 : rd->complete(err, 0);
185 : else
186 7851 : rd->perform_io();
187 :
188 7853 : if (rd->errn == EAGAIN || rd->errn == EWOULDBLOCK)
189 : {
190 355 : rd->errn = 0;
191 : }
192 : else
193 : {
194 7498 : read_op = nullptr;
195 7498 : local_ops.push(rd);
196 : }
197 : }
198 : else
199 : {
200 450597 : 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 458450 : if (wait_read_op)
209 : {
210 22 : auto* wo = wait_read_op;
211 22 : if (err)
212 MIS 0 : wo->complete(err, 0);
213 : else
214 HIT 22 : wo->perform_io();
215 :
216 22 : if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
217 : {
218 4 : wo->errn = 0;
219 : }
220 : else
221 : {
222 18 : wait_read_op = nullptr;
223 18 : local_ops.push(wo);
224 : }
225 : }
226 : }
227 499784 : if (ev & reactor_event_write)
228 : {
229 55919 : 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 55919 : if (connect_op)
235 : {
236 7014 : auto* cn = connect_op;
237 7014 : if (err)
238 8 : cn->complete(err, 0);
239 : else
240 7006 : cn->perform_io();
241 :
242 7014 : if (cn->errn == EAGAIN || cn->errn == EWOULDBLOCK)
243 : {
244 MIS 0 : cn->errn = 0;
245 : }
246 : else
247 : {
248 HIT 7014 : connect_op = nullptr;
249 7014 : local_ops.push(cn);
250 : }
251 : }
252 55919 : if (write_op)
253 : {
254 190 : auto* wr = write_op;
255 190 : if (err)
256 MIS 0 : wr->complete(err, 0);
257 : else
258 HIT 190 : wr->perform_io();
259 :
260 190 : if (wr->errn == EAGAIN || wr->errn == EWOULDBLOCK)
261 : {
262 1 : wr->errn = 0;
263 : }
264 : else
265 : {
266 189 : write_op = nullptr;
267 189 : local_ops.push(wr);
268 : }
269 : }
270 55919 : if (!had_write_op)
271 48715 : write_ready = true;
272 :
273 : // Same re-probe discipline as the wait-for-read dispatch.
274 55919 : if (wait_write_op)
275 : {
276 4 : auto* wo = wait_write_op;
277 4 : if (err)
278 MIS 0 : wo->complete(err, 0);
279 : else
280 HIT 4 : wo->perform_io();
281 :
282 4 : if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
283 : {
284 MIS 0 : wo->errn = 0;
285 : }
286 : else
287 : {
288 HIT 4 : wait_write_op = nullptr;
289 4 : local_ops.push(wo);
290 : }
291 : }
292 : }
293 : // Complete a parked wait-for-error on any error condition.
294 499784 : if ((ev & reactor_event_error) || err)
295 : {
296 13 : if (wait_error_op)
297 : {
298 MIS 0 : wait_error_op->complete(err, 0);
299 0 : local_ops.push(std::exchange(wait_error_op, nullptr));
300 : }
301 : }
302 HIT 499784 : if (err)
303 : {
304 13 : if (read_op)
305 : {
306 MIS 0 : read_op->complete(err, 0);
307 0 : local_ops.push(std::exchange(read_op, nullptr));
308 : }
309 HIT 13 : if (write_op)
310 : {
311 MIS 0 : write_op->complete(err, 0);
312 0 : local_ops.push(std::exchange(write_op, nullptr));
313 : }
314 HIT 13 : if (connect_op)
315 : {
316 MIS 0 : connect_op->complete(err, 0);
317 0 : local_ops.push(std::exchange(connect_op, nullptr));
318 : }
319 HIT 13 : if (wait_read_op)
320 : {
321 MIS 0 : wait_read_op->complete(err, 0);
322 0 : local_ops.push(std::exchange(wait_read_op, nullptr));
323 : }
324 HIT 13 : if (wait_write_op)
325 : {
326 MIS 0 : wait_write_op->complete(err, 0);
327 0 : local_ops.push(std::exchange(wait_write_op, nullptr));
328 : }
329 : }
330 HIT 499788 : }
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 499784 : scheduler_op* first = ready_as_op(local_ops.pop());
336 499784 : if (first)
337 : {
338 14723 : scheduler_->post_deferred_completions(local_ops);
339 14723 : (*first)();
340 : }
341 : else
342 : {
343 485061 : scheduler_->compensating_work_started();
344 : }
345 499788 : }
346 :
347 : } // namespace boost::corosio::detail
348 :
349 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_STATE_HPP
|