100.00% Lines (44/44) 100.00% Functions (15/15)
TLA Baseline Branch
Line Hits Code Line Hits Code
1   // 1   //
2   // Copyright (c) 2026 Steve Gerbino 2   // Copyright (c) 2026 Steve Gerbino
3   // Copyright (c) 2026 Michael Vandeberg 3   // Copyright (c) 2026 Michael Vandeberg
4   // 4   //
5   // Distributed under the Boost Software License, Version 1.0. (See accompanying 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) 6   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7   // 7   //
8   // Official repository: https://github.com/cppalliance/corosio 8   // Official repository: https://github.com/cppalliance/corosio
9   // 9   //
10   10  
11   #ifndef BOOST_COROSIO_NATIVE_NATIVE_LOCAL_STREAM_SOCKET_HPP 11   #ifndef BOOST_COROSIO_NATIVE_NATIVE_LOCAL_STREAM_SOCKET_HPP
12   #define BOOST_COROSIO_NATIVE_NATIVE_LOCAL_STREAM_SOCKET_HPP 12   #define BOOST_COROSIO_NATIVE_NATIVE_LOCAL_STREAM_SOCKET_HPP
13   13  
14   #include <boost/corosio/local_stream_socket.hpp> 14   #include <boost/corosio/local_stream_socket.hpp>
15   #include <boost/corosio/backend.hpp> 15   #include <boost/corosio/backend.hpp>
  16 + #include <boost/corosio/detail/op_base.hpp>
