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

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