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

92.5% Lines (98 / 106) 100.0% Functions (4 / 4)
reactor_descriptor_state.hpp
f(x) Functions (4)
Line TLA Hits 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 351466x void add_ready_events(std::uint32_t ev) noexcept
99 {
100 351466x ready_events_.fetch_or(ev, std::memory_order_release);
101 351466x }
102
103 /// Invoke deferred I/O and dispatch completions.
104 351157x void operator()() override
105 {
106 351157x invoke_deferred_io();
107 351157x }
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 309x void destroy() override
113 {
114 309x impl_ref_.reset();
115 309x }
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 351157x reactor_descriptor_state::invoke_deferred_io()
128 {
129 351157x std::shared_ptr<void> prevent_impl_destruction;
130 351157x ready_queue local_ops;
131
132 {
133 351157x 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 351157x prevent_impl_destruction = std::move(impl_ref_);
141 351157x is_enqueued_.store(false, std::memory_order_release);
142
143 351157x std::uint32_t ev = ready_events_.exchange(0, std::memory_order_acquire);
144 351157x if (ev == 0)
145 {
146 // Mutex unlocks here; compensate for work_cleanup's decrement
147 4x scheduler_->compensating_work_started();
148 4x return;
149 }
150
151 351153x int err = 0;
152 351153x if (ev & reactor_event_error)
153 {
154 33x socklen_t len = sizeof(err);
155 33x if (::getsockopt(fd, SOL_SOCKET, SO_ERROR, &err, &len) < 0)
156 {
157 12x 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 2x err = 0;
174 2x ev |= reactor_event_read | reactor_event_write;
175 }
176 else
177 {
178 10x 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 351153x if (ev & reactor_event_read)
190 {
191 325877x if (read_op)
192 {
193 5343x auto* rd = read_op;
194 5343x if (err)
195 3x rd->complete(err, 0);
196 else
197 5340x rd->perform_io();
198
199 5343x if (rd->errn == EAGAIN || rd->errn == EWOULDBLOCK)
200 {
201 370x rd->errn = 0;
202 }
203 else
204 {
205 4973x read_op = nullptr;
206 4973x local_ops.push(rd);
207 }
208 }
209 else
210 {
211 320534x 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 325877x if (wait_read_op)
220 {
221 31x auto* wo = wait_read_op;
222 31x if (err)
223 1x wo->complete(err, 0);
224 else
225 30x wo->perform_io();
226
227 31x if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
228 {
229 4x wo->errn = 0;
230 }
231 else
232 {
233 27x wait_read_op = nullptr;
234 27x local_ops.push(wo);
235 }
236 }
237 }
238 351153x if (ev & reactor_event_write)
239 {
240 33417x 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 33417x if (connect_op)
246 {
247 4500x auto* cn = connect_op;
248 4500x if (err)
249 8x cn->complete(err, 0);
250 else
251 4492x cn->perform_io();
252
253 4500x if (cn->errn == EAGAIN || cn->errn == EWOULDBLOCK)
254 {
255 ✗ cn->errn = 0;
256 }
257 else
258 {
259 4500x connect_op = nullptr;
260 4500x local_ops.push(cn);
261 }
262 }
263 33417x if (write_op)
264 {
265 198x auto* wr = write_op;
266 198x if (err)
267 2x wr->complete(err, 0);
268 else
269 196x wr->perform_io();
270
271 198x if (wr->errn == EAGAIN || wr->errn == EWOULDBLOCK)
272 {
273 1x wr->errn = 0;
274 }
275 else
276 {
277 197x write_op = nullptr;
278 197x local_ops.push(wr);
279 }
280 }
281 33417x if (!had_write_op)
282 28719x write_ready = true;
283
284 // Same re-probe discipline as the wait-for-read dispatch.
285 33417x if (wait_write_op)
286 {
287 9x auto* wo = wait_write_op;
288 9x if (err)
289 2x wo->complete(err, 0);
290 else
291 7x wo->perform_io();
292
293 9x if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
294 {
295 ✗ wo->errn = 0;
296 }
297 else
298 {
299 9x wait_write_op = nullptr;
300 9x local_ops.push(wo);
301 }
302 }
303 }
304 // Complete a parked wait-for-error on any error condition.
305 351153x if ((ev & reactor_event_error) || err)
306 {
307 33x 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 4x int const werr = err ? err : EIO;
313 4x wait_error_op->complete(werr, 0);
314 4x local_ops.push(std::exchange(wait_error_op, nullptr));
315 }
316 }
317 351153x if (err)
318 {
319 26x if (read_op)
320 {
321 1x read_op->complete(err, 0);
322 1x local_ops.push(std::exchange(read_op, nullptr));
323 }
324 26x if (write_op)
325 {
326 ✗ write_op->complete(err, 0);
327 ✗ local_ops.push(std::exchange(write_op, nullptr));
328 }
329 26x if (connect_op)
330 {
331 ✗ connect_op->complete(err, 0);
332 ✗ local_ops.push(std::exchange(connect_op, nullptr));
333 }
334 26x if (wait_read_op)
335 {
336 1x wait_read_op->complete(err, 0);
337 1x local_ops.push(std::exchange(wait_read_op, nullptr));
338 }
339 26x if (wait_write_op)
340 {
341 ✗ wait_write_op->complete(err, 0);
342 ✗ local_ops.push(std::exchange(wait_write_op, nullptr));
343 }
344 }
345 351157x }
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 351153x scheduler_op* first = ready_as_op(local_ops.pop());
351 351153x if (first)
352 {
353 9710x scheduler_->post_deferred_completions(local_ops);
354 9710x (*first)();
355 }
356 else
357 {
358 341443x scheduler_->compensating_work_started();
359 }
360 351157x }
361
362 } // namespace boost::corosio::detail
363
364 #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_STATE_HPP
365