80.81% Lines (160/198) 100.00% Functions (28/28)
TLA Baseline Branch
Line Hits Code Line Hits Code
1   // 1   //
2   // Copyright (c) 2025 Vinnie Falco (vinnie.falco@gmail.com) 2   // Copyright (c) 2025 Vinnie Falco (vinnie.falco@gmail.com)
3   // Copyright (c) 2026 Steve Gerbino 3   // Copyright (c) 2026 Steve Gerbino
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_TEST_MOCKET_HPP 11   #ifndef BOOST_COROSIO_TEST_MOCKET_HPP
12   #define BOOST_COROSIO_TEST_MOCKET_HPP 12   #define BOOST_COROSIO_TEST_MOCKET_HPP
13   13  
14   #include <boost/corosio/detail/except.hpp> 14   #include <boost/corosio/detail/except.hpp>
15   #include <boost/corosio/io_context.hpp> 15   #include <boost/corosio/io_context.hpp>
16   #include <boost/corosio/socket_option.hpp> 16   #include <boost/corosio/socket_option.hpp>
17   #include <boost/corosio/tcp_acceptor.hpp> 17   #include <boost/corosio/tcp_acceptor.hpp>
18   #include <boost/corosio/tcp_socket.hpp> 18   #include <boost/corosio/tcp_socket.hpp>
19   #include <boost/capy/buffers/buffer_copy.hpp> 19   #include <boost/capy/buffers/buffer_copy.hpp>
20   #include <boost/capy/buffers/make_buffer.hpp> 20   #include <boost/capy/buffers/make_buffer.hpp>
21   #include <boost/capy/error.hpp> 21   #include <boost/capy/error.hpp>
  22 + #include <boost/capy/ex/io_env.hpp>
