TLA Line data Source code
1 : //
2 : // Copyright (c) 2026 Michael Vandeberg
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_IO_CORE_HPP
11 : #define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_IO_CORE_HPP
12 :
13 : #include <boost/corosio/detail/config.hpp>
14 : #include <boost/corosio/native/detail/reactor/reactor_op_base.hpp>
15 :
16 : #include <atomic>
17 : #include <memory>
18 : #include <mutex>
19 : #include <system_error>
20 : #include <utility>
21 :
22 : #include <errno.h>
23 :
24 : /* Shared reactor I/O protocol.
25 :
26 : One implementation of the register/park/cancel/teardown protocol
27 : for every reactor-backed object -- stream and datagram sockets,
28 : acceptors, and posix descriptors -- over a backend's
29 : descriptor_state. Asio keeps the same protocol in its reactor
30 : (start_op, cancel_ops, deregister_descriptor); services and the
31 : objects' own verbs stay per type.
32 :
33 : Derived supplies its op slots through for_each_op,
34 : for_each_desc_entry and op_to_desc_slot, and must derive from
35 : std::enable_shared_from_this<Derived>.
36 : */
37 :
38 : namespace boost::corosio::detail {
39 :
40 : /** CRTP base holding the reactor parking and cancel protocol.
41 :
42 : @tparam Derived The concrete object type (CRTP).
43 : @tparam Service The backend service that owns Derived.
44 : @tparam DescState The backend's descriptor_state type.
45 : */
46 : template<class Derived, class Service, class DescState>
47 : class reactor_io_core
48 : {
49 HIT 153565 : Derived* self_ptr() noexcept
50 : {
51 153565 : return static_cast<Derived*>(this);
52 : }
53 :
54 : protected:
55 : // NOLINTNEXTLINE(bugprone-crtp-constructor-accessibility)
56 16003 : explicit reactor_io_core(Service& svc) noexcept : svc_(svc) {}
57 :
58 : Service& svc_;
59 :
60 : public:
61 : /// Per-descriptor state for persistent reactor registration.
62 : DescState desc_state_;
63 :
64 : /** Cancel a single pending operation.
65 :
66 : Claims the operation from its descriptor_state slot under
67 : the mutex and posts it to the scheduler as cancelled.
68 : */
69 : template<class Op>
70 317 : void cancel_single_op(Op& op) noexcept
71 : {
72 317 : auto self = self_ptr()->weak_from_this().lock();
73 317 : if (!self)
74 MIS 0 : return;
75 :
76 HIT 317 : op.request_cancel();
77 :
78 317 : reactor_op_base** desc_op_ptr = self_ptr()->op_to_desc_slot(op);
79 317 : if (!desc_op_ptr)
80 MIS 0 : return;
81 :
82 HIT 317 : reactor_op_base* claimed = nullptr;
83 : {
84 317 : std::lock_guard lock(desc_state_.mutex);
85 317 : if (*desc_op_ptr == &op)
86 305 : claimed = std::exchange(*desc_op_ptr, nullptr);
87 : // Not in the slot: request_cancel() above already set
88 : // op.cancelled, which register_op consults before parking
89 : // and the completion decode consults on delivery. Latching
90 : // a descriptor flag here instead would outlive this op and
91 : // cancel the next wait in the same direction.
92 317 : }
93 317 : if (claimed)
94 : {
95 305 : op.impl_ptr = self;
96 305 : svc_.post(&op);
97 305 : svc_.work_finished();
98 : }
99 317 : }
100 :
101 : protected:
102 : /** Clear every slot and register @a fd with the reactor.
103 :
104 : @return The reactor's refusal, in which case the state is left
105 : unregistered.
106 : */
107 5598 : std::error_code register_fd(int fd) noexcept
108 : {
109 5598 : desc_state_.fd = fd;
110 : {
111 5598 : std::lock_guard lock(desc_state_.mutex);
112 5598 : self_ptr()->for_each_desc_entry(
113 34566 : [](auto&, reactor_op_base*& slot) { slot = nullptr; });
114 5598 : }
115 5598 : if (auto ec = svc_.scheduler().register_descriptor(fd, &desc_state_))
116 : {
117 4 : desc_state_.fd = -1;
118 4 : desc_state_.registered_events = 0;
119 4 : return ec;
120 : }
121 5594 : return {};
122 : }
123 :
124 : /** Register an op with the reactor.
125 :
126 : Handles cached edge events. Called on the EAGAIN/EINPROGRESS
127 : path when speculative I/O failed.
128 : */
129 : template<class Op>
130 5918 : void register_op(
131 : Op& op,
132 : reactor_op_base*& desc_slot,
133 : bool& ready_flag,
134 : bool is_write_direction = false) noexcept
135 : {
136 5918 : svc_.work_started();
137 :
138 5918 : std::lock_guard lock(desc_state_.mutex);
139 5918 : bool io_done = false;
140 5918 : if (ready_flag)
141 : {
142 225 : ready_flag = false;
143 225 : op.perform_io();
144 225 : io_done = (op.errn != EAGAIN && op.errn != EWOULDBLOCK);
145 225 : if (!io_done)
146 219 : op.errn = 0;
147 : }
148 :
149 5918 : if (io_done || op.cancelled.load(std::memory_order_acquire))
150 : {
151 6 : svc_.post(&op);
152 6 : svc_.work_finished();
153 6 : return;
154 : }
155 :
156 5912 : if (desc_state_.unpollable)
157 : {
158 : // Nothing will ever report readiness for this fd.
159 1 : op.complete(EOPNOTSUPP, 0);
160 1 : svc_.post(&op);
161 1 : svc_.work_finished();
162 1 : return;
163 : }
164 :
165 5911 : if (is_write_direction)
166 : {
167 4858 : if (auto ec = svc_.scheduler().ensure_write_registered(
168 : desc_state_.fd, &desc_state_))
169 : {
170 MIS 0 : op.complete(ec.value(), 0);
171 0 : svc_.post(&op);
172 0 : svc_.work_finished();
173 0 : return;
174 : }
175 : }
176 :
177 HIT 5911 : desc_slot = &op;
178 :
179 : // Select rebuilds its fd_sets from parked ops only, so parking
180 : // must wake it. Compiled away for epoll and kqueue.
181 : if constexpr (Service::needs_park_notification)
182 2779 : svc_.scheduler().notify_reactor();
183 5918 : }
184 :
185 : /// Cancel every pending operation.
186 306 : void cancel_all() noexcept
187 : {
188 306 : auto self = self_ptr()->weak_from_this().lock();
189 306 : if (!self)
190 MIS 0 : return;
191 :
192 HIT 2197 : self_ptr()->for_each_op([](auto& op) { op.request_cancel(); });
193 :
194 : reactor_op_base* claimed[max_claimed];
195 306 : int count = 0;
196 : {
197 306 : std::lock_guard lock(desc_state_.mutex);
198 306 : self_ptr()->for_each_desc_entry(
199 3782 : [&](auto& op, reactor_op_base*& desc_slot) {
200 1891 : if (desc_slot == &op)
201 : {
202 213 : BOOST_COROSIO_ASSERT(count < max_claimed);
203 213 : claimed[count++] = std::exchange(desc_slot, nullptr);
204 : }
205 : });
206 306 : }
207 306 : post_claimed(claimed, count, self);
208 306 : }
209 :
210 : /** Cancel every operation and claim every parked one for teardown.
211 :
212 : Also clears the cached edge flags and, if the state is queued
213 : in the scheduler, pins the object alive until it is drained.
214 : */
215 48805 : void abandon_all() noexcept
216 : {
217 48805 : auto self = self_ptr()->weak_from_this().lock();
218 48805 : if (!self)
219 MIS 0 : return;
220 :
221 HIT 339705 : self_ptr()->for_each_op([](auto& op) { op.request_cancel(); });
222 :
223 : reactor_op_base* claimed[max_claimed];
224 48805 : int count = 0;
225 : {
226 48805 : std::lock_guard lock(desc_state_.mutex);
227 48805 : self_ptr()->for_each_desc_entry(
228 581800 : [&](auto& /*op*/, reactor_op_base*& desc_slot) {
229 290900 : if (auto* c = std::exchange(desc_slot, nullptr))
230 : {
231 106 : BOOST_COROSIO_ASSERT(count < max_claimed);
232 106 : claimed[count++] = c;
233 : }
234 : });
235 48805 : desc_state_.read_ready = false;
236 48805 : desc_state_.write_ready = false;
237 :
238 : // Must be set under the same lock that invoke_deferred_io
239 : // clears is_enqueued_ under, or the object could be destroyed
240 : // while the scheduler still holds the queued descriptor_state.
241 48805 : if (desc_state_.is_enqueued_.load(std::memory_order_acquire))
242 712 : desc_state_.impl_ref_ = self;
243 48805 : }
244 48805 : post_claimed(claimed, count, self);
245 48805 : }
246 :
247 : /// Drop the reactor registration of @a fd and reset the state.
248 48805 : void unregister_fd(int fd) noexcept
249 : {
250 48805 : if (fd >= 0 && desc_state_.registered_events != 0)
251 10839 : svc_.scheduler().deregister_descriptor(fd);
252 48805 : desc_state_.fd = -1;
253 48805 : desc_state_.registered_events = 0;
254 48805 : desc_state_.unpollable = false;
255 48805 : }
256 :
257 : private:
258 : // A claim empties its slot, so no more ops are claimed than
259 : // descriptor_state has slots (read, write, connect, wait_read,
260 : // wait_write, wait_error), however many ops share them.
261 : static constexpr int max_claimed = 6;
262 :
263 49111 : void post_claimed(
264 : reactor_op_base** claimed,
265 : int count,
266 : std::shared_ptr<Derived> const& self) noexcept
267 : {
268 49430 : for (int i = 0; i < count; ++i)
269 : {
270 319 : claimed[i]->impl_ptr = self;
271 319 : svc_.post(claimed[i]);
272 319 : svc_.work_finished();
273 : }
274 49111 : }
275 : };
276 :
277 : } // namespace boost::corosio::detail
278 :
279 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_IO_CORE_HPP
|