16   17  
17   #ifndef BOOST_COROSIO_MRDOCS 18   #ifndef BOOST_COROSIO_MRDOCS
18   #if BOOST_COROSIO_HAS_EPOLL 19   #if BOOST_COROSIO_HAS_EPOLL
19   #include <boost/corosio/native/detail/epoll/epoll_types.hpp> 20   #include <boost/corosio/native/detail/epoll/epoll_types.hpp>
20   #endif 21   #endif
21   22  
22   #if BOOST_COROSIO_HAS_SELECT 23   #if BOOST_COROSIO_HAS_SELECT
23   #include <boost/corosio/native/detail/select/select_types.hpp> 24   #include <boost/corosio/native/detail/select/select_types.hpp>
24   #endif 25   #endif
25   26  
26   #if BOOST_COROSIO_HAS_KQUEUE 27   #if BOOST_COROSIO_HAS_KQUEUE
27   #include <boost/corosio/native/detail/kqueue/kqueue_types.hpp> 28   #include <boost/corosio/native/detail/kqueue/kqueue_types.hpp>
28   #endif 29   #endif
29   30  
30   #if BOOST_COROSIO_HAS_URING 31   #if BOOST_COROSIO_HAS_URING
31   #include <boost/corosio/native/detail/uring/uring_types.hpp> 32   #include <boost/corosio/native/detail/uring/uring_types.hpp>
32   #endif 33   #endif
33   34  
34   #if BOOST_COROSIO_HAS_IOCP 35   #if BOOST_COROSIO_HAS_IOCP
35   #include <boost/corosio/native/detail/iocp/win_local_stream_service.hpp> 36   #include <boost/corosio/native/detail/iocp/win_local_stream_service.hpp>
36   #endif 37   #endif
37   #endif // !BOOST_COROSIO_MRDOCS 38   #endif // !BOOST_COROSIO_MRDOCS
38   39  
39   namespace boost::corosio { 40   namespace boost::corosio {
40   41  
41   /** An asynchronous Unix stream socket with devirtualized I/O operations. 42   /** An asynchronous Unix stream socket with devirtualized I/O operations.
42   43  
43   This class template inherits from @ref local_stream_socket and 44   This class template inherits from @ref local_stream_socket and
44   shadows the async operations (`read_some`, `write_some`, 45   shadows the async operations (`read_some`, `write_some`,
45   `connect`) with versions that call the backend implementation 46   `connect`) with versions that call the backend implementation
46   directly, allowing the compiler to inline through the entire 47   directly, allowing the compiler to inline through the entire
47   call chain. 48   call chain.
48   49  
49   Non-async operations (`open`, `close`, `cancel`, socket options) 50   Non-async operations (`open`, `close`, `cancel`, socket options)
50   remain unchanged and dispatch through the compiled library. 51   remain unchanged and dispatch through the compiled library.
51   52  
52   A `native_local_stream_socket` IS-A `local_stream_socket` and 53   A `native_local_stream_socket` IS-A `local_stream_socket` and
53   can be passed to any function expecting `local_stream_socket&` 54   can be passed to any function expecting `local_stream_socket&`
54   or `io_stream&`, in which case virtual dispatch is used 55   or `io_stream&`, in which case virtual dispatch is used
55   transparently. 56   transparently.
56   57  
57   @tparam Backend A backend tag value (e.g., `epoll`) whose type 58   @tparam Backend A backend tag value (e.g., `epoll`) whose type
58   provides the concrete implementation types. 59   provides the concrete implementation types.
59   60  
60   @par Thread Safety 61   @par Thread Safety
61   Same as @ref local_stream_socket. 62   Same as @ref local_stream_socket.
62   63  
63   @par Example 64   @par Example
64   @par !example connect 65   @par !example connect
65   66  
66   @see local_stream_socket, epoll_t, iocp_t 67   @see local_stream_socket, epoll_t, iocp_t
67   */ 68   */
68   template<auto Backend> 69   template<auto Backend>
69   class native_local_stream_socket : public local_stream_socket 70   class native_local_stream_socket : public local_stream_socket
70   { 71   {
71   using backend_type = decltype(Backend); 72   using backend_type = decltype(Backend);
72   using impl_type = typename backend_type::local_stream_socket_type; 73   using impl_type = typename backend_type::local_stream_socket_type;
73   using service_type = typename backend_type::local_stream_service_type; 74   using service_type = typename backend_type::local_stream_service_type;
74   75  
HITCBC 75   34 impl_type& get_impl() noexcept 76   26 impl_type& get_impl() noexcept
76   { 77   {
HITCBC 77   34 return *static_cast<impl_type*>(h_.get()); 78   26 return *static_cast<impl_type*>(h_.get());
78   } 79   }
79   80  
80   template<class MutableBufferSequence> 81   template<class MutableBufferSequence>
81   struct native_read_awaitable 82   struct native_read_awaitable
  83 + : detail::bytes_op_base<native_read_awaitable<MutableBufferSequence>>
82   { 84   {
83   native_local_stream_socket& self_; 85   native_local_stream_socket& self_;
84 - std::stop_token token_;  
85 - mutable std::error_code ec_;  
86 - mutable std::size_t bytes_transferred_ = 0;  
87   MutableBufferSequence buffers_; 86   MutableBufferSequence buffers_;
88   87  
HITCBC 89   8 native_read_awaitable( 88   8 native_read_awaitable(
90   native_local_stream_socket& self, 89   native_local_stream_socket& self,
91   MutableBufferSequence buffers) noexcept 90   MutableBufferSequence buffers) noexcept
HITCBC 92   8 : self_(self) 91   8 : self_(self)
HITCBC 93   8 , buffers_(std::move(buffers)) 92   8 , buffers_(std::move(buffers))
94   { 93   {
HITCBC 95   8 } 94   8 }
96   95  
ECB 97 - 8 bool await_ready() const noexcept 96 + std::coroutine_handle<>
HITGIC 98 - { 97 + 6 dispatch(std::coroutine_handle<> h, capy::executor_ref ex) const
99 - // A pre-set ec_ means the initiator failed before  
100 - // dispatch (e.g. a closed object).  
DCB 101 - 8 return static_cast<bool>(ec_) || token_.stop_requested();  
102 - }  
103 -  
DCB 104 - 8 [[nodiscard]] capy::io_result<std::size_t> await_resume() const noexcept  
105 - {  
DCB 106 - 8 if (token_.stop_requested())  
DCB 107 - 2 return {make_error_code(std::errc::operation_canceled), 0};  
DCB 108 - 6 return {ec_, bytes_transferred_};  
109 - }  
110 -  
DCB 111 - 8 auto await_suspend(std::coroutine_handle<> h, capy::io_env const* env)  
112 - -> std::coroutine_handle<>  
113 - token_ = env->stop_token;  
ECB 114   8 { 98   {
HITCBC 115   24 return self_.get_impl().read_some( 99   18 return self_.get_impl().read_some(
HITCBC 116 - 24 h, env->executor, buffers_, token_, &ec_, &bytes_transferred_); 100 + 18 h, ex, buffers_, this->token_, &this->ec_, &this->bytes_);
117   } 101   }
118   }; 102   };
119   103  
120   template<class ConstBufferSequence> 104   template<class ConstBufferSequence>
121   struct native_write_awaitable 105   struct native_write_awaitable
  106 + : detail::bytes_op_base<native_write_awaitable<ConstBufferSequence>>
122   { 107   {
123   native_local_stream_socket& self_; 108   native_local_stream_socket& self_;
124 - std::stop_token token_;  
125 - mutable std::error_code ec_;  
126 - mutable std::size_t bytes_transferred_ = 0;  
127   ConstBufferSequence buffers_; 109   ConstBufferSequence buffers_;
128   110  
HITCBC 129   8 native_write_awaitable( 111   8 native_write_awaitable(
130   native_local_stream_socket& self, 112   native_local_stream_socket& self,
131   ConstBufferSequence buffers) noexcept 113   ConstBufferSequence buffers) noexcept
HITCBC 132   8 : self_(self) 114   8 : self_(self)
HITCBC 133   8 , buffers_(std::move(buffers)) 115   8 , buffers_(std::move(buffers))
134   { 116   {
HITCBC 135   8 } 117   8 }
136   118  
ECB 137 - 8 bool await_ready() const noexcept 119 + std::coroutine_handle<>
HITGIC 138 - { 120 + 6 dispatch(std::coroutine_handle<> h, capy::executor_ref ex) const
139 - // A pre-set ec_ means the initiator failed before  
140 - // dispatch (e.g. a closed object).  
DCB 141 - 8 return static_cast<bool>(ec_) || token_.stop_requested();  
142 - }  
143 -  
DCB 144 - 8 [[nodiscard]] capy::io_result<std::size_t> await_resume() const noexcept  
145 - {  
DCB 146 - 8 if (token_.stop_requested())  
DCB 147 - 2 return {make_error_code(std::errc::operation_canceled), 0};  
DCB 148 - 6 return {ec_, bytes_transferred_};  
149 - }  
150 -  
DCB 151 - 8 auto await_suspend(std::coroutine_handle<> h, capy::io_env const* env)  
152 - -> std::coroutine_handle<>  
153 - token_ = env->stop_token;  
ECB 154   8 { 121   {
HITCBC 155   24 return self_.get_impl().write_some( 122   18 return self_.get_impl().write_some(
HITCBC 156 - 24 h, env->executor, buffers_, token_, &ec_, &bytes_transferred_); 123 + 18 h, ex, buffers_, this->token_, &this->ec_, &this->bytes_);
157   } 124   }
158   }; 125   };
159   126  
160 - struct native_wait_awaitable 127 + struct native_wait_awaitable : detail::void_op_base<native_wait_awaitable>
161   { 128   {
162   native_local_stream_socket& self_; 129   native_local_stream_socket& self_;
163 - std::stop_token token_;  
164 - mutable std::error_code ec_;  
165   wait_type w_; 130   wait_type w_;
166   131  
HITCBC 167   6 native_wait_awaitable( 132   6 native_wait_awaitable(
168   native_local_stream_socket& self, wait_type w) noexcept 133   native_local_stream_socket& self, wait_type w) noexcept
HITCBC 169   6 : self_(self) 134   6 : self_(self)
HITCBC 170   6 , w_(w) 135   6 , w_(w)
171   { 136   {
HITCBC 172   6 } 137   6 }
173   138  
ECB 174 - 6 bool await_ready() const noexcept 139 + std::coroutine_handle<>
HITGIC 175 - { 140 + 4 dispatch(std::coroutine_handle<> h, capy::executor_ref ex) const
176 - // A pre-set ec_ means the initiator failed before  
177 - // dispatch (e.g. auto-open).  
DCB 178 - 6 return static_cast<bool>(ec_) || token_.stop_requested();  
179 - }  
180 -  
DCB 181 - 6 [[nodiscard]] capy::io_result<> await_resume() const noexcept  
182 - {  
DCB 183 - 6 if (token_.stop_requested())  
DCB 184 - 2 return {make_error_code(std::errc::operation_canceled)};  
DCB 185 - 4 return {ec_};  
186 - }  
187 -  
DCB 188 - 6 auto await_suspend(std::coroutine_handle<> h, capy::io_env const* env)  
189 - -> std::coroutine_handle<>  
190   { 141   {
HITCBC 191 - 6 token_ = env->stop_token; 142 + 4 return self_.get_impl().wait(h, ex, w_, this->token_, &this->ec_);
DCB 192 - 6 return self_.get_impl().wait(h, env->executor, w_, token_, &ec_);  
193   } 143   }
194   }; 144   };
195   145  
196   struct native_connect_awaitable 146   struct native_connect_awaitable
  147 + : detail::void_op_base<native_connect_awaitable>
197   { 148   {
198   native_local_stream_socket& self_; 149   native_local_stream_socket& self_;
199 - std::stop_token token_;  
200 - mutable std::error_code ec_;  
201   corosio::local_endpoint endpoint_; 150   corosio::local_endpoint endpoint_;
202   151  
HITCBC 203   12 native_connect_awaitable( 152   12 native_connect_awaitable(
204   native_local_stream_socket& self, 153   native_local_stream_socket& self,
205   corosio::local_endpoint ep) noexcept 154   corosio::local_endpoint ep) noexcept
HITCBC 206   12 : self_(self) 155   12 : self_(self)
HITCBC 207   12 , endpoint_(ep) 156   12 , endpoint_(ep)
208   { 157   {
HITCBC 209   12 } 158   12 }
210   159  
ECB 211 - 12 bool await_ready() const noexcept 160 + std::coroutine_handle<>
HITGIC 212 - { 161 + 10 dispatch(std::coroutine_handle<> h, capy::executor_ref ex) const
213 - // A pre-set ec_ means the initiator failed before  
214 - // dispatch (e.g. a closed object).  
DCB 215 - 12 return static_cast<bool>(ec_) || token_.stop_requested();  
216 - }  
217 -  
DCB 218 - 12 [[nodiscard]] capy::io_result<> await_resume() const noexcept  
219 - {  
DCB 220 - 12 if (token_.stop_requested())  
DCB 221 - 2 return {make_error_code(std::errc::operation_canceled)};  
DCB 222 - 10 return {ec_};  
223 - }  
224 -  
DCB 225 - 12 auto await_suspend(std::coroutine_handle<> h, capy::io_env const* env)  
226 - -> std::coroutine_handle<>  
227 - token_ = env->stop_token;  
ECB 228   12 { 162   {
HITCBC 229   36 return self_.get_impl().connect( 163   30 return self_.get_impl().connect(
HITCBC 230 - 36 h, env->executor, endpoint_, token_, &ec_); 164 + 30 h, ex, endpoint_, this->token_, &this->ec_);
231   } 165   }
232   }; 166   };
233   167  
234   public: 168   public:
235   /** Construct a native socket from an execution context. 169   /** Construct a native socket from an execution context.
236   170  
237   @param ctx The execution context that will own this socket. 171   @param ctx The execution context that will own this socket.
238   */ 172   */
HITCBC 239   40 explicit native_local_stream_socket(capy::execution_context& ctx) 173   40 explicit native_local_stream_socket(capy::execution_context& ctx)
HITCBC 240   40 : io_object(create_handle<service_type>(ctx)) 174   40 : io_object(create_handle<service_type>(ctx))
241   { 175   {
HITCBC 242   40 } 176   40 }
243   177  
244   /** Construct a native socket from an executor. 178   /** Construct a native socket from an executor.
245   179  
246   @param ex The executor whose context will own the socket. 180   @param ex The executor whose context will own the socket.
247   */ 181   */
248   template<class Ex> 182   template<class Ex>
249   requires(!std::same_as< 183   requires(!std::same_as<
250   std::remove_cvref_t<Ex>, 184   std::remove_cvref_t<Ex>,
251   native_local_stream_socket>) && 185   native_local_stream_socket>) &&
252   capy::Executor<Ex> 186   capy::Executor<Ex>
253   explicit native_local_stream_socket(Ex const& ex) 187   explicit native_local_stream_socket(Ex const& ex)
254   : native_local_stream_socket(ex.context()) 188   : native_local_stream_socket(ex.context())
255   { 189   {
256   } 190   }
257   191  
258   /// Move construct. 192   /// Move construct.
HITCBC 259   6 native_local_stream_socket(native_local_stream_socket&&) noexcept = default; 193   6 native_local_stream_socket(native_local_stream_socket&&) noexcept = default;
260   194  
261   /// Move assign. 195   /// Move assign.
262   native_local_stream_socket& 196   native_local_stream_socket&
263   operator=(native_local_stream_socket&&) noexcept = default; 197   operator=(native_local_stream_socket&&) noexcept = default;
264   198  
265   native_local_stream_socket(native_local_stream_socket const&) = delete; 199   native_local_stream_socket(native_local_stream_socket const&) = delete;
266   native_local_stream_socket& 200   native_local_stream_socket&
267   operator=(native_local_stream_socket const&) = delete; 201   operator=(native_local_stream_socket const&) = delete;
268   202  
269   /** Asynchronously read data from the socket. 203   /** Asynchronously read data from the socket.
270   204  
271   Calls the backend implementation directly, bypassing virtual 205   Calls the backend implementation directly, bypassing virtual
272   dispatch. Otherwise identical to @ref io_stream::read_some. 206   dispatch. Otherwise identical to @ref io_stream::read_some.
273   207  
274   @param buffers The buffer sequence to read into. 208   @param buffers The buffer sequence to read into.
275   209  
276   @return An awaitable yielding `(error_code, std::size_t)`. 210   @return An awaitable yielding `(error_code, std::size_t)`.
277   */ 211   */
278   template<capy::MutableBufferSequence MB> 212   template<capy::MutableBufferSequence MB>
HITCBC 279   8 [[nodiscard]] auto read_some(MB const& buffers) 213   8 [[nodiscard]] auto read_some(MB const& buffers)
280   { 214   {
HITCBC 281   8 return native_read_awaitable<MB>(*this, buffers); 215   8 return native_read_awaitable<MB>(*this, buffers);
282   } 216   }
283   217  
284   /** Asynchronously write data to the socket. 218   /** Asynchronously write data to the socket.
285   219  
286   Calls the backend implementation directly, bypassing virtual 220   Calls the backend implementation directly, bypassing virtual
287   dispatch. Otherwise identical to @ref io_stream::write_some. 221   dispatch. Otherwise identical to @ref io_stream::write_some.
288   222  
289   @param buffers The buffer sequence to write from. 223   @param buffers The buffer sequence to write from.
290   224  
291   @return An awaitable yielding `(error_code, std::size_t)`. 225   @return An awaitable yielding `(error_code, std::size_t)`.
292   */ 226   */
293   template<capy::ConstBufferSequence CB> 227   template<capy::ConstBufferSequence CB>
HITCBC 294   8 [[nodiscard]] auto write_some(CB const& buffers) 228   8 [[nodiscard]] auto write_some(CB const& buffers)
295   { 229   {
HITCBC 296   8 return native_write_awaitable<CB>(*this, buffers); 230   8 return native_write_awaitable<CB>(*this, buffers);
297   } 231   }
298   232  
299   /** Asynchronously connect to a remote endpoint. 233   /** Asynchronously connect to a remote endpoint.
300   234  
301   Calls the backend implementation directly, bypassing virtual 235   Calls the backend implementation directly, bypassing virtual
302   dispatch. Otherwise identical to @ref local_stream_socket::connect. 236   dispatch. Otherwise identical to @ref local_stream_socket::connect.
303   237  
304   If the socket is not already open, it is opened automatically. 238   If the socket is not already open, it is opened automatically.
305   239  
306   @param ep The local endpoint (path) to connect to. 240   @param ep The local endpoint (path) to connect to.
307   241  
308   @return An awaitable yielding `io_result<>`. 242   @return An awaitable yielding `io_result<>`.
309   243  
310   If the socket needs to be opened and the open fails, the 244   If the socket needs to be opened and the open fails, the
311   awaitable completes immediately with that error. 245   awaitable completes immediately with that error.
312   */ 246   */
HITCBC 313   12 [[nodiscard]] auto connect(corosio::local_endpoint ep) 247   12 [[nodiscard]] auto connect(corosio::local_endpoint ep)
314   { 248   {
HITCBC 315   12 native_connect_awaitable aw(*this, ep); 249   12 native_connect_awaitable aw(*this, ep);
HITCBC 316   12 if (!is_open()) 250   12 if (!is_open())
HITCBC 317   10 aw.ec_ = open(); 251   10 aw.ec_ = open();
HITCBC 318   12 return aw; 252   12 return aw;
319   } 253   }
320   254  
321   /** Asynchronously wait for the socket to be ready. 255   /** Asynchronously wait for the socket to be ready.
322   256  
323   Calls the backend implementation directly, bypassing virtual 257   Calls the backend implementation directly, bypassing virtual
324   dispatch. Otherwise identical to @ref local_stream_socket::wait. 258   dispatch. Otherwise identical to @ref local_stream_socket::wait.
325   259  
326   @param w The wait direction (read, write, or error). 260   @param w The wait direction (read, write, or error).
327   261  
328   @return An awaitable yielding `io_result<>`. 262   @return An awaitable yielding `io_result<>`.
329   */ 263   */
HITCBC 330   6 [[nodiscard]] auto wait(wait_type w) 264   6 [[nodiscard]] auto wait(wait_type w)
331   { 265   {
HITCBC 332   6 return native_wait_awaitable(*this, w); 266   6 return native_wait_awaitable(*this, w);
333   } 267   }
334   }; 268   };
335   269  
336   } // namespace boost::corosio 270   } // namespace boost::corosio
337   271  
338   #endif // BOOST_COROSIO_NATIVE_NATIVE_LOCAL_STREAM_SOCKET_HPP 272   #endif // BOOST_COROSIO_NATIVE_NATIVE_LOCAL_STREAM_SOCKET_HPP