include/boost/corosio/native/detail/io_uring/io_uring_multishot_acceptor.hpp

80.2% Lines (296/0/369) 95.1% List of functions (39/0/41)
io_uring_multishot_acceptor.hpp
f(x) Functions (41)
Function Calls Lines Blocks
boost::corosio::detail::fd_is_listening(int) :55 8x 80.0% 83.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::io_uring_multishot_acceptor_base(boost::corosio::detail::io_uring_scheduler&, boost::corosio::detail::io_uring_local_stream_service&) :131 45x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::io_uring_multishot_acceptor_base(boost::corosio::detail::io_uring_scheduler&, boost::corosio::detail::io_uring_tcp_service&) :131 165x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::~io_uring_multishot_acceptor_base() :139 45x 85.0% 90.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::~io_uring_multishot_acceptor_base() :139 165x 85.0% 90.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::local_endpoint() const :177 5x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::local_endpoint() const :177 129x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::is_open() const :182 204x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::is_open() const :182 3259x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::native_handle() const :187 3x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::native_handle() const :187 11x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::release_socket() :192 5x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::release_socket() :192 4x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::cancel() :213 3x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::cancel() :213 4x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::drain_waiters_only() :224 40x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::drain_waiters_only() :224 161x 90.9% 92.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::park_read_wait(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::stop_token const&, std::error_code*) :265 2x 53.3% 52.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::park_read_wait(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::stop_token const&, std::error_code*) :265 4x 71.1% 65.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::set_option(int, int, void const*, unsigned long) :342 2x 100.0% 89.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::set_option(int, int, void const*, unsigned long) :342 147x 100.0% 89.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::get_option(int, int, void*, unsigned long*) const :354 2x 100.0% 90.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::get_option(int, int, void*, unsigned long*) const :354 4x 100.0% 90.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::retire_multishot() :387 27x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::retire_multishot() :387 138x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::adopt_listening_fd(int) :419 3x 75.0% 60.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::adopt_listening_fd(int) :419 5x 75.0% 60.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::prepare_listen_arm() :455 24x 75.0% 68.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::prepare_listen_arm() :455 134x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::start_multishot() :480 27x 71.4% 63.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::start_multishot() :480 138x 71.4% 63.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_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**) :512 16x 71.6% 72.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_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**) :512 2211x 71.6% 72.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::cancel_waiter(boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::waiter_node*) :652 2x 85.0% 83.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::cancel_waiter(boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::waiter_node*) :652 12x 85.0% 83.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::on_accept_cqe(void*, int, int, bool) :686 19x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::on_accept_cqe(void*, int, int, bool) :686 2202x 100.0% 100.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::on_accept_cqe_impl(int, int, bool) :694 19x 73.7% 57.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::on_accept_cqe_impl(int, int, bool) :694 2202x 73.7% 57.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::io_uring_local_stream_service>::waiter_canceller::operator()() const :860 0 0.0% 0.0% boost::corosio::detail::io_uring_multishot_acceptor_base<boost::corosio::detail::io_uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::io_uring_tcp_service>::waiter_canceller::operator()() const :860 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_IO_URING_IO_URING_MULTISHOT_ACCEPTOR_HPP
11 #define BOOST_COROSIO_NATIVE_DETAIL_IO_URING_IO_URING_MULTISHOT_ACCEPTOR_HPP
12
13 #include <boost/corosio/detail/platform.hpp>
14
15 #if BOOST_COROSIO_HAS_IO_URING
16
17 #include <liburing.h>
18
19 #include <boost/corosio/detail/intrusive.hpp>
20 #include <boost/corosio/native/detail/io_uring/io_uring_acceptor_ops.hpp>
21 #include <boost/corosio/native/detail/io_uring/io_uring_buffer.hpp>
22 #include <boost/corosio/native/detail/io_uring/io_uring_op.hpp>
23 #include <boost/corosio/native/detail/io_uring/io_uring_scheduler.hpp>
24 #include <boost/corosio/native/detail/io_uring/io_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 8x fd_is_listening(int fd) noexcept
56 {
57 8x int accepting = 0;
58 8x socklen_t alen = sizeof(accepting);
59 8x if (::getsockopt(fd, SOL_SOCKET, SO_ACCEPTCONN, &accepting, &alen) != 0)
60 return true;
61 8x return accepting != 0;
62 }
63
64 template<class Derived, class ImplBase, class Endpoint, class PeerService>
65 class io_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 io_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 /// Bumped whenever an arming is retired. A re-arm posted for an
118 /// earlier generation must not resubmit: `multi_op_` now names a
119 /// different op, and resubmitting a live one would alias a single
120 /// `user_data` across two kernel armings. The re-arm's check is
121 /// not atomic with its submit — a retirement landing between them
122 /// (concurrent `assign()` on another thread) can still double-arm;
123 /// closing that needs a generation-aware submit.
124 std::atomic<std::uint64_t> arm_generation_{0};
125
126 private:
127 // CRTP ctor private + Derived friended so the base cannot be
128 // constructed except as a CRTP base of Derived
129 // (clang-tidy bugprone-crtp-constructor-accessibility).
130 friend Derived;
131 210x io_uring_multishot_acceptor_base(
132 io_uring_scheduler& sched, PeerService& peer_svc) noexcept
133 210x : sched_(&sched)
134 210x , peer_service_(&peer_svc)
135 210x {}
136
137 public:
138
139 210x ~io_uring_multishot_acceptor_base() override
140 {
141 {
142 210x std::lock_guard lk(mutex_);
143 210x closing_ = true;
144 210x }
145 210x if (fd_ >= 0)
146 {
147 sched_->submit_cancel_by_fd(fd_);
148 ::close(fd_);
149 fd_ = -1;
150 }
151
152 // Drain parked accepted-connection fds unconditionally. These are
153 // distinct from the listener fd and can be present even when the
154 // service close() path already closed and cleared fd_ — that path
155 // does not touch ready_fds_, so the drain must run here.
156 210x intrusive_list<ready_fd_node> drained;
157 {
158 210x std::lock_guard lk(mutex_);
159 212x while (auto* r = ready_fds_.pop_front())
160 2x drained.push_back(r);
161 210x }
162 212x while (auto* r = drained.pop_front())
163 {
164 2x ::close(r->fd);
165 2x delete r;
166 }
167
168 // Break the multi_op_ → impl_ptr (shared_ptr<this>) cycle and
169 // drain pending CQEs so unique_ptr<multi_op_> can free safely.
170 210x if (multi_op_)
171 {
172 158x multi_op_->impl_ptr.reset();
173 158x sched_->drain_cqes_for(multi_op_.get());
174 }
175 420x }
176
177 134x Endpoint local_endpoint() const noexcept override
178 {
179 134x return local_endpoint_;
180 }
181
182 3463x bool is_open() const noexcept override
183 {
184 3463x return fd_ >= 0;
185 }
186
187 14x native_handle_type native_handle() const noexcept override
188 {
189 14x return fd_;
190 }
191
192 9x native_handle_type release_socket() noexcept override
193 {
194 // Mirror the service close() path: cancel the multishot SQE and
195 // break the multi_op_ -> impl_ptr (shared_ptr<this>) cycle that
196 // start_multishot established. Without this, the cycle keeps the
197 // acceptor and its multi_op_ alive after the caller takes the fd,
198 // which LeakSanitizer reports on process exit. Caller still owns
199 // the returned fd, so we do NOT ::close it here.
200 9x if (fd_ >= 0)
201 {
202 9x sched_->cancel_and_flush(fd_);
203 9x drain_waiters_only();
204 9x if (multi_op_)
205 9x multi_op_->impl_ptr.reset();
206 }
207 9x int fd = fd_;
208 9x fd_ = -1;
209 9x local_endpoint_ = Endpoint{};
210 9x return fd;
211 }
212
213 7x void cancel() noexcept override
214 {
215 7x drain_waiters_only();
216 7x if (fd_ >= 0)
217 7x sched_->submit_cancel_by_fd(fd_);
218 7x }
219
220 /// Drain queued waiters with operation_aborted but do NOT submit
221 /// any kernel cancel for the fd. Used by service close() paths
222 /// that have already submitted (or are about to submit) the
223 /// cancel-by-fd themselves via `cancel_and_flush`.
224 201x void drain_waiters_only() noexcept
225 {
226 201x intrusive_list<waiter_node> drained;
227 {
228 201x std::lock_guard lk(mutex_);
229 201x closing_ = true;
230 // Drain under the lock — the kernel cancel may not produce
231 // a !more CQE before the fd is closed, so we can't rely on
232 // on_accept_cqe_impl to surface operation_aborted.
233 204x while (auto* w = waiters_.pop_front())
234 3x drained.push_back(w);
235 201x if (read_wait_)
236 {
237 1x drained.push_back(read_wait_);
238 1x read_wait_ = nullptr;
239 }
240 201x }
241
242 205x while (auto* w = drained.pop_front())
243 {
244 4x w->stop_cb.reset();
245 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept destructor path: OOM => std::terminate is the intended behavior
246 4x auto* op = new uring_accept_op();
247 4x op->h = w->h;
248 4x op->ex = w->ex;
249 4x op->ec_out = w->ec_out;
250 4x op->impl_out = w->impl_out;
251 4x op->cancelled.store(true, std::memory_order_release);
252 4x delete w;
253 4x sched_->post(op);
254 4x sched_->work_finished();
255 }
256 201x }
257
258 /** Park a readiness wait, or complete it if a connection is
259 already queued.
260
261 Multishot accepting consumes the kernel queue as connections
262 arrive, so a poll on the listener never reports it readable;
263 readiness is the impl's ready queue plus future deliveries.
264 */
265 6x void park_read_wait(
266 std::coroutine_handle<> h,
267 capy::executor_ref ex,
268 std::stop_token const& token,
269 std::error_code* ec) noexcept
270 {
271 6x bool ready = false;
272 6x bool aborted = false;
273 {
274 6x std::lock_guard lk(mutex_);
275 6x if (closing_)
276 {
277 aborted = true;
278 }
279 6x else if (!ready_fds_.empty())
280 {
281 2x ready = true;
282 }
283 6x }
284 6x if (ready || aborted)
285 {
286 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept initiation path: OOM => std::terminate is the intended behavior
287 2x auto* op = new uring_accept_op();
288 2x op->h = h;
289 2x op->ex = ex;
290 2x op->ec_out = ec;
291 2x if (aborted)
292 op->cancelled.store(true, std::memory_order_release);
293 2x sched_->post(op);
294 2x return;
295 }
296
297 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept initiation path: OOM => std::terminate is the intended behavior
298 4x auto* w = new waiter_node{};
299 4x w->h = h;
300 4x w->ex = ex;
301 4x w->ec_out = ec;
302 4x w->owner = static_cast<Derived*>(this);
303 4x w->peek = true;
304
305 // Same protocol as accept parking: arm the callback before
306 // the node is visible and outside `mutex_` (a pre-stopped
307 // token invokes the canceller synchronously, and the
308 // canceller takes `mutex_`).
309 4x if (token.stop_possible())
310 w->stop_cb.emplace(token, waiter_canceller{w});
311
312 4x bool was_cancelled = false;
313 {
314 4x std::lock_guard lk(mutex_);
315 4x if (w->cancelled.load(std::memory_order_acquire) || closing_)
316 {
317 was_cancelled = true;
318 }
319 4x else if (ready_fds_.empty())
320 {
321 4x w->queued = true;
322 4x sched_->work_started();
323 4x read_wait_ = w;
324 4x return;
325 }
326 // else: a connection arrived while the callback was armed;
327 // complete as ready below.
328 4x }
329
330 w->stop_cb.reset();
331 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept initiation path: OOM => std::terminate is the intended behavior
332 auto* op = new uring_accept_op();
333 op->h = w->h;
334 op->ex = w->ex;
335 op->ec_out = w->ec_out;
336 if (was_cancelled)
337 op->cancelled.store(true, std::memory_order_release);
338 delete w;
339 sched_->post(op);
340 }
341
342 149x std::error_code set_option(
343 int level, int optname,
344 void const* data, std::size_t size) noexcept override
345 {
346 149x if (fd_ < 0) return make_err(EBADF);
347 149x if (::setsockopt(fd_, level, optname,
348 reinterpret_cast<char const*>(data),
349 149x static_cast<socklen_t>(size)) < 0)
350 2x return make_err(errno);
351 147x return {};
352 }
353
354 6x std::error_code get_option(
355 int level, int optname,
356 void* data, std::size_t* size) const noexcept override
357 {
358 6x if (fd_ < 0) return make_err(EBADF);
359 6x socklen_t len = static_cast<socklen_t>(*size);
360 6x if (::getsockopt(fd_, level, optname,
361 6x reinterpret_cast<char*>(data), &len) < 0)
362 2x return make_err(errno);
363 4x *size = static_cast<std::size_t>(len);
364 4x return {};
365 }
366
367 /** Retire the multishot op before the acceptor changes descriptor.
368
369 `cancel_and_flush` and `submit_cancel_by_fd` only submit: the
370 terminating CQE for the previous arming is still queued when
371 the caller returns. Left alone, a subsequent `start_multishot`
372 would alias one `user_data` across two kernel ops, and the
373 stale `!more` CQE would observe a cleared `closing_` and take
374 the re-arm branch.
375
376 Ownership therefore moves to the scheduler rather than being
377 drained here. Draining is a teardown-only tool — it consumes
378 CQEs without dispatching them, so on a live context it would
379 swallow unrelated ops' completions and park their coroutines
380 forever. Handing the op over keeps the normal run loop in
381 charge of every CQE, and the op stays allocated (so its
382 `user_data` stays reserved) until the kernel is done with it.
383
384 Safe with no op armed, and safe after `release_socket` left
385 `fd_` cleared with the op still in flight.
386 */
387 165x void retire_multishot() noexcept
388 {
389 // Every field below is read by the leader mid-dispatch, so all
390 // of it is published inside retire_op's ring_mutex_ critical
391 // section — including moving multi_op_ out, since
392 // on_accept_cqe_impl dereferences it for peer_storage. Taking
393 // the acceptor mutex_ too keeps the generation bump ordered
394 // against the re-arm path's check. Lock order is
395 // ring_mutex_ -> mutex_, the same order the dispatch path
396 // acquires them in.
397 165x sched_->retire_op(
398 165x multi_op_, [this](uring_multi_accept_op& op) noexcept {
399 std::lock_guard lk(mutex_);
400 op.impl_ptr.reset();
401 op.acceptor_impl = nullptr;
402 op.on_cqe = nullptr;
403 op.retire_func = &uring_multi_accept_op::do_retired_cqe;
404 arm_generation_.fetch_add(1, std::memory_order_acq_rel);
405 });
406 165x }
407
408 /** Take over an already-listening descriptor.
409
410 Clears the shutdown latch a previous release left behind and
411 discards connections parked from the replaced descriptor:
412 those belong to the socket the caller is handing away.
413
414 @pre `retire_multishot` has run, so no arming from a previous
415 descriptor is still in flight.
416
417 @param fd The adopted descriptor.
418 */
419 8x void adopt_listening_fd(int fd) noexcept
420 {
421 8x intrusive_list<ready_fd_node> stale;
422 {
423 8x std::lock_guard lk(mutex_);
424 8x fd_ = fd;
425 8x closing_ = false;
426 8x while (auto* r = ready_fds_.pop_front())
427 stale.push_back(r);
428 8x }
429 8x while (auto* r = stale.pop_front())
430 {
431 ::close(r->fd);
432 delete r;
433 }
434 8x }
435
436 /** Ready an acceptor whose own descriptor is about to be armed.
437
438 The open/bind/listen path reaches `start_multishot` without
439 going through `adopt_listening_fd`, and a released acceptor may
440 be opened and listened on again. Both of the things that path
441 would otherwise inherit belong to the descriptor the caller
442 already took away: the shutdown latch, which would have every
443 connection the kernel hands back closed on arrival, and the
444 arming still owed a terminal CQE, which a fresh submission
445 would alias by `user_data`.
446 */
447 /** Ready the acceptor for a listen-time arming.
448
449 Returns `false` when a live arming already covers the
450 descriptor: a re-listen only changes the backlog, and retiring
451 a live arming here would leave it un-cancelled in the kernel —
452 two armings on one listener, with the retired one's deliveries
453 closed on arrival.
454 */
455 158x bool prepare_listen_arm() noexcept
456 {
457 {
458 158x std::lock_guard lk(mutex_);
459 158x if (multi_op_ && !closing_)
460 1x return false;
461 158x }
462 157x retire_multishot();
463 157x intrusive_list<ready_fd_node> stale;
464 {
465 157x std::lock_guard lk(mutex_);
466 157x closing_ = false;
467 // Deliveries queued by a released descriptor belong to
468 // the socket the caller took away, not to this listener.
469 158x while (auto* r = ready_fds_.pop_front())
470 1x stale.push_back(r);
471 157x }
472 158x while (auto* r = stale.pop_front())
473 {
474 1x ::close(r->fd);
475 1x delete r;
476 }
477 157x return true;
478 }
479
480 165x void start_multishot()
481 {
482 165x if (!multi_op_)
483 {
484 165x multi_op_ = std::make_unique<uring_multi_accept_op>();
485 165x multi_op_->listen_fd = fd_;
486 165x multi_op_->acceptor_impl = this;
487 165x multi_op_->on_cqe =
488 &io_uring_multishot_acceptor_base::on_accept_cqe;
489 165x multi_op_->impl_ptr = this->shared_from_this();
490 }
491 else
492 {
493 // Reuse the existing op (re-arm path). Reset peer scratch
494 // so the kernel writes into a clean slot. listen_fd and
495 // impl_ptr are re-seeded so the op can never carry state
496 // from an arming that has since been torn down.
497 multi_op_->peer_storage = sockaddr_storage{};
498 multi_op_->peer_len = sizeof(sockaddr_storage);
499 multi_op_->listen_fd = fd_;
500 multi_op_->impl_ptr = this->shared_from_this();
501 }
502
503 165x auto* op = multi_op_.get();
504 165x io_uring_submit_op(*sched_, op);
505 // Deliberately no work_started(): the multishot SQE is a persistent
506 // internal mechanism. User-visible work is tracked per-accept call.
507 165x }
508
509 /// Pull a parked fd or queue a waiter — used by Derived::accept().
510 /// Either case ends with the calling coroutine suspending; the
511 /// caller returns `std::noop_coroutine()` unconditionally.
512 2227x void dispatch_or_queue(
513 std::coroutine_handle<> h,
514 capy::executor_ref ex,
515 std::stop_token const& token,
516 std::error_code* ec,
517 io_object::implementation** impl_out)
518 {
519 2227x sockaddr_storage peer_storage{};
520 2227x socklen_t peer_len = sizeof(peer_storage);
521 2227x int accepted_fd = ::accept4(fd_,
522 reinterpret_cast<sockaddr*>(&peer_storage), &peer_len,
523 SOCK_NONBLOCK | SOCK_CLOEXEC);
524 2227x if (accepted_fd >= 0)
525 {
526 auto* op = new uring_accept_op();
527 op->h = h;
528 op->ex = ex;
529 op->ec_out = ec;
530 op->impl_out = impl_out;
531 op->peer_service = peer_service_;
532 op->adopt_fn = &Derived::adopt_thunk;
533 op->accepted_fd = accepted_fd;
534 op->peer_storage = peer_storage;
535 op->peer_len = peer_len;
536 sched_->post(op);
537 2227x return;
538 }
539 // accept4 returned <0 — only EAGAIN/EWOULDBLOCK should fall
540 // through to the parked/waiter path. Other errors (EBADF, etc.)
541 // surface through the existing scheduler-completion path so the
542 // user sees them via the op's ec_out. Build an op with `err`
543 // set so do_handler delivers make_err(err).
544 2227x if (errno != EAGAIN && errno != EWOULDBLOCK)
545 {
546 2x int saved_errno = errno;
547 2x auto* op = new uring_accept_op();
548 2x op->h = h;
549 2x op->ex = ex;
550 2x op->ec_out = ec;
551 2x op->impl_out = impl_out;
552 2x op->err = saved_errno;
553 2x sched_->post(op);
554 2x return;
555 }
556
557 2225x uring_accept_op* ready_op = nullptr;
558 {
559 2225x std::lock_guard lk(mutex_);
560 2225x if (auto* r = ready_fds_.pop_front())
561 {
562 13x ready_op = new uring_accept_op();
563 13x ready_op->h = h;
564 13x ready_op->ex = ex;
565 13x ready_op->ec_out = ec;
566 13x ready_op->impl_out = impl_out;
567 13x ready_op->peer_service = peer_service_;
568 13x ready_op->adopt_fn = &Derived::adopt_thunk;
569 13x ready_op->accepted_fd = r->fd;
570 13x ready_op->peer_storage = r->peer;
571 13x ready_op->peer_len = r->peer_len;
572 13x delete r;
573 }
574 2225x }
575 2225x if (ready_op)
576 {
577 // Post outside the lock — acceptor mutex_ must never be
578 // held while dispatch_mutex_ is acquired by sched_->post().
579 13x sched_->post(ready_op);
580 13x return;
581 }
582
583 2212x auto* w = new waiter_node{};
584 2212x w->h = h;
585 2212x w->ex = ex;
586 2212x w->ec_out = ec;
587 2212x w->impl_out = impl_out;
588 2212x w->owner = static_cast<Derived*>(this);
589
590 // Arm the stop callback before the node is visible in
591 // `waiters_` and outside `mutex_`: an already-stopped token
592 // invokes the canceller synchronously from emplace, and
593 // cancel_waiter takes `mutex_` (self-deadlock if held).
594 // Arming pre-queue also keeps the CQE handler from claiming
595 // and deleting a node whose callback is not yet constructed.
596 2212x if (token.stop_possible())
597 25x w->stop_cb.emplace(token, waiter_canceller{w});
598
599 2212x bool was_cancelled = false;
600 {
601 2212x std::lock_guard lk(mutex_);
602 2212x if (w->cancelled.load(std::memory_order_acquire))
603 {
604 // Canceller already fired (pre-stopped token); it saw
605 // queued == false and left completion to us.
606 2x was_cancelled = true;
607 }
608 2210x else if (auto* r = ready_fds_.pop_front())
609 {
610 // A connection arrived while the callback was armed;
611 // prefer it over parking the waiter behind it.
612 ready_op = new uring_accept_op();
613 ready_op->h = h;
614 ready_op->ex = ex;
615 ready_op->ec_out = ec;
616 ready_op->impl_out = impl_out;
617 ready_op->peer_service = peer_service_;
618 ready_op->adopt_fn = &Derived::adopt_thunk;
619 ready_op->accepted_fd = r->fd;
620 ready_op->peer_storage = r->peer;
621 ready_op->peer_len = r->peer_len;
622 delete r;
623 }
624 else
625 {
626 2210x w->queued = true;
627 2210x sched_->work_started();
628 2210x waiters_.push_back(w);
629 2210x return;
630 }
631 2212x }
632
633 2x if (was_cancelled)
634 {
635 2x auto* op = new uring_accept_op();
636 2x op->h = w->h;
637 2x op->ex = w->ex;
638 2x op->ec_out = w->ec_out;
639 2x op->impl_out = w->impl_out;
640 2x op->cancelled.store(true, std::memory_order_release);
641 2x w->stop_cb.reset();
642 2x delete w;
643 2x sched_->post(op);
644 2x return;
645 }
646
647 w->stop_cb.reset();
648 delete w;
649 sched_->post(ready_op);
650 }
651
652 14x void cancel_waiter(waiter_node* w) noexcept
653 {
654 {
655 14x std::lock_guard lk(mutex_);
656 14x if (closing_) return; // on_accept_cqe_impl will drain with closing_ set
657 14x if (!w->queued)
658 2x return; // not queued yet; the parking path observes
659 // `cancelled` and completes the op
660 12x if (w->peek)
661 {
662 if (read_wait_ != w)
663 return; // already claimed by a delivery
664 read_wait_ = nullptr;
665 }
666 else
667 {
668 12x waiters_.remove(w);
669 }
670 14x }
671 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — stop-token callback: noexcept, OOM => std::terminate is the intended behavior
672 12x auto* op = new uring_accept_op();
673 12x op->h = w->h;
674 12x op->ex = w->ex;
675 12x op->ec_out = w->ec_out;
676 12x op->impl_out = w->impl_out;
677 12x op->cancelled.store(true, std::memory_order_release);
678 12x delete w;
679 // post() increments outstanding_work_; balances the work_started()
680 // from accept() when the waiter was queued.
681 12x sched_->post(op);
682 12x sched_->work_finished(); // balance the work_started() from accept()
683 }
684
685 private:
686 2221x static void on_accept_cqe(
687 void* self_ptr, int new_fd, int err, bool more) noexcept
688 {
689 static_cast<Derived*>(self_ptr)
690 2221x ->on_accept_cqe_impl(new_fd, err, more);
691 2221x }
692
693 protected:
694 2221x void on_accept_cqe_impl(int new_fd, int err, bool more) noexcept
695 {
696 2221x bool was_closing = false;
697 2221x waiter_node* matched = nullptr;
698 2221x waiter_node* claimed_peek = nullptr;
699 2221x intrusive_list<waiter_node> closing_waiters;
700 {
701 2221x std::lock_guard lk(mutex_);
702 2221x was_closing = closing_;
703 2224x if (!was_closing && new_fd >= 0 && read_wait_ &&
704 3x !read_wait_->cancelled.exchange(
705 true, std::memory_order_acq_rel))
706 {
707 // A parked readiness wait observes the delivery
708 // without consuming it; the connection still flows
709 // to a waiter or the ready queue below.
710 3x claimed_peek = read_wait_;
711 3x read_wait_ = nullptr;
712 }
713 2221x if (was_closing)
714 {
715 10x if (new_fd >= 0)
716 ::close(new_fd);
717 10x if (!more)
718 {
719 // Collect waiters to drain after the lock is released.
720 10x while (auto* w = waiters_.pop_front())
721 closing_waiters.push_back(w);
722 }
723 }
724 2211x else if (!waiters_.empty())
725 {
726 // Claim the head waiter atomically. If the canceller
727 // already won the race (cancelled was already true),
728 // leave the waiter in the list for cancel_waiter to
729 // remove and dispatch with operation_aborted; park the
730 // new_fd so the next waiter consumes it.
731 2195x auto* head_w = waiters_.front();
732 2195x if (!head_w->cancelled.exchange(
733 true, std::memory_order_acq_rel))
734 {
735 2195x waiters_.pop_front();
736 2195x matched = head_w;
737 }
738 else if (new_fd >= 0)
739 {
740 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — CQE handler: noexcept, OOM => std::terminate is the intended behavior
741 auto* node = new ready_fd_node{};
742 node->fd = new_fd;
743 node->peer = multi_op_->peer_storage;
744 node->peer_len = multi_op_->peer_len;
745 ready_fds_.push_back(node);
746 }
747 }
748 16x else if (new_fd >= 0)
749 {
750 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — CQE handler: noexcept, OOM => std::terminate is the intended behavior
751 16x auto* node = new ready_fd_node{};
752 16x node->fd = new_fd;
753 16x node->peer = multi_op_->peer_storage;
754 16x node->peer_len = multi_op_->peer_len;
755 16x ready_fds_.push_back(node);
756 }
757 2221x }
758
759 2221x if (claimed_peek)
760 {
761 3x claimed_peek->stop_cb.reset();
762 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — CQE handler: noexcept, OOM => std::terminate is the intended behavior
763 3x auto* op = new uring_accept_op();
764 3x op->h = claimed_peek->h;
765 3x op->ex = claimed_peek->ex;
766 3x op->ec_out = claimed_peek->ec_out;
767 3x delete claimed_peek;
768 3x sched_->post(op);
769 3x sched_->work_finished(); // balance the parking work_started
770 }
771
772 2221x if (matched)
773 {
774 2195x matched->stop_cb.reset();
775 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — CQE handler: noexcept, OOM => std::terminate is the intended behavior
776 2195x auto* op = new uring_accept_op();
777 2195x op->h = matched->h;
778 2195x op->ex = matched->ex;
779 2195x op->ec_out = matched->ec_out;
780 2195x op->impl_out = matched->impl_out;
781 2195x op->peer_service = peer_service_;
782 2195x op->adopt_fn = &Derived::adopt_thunk;
783 2195x if (err)
784 {
785 op->err = err;
786 }
787 2195x else if (new_fd >= 0)
788 {
789 2195x op->accepted_fd = new_fd;
790 2195x op->peer_storage = multi_op_->peer_storage;
791 2195x op->peer_len = multi_op_->peer_len;
792 }
793 2195x delete matched;
794 2195x sched_->post(op);
795 2195x sched_->work_finished(); // balance waiter's work_started
796 }
797
798 2221x while (auto* w = closing_waiters.pop_front())
799 {
800 w->stop_cb.reset();
801 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — CQE handler shutdown path: noexcept, OOM => std::terminate is the intended behavior
802 auto* op = new uring_accept_op();
803 op->h = w->h;
804 op->ex = w->ex;
805 op->ec_out = w->ec_out;
806 op->impl_out = w->impl_out;
807 op->cancelled.store(true, std::memory_order_release);
808 delete w;
809 sched_->post(op);
810 sched_->work_finished(); // balance waiter's work_started
811 }
812
813 2221x if (!more && !was_closing)
814 {
815 // Re-arm: kernel terminated multishot non-fatally.
816 struct rearm_op final : scheduler_op
817 {
818 std::shared_ptr<Derived> self_;
819 std::uint64_t generation_;
820 rearm_op(
821 std::shared_ptr<Derived> s,
822 std::uint64_t generation) noexcept
823 : self_(std::move(s))
824 , generation_(generation) {}
825
826 void operator()() override
827 {
828 auto self = std::move(self_);
829 auto generation = generation_;
830 delete this;
831 {
832 std::lock_guard lk(self->mutex_);
833 if (self->closing_)
834 return;
835 // The arming this was posted for may have been
836 // retired by assign() in the meantime. multi_op_
837 // then names a different, already-armed op, and
838 // resubmitting it would alias one user_data
839 // across two kernel armings — a use-after-free
840 // once the first terminal CQE frees the op.
841 if (self->arm_generation_.load(
842 std::memory_order_acquire) != generation)
843 return;
844 }
845 self->start_multishot();
846 }
847
848 void destroy() override { delete this; }
849 };
850 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — CQE handler re-arm: noexcept, OOM => std::terminate is the intended behavior
851 sched_->post(new rearm_op(
852 this->shared_from_this(),
853 arm_generation_.load(std::memory_order_acquire)));
854 }
855 2221x }
856 };
857
858 template<class Derived, class ImplBase, class Endpoint, class PeerService>
859 inline void
860 14x io_uring_multishot_acceptor_base<Derived, ImplBase, Endpoint, PeerService>
861 ::waiter_canceller::operator()() const noexcept
862 {
863 14x if (w->cancelled.exchange(true, std::memory_order_acq_rel))
864 return;
865 14x w->owner->cancel_waiter(w);
866 }
867
868 } // namespace boost::corosio::detail
869
870 #endif // BOOST_COROSIO_HAS_IO_URING
871
872 #endif // BOOST_COROSIO_NATIVE_DETAIL_IO_URING_IO_URING_MULTISHOT_ACCEPTOR_HPP
873