TLA Line data 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 : 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 : conditionally_enabled_mutex mutex{true};
51 :
52 : /// Pending read operation (guarded by `mutex`).
53 : reactor_op_base* read_op = nullptr;
54 :
55 : /// Pending write operation (guarded by `mutex`).
56 : reactor_op_base* write_op = nullptr;
57 :
58 : /// Pending connect operation (guarded by `mutex`).
59 : reactor_op_base* connect_op = nullptr;
60 :
61 : /// Pending wait-for-read operation (guarded by `mutex`).
62 : reactor_op_base* wait_read_op = nullptr;
63 :
64 : /// Pending wait-for-write operation (guarded by `mutex`).
65 : reactor_op_base* wait_write_op = nullptr;
66 :
67 : /// Pending wait-for-error operation (guarded by `mutex`).
68 : reactor_op_base* wait_error_op = nullptr;
69 :
70 : /// True if a read edge event arrived before an op was registered.
71 : bool read_ready = false;
72 :
73 : /// True if a write edge event arrived before an op was registered.
74 : bool write_ready = false;
75 :
76 : /// Event mask set during registration (no mutex needed).
77 : std::uint32_t registered_events = 0;
78 :
79 : /// File descriptor this state tracks.
80 : int fd = -1;
81 :
82 : /// Accumulated ready events (set by reactor, read by scheduler).
83 : std::atomic<std::uint32_t> ready_events_{0};
84 :
85 : /// True while this state is queued in the scheduler's completed_ops.
86 : std::atomic<bool> is_enqueued_{false};
87 :
88 : /// Owning scheduler for posting completions.
89 : reactor_scheduler const* scheduler_ = nullptr;
90 :
91 : /// Prevents impl destruction while queued in the scheduler.
92 : std::shared_ptr<void> impl_ref_;
93 :
94 : /// Add ready events atomically.
95 : /// Release pairs with the consumer's acquire exchange on
96 : /// ready_events_ so the consumer sees all flags. On x86 (TSO)
97 : /// this compiles to the same LOCK OR as relaxed.
98 HIT 351466 : void add_ready_events(std::uint32_t ev) noexcept
99 : {
100 351466 : ready_events_.fetch_or(ev, std::memory_order_release);
101 351466 : }
102 :
103 : /// Invoke deferred I/O and dispatch completions.
104 351157 : void operator()() override
105 : {
106 351157 : invoke_deferred_io();
107 351157 : }
108 :
109 : /// Destroy without invoking.
110 : /// Called during scheduler::shutdown() drain. Clear impl_ref_ to break
111 : /// the self-referential cycle set by close_socket().
112 309 : void destroy() override
113 : {
114 309 : impl_ref_.reset();
115 309 : }
116 :
117 : /** Perform deferred I/O and queue completions.
118 :
119 : Performs I/O under the mutex and queues completed ops. EAGAIN
120 : ops stay parked in their slot for re-delivery on the next
121 : edge event.
122 : */
123 : void invoke_deferred_io();
124 : };
125 :
126 : inline void
127 351157 : reactor_descriptor_state::invoke_deferred_io()
128 : {
129 351157 : std::shared_ptr<void> prevent_impl_destruction;
130 351157 : ready_queue local_ops;
131 :
132 : {
133 351157 : conditionally_enabled_mutex::scoped_lock lock(mutex);
134 :
135 : // Must clear is_enqueued_ and move impl_ref_ under the same
136 : // lock that processes I/O. close_socket() checks is_enqueued_
137 : // under this mutex — without atomicity between the flag store
138 : // and the ref move, close_socket() could see is_enqueued_==false,
139 : // skip setting impl_ref_, and destroy the impl under us.
140 351157 : prevent_impl_destruction = std::move(impl_ref_);
141 351157 : is_enqueued_.store(false, std::memory_order_release);
142 :
143 351157 : std::uint32_t ev = ready_events_.exchange(0, std::memory_order_acquire);
144 351157 : if (ev == 0)
145 : {
146 : // Mutex unlocks here; compensate for work_cleanup's decrement
147 4 : scheduler_->compensating_work_started();
148 4 : return;
149 : }
150 :
151 351153 : int err = 0;
152 351153 : if (ev & reactor_event_error)
153 : {
154 33 : socklen_t len = sizeof(err);
155 33 : if (::getsockopt(fd, SOL_SOCKET, SO_ERROR, &err, &len) < 0)
156 : {
157 12 : if (errno == ENOTSOCK)
158 : {
159 : // Non-socket fd (pipe, chardev, ...): no SO_ERROR, so
160 : // let the op's own syscall name the real failure.
161 : // Also force the read/write dispatch below to run: an
162 : // edge-triggered EPOLLERR can arrive alone, and
163 : // without this a parked op never calls perform_io()
164 : // and, the edge being one-shot, never gets another
165 : // chance -- a permanent hang.
166 : //
167 : // Assumes at least one parked op's own syscall makes
168 : // non-EAGAIN progress; if every op re-parks with
169 : // EAGAIN this sticky error is never redelivered and
170 : // they hang. No such case is known for pipes -- a
171 : // future non-socket type that hits one should be
172 : // handled here.
173 2 : err = 0;
174 2 : ev |= reactor_event_read | reactor_event_write;
175 : }
176 : else
177 : {
178 10 : err = errno;
179 : }
180 : }
181 : // select raises its exceptional set for out-of-band/urgent
182 : // data as well as for genuine faults; on a healthy socket the
183 : // probe then reads SO_ERROR == 0. Faulting a pending read or
184 : // write on that is wrong, so an I/O operation completes only
185 : // on a real (non-zero) error. wait(error) still names a code
186 : // below.
187 : }
188 :
189 351153 : if (ev & reactor_event_read)
190 : {
191 325877 : if (read_op)
192 : {
193 5343 : auto* rd = read_op;
194 5343 : if (err)
195 3 : rd->complete(err, 0);
196 : else
197 5340 : rd->perform_io();
198 :
199 5343 : if (rd->errn == EAGAIN || rd->errn == EWOULDBLOCK)
200 : {
201 370 : rd->errn = 0;
202 : }
203 : else
204 : {
205 4973 : read_op = nullptr;
206 4973 : local_ops.push(rd);
207 : }
208 : }
209 : else
210 : {
211 320534 : read_ready = true;
212 : }
213 :
214 : // The event does not prove the socket is still readable: a
215 : // parked read op above may have drained it, or a speculative
216 : // read consumed the data before this dispatch ran. The wait
217 : // op's perform_io() re-probes and reports EAGAIN to stay
218 : // parked.
219 325877 : if (wait_read_op)
220 : {
221 31 : auto* wo = wait_read_op;
222 31 : if (err)
223 1 : wo->complete(err, 0);
224 : else
225 30 : wo->perform_io();
226 :
227 31 : if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
228 : {
229 4 : wo->errn = 0;
230 : }
231 : else
232 : {
233 27 : wait_read_op = nullptr;
234 27 : local_ops.push(wo);
235 : }
236 : }
237 : }
238 351153 : if (ev & reactor_event_write)
239 : {
240 33417 : bool had_write_op = (connect_op || write_op);
241 : // A writable event on a socket still in SYN_SENT (e.g. the
242 : // spurious pre-connect readiness of a fresh socket) must
243 : // not complete the connect; perform_io() reports EAGAIN
244 : // until a peer is actually established.
245 33417 : if (connect_op)
246 : {
247 4500 : auto* cn = connect_op;
248 4500 : if (err)
249 8 : cn->complete(err, 0);
250 : else
251 4492 : cn->perform_io();
252 :
253 4500 : if (cn->errn == EAGAIN || cn->errn == EWOULDBLOCK)
254 : {
255 MIS 0 : cn->errn = 0;
256 : }
257 : else
258 : {
259 HIT 4500 : connect_op = nullptr;
260 4500 : local_ops.push(cn);
261 : }
262 : }
263 33417 : if (write_op)
264 : {
265 198 : auto* wr = write_op;
266 198 : if (err)
267 2 : wr->complete(err, 0);
268 : else
269 196 : wr->perform_io();
270 :
271 198 : if (wr->errn == EAGAIN || wr->errn == EWOULDBLOCK)
272 : {
273 1 : wr->errn = 0;
274 : }
275 : else
276 : {
277 197 : write_op = nullptr;
278 197 : local_ops.push(wr);
279 : }
280 : }
281 33417 : if (!had_write_op)
282 28719 : write_ready = true;
283 :
284 : // Same re-probe discipline as the wait-for-read dispatch.
285 33417 : if (wait_write_op)
286 : {
287 9 : auto* wo = wait_write_op;
288 9 : if (err)
289 2 : wo->complete(err, 0);
290 : else
291 7 : wo->perform_io();
292 :
293 9 : if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
294 : {
295 MIS 0 : wo->errn = 0;
296 : }
297 : else
298 : {
299 HIT 9 : wait_write_op = nullptr;
300 9 : local_ops.push(wo);
301 : }
302 : }
303 : }
304 : // Complete a parked wait-for-error on any error condition.
305 351153 : if ((ev & reactor_event_error) || err)
306 : {
307 33 : if (wait_error_op)
308 : {
309 : // wait(error) fired on the exceptional condition; name a
310 : // code even when the kernel exposed none (e.g. urgent
311 : // data leaves SO_ERROR == 0).
312 4 : int const werr = err ? err : EIO;
313 4 : wait_error_op->complete(werr, 0);
314 4 : local_ops.push(std::exchange(wait_error_op, nullptr));
315 : }
316 : }
317 351153 : if (err)
318 : {
319 26 : if (read_op)
320 : {
321 1 : read_op->complete(err, 0);
322 1 : local_ops.push(std::exchange(read_op, nullptr));
323 : }
324 26 : if (write_op)
325 : {
326 MIS 0 : write_op->complete(err, 0);
327 0 : local_ops.push(std::exchange(write_op, nullptr));
328 : }
329 HIT 26 : if (connect_op)
330 : {
331 MIS 0 : connect_op->complete(err, 0);
332 0 : local_ops.push(std::exchange(connect_op, nullptr));
333 : }
334 HIT 26 : if (wait_read_op)
335 : {
336 1 : wait_read_op->complete(err, 0);
337 1 : local_ops.push(std::exchange(wait_read_op, nullptr));
338 : }
339 26 : if (wait_write_op)
340 : {
341 MIS 0 : wait_write_op->complete(err, 0);
342 0 : local_ops.push(std::exchange(wait_write_op, nullptr));
343 : }
344 : }
345 HIT 351157 : }
346 :
347 : // Execute first handler inline — the scheduler's work_cleanup
348 : // accounts for this as the "consumed" work item. local_ops holds
349 : // only ops, so the popped entry decodes directly.
350 351153 : scheduler_op* first = ready_as_op(local_ops.pop());
351 351153 : if (first)
352 : {
353 9710 : scheduler_->post_deferred_completions(local_ops);
354 9710 : (*first)();
355 : }
356 : else
357 : {
358 341443 : scheduler_->compensating_work_started();
359 : }
360 351157 : }
361 :
362 : } // namespace boost::corosio::detail
363 :
364 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_STATE_HPP
|