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

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