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 : /// 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 : bool unpollable = false;
82 :
83 : /// File descriptor this state tracks.
84 : int fd = -1;
85 :
86 : /// Accumulated ready events (set by reactor, read by scheduler).
87 : std::atomic<std::uint32_t> ready_events_{0};
88 :
89 : /// True while this state is queued in the scheduler's completed_ops.
90 : std::atomic<bool> is_enqueued_{false};
91 :
92 : /// Owning scheduler for posting completions.
93 : 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 HIT 40459 : void add_ready_events(std::uint32_t ev) noexcept
103 : {
104 40459 : ready_events_.fetch_or(ev, std::memory_order_release);
105 40459 : }
106 :
107 : /// Invoke deferred I/O and dispatch completions.
108 40254 : void operator()() override
109 : {
110 40254 : invoke_deferred_io();
111 40254 : }
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 205 : void destroy() override
117 : {
118 205 : impl_ref_.reset();
119 205 : }
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 40254 : reactor_descriptor_state::invoke_deferred_io()
132 : {
133 40254 : std::shared_ptr<void> prevent_impl_destruction;
134 40254 : ready_queue local_ops;
135 :
136 : {
137 40254 : 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 40254 : prevent_impl_destruction = std::move(impl_ref_);
145 40254 : is_enqueued_.store(false, std::memory_order_release);
146 :
147 40254 : std::uint32_t ev = ready_events_.exchange(0, std::memory_order_acquire);
148 40254 : if (ev == 0)
149 : {
150 : // Mutex unlocks here; compensate for work_cleanup's decrement
151 MIS 0 : scheduler_->compensating_work_started();
152 0 : return;
153 : }
154 :
155 HIT 40254 : int err = 0;
156 40254 : 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 43 : 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 43 : bool const reports_error =
177 43 : read_op || write_op || connect_op || wait_error_op;
178 43 : socklen_t len = sizeof(err);
179 68 : if (reports_error &&
180 25 : ::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 11 : err = (errno == ENOTSOCK) ? 0 : errno;
185 : }
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 : }
193 :
194 40254 : if (ev & reactor_event_read)
195 : {
196 15185 : if (read_op)
197 : {
198 5435 : auto* rd = read_op;
199 5435 : if (err)
200 4 : rd->complete(err, 0);
201 : else
202 5431 : rd->perform_io();
203 :
204 5435 : if (rd->errn == EAGAIN || rd->errn == EWOULDBLOCK)
205 : {
206 319 : rd->errn = 0;
207 : }
208 : else
209 : {
210 5116 : read_op = nullptr;
211 5116 : local_ops.push(rd);
212 : }
213 : }
214 : else
215 : {
216 9750 : 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 15185 : if (wait_read_op)
225 : {
226 36 : auto* wo = wait_read_op;
227 36 : wo->perform_io();
228 :
229 36 : if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
230 : {
231 2 : wo->errn = 0;
232 : }
233 : else
234 : {
235 34 : wait_read_op = nullptr;
236 34 : local_ops.push(wo);
237 : }
238 : }
239 : }
240 40254 : if (ev & reactor_event_write)
241 : {
242 35039 : 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 35039 : if (connect_op)
248 : {
249 4622 : auto* cn = connect_op;
250 4622 : if (err)
251 9 : cn->complete(err, 0);
252 : else
253 4613 : cn->perform_io();
254 :
255 4622 : if (cn->errn == EAGAIN || cn->errn == EWOULDBLOCK)
256 : {
257 MIS 0 : cn->errn = 0;
258 : }
259 : else
260 : {
261 HIT 4622 : connect_op = nullptr;
262 4622 : local_ops.push(cn);
263 : }
264 : }
265 35039 : if (write_op)
266 : {
267 200 : auto* wr = write_op;
268 200 : if (err)
269 2 : wr->complete(err, 0);
270 : else
271 198 : wr->perform_io();
272 :
273 200 : if (wr->errn == EAGAIN || wr->errn == EWOULDBLOCK)
274 : {
275 1 : wr->errn = 0;
276 : }
277 : else
278 : {
279 199 : write_op = nullptr;
280 199 : local_ops.push(wr);
281 : }
282 : }
283 35039 : if (!had_write_op)
284 30217 : write_ready = true;
285 :
286 : // Same re-probe discipline as the wait-for-read dispatch.
287 35039 : if (wait_write_op)
288 : {
289 11 : auto* wo = wait_write_op;
290 11 : wo->perform_io();
291 :
292 11 : if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
293 : {
294 MIS 0 : wo->errn = 0;
295 : }
296 : else
297 : {
298 HIT 11 : wait_write_op = nullptr;
299 11 : local_ops.push(wo);
300 : }
301 : }
302 : }
303 : // Complete a parked wait-for-error on any error condition.
304 40254 : if (ev & reactor_event_error)
305 : {
306 43 : 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 6 : int const werr = err ? err : EIO;
312 6 : wait_error_op->complete(werr, 0);
313 6 : local_ops.push(std::exchange(wait_error_op, nullptr));
314 : }
315 : }
316 40254 : if (err)
317 : {
318 18 : if (read_op)
319 : {
320 MIS 0 : read_op->complete(err, 0);
321 0 : local_ops.push(std::exchange(read_op, nullptr));
322 : }
323 HIT 18 : if (write_op)
324 : {
325 MIS 0 : write_op->complete(err, 0);
326 0 : local_ops.push(std::exchange(write_op, nullptr));
327 : }
328 HIT 18 : if (connect_op)
329 : {
330 MIS 0 : connect_op->complete(err, 0);
331 0 : local_ops.push(std::exchange(connect_op, nullptr));
332 : }
333 : }
334 HIT 40254 : }
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 40254 : scheduler_op* first = ready_as_op(local_ops.pop());
340 40254 : if (first)
341 : {
342 9986 : scheduler_->post_deferred_completions(local_ops);
343 9986 : (*first)();
344 : }
345 : else
346 : {
347 30268 : scheduler_->compensating_work_started();
348 : }
349 40254 : }
350 :
351 : } // namespace boost::corosio::detail
352 :
353 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_STATE_HPP
|