include/boost/corosio/native/detail/uring/uring_multishot_acceptor.hpp

83.3% Lines (349/0/419) 93.0% List of functions (40/0/43)
uring_multishot_acceptor.hpp
f(x) Functions (43)
Function Calls Lines Blocks
boost::corosio::detail::fd_is_listening(int) :55 9x 90.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::uring_multishot_acceptor_base(boost::corosio::detail::uring_scheduler&, boost::corosio::detail::uring_local_stream_service&) :137 55x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::uring_multishot_acceptor_base(boost::corosio::detail::uring_scheduler&, boost::corosio::detail::uring_tcp_service&) :137 217x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::~uring_multishot_acceptor_base() :145 54x 85.0% 90.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::~uring_multishot_acceptor_base() :145 215x 85.0% 90.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::local_endpoint() const :183 5x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::local_endpoint() const :183 165x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::is_open() const :188 245x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::is_open() const :188 4630x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::native_handle() const :193 4x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::native_handle() const :193 12x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::release_socket() :198 5x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::release_socket() :198 4x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::cancel() :219 4x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::cancel() :219 6x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::drain_waiters_only() :230 48x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::drain_waiters_only() :230 213x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::park_read_wait(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::stop_token const&, std::error_code*) :271 3x 79.6% 67.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::park_read_wait(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::stop_token const&, std::error_code*) :271 10x 95.9% 90.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::set_option(int, int, void const*, unsigned long) :357 2x 85.7% 89.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::set_option(int, int, void const*, unsigned long) :357 183x 85.7% 89.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::get_option(int, int, void*, unsigned long*) const :373 2x 88.9% 90.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::get_option(int, int, void*, unsigned long*) const :373 4x 88.9% 90.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::retire_multishot() :406 31x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::retire_multishot() :406 187x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::adopt_listening_fd(int) :438 3x 75.0% 60.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::adopt_listening_fd(int) :438 6x 75.0% 60.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::prepare_listen_arm() :478 28x 78.9% 67.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::prepare_listen_arm() :478 183x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::start_multishot() :512 31x 68.4% 62.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::start_multishot() :512 191x 100.0% 91.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::fail_arm(int) :564 0 0.0% 0.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::fail_arm(int) :564 5x 39.3% 42.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::dispatch_or_queue(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::stop_token const&, std::error_code*, boost::corosio::io_object::implementation**) :612 20x 67.4% 70.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::dispatch_or_queue(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::stop_token const&, std::error_code*, boost::corosio::io_object::implementation**) :612 3266x 76.8% 82.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::cancel_waiter(boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::waiter_node*) :764 6x 81.0% 83.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::cancel_waiter(boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::waiter_node*) :764 18x 90.5% 92.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::on_accept_cqe(void*, int, int, bool) :800 20x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::on_accept_cqe(void*, int, int, bool) :800 3253x 100.0% 100.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::on_accept_cqe_impl(int, int, bool) :806 20x 73.7% 57.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::on_accept_cqe_impl(int, int, bool) :806 3253x 76.3% 75.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::waiter_canceller::operator()() const :977 0 0.0% 0.0% boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::waiter_canceller::operator()() const :977 0 0.0% 0.0%
Line TLA Hits Source Code
1 //
2 // Copyright (c) 2026 Steve Gerbino
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_URING_URING_MULTISHOT_ACCEPTOR_HPP
11 #define BOOST_COROSIO_NATIVE_DETAIL_URING_URING_MULTISHOT_ACCEPTOR_HPP
12
13 #include <boost/corosio/detail/platform.hpp>
14
15 #if BOOST_COROSIO_HAS_URING
16
17 #include <liburing.h>
18
19 #include <boost/corosio/detail/intrusive.hpp>
20 #include <boost/corosio/native/detail/uring/uring_acceptor_ops.hpp>
21 #include <boost/corosio/native/detail/uring/uring_buffer.hpp>
22 #include <boost/corosio/native/detail/uring/uring_op.hpp>
23 #include <boost/corosio/native/detail/uring/uring_scheduler.hpp>
24 #include <boost/corosio/native/detail/uring/uring_socket_ops.hpp>
25 #include <boost/corosio/native/detail/make_err.hpp>
26 #include <boost/corosio/detail/native_handle.hpp>
27 #include <boost/corosio/io/io_object.hpp>
28
29 #include <atomic>
30 #include <cstdint>
31 #include <coroutine>
32 #include <memory>
33 #include <mutex>
34 #include <optional>
35 #include <stop_token>
36 #include <system_error>
37
38 #include <netinet/in.h>
39 #include <sys/socket.h>
40 #include <unistd.h>
41
42 namespace boost::corosio::detail {
43
44 /** Check whether a descriptor is in the listening state.
45
46 Multishot accept fails immediately on a non-listening socket and
47 the re-arm path would spin on that failure, so an adopted fd is
48 only armed when the kernel reports a listener. A query the
49 platform refuses is treated as listening so adoption still works.
50
51 @param fd The descriptor to probe.
52 @return True unless the kernel positively reports a non-listener.
53 */
54 inline bool
55 9x fd_is_listening(int fd) noexcept
56 {
57 9x int accepting = 0;
58 9x socklen_t alen = sizeof(accepting);
59 9x if (::getsockopt(fd, SOL_SOCKET, SO_ACCEPTCONN, &accepting, &alen) != 0)
60 1x return true;
61 8x return accepting != 0;
62 }
63
64 template<class Derived, class ImplBase, class Endpoint, class PeerService>
65 class uring_multishot_acceptor_base
66 : public ImplBase
67 , public std::enable_shared_from_this<Derived>
68 {
69 protected:
70 struct ready_fd_node : intrusive_list<ready_fd_node>::node
71 {
72 int fd = -1;
73 sockaddr_storage peer{};
74 socklen_t peer_len = 0;
75 };
76
77 struct waiter_node;
78
79 struct waiter_canceller
80 {
81 waiter_node* w;
82 void operator()() const noexcept;
83 };
84
85 struct waiter_node : intrusive_list<waiter_node>::node
86 {
87 std::coroutine_handle<> h;
88 capy::executor_ref ex;
89 std::error_code* ec_out = nullptr;
90 io_object::implementation** impl_out = nullptr;
91 Derived* owner = nullptr;
92 std::atomic<bool> cancelled{false};
93 /// True once linked into `waiters_` (guarded by `mutex_`).
94 /// The stop callback is armed before the node is queued, so
95 /// cancel_waiter must not unlink a node it never queued.
96 bool queued = false;
97 /// A readiness wait rather than an accept: completion
98 /// observes a pending connection without consuming it.
99 bool peek = false;
100 std::optional<std::stop_callback<waiter_canceller>> stop_cb;
101 };
102
103 int fd_ = -1;
104 uring_scheduler* sched_;
105 PeerService* peer_service_;
106 Endpoint local_endpoint_{};
107 mutable std::mutex mutex_;
108 intrusive_list<ready_fd_node> ready_fds_;
109 intrusive_list<waiter_node> waiters_;
110 /// Single parked readiness wait (guarded by `mutex_`). Multishot
111 /// accepting drains the kernel queue instantly, so a listener's
112 /// readiness lives in `ready_fds_`, not in `poll()`; the wait is
113 /// completed by the next delivery instead of a kernel poll.
114 waiter_node* read_wait_ = nullptr;
115 std::unique_ptr<uring_multi_accept_op> multi_op_;
116 bool closing_ = false;
117 /// Non-zero once an arming failed to reach the kernel (guarded by
118 /// `mutex_`). Nothing will ever deliver a connection through an
119 /// SQE the ring never took, so an accept reports this instead of
120 /// parking on a delivery that cannot come. Cleared by the next
121 /// arming that does reach the kernel.
122 int arm_err_ = 0;
123 /// Bumped whenever an arming is retired. A re-arm posted for an
124 /// earlier generation must not resubmit: `multi_op_` now names a
125 /// different op, and resubmitting a live one would alias a single
126 /// `user_data` across two kernel armings. The re-arm's check is
127 /// not atomic with its submit — a retirement landing between them
128 /// (concurrent `assign()` on another thread) can still double-arm;
129 /// closing that needs a generation-aware submit.
130 std::atomic<std::uint64_t> arm_generation_{0};
131
132 private:
133 // CRTP ctor private + Derived friended so the base cannot be
134 // constructed except as a CRTP base of Derived
135 // (clang-tidy bugprone-crtp-constructor-accessibility).
136 friend Derived;
137 272x uring_multishot_acceptor_base(
138 uring_scheduler& sched, PeerService& peer_svc) noexcept
139 272x : sched_(&sched)
140 272x , peer_service_(&peer_svc)
141 {
142 272x }
143
144 public:
145 269x ~uring_multishot_acceptor_base() override
146 {
147 {
148 269x std::lock_guard lk(mutex_);
149 269x closing_ = true;
150 269x }
151 269x if (fd_ >= 0)
152 {
153 sched_->submit_cancel_by_fd(fd_);
154 ::close(fd_);
155 fd_ = -1;
156 }
157
158 // Drain parked accepted-connection fds unconditionally. These are
159 // distinct from the listener fd and can be present even when the
160 // service close() path already closed and cleared fd_ — that path
161 // does not touch ready_fds_, so the drain must run here.
162 269x intrusive_list<ready_fd_node> drained;
163 {
164 269x std::lock_guard lk(mutex_);
165 278x while (auto* r = ready_fds_.pop_front())
166 9x drained.push_back(r);
167 269x }
168 278x while (auto* r = drained.pop_front())
169 {
170 9x ::close(r->fd);
171 9x delete r;
172 }
173
174 // Break the multi_op_ → impl_ptr (shared_ptr<this>) cycle and
175 // drain pending CQEs so unique_ptr<multi_op_> can free safely.
176 269x if (multi_op_)
177 {
178 208x multi_op_->impl_ptr.reset();
179 208x sched_->drain_cqes_for(multi_op_.get());
180 }
181 538x }
182
183 170x Endpoint local_endpoint() const noexcept override
184 {
185 170x return local_endpoint_;
186 }
187
188 4875x bool is_open() const noexcept override
189 {
190 4875x return fd_ >= 0;
191 }
192
193 16x native_handle_type native_handle() const noexcept override
194 {
195 16x return fd_;
196 }
197
198 9x native_handle_type release_socket() noexcept override
199 {
200 // Mirror the service close() path: cancel the multishot SQE and
201 // break the multi_op_ -> impl_ptr (shared_ptr<this>) cycle that
202 // start_multishot established. Without this, the cycle keeps the
203 // acceptor and its multi_op_ alive after the caller takes the fd,
204 // which LeakSanitizer reports on process exit. Caller still owns
205 // the returned fd, so we do NOT ::close it here.
206 9x if (fd_ >= 0)
207 {
208 9x sched_->cancel_and_flush(fd_);
209 9x drain_waiters_only();
210 9x if (multi_op_)
211 9x multi_op_->impl_ptr.reset();
212 }
213 9x int fd = fd_;
214 9x fd_ = -1;
215 9x local_endpoint_ = Endpoint{};
216 9x return fd;
217 }
218
219 10x void cancel() noexcept override
220 {
221 10x drain_waiters_only();
222 10x if (fd_ >= 0)
223 10x sched_->submit_cancel_by_fd(fd_);
224 10x }
225
226 /// Drain queued waiters with operation_aborted but do NOT submit
227 /// any kernel cancel for the fd. Used by service close() paths
228 /// that have already submitted (or are about to submit) the
229 /// cancel-by-fd themselves via `cancel_and_flush`.
230 261x void drain_waiters_only() noexcept
231 {
232 261x intrusive_list<waiter_node> drained;
233 {
234 261x std::lock_guard lk(mutex_);
235 261x closing_ = true;
236 // Drain under the lock — the kernel cancel may not produce
237 // a !more CQE before the fd is closed, so we can't rely on
238 // on_accept_cqe_impl to surface operation_aborted.
239 273x while (auto* w = waiters_.pop_front())
240 12x drained.push_back(w);
241 261x if (read_wait_)
242 {
243 2x drained.push_back(read_wait_);
244 2x read_wait_ = nullptr;
245 }
246 261x }
247
248 275x while (auto* w = drained.pop_front())
249 {
250 14x w->stop_cb.reset();
251 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept destructor path: OOM => std::terminate is the intended behavior
252 14x auto* op = new uring_accept_op();
253 14x op->h = w->h;
254 14x op->ex = w->ex;
255 14x op->ec_out = w->ec_out;
256 14x op->impl_out = w->impl_out;
257 14x op->cancelled.store(true, std::memory_order_release);
258 14x delete w;
259 14x sched_->post(op);
260 14x sched_->work_finished();
261 }
262 261x }
263
264 /** Park a readiness wait, or complete it if a connection is
265 already queued.
266
267 Multishot accepting consumes the kernel queue as connections
268 arrive, so a poll on the listener never reports it readable;
269 readiness is the impl's ready queue plus future deliveries.
270 */
271 13x void park_read_wait(
272 std::coroutine_handle<> h,
273 capy::executor_ref ex,
274 std::stop_token const& token,
275 std::error_code* ec) noexcept
276 {
277 13x bool ready = false;
278 13x bool aborted = false;
279 {
280 13x std::lock_guard lk(mutex_);
281 13x if (closing_)
282 {
283 aborted = true;
284 }
285 13x else if (!ready_fds_.empty())
286 {
287 4x ready = true;
288 }
289 13x }
290 13x if (ready || aborted)
291 {
292 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept initiation path: OOM => std::terminate is the intended behavior
293 4x auto* op = new uring_accept_op();
294 4x op->h = h;
295 4x op->ex = ex;
296 4x op->ec_out = ec;
297 4x if (aborted)
298 op->cancelled.store(true, std::memory_order_release);
299 4x sched_->post(op);
300 4x return;
301 }
302
303 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept initiation path: OOM => std::terminate is the intended behavior
304 9x auto* w = new waiter_node{};
305 9x w->h = h;
306 9x w->ex = ex;
307 9x w->ec_out = ec;
308 9x w->owner = static_cast<Derived*>(this);
309 9x w->peek = true;
310
311 // Same protocol as accept parking: arm the callback before
312 // the node is visible and outside `mutex_` (a pre-stopped
313 // token invokes the canceller synchronously, and the
314 // canceller takes `mutex_`).
315 9x if (token.stop_possible())
316 3x w->stop_cb.emplace(token, waiter_canceller{w});
317
318 9x bool was_cancelled = false;
319 9x int arm_err = 0;
320 {
321 9x std::lock_guard lk(mutex_);
322 9x if (w->cancelled.load(std::memory_order_acquire) || closing_)
323 {
324 2x was_cancelled = true;
325 }
326 7x else if (ready_fds_.empty())
327 {
328 // Readiness here means a future delivery, which a failed
329 // arming has already ruled out: report it rather than
330 // wait on one.
331 7x arm_err = arm_err_;
332 7x if (arm_err == 0)
333 {
334 6x w->queued = true;
335 6x sched_->work_started();
336 6x read_wait_ = w;
337 6x return;
338 }
339 }
340 // else: a connection arrived while the callback was armed;
341 // complete as ready below.
342 9x }
343
344 3x w->stop_cb.reset();
345 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept initiation path: OOM => std::terminate is the intended behavior
346 3x auto* op = new uring_accept_op();
347 3x op->h = w->h;
348 3x op->ex = w->ex;
349 3x op->ec_out = w->ec_out;
350 3x op->err = arm_err;
351 3x if (was_cancelled)
352 2x op->cancelled.store(true, std::memory_order_release);
353 3x delete w;
354 3x sched_->post(op);
355 }
356
357 185x std::error_code set_option(
358 int level,
359 int optname,
360 void const* data,
361 std::size_t size) noexcept override
362 {
363 185x if (fd_ < 0)
364 return make_err(EBADF);
365 185x if (::setsockopt(
366 fd_, level, optname, reinterpret_cast<char const*>(data),
367 185x static_cast<socklen_t>(size)) < 0)
368 2x return make_err(errno);
369 183x return {};
370 }
371
372 std::error_code
373 6x get_option(int level, int optname, void* data, std::size_t* size)
374 const noexcept override
375 {
376 6x if (fd_ < 0)
377 return make_err(EBADF);
378 6x socklen_t len = static_cast<socklen_t>(*size);
379 12x if (::getsockopt(
380 6x fd_, level, optname, reinterpret_cast<char*>(data), &len) < 0)
381 2x return make_err(errno);
382 4x *size = static_cast<std::size_t>(len);
383 4x return {};
384 }
385
386 /** Retire the multishot op before the acceptor changes descriptor.
387
388 `cancel_and_flush` and `submit_cancel_by_fd` only submit: the
389 terminating CQE for the previous arming is still queued when
390 the caller returns. Left alone, a subsequent `start_multishot`
391 would alias one `user_data` across two kernel ops, and the
392 stale `!more` CQE would observe a cleared `closing_` and take
393 the re-arm branch.
394
395 Ownership therefore moves to the scheduler rather than being
396 drained here. Draining is a teardown-only tool — it consumes
397 CQEs without dispatching them, so on a live context it would
398 swallow unrelated ops' completions and park their coroutines
399 forever. Handing the op over keeps the normal run loop in
400 charge of every CQE, and the op stays allocated (so its
401 `user_data` stays reserved) until the kernel is done with it.
402
403 Safe with no op armed, and safe after `release_socket` left
404 `fd_` cleared with the op still in flight.
405 */
406 218x void retire_multishot() noexcept
407 {
408 // Every field below is read by the leader mid-dispatch, so all
409 // of it is published inside retire_op's ring_mutex_ critical
410 // section — including moving multi_op_ out, since
411 // on_accept_cqe_impl dereferences it for peer_storage. Taking
412 // the acceptor mutex_ too keeps the generation bump ordered
413 // against the re-arm path's check. Lock order is
414 // ring_mutex_ -> mutex_, the same order the dispatch path
415 // acquires them in.
416 218x sched_->retire_op(
417 218x multi_op_, [this](uring_multi_accept_op& op) noexcept {
418 std::lock_guard lk(mutex_);
419 op.impl_ptr.reset();
420 op.acceptor_impl = nullptr;
421 op.on_cqe = nullptr;
422 op.retire_func = &uring_multi_accept_op::do_retired_cqe;
423 arm_generation_.fetch_add(1, std::memory_order_acq_rel);
424 });
425 218x }
426
427 /** Take over an already-listening descriptor.
428
429 Clears the shutdown latch a previous release left behind and
430 discards connections parked from the replaced descriptor:
431 those belong to the socket the caller is handing away.
432
433 @pre `retire_multishot` has run, so no arming from a previous
434 descriptor is still in flight.
435
436 @param fd The adopted descriptor.
437 */
438 9x void adopt_listening_fd(int fd) noexcept
439 {
440 9x intrusive_list<ready_fd_node> stale;
441 {
442 9x std::lock_guard lk(mutex_);
443 9x fd_ = fd;
444 9x closing_ = false;
445 9x while (auto* r = ready_fds_.pop_front())
446 stale.push_back(r);
447 9x }
448 9x while (auto* r = stale.pop_front())
449 {
450 ::close(r->fd);
451 delete r;
452 }
453 9x }
454
455 /** Ready an acceptor whose own descriptor is about to be armed.
456
457 The open/bind/listen path reaches `start_multishot` without
458 going through `adopt_listening_fd`, and a released acceptor may
459 be opened and listened on again. Both of the things that path
460 would otherwise inherit belong to the descriptor the caller
461 already took away: the shutdown latch, which would have every
462 connection the kernel hands back closed on arrival, and the
463 arming still owed a terminal CQE, which a fresh submission
464 would alias by `user_data`.
465 */
466 /** Ready the acceptor for a listen-time arming.
467
468 Returns `false` when a live arming already covers the
469 descriptor: a re-listen only changes the backlog, and retiring
470 a live arming here would leave it un-cancelled in the kernel —
471 two armings on one listener, with the retired one's deliveries
472 closed on arrival.
473
474 An arming that failed to reach the kernel is not one of those.
475 It covers nothing, so a re-listen is the caller's way back and
476 has to be allowed through.
477 */
478 211x bool prepare_listen_arm() noexcept
479 {
480 211x bool submitted = true;
481 {
482 211x std::lock_guard lk(mutex_);
483 211x if (multi_op_ && !closing_ && arm_err_ == 0)
484 1x return false;
485 // The op behind a failed arming was never handed to the
486 // ring, so no CQE is owed for it and it is still ours.
487 210x submitted = (arm_err_ == 0);
488 211x }
489 // Retiring is for an op the kernel still holds: it parks the op
490 // in the scheduler until a terminal CQE releases it. One that
491 // was never submitted would wait there for a completion that
492 // cannot come, so it is reused in place instead.
493 210x if (submitted)
494 209x retire_multishot();
495 210x intrusive_list<ready_fd_node> stale;
496 {
497 210x std::lock_guard lk(mutex_);
498 210x closing_ = false;
499 // Deliveries queued by a released descriptor belong to
500 // the socket the caller took away, not to this listener.
501 211x while (auto* r = ready_fds_.pop_front())
502 1x stale.push_back(r);
503 210x }
504 211x while (auto* r = stale.pop_front())
505 {
506 1x ::close(r->fd);
507 1x delete r;
508 }
509 210x return true;
510 }
511
512 222x void start_multishot()
513 {
514 222x if (!multi_op_)
515 {
516 218x multi_op_ = std::make_unique<uring_multi_accept_op>();
517 218x multi_op_->listen_fd = fd_;
518 218x multi_op_->acceptor_impl = this;
519 218x multi_op_->on_cqe = &uring_multishot_acceptor_base::on_accept_cqe;
520 218x multi_op_->impl_ptr = this->shared_from_this();
521 }
522 else
523 {
524 // Reuse the existing op (re-arm path). Reset peer scratch
525 // so the kernel writes into a clean slot. listen_fd and
526 // impl_ptr are re-seeded so the op can never carry state
527 // from an arming that has since been torn down. `res` is
528 // one of those: an arming that failed to submit left
529 // -EAGAIN there, and the reader of `res` cannot tell a
530 // result the kernel wrote from one it did not.
531 4x multi_op_->peer_storage = sockaddr_storage{};
532 4x multi_op_->peer_len = sizeof(sockaddr_storage);
533 4x multi_op_->res = 0;
534 4x multi_op_->listen_fd = fd_;
535 4x multi_op_->impl_ptr = this->shared_from_this();
536 }
537
538 222x auto* op = multi_op_.get();
539 // Deliberately no work_started(): the multishot SQE is a persistent
540 // internal mechanism. User-visible work is tracked per-accept call.
541 // The try_ spelling is what says so: it keeps a failed submission
542 // off the scheduler's completion queue, which spends a
543 // work_finished() on everything it dispatches.
544 222x if (uring_try_submit_op(*sched_, op))
545 {
546 217x std::lock_guard lk(mutex_);
547 217x arm_err_ = 0;
548 217x return;
549 217x }
550 5x fail_arm(EAGAIN);
551 }
552
553 /** Report an arming that never reached the kernel.
554
555 No CQE can arrive for an SQE the ring never took, so an accept
556 parked on this acceptor would park for good. The error is
557 remembered for the accepts still to come and delivered now to
558 whoever is already parked, which is the same EAGAIN every other
559 op reports when the SQ stays full.
560
561 @param err The code to report to the accepts this arming can no
562 longer serve.
563 */
564 5x void fail_arm(int err) noexcept
565 {
566 5x intrusive_list<waiter_node> claimed;
567 {
568 5x std::lock_guard lk(mutex_);
569 5x arm_err_ = err;
570 // Claim each waiter the way a delivery does. One the
571 // canceller already claimed belongs to it: cancel_waiter
572 // is waiting on this mutex to unlink the node itself, so
573 // it has to still be in the list when it gets in.
574 5x intrusive_list<waiter_node> keep;
575 5x while (auto* w = waiters_.pop_front())
576 {
577 if (!w->cancelled.exchange(true, std::memory_order_acq_rel))
578 claimed.push_back(w);
579 else
580 keep.push_back(w);
581 }
582 5x while (auto* w = keep.pop_front())
583 waiters_.push_back(w);
584 5x if (read_wait_ &&
585 !read_wait_->cancelled.exchange(
586 true, std::memory_order_acq_rel))
587 {
588 claimed.push_back(read_wait_);
589 read_wait_ = nullptr;
590 }
591 5x }
592
593 5x while (auto* w = claimed.pop_front())
594 {
595 w->stop_cb.reset();
596 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept arming path: OOM => std::terminate is the intended behavior
597 auto* op = new uring_accept_op();
598 op->h = w->h;
599 op->ex = w->ex;
600 op->ec_out = w->ec_out;
601 op->impl_out = w->impl_out;
602 op->err = err;
603 delete w;
604 sched_->post(op);
605 sched_->work_finished(); // balance the waiter's work_started
606 }
607 5x }
608
609 /// Pull a parked fd or queue a waiter — used by Derived::accept().
610 /// Either case ends with the calling coroutine suspending; the
611 /// caller returns `std::noop_coroutine()` unconditionally.
612 3286x void dispatch_or_queue(
613 std::coroutine_handle<> h,
614 capy::executor_ref ex,
615 std::stop_token const& token,
616 std::error_code* ec,
617 io_object::implementation** impl_out)
618 {
619 3286x sockaddr_storage peer_storage{};
620 3286x socklen_t peer_len = sizeof(peer_storage);
621 3286x int accepted_fd = ::accept4(
622 fd_, reinterpret_cast<sockaddr*>(&peer_storage), &peer_len,
623 SOCK_NONBLOCK | SOCK_CLOEXEC);
624 3286x if (accepted_fd >= 0)
625 {
626 auto* op = new uring_accept_op();
627 op->h = h;
628 op->ex = ex;
629 op->ec_out = ec;
630 op->impl_out = impl_out;
631 op->peer_service = peer_service_;
632 op->adopt_fn = &Derived::adopt_thunk;
633 op->accepted_fd = accepted_fd;
634 op->peer_storage = peer_storage;
635 op->peer_len = peer_len;
636 sched_->post(op);
637 3285x return;
638 }
639 // accept4 returned <0 — only EAGAIN/EWOULDBLOCK should fall
640 // through to the parked/waiter path. Other errors (EBADF, etc.)
641 // surface through the existing scheduler-completion path so the
642 // user sees them via the op's ec_out. Build an op with `err`
643 // set so do_handler delivers make_err(err).
644 3286x if (errno != EAGAIN && errno != EWOULDBLOCK)
645 {
646 2x int saved_errno = errno;
647 2x auto* op = new uring_accept_op();
648 2x op->h = h;
649 2x op->ex = ex;
650 2x op->ec_out = ec;
651 2x op->impl_out = impl_out;
652 2x op->err = saved_errno;
653 2x sched_->post(op);
654 2x return;
655 }
656
657 3284x uring_accept_op* ready_op = nullptr;
658 {
659 3284x std::lock_guard lk(mutex_);
660 3284x if (auto* r = ready_fds_.pop_front())
661 {
662 17x ready_op = new uring_accept_op();
663 17x ready_op->h = h;
664 17x ready_op->ex = ex;
665 17x ready_op->ec_out = ec;
666 17x ready_op->impl_out = impl_out;
667 17x ready_op->peer_service = peer_service_;
668 17x ready_op->adopt_fn = &Derived::adopt_thunk;
669 17x ready_op->accepted_fd = r->fd;
670 17x ready_op->peer_storage = r->peer;
671 17x ready_op->peer_len = r->peer_len;
672 17x delete r;
673 }
674 3284x }
675 3284x if (ready_op)
676 {
677 // Post outside the lock — acceptor mutex_ must never be
678 // held while dispatch_mutex_ is acquired by sched_->post().
679 17x sched_->post(ready_op);
680 17x return;
681 }
682
683 3267x auto* w = new waiter_node{};
684 3267x w->h = h;
685 3267x w->ex = ex;
686 3267x w->ec_out = ec;
687 3267x w->impl_out = impl_out;
688 3267x w->owner = static_cast<Derived*>(this);
689
690 // Arm the stop callback before the node is visible in
691 // `waiters_` and outside `mutex_`: an already-stopped token
692 // invokes the canceller synchronously from emplace, and
693 // cancel_waiter takes `mutex_` (self-deadlock if held).
694 // Arming pre-queue also keeps the CQE handler from claiming
695 // and deleting a node whose callback is not yet constructed.
696 3267x if (token.stop_possible())
697 33x w->stop_cb.emplace(token, waiter_canceller{w});
698
699 3267x bool was_cancelled = false;
700 {
701 3267x std::lock_guard lk(mutex_);
702 3267x if (w->cancelled.load(std::memory_order_acquire))
703 {
704 // Canceller already fired (pre-stopped token); it saw
705 // queued == false and left completion to us.
706 6x was_cancelled = true;
707 }
708 3261x else if (auto* r = ready_fds_.pop_front())
709 {
710 // A connection arrived while the callback was armed;
711 // prefer it over parking the waiter behind it.
712 ready_op = new uring_accept_op();
713 ready_op->h = h;
714 ready_op->ex = ex;
715 ready_op->ec_out = ec;
716 ready_op->impl_out = impl_out;
717 ready_op->peer_service = peer_service_;
718 ready_op->adopt_fn = &Derived::adopt_thunk;
719 ready_op->accepted_fd = r->fd;
720 ready_op->peer_storage = r->peer;
721 ready_op->peer_len = r->peer_len;
722 delete r;
723 }
724 3261x else if (arm_err_ != 0)
725 {
726 // No arming reached the kernel, so no CQE will deliver
727 // a connection: parking here would park for good.
728 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept accept path: OOM => std::terminate is the intended behavior
729 1x ready_op = new uring_accept_op();
730 1x ready_op->h = h;
731 1x ready_op->ex = ex;
732 1x ready_op->ec_out = ec;
733 1x ready_op->impl_out = impl_out;
734 1x ready_op->err = arm_err_;
735 }
736 else
737 {
738 3260x w->queued = true;
739 3260x sched_->work_started();
740 3260x waiters_.push_back(w);
741 3260x return;
742 }
743 3267x }
744
745 7x if (was_cancelled)
746 {
747 6x auto* op = new uring_accept_op();
748 6x op->h = w->h;
749 6x op->ex = w->ex;
750 6x op->ec_out = w->ec_out;
751 6x op->impl_out = w->impl_out;
752 6x op->cancelled.store(true, std::memory_order_release);
753 6x w->stop_cb.reset();
754 6x delete w;
755 6x sched_->post(op);
756 6x return;
757 }
758
759 1x w->stop_cb.reset();
760 1x delete w;
761 1x sched_->post(ready_op);
762 }
763
764 24x void cancel_waiter(waiter_node* w) noexcept
765 {
766 {
767 24x std::lock_guard lk(mutex_);
768 24x if (closing_)
769 return; // on_accept_cqe_impl will drain with closing_ set
770 24x if (!w->queued)
771 8x return; // not queued yet; the parking path observes
772 // `cancelled` and completes the op
773 16x if (w->peek)
774 {
775 1x if (read_wait_ != w)
776 return; // already claimed by a delivery
777 1x read_wait_ = nullptr;
778 }
779 else
780 {
781 15x waiters_.remove(w);
782 }
783 24x }
784 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — stop-token callback: noexcept, OOM => std::terminate is the intended behavior
785 16x auto* op = new uring_accept_op();
786 16x op->h = w->h;
787 16x op->ex = w->ex;
788 16x op->ec_out = w->ec_out;
789 16x op->impl_out = w->impl_out;
790 16x op->cancelled.store(true, std::memory_order_release);
791 16x delete w;
792 // post() increments outstanding_work_; balances the work_started()
793 // from accept() when the waiter was queued.
794 16x sched_->post(op);
795 16x sched_->work_finished(); // balance the work_started() from accept()
796 }
797
798 private:
799 static void
800 3273x on_accept_cqe(void* self_ptr, int new_fd, int err, bool more) noexcept
801 {
802 3273x static_cast<Derived*>(self_ptr)->on_accept_cqe_impl(new_fd, err, more);
803 3273x }
804
805 protected:
806 3273x void on_accept_cqe_impl(int new_fd, int err, bool more) noexcept
807 {
808 3273x bool was_closing = false;
809 3273x waiter_node* matched = nullptr;
810 3273x waiter_node* claimed_peek = nullptr;
811 3273x intrusive_list<waiter_node> closing_waiters;
812 {
813 3273x std::lock_guard lk(mutex_);
814 3273x was_closing = closing_;
815 3276x if (!was_closing && new_fd >= 0 && read_wait_ &&
816 3x !read_wait_->cancelled.exchange(
817 true, std::memory_order_acq_rel))
818 {
819 // A parked readiness wait observes the delivery
820 // without consuming it; the connection still flows
821 // to a waiter or the ready queue below.
822 3x claimed_peek = read_wait_;
823 3x read_wait_ = nullptr;
824 }
825 3273x if (was_closing)
826 {
827 13x if (new_fd >= 0)
828 ::close(new_fd);
829 13x if (!more)
830 {
831 // Collect waiters to drain after the lock is released.
832 13x while (auto* w = waiters_.pop_front())
833 closing_waiters.push_back(w);
834 }
835 }
836 3260x else if (!waiters_.empty())
837 {
838 // Claim the head waiter atomically. If the canceller
839 // already won the race (cancelled was already true),
840 // leave the waiter in the list for cancel_waiter to
841 // remove and dispatch with operation_aborted; park the
842 // new_fd so the next waiter consumes it.
843 3233x auto* head_w = waiters_.front();
844 3233x if (!head_w->cancelled.exchange(
845 true, std::memory_order_acq_rel))
846 {
847 3233x waiters_.pop_front();
848 3233x matched = head_w;
849 }
850 else if (new_fd >= 0)
851 {
852 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — CQE handler: noexcept, OOM => std::terminate is the intended behavior
853 auto* node = new ready_fd_node{};
854 node->fd = new_fd;
855 node->peer = multi_op_->peer_storage;
856 node->peer_len = multi_op_->peer_len;
857 ready_fds_.push_back(node);
858 }
859 }
860 27x else if (new_fd >= 0)
861 {
862 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — CQE handler: noexcept, OOM => std::terminate is the intended behavior
863 27x auto* node = new ready_fd_node{};
864 27x node->fd = new_fd;
865 27x node->peer = multi_op_->peer_storage;
866 27x node->peer_len = multi_op_->peer_len;
867 27x ready_fds_.push_back(node);
868 }
869 3273x }
870
871 3273x if (claimed_peek)
872 {
873 3x claimed_peek->stop_cb.reset();
874 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — CQE handler: noexcept, OOM => std::terminate is the intended behavior
875 3x auto* op = new uring_accept_op();
876 3x op->h = claimed_peek->h;
877 3x op->ex = claimed_peek->ex;
878 3x op->ec_out = claimed_peek->ec_out;
879 3x delete claimed_peek;
880 3x sched_->post(op);
881 3x sched_->work_finished(); // balance the parking work_started
882 }
883
884 3273x if (matched)
885 {
886 3233x matched->stop_cb.reset();
887 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — CQE handler: noexcept, OOM => std::terminate is the intended behavior
888 3233x auto* op = new uring_accept_op();
889 3233x op->h = matched->h;
890 3233x op->ex = matched->ex;
891 3233x op->ec_out = matched->ec_out;
892 3233x op->impl_out = matched->impl_out;
893 3233x op->peer_service = peer_service_;
894 3233x op->adopt_fn = &Derived::adopt_thunk;
895 3233x if (err)
896 {
897 4x op->err = err;
898 }
899 3229x else if (new_fd >= 0)
900 {
901 3229x op->accepted_fd = new_fd;
902 3229x op->peer_storage = multi_op_->peer_storage;
903 3229x op->peer_len = multi_op_->peer_len;
904 }
905 3233x delete matched;
906 3233x sched_->post(op);
907 3233x sched_->work_finished(); // balance waiter's work_started
908 }
909
910 3273x while (auto* w = closing_waiters.pop_front())
911 {
912 w->stop_cb.reset();
913 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — CQE handler shutdown path: noexcept, OOM => std::terminate is the intended behavior
914 auto* op = new uring_accept_op();
915 op->h = w->h;
916 op->ex = w->ex;
917 op->ec_out = w->ec_out;
918 op->impl_out = w->impl_out;
919 op->cancelled.store(true, std::memory_order_release);
920 delete w;
921 sched_->post(op);
922 sched_->work_finished(); // balance waiter's work_started
923 }
924
925 3273x if (!more && !was_closing)
926 {
927 // Re-arm: kernel terminated multishot non-fatally.
928 struct rearm_op final : scheduler_op
929 {
930 std::shared_ptr<Derived> self_;
931 std::uint64_t generation_;
932 rearm_op(
933 std::shared_ptr<Derived> s,
934 std::uint64_t generation) noexcept
935 : self_(std::move(s))
936 , generation_(generation)
937 {
938 }
939
940 void operator()() override
941 {
942 auto self = std::move(self_);
943 auto generation = generation_;
944 delete this;
945 {
946 std::lock_guard lk(self->mutex_);
947 if (self->closing_)
948 return;
949 // The arming this was posted for may have been
950 // retired by assign() in the meantime. multi_op_
951 // then names a different, already-armed op, and
952 // resubmitting it would alias one user_data
953 // across two kernel armings — a use-after-free
954 // once the first terminal CQE frees the op.
955 if (self->arm_generation_.load(
956 std::memory_order_acquire) != generation)
957 return;
958 }
959 self->start_multishot();
960 }
961
962 void destroy() override
963 {
964 delete this;
965 }
966 };
967 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — CQE handler re-arm: noexcept, OOM => std::terminate is the intended behavior
968 6x sched_->post(new rearm_op(
969 this->shared_from_this(),
970 arm_generation_.load(std::memory_order_acquire)));
971 }
972 3273x }
973 };
974
975 template<class Derived, class ImplBase, class Endpoint, class PeerService>
976 inline void
977 24x uring_multishot_acceptor_base<Derived, ImplBase, Endpoint, PeerService>::
978 waiter_canceller::operator()() const noexcept
979 {
980 24x if (w->cancelled.exchange(true, std::memory_order_acq_rel))
981 return;
982 24x w->owner->cancel_waiter(w);
983 }
984
985 } // namespace boost::corosio::detail
986
987 #endif // BOOST_COROSIO_HAS_URING
988
989 #endif // BOOST_COROSIO_NATIVE_DETAIL_URING_URING_MULTISHOT_ACCEPTOR_HPP
990