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

94.4% Lines (1169/29/1267) 100.0% List of functions (98/1/99)
uring_types.hpp
f(x) Functions (99)
Function Calls Lines Blocks
boost::corosio::detail::uring_tcp_socket::uring_tcp_socket(boost::corosio::detail::uring_tcp_service&, boost::corosio::detail::uring_scheduler&) :138 9842x 100.0% 100.0% boost::corosio::detail::uring_tcp_socket::~uring_tcp_socket() :145 9839x 80.0% 88.0% boost::corosio::detail::uring_tcp_socket::read_some(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :156 186679x 91.3% 84.0% boost::corosio::detail::uring_tcp_socket::write_some(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :239 186577x 91.3% 83.0% boost::corosio::detail::uring_tcp_socket::connect(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::endpoint, std::stop_token, std::error_code*) :322 3267x 85.7% 78.0% boost::corosio::detail::uring_tcp_socket::wait(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::wait_type, std::stop_token, std::error_code*) :369 20x 100.0% 93.0% boost::corosio::detail::uring_tcp_socket::shutdown(boost::corosio::shutdown_type) :402 4x 100.0% 100.0% boost::corosio::detail::uring_tcp_socket::release_socket() :412 1x 100.0% 100.0% boost::corosio::detail::uring_tcp_socket::cancel() :428 75x 100.0% 100.0% boost::corosio::detail::uring_tcp_socket::close_socket() :438 16355x 100.0% 100.0% boost::corosio::detail::uring_tcp_socket::local_endpoint() const :452 29x 100.0% 100.0% boost::corosio::detail::uring_tcp_socket::remote_endpoint() const :476 30x 100.0% 100.0% boost::corosio::detail::uring_tcp_service::uring_tcp_service(boost::capy::execution_context&) :515 799x 100.0% 100.0% boost::corosio::detail::uring_tcp_service::open_socket(boost::corosio::tcp_socket::implementation&, int, int, int) :532 3285x 100.0% 85.0% boost::corosio::detail::uring_tcp_service::assign_socket(boost::corosio::tcp_socket::implementation&, int) :573 13x 100.0% 100.0% boost::corosio::detail::uring_tcp_service::bind_socket(boost::corosio::tcp_socket::implementation&, boost::corosio::endpoint) :622 9x 100.0% 100.0% boost::corosio::detail::uring_tcp_service::adopt_fd(int, boost::corosio::endpoint const&) :651 3233x 100.0% 74.0% boost::corosio::detail::uring_tcp_acceptor::uring_tcp_acceptor(boost::corosio::detail::uring_tcp_acceptor_service&, boost::corosio::detail::uring_scheduler&, boost::corosio::detail::uring_tcp_service&) :698 217x 100.0% 100.0% boost::corosio::detail::uring_tcp_acceptor::accept(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::stop_token, std::error_code*, boost::corosio::io_object::implementation**) :706 3266x 100.0% 100.0% boost::corosio::detail::uring_tcp_acceptor::wait(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::wait_type, std::stop_token, std::error_code*) :717 13x 86.7% 64.0% boost::corosio::detail::uring_tcp_acceptor::adopt_thunk(void*, int, sockaddr_storage const&, unsigned int) :774 3233x 100.0% 100.0% boost::corosio::detail::uring_tcp_acceptor_service::uring_tcp_acceptor_service(boost::capy::execution_context&) :811 799x 100.0% 88.0% boost::corosio::detail::uring_tcp_acceptor_service::shutdown() :817 799x 100.0% 80.0% boost::corosio::detail::uring_tcp_acceptor_service::construct() :832 217x 100.0% 71.0% boost::corosio::detail::uring_tcp_acceptor_service::destroy(boost::corosio::io_object::implementation*) :842 215x 83.3% 64.0% boost::corosio::detail::uring_tcp_acceptor_service::close(boost::corosio::io_object::handle&) :852 417x 100.0% 100.0% boost::corosio::detail::uring_tcp_acceptor_service::open_acceptor_socket(boost::corosio::tcp_acceptor::implementation&, int, int, int) :884 204x 100.0% 85.0% boost::corosio::detail::uring_tcp_acceptor_service::assign_socket(boost::corosio::tcp_acceptor::implementation&, int) :920 10x 100.0% 100.0% boost::corosio::detail::uring_tcp_acceptor_service::bind_acceptor(boost::corosio::tcp_acceptor::implementation&, boost::corosio::endpoint) :963 197x 100.0% 100.0% boost::corosio::detail::uring_tcp_acceptor_service::listen_acceptor(boost::corosio::tcp_acceptor::implementation&, int) :989 186x 100.0% 100.0% boost::corosio::detail::uring_local_stream_socket::uring_local_stream_socket(boost::corosio::detail::uring_local_stream_service&, boost::corosio::detail::uring_scheduler&) :1061 159x 100.0% 100.0% boost::corosio::detail::uring_local_stream_socket::~uring_local_stream_socket() :1068 157x 80.0% 88.0% boost::corosio::detail::uring_local_stream_socket::read_some(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :1079 64x 91.3% 84.0% boost::corosio::detail::uring_local_stream_socket::write_some(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :1162 60x 91.3% 83.0% boost::corosio::detail::uring_local_stream_socket::connect(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::local_endpoint, std::stop_token, std::error_code*) :1245 38x 85.7% 78.0% boost::corosio::detail::uring_local_stream_socket::wait(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::wait_type, std::stop_token, std::error_code*) :1292 10x 100.0% 89.0% boost::corosio::detail::uring_local_stream_socket::shutdown(boost::corosio::shutdown_type) :1326 3x 100.0% 100.0% boost::corosio::detail::uring_local_stream_socket::release_socket() :1336 2x 100.0% 100.0% boost::corosio::detail::uring_local_stream_socket::cancel() :1350 3x 100.0% 100.0% boost::corosio::detail::uring_local_stream_socket::close_socket() :1358 258x 100.0% 100.0% boost::corosio::detail::uring_local_stream_socket::remote_endpoint() const :1370 2x 100.0% 100.0% boost::corosio::detail::uring_local_stream_service::uring_local_stream_service(boost::capy::execution_context&) :1410 799x 100.0% 100.0% boost::corosio::detail::uring_local_stream_service::open_socket(boost::corosio::local_stream_socket::implementation&, int, int, int) :1429 26x 100.0% 80.0% boost::corosio::detail::uring_local_stream_service::assign_socket(boost::corosio::local_stream_socket::implementation&, int) :1461 74x 100.0% 100.0% boost::corosio::detail::uring_local_stream_service::adopt_fd(int, boost::corosio::local_endpoint const&) :1507 12x 100.0% 77.0% boost::corosio::detail::uring_local_stream_acceptor::uring_local_stream_acceptor(boost::corosio::detail::uring_local_stream_acceptor_service&, boost::corosio::detail::uring_scheduler&, boost::corosio::detail::uring_local_stream_service&) :1549 55x 100.0% 100.0% boost::corosio::detail::uring_local_stream_acceptor::accept(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::stop_token, std::error_code*, boost::corosio::io_object::implementation**) :1557 20x 100.0% 100.0% boost::corosio::detail::uring_local_stream_acceptor::wait(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::wait_type, std::stop_token, std::error_code*) :1568 6x 86.7% 64.0% boost::corosio::detail::uring_local_stream_acceptor::adopt_thunk(void*, int, sockaddr_storage const&, unsigned int) :1625 12x 100.0% 100.0% boost::corosio::detail::uring_local_stream_acceptor_service::uring_local_stream_acceptor_service(boost::capy::execution_context&) :1663 799x 100.0% 88.0% boost::corosio::detail::uring_local_stream_acceptor_service::shutdown() :1669 799x 100.0% 80.0% boost::corosio::detail::uring_local_stream_acceptor_service::construct() :1684 55x 100.0% 71.0% boost::corosio::detail::uring_local_stream_acceptor_service::destroy(boost::corosio::io_object::implementation*) :1694 54x 83.3% 64.0% boost::corosio::detail::uring_local_stream_acceptor_service::close(boost::corosio::io_object::handle&) :1704 92x 100.0% 100.0% boost::corosio::detail::uring_local_stream_acceptor_service::open_acceptor_socket(boost::corosio::local_stream_acceptor::implementation&, int, int, int) :1732 45x 100.0% 80.0% boost::corosio::detail::uring_local_stream_acceptor_service::assign_socket(boost::corosio::local_stream_acceptor::implementation&, int) :1761 5x 100.0% 100.0% boost::corosio::detail::uring_local_stream_acceptor_service::bind_acceptor(boost::corosio::local_stream_acceptor::implementation&, boost::corosio::local_endpoint) :1804 38x 100.0% 100.0% boost::corosio::detail::uring_local_stream_acceptor_service::listen_acceptor(boost::corosio::local_stream_acceptor::implementation&, int) :1831 31x 100.0% 100.0% boost::corosio::detail::uring_udp_socket::uring_udp_socket(boost::corosio::detail::uring_udp_service&, boost::corosio::detail::uring_scheduler&) :1907 132x 100.0% 100.0% boost::corosio::detail::uring_udp_socket::~uring_udp_socket() :1914 132x 80.0% 88.0% boost::corosio::detail::uring_udp_socket::send_to(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, boost::corosio::endpoint, int, std::stop_token, std::error_code*, unsigned long*) :1925 50x 100.0% 100.0% boost::corosio::detail::uring_udp_socket::recv_from(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, boost::corosio::endpoint*, int, std::stop_token, std::error_code*, unsigned long*) :1940 97x 100.0% 100.0% boost::corosio::detail::uring_udp_socket::send(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, int, std::stop_token, std::error_code*, unsigned long*) :1954 149x 100.0% 100.0% boost::corosio::detail::uring_udp_socket::recv(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, int, std::stop_token, std::error_code*, unsigned long*) :1967 67x 100.0% 100.0% boost::corosio::detail::uring_udp_socket::connect(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::endpoint, std::stop_token, std::error_code*) :1980 61x 85.7% 78.0% boost::corosio::detail::uring_udp_socket::wait(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::wait_type, std::stop_token, std::error_code*) :2027 11x 100.0% 93.0% boost::corosio::detail::uring_udp_socket::shutdown(boost::corosio::shutdown_type) :2063 2x 100.0% 100.0% boost::corosio::detail::uring_udp_socket::release_socket() :2070 2x 100.0% 100.0% boost::corosio::detail::uring_udp_socket::cancel() :2084 6x 100.0% 100.0% boost::corosio::detail::uring_udp_socket::close_socket() :2092 249x 100.0% 100.0% boost::corosio::detail::uring_udp_socket::remote_endpoint() const :2104 3x 100.0% 100.0% boost::corosio::detail::uring_udp_socket::submit_send(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, unsigned int, sockaddr_storage const&, int, std::stop_token const&, std::error_code*, unsigned long*) :2110 199x 92.2% 84.0% boost::corosio::detail::uring_udp_socket::submit_recv(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, bool, boost::corosio::endpoint*, int, std::stop_token const&, std::error_code*, unsigned long*) :2200 164x 93.2% 86.0% boost::corosio::detail::uring_udp_socket::write_ip_source(void*, sockaddr_storage const&, unsigned int) :2306 6x 100.0% 100.0% boost::corosio::detail::uring_udp_service::uring_udp_service(boost::capy::execution_context&) :2348 799x 100.0% 100.0% boost::corosio::detail::uring_udp_service::open_datagram_socket(boost::corosio::udp_socket::implementation&, int, int, int) :2365 117x 100.0% 85.0% boost::corosio::detail::uring_udp_service::assign_socket(boost::corosio::udp_socket::implementation&, int) :2403 11x 100.0% 100.0% boost::corosio::detail::uring_udp_service::bind_datagram(boost::corosio::udp_socket::implementation&, boost::corosio::endpoint) :2449 69x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_socket::uring_local_datagram_socket(boost::corosio::detail::uring_local_datagram_service&, boost::corosio::detail::uring_scheduler&) :2514 126x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_socket::~uring_local_datagram_socket() :2521 125x 80.0% 88.0% boost::corosio::detail::uring_local_datagram_socket::send_to(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, boost::corosio::local_endpoint, int, std::stop_token, std::error_code*, unsigned long*) :2532 77x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_socket::recv_from(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, boost::corosio::local_endpoint*, int, std::stop_token, std::error_code*, unsigned long*) :2547 67x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_socket::send(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, int, std::stop_token, std::error_code*, unsigned long*) :2561 68x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_socket::recv(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, int, std::stop_token, std::error_code*, unsigned long*) :2574 69x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_socket::connect(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::local_endpoint, std::stop_token, std::error_code*) :2587 47x 85.7% 78.0% boost::corosio::detail::uring_local_datagram_socket::wait(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::wait_type, std::stop_token, std::error_code*) :2634 7x 100.0% 93.0% boost::corosio::detail::uring_local_datagram_socket::shutdown(boost::corosio::shutdown_type) :2668 3x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_socket::release_socket() :2678 2x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_socket::cancel() :2692 2x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_socket::close_socket() :2700 227x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_socket::remote_endpoint() const :2712 1x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_socket::bind(boost::corosio::local_endpoint) :2720 boost::corosio::detail::uring_local_datagram_socket::submit_send(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, unsigned int, sockaddr_storage const&, int, std::stop_token const&, std::error_code*, unsigned long*) :2737 145x 90.2% 81.0% boost::corosio::detail::uring_local_datagram_socket::submit_recv(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, bool, boost::corosio::local_endpoint*, int, std::stop_token const&, std::error_code*, unsigned long*) :2827 136x 83.1% 78.0% boost::corosio::detail::uring_local_datagram_socket::write_local_source(void*, sockaddr_storage const&, unsigned int) :2935 22x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_service::uring_local_datagram_service(boost::capy::execution_context&) :2977 799x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_service::open_socket(boost::corosio::local_datagram_socket::implementation&, int, int, int) :2996 50x 100.0% 80.0% boost::corosio::detail::uring_local_datagram_service::assign_socket(boost::corosio::local_datagram_socket::implementation&, int) :3028 60x 100.0% 100.0% boost::corosio::detail::uring_local_datagram_service::bind_socket(boost::corosio::local_datagram_socket::implementation&, boost::corosio::local_endpoint) :3069 37x 100.0% 100.0%
Line TLA Hits Source Code
1 //
2 // Copyright (c) 2026 Steve Gerbino
3 // Copyright (c) 2026 Michael Vandeberg
4 //
5 // Distributed under the Boost Software License, Version 1.0. (See accompanying
6 // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7 //
8 // Official repository: https://github.com/cppalliance/corosio
9 //
10
11 #ifndef BOOST_COROSIO_NATIVE_DETAIL_URING_URING_TYPES_HPP
12 #define BOOST_COROSIO_NATIVE_DETAIL_URING_URING_TYPES_HPP
13
14 #include <boost/corosio/detail/platform.hpp>
15
16 #if BOOST_COROSIO_HAS_URING
17
18 #include <boost/corosio/detail/intrusive.hpp>
19 #include <boost/corosio/native/detail/uring/uring_acceptor_ops.hpp>
20 #include <boost/corosio/native/detail/uring/uring_buffer.hpp>
21 #include <boost/corosio/native/detail/uring/uring_dgram_ops.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_multishot_acceptor.hpp>
25 #include <boost/corosio/native/detail/uring/uring_socket_ops.hpp>
26 #include <boost/corosio/native/detail/uring/uring_socket_service_base.hpp>
27 #include <boost/corosio/native/detail/native_socket_base.hpp>
28 #include <boost/corosio/native/detail/make_err.hpp>
29 #include <boost/corosio/native/detail/msg_flags.hpp>
30 #include <boost/corosio/native/detail/validate_fd.hpp>
31 #include <boost/corosio/detail/local_datagram_service.hpp>
32 #include <boost/corosio/detail/local_stream_acceptor_service.hpp>
33 #include <boost/corosio/detail/local_stream_service.hpp>
34 #include <boost/corosio/detail/tcp_acceptor_service.hpp>
35 #include <boost/corosio/detail/tcp_service.hpp>
36 #include <boost/corosio/detail/udp_service.hpp>
37 #include <boost/corosio/local_endpoint.hpp>
38 #include <boost/corosio/local_datagram_socket.hpp>
39 #include <boost/corosio/local_stream_acceptor.hpp>
40 #include <boost/corosio/local_stream_socket.hpp>
41 #include <boost/corosio/tcp_acceptor.hpp>
42 #include <boost/corosio/tcp_socket.hpp>
43 #include <boost/corosio/udp_socket.hpp>
44
45 #include <memory>
46 #include <mutex>
47 #include <optional>
48 #include <unordered_map>
49 #include <vector>
50
51 #include <fcntl.h>
52 #include <netinet/in.h>
53 #include <sys/socket.h>
54 #include <sys/un.h>
55 #include <unistd.h>
56
57 namespace boost::corosio::detail {
58
59 class uring_tcp_service;
60 class uring_tcp_acceptor_service; // Task 18
61 class uring_local_stream_service;
62 class uring_local_stream_acceptor_service;
63 class uring_udp_service;
64 class uring_local_datagram_service;
65
66 /** TCP socket implementation for io_uring.
67
68 Implements `tcp_socket::implementation` using a proactor model:
69 read, write, and connect operations are submitted to the kernel
70 via `uring_submit_op` and complete through the ring's CQE path.
71
72 The object is always owned by a `shared_ptr` managed by the service.
73 In-flight ops hold an additional `shared_ptr` copy (`impl_ptr`) so
74 the kernel's user-data pointer remains valid until the CQE arrives.
75
76 @par Thread Safety
77 Distinct objects: Safe.
78 Shared objects: Unsafe. A socket must not have two operations of
79 the same type in flight simultaneously.
80 */
81 class BOOST_COROSIO_DECL uring_tcp_socket final
82 : public native_socket_base<
83 uring_tcp_socket,
84 tcp_socket::implementation,
85 endpoint>
86 {
87 friend uring_tcp_service;
88
89 int family_ = AF_UNSPEC; // cached at open_socket
90 uring_scheduler* sched_ = nullptr;
91 [[maybe_unused]] uring_tcp_service* svc_ = nullptr;
92
93 // fd_ and local_endpoint_ are provided by native_socket_base (the
94 // readiness/completion-agnostic socket base shared with the reactor
95 // sockets). native_handle()/is_open()/set_option()/get_option() come
96 // from there too; local_endpoint() is overridden below for lazy
97 // getsockname resolution.
98 // Three-state machine for the local endpoint:
99 // unresolved — never set; accessor returns default endpoint
100 // (open-but-unbound socket, failed-connect, etc.)
101 // lazy_pending — set by adopt_fd to signal "this socket has an
102 // authoritative local endpoint that hasn't been
103 // fetched yet"; accessor will getsockname on
104 // first read
105 // resolved — local_endpoint_ is authoritative; accessor
106 // returns the cached value
107 enum class endpoint_state : int
108 {
109 unresolved,
110 lazy_pending,
111 resolved
112 };
113 mutable std::atomic<endpoint_state> local_endpoint_state_{
114 endpoint_state::unresolved};
115 endpoint remote_endpoint_;
116
117 // Per-fd op slots — embedded to eliminate per-call heap allocation.
118 // Single-pending invariant per slot: at most one read, write, or
119 // connect in flight on this socket at any time (the awaitable
120 // contract).
121 uring_read_op rd_;
122 uring_write_op wr_;
123 uring_connect_op conn_;
124 uring_wait_op wait_op_;
125
126 mutable detail::speculative_state spec_;
127
128 public:
129 /** Construct with service and scheduler references.
130
131 Both refs must outlive this socket. `sched_` and `svc_` are
132 intentionally separate so service subclasses can pass a
133 different scheduler if needed.
134
135 @param svc The owning service (Task 13).
136 @param sched The io_uring scheduler owned by the context.
137 */
138 9842x explicit uring_tcp_socket(
139 uring_tcp_service& svc, uring_scheduler& sched) noexcept
140 19684x : sched_(&sched)
141 9842x , svc_(&svc)
142 {
143 9842x }
144
145 9839x ~uring_tcp_socket() override
146 9839x {
147 9839x if (fd_ >= 0)
148 ::close(
149 fd_); // LCOV_EXCL_LINE backstop: close_socket() clears fd_ before destroy
150 9839x }
151
152 // ----------------------------------------------------------------
153 // io_stream::implementation
154 // ----------------------------------------------------------------
155
156 186679x std::coroutine_handle<> read_some(
157 std::coroutine_handle<> h,
158 capy::executor_ref ex,
159 buffer_param buffers,
160 std::stop_token token,
161 std::error_code* ec,
162 std::size_t* bytes) override
163 {
164 iovec iovecs[uring_max_iov];
165 186679x int iovec_count = copy_to_iovec(buffers, iovecs);
166 186679x bool stop_now = token.stop_possible() && token.stop_requested();
167 186679x bool empty_buf = (iovec_count == 0);
168
169 186679x ssize_t n = 0;
170 186679x int err = 0;
171 186679x bool have_sync_res = stop_now || empty_buf;
172 186679x if (!have_sync_res && spec_.may_speculate_read())
173 {
174 do
175 {
176 186466x n = ::readv(fd_, iovecs, iovec_count);
177 }
178 186466x while (n < 0 && errno == EINTR);
179 186466x if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
180 {
181 186426x have_sync_res = true;
182 186426x if (n < 0)
183 2x err = errno;
184 // Speculative read produced a definitive answer (data
185 // or non-EAGAIN error); reset the failure streak so a
186 // burst of past EAGAINs doesn't latch perma-off when
187 // the workload is in fact speculation-friendly.
188 186426x if (n >= 0)
189 186424x spec_.on_read_success();
190 }
191 else
192 {
193 40x spec_.on_read_exhausted();
194 }
195 }
196
197 186679x if (have_sync_res)
198 {
199 186450x if (sched_->try_consume_inline_budget())
200 {
201 175487x decode_io_result(
202 2x ec, stop_now, err ? make_err(err) : std::error_code{},
203 175487x /*is_read=*/true, n < 0 ? 0u : static_cast<std::size_t>(n),
204 empty_buf);
205 175487x if (bytes)
206 175487x *bytes = (n < 0) ? 0u : static_cast<std::size_t>(n);
207 175487x rd_.cont.h = h;
208 175487x return dispatch_coro(ex, rd_.cont);
209 }
210 10963x rd_.prepare(
211 21926x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_,
212 buffers, token);
213 10963x if (stop_now)
214 1x rd_.cancelled.store(true, std::memory_order_release);
215 else
216 10962x rd_.res = (n < 0) ? -err : static_cast<int>(n);
217 10963x sched_->work_started();
218 {
219 10963x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
220 10963x sched_->push_completed_locked(&rd_);
221 10963x }
222 10963x return std::noop_coroutine();
223 }
224
225 229x rd_.prepare(
226 458x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_, buffers,
227 token);
228 229x sched_->work_started();
229 229x if (rd_.cancelled.load(std::memory_order_acquire))
230 {
231 uring_scheduler::lock_type lock(sched_->dispatch_mutex());
232 sched_->push_completed_locked(&rd_);
233 return std::noop_coroutine();
234 }
235 229x uring_submit_op(*sched_, &rd_);
236 229x return std::noop_coroutine();
237 }
238
239 186577x std::coroutine_handle<> write_some(
240 std::coroutine_handle<> h,
241 capy::executor_ref ex,
242 buffer_param buffers,
243 std::stop_token token,
244 std::error_code* ec,
245 std::size_t* bytes) override
246 {
247 iovec iovecs[uring_max_iov];
248 186577x int iovec_count = copy_to_iovec(buffers, iovecs);
249 186577x bool stop_now = token.stop_possible() && token.stop_requested();
250 186577x bool empty_buf = (iovec_count == 0);
251
252 186577x ssize_t n = 0;
253 186577x int err = 0;
254 186577x bool have_sync_res = stop_now || empty_buf;
255 186577x if (!have_sync_res && spec_.may_speculate_write())
256 {
257 186553x msghdr msg{};
258 186553x msg.msg_iov = iovecs;
259 186553x msg.msg_iovlen = static_cast<decltype(msg.msg_iovlen)>(iovec_count);
260 do
261 {
262 186553x n = ::sendmsg(fd_, &msg, MSG_NOSIGNAL);
263 }
264 186553x while (n < 0 && errno == EINTR);
265 186553x if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
266 {
267 186549x have_sync_res = true;
268 186549x if (n < 0)
269 4x err = errno;
270 }
271 else
272 {
273 4x spec_.on_write_exhausted();
274 }
275 }
276
277 186577x if (have_sync_res)
278 {
279 186572x if (sched_->try_consume_inline_budget())
280 {
281 175611x decode_io_result(
282 4x ec, stop_now, err ? make_err(err) : std::error_code{},
283 /*is_read=*/false, /*bytes=*/0, /*empty_buffer=*/false);
284 175607x if (bytes)
285 175607x *bytes = (n < 0) ? 0u : static_cast<std::size_t>(n);
286 175607x wr_.cont.h = h;
287 175607x return dispatch_coro(ex, wr_.cont);
288 }
289 10965x wr_.prepare(
290 21930x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_,
291 buffers, token);
292 10965x if (stop_now)
293 1x wr_.cancelled.store(true, std::memory_order_release);
294 else
295 10964x wr_.res = (n < 0) ? -err : static_cast<int>(n);
296 10965x sched_->work_started();
297 {
298 10965x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
299 10965x sched_->push_completed_locked(&wr_);
300 10965x }
301 10965x return std::noop_coroutine();
302 }
303
304 5x wr_.prepare(
305 10x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_, buffers,
306 token);
307 5x sched_->work_started();
308 5x if (wr_.cancelled.load(std::memory_order_acquire))
309 {
310 uring_scheduler::lock_type lock(sched_->dispatch_mutex());
311 sched_->push_completed_locked(&wr_);
312 return std::noop_coroutine();
313 }
314 5x uring_submit_op(*sched_, &wr_);
315 5x return std::noop_coroutine();
316 }
317
318 // ----------------------------------------------------------------
319 // tcp_socket::implementation
320 // ----------------------------------------------------------------
321
322 3267x std::coroutine_handle<> connect(
323 std::coroutine_handle<> h,
324 capy::executor_ref ex,
325 endpoint ep,
326 std::stop_token token,
327 std::error_code* ec) override
328 {
329 3267x bool stop_now = token.stop_possible() && token.stop_requested();
330 3267x if (stop_now)
331 {
332 22x if (sched_->try_consume_inline_budget())
333 {
334 21x if (ec)
335 21x *ec = capy::error::canceled;
336 21x conn_.cont.h = h;
337 21x return dispatch_coro(ex, conn_.cont);
338 }
339 1x conn_.addrlen = to_sockaddr(ep, family_, conn_.addr);
340 1x conn_.prepare(
341 2x h, ex, ec, fd_, sched_, shared_from_this(), ep,
342 &remote_endpoint_, &local_endpoint_, token);
343 1x conn_.cancelled.store(true, std::memory_order_release);
344 1x sched_->work_started();
345 {
346 1x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
347 1x sched_->push_completed_locked(&conn_);
348 1x }
349 1x return std::noop_coroutine();
350 }
351
352 // A speculative ::connect would leave the fd in EINPROGRESS and
353 // a subsequent IORING_OP_CONNECT would see EALREADY — avoid.
354 3245x conn_.addrlen = to_sockaddr(ep, family_, conn_.addr);
355 3245x conn_.prepare(
356 6490x h, ex, ec, fd_, sched_, shared_from_this(), ep, &remote_endpoint_,
357 &local_endpoint_, token);
358 3245x sched_->work_started();
359 3245x if (conn_.cancelled.load(std::memory_order_acquire))
360 {
361 uring_scheduler::lock_type lock(sched_->dispatch_mutex());
362 sched_->push_completed_locked(&conn_);
363 return std::noop_coroutine();
364 }
365 3245x uring_submit_op(*sched_, &conn_);
366 3245x return std::noop_coroutine();
367 }
368
369 20x std::coroutine_handle<> wait(
370 std::coroutine_handle<> h,
371 capy::executor_ref ex,
372 wait_type w,
373 std::stop_token token,
374 std::error_code* ec) override
375 {
376 20x int poll_flags = 0;
377 20x switch (w)
378 {
379 11x case wait_type::read:
380 11x poll_flags = POLLIN;
381 11x break;
382 7x case wait_type::write:
383 7x poll_flags = POLLOUT;
384 7x break;
385 2x case wait_type::error:
386 2x poll_flags = POLLPRI | POLLERR | POLLHUP;
387 2x break;
388 }
389 20x wait_op_.prepare(
390 40x h, ex, ec, fd_, sched_, shared_from_this(), poll_flags, token);
391 20x sched_->work_started();
392 20x if (wait_op_.cancelled.load(std::memory_order_acquire))
393 {
394 3x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
395 3x sched_->push_completed_locked(&wait_op_);
396 3x return std::noop_coroutine();
397 3x }
398 17x uring_submit_op(*sched_, &wait_op_);
399 17x return std::noop_coroutine();
400 }
401
402 4x std::error_code shutdown(tcp_socket::shutdown_type what) noexcept override
403 {
404 4x if (::shutdown(fd_, static_cast<int>(what)) != 0)
405 1x return make_err(errno);
406 3x return {};
407 }
408
409 // native_handle() / set_option() / get_option() are inherited from
410 // native_socket_base.
411
412 1x native_handle_type release_socket() noexcept override
413 {
414 // Flush while the fd is still open so the kernel resolves
415 // pending SQEs before the caller can close and recycle the
416 // number (same reasoning as close_socket).
417 1x if (fd_ >= 0)
418 1x sched_->cancel_and_flush(fd_);
419 1x int fd = fd_;
420 1x fd_ = -1;
421 1x local_endpoint_ = endpoint{};
422 1x remote_endpoint_ = endpoint{};
423 1x local_endpoint_state_.store(
424 endpoint_state::unresolved, std::memory_order_release);
425 1x return fd;
426 }
427
428 75x void cancel() noexcept override
429 {
430 75x if (fd_ >= 0)
431 73x sched_->submit_cancel_by_fd(fd_);
432 75x }
433
434 /// Cancel in-flight ops, close the fd, and reset cached endpoints.
435 /// Called by the service on close()/teardown. cancel_and_flush submits
436 /// the cancel SQE while the fd is still open so IORING_ASYNC_CANCEL_FD
437 /// resolves before the fd number can be recycled.
438 16355x void close_socket() noexcept
439 {
440 16355x if (fd_ >= 0)
441 {
442 6516x sched_->cancel_and_flush(fd_);
443 6516x ::close(fd_);
444 6516x fd_ = -1;
445 }
446 16355x local_endpoint_ = endpoint{};
447 16355x remote_endpoint_ = endpoint{};
448 16355x local_endpoint_state_.store(
449 endpoint_state::unresolved, std::memory_order_release);
450 16355x }
451
452 29x endpoint local_endpoint() const noexcept override
453 {
454 // Lazy resolution: only fire the getsockname syscall when
455 // adopt_fd marked the endpoint as "lazy_pending". For
456 // unbound/disconnected sockets the state remains unresolved
457 // and the accessor returns the default endpoint without a
458 // syscall. The mutable update races benignly with concurrent
459 // readers — both threads would compute the same value from
460 // the same fd.
461 29x if (local_endpoint_state_.load(std::memory_order_acquire) ==
462 39x endpoint_state::lazy_pending &&
463 10x fd_ >= 0)
464 {
465 10x sockaddr_storage local{};
466 10x socklen_t len = sizeof(local);
467 10x if (::getsockname(fd_, reinterpret_cast<sockaddr*>(&local), &len) ==
468 0)
469 10x local_endpoint_ = sockaddr_to_endpoint(local);
470 10x local_endpoint_state_.store(
471 endpoint_state::resolved, std::memory_order_release);
472 }
473 29x return local_endpoint_;
474 }
475
476 30x endpoint remote_endpoint() const noexcept override
477 {
478 30x return remote_endpoint_;
479 }
480 };
481
482 /** TCP socket service for io_uring.
483
484 Owns all `uring_tcp_socket` implementations for an `io_context`.
485 Satisfies the `tcp_service` interface so the generic `tcp_socket`
486 front-end can call `open_socket` and `bind_socket` transparently.
487
488 Socket impls are reference-counted inside the service map; raw
489 pointers returned from `construct()` remain valid until `destroy()`
490 or `shutdown()` is called.
491
492 @par Thread Safety
493 All public member functions are thread-safe.
494 */
495 class BOOST_COROSIO_DECL uring_tcp_service final
496 : public uring_socket_service_base<
497 uring_tcp_service,
498 tcp_service,
499 uring_tcp_socket>
500 {
501 using base_service = uring_socket_service_base<
502 uring_tcp_service,
503 tcp_service,
504 uring_tcp_socket>;
505
506 public:
507 /// Identifies this service for `execution_context` lookup.
508 using key_type = tcp_service;
509
510 /** Construct the TCP service.
511
512 @param ctx The owning execution context. The io_uring scheduler
513 must already be registered.
514 */
515 799x explicit uring_tcp_service(capy::execution_context& ctx) : base_service(ctx)
516 {
517 799x }
518
519 // construct / destroy / shutdown / close / scheduler() are inherited
520 // from uring_socket_service_base. The methods below are TCP-specific.
521
522 /** Open a socket fd and associate it with an impl.
523
524 Creates a non-blocking, close-on-exec socket via `socket(2)`.
525
526 @param impl The socket implementation to initialise.
527 @param family Address family (e.g. `AF_INET`, `AF_INET6`).
528 @param type Socket type (e.g. `SOCK_STREAM`).
529 @param protocol Protocol number (e.g. `IPPROTO_TCP`).
530 @return Error code on failure, empty on success.
531 */
532 3285x std::error_code open_socket(
533 tcp_socket::implementation& impl,
534 int family,
535 int type,
536 int protocol) override
537 {
538 3285x auto& sock = static_cast<uring_tcp_socket&>(impl);
539 int fd =
540 3285x ::socket(family, type | SOCK_NONBLOCK | SOCK_CLOEXEC, protocol);
541 3285x if (fd < 0)
542 2x return make_err(errno);
543 // LCOV_EXCL_START: dead — open() guards is_open(), so open_socket
544 // never runs against an already-open fd (assign_socket handles that).
545 if (sock.fd_ >= 0)
546 {
547 sched_->submit_cancel_by_fd(sock.fd_);
548 ::close(sock.fd_);
549 }
550 // LCOV_EXCL_STOP
551 3283x sock.fd_ = fd;
552 3283x sock.family_ = family;
553 // Mirror epoll/select: IPv6 sockets default to v6-only so they
554 // behave consistently across platforms regardless of the kernel
555 // default for /proc/sys/net/ipv6/bindv6only.
556 3283x if (family == AF_INET6)
557 {
558 9x int one = 1;
559 9x ::setsockopt(fd, IPPROTO_IPV6, IPV6_V6ONLY, &one, sizeof(one));
560 }
561 3283x return {};
562 }
563
564 /** Adopt a pre-created fd into an impl.
565
566 Takes ownership of `fd` on success; the caller retains
567 ownership on failure.
568
569 @param impl The socket implementation to assign to.
570 @param fd A valid, open, non-blocking IP stream fd.
571 @return Error code on failure, empty on success.
572 */
573 13x std::error_code assign_socket(
574 tcp_socket::implementation& impl, native_handle_type fd) override
575 {
576 13x auto& sock = static_cast<uring_tcp_socket&>(impl);
577 13x int nfd = static_cast<int>(fd);
578 13x if (nfd >= 0 && nfd == sock.fd_)
579 1x return std::make_error_code(std::errc::invalid_argument);
580 12x if (auto ec = validate_socket_fd(nfd, SOCK_STREAM, true))
581 9x return ec;
582
583 3x if (sock.fd_ >= 0)
584 {
585 1x sched_->cancel_and_flush(sock.fd_);
586 1x ::close(sock.fd_);
587 }
588 3x sock.fd_ = nfd;
589
590 3x sock.local_endpoint_ = endpoint{};
591 3x sock.remote_endpoint_ = endpoint{};
592
593 3x sockaddr_storage local{};
594 3x socklen_t local_len = sizeof(local);
595 3x if (::getsockname(
596 3x sock.fd_, reinterpret_cast<sockaddr*>(&local), &local_len) == 0)
597 {
598 3x sock.local_endpoint_ = sockaddr_to_endpoint(local);
599 3x sock.family_ = local.ss_family;
600 }
601 3x sock.local_endpoint_state_.store(
602 uring_tcp_socket::endpoint_state::resolved,
603 std::memory_order_release);
604
605 3x sockaddr_storage remote{};
606 3x socklen_t remote_len = sizeof(remote);
607 3x if (::getpeername(
608 3x sock.fd_, reinterpret_cast<sockaddr*>(&remote), &remote_len) ==
609 0)
610 3x sock.remote_endpoint_ = sockaddr_to_endpoint(remote);
611
612 3x return {};
613 }
614
615 /** Bind the socket and capture the local endpoint via `getsockname`.
616
617 @param impl The socket implementation to bind.
618 @param ep The local endpoint to bind to.
619 @return Error code on failure, empty on success.
620 */
621 std::error_code
622 9x bind_socket(tcp_socket::implementation& impl, endpoint ep) override
623 {
624 9x auto& sock = static_cast<uring_tcp_socket&>(impl);
625 9x sockaddr_storage addr{};
626 9x socklen_t len = endpoint_to_sockaddr(ep, addr);
627 9x if (::bind(sock.fd_, reinterpret_cast<sockaddr*>(&addr), len) < 0)
628 4x return make_err(errno);
629
630 5x sockaddr_storage local{};
631 5x socklen_t local_len = sizeof(local);
632 5x if (::getsockname(
633 5x sock.fd_, reinterpret_cast<sockaddr*>(&local), &local_len) == 0)
634 5x sock.local_endpoint_ = sockaddr_to_endpoint(local);
635 5x sock.local_endpoint_state_.store(
636 uring_tcp_socket::endpoint_state::resolved,
637 std::memory_order_release);
638 5x return {};
639 }
640
641 /** Wrap an already-accepted fd as a new socket impl.
642
643 Called by the acceptor service (Task 17) after `accept(2)`
644 returns a connected fd. Captures both endpoints via the provided
645 peer address and a `getsockname` call.
646
647 @param fd Accepted file descriptor (must be non-blocking).
648 @param peer Peer endpoint from `accept(2)`.
649 @return Raw pointer to the registered impl.
650 */
651 3233x uring_tcp_socket* adopt_fd(int fd, endpoint const& peer)
652 {
653 3233x auto p = std::make_shared<uring_tcp_socket>(*this, *sched_);
654 3233x p->fd_ = fd;
655 3233x p->remote_endpoint_ = peer;
656 // Mark the local endpoint as authoritative-but-unresolved.
657 // The accessor will fetch it via getsockname on first call.
658 // Accept-heavy workloads that never query the local endpoint
659 // skip the syscall entirely.
660 3233x p->local_endpoint_state_.store(
661 uring_tcp_socket::endpoint_state::lazy_pending,
662 std::memory_order_release);
663
664 6466x return this->register_impl(std::move(p));
665 3233x }
666 };
667
668 /** TCP acceptor implementation for io_uring.
669
670 Inherits the multishot machinery (parked-fd queue, waiter queue,
671 CQE drain on destruction) from `uring_multishot_acceptor_base`.
672 This class adds only the `accept()` override (matching
673 `tcp_acceptor::implementation`'s exact signature) and the
674 `adopt_thunk` static that wraps an accepted fd via
675 `uring_tcp_service::adopt_fd`.
676 */
677 class BOOST_COROSIO_DECL uring_tcp_acceptor final
678 : public uring_multishot_acceptor_base<
679 uring_tcp_acceptor,
680 tcp_acceptor::implementation,
681 endpoint,
682 uring_tcp_service>
683 {
684 friend uring_tcp_acceptor_service;
685
686 using base_type = uring_multishot_acceptor_base<
687 uring_tcp_acceptor,
688 tcp_acceptor::implementation,
689 endpoint,
690 uring_tcp_service>;
691
692 // Readiness-wait slot. The multishot accept op delivers accepted
693 // fds, but `wait()` reports raw poll readiness on the listening fd
694 // without consuming a connection — see the wait() override.
695 uring_wait_op wait_op_;
696
697 public:
698 217x explicit uring_tcp_acceptor(
699 uring_tcp_acceptor_service&,
700 uring_scheduler& sched,
701 uring_tcp_service& peer_svc) noexcept
702 217x : base_type(sched, peer_svc)
703 {
704 217x }
705
706 3266x std::coroutine_handle<> accept(
707 std::coroutine_handle<> h,
708 capy::executor_ref ex,
709 std::stop_token token,
710 std::error_code* ec,
711 io_object::implementation** impl_out) override
712 {
713 3266x base_type::dispatch_or_queue(h, ex, token, ec, impl_out);
714 3266x return std::noop_coroutine();
715 }
716
717 13x std::coroutine_handle<> wait(
718 std::coroutine_handle<> h,
719 capy::executor_ref ex,
720 wait_type w,
721 std::stop_token token,
722 std::error_code* ec) override
723 {
724 // Closed-object contract: complete with bad_file_descriptor
725 // instead of parking a waiter no accept machinery will signal.
726 13x if (this->fd_ < 0)
727 {
728 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept-adjacent initiation path: OOM => std::terminate is the intended behavior
729 1x auto* op = new uring_accept_op();
730 1x op->h = h;
731 1x op->ex = ex;
732 1x op->ec_out = ec;
733 1x op->err = EBADF;
734 1x this->sched_->post(op);
735 1x return std::noop_coroutine();
736 }
737 // Multishot accepting drains the kernel queue as connections
738 // arrive, so a poll on the listener never reports it
739 // readable; read waits complete from the delivery queue.
740 12x if (w == wait_type::read)
741 {
742 10x this->park_read_wait(h, ex, token, ec);
743 10x return std::noop_coroutine();
744 }
745 // Writability carries no meaning for a listening socket;
746 // fail uniformly instead of never completing.
747 2x if (w == wait_type::write)
748 {
749 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept-adjacent initiation path: OOM => std::terminate is the intended behavior
750 1x auto* op = new uring_accept_op();
751 1x op->h = h;
752 1x op->ex = ex;
753 1x op->ec_out = ec;
754 1x op->err = ENOTSUP;
755 1x this->sched_->post(op);
756 1x return std::noop_coroutine();
757 }
758 // Errors are not consumed by the accept machinery, so the
759 // error wait still polls the descriptor.
760 1x wait_op_.prepare(
761 2x h, ex, ec, this->fd_, this->sched_, this->shared_from_this(),
762 POLLPRI | POLLERR | POLLHUP, token);
763 1x this->sched_->work_started();
764 1x if (wait_op_.cancelled.load(std::memory_order_acquire))
765 {
766 uring_scheduler::lock_type lock(this->sched_->dispatch_mutex());
767 this->sched_->push_completed_locked(&wait_op_);
768 return std::noop_coroutine();
769 }
770 1x uring_submit_op(*this->sched_, &wait_op_);
771 1x return std::noop_coroutine();
772 }
773
774 3233x static io_object::implementation* adopt_thunk(
775 void* peer_service,
776 int fd,
777 sockaddr_storage const& peer,
778 socklen_t /*peer_len*/) noexcept
779 {
780 3233x auto* svc = static_cast<uring_tcp_service*>(peer_service);
781 3233x return svc->adopt_fd(fd, sockaddr_to_endpoint(peer));
782 }
783 };
784
785 /** TCP acceptor service for io_uring.
786
787 Owns all `uring_tcp_acceptor` implementations for an `io_context`.
788 Satisfies the `tcp_acceptor_service` interface so the generic
789 `tcp_acceptor` front-end can call `open_acceptor_socket`,
790 `bind_acceptor`, and `listen_acceptor` transparently.
791
792 Acceptor impls are reference-counted inside the service map; raw
793 pointers returned from `construct()` remain valid until `destroy()`
794 or `shutdown()` is called.
795
796 @par Thread Safety
797 All public member functions are thread-safe.
798 */
799 class BOOST_COROSIO_DECL uring_tcp_acceptor_service final
800 : public tcp_acceptor_service
801 {
802 public:
803 /// Identifies this service for `execution_context` lookup.
804 using key_type = tcp_acceptor_service;
805
806 /** Construct the TCP acceptor service.
807
808 @param ctx The owning execution context. Both the io_uring scheduler
809 and the TCP socket service must already be registered.
810 */
811 799x explicit uring_tcp_acceptor_service(capy::execution_context& ctx)
812 2397x : sched_(&ctx.use_service<uring_scheduler>())
813 799x , peer_svc_(&ctx.use_service<uring_tcp_service>())
814 {
815 799x }
816
817 799x void shutdown() override
818 {
819 799x std::vector<std::shared_ptr<uring_tcp_acceptor>> live;
820 {
821 799x std::lock_guard lk(mutex_);
822 799x live.reserve(impls_.size());
823 801x for (auto& [_, p] : impls_)
824 2x live.push_back(p);
825 799x }
826 // Cancel without the lock held to avoid inversion if cancel()
827 // re-enters the service.
828 801x for (auto& p : live)
829 2x p->cancel();
830 799x }
831
832 217x io_object::implementation* construct() override
833 {
834 auto p =
835 217x std::make_shared<uring_tcp_acceptor>(*this, *sched_, *peer_svc_);
836 217x auto* raw = p.get();
837 217x std::lock_guard lk(mutex_);
838 217x impls_.emplace(raw, std::move(p));
839 217x return raw;
840 217x }
841
842 215x void destroy(io_object::implementation* p) override
843 {
844 215x if (!p)
845 return;
846 215x std::lock_guard lk(mutex_);
847 215x impls_.erase(static_cast<uring_tcp_acceptor*>(p));
848 215x }
849
850 // Close the fd eagerly when tcp_acceptor::close() is called, before
851 // destroy() drops the shared_ptr and the destructor runs.
852 417x void close(io_object::handle& h) override
853 {
854 417x auto* acc = static_cast<uring_tcp_acceptor*>(h.get());
855 417x if (acc && acc->fd_ >= 0)
856 {
857 // Flush the cancel SQE before closing the fd so the kernel
858 // resolves the file from the fd number while it is still
859 // valid. drain_waiters_only avoids submitting cancel-by-fd
860 // a second time (cancel_and_flush already did it).
861 202x sched_->cancel_and_flush(acc->fd_);
862 202x acc->drain_waiters_only();
863 202x ::close(acc->fd_);
864 202x acc->fd_ = -1;
865
866 // Break the multi_op_ -> impl_ptr (shared_ptr<this>) cycle
867 // start_multishot established. The acceptor destructor's
868 // drain_cqes_for(multi_op_.get()) is the safety net; here
869 // we just drop the cycle so the impl can be released when
870 // the user's last shared_ptr does.
871 202x if (acc->multi_op_)
872 180x acc->multi_op_->impl_ptr.reset();
873 }
874 417x }
875
876 /** Create a non-blocking, close-on-exec socket for accepting.
877
878 @param impl The acceptor implementation to initialise.
879 @param family Address family (e.g. `AF_INET`, `AF_INET6`).
880 @param type Socket type (e.g. `SOCK_STREAM`).
881 @param protocol Protocol number (e.g. `IPPROTO_TCP`).
882 @return Error code on failure, empty on success.
883 */
884 204x std::error_code open_acceptor_socket(
885 tcp_acceptor::implementation& impl,
886 int family,
887 int type,
888 int protocol) override
889 {
890 204x auto& acc = static_cast<uring_tcp_acceptor&>(impl);
891 int fd =
892 204x ::socket(family, type | SOCK_NONBLOCK | SOCK_CLOEXEC, protocol);
893 204x if (fd < 0)
894 1x return make_err(errno);
895 // LCOV_EXCL_START: dead — open() guards is_open(), so open_socket
896 // never runs against an already-open fd (assign_socket handles that).
897 if (acc.fd_ >= 0)
898 {
899 sched_->submit_cancel_by_fd(acc.fd_);
900 ::close(acc.fd_);
901 }
902 // LCOV_EXCL_STOP
903 203x acc.fd_ = fd;
904 // Match epoll/select: IPv6 acceptors default to dual-stack
905 // (v6-only=false) so they accept both IPv4 and IPv6 connections.
906 203x if (family == AF_INET6)
907 {
908 11x int zero = 0;
909 11x ::setsockopt(fd, IPPROTO_IPV6, IPV6_V6ONLY, &zero, sizeof(zero));
910 }
911 203x return {};
912 }
913
914 /** Adopt an already-listening descriptor.
915
916 @param impl The acceptor implementation to assign to.
917 @param fd The native socket to adopt.
918 @return Error code on failure, empty on success.
919 */
920 10x std::error_code assign_socket(
921 tcp_acceptor::implementation& impl, native_handle_type fd) override
922 {
923 10x auto& acc = static_cast<uring_tcp_acceptor&>(impl);
924 10x int nfd = static_cast<int>(fd);
925 10x if (nfd >= 0 && nfd == acc.fd_)
926 1x return std::make_error_code(std::errc::invalid_argument);
927 9x if (auto ec = validate_socket_fd(nfd, SOCK_STREAM, true))
928 3x return ec;
929
930 6x if (acc.fd_ >= 0)
931 {
932 1x sched_->cancel_and_flush(acc.fd_);
933 1x acc.drain_waiters_only();
934 1x ::close(acc.fd_);
935 1x acc.fd_ = -1;
936 }
937
938 // Unconditional: release_socket() also leaves the op in flight,
939 // and it clears fd_ before returning.
940 6x acc.retire_multishot();
941
942 6x acc.adopt_listening_fd(nfd);
943
944 6x acc.local_endpoint_ = endpoint{};
945 6x sockaddr_storage local{};
946 6x socklen_t local_len = sizeof(local);
947 6x if (::getsockname(
948 6x nfd, reinterpret_cast<sockaddr*>(&local), &local_len) == 0)
949 6x acc.local_endpoint_ = sockaddr_to_endpoint(local);
950
951 6x if (fd_is_listening(nfd))
952 6x acc.start_multishot();
953 6x return {};
954 }
955
956 /** Bind an open acceptor and capture the local endpoint.
957
958 @param impl The acceptor implementation to bind.
959 @param ep The local endpoint to bind to.
960 @return Error code on failure, empty on success.
961 */
962 std::error_code
963 197x bind_acceptor(tcp_acceptor::implementation& impl, endpoint ep) override
964 {
965 197x auto& acc = static_cast<uring_tcp_acceptor&>(impl);
966 197x sockaddr_storage addr{};
967 197x socklen_t len = endpoint_to_sockaddr(ep, addr);
968 197x if (::bind(acc.fd_, reinterpret_cast<sockaddr*>(&addr), len) < 0)
969 5x return make_err(errno);
970
971 192x sockaddr_storage local{};
972 192x socklen_t local_len = sizeof(local);
973 192x if (::getsockname(
974 192x acc.fd_, reinterpret_cast<sockaddr*>(&local), &local_len) == 0)
975 192x acc.local_endpoint_ = sockaddr_to_endpoint(local);
976 192x return {};
977 }
978
979 /** Start listening and submit the multishot accept SQE.
980
981 Calls `::listen(2)` then arms the io_uring multishot accept
982 operation that delivers one CQE per accepted connection.
983
984 @param impl The acceptor implementation to listen on.
985 @param backlog Maximum pending-connection queue length.
986 @return Error code on failure, empty on success.
987 */
988 std::error_code
989 186x listen_acceptor(tcp_acceptor::implementation& impl, int backlog) override
990 {
991 186x auto& acc = static_cast<uring_tcp_acceptor&>(impl);
992 186x if (::listen(acc.fd_, backlog) < 0)
993 3x return make_err(errno);
994 183x if (acc.prepare_listen_arm())
995 182x acc.start_multishot();
996 183x return {};
997 }
998
999 /// Return the scheduler used by acceptors created by this service.
1000 uring_scheduler& scheduler() noexcept
1001 {
1002 return *sched_;
1003 }
1004
1005 private:
1006 uring_scheduler* sched_;
1007 uring_tcp_service* peer_svc_;
1008 std::mutex mutex_;
1009 std::unordered_map<uring_tcp_acceptor*, std::shared_ptr<uring_tcp_acceptor>>
1010 impls_;
1011 };
1012
1013 /** Unix domain stream socket implementation for io_uring.
1014
1015 Implements `local_stream_socket::implementation` using a proactor
1016 model: read, write, and connect operations are submitted to the
1017 kernel via `uring_submit_op` and complete through the ring's
1018 CQE path.
1019
1020 The object is always owned by a `shared_ptr` managed by the service.
1021 In-flight ops hold an additional `shared_ptr` copy (`impl_ptr`) so
1022 the kernel's user-data pointer remains valid until the CQE arrives.
1023
1024 @par Thread Safety
1025 Distinct objects: Safe.
1026 Shared objects: Unsafe. A socket must not have two operations of
1027 the same type in flight simultaneously.
1028 */
1029 class BOOST_COROSIO_DECL uring_local_stream_socket final
1030 : public native_socket_base<
1031 uring_local_stream_socket,
1032 local_stream_socket::implementation,
1033 corosio::local_endpoint>
1034 {
1035 friend uring_local_stream_service;
1036
1037 uring_scheduler* sched_ = nullptr;
1038 [[maybe_unused]] uring_local_stream_service* svc_ = nullptr;
1039
1040 // fd_ and local_endpoint_ live in native_socket_base, which also
1041 // provides native_handle/is_open/set_option/get_option/local_endpoint.
1042 corosio::local_endpoint remote_endpoint_;
1043
1044 // Per-fd op slots — embedded to eliminate per-call heap allocation.
1045 // Single-pending invariant per slot.
1046 uring_read_op rd_;
1047 uring_write_op wr_;
1048 uring_local_connect_op conn_;
1049 uring_wait_op wait_op_;
1050
1051 mutable detail::speculative_state spec_;
1052
1053 public:
1054 /** Construct with service and scheduler references.
1055
1056 Both refs must outlive this socket.
1057
1058 @param svc The owning service.
1059 @param sched The io_uring scheduler owned by the context.
1060 */
1061 159x explicit uring_local_stream_socket(
1062 uring_local_stream_service& svc, uring_scheduler& sched) noexcept
1063 318x : sched_(&sched)
1064 159x , svc_(&svc)
1065 {
1066 159x }
1067
1068 157x ~uring_local_stream_socket() override
1069 157x {
1070 157x if (fd_ >= 0)
1071 ::close(
1072 fd_); // LCOV_EXCL_LINE backstop: close_socket() clears fd_ before destroy
1073 157x }
1074
1075 // ----------------------------------------------------------------
1076 // io_stream::implementation
1077 // ----------------------------------------------------------------
1078
1079 64x std::coroutine_handle<> read_some(
1080 std::coroutine_handle<> h,
1081 capy::executor_ref ex,
1082 buffer_param buffers,
1083 std::stop_token token,
1084 std::error_code* ec,
1085 std::size_t* bytes) override
1086 {
1087 iovec iovecs[uring_max_iov];
1088 64x int iovec_count = copy_to_iovec(buffers, iovecs);
1089 64x bool stop_now = token.stop_possible() && token.stop_requested();
1090 64x bool empty_buf = (iovec_count == 0);
1091
1092 64x ssize_t n = 0;
1093 64x int err = 0;
1094 64x bool have_sync_res = stop_now || empty_buf;
1095 64x if (!have_sync_res && spec_.may_speculate_read())
1096 {
1097 do
1098 {
1099 39x n = ::readv(fd_, iovecs, iovec_count);
1100 }
1101 39x while (n < 0 && errno == EINTR);
1102 39x if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
1103 {
1104 24x have_sync_res = true;
1105 24x if (n < 0)
1106 2x err = errno;
1107 // Speculative read produced a definitive answer (data
1108 // or non-EAGAIN error); reset the failure streak so a
1109 // burst of past EAGAINs doesn't latch perma-off when
1110 // the workload is in fact speculation-friendly.
1111 24x if (n >= 0)
1112 22x spec_.on_read_success();
1113 }
1114 else
1115 {
1116 15x spec_.on_read_exhausted();
1117 }
1118 }
1119
1120 64x if (have_sync_res)
1121 {
1122 47x if (sched_->try_consume_inline_budget())
1123 {
1124 45x decode_io_result(
1125 2x ec, stop_now, err ? make_err(err) : std::error_code{},
1126 45x /*is_read=*/true, n < 0 ? 0u : static_cast<std::size_t>(n),
1127 empty_buf);
1128 45x if (bytes)
1129 45x *bytes = (n < 0) ? 0u : static_cast<std::size_t>(n);
1130 45x rd_.cont.h = h;
1131 45x return dispatch_coro(ex, rd_.cont);
1132 }
1133 2x rd_.prepare(
1134 4x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_,
1135 buffers, token);
1136 2x if (stop_now)
1137 1x rd_.cancelled.store(true, std::memory_order_release);
1138 else
1139 1x rd_.res = (n < 0) ? -err : static_cast<int>(n);
1140 2x sched_->work_started();
1141 {
1142 2x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
1143 2x sched_->push_completed_locked(&rd_);
1144 2x }
1145 2x return std::noop_coroutine();
1146 }
1147
1148 17x rd_.prepare(
1149 34x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_, buffers,
1150 token);
1151 17x sched_->work_started();
1152 17x if (rd_.cancelled.load(std::memory_order_acquire))
1153 {
1154 uring_scheduler::lock_type lock(sched_->dispatch_mutex());
1155 sched_->push_completed_locked(&rd_);
1156 return std::noop_coroutine();
1157 }
1158 17x uring_submit_op(*sched_, &rd_);
1159 17x return std::noop_coroutine();
1160 }
1161
1162 60x std::coroutine_handle<> write_some(
1163 std::coroutine_handle<> h,
1164 capy::executor_ref ex,
1165 buffer_param buffers,
1166 std::stop_token token,
1167 std::error_code* ec,
1168 std::size_t* bytes) override
1169 {
1170 iovec iovecs[uring_max_iov];
1171 60x int iovec_count = copy_to_iovec(buffers, iovecs);
1172 60x bool stop_now = token.stop_possible() && token.stop_requested();
1173 60x bool empty_buf = (iovec_count == 0);
1174
1175 60x ssize_t n = 0;
1176 60x int err = 0;
1177 60x bool have_sync_res = stop_now || empty_buf;
1178 60x if (!have_sync_res && spec_.may_speculate_write())
1179 {
1180 37x msghdr msg{};
1181 37x msg.msg_iov = iovecs;
1182 37x msg.msg_iovlen = static_cast<decltype(msg.msg_iovlen)>(iovec_count);
1183 do
1184 {
1185 37x n = ::sendmsg(fd_, &msg, MSG_NOSIGNAL);
1186 }
1187 37x while (n < 0 && errno == EINTR);
1188 37x if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
1189 {
1190 36x have_sync_res = true;
1191 36x if (n < 0)
1192 2x err = errno;
1193 }
1194 else
1195 {
1196 1x spec_.on_write_exhausted();
1197 }
1198 }
1199
1200 60x if (have_sync_res)
1201 {
1202 59x if (sched_->try_consume_inline_budget())
1203 {
1204 59x decode_io_result(
1205 2x ec, stop_now, err ? make_err(err) : std::error_code{},
1206 /*is_read=*/false, /*bytes=*/0, /*empty_buffer=*/false);
1207 57x if (bytes)
1208 57x *bytes = (n < 0) ? 0u : static_cast<std::size_t>(n);
1209 57x wr_.cont.h = h;
1210 57x return dispatch_coro(ex, wr_.cont);
1211 }
1212 2x wr_.prepare(
1213 4x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_,
1214 buffers, token);
1215 2x if (stop_now)
1216 1x wr_.cancelled.store(true, std::memory_order_release);
1217 else
1218 1x wr_.res = (n < 0) ? -err : static_cast<int>(n);
1219 2x sched_->work_started();
1220 {
1221 2x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
1222 2x sched_->push_completed_locked(&wr_);
1223 2x }
1224 2x return std::noop_coroutine();
1225 }
1226
1227 1x wr_.prepare(
1228 2x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_, buffers,
1229 token);
1230 1x sched_->work_started();
1231 1x if (wr_.cancelled.load(std::memory_order_acquire))
1232 {
1233 uring_scheduler::lock_type lock(sched_->dispatch_mutex());
1234 sched_->push_completed_locked(&wr_);
1235 return std::noop_coroutine();
1236 }
1237 1x uring_submit_op(*sched_, &wr_);
1238 1x return std::noop_coroutine();
1239 }
1240
1241 // ----------------------------------------------------------------
1242 // local_stream_socket::implementation
1243 // ----------------------------------------------------------------
1244
1245 38x std::coroutine_handle<> connect(
1246 std::coroutine_handle<> h,
1247 capy::executor_ref ex,
1248 corosio::local_endpoint ep,
1249 std::stop_token token,
1250 std::error_code* ec) override
1251 {
1252 38x bool stop_now = token.stop_possible() && token.stop_requested();
1253 38x if (stop_now)
1254 {
1255 22x if (sched_->try_consume_inline_budget())
1256 {
1257 21x if (ec)
1258 21x *ec = capy::error::canceled;
1259 21x conn_.cont.h = h;
1260 21x return dispatch_coro(ex, conn_.cont);
1261 }
1262 1x conn_.addrlen = to_sockaddr(ep, conn_.addr);
1263 1x conn_.prepare(
1264 2x h, ex, ec, fd_, sched_, shared_from_this(), ep,
1265 &remote_endpoint_, &local_endpoint_, token);
1266 1x conn_.cancelled.store(true, std::memory_order_release);
1267 1x sched_->work_started();
1268 {
1269 1x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
1270 1x sched_->push_completed_locked(&conn_);
1271 1x }
1272 1x return std::noop_coroutine();
1273 }
1274
1275 // A speculative ::connect would leave the fd in EINPROGRESS and
1276 // a subsequent IORING_OP_CONNECT would see EALREADY — avoid.
1277 16x conn_.addrlen = to_sockaddr(ep, conn_.addr);
1278 16x conn_.prepare(
1279 32x h, ex, ec, fd_, sched_, shared_from_this(), ep, &remote_endpoint_,
1280 &local_endpoint_, token);
1281 16x sched_->work_started();
1282 16x if (conn_.cancelled.load(std::memory_order_acquire))
1283 {
1284 uring_scheduler::lock_type lock(sched_->dispatch_mutex());
1285 sched_->push_completed_locked(&conn_);
1286 return std::noop_coroutine();
1287 }
1288 16x uring_submit_op(*sched_, &conn_);
1289 16x return std::noop_coroutine();
1290 }
1291
1292 10x std::coroutine_handle<> wait(
1293 std::coroutine_handle<> h,
1294 capy::executor_ref ex,
1295 wait_type w,
1296 std::stop_token token,
1297 std::error_code* ec) override
1298 {
1299 10x int poll_flags = 0;
1300 10x switch (w)
1301 {
1302 8x case wait_type::read:
1303 8x poll_flags = POLLIN;
1304 8x break;
1305 1x case wait_type::write:
1306 1x poll_flags = POLLOUT;
1307 1x break;
1308 1x case wait_type::error:
1309 1x poll_flags = POLLPRI | POLLERR | POLLHUP;
1310 1x break;
1311 }
1312 10x wait_op_.prepare(
1313 20x h, ex, ec, fd_, sched_, shared_from_this(), poll_flags, token);
1314 10x sched_->work_started();
1315 10x if (wait_op_.cancelled.load(std::memory_order_acquire))
1316 {
1317 2x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
1318 2x sched_->push_completed_locked(&wait_op_);
1319 2x return std::noop_coroutine();
1320 2x }
1321 8x uring_submit_op(*sched_, &wait_op_);
1322 8x return std::noop_coroutine();
1323 }
1324
1325 std::error_code
1326 3x shutdown(local_stream_socket::shutdown_type what) noexcept override
1327 {
1328 3x if (::shutdown(fd_, static_cast<int>(what)) != 0)
1329 1x return make_err(errno);
1330 2x return {};
1331 }
1332
1333 // native_handle / is_open / set_option / get_option / local_endpoint
1334 // are inherited from native_socket_base.
1335
1336 2x native_handle_type release_socket() noexcept override
1337 {
1338 // Flush while the fd is still open so the kernel resolves
1339 // pending SQEs before the caller can close and recycle the
1340 // number (same reasoning as close_socket).
1341 2x if (fd_ >= 0)
1342 2x sched_->cancel_and_flush(fd_);
1343 2x int fd = fd_;
1344 2x fd_ = -1;
1345 2x local_endpoint_ = corosio::local_endpoint{};
1346 2x remote_endpoint_ = corosio::local_endpoint{};
1347 2x return fd;
1348 }
1349
1350 3x void cancel() noexcept override
1351 {
1352 3x if (fd_ >= 0)
1353 2x sched_->submit_cancel_by_fd(fd_);
1354 3x }
1355
1356 /// Cancel in-flight ops, close the fd, and reset cached endpoints.
1357 /// Called by the service on close()/teardown.
1358 258x void close_socket() noexcept
1359 {
1360 258x if (fd_ >= 0)
1361 {
1362 100x sched_->cancel_and_flush(fd_);
1363 100x ::close(fd_);
1364 100x fd_ = -1;
1365 }
1366 258x local_endpoint_ = corosio::local_endpoint{};
1367 258x remote_endpoint_ = corosio::local_endpoint{};
1368 258x }
1369
1370 2x corosio::local_endpoint remote_endpoint() const noexcept override
1371 {
1372 2x return remote_endpoint_;
1373 }
1374 };
1375
1376 /** Unix domain stream socket service for io_uring.
1377
1378 Owns all `uring_local_stream_socket` implementations for an
1379 `io_context`. Satisfies the `local_stream_service` interface so the
1380 generic `local_stream_socket` front-end can call `open_socket` and
1381 `assign_socket` transparently.
1382
1383 Socket impls are reference-counted inside the service map; raw
1384 pointers returned from `construct()` remain valid until `destroy()`
1385 or `shutdown()` is called.
1386
1387 @par Thread Safety
1388 All public member functions are thread-safe.
1389 */
1390 class BOOST_COROSIO_DECL uring_local_stream_service final
1391 : public uring_socket_service_base<
1392 uring_local_stream_service,
1393 local_stream_service,
1394 uring_local_stream_socket>
1395 {
1396 using base_service = uring_socket_service_base<
1397 uring_local_stream_service,
1398 local_stream_service,
1399 uring_local_stream_socket>;
1400
1401 public:
1402 /// Identifies this service for `execution_context` lookup.
1403 using key_type = local_stream_service;
1404
1405 /** Construct the local stream service.
1406
1407 @param ctx The owning execution context. The io_uring scheduler
1408 must already be registered.
1409 */
1410 799x explicit uring_local_stream_service(capy::execution_context& ctx)
1411 799x : base_service(ctx)
1412 {
1413 799x }
1414
1415 // construct / destroy / shutdown / close / scheduler() are inherited
1416 // from uring_socket_service_base.
1417
1418 /** Open an AF_UNIX stream socket and associate it with an impl.
1419
1420 Creates a non-blocking, close-on-exec socket via `socket(2)`.
1421 `family` is always `AF_UNIX` for local stream sockets.
1422
1423 @param impl The socket implementation to initialise.
1424 @param family Address family (`AF_UNIX`).
1425 @param type Socket type (`SOCK_STREAM`).
1426 @param protocol Protocol number (typically 0).
1427 @return Error code on failure, empty on success.
1428 */
1429 26x std::error_code open_socket(
1430 local_stream_socket::implementation& impl,
1431 int family,
1432 int type,
1433 int protocol) override
1434 {
1435 26x auto& sock = static_cast<uring_local_stream_socket&>(impl);
1436 int fd =
1437 26x ::socket(family, type | SOCK_NONBLOCK | SOCK_CLOEXEC, protocol);
1438 26x if (fd < 0)
1439 1x return make_err(errno);
1440 // LCOV_EXCL_START: dead — open() guards is_open(), so open_socket
1441 // never runs against an already-open fd (assign_socket handles that).
1442 if (sock.fd_ >= 0)
1443 {
1444 sched_->submit_cancel_by_fd(sock.fd_);
1445 ::close(sock.fd_);
1446 }
1447 // LCOV_EXCL_STOP
1448 25x sock.fd_ = fd;
1449 25x return {};
1450 }
1451
1452 /** Adopt a pre-created fd into an impl (e.g. from `socketpair`).
1453
1454 Takes ownership of `fd` on success; the caller retains ownership
1455 on failure.
1456
1457 @param impl The socket implementation to assign to.
1458 @param fd A valid, open, non-blocking AF_UNIX stream fd.
1459 @return Error code on failure, empty on success.
1460 */
1461 74x std::error_code assign_socket(
1462 local_stream_socket::implementation& impl,
1463 native_handle_type fd) override
1464 {
1465 74x auto& sock = static_cast<uring_local_stream_socket&>(impl);
1466 74x int nfd = static_cast<int>(fd);
1467 74x if (nfd >= 0 && nfd == sock.fd_)
1468 1x return std::make_error_code(std::errc::invalid_argument);
1469 73x if (auto ec = validate_socket_fd(nfd, SOCK_STREAM, false))
1470 7x return ec;
1471
1472 66x if (sock.fd_ >= 0)
1473 {
1474 1x sched_->cancel_and_flush(sock.fd_);
1475 1x ::close(sock.fd_);
1476 }
1477 66x sock.fd_ = nfd;
1478
1479 66x sockaddr_storage local{};
1480 66x socklen_t local_len = sizeof(local);
1481 66x if (::getsockname(
1482 66x sock.fd_, reinterpret_cast<sockaddr*>(&local), &local_len) == 0)
1483 66x sock.local_endpoint_ = sockaddr_to_local_endpoint(local, local_len);
1484
1485 66x sockaddr_storage remote{};
1486 66x socklen_t remote_len = sizeof(remote);
1487 66x if (::getpeername(
1488 66x sock.fd_, reinterpret_cast<sockaddr*>(&remote), &remote_len) ==
1489 0)
1490 sock.remote_endpoint_ =
1491 66x sockaddr_to_local_endpoint(remote, remote_len);
1492
1493 66x return {};
1494 }
1495
1496 /** Wrap an already-accepted fd as a new socket impl.
1497
1498 Called by the acceptor service after `accept(2)` returns a
1499 connected fd. Captures both endpoints via the provided peer
1500 address and a `getsockname` call.
1501
1502 @param fd Accepted file descriptor (must be non-blocking).
1503 @param peer Peer endpoint from `accept(2)`.
1504 @return Raw pointer to the registered impl.
1505 */
1506 uring_local_stream_socket*
1507 12x adopt_fd(int fd, corosio::local_endpoint const& peer)
1508 {
1509 12x auto p = std::make_shared<uring_local_stream_socket>(*this, *sched_);
1510 12x p->fd_ = fd;
1511 12x p->remote_endpoint_ = peer;
1512
1513 12x sockaddr_storage local{};
1514 12x socklen_t len = sizeof(local);
1515 12x if (::getsockname(fd, reinterpret_cast<sockaddr*>(&local), &len) == 0)
1516 12x p->local_endpoint_ = sockaddr_to_local_endpoint(local, len);
1517
1518 24x return this->register_impl(std::move(p));
1519 12x }
1520 };
1521
1522 /** Local-stream (Unix domain) acceptor for io_uring.
1523
1524 Inherits all multishot machinery (parked-fd queue, waiter queue,
1525 descriptor release, CQE drain on destruction) from
1526 `uring_multishot_acceptor_base`. Adds only the `accept()`
1527 override and the `adopt_thunk` static that wraps an accepted fd
1528 via `uring_local_stream_service::adopt_fd`.
1529 */
1530 class BOOST_COROSIO_DECL uring_local_stream_acceptor final
1531 : public uring_multishot_acceptor_base<
1532 uring_local_stream_acceptor,
1533 local_stream_acceptor::implementation,
1534 corosio::local_endpoint,
1535 uring_local_stream_service>
1536 {
1537 friend uring_local_stream_acceptor_service;
1538
1539 using base_type = uring_multishot_acceptor_base<
1540 uring_local_stream_acceptor,
1541 local_stream_acceptor::implementation,
1542 corosio::local_endpoint,
1543 uring_local_stream_service>;
1544
1545 // Readiness-wait slot. See uring_tcp_acceptor::wait_op_.
1546 uring_wait_op wait_op_;
1547
1548 public:
1549 55x explicit uring_local_stream_acceptor(
1550 uring_local_stream_acceptor_service&,
1551 uring_scheduler& sched,
1552 uring_local_stream_service& peer_svc) noexcept
1553 55x : base_type(sched, peer_svc)
1554 {
1555 55x }
1556
1557 20x std::coroutine_handle<> accept(
1558 std::coroutine_handle<> h,
1559 capy::executor_ref ex,
1560 std::stop_token token,
1561 std::error_code* ec,
1562 io_object::implementation** impl_out) override
1563 {
1564 20x base_type::dispatch_or_queue(h, ex, token, ec, impl_out);
1565 20x return std::noop_coroutine();
1566 }
1567
1568 6x std::coroutine_handle<> wait(
1569 std::coroutine_handle<> h,
1570 capy::executor_ref ex,
1571 wait_type w,
1572 std::stop_token token,
1573 std::error_code* ec) override
1574 {
1575 // Closed-object contract: complete with bad_file_descriptor
1576 // instead of parking a waiter no accept machinery will signal.
1577 6x if (this->fd_ < 0)
1578 {
1579 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept-adjacent initiation path: OOM => std::terminate is the intended behavior
1580 1x auto* op = new uring_accept_op();
1581 1x op->h = h;
1582 1x op->ex = ex;
1583 1x op->ec_out = ec;
1584 1x op->err = EBADF;
1585 1x this->sched_->post(op);
1586 1x return std::noop_coroutine();
1587 }
1588 // Multishot accepting drains the kernel queue as connections
1589 // arrive, so a poll on the listener never reports it
1590 // readable; read waits complete from the delivery queue.
1591 5x if (w == wait_type::read)
1592 {
1593 3x this->park_read_wait(h, ex, token, ec);
1594 3x return std::noop_coroutine();
1595 }
1596 // Writability carries no meaning for a listening socket;
1597 // fail uniformly instead of never completing.
1598 2x if (w == wait_type::write)
1599 {
1600 // NOLINTNEXTLINE(bugprone-unhandled-exception-at-new) — noexcept-adjacent initiation path: OOM => std::terminate is the intended behavior
1601 1x auto* op = new uring_accept_op();
1602 1x op->h = h;
1603 1x op->ex = ex;
1604 1x op->ec_out = ec;
1605 1x op->err = ENOTSUP;
1606 1x this->sched_->post(op);
1607 1x return std::noop_coroutine();
1608 }
1609 // Errors are not consumed by the accept machinery, so the
1610 // error wait still polls the descriptor.
1611 1x wait_op_.prepare(
1612 2x h, ex, ec, this->fd_, this->sched_, this->shared_from_this(),
1613 POLLPRI | POLLERR | POLLHUP, token);
1614 1x this->sched_->work_started();
1615 1x if (wait_op_.cancelled.load(std::memory_order_acquire))
1616 {
1617 uring_scheduler::lock_type lock(this->sched_->dispatch_mutex());
1618 this->sched_->push_completed_locked(&wait_op_);
1619 return std::noop_coroutine();
1620 }
1621 1x uring_submit_op(*this->sched_, &wait_op_);
1622 1x return std::noop_coroutine();
1623 }
1624
1625 12x static io_object::implementation* adopt_thunk(
1626 void* peer_service,
1627 int fd,
1628 sockaddr_storage const& peer,
1629 socklen_t peer_len) noexcept
1630 {
1631 12x auto* svc = static_cast<uring_local_stream_service*>(peer_service);
1632 12x return svc->adopt_fd(fd, sockaddr_to_local_endpoint(peer, peer_len));
1633 }
1634 };
1635
1636 /** Unix domain stream acceptor service for io_uring.
1637
1638 Owns all `uring_local_stream_acceptor` implementations for an
1639 `io_context`. Satisfies the `local_stream_acceptor_service` interface
1640 so the generic `local_stream_acceptor` front-end can call
1641 `open_acceptor_socket`, `bind_acceptor`, and `listen_acceptor`
1642 transparently.
1643
1644 Acceptor impls are reference-counted inside the service map; raw
1645 pointers returned from `construct()` remain valid until `destroy()`
1646 or `shutdown()` is called.
1647
1648 @par Thread Safety
1649 All public member functions are thread-safe.
1650 */
1651 class BOOST_COROSIO_DECL uring_local_stream_acceptor_service final
1652 : public local_stream_acceptor_service
1653 {
1654 public:
1655 /// Identifies this service for `execution_context` lookup.
1656 using key_type = local_stream_acceptor_service;
1657
1658 /** Construct the local stream acceptor service.
1659
1660 @param ctx The owning execution context. Both the io_uring scheduler
1661 and the local stream socket service must already be registered.
1662 */
1663 799x explicit uring_local_stream_acceptor_service(capy::execution_context& ctx)
1664 2397x : sched_(&ctx.use_service<uring_scheduler>())
1665 799x , peer_svc_(&ctx.use_service<uring_local_stream_service>())
1666 {
1667 799x }
1668
1669 799x void shutdown() override
1670 {
1671 799x std::vector<std::shared_ptr<uring_local_stream_acceptor>> live;
1672 {
1673 799x std::lock_guard lk(mutex_);
1674 799x live.reserve(impls_.size());
1675 800x for (auto& [_, p] : impls_)
1676 1x live.push_back(p);
1677 799x }
1678 // Cancel without the lock held to avoid inversion if cancel()
1679 // re-enters the service.
1680 800x for (auto& p : live)
1681 1x p->cancel();
1682 799x }
1683
1684 55x io_object::implementation* construct() override
1685 {
1686 auto p = std::make_shared<uring_local_stream_acceptor>(
1687 55x *this, *sched_, *peer_svc_);
1688 55x auto* raw = p.get();
1689 55x std::lock_guard lk(mutex_);
1690 55x impls_.emplace(raw, std::move(p));
1691 55x return raw;
1692 55x }
1693
1694 54x void destroy(io_object::implementation* p) override
1695 {
1696 54x if (!p)
1697 return;
1698 54x std::lock_guard lk(mutex_);
1699 54x impls_.erase(static_cast<uring_local_stream_acceptor*>(p));
1700 54x }
1701
1702 // Close the fd eagerly when local_stream_acceptor::close() is called,
1703 // before destroy() drops the shared_ptr and the destructor runs.
1704 92x void close(io_object::handle& h) override
1705 {
1706 92x auto* acc = static_cast<uring_local_stream_acceptor*>(h.get());
1707 92x if (acc && acc->fd_ >= 0)
1708 {
1709 // cancel_and_flush submits cancel-by-fd; drain_waiters_only
1710 // drains queued waiters without re-submitting it.
1711 38x sched_->cancel_and_flush(acc->fd_);
1712 38x acc->drain_waiters_only();
1713 38x ::close(acc->fd_);
1714 38x acc->fd_ = -1;
1715
1716 // Break the multi_op_ -> impl_ptr (shared_ptr<this>) cycle
1717 // start_multishot established. See the symmetric comment
1718 // in uring_tcp_acceptor_service::close.
1719 38x if (acc->multi_op_)
1720 24x acc->multi_op_->impl_ptr.reset();
1721 }
1722 92x }
1723
1724 /** Create a non-blocking, close-on-exec AF_UNIX socket for accepting.
1725
1726 @param impl The acceptor implementation to initialise.
1727 @param family Address family (`AF_UNIX`).
1728 @param type Socket type (`SOCK_STREAM`).
1729 @param protocol Protocol number (typically 0).
1730 @return Error code on failure, empty on success.
1731 */
1732 45x std::error_code open_acceptor_socket(
1733 local_stream_acceptor::implementation& impl,
1734 int family,
1735 int type,
1736 int protocol) override
1737 {
1738 45x auto& acc = static_cast<uring_local_stream_acceptor&>(impl);
1739 int fd =
1740 45x ::socket(family, type | SOCK_NONBLOCK | SOCK_CLOEXEC, protocol);
1741 45x if (fd < 0)
1742 3x return make_err(errno);
1743 // LCOV_EXCL_START: dead — open() guards is_open(), so open_socket
1744 // never runs against an already-open fd (assign_socket handles that).
1745 if (acc.fd_ >= 0)
1746 {
1747 sched_->submit_cancel_by_fd(acc.fd_);
1748 ::close(acc.fd_);
1749 }
1750 // LCOV_EXCL_STOP
1751 42x acc.fd_ = fd;
1752 42x return {};
1753 }
1754
1755 /** Adopt an already-listening descriptor.
1756
1757 @param impl The acceptor implementation to assign to.
1758 @param fd The native socket to adopt.
1759 @return Error code on failure, empty on success.
1760 */
1761 5x std::error_code assign_socket(
1762 local_stream_acceptor::implementation& impl,
1763 native_handle_type fd) override
1764 {
1765 5x auto& acc = static_cast<uring_local_stream_acceptor&>(impl);
1766 5x int nfd = static_cast<int>(fd);
1767 5x if (nfd >= 0 && nfd == acc.fd_)
1768 1x return std::make_error_code(std::errc::invalid_argument);
1769 4x if (auto ec = validate_socket_fd(nfd, SOCK_STREAM, false))
1770 1x return ec;
1771
1772 3x if (acc.fd_ >= 0)
1773 {
1774 1x sched_->cancel_and_flush(acc.fd_);
1775 1x acc.drain_waiters_only();
1776 1x ::close(acc.fd_);
1777 1x acc.fd_ = -1;
1778 }
1779
1780 // Unconditional: release_socket() also leaves the op in flight,
1781 // and it clears fd_ before returning.
1782 3x acc.retire_multishot();
1783
1784 3x acc.adopt_listening_fd(nfd);
1785
1786 3x acc.local_endpoint_ = corosio::local_endpoint{};
1787 3x sockaddr_storage local{};
1788 3x socklen_t local_len = sizeof(local);
1789 3x if (::getsockname(
1790 3x nfd, reinterpret_cast<sockaddr*>(&local), &local_len) == 0)
1791 3x acc.local_endpoint_ = sockaddr_to_local_endpoint(local, local_len);
1792
1793 3x if (fd_is_listening(nfd))
1794 3x acc.start_multishot();
1795 3x return {};
1796 }
1797
1798 /** Bind an open acceptor and capture the local endpoint.
1799
1800 @param impl The acceptor implementation to bind.
1801 @param ep The local endpoint (path) to bind to.
1802 @return Error code on failure, empty on success.
1803 */
1804 38x std::error_code bind_acceptor(
1805 local_stream_acceptor::implementation& impl,
1806 corosio::local_endpoint ep) override
1807 {
1808 38x auto& acc = static_cast<uring_local_stream_acceptor&>(impl);
1809 38x sockaddr_storage addr{};
1810 38x socklen_t len = endpoint_to_sockaddr(ep, addr);
1811 38x if (::bind(acc.fd_, reinterpret_cast<sockaddr*>(&addr), len) < 0)
1812 2x return make_err(errno);
1813
1814 36x sockaddr_storage local{};
1815 36x socklen_t local_len = sizeof(local);
1816 36x if (::getsockname(
1817 36x acc.fd_, reinterpret_cast<sockaddr*>(&local), &local_len) == 0)
1818 36x acc.local_endpoint_ = sockaddr_to_local_endpoint(local, local_len);
1819 36x return {};
1820 }
1821
1822 /** Start listening and submit the multishot accept SQE.
1823
1824 Calls `::listen(2)` then arms the io_uring multishot accept
1825 operation that delivers one CQE per accepted connection.
1826
1827 @param impl The acceptor implementation to listen on.
1828 @param backlog Maximum pending-connection queue length.
1829 @return Error code on failure, empty on success.
1830 */
1831 31x std::error_code listen_acceptor(
1832 local_stream_acceptor::implementation& impl, int backlog) override
1833 {
1834 31x auto& acc = static_cast<uring_local_stream_acceptor&>(impl);
1835 31x if (::listen(acc.fd_, backlog) < 0)
1836 3x return make_err(errno);
1837 28x if (acc.prepare_listen_arm())
1838 28x acc.start_multishot();
1839 28x return {};
1840 }
1841
1842 /// Return the scheduler used by acceptors created by this service.
1843 uring_scheduler& scheduler() noexcept
1844 {
1845 return *sched_;
1846 }
1847
1848 private:
1849 uring_scheduler* sched_;
1850 uring_local_stream_service* peer_svc_;
1851 std::mutex mutex_;
1852 std::unordered_map<
1853 uring_local_stream_acceptor*,
1854 std::shared_ptr<uring_local_stream_acceptor>>
1855 impls_;
1856 };
1857
1858 /** UDP socket implementation for io_uring.
1859
1860 Implements `udp_socket::implementation` using a proactor model:
1861 send_to, recv_from, send, recv, and connect operations are submitted
1862 to the kernel via `uring_submit_op` and complete through the ring's
1863 CQE path.
1864
1865 The object is always owned by a `shared_ptr` managed by the service.
1866 In-flight ops hold an additional `shared_ptr` copy (`impl_ptr`) so
1867 the kernel's user-data pointer remains valid until the CQE arrives.
1868
1869 @par Thread Safety
1870 Distinct objects: Safe.
1871 Shared objects: Unsafe. One send and one recv may be in flight
1872 simultaneously, but two sends or two recvs must not overlap.
1873 */
1874 class BOOST_COROSIO_DECL uring_udp_socket final
1875 : public native_socket_base<
1876 uring_udp_socket,
1877 udp_socket::implementation,
1878 corosio::endpoint>
1879 {
1880 friend uring_udp_service;
1881
1882 int family_ = AF_UNSPEC; // cached at open_socket
1883 uring_scheduler* sched_ = nullptr;
1884 [[maybe_unused]] uring_udp_service* svc_ = nullptr;
1885
1886 // fd_ and local_endpoint_ live in native_socket_base, which also
1887 // provides native_handle/is_open/set_option/get_option/local_endpoint.
1888 corosio::endpoint remote_endpoint_;
1889
1890 // Per-fd op slots — embedded to eliminate per-call heap allocation.
1891 // Single-pending invariant per slot.
1892 uring_connect_op conn_;
1893 uring_dgram_send_op send_;
1894 uring_dgram_recv_op recv_;
1895 uring_wait_op wait_op_;
1896
1897 mutable detail::speculative_state spec_;
1898
1899 public:
1900 /** Construct with service and scheduler references.
1901
1902 Both refs must outlive this socket.
1903
1904 @param svc The owning service.
1905 @param sched The io_uring scheduler owned by the context.
1906 */
1907 132x explicit uring_udp_socket(
1908 uring_udp_service& svc, uring_scheduler& sched) noexcept
1909 264x : sched_(&sched)
1910 132x , svc_(&svc)
1911 {
1912 132x }
1913
1914 132x ~uring_udp_socket() override
1915 132x {
1916 132x if (fd_ >= 0)
1917 ::close(
1918 fd_); // LCOV_EXCL_LINE backstop: close_socket() clears fd_ before destroy
1919 132x }
1920
1921 // ----------------------------------------------------------------
1922 // udp_socket::implementation
1923 // ----------------------------------------------------------------
1924
1925 50x std::coroutine_handle<> send_to(
1926 std::coroutine_handle<> h,
1927 capy::executor_ref ex,
1928 buffer_param buf,
1929 endpoint dest,
1930 int flags,
1931 std::stop_token token,
1932 std::error_code* ec,
1933 std::size_t* bytes_out) override
1934 {
1935 50x sockaddr_storage addr{};
1936 50x socklen_t len = endpoint_to_sockaddr(dest, addr);
1937 100x return submit_send(h, ex, buf, len, addr, flags, token, ec, bytes_out);
1938 }
1939
1940 97x std::coroutine_handle<> recv_from(
1941 std::coroutine_handle<> h,
1942 capy::executor_ref ex,
1943 buffer_param buf,
1944 endpoint* source,
1945 int flags,
1946 std::stop_token token,
1947 std::error_code* ec,
1948 std::size_t* bytes_out) override
1949 {
1950 97x return submit_recv(
1951 97x h, ex, buf, source != nullptr, source, flags, token, ec, bytes_out);
1952 }
1953
1954 149x std::coroutine_handle<> send(
1955 std::coroutine_handle<> h,
1956 capy::executor_ref ex,
1957 buffer_param buf,
1958 int flags,
1959 std::stop_token token,
1960 std::error_code* ec,
1961 std::size_t* bytes_out) override
1962 {
1963 149x sockaddr_storage empty{};
1964 298x return submit_send(h, ex, buf, 0, empty, flags, token, ec, bytes_out);
1965 }
1966
1967 67x std::coroutine_handle<> recv(
1968 std::coroutine_handle<> h,
1969 capy::executor_ref ex,
1970 buffer_param buf,
1971 int flags,
1972 std::stop_token token,
1973 std::error_code* ec,
1974 std::size_t* bytes_out) override
1975 {
1976 67x return submit_recv(
1977 67x h, ex, buf, false, nullptr, flags, token, ec, bytes_out);
1978 }
1979
1980 61x std::coroutine_handle<> connect(
1981 std::coroutine_handle<> h,
1982 capy::executor_ref ex,
1983 endpoint ep,
1984 std::stop_token token,
1985 std::error_code* ec) override
1986 {
1987 61x bool stop_now = token.stop_possible() && token.stop_requested();
1988 61x if (stop_now)
1989 {
1990 22x if (sched_->try_consume_inline_budget())
1991 {
1992 21x if (ec)
1993 21x *ec = capy::error::canceled;
1994 21x conn_.cont.h = h;
1995 21x return dispatch_coro(ex, conn_.cont);
1996 }
1997 1x conn_.addrlen = to_sockaddr(ep, family_, conn_.addr);
1998 1x conn_.prepare(
1999 2x h, ex, ec, fd_, sched_, shared_from_this(), ep,
2000 &remote_endpoint_, &local_endpoint_, token);
2001 1x conn_.cancelled.store(true, std::memory_order_release);
2002 1x sched_->work_started();
2003 {
2004 1x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2005 1x sched_->push_completed_locked(&conn_);
2006 1x }
2007 1x return std::noop_coroutine();
2008 }
2009
2010 // io_uring's IORING_OP_CONNECT re-invokes connect(2) internally;
2011 // a prior speculative ::connect would leave EINPROGRESS → EALREADY.
2012 39x conn_.addrlen = to_sockaddr(ep, family_, conn_.addr);
2013 39x conn_.prepare(
2014 78x h, ex, ec, fd_, sched_, shared_from_this(), ep, &remote_endpoint_,
2015 &local_endpoint_, token);
2016 39x sched_->work_started();
2017 39x if (conn_.cancelled.load(std::memory_order_acquire))
2018 {
2019 uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2020 sched_->push_completed_locked(&conn_);
2021 return std::noop_coroutine();
2022 }
2023 39x uring_submit_op(*sched_, &conn_);
2024 39x return std::noop_coroutine();
2025 }
2026
2027 11x std::coroutine_handle<> wait(
2028 std::coroutine_handle<> h,
2029 capy::executor_ref ex,
2030 wait_type w,
2031 std::stop_token token,
2032 std::error_code* ec) override
2033 {
2034 11x int poll_flags = 0;
2035 11x switch (w)
2036 {
2037 9x case wait_type::read:
2038 9x poll_flags = POLLIN;
2039 9x break;
2040 1x case wait_type::write:
2041 1x poll_flags = POLLOUT;
2042 1x break;
2043 1x case wait_type::error:
2044 1x poll_flags = POLLPRI | POLLERR | POLLHUP;
2045 1x break;
2046 }
2047 11x wait_op_.prepare(
2048 22x h, ex, ec, fd_, sched_, shared_from_this(), poll_flags, token);
2049 11x sched_->work_started();
2050 11x if (wait_op_.cancelled.load(std::memory_order_acquire))
2051 {
2052 2x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2053 2x sched_->push_completed_locked(&wait_op_);
2054 2x return std::noop_coroutine();
2055 2x }
2056 9x uring_submit_op(*sched_, &wait_op_);
2057 9x return std::noop_coroutine();
2058 }
2059
2060 // native_handle / is_open / set_option / get_option / local_endpoint
2061 // are inherited from native_socket_base.
2062
2063 2x std::error_code shutdown(udp_socket::shutdown_type what) noexcept override
2064 {
2065 2x if (::shutdown(fd_, static_cast<int>(what)) != 0)
2066 1x return make_err(errno);
2067 1x return {};
2068 }
2069
2070 2x native_handle_type release_socket() noexcept override
2071 {
2072 // Flush while the fd is still open so the kernel resolves
2073 // pending SQEs before the caller can close and recycle the
2074 // number (same reasoning as close_socket).
2075 2x if (fd_ >= 0)
2076 2x sched_->cancel_and_flush(fd_);
2077 2x int fd = fd_;
2078 2x fd_ = -1;
2079 2x local_endpoint_ = endpoint{};
2080 2x remote_endpoint_ = endpoint{};
2081 2x return fd;
2082 }
2083
2084 6x void cancel() noexcept override
2085 {
2086 6x if (fd_ >= 0)
2087 6x sched_->submit_cancel_by_fd(fd_);
2088 6x }
2089
2090 /// Cancel in-flight ops, close the fd, and reset cached endpoints.
2091 /// Called by the service on close()/teardown.
2092 249x void close_socket() noexcept
2093 {
2094 249x if (fd_ >= 0)
2095 {
2096 117x sched_->cancel_and_flush(fd_);
2097 117x ::close(fd_);
2098 117x fd_ = -1;
2099 }
2100 249x local_endpoint_ = endpoint{};
2101 249x remote_endpoint_ = endpoint{};
2102 249x }
2103
2104 3x endpoint remote_endpoint() const noexcept override
2105 {
2106 3x return remote_endpoint_;
2107 }
2108
2109 private:
2110 199x std::coroutine_handle<> submit_send(
2111 std::coroutine_handle<> h,
2112 capy::executor_ref ex,
2113 buffer_param buffers,
2114 socklen_t dest_len,
2115 sockaddr_storage const& dest_storage,
2116 int flags,
2117 std::stop_token const& token,
2118 std::error_code* ec,
2119 std::size_t* bytes)
2120 {
2121 iovec iovecs[uring_max_iov];
2122 199x int iovec_count = copy_to_iovec(buffers, iovecs);
2123 199x bool stop_now = token.stop_possible() && token.stop_requested();
2124 199x bool empty_buf = (iovec_count == 0);
2125
2126 199x ssize_t n = 0;
2127 199x int err = 0;
2128 199x bool have_sync_res = stop_now || empty_buf;
2129 199x if (!have_sync_res && spec_.may_speculate_write())
2130 {
2131 154x msghdr msg{};
2132 154x msg.msg_iov = iovecs;
2133 154x msg.msg_iovlen = static_cast<decltype(msg.msg_iovlen)>(iovec_count);
2134 154x sockaddr_storage dest_copy = dest_storage;
2135 154x if (dest_len > 0)
2136 {
2137 27x msg.msg_name = &dest_copy;
2138 27x msg.msg_namelen = dest_len;
2139 }
2140 154x int native_flags = to_native_msg_flags(flags) | MSG_NOSIGNAL;
2141 do
2142 {
2143 154x n = ::sendmsg(fd_, &msg, native_flags);
2144 }
2145 154x while (n < 0 && errno == EINTR);
2146 154x if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
2147 {
2148 153x have_sync_res = true;
2149 153x if (n < 0)
2150 1x err = errno;
2151 }
2152 else
2153 {
2154 1x spec_.on_write_exhausted();
2155 }
2156 }
2157
2158 199x if (have_sync_res)
2159 {
2160 198x if (sched_->try_consume_inline_budget())
2161 {
2162 191x decode_io_result(
2163 1x ec, stop_now, err ? make_err(err) : std::error_code{},
2164 /*is_read=*/false, /*bytes=*/0, /*empty_buffer=*/false);
2165 190x if (bytes)
2166 190x *bytes = (n < 0) ? 0u : static_cast<std::size_t>(n);
2167 190x send_.cont.h = h;
2168 190x return dispatch_coro(ex, send_.cont);
2169 }
2170 16x send_.prepare(
2171 16x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_,
2172 buffers, dest_len, dest_storage, to_native_msg_flags(flags),
2173 token);
2174 8x if (stop_now)
2175 2x send_.cancelled.store(true, std::memory_order_release);
2176 else
2177 6x send_.res = (n < 0) ? -err : static_cast<int>(n);
2178 8x sched_->work_started();
2179 {
2180 8x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2181 8x sched_->push_completed_locked(&send_);
2182 8x }
2183 8x return std::noop_coroutine();
2184 }
2185
2186 2x send_.prepare(
2187 2x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_, buffers,
2188 dest_len, dest_storage, to_native_msg_flags(flags), token);
2189 1x sched_->work_started();
2190 1x if (send_.cancelled.load(std::memory_order_acquire))
2191 {
2192 uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2193 sched_->push_completed_locked(&send_);
2194 return std::noop_coroutine();
2195 }
2196 1x uring_submit_op(*sched_, &send_);
2197 1x return std::noop_coroutine();
2198 }
2199
2200 164x std::coroutine_handle<> submit_recv(
2201 std::coroutine_handle<> h,
2202 capy::executor_ref ex,
2203 buffer_param buffers,
2204 bool want_source,
2205 corosio::endpoint* source_out,
2206 int flags,
2207 std::stop_token const& token,
2208 std::error_code* ec,
2209 std::size_t* bytes)
2210 {
2211 iovec iovecs[uring_max_iov];
2212 164x int iovec_count = copy_to_iovec(buffers, iovecs);
2213 164x bool stop_now = token.stop_possible() && token.stop_requested();
2214 164x bool empty_buf = (iovec_count == 0);
2215
2216 164x ssize_t n = 0;
2217 164x int err = 0;
2218 164x bool have_sync_res = stop_now || empty_buf;
2219 164x sockaddr_storage src_storage{};
2220 164x socklen_t src_namelen = 0;
2221 164x if (!have_sync_res && spec_.may_speculate_read())
2222 {
2223 119x msghdr msg{};
2224 119x msg.msg_iov = iovecs;
2225 119x msg.msg_iovlen = static_cast<decltype(msg.msg_iovlen)>(iovec_count);
2226 119x if (want_source)
2227 {
2228 74x msg.msg_name = &src_storage;
2229 74x msg.msg_namelen = sizeof(src_storage);
2230 }
2231 119x int native_flags = to_native_msg_flags(flags);
2232 do
2233 {
2234 119x n = ::recvmsg(fd_, &msg, native_flags);
2235 }
2236 119x while (n < 0 && errno == EINTR);
2237 119x if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
2238 {
2239 104x have_sync_res = true;
2240 104x if (n < 0)
2241 1x err = errno;
2242 104x src_namelen = (n >= 0) ? msg.msg_namelen : 0;
2243 }
2244 else
2245 {
2246 15x spec_.on_read_exhausted();
2247 }
2248 }
2249
2250 164x if (have_sync_res)
2251 {
2252 149x if (sched_->try_consume_inline_budget())
2253 {
2254 144x decode_io_result(
2255 1x ec, stop_now, err ? make_err(err) : std::error_code{},
2256 /*is_read=*/false, /*bytes=*/0, /*empty_buffer=*/false);
2257 143x if (bytes)
2258 143x *bytes = (n < 0) ? 0u : static_cast<std::size_t>(n);
2259 143x if (n >= 0 && want_source && source_out && !empty_buf)
2260 80x *source_out = sockaddr_to_endpoint(src_storage);
2261 143x recv_.cont.h = h;
2262 143x return dispatch_coro(ex, recv_.cont);
2263 }
2264 12x recv_.prepare(
2265 12x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_,
2266 buffers, source_out, want_source ? &write_ip_source : nullptr,
2267 to_native_msg_flags(flags), token);
2268 6x if (stop_now)
2269 2x recv_.cancelled.store(true, std::memory_order_release);
2270 else
2271 {
2272 4x recv_.res = (n < 0) ? -err : static_cast<int>(n);
2273 // Hand the speculative source over to do_handler's
2274 // source_writer so it translates into source_out the same
2275 // way the kernel-completed path would.
2276 4x if (n >= 0 && want_source)
2277 {
2278 2x recv_.source_storage = src_storage;
2279 2x recv_.source_len = src_namelen;
2280 }
2281 }
2282 6x sched_->work_started();
2283 {
2284 6x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2285 6x sched_->push_completed_locked(&recv_);
2286 6x }
2287 6x return std::noop_coroutine();
2288 }
2289
2290 30x recv_.prepare(
2291 30x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_, buffers,
2292 source_out, want_source ? &write_ip_source : nullptr,
2293 to_native_msg_flags(flags), token);
2294 15x sched_->work_started();
2295 30x if (recv_.iovec_count == 0 ||
2296 15x recv_.cancelled.load(std::memory_order_acquire))
2297 {
2298 uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2299 sched_->push_completed_locked(&recv_);
2300 return std::noop_coroutine();
2301 }
2302 15x uring_submit_op(*sched_, &recv_);
2303 15x return std::noop_coroutine();
2304 }
2305
2306 6x static void write_ip_source(
2307 void* ctx, sockaddr_storage const& s, socklen_t /*len*/) noexcept
2308 {
2309 6x if (auto* out = static_cast<corosio::endpoint*>(ctx))
2310 6x *out = sockaddr_to_endpoint(s);
2311 6x }
2312 };
2313
2314 /** UDP socket service for io_uring.
2315
2316 Owns all `uring_udp_socket` implementations for an `io_context`.
2317 Satisfies the `udp_service` interface so the generic `udp_socket`
2318 front-end can call `open_datagram_socket` and `bind_datagram`
2319 transparently.
2320
2321 Socket impls are reference-counted inside the service map; raw
2322 pointers returned from `construct()` remain valid until `destroy()`
2323 or `shutdown()` is called.
2324
2325 @par Thread Safety
2326 All public member functions are thread-safe.
2327 */
2328 class BOOST_COROSIO_DECL uring_udp_service final
2329 : public uring_socket_service_base<
2330 uring_udp_service,
2331 udp_service,
2332 uring_udp_socket>
2333 {
2334 using base_service = uring_socket_service_base<
2335 uring_udp_service,
2336 udp_service,
2337 uring_udp_socket>;
2338
2339 public:
2340 /// Identifies this service for `execution_context` lookup.
2341 using key_type = udp_service;
2342
2343 /** Construct the UDP service.
2344
2345 @param ctx The owning execution context. The io_uring scheduler
2346 must already be registered.
2347 */
2348 799x explicit uring_udp_service(capy::execution_context& ctx) : base_service(ctx)
2349 {
2350 799x }
2351
2352 // construct / destroy / shutdown / close / scheduler() are inherited
2353 // from uring_socket_service_base.
2354
2355 /** Open a datagram socket and associate it with an impl.
2356
2357 Creates a non-blocking, close-on-exec socket via `socket(2)`.
2358
2359 @param impl The socket implementation to initialise.
2360 @param family Address family (e.g. `AF_INET`, `AF_INET6`).
2361 @param type Socket type (`SOCK_DGRAM`).
2362 @param protocol Protocol number (`IPPROTO_UDP`).
2363 @return Error code on failure, empty on success.
2364 */
2365 117x std::error_code open_datagram_socket(
2366 udp_socket::implementation& impl,
2367 int family,
2368 int type,
2369 int protocol) override
2370 {
2371 117x auto& sock = static_cast<uring_udp_socket&>(impl);
2372 int fd =
2373 117x ::socket(family, type | SOCK_NONBLOCK | SOCK_CLOEXEC, protocol);
2374 117x if (fd < 0)
2375 1x return make_err(errno);
2376 // LCOV_EXCL_START: dead — open() guards is_open(), so open_socket
2377 // never runs against an already-open fd (assign_socket handles that).
2378 if (sock.fd_ >= 0)
2379 {
2380 sched_->submit_cancel_by_fd(sock.fd_);
2381 ::close(sock.fd_);
2382 }
2383 // LCOV_EXCL_STOP
2384 116x sock.fd_ = fd;
2385 116x sock.family_ = family;
2386 116x if (family == AF_INET6)
2387 {
2388 12x int one = 1;
2389 12x ::setsockopt(fd, IPPROTO_IPV6, IPV6_V6ONLY, &one, sizeof(one));
2390 }
2391 116x return {};
2392 }
2393
2394 /** Adopt a pre-created fd into an impl.
2395
2396 Takes ownership of `fd` on success; the caller retains
2397 ownership on failure.
2398
2399 @param impl The socket implementation to assign to.
2400 @param fd A valid, open, non-blocking IP datagram fd.
2401 @return Error code on failure, empty on success.
2402 */
2403 11x std::error_code assign_socket(
2404 udp_socket::implementation& impl, native_handle_type fd) override
2405 {
2406 11x auto& sock = static_cast<uring_udp_socket&>(impl);
2407 11x int nfd = static_cast<int>(fd);
2408 11x if (nfd >= 0 && nfd == sock.fd_)
2409 2x return std::make_error_code(std::errc::invalid_argument);
2410 9x if (auto ec = validate_socket_fd(nfd, SOCK_DGRAM, true))
2411 5x return ec;
2412
2413 4x if (sock.fd_ >= 0)
2414 {
2415 1x sched_->cancel_and_flush(sock.fd_);
2416 1x ::close(sock.fd_);
2417 }
2418 4x sock.fd_ = nfd;
2419
2420 4x sock.local_endpoint_ = endpoint{};
2421 4x sock.remote_endpoint_ = endpoint{};
2422
2423 4x sockaddr_storage local{};
2424 4x socklen_t local_len = sizeof(local);
2425 4x if (::getsockname(
2426 4x sock.fd_, reinterpret_cast<sockaddr*>(&local), &local_len) == 0)
2427 {
2428 4x sock.local_endpoint_ = sockaddr_to_endpoint(local);
2429 4x sock.family_ = local.ss_family;
2430 }
2431
2432 4x sockaddr_storage remote{};
2433 4x socklen_t remote_len = sizeof(remote);
2434 4x if (::getpeername(
2435 4x sock.fd_, reinterpret_cast<sockaddr*>(&remote), &remote_len) ==
2436 0)
2437 1x sock.remote_endpoint_ = sockaddr_to_endpoint(remote);
2438
2439 4x return {};
2440 }
2441
2442 /** Bind the socket and capture the local endpoint via `getsockname`.
2443
2444 @param impl The socket implementation to bind.
2445 @param ep The local endpoint to bind to.
2446 @return Error code on failure, empty on success.
2447 */
2448 std::error_code
2449 69x bind_datagram(udp_socket::implementation& impl, endpoint ep) override
2450 {
2451 69x auto& sock = static_cast<uring_udp_socket&>(impl);
2452 69x sockaddr_storage addr{};
2453 69x socklen_t len = endpoint_to_sockaddr(ep, addr);
2454 69x if (::bind(sock.fd_, reinterpret_cast<sockaddr*>(&addr), len) < 0)
2455 2x return make_err(errno);
2456
2457 67x sockaddr_storage local{};
2458 67x socklen_t local_len = sizeof(local);
2459 67x if (::getsockname(
2460 67x sock.fd_, reinterpret_cast<sockaddr*>(&local), &local_len) == 0)
2461 67x sock.local_endpoint_ = sockaddr_to_endpoint(local);
2462 67x return {};
2463 }
2464 };
2465
2466 /** Unix domain datagram socket implementation for io_uring.
2467
2468 Implements `local_datagram_socket::implementation` using a proactor
2469 model: send_to, recv_from, send, recv, and connect operations are
2470 submitted to the kernel via `uring_submit_op` and complete through
2471 the ring's CQE path.
2472
2473 The object is always owned by a `shared_ptr` managed by the service.
2474 In-flight ops hold an additional `shared_ptr` copy (`impl_ptr`) so
2475 the kernel's user-data pointer remains valid until the CQE arrives.
2476
2477 @par Thread Safety
2478 Distinct objects: Safe.
2479 Shared objects: Unsafe. One send and one recv may be in flight
2480 simultaneously, but two sends or two recvs must not overlap.
2481 */
2482 class BOOST_COROSIO_DECL uring_local_datagram_socket final
2483 : public native_socket_base<
2484 uring_local_datagram_socket,
2485 local_datagram_socket::implementation,
2486 corosio::local_endpoint>
2487 {
2488 friend uring_local_datagram_service;
2489
2490 uring_scheduler* sched_ = nullptr;
2491 [[maybe_unused]] uring_local_datagram_service* svc_ = nullptr;
2492
2493 // fd_ and local_endpoint_ live in native_socket_base, which also
2494 // provides native_handle/is_open/set_option/get_option/local_endpoint.
2495 corosio::local_endpoint remote_endpoint_;
2496
2497 // Per-fd op slots — embedded to eliminate per-call heap allocation.
2498 // Single-pending invariant per slot.
2499 uring_local_connect_op conn_;
2500 uring_dgram_send_op send_;
2501 uring_dgram_recv_op recv_;
2502 uring_wait_op wait_op_;
2503
2504 mutable detail::speculative_state spec_;
2505
2506 public:
2507 /** Construct with service and scheduler references.
2508
2509 Both refs must outlive this socket.
2510
2511 @param svc The owning service.
2512 @param sched The io_uring scheduler owned by the context.
2513 */
2514 126x explicit uring_local_datagram_socket(
2515 uring_local_datagram_service& svc, uring_scheduler& sched) noexcept
2516 252x : sched_(&sched)
2517 126x , svc_(&svc)
2518 {
2519 126x }
2520
2521 125x ~uring_local_datagram_socket() override
2522 125x {
2523 125x if (fd_ >= 0)
2524 ::close(
2525 fd_); // LCOV_EXCL_LINE backstop: close_socket() clears fd_ before destroy
2526 125x }
2527
2528 // ----------------------------------------------------------------
2529 // local_datagram_socket::implementation
2530 // ----------------------------------------------------------------
2531
2532 77x std::coroutine_handle<> send_to(
2533 std::coroutine_handle<> h,
2534 capy::executor_ref ex,
2535 buffer_param buf,
2536 corosio::local_endpoint dest,
2537 int flags,
2538 std::stop_token token,
2539 std::error_code* ec,
2540 std::size_t* bytes_out) override
2541 {
2542 77x sockaddr_storage addr{};
2543 77x socklen_t len = endpoint_to_sockaddr(dest, addr);
2544 154x return submit_send(h, ex, buf, len, addr, flags, token, ec, bytes_out);
2545 }
2546
2547 67x std::coroutine_handle<> recv_from(
2548 std::coroutine_handle<> h,
2549 capy::executor_ref ex,
2550 buffer_param buf,
2551 corosio::local_endpoint* source,
2552 int flags,
2553 std::stop_token token,
2554 std::error_code* ec,
2555 std::size_t* bytes_out) override
2556 {
2557 67x return submit_recv(
2558 67x h, ex, buf, source != nullptr, source, flags, token, ec, bytes_out);
2559 }
2560
2561 68x std::coroutine_handle<> send(
2562 std::coroutine_handle<> h,
2563 capy::executor_ref ex,
2564 buffer_param buf,
2565 int flags,
2566 std::stop_token token,
2567 std::error_code* ec,
2568 std::size_t* bytes_out) override
2569 {
2570 68x sockaddr_storage empty{};
2571 136x return submit_send(h, ex, buf, 0, empty, flags, token, ec, bytes_out);
2572 }
2573
2574 69x std::coroutine_handle<> recv(
2575 std::coroutine_handle<> h,
2576 capy::executor_ref ex,
2577 buffer_param buf,
2578 int flags,
2579 std::stop_token token,
2580 std::error_code* ec,
2581 std::size_t* bytes_out) override
2582 {
2583 69x return submit_recv(
2584 69x h, ex, buf, false, nullptr, flags, token, ec, bytes_out);
2585 }
2586
2587 47x std::coroutine_handle<> connect(
2588 std::coroutine_handle<> h,
2589 capy::executor_ref ex,
2590 corosio::local_endpoint ep,
2591 std::stop_token token,
2592 std::error_code* ec) override
2593 {
2594 47x bool stop_now = token.stop_possible() && token.stop_requested();
2595 47x if (stop_now)
2596 {
2597 21x if (sched_->try_consume_inline_budget())
2598 {
2599 20x if (ec)
2600 20x *ec = capy::error::canceled;
2601 20x conn_.cont.h = h;
2602 20x return dispatch_coro(ex, conn_.cont);
2603 }
2604 1x conn_.addrlen = to_sockaddr(ep, conn_.addr);
2605 1x conn_.prepare(
2606 2x h, ex, ec, fd_, sched_, shared_from_this(), ep,
2607 &remote_endpoint_, &local_endpoint_, token);
2608 1x conn_.cancelled.store(true, std::memory_order_release);
2609 1x sched_->work_started();
2610 {
2611 1x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2612 1x sched_->push_completed_locked(&conn_);
2613 1x }
2614 1x return std::noop_coroutine();
2615 }
2616
2617 // io_uring's IORING_OP_CONNECT re-invokes connect(2) internally;
2618 // a prior speculative ::connect would leave EINPROGRESS → EALREADY.
2619 26x conn_.addrlen = to_sockaddr(ep, conn_.addr);
2620 26x conn_.prepare(
2621 52x h, ex, ec, fd_, sched_, shared_from_this(), ep, &remote_endpoint_,
2622 &local_endpoint_, token);
2623 26x sched_->work_started();
2624 26x if (conn_.cancelled.load(std::memory_order_acquire))
2625 {
2626 uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2627 sched_->push_completed_locked(&conn_);
2628 return std::noop_coroutine();
2629 }
2630 26x uring_submit_op(*sched_, &conn_);
2631 26x return std::noop_coroutine();
2632 }
2633
2634 7x std::coroutine_handle<> wait(
2635 std::coroutine_handle<> h,
2636 capy::executor_ref ex,
2637 wait_type w,
2638 std::stop_token token,
2639 std::error_code* ec) override
2640 {
2641 7x int poll_flags = 0;
2642 7x switch (w)
2643 {
2644 5x case wait_type::read:
2645 5x poll_flags = POLLIN;
2646 5x break;
2647 1x case wait_type::write:
2648 1x poll_flags = POLLOUT;
2649 1x break;
2650 1x case wait_type::error:
2651 1x poll_flags = POLLPRI | POLLERR | POLLHUP;
2652 1x break;
2653 }
2654 7x wait_op_.prepare(
2655 14x h, ex, ec, fd_, sched_, shared_from_this(), poll_flags, token);
2656 7x sched_->work_started();
2657 7x if (wait_op_.cancelled.load(std::memory_order_acquire))
2658 {
2659 2x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2660 2x sched_->push_completed_locked(&wait_op_);
2661 2x return std::noop_coroutine();
2662 2x }
2663 5x uring_submit_op(*sched_, &wait_op_);
2664 5x return std::noop_coroutine();
2665 }
2666
2667 std::error_code
2668 3x shutdown(local_datagram_socket::shutdown_type what) noexcept override
2669 {
2670 3x if (::shutdown(fd_, static_cast<int>(what)) != 0)
2671 1x return make_err(errno);
2672 2x return {};
2673 }
2674
2675 // native_handle / is_open / set_option / get_option / local_endpoint
2676 // are inherited from native_socket_base.
2677
2678 2x native_handle_type release_socket() noexcept override
2679 {
2680 // Flush while the fd is still open so the kernel resolves
2681 // pending SQEs before the caller can close and recycle the
2682 // number (same reasoning as close_socket).
2683 2x if (fd_ >= 0)
2684 2x sched_->cancel_and_flush(fd_);
2685 2x int fd = fd_;
2686 2x fd_ = -1;
2687 2x local_endpoint_ = corosio::local_endpoint{};
2688 2x remote_endpoint_ = corosio::local_endpoint{};
2689 2x return fd;
2690 }
2691
2692 2x void cancel() noexcept override
2693 {
2694 2x if (fd_ >= 0)
2695 2x sched_->submit_cancel_by_fd(fd_);
2696 2x }
2697
2698 /// Cancel in-flight ops, close the fd, and reset cached endpoints.
2699 /// Called by the service on close()/teardown.
2700 227x void close_socket() noexcept
2701 {
2702 227x if (fd_ >= 0)
2703 {
2704 102x sched_->cancel_and_flush(fd_);
2705 102x ::close(fd_);
2706 102x fd_ = -1;
2707 }
2708 227x local_endpoint_ = corosio::local_endpoint{};
2709 227x remote_endpoint_ = corosio::local_endpoint{};
2710 227x }
2711
2712 1x corosio::local_endpoint remote_endpoint() const noexcept override
2713 {
2714 1x return remote_endpoint_;
2715 }
2716
2717 // LCOV_EXCL_START: the public bind routes through the
2718 // service's bind_socket; nothing calls the implementation
2719 // interface's bind on this backend.
2720 std::error_code bind(corosio::local_endpoint ep) noexcept override
2721 {
2722 sockaddr_storage addr{};
2723 socklen_t len = endpoint_to_sockaddr(ep, addr);
2724 if (::bind(fd_, reinterpret_cast<sockaddr*>(&addr), len) != 0)
2725 return make_err(errno);
2726
2727 sockaddr_storage local{};
2728 socklen_t local_len = sizeof(local);
2729 if (::getsockname(
2730 fd_, reinterpret_cast<sockaddr*>(&local), &local_len) == 0)
2731 local_endpoint_ = sockaddr_to_local_endpoint(local, local_len);
2732 return {};
2733 }
2734 // LCOV_EXCL_STOP
2735
2736 private:
2737 145x std::coroutine_handle<> submit_send(
2738 std::coroutine_handle<> h,
2739 capy::executor_ref ex,
2740 buffer_param buffers,
2741 socklen_t dest_len,
2742 sockaddr_storage const& dest_storage,
2743 int flags,
2744 std::stop_token const& token,
2745 std::error_code* ec,
2746 std::size_t* bytes)
2747 {
2748 iovec iovecs[uring_max_iov];
2749 145x int iovec_count = copy_to_iovec(buffers, iovecs);
2750 145x bool stop_now = token.stop_possible() && token.stop_requested();
2751 145x bool empty_buf = (iovec_count == 0);
2752
2753 145x ssize_t n = 0;
2754 145x int err = 0;
2755 145x bool have_sync_res = stop_now || empty_buf;
2756 145x if (!have_sync_res && spec_.may_speculate_write())
2757 {
2758 101x msghdr msg{};
2759 101x msg.msg_iov = iovecs;
2760 101x msg.msg_iovlen = static_cast<decltype(msg.msg_iovlen)>(iovec_count);
2761 101x sockaddr_storage dest_copy = dest_storage;
2762 101x if (dest_len > 0)
2763 {
2764 55x msg.msg_name = &dest_copy;
2765 55x msg.msg_namelen = dest_len;
2766 }
2767 101x int native_flags = to_native_msg_flags(flags) | MSG_NOSIGNAL;
2768 do
2769 {
2770 101x n = ::sendmsg(fd_, &msg, native_flags);
2771 }
2772 101x while (n < 0 && errno == EINTR);
2773 101x if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
2774 {
2775 78x have_sync_res = true;
2776 78x if (n < 0)
2777 1x err = errno;
2778 }
2779 else
2780 {
2781 23x spec_.on_write_exhausted();
2782 }
2783 }
2784
2785 145x if (have_sync_res)
2786 {
2787 122x if (sched_->try_consume_inline_budget())
2788 {
2789 121x decode_io_result(
2790 1x ec, stop_now, err ? make_err(err) : std::error_code{},
2791 /*is_read=*/false, /*bytes=*/0, /*empty_buffer=*/false);
2792 120x if (bytes)
2793 120x *bytes = (n < 0) ? 0u : static_cast<std::size_t>(n);
2794 120x send_.cont.h = h;
2795 120x return dispatch_coro(ex, send_.cont);
2796 }
2797 4x send_.prepare(
2798 4x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_,
2799 buffers, dest_len, dest_storage, to_native_msg_flags(flags),
2800 token);
2801 2x if (stop_now)
2802 2x send_.cancelled.store(true, std::memory_order_release);
2803 else
2804 send_.res = (n < 0) ? -err : static_cast<int>(n);
2805 2x sched_->work_started();
2806 {
2807 2x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2808 2x sched_->push_completed_locked(&send_);
2809 2x }
2810 2x return std::noop_coroutine();
2811 }
2812
2813 46x send_.prepare(
2814 46x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_, buffers,
2815 dest_len, dest_storage, to_native_msg_flags(flags), token);
2816 23x sched_->work_started();
2817 23x if (send_.cancelled.load(std::memory_order_acquire))
2818 {
2819 uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2820 sched_->push_completed_locked(&send_);
2821 return std::noop_coroutine();
2822 }
2823 23x uring_submit_op(*sched_, &send_);
2824 23x return std::noop_coroutine();
2825 }
2826
2827 136x std::coroutine_handle<> submit_recv(
2828 std::coroutine_handle<> h,
2829 capy::executor_ref ex,
2830 buffer_param buffers,
2831 bool want_source,
2832 corosio::local_endpoint* source_out,
2833 int flags,
2834 std::stop_token const& token,
2835 std::error_code* ec,
2836 std::size_t* bytes)
2837 {
2838 iovec iovecs[uring_max_iov];
2839 136x int iovec_count = copy_to_iovec(buffers, iovecs);
2840 136x bool stop_now = token.stop_possible() && token.stop_requested();
2841 136x bool empty_buf = (iovec_count == 0);
2842
2843 136x ssize_t n = 0;
2844 136x int err = 0;
2845 136x bool have_sync_res = stop_now || empty_buf;
2846 136x sockaddr_storage src_storage{};
2847 136x socklen_t src_namelen = 0;
2848 136x if (!have_sync_res && spec_.may_speculate_read())
2849 {
2850 59x msghdr msg{};
2851 59x msg.msg_iov = iovecs;
2852 59x msg.msg_iovlen = static_cast<decltype(msg.msg_iovlen)>(iovec_count);
2853 59x if (want_source)
2854 {
2855 29x msg.msg_name = &src_storage;
2856 29x msg.msg_namelen = sizeof(src_storage);
2857 }
2858 59x int native_flags = to_native_msg_flags(flags);
2859 do
2860 {
2861 59x n = ::recvmsg(fd_, &msg, native_flags);
2862 }
2863 59x while (n < 0 && errno == EINTR);
2864 59x if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
2865 {
2866 43x have_sync_res = true;
2867 43x if (n < 0)
2868 err = errno;
2869 43x src_namelen = (n >= 0) ? msg.msg_namelen : 0;
2870 }
2871 else
2872 {
2873 16x spec_.on_read_exhausted();
2874 }
2875 }
2876
2877 136x if (have_sync_res)
2878 {
2879 87x if (sched_->try_consume_inline_budget())
2880 {
2881 85x decode_io_result(
2882 ec, stop_now, err ? make_err(err) : std::error_code{},
2883 /*is_read=*/false, /*bytes=*/0, /*empty_buffer=*/false);
2884 85x if (bytes)
2885 85x *bytes = (n < 0) ? 0u : static_cast<std::size_t>(n);
2886 85x if (n >= 0 && want_source && source_out && !empty_buf)
2887 *source_out =
2888 42x sockaddr_to_local_endpoint(src_storage, src_namelen);
2889 85x recv_.cont.h = h;
2890 85x return dispatch_coro(ex, recv_.cont);
2891 }
2892 4x recv_.prepare(
2893 4x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_,
2894 buffers, source_out,
2895 want_source ? &write_local_source : nullptr,
2896 to_native_msg_flags(flags), token);
2897 2x if (stop_now)
2898 2x recv_.cancelled.store(true, std::memory_order_release);
2899 else
2900 {
2901 recv_.res = (n < 0) ? -err : static_cast<int>(n);
2902 // Hand the speculative source over to do_handler's
2903 // source_writer so it translates into source_out the same
2904 // way the kernel-completed path would.
2905 if (n >= 0 && want_source)
2906 {
2907 recv_.source_storage = src_storage;
2908 recv_.source_len = src_namelen;
2909 }
2910 }
2911 2x sched_->work_started();
2912 {
2913 2x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2914 2x sched_->push_completed_locked(&recv_);
2915 2x }
2916 2x return std::noop_coroutine();
2917 }
2918
2919 98x recv_.prepare(
2920 98x h, ex, ec, bytes, fd_, sched_, shared_from_this(), &spec_, buffers,
2921 source_out, want_source ? &write_local_source : nullptr,
2922 to_native_msg_flags(flags), token);
2923 49x sched_->work_started();
2924 98x if (recv_.iovec_count == 0 ||
2925 49x recv_.cancelled.load(std::memory_order_acquire))
2926 {
2927 uring_scheduler::lock_type lock(sched_->dispatch_mutex());
2928 sched_->push_completed_locked(&recv_);
2929 return std::noop_coroutine();
2930 }
2931 49x uring_submit_op(*sched_, &recv_);
2932 49x return std::noop_coroutine();
2933 }
2934
2935 22x static void write_local_source(
2936 void* ctx, sockaddr_storage const& s, socklen_t len) noexcept
2937 {
2938 22x if (auto* out = static_cast<corosio::local_endpoint*>(ctx))
2939 22x *out = sockaddr_to_local_endpoint(s, len);
2940 22x }
2941 };
2942
2943 /** Unix domain datagram socket service for io_uring.
2944
2945 Owns all `uring_local_datagram_socket` implementations for an
2946 `io_context`. Satisfies the `local_datagram_service` interface so the
2947 generic `local_datagram_socket` front-end can call `open_socket` and
2948 `bind_socket` transparently.
2949
2950 Socket impls are reference-counted inside the service map; raw
2951 pointers returned from `construct()` remain valid until `destroy()`
2952 or `shutdown()` is called.
2953
2954 @par Thread Safety
2955 All public member functions are thread-safe.
2956 */
2957 class BOOST_COROSIO_DECL uring_local_datagram_service final
2958 : public uring_socket_service_base<
2959 uring_local_datagram_service,
2960 local_datagram_service,
2961 uring_local_datagram_socket>
2962 {
2963 using base_service = uring_socket_service_base<
2964 uring_local_datagram_service,
2965 local_datagram_service,
2966 uring_local_datagram_socket>;
2967
2968 public:
2969 /// Identifies this service for `execution_context` lookup.
2970 using key_type = local_datagram_service;
2971
2972 /** Construct the local datagram service.
2973
2974 @param ctx The owning execution context. The io_uring scheduler
2975 must already be registered.
2976 */
2977 799x explicit uring_local_datagram_service(capy::execution_context& ctx)
2978 799x : base_service(ctx)
2979 {
2980 799x }
2981
2982 // construct / destroy / shutdown / close / scheduler() are inherited
2983 // from uring_socket_service_base.
2984
2985 /** Open an AF_UNIX datagram socket and associate it with an impl.
2986
2987 Creates a non-blocking, close-on-exec socket via `socket(2)`.
2988 `family` is always `AF_UNIX` for local datagram sockets.
2989
2990 @param impl The socket implementation to initialise.
2991 @param family Address family (`AF_UNIX`).
2992 @param type Socket type (`SOCK_DGRAM`).
2993 @param protocol Protocol number (typically 0).
2994 @return Error code on failure, empty on success.
2995 */
2996 50x std::error_code open_socket(
2997 local_datagram_socket::implementation& impl,
2998 int family,
2999 int type,
3000 int protocol) override
3001 {
3002 50x auto& sock = static_cast<uring_local_datagram_socket&>(impl);
3003 int fd =
3004 50x ::socket(family, type | SOCK_NONBLOCK | SOCK_CLOEXEC, protocol);
3005 50x if (fd < 0)
3006 1x return make_err(errno);
3007 // LCOV_EXCL_START: dead — open() guards is_open(), so open_socket
3008 // never runs against an already-open fd (assign_socket handles that).
3009 if (sock.fd_ >= 0)
3010 {
3011 sched_->submit_cancel_by_fd(sock.fd_);
3012 ::close(sock.fd_);
3013 }
3014 // LCOV_EXCL_STOP
3015 49x sock.fd_ = fd;
3016 49x return {};
3017 }
3018
3019 /** Adopt a pre-created fd into an impl (e.g. from `socketpair`).
3020
3021 Takes ownership of `fd` on success; the caller retains ownership
3022 on failure.
3023
3024 @param impl The socket implementation to assign to.
3025 @param fd A valid, open, non-blocking AF_UNIX datagram fd.
3026 @return Error code on failure, empty on success.
3027 */
3028 60x std::error_code assign_socket(
3029 local_datagram_socket::implementation& impl,
3030 native_handle_type fd) override
3031 {
3032 60x auto& sock = static_cast<uring_local_datagram_socket&>(impl);
3033 60x int nfd = static_cast<int>(fd);
3034 60x if (nfd >= 0 && nfd == sock.fd_)
3035 1x return std::make_error_code(std::errc::invalid_argument);
3036 59x if (auto ec = validate_socket_fd(nfd, SOCK_DGRAM, false))
3037 2x return ec;
3038
3039 57x if (sock.fd_ >= 0)
3040 {
3041 1x sched_->cancel_and_flush(sock.fd_);
3042 1x ::close(sock.fd_);
3043 }
3044 57x sock.fd_ = nfd;
3045
3046 57x sockaddr_storage local{};
3047 57x socklen_t local_len = sizeof(local);
3048 57x if (::getsockname(
3049 57x sock.fd_, reinterpret_cast<sockaddr*>(&local), &local_len) == 0)
3050 57x sock.local_endpoint_ = sockaddr_to_local_endpoint(local, local_len);
3051
3052 57x sockaddr_storage remote{};
3053 57x socklen_t remote_len = sizeof(remote);
3054 57x if (::getpeername(
3055 57x sock.fd_, reinterpret_cast<sockaddr*>(&remote), &remote_len) ==
3056 0)
3057 sock.remote_endpoint_ =
3058 57x sockaddr_to_local_endpoint(remote, remote_len);
3059
3060 57x return {};
3061 }
3062
3063 /** Bind the socket and capture the local endpoint via `getsockname`.
3064
3065 @param impl The socket implementation to bind.
3066 @param ep The local endpoint (path) to bind to.
3067 @return Error code on failure, empty on success.
3068 */
3069 37x std::error_code bind_socket(
3070 local_datagram_socket::implementation& impl,
3071 corosio::local_endpoint ep) override
3072 {
3073 37x auto& sock = static_cast<uring_local_datagram_socket&>(impl);
3074 37x sockaddr_storage addr{};
3075 37x socklen_t len = endpoint_to_sockaddr(ep, addr);
3076 37x if (::bind(sock.fd_, reinterpret_cast<sockaddr*>(&addr), len) < 0)
3077 4x return make_err(errno);
3078
3079 33x sockaddr_storage local{};
3080 33x socklen_t local_len = sizeof(local);
3081 33x if (::getsockname(
3082 33x sock.fd_, reinterpret_cast<sockaddr*>(&local), &local_len) == 0)
3083 33x sock.local_endpoint_ = sockaddr_to_local_endpoint(local, local_len);
3084 33x return {};
3085 }
3086 };
3087
3088 } // namespace boost::corosio::detail
3089
3090 #endif // BOOST_COROSIO_HAS_URING
3091
3092 #endif // BOOST_COROSIO_NATIVE_DETAIL_URING_URING_TYPES_HPP
3093