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

89.7% Lines (87 / 97) 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 /// 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 53376x void add_ready_events(std::uint32_t ev) noexcept
103 {
104 53376x ready_events_.fetch_or(ev, std::memory_order_release);
105 53376x }
106
107 /// Invoke deferred I/O and dispatch completions.
108 51958x void operator()() override
109 {
110 51958x invoke_deferred_io();
111 51958x }
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 254x void destroy() override
117 {
118 254x impl_ref_.reset();
119 254x }
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 51958x reactor_descriptor_state::invoke_deferred_io()
132 {
133 51958x std::shared_ptr<void> prevent_impl_destruction;
134 51958x ready_queue local_ops;
135
136 {
137 51958x 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 51958x prevent_impl_destruction = std::move(impl_ref_);
145 51958x is_enqueued_.store(false, std::memory_order_release);
146
147 51958x std::uint32_t ev = ready_events_.exchange(0, std::memory_order_acquire);
148 51958x if (ev == 0)
149 {
150 // Mutex unlocks here; compensate for work_cleanup's decrement
151 ✗ scheduler_->compensating_work_started();
152 ✗ return;
153 }
154
155 51958x int err = 0;
156 51958x 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 49x 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 49x bool const reports_error =
177 49x read_op || write_op || connect_op || wait_error_op;
178 49x socklen_t len = sizeof(err);
179 76x if (reports_error &&
180 27x ::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 11x 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 51958x if (ev & reactor_event_read)
195 {
196 18925x if (read_op)
197 {
198 8582x auto* rd = read_op;
199 8582x if (err)
200 6x rd->complete(err, 0);
201 else
202 8576x rd->perform_io();
203
204 8582x if (rd->errn == EAGAIN || rd->errn == EWOULDBLOCK)
205 {
206 319x rd->errn = 0;
207 }
208 else
209 {
210 8263x read_op = nullptr;
211 8263x local_ops.push(rd);
212 }
213 }
214 else
215 {
216 10343x 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 18925x if (wait_read_op)
225 {
226 36x auto* wo = wait_read_op;
227 36x wo->perform_io();
228
229 36x if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
230 {
231 2x wo->errn = 0;
232 }
233 else
234 {
235 34x wait_read_op = nullptr;
236 34x local_ops.push(wo);
237 }
238 }
239 }
240 51958x if (ev & reactor_event_write)
241 {
242 44925x 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 44925x if (connect_op)
248 {
249 6450x auto* cn = connect_op;
250 6450x if (err)
251 8x cn->complete(err, 0);
252 else
253 6442x cn->perform_io();
254
255 6450x if (cn->errn == EAGAIN || cn->errn == EWOULDBLOCK)
256 {
257 ✗ cn->errn = 0;
258 }
259 else
260 {
261 6450x connect_op = nullptr;
262 6450x local_ops.push(cn);
263 }
264 }
265 44925x if (write_op)
266 {
267 200x auto* wr = write_op;
268 200x if (err)
269 2x wr->complete(err, 0);
270 else
271 198x wr->perform_io();
272
273 200x if (wr->errn == EAGAIN || wr->errn == EWOULDBLOCK)
274 {
275 1x wr->errn = 0;
276 }
277 else
278 {
279 199x write_op = nullptr;
280 199x local_ops.push(wr);
281 }
282 }
283 44925x if (!had_write_op)
284 38275x write_ready = true;
285
286 // Same re-probe discipline as the wait-for-read dispatch.
287 44925x if (wait_write_op)
288 {
289 11x auto* wo = wait_write_op;
290 11x wo->perform_io();
291
292 11x if (wo->errn == EAGAIN || wo->errn == EWOULDBLOCK)
293 {
294 ✗ wo->errn = 0;
295 }
296 else
297 {
298 11x wait_write_op = nullptr;
299 11x local_ops.push(wo);
300 }
301 }
302 }
303 // Complete a parked wait-for-error on any error condition.
304 51958x if (ev & reactor_event_error)
305 {
306 49x 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 7x int const werr = err ? err : EIO;
312 7x wait_error_op->complete(werr, 0);
313 7x local_ops.push(std::exchange(wait_error_op, nullptr));
314 }
315 }
316 51958x if (err)
317 {
318 20x if (read_op)
319 {
320 ✗ read_op->complete(err, 0);
321 ✗ local_ops.push(std::exchange(read_op, nullptr));
322 }
323 20x if (write_op)
324 {
325 ✗ write_op->complete(err, 0);
326 ✗ local_ops.push(std::exchange(write_op, nullptr));
327 }
328 20x if (connect_op)
329 {
330 ✗ connect_op->complete(err, 0);
331 ✗ local_ops.push(std::exchange(connect_op, nullptr));
332 }
333 }
334 51958x }
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 51958x scheduler_op* first = ready_as_op(local_ops.pop());
340 51958x if (first)
341 {
342 14962x scheduler_->post_deferred_completions(local_ops);
343 14962x (*first)();
344 }
345 else
346 {
347 36996x scheduler_->compensating_work_started();
348 }
349 51958x }
350
351 } // namespace boost::corosio::detail
352
353 #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_STATE_HPP
354