22   #include <boost/capy/ex/run_async.hpp> 23   #include <boost/capy/ex/run_async.hpp>
23   #include <boost/capy/io_result.hpp> 24   #include <boost/capy/io_result.hpp>
24   #include <boost/capy/task.hpp> 25   #include <boost/capy/task.hpp>
25   #include <boost/capy/test/fuse.hpp> 26   #include <boost/capy/test/fuse.hpp>
26   27  
27   #include <cstddef> 28   #include <cstddef>
28   #include <cstdio> 29   #include <cstdio>
29   #include <cstring> 30   #include <cstring>
30   #include <stdexcept> 31   #include <stdexcept>
31   #include <string> 32   #include <string>
32   #include <system_error> 33   #include <system_error>
33   #include <tuple> 34   #include <tuple>
34   #include <utility> 35   #include <utility>
35   36  
36   namespace boost::corosio::test { 37   namespace boost::corosio::test {
37   38  
38   /** A mock socket for testing I/O operations. 39   /** A mock socket for testing I/O operations.
39   40  
40   This class provides a testable socket-like interface where data 41   This class provides a testable socket-like interface where data
41   can be staged for reading and expected data can be validated on 42   can be staged for reading and expected data can be validated on
42   writes. A mocket is paired with a regular socket using 43   writes. A mocket is paired with a regular socket using
43   @ref make_mocket_pair, allowing bidirectional communication testing. 44   @ref make_mocket_pair, allowing bidirectional communication testing.
44   45  
45   When reading, data comes from the `provide()` buffer first. 46   When reading, data comes from the `provide()` buffer first.
46   When writing, data is validated against the `expect()` buffer. 47   When writing, data is validated against the `expect()` buffer.
47   Once buffers are exhausted, I/O passes through to the underlying 48   Once buffers are exhausted, I/O passes through to the underlying
48   socket connection. 49   socket connection.
49   50  
50   Satisfies the `capy::Stream` concept. 51   Satisfies the `capy::Stream` concept.
51   52  
52   @tparam Socket The underlying socket type (default `tcp_socket`). 53   @tparam Socket The underlying socket type (default `tcp_socket`).
53   54  
54   @par Thread Safety 55   @par Thread Safety
55   Not thread-safe. All operations must occur on a single thread. 56   Not thread-safe. All operations must occur on a single thread.
56   All coroutines using the mocket must be suspended when calling 57   All coroutines using the mocket must be suspended when calling
57   `expect()` or `provide()`. 58   `expect()` or `provide()`.
58   59  
59   @see make_mocket_pair 60   @see make_mocket_pair
60   */ 61   */
61   template<class Socket = tcp_socket> 62   template<class Socket = tcp_socket>
62   class basic_mocket 63   class basic_mocket
63   { 64   {
64   Socket sock_; 65   Socket sock_;
65   std::string provide_; 66   std::string provide_;
66   std::string expect_; 67   std::string expect_;
67   capy::test::fuse fuse_; 68   capy::test::fuse fuse_;
68   std::size_t max_read_size_; 69   std::size_t max_read_size_;
69   std::size_t max_write_size_; 70   std::size_t max_write_size_;
70   71  
71   template<class MutableBufferSequence> 72   template<class MutableBufferSequence>
72   std::size_t consume_provide(MutableBufferSequence const& buffers) noexcept; 73   std::size_t consume_provide(MutableBufferSequence const& buffers) noexcept;
73   74  
74   template<class ConstBufferSequence> 75   template<class ConstBufferSequence>
75   bool validate_expect( 76   bool validate_expect(
76   ConstBufferSequence const& buffers, std::size_t& bytes_written); 77   ConstBufferSequence const& buffers, std::size_t& bytes_written);
77   78  
78   public: 79   public:
79   template<class MutableBufferSequence> 80   template<class MutableBufferSequence>
80   class read_some_awaitable; 81   class read_some_awaitable;
81   82  
82   template<class ConstBufferSequence> 83   template<class ConstBufferSequence>
83   class write_some_awaitable; 84   class write_some_awaitable;
84   85  
85   /** Destructor. 86   /** Destructor.
86   */ 87   */
HITCBC 87   36 ~basic_mocket() = default; 88   40 ~basic_mocket() = default;
88   89  
89   /** Construct a mocket. 90   /** Construct a mocket.
90   91  
91   @param ctx The execution context for the socket. 92   @param ctx The execution context for the socket.
92   @param f The fuse for error injection testing. 93   @param f The fuse for error injection testing.
93   @param max_read_size Maximum bytes per read operation. 94   @param max_read_size Maximum bytes per read operation.
94   @param max_write_size Maximum bytes per write operation. 95   @param max_write_size Maximum bytes per write operation.
95   */ 96   */
HITCBC 96   18 basic_mocket( 97   20 basic_mocket(
97   capy::execution_context& ctx, 98   capy::execution_context& ctx,
98   capy::test::fuse f = {}, 99   capy::test::fuse f = {},
99   std::size_t max_read_size = std::size_t(-1), 100   std::size_t max_read_size = std::size_t(-1),
100   std::size_t max_write_size = std::size_t(-1)) 101   std::size_t max_write_size = std::size_t(-1))
HITCBC 101   18 : sock_(ctx) 102   20 : sock_(ctx)
HITCBC 102   18 , fuse_(std::move(f)) 103   20 , fuse_(std::move(f))
HITCBC 103   18 , max_read_size_(max_read_size) 104   20 , max_read_size_(max_read_size)
HITCBC 104   18 , max_write_size_(max_write_size) 105   20 , max_write_size_(max_write_size)
105   { 106   {
HITCBC 106   18 if (max_read_size == 0) 107   20 if (max_read_size == 0)
MISUBC 107   ✗ detail::throw_logic_error("mocket: max_read_size cannot be 0"); 108   ✗ detail::throw_logic_error("mocket: max_read_size cannot be 0");
HITCBC 108   18 if (max_write_size == 0) 109   20 if (max_write_size == 0)
MISUBC 109   ✗ detail::throw_logic_error("mocket: max_write_size cannot be 0"); 110   ✗ detail::throw_logic_error("mocket: max_write_size cannot be 0");
HITCBC 110   18 } 111   20 }
111   112  
112   /** Move constructor. 113   /** Move constructor.
113   */ 114   */
HITCBC 114   18 basic_mocket(basic_mocket&& other) noexcept 115   20 basic_mocket(basic_mocket&& other) noexcept
HITCBC 115   18 : sock_(std::move(other.sock_)) 116   20 : sock_(std::move(other.sock_))
HITCBC 116   18 , provide_(std::move(other.provide_)) 117   20 , provide_(std::move(other.provide_))
HITCBC 117   18 , expect_(std::move(other.expect_)) 118   20 , expect_(std::move(other.expect_))
HITCBC 118   18 , fuse_(std::move(other.fuse_)) 119   20 , fuse_(std::move(other.fuse_))
HITCBC 119   18 , max_read_size_(other.max_read_size_) 120   20 , max_read_size_(other.max_read_size_)
HITCBC 120   18 , max_write_size_(other.max_write_size_) 121   20 , max_write_size_(other.max_write_size_)
121   { 122   {
HITCBC 122   18 } 123   20 }
123   124  
124   /** Move assignment. 125   /** Move assignment.
125   */ 126   */
126   basic_mocket& operator=(basic_mocket&& other) noexcept 127   basic_mocket& operator=(basic_mocket&& other) noexcept
127   { 128   {
128   if (this != &other) 129   if (this != &other)
129   { 130   {
130   sock_ = std::move(other.sock_); 131   sock_ = std::move(other.sock_);
131   provide_ = std::move(other.provide_); 132   provide_ = std::move(other.provide_);
132   expect_ = std::move(other.expect_); 133   expect_ = std::move(other.expect_);
133   fuse_ = other.fuse_; 134   fuse_ = other.fuse_;
134   max_read_size_ = other.max_read_size_; 135   max_read_size_ = other.max_read_size_;
135   max_write_size_ = other.max_write_size_; 136   max_write_size_ = other.max_write_size_;
136   } 137   }
137   return *this; 138   return *this;
138   } 139   }
139   140  
140   basic_mocket(basic_mocket const&) = delete; 141   basic_mocket(basic_mocket const&) = delete;
141   basic_mocket& operator=(basic_mocket const&) = delete; 142   basic_mocket& operator=(basic_mocket const&) = delete;
142   143  
143   /** Return the execution context. 144   /** Return the execution context.
144   145  
145   @return Reference to the execution context that owns this mocket. 146   @return Reference to the execution context that owns this mocket.
146   */ 147   */
147   capy::execution_context& context() const noexcept 148   capy::execution_context& context() const noexcept
148   { 149   {
149   return sock_.context(); 150   return sock_.context();
150   } 151   }
151   152  
152   /** Return the underlying socket. 153   /** Return the underlying socket.
153   154  
154   @return Reference to the underlying socket. 155   @return Reference to the underlying socket.
155   */ 156   */
HITCBC 156   20 Socket& socket() noexcept 157   22 Socket& socket() noexcept
157   { 158   {
HITCBC 158   20 return sock_; 159   22 return sock_;
159   } 160   }
160   161  
161   /** Stage data for reads. 162   /** Stage data for reads.
162   163  
163   Appends the given string to this mocket's provide buffer. 164   Appends the given string to this mocket's provide buffer.
164   When `read_some` is called, it will receive this data first 165   When `read_some` is called, it will receive this data first
165   before reading from the underlying socket. 166   before reading from the underlying socket.
166   167  
167   @param s The data to provide. 168   @param s The data to provide.
168   169  
169   @pre All coroutines using this mocket must be suspended. 170   @pre All coroutines using this mocket must be suspended.
170   */ 171   */
HITCBC 171   9 void provide(std::string const& s) 172   10 void provide(std::string const& s)
172   { 173   {
HITCBC 173   9 provide_.append(s); 174   10 provide_.append(s);
HITCBC 174   9 } 175   10 }
175   176  
176   /** Set expected data for writes. 177   /** Set expected data for writes.
177   178  
178   Appends the given string to this mocket's expect buffer. 179   Appends the given string to this mocket's expect buffer.
179   When the caller writes to this mocket, the written data 180   When the caller writes to this mocket, the written data
180   must match the expected data. On mismatch, `fuse::fail()` 181   must match the expected data. On mismatch, `fuse::fail()`
181   is called. 182   is called.
182   183  
183   @param s The expected data. 184   @param s The expected data.
184   185  
185   @pre All coroutines using this mocket must be suspended. 186   @pre All coroutines using this mocket must be suspended.
186   */ 187   */
HITCBC 187   8 void expect(std::string const& s) 188   10 void expect(std::string const& s)
188   { 189   {
HITCBC 189   8 expect_.append(s); 190   10 expect_.append(s);
HITCBC 190   8 } 191   10 }
191   192  
192   /** Check that every test expectation was consumed. 193   /** Check that every test expectation was consumed.
193   194  
194   Verifies that both the `expect()` and `provide()` buffers are 195   Verifies that both the `expect()` and `provide()` buffers are
195   empty. An unmet expectation also trips the fuse, so even a 196   empty. An unmet expectation also trips the fuse, so even a
196   discarded result still fails the test. 197   discarded result still fails the test.
197   198  
198   @return `error::test_failure` if either buffer holds 199   @return `error::test_failure` if either buffer holds
199   unconsumed data; empty otherwise. 200   unconsumed data; empty otherwise.
200   */ 201   */
HITCBC 201   36 [[nodiscard]] std::error_code verify() noexcept 202   40 [[nodiscard]] std::error_code verify() noexcept
202   { 203   {
HITCBC 203   36 if (expect_.empty() && provide_.empty()) 204   40 if (expect_.empty() && provide_.empty())
HITCBC 204   28 return {}; 205   30 return {};
HITCBC 205   8 fuse_.fail(); 206   10 fuse_.fail();
HITCBC 206   8 return capy::error::test_failure; 207   10 return capy::error::test_failure;
207   } 208   }
208   209  
209   /** Close the mocket. 210   /** Close the mocket.
210   211  
211   Idempotent, like every `close()` in the library. Unconsumed 212   Idempotent, like every `close()` in the library. Unconsumed
212   `expect()`/`provide()` data trips the fuse on the way out; use 213   `expect()`/`provide()` data trips the fuse on the way out; use
213   @ref verify to inspect the outcome as a code. 214   @ref verify to inspect the outcome as a code.
214   */ 215   */
HITCBC 215   18 void close() noexcept 216   20 void close() noexcept
216   { 217   {
HITCBC 217   18 if (!sock_.is_open()) 218   20 if (!sock_.is_open())
MISUBC 218   ✗ return; 219   ✗ return;
219   220  
220   // Discarded on purpose: the fuse reports unmet expectations. 221   // Discarded on purpose: the fuse reports unmet expectations.
HITCBC 221   18 std::ignore = verify(); 222   20 std::ignore = verify();
HITCBC 222   18 sock_.close(); 223   20 sock_.close();
223   } 224   }
224   225  
225   /** Cancel pending I/O operations. 226   /** Cancel pending I/O operations.
226   227  
227   Cancels any pending asynchronous operations on the underlying 228   Cancels any pending asynchronous operations on the underlying
228   socket. Outstanding operations complete with `cond::canceled`. 229   socket. Outstanding operations complete with `cond::canceled`.
229   */ 230   */
230   void cancel() noexcept 231   void cancel() noexcept
231   { 232   {
232   sock_.cancel(); 233   sock_.cancel();
233   } 234   }
234   235  
235   /** Check if the mocket is open. 236   /** Check if the mocket is open.
236   237  
237   @return `true` if the mocket is open. 238   @return `true` if the mocket is open.
238   */ 239   */
HITCBC 239   5 bool is_open() const noexcept 240   5 bool is_open() const noexcept
240   { 241   {
HITCBC 241   5 return sock_.is_open(); 242   5 return sock_.is_open();
242   } 243   }
243   244  
244   /** Initiate an asynchronous read operation. 245   /** Initiate an asynchronous read operation.
245   246  
246   Reads available data into the provided buffer sequence. If the 247   Reads available data into the provided buffer sequence. If the
247   provide buffer has data, it is consumed first. Otherwise, the 248   provide buffer has data, it is consumed first. Otherwise, the
248   operation delegates to the underlying socket. 249   operation delegates to the underlying socket.
249   250  
250   @param buffers The buffer sequence to read data into. 251   @param buffers The buffer sequence to read data into.
251   252  
252   @return An awaitable yielding `(error_code, std::size_t)`. 253   @return An awaitable yielding `(error_code, std::size_t)`.
253   */ 254   */
254   template<class MutableBufferSequence> 255   template<class MutableBufferSequence>
HITCBC 255   11 [[nodiscard]] auto read_some(MutableBufferSequence const& buffers) 256   12 [[nodiscard]] auto read_some(MutableBufferSequence const& buffers)
256   { 257   {
HITCBC 257   11 return read_some_awaitable<MutableBufferSequence>(*this, buffers); 258   12 return read_some_awaitable<MutableBufferSequence>(*this, buffers);
258   } 259   }
259   260  
260   /** Initiate an asynchronous write operation. 261   /** Initiate an asynchronous write operation.
261   262  
262   Writes data from the provided buffer sequence. If the expect 263   Writes data from the provided buffer sequence. If the expect
263   buffer has data, it is validated. Otherwise, the operation 264   buffer has data, it is validated. Otherwise, the operation
264   delegates to the underlying socket. 265   delegates to the underlying socket.
265   266  
266   @param buffers The buffer sequence containing data to write. 267   @param buffers The buffer sequence containing data to write.
267   268  
268   @return An awaitable yielding `(error_code, std::size_t)`. 269   @return An awaitable yielding `(error_code, std::size_t)`.
269   */ 270   */
270   template<class ConstBufferSequence> 271   template<class ConstBufferSequence>
HITCBC 271   8 [[nodiscard]] auto write_some(ConstBufferSequence const& buffers) 272   10 [[nodiscard]] auto write_some(ConstBufferSequence const& buffers)
272   { 273   {
HITCBC 273   8 return write_some_awaitable<ConstBufferSequence>(*this, buffers); 274   10 return write_some_awaitable<ConstBufferSequence>(*this, buffers);
274   } 275   }
275   }; 276   };
276   277  
277   /// Default mocket type using `tcp_socket`. 278   /// Default mocket type using `tcp_socket`.
278   using mocket = basic_mocket<>; 279   using mocket = basic_mocket<>;
279   280  
280   template<class Socket> 281   template<class Socket>
281   template<class MutableBufferSequence> 282   template<class MutableBufferSequence>
282   std::size_t 283   std::size_t
HITCBC 283   10 basic_mocket<Socket>::consume_provide( 284   10 basic_mocket<Socket>::consume_provide(
284   MutableBufferSequence const& buffers) noexcept 285   MutableBufferSequence const& buffers) noexcept
285   { 286   {
286   auto n = 287   auto n =
HITCBC 287   10 capy::buffer_copy(buffers, capy::make_buffer(provide_), max_read_size_); 288   10 capy::buffer_copy(buffers, capy::make_buffer(provide_), max_read_size_);
HITCBC 288   10 provide_.erase(0, n); 289   10 provide_.erase(0, n);
HITCBC 289   10 return n; 290   10 return n;
290   } 291   }
291   292  
292   template<class Socket> 293   template<class Socket>
293   template<class ConstBufferSequence> 294   template<class ConstBufferSequence>
294   bool 295   bool
HITCBC 295   7 basic_mocket<Socket>::validate_expect( 296   8 basic_mocket<Socket>::validate_expect(
296   ConstBufferSequence const& buffers, std::size_t& bytes_written) 297   ConstBufferSequence const& buffers, std::size_t& bytes_written)
297   { 298   {
HITCBC 298   7 if (expect_.empty()) 299   8 if (expect_.empty())
MISUBC 299   ✗ return true; 300   ✗ return true;
300   301  
301   // Build the write data up to max_write_size_ 302   // Build the write data up to max_write_size_
HITCBC 302   7 std::string written; 303   8 std::string written;
HITCBC 303   7 auto total = capy::buffer_size(buffers); 304   8 auto total = capy::buffer_size(buffers);
HITCBC 304   7 if (total > max_write_size_) 305   8 if (total > max_write_size_)
HITCBC 305   1 total = max_write_size_; 306   1 total = max_write_size_;
HITCBC 306   7 written.resize(total); 307   8 written.resize(total);
HITCBC 307   7 capy::buffer_copy(capy::make_buffer(written), buffers, max_write_size_); 308   8 capy::buffer_copy(capy::make_buffer(written), buffers, max_write_size_);
308   309  
309   // Check if written data matches expect prefix 310   // Check if written data matches expect prefix
HITCBC 310   7 auto const match_size = (std::min)(written.size(), expect_.size()); 311   8 auto const match_size = (std::min)(written.size(), expect_.size());
HITCBC 311   7 if (std::memcmp(written.data(), expect_.data(), match_size) != 0) 312   8 if (std::memcmp(written.data(), expect_.data(), match_size) != 0)
312   { 313   {
MISUBC 313   ✗ fuse_.fail(); 314   ✗ fuse_.fail();
MISUBC 314   ✗ bytes_written = 0; 315   ✗ bytes_written = 0;
MISUBC 315   ✗ return false; 316   ✗ return false;
316   } 317   }
317   318  
318 - // Consume matched portion 319 + // Only the validated prefix counts as written — a longer request
  320 + // is a partial write, per WriteStream.
HITCBC 319   7 expect_.erase(0, match_size); 321   8 expect_.erase(0, match_size);
HITCBC 320 - 7 bytes_written = written.size(); 322 + 8 bytes_written = match_size;
HITCBC 321   7 return true; 323   8 return true;
HITCBC 322   7 } 324   8 }
323   325  
324   template<class Socket> 326   template<class Socket>
325   template<class MutableBufferSequence> 327   template<class MutableBufferSequence>
326   class basic_mocket<Socket>::read_some_awaitable 328   class basic_mocket<Socket>::read_some_awaitable
327   { 329   {
328   using sock_awaitable = decltype(std::declval<Socket&>().read_some( 330   using sock_awaitable = decltype(std::declval<Socket&>().read_some(
329   std::declval<MutableBufferSequence>())); 331   std::declval<MutableBufferSequence>()));
330   332  
331   basic_mocket* m_; 333   basic_mocket* m_;
332   MutableBufferSequence buffers_; 334   MutableBufferSequence buffers_;
333   std::size_t n_ = 0; 335   std::size_t n_ = 0;
334   std::error_code ec_; 336   std::error_code ec_;
335   union 337   union
336   { 338   {
337   char dummy_; 339   char dummy_;
338   sock_awaitable underlying_; 340   sock_awaitable underlying_;
339   }; 341   };
340   bool sync_ = true; 342   bool sync_ = true;
341   343  
342   public: 344   public:
HITCBC 343   11 read_some_awaitable(basic_mocket& m, MutableBufferSequence buffers) noexcept 345   12 read_some_awaitable(basic_mocket& m, MutableBufferSequence buffers) noexcept
HITCBC 344   11 : m_(&m) 346   12 : m_(&m)
HITCBC 345   11 , buffers_(std::move(buffers)) 347   12 , buffers_(std::move(buffers))
346   { 348   {
HITCBC 347   11 } 349   12 }
348   350  
HITCBC 349   22 ~read_some_awaitable() 351   24 ~read_some_awaitable()
350   { 352   {
HITCBC 351   22 if (!sync_) 353   24 if (!sync_)
HITCBC 352   1 underlying_.~sock_awaitable(); 354   1 underlying_.~sock_awaitable();
HITCBC 353   22 } 355   24 }
354   356  
HITCBC 355   11 read_some_awaitable(read_some_awaitable&& other) noexcept 357   12 read_some_awaitable(read_some_awaitable&& other) noexcept
HITCBC 356   11 : m_(other.m_) 358   12 : m_(other.m_)
HITCBC 357   11 , buffers_(std::move(other.buffers_)) 359   12 , buffers_(std::move(other.buffers_))
HITCBC 358   11 , n_(other.n_) 360   12 , n_(other.n_)
HITCBC 359   11 , ec_(other.ec_) 361   12 , ec_(other.ec_)
HITCBC 360   11 , sync_(other.sync_) 362   12 , sync_(other.sync_)
361   { 363   {
HITCBC 362   11 if (!sync_) 364   12 if (!sync_)
363   { 365   {
MISUBC 364   ✗ new (&underlying_) sock_awaitable(std::move(other.underlying_)); 366   ✗ new (&underlying_) sock_awaitable(std::move(other.underlying_));
MISUBC 365   ✗ other.underlying_.~sock_awaitable(); 367   ✗ other.underlying_.~sock_awaitable();
MISUBC 366   ✗ other.sync_ = true; 368   ✗ other.sync_ = true;
367   } 369   }
HITCBC 368   11 } 370   12 }
369   371  
370   read_some_awaitable(read_some_awaitable const&) = delete; 372   read_some_awaitable(read_some_awaitable const&) = delete;
371   read_some_awaitable& operator=(read_some_awaitable const&) = delete; 373   read_some_awaitable& operator=(read_some_awaitable const&) = delete;
372   read_some_awaitable& operator=(read_some_awaitable&&) = delete; 374   read_some_awaitable& operator=(read_some_awaitable&&) = delete;
373   375  
ECB 374 - 11 bool await_ready() 376 + // All decisions wait for await_suspend, where the io_env (and thus
  377 + // the stop token) is available — a pre-stopped token must
  378 + // short-circuit before any staged data is consumed.
HITGNC   379 + 12 bool await_ready() const noexcept
  380 + {
HITGNC   381 + 12 return false;
  382 + }
  383 +
HITGNC   384 + 12 auto await_suspend(std::coroutine_handle<> h, capy::io_env const* env)
  385 + -> std::coroutine_handle<>
375   { 386   {
HITGNC   387 + 12 if (env->stop_token.stop_requested())
  388 + {
HITGNC   389 + 1 ec_ = capy::error::canceled;
HITGNC   390 + 1 n_ = 0;
HITGNC   391 + 1 return h;
  392 + }
376   // Fuse injection point: an armed fuse fails this read as if the 393   // Fuse injection point: an armed fuse fails this read as if the
377   // transport did, so a fault-injection sweep exercises the error 394   // transport did, so a fault-injection sweep exercises the error
378   // path of every read the caller issues. Inert outside armed(). 395   // path of every read the caller issues. Inert outside armed().
379   // A transport reports failure through the result, never by 396   // A transport reports failure through the result, never by
380   // throwing from read_some, so the fuse's exception phase is 397   // throwing from read_some, so the fuse's exception phase is
381   // converted to the same error code its error-code phase yields. 398   // converted to the same error code its error-code phase yields.
HITCBC 382   11 std::error_code fec; 399   11 std::error_code fec;
383   try 400   try
384   { 401   {
HITCBC 385   11 fec = m_->fuse_.maybe_fail(); 402   11 fec = m_->fuse_.maybe_fail();
386   } 403   }
MISUBC 387   ✗ catch (std::system_error const& e) 404   ✗ catch (std::system_error const& e)
388   { 405   {
MISUBC 389   ✗ fec = e.code(); 406   ✗ fec = e.code();
390   } 407   }
HITCBC 391   11 if (fec) 408   11 if (fec)
392   { 409   {
MISUBC 393   ✗ ec_ = fec; 410   ✗ ec_ = fec;
MISUBC 394   ✗ n_ = 0; 411   ✗ n_ = 0;
MISUBC 395 - ✗ return true; 412 + ✗ return h;
396   } 413   }
HITCBC 397   11 if (!m_->provide_.empty()) 414   11 if (!m_->provide_.empty())
398   { 415   {
HITCBC 399   10 n_ = m_->consume_provide(buffers_); 416   10 n_ = m_->consume_provide(buffers_);
HITCBC 400 - 10 return true; 417 + 10 return h;
401   } 418   }
HITCBC 402   1 new (&underlying_) sock_awaitable(m_->sock_.read_some(buffers_)); 419   1 new (&underlying_) sock_awaitable(m_->sock_.read_some(buffers_));
HITCBC 403   1 sync_ = false; 420   1 sync_ = false;
HITCBC 404 - 1 return underlying_.await_ready(); 421 + 1 if (underlying_.await_ready())
MISUIC 405 - } 422 + ✗ return h;
HITGIC 406 - 423 + 1 return underlying_.await_suspend(h, env);
407 - template<class... Args>  
DCB 408 - 1 auto await_suspend(Args&&... args)  
409 - {  
DCB 410 - 1 return underlying_.await_suspend(std::forward<Args>(args)...);  
411   } 424   }
412   425  
HITCBC 413   11 [[nodiscard]] capy::io_result<std::size_t> await_resume() 426   12 [[nodiscard]] capy::io_result<std::size_t> await_resume()
414   { 427   {
HITCBC 415   11 if (sync_) 428   12 if (sync_)
HITCBC 416   10 return {ec_, n_}; 429   11 return {ec_, n_};
HITCBC 417   1 return underlying_.await_resume(); 430   1 return underlying_.await_resume();
418   } 431   }
419   }; 432   };
420   433  
421   template<class Socket> 434   template<class Socket>
422   template<class ConstBufferSequence> 435   template<class ConstBufferSequence>
423   class basic_mocket<Socket>::write_some_awaitable 436   class basic_mocket<Socket>::write_some_awaitable
424   { 437   {
425   using sock_awaitable = decltype(std::declval<Socket&>().write_some( 438   using sock_awaitable = decltype(std::declval<Socket&>().write_some(
426   std::declval<ConstBufferSequence>())); 439   std::declval<ConstBufferSequence>()));
427   440  
428   basic_mocket* m_; 441   basic_mocket* m_;
429   ConstBufferSequence buffers_; 442   ConstBufferSequence buffers_;
430   std::size_t n_ = 0; 443   std::size_t n_ = 0;
431   std::error_code ec_; 444   std::error_code ec_;
432   union 445   union
433   { 446   {
434   char dummy_; 447   char dummy_;
435   sock_awaitable underlying_; 448   sock_awaitable underlying_;
436   }; 449   };
437   bool sync_ = true; 450   bool sync_ = true;
438   451  
439   public: 452   public:
HITCBC 440   8 write_some_awaitable(basic_mocket& m, ConstBufferSequence buffers) noexcept 453   10 write_some_awaitable(basic_mocket& m, ConstBufferSequence buffers) noexcept
HITCBC 441   8 : m_(&m) 454   10 : m_(&m)
HITCBC 442   8 , buffers_(std::move(buffers)) 455   10 , buffers_(std::move(buffers))
443   { 456   {
HITCBC 444   8 } 457   10 }
445   458  
HITCBC 446   16 ~write_some_awaitable() 459   20 ~write_some_awaitable()
447   { 460   {
HITCBC 448   16 if (!sync_) 461   20 if (!sync_)
HITCBC 449   1 underlying_.~sock_awaitable(); 462   1 underlying_.~sock_awaitable();
HITCBC 450   16 } 463   20 }
451   464  
HITCBC 452   8 write_some_awaitable(write_some_awaitable&& other) noexcept 465   10 write_some_awaitable(write_some_awaitable&& other) noexcept
HITCBC 453   8 : m_(other.m_) 466   10 : m_(other.m_)
HITCBC 454   8 , buffers_(std::move(other.buffers_)) 467   10 , buffers_(std::move(other.buffers_))
HITCBC 455   8 , n_(other.n_) 468   10 , n_(other.n_)
HITCBC 456   8 , ec_(other.ec_) 469   10 , ec_(other.ec_)
HITCBC 457   8 , sync_(other.sync_) 470   10 , sync_(other.sync_)
458   { 471   {
HITCBC 459   8 if (!sync_) 472   10 if (!sync_)
460   { 473   {
MISUBC 461   ✗ new (&underlying_) sock_awaitable(std::move(other.underlying_)); 474   ✗ new (&underlying_) sock_awaitable(std::move(other.underlying_));
MISUBC 462   ✗ other.underlying_.~sock_awaitable(); 475   ✗ other.underlying_.~sock_awaitable();
MISUBC 463   ✗ other.sync_ = true; 476   ✗ other.sync_ = true;
464   } 477   }
HITCBC 465   8 } 478   10 }
466   479  
467   write_some_awaitable(write_some_awaitable const&) = delete; 480   write_some_awaitable(write_some_awaitable const&) = delete;
468   write_some_awaitable& operator=(write_some_awaitable const&) = delete; 481   write_some_awaitable& operator=(write_some_awaitable const&) = delete;
469   write_some_awaitable& operator=(write_some_awaitable&&) = delete; 482   write_some_awaitable& operator=(write_some_awaitable&&) = delete;
470   483  
ECB 471 - 8 bool await_ready() 484 + // All decisions wait for await_suspend, where the io_env (and thus
  485 + // the stop token) is available — a pre-stopped token must
  486 + // short-circuit before any of the expect script is consumed.
HITGNC   487 + 10 bool await_ready() const noexcept
472   { 488   {
HITGNC   489 + 10 return false;
  490 + }
  491 +
HITGNC   492 + 10 auto await_suspend(std::coroutine_handle<> h, capy::io_env const* env)
  493 + -> std::coroutine_handle<>
  494 + {
HITGNC   495 + 10 if (env->stop_token.stop_requested())
  496 + {
HITGNC   497 + 1 ec_ = capy::error::canceled;
HITGNC   498 + 1 n_ = 0;
HITGNC   499 + 1 return h;
  500 + }
473   // Fuse injection point: an armed fuse fails this write as if the 501   // Fuse injection point: an armed fuse fails this write as if the
474   // transport did, so a fault-injection sweep exercises the error 502   // transport did, so a fault-injection sweep exercises the error
475   // path of every write the caller issues. Inert outside armed(). 503   // path of every write the caller issues. Inert outside armed().
476   // A transport reports failure through the result, never by 504   // A transport reports failure through the result, never by
477   // throwing from write_some, so the fuse's exception phase is 505   // throwing from write_some, so the fuse's exception phase is
478   // converted to the same error code its error-code phase yields. 506   // converted to the same error code its error-code phase yields.
HITCBC 479   8 std::error_code fec; 507   9 std::error_code fec;
480   try 508   try
481   { 509   {
HITCBC 482   8 fec = m_->fuse_.maybe_fail(); 510   9 fec = m_->fuse_.maybe_fail();
483   } 511   }
MISUBC 484   ✗ catch (std::system_error const& e) 512   ✗ catch (std::system_error const& e)
485   { 513   {
MISUBC 486   ✗ fec = e.code(); 514   ✗ fec = e.code();
487   } 515   }
HITCBC 488   8 if (fec) 516   9 if (fec)
489   { 517   {
MISUBC 490   ✗ ec_ = fec; 518   ✗ ec_ = fec;
MISUBC 491   ✗ n_ = 0; 519   ✗ n_ = 0;
MISUBC 492 - ✗ return true; 520 + ✗ return h;
493   } 521   }
HITCBC 494   8 if (!m_->expect_.empty()) 522   9 if (!m_->expect_.empty())
495   { 523   {
HITCBC 496   7 if (!m_->validate_expect(buffers_, n_)) 524   8 if (!m_->validate_expect(buffers_, n_))
497   { 525   {
MISUBC 498   ✗ ec_ = capy::error::test_failure; 526   ✗ ec_ = capy::error::test_failure;
MISUBC 499   ✗ n_ = 0; 527   ✗ n_ = 0;
500   } 528   }
HITCBC 501 - 7 return true; 529 + 8 return h;
502   } 530   }
HITCBC 503   1 new (&underlying_) sock_awaitable(m_->sock_.write_some(buffers_)); 531   1 new (&underlying_) sock_awaitable(m_->sock_.write_some(buffers_));
HITCBC 504   1 sync_ = false; 532   1 sync_ = false;
HITCBC 505 - 1 return underlying_.await_ready(); 533 + 1 if (underlying_.await_ready())
MISUIC 506 - } 534 + ✗ return h;
HITGIC 507 - 535 + 1 return underlying_.await_suspend(h, env);
508 - template<class... Args>  
DCB 509 - 1 auto await_suspend(Args&&... args)  
510 - {  
DCB 511 - 1 return underlying_.await_suspend(std::forward<Args>(args)...);  
512   } 536   }
513   537  
HITCBC 514   8 [[nodiscard]] capy::io_result<std::size_t> await_resume() 538   10 [[nodiscard]] capy::io_result<std::size_t> await_resume()
515   { 539   {
HITCBC 516   8 if (sync_) 540   10 if (sync_)
HITCBC 517   7 return {ec_, n_}; 541   9 return {ec_, n_};
HITCBC 518   1 return underlying_.await_resume(); 542   1 return underlying_.await_resume();
519   } 543   }
520   }; 544   };
521   545  
522   /** Create a mocket paired with a socket. 546   /** Create a mocket paired with a socket.
523   547  
524   Creates a mocket and a socket connected via loopback. 548   Creates a mocket and a socket connected via loopback.
525   Data written to one can be read from the other. 549   Data written to one can be read from the other.
526   550  
527   The mocket has fuse checks enabled via `maybe_fail()` and 551   The mocket has fuse checks enabled via `maybe_fail()` and
528   supports provide/expect buffers for test instrumentation. 552   supports provide/expect buffers for test instrumentation.
529   The socket is the "peer" end with no test instrumentation. 553   The socket is the "peer" end with no test instrumentation.
530   554  
531   Optional max_read_size and max_write_size parameters limit the 555   Optional max_read_size and max_write_size parameters limit the
532   number of bytes transferred per I/O operation on the mocket, 556   number of bytes transferred per I/O operation on the mocket,
533   simulating chunked network delivery for testing purposes. 557   simulating chunked network delivery for testing purposes.
534   558  
535   @tparam Socket The socket type (default `tcp_socket`). 559   @tparam Socket The socket type (default `tcp_socket`).
536   @tparam Acceptor The acceptor type (default `tcp_acceptor`). 560   @tparam Acceptor The acceptor type (default `tcp_acceptor`).
537   561  
538   @param ctx The I/O context for the sockets. 562   @param ctx The I/O context for the sockets.
539   @param f The fuse for error injection testing. 563   @param f The fuse for error injection testing.
540   @param max_read_size Maximum bytes per read operation (default unlimited). 564   @param max_read_size Maximum bytes per read operation (default unlimited).
541   @param max_write_size Maximum bytes per write operation (default unlimited). 565   @param max_write_size Maximum bytes per write operation (default unlimited).
542   566  
543   @return A pair of (mocket, socket). 567   @return A pair of (mocket, socket).
544   568  
545   @note Mockets are not thread-safe and must be used in a 569   @note Mockets are not thread-safe and must be used in a
546   single-threaded, deterministic context. 570   single-threaded, deterministic context.
547   */ 571   */
548   template<class Socket = tcp_socket, class Acceptor = tcp_acceptor> 572   template<class Socket = tcp_socket, class Acceptor = tcp_acceptor>
549   std::pair<basic_mocket<Socket>, Socket> 573   std::pair<basic_mocket<Socket>, Socket>
HITCBC 550   18 make_mocket_pair( 574   20 make_mocket_pair(
551   io_context& ctx, 575   io_context& ctx,
552   capy::test::fuse f = {}, 576   capy::test::fuse f = {},
553   std::size_t max_read_size = std::size_t(-1), 577   std::size_t max_read_size = std::size_t(-1),
554   std::size_t max_write_size = std::size_t(-1)) 578   std::size_t max_write_size = std::size_t(-1))
555   { 579   {
HITCBC 556   18 auto ex = ctx.get_executor(); 580   20 auto ex = ctx.get_executor();
557   581  
HITCBC 558   18 basic_mocket<Socket> m(ctx, std::move(f), max_read_size, max_write_size); 582   20 basic_mocket<Socket> m(ctx, std::move(f), max_read_size, max_write_size);
559   583  
HITCBC 560   18 Socket peer(ctx); 584   20 Socket peer(ctx);
561   585  
HITCBC 562   18 std::error_code accept_ec; 586   20 std::error_code accept_ec;
HITCBC 563   18 std::error_code connect_ec; 587   20 std::error_code connect_ec;
HITCBC 564   18 bool accept_done = false; 588   20 bool accept_done = false;
HITCBC 565   18 bool connect_done = false; 589   20 bool connect_done = false;
566   590  
HITCBC 567   18 Acceptor acc(ctx); 591   20 Acceptor acc(ctx);
HITCBC 568   18 if (auto open_ec = acc.open()) 592   20 if (auto open_ec = acc.open())
MISUBC 569   ✗ throw std::runtime_error("mocket open failed: " + open_ec.message()); 593   ✗ throw std::runtime_error("mocket open failed: " + open_ec.message());
HITCBC 570   18 acc.set_option(socket_option::reuse_address(true)); 594   20 acc.set_option(socket_option::reuse_address(true));
HITCBC 571   18 if (auto bind_ec = acc.bind(endpoint(ipv4_address::loopback(), 0))) 595   20 if (auto bind_ec = acc.bind(endpoint(ipv4_address::loopback(), 0)))
MISUBC 572   ✗ throw std::runtime_error("mocket bind failed: " + bind_ec.message()); 596   ✗ throw std::runtime_error("mocket bind failed: " + bind_ec.message());
HITCBC 573   18 if (auto listen_ec = acc.listen()) 597   20 if (auto listen_ec = acc.listen())
MISUBC 574   ✗ throw std::runtime_error( 598   ✗ throw std::runtime_error(
575   "mocket listen failed: " + listen_ec.message()); 599   "mocket listen failed: " + listen_ec.message());
HITCBC 576   18 auto port = acc.local_endpoint().port(); 600   20 auto port = acc.local_endpoint().port();
577   601  
HITCBC 578   18 if (auto open_ec = peer.open()) 602   20 if (auto open_ec = peer.open())
MISUBC 579   ✗ throw std::runtime_error("mocket open failed: " + open_ec.message()); 603   ✗ throw std::runtime_error("mocket open failed: " + open_ec.message());
580   604  
HITCBC 581   18 Socket accepted_socket(ctx); 605   20 Socket accepted_socket(ctx);
582   606  
HITCBC 583   18 capy::run_async(ex)( 607   20 capy::run_async(ex)(
HITCBC 584   36 [](Acceptor& a, Socket& s, std::error_code& ec_out, 608   40 [](Acceptor& a, Socket& s, std::error_code& ec_out,
585   bool& done_out) -> capy::task<> { 609   bool& done_out) -> capy::task<> {
586   auto [ec] = co_await a.accept(s); 610   auto [ec] = co_await a.accept(s);
587   ec_out = ec; 611   ec_out = ec;
588   done_out = true; 612   done_out = true;
589   }(acc, accepted_socket, accept_ec, accept_done)); 613   }(acc, accepted_socket, accept_ec, accept_done));
590   614  
HITCBC 591   18 capy::run_async(ex)( 615   20 capy::run_async(ex)(
HITCBC 592   36 [](Socket& s, endpoint ep, std::error_code& ec_out, 616   40 [](Socket& s, endpoint ep, std::error_code& ec_out,
593   bool& done_out) -> capy::task<> { 617   bool& done_out) -> capy::task<> {
594   auto [ec] = co_await s.connect(ep); 618   auto [ec] = co_await s.connect(ep);
595   ec_out = ec; 619   ec_out = ec;
596   done_out = true; 620   done_out = true;
597   }(peer, endpoint(ipv4_address::loopback(), port), connect_ec, 621   }(peer, endpoint(ipv4_address::loopback(), port), connect_ec,
598   connect_done)); 622   connect_done));
599   623  
HITCBC 600   18 ctx.run(); 624   20 ctx.run();
HITCBC 601   18 ctx.restart(); 625   20 ctx.restart();
602   626  
HITCBC 603   18 if (!accept_done || accept_ec) 627   20 if (!accept_done || accept_ec)
604   { 628   {
MISUBC 605   ✗ std::fprintf( 629   ✗ std::fprintf(
606   stderr, "make_mocket_pair: accept failed (done=%d, ec=%s)\n", 630   stderr, "make_mocket_pair: accept failed (done=%d, ec=%s)\n",
607   accept_done, accept_ec.message().c_str()); 631   accept_done, accept_ec.message().c_str());
MISUBC 608   ✗ acc.close(); 632   ✗ acc.close();
MISUBC 609   ✗ throw std::runtime_error("mocket accept failed"); 633   ✗ throw std::runtime_error("mocket accept failed");
610   } 634   }
611   635  
HITCBC 612   18 if (!connect_done || connect_ec) 636   20 if (!connect_done || connect_ec)
613   { 637   {
MISUBC 614   ✗ std::fprintf( 638   ✗ std::fprintf(
615   stderr, "make_mocket_pair: connect failed (done=%d, ec=%s)\n", 639   stderr, "make_mocket_pair: connect failed (done=%d, ec=%s)\n",
616   connect_done, connect_ec.message().c_str()); 640   connect_done, connect_ec.message().c_str());
MISUBC 617   ✗ acc.close(); 641   ✗ acc.close();
MISUBC 618   ✗ accepted_socket.close(); 642   ✗ accepted_socket.close();
MISUBC 619   ✗ throw std::runtime_error("mocket connect failed"); 643   ✗ throw std::runtime_error("mocket connect failed");
620   } 644   }
621   645  
HITCBC 622   18 m.socket() = std::move(accepted_socket); 646   20 m.socket() = std::move(accepted_socket);
623   647  
HITCBC 624   18 acc.close(); 648   20 acc.close();
625   649  
HITCBC 626   36 return {std::move(m), std::move(peer)}; 650   40 return {std::move(m), std::move(peer)};
HITCBC 627   18 } 651   20 }
628   652  
629   } // namespace boost::corosio::test 653   } // namespace boost::corosio::test
630   654  
631   #endif 655   #endif