99.38% Lines (159/160) 100.00% Functions (26/26)
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 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/capy 8   // Official repository: https://github.com/cppalliance/capy
9   // 9   //
10   10  
11   #ifndef BOOST_CAPY_TEST_STREAM_HPP 11   #ifndef BOOST_CAPY_TEST_STREAM_HPP
12   #define BOOST_CAPY_TEST_STREAM_HPP 12   #define BOOST_CAPY_TEST_STREAM_HPP
13   13  
14   #include <boost/capy/detail/config.hpp> 14   #include <boost/capy/detail/config.hpp>
15   #include <boost/capy/buffers.hpp> 15   #include <boost/capy/buffers.hpp>
16   #include <boost/capy/buffers/buffer_copy.hpp> 16   #include <boost/capy/buffers/buffer_copy.hpp>
17   #include <boost/capy/buffers/make_buffer.hpp> 17   #include <boost/capy/buffers/make_buffer.hpp>
18   #include <boost/capy/continuation.hpp> 18   #include <boost/capy/continuation.hpp>
19   #include <coroutine> 19   #include <coroutine>
20   #include <boost/capy/ex/io_env.hpp> 20   #include <boost/capy/ex/io_env.hpp>
21   #include <boost/capy/io_result.hpp> 21   #include <boost/capy/io_result.hpp>
22   #include <boost/capy/error.hpp> 22   #include <boost/capy/error.hpp>
23   #include <boost/capy/read.hpp> 23   #include <boost/capy/read.hpp>
24   #include <boost/capy/task.hpp> 24   #include <boost/capy/task.hpp>
25   #include <boost/capy/test/fuse.hpp> 25   #include <boost/capy/test/fuse.hpp>
26   #include <boost/capy/test/run_blocking.hpp> 26   #include <boost/capy/test/run_blocking.hpp>
27   27  
28   #include <atomic> 28   #include <atomic>
29   #include <memory> 29   #include <memory>
30   #include <new> 30   #include <new>
31   #include <stop_token> 31   #include <stop_token>
32   #include <string> 32   #include <string>
33   #include <string_view> 33   #include <string_view>
34   #include <utility> 34   #include <utility>
35   35  
36   namespace boost { 36   namespace boost {
37   namespace capy { 37   namespace capy {
38   namespace test { 38   namespace test {
39   39  
40   /** Suspends a reader until its paired end writes, or the shared fuse injects an error. 40   /** Suspends a reader until its paired end writes, or the shared fuse injects an error.
41   41  
42   Streams are created in pairs via @ref make_stream_pair. 42   Streams are created in pairs via @ref make_stream_pair.
43   Data written to one end becomes available for reading on 43   Data written to one end becomes available for reading on
44   the other. If no data is available when @ref read_some 44   the other. If no data is available when @ref read_some
45   is called, the calling coroutine suspends until the peer 45   is called, the calling coroutine suspends until the peer
46   calls @ref write_some. The shared @ref fuse enables error 46   calls @ref write_some. The shared @ref fuse enables error
47   injection at controlled points in both directions. 47   injection at controlled points in both directions.
48   48  
49   When the fuse injects an error or throws on one end, the 49   When the fuse injects an error or throws on one end, the
50   pair is automatically closed. Any suspended reader on 50   pair is automatically closed. Any suspended reader on
51   either end is resumed with `error::eof`, and subsequent 51   either end is resumed with `error::eof`, and subsequent
52   operations on both ends return `error::eof`. Calling 52   operations on both ends return `error::eof`. Calling
53   @ref close on one end signals eof to the peer's reads 53   @ref close on one end signals eof to the peer's reads
54   after draining any buffered data, while the peer may 54   after draining any buffered data, while the peer may
55   still write. 55   still write.
56   56  
57   @par Thread Safety 57   @par Thread Safety
58   Single-threaded only. Both ends of the pair must be 58   Single-threaded only. Both ends of the pair must be
59   accessed from the same thread. Concurrent access is 59   accessed from the same thread. Concurrent access is
60   undefined behavior. 60   undefined behavior.
61   61  
62   @par Example 62   @par Example
63   @code 63   @code
64   fuse f; 64   fuse f;
65   65  
66   auto r = f.armed( [&]( fuse& ) -> task<> { 66   auto r = f.armed( [&]( fuse& ) -> task<> {
67   // Constructed inside the lambda: armed() re-invokes this 67   // Constructed inside the lambda: armed() re-invokes this
68   // function once per injected failure point, and a stream 68   // function once per injected failure point, and a stream
69   // pair constructed outside would carry buffered state 69   // pair constructed outside would carry buffered state
70   // across those rounds. 70   // across those rounds.
71   auto [a, b] = make_stream_pair( f ); 71   auto [a, b] = make_stream_pair( f );
72   72  
73   auto [ec, n] = co_await a.write_some( 73   auto [ec, n] = co_await a.write_some(
74   const_buffer( "hello", 5 ) ); 74   const_buffer( "hello", 5 ) );
75   if( ec ) 75   if( ec )
76   co_return; 76   co_return;
77   77  
78   char buf[32]; 78   char buf[32];
79   auto [ec2, n2] = co_await b.read_some( 79   auto [ec2, n2] = co_await b.read_some(
80   mutable_buffer( buf, sizeof( buf ) ) ); 80   mutable_buffer( buf, sizeof( buf ) ) );
81   if( ec2 ) 81   if( ec2 )
82   co_return; 82   co_return;
83   // buf contains "hello" 83   // buf contains "hello"
84   } ); 84   } );
85   @endcode 85   @endcode
86   86  
87   @see make_stream_pair, fuse 87   @see make_stream_pair, fuse
88   */ 88   */
89   class stream 89   class stream
90   { 90   {
91   // Single-threaded only. No concurrent access to either 91   // Single-threaded only. No concurrent access to either
92   // end of the pair. Both streams and all operations must 92   // end of the pair. Both streams and all operations must
93   // run on the same thread. 93   // run on the same thread.
94   94  
95   struct half 95   struct half
96   { 96   {
97   std::string buf; 97   std::string buf;
98   std::size_t max_read_size = std::size_t(-1); 98   std::size_t max_read_size = std::size_t(-1);
99   continuation pending_cont_; 99   continuation pending_cont_;
100   executor_ref pending_ex; 100   executor_ref pending_ex;
101   // Points at the suspended reader's claim flag (owned by the 101   // Points at the suspended reader's claim flag (owned by the
102   // read awaitable). Lets a peer wake coordinate with a stop 102   // read awaitable). Lets a peer wake coordinate with a stop
103   // callback so the parked read is resumed exactly once. 103   // callback so the parked read is resumed exactly once.
104   std::atomic<bool>* pending_claimed = nullptr; 104   std::atomic<bool>* pending_claimed = nullptr;
105   bool eof = false; 105   bool eof = false;
106   }; 106   };
107   107  
108   struct state 108   struct state
109   { 109   {
110   fuse f; 110   fuse f;
111   bool closed = false; 111   bool closed = false;
112   half sides[2]; 112   half sides[2];
113   113  
HITCBC 114   315 explicit state(fuse f_) noexcept 114   315 explicit state(fuse f_) noexcept
HITCBC 115   945 : f(std::move(f_)) 115   945 : f(std::move(f_))
116   { 116   {
HITCBC 117   315 } 117   315 }
118   118  
119   // Resume a suspended reader on this side, if any. Claims the 119   // Resume a suspended reader on this side, if any. Claims the
120   // reader's atomic so it is never double-resumed by a racing 120   // reader's atomic so it is never double-resumed by a racing
121   // stop callback; the loser of the race skips the post. 121   // stop callback; the loser of the race skips the post.
HITCBC 122   704 static void wake(half& side) 122   704 static void wake(half& side)
123   { 123   {
HITCBC 124   704 if(! side.pending_cont_.h) 124   704 if(! side.pending_cont_.h)
HITCBC 125   679 return; 125   679 return;
HITCBC 126   50 if(! side.pending_claimed || 126   50 if(! side.pending_claimed ||
HITCBC 127   25 ! side.pending_claimed->exchange( 127   25 ! side.pending_claimed->exchange(
128   true, std::memory_order_acq_rel)) 128   true, std::memory_order_acq_rel))
129   { 129   {
HITCBC 130   25 side.pending_ex.post(side.pending_cont_); 130   25 side.pending_ex.post(side.pending_cont_);
131   } 131   }
HITCBC 132   25 side.pending_cont_.h = {}; 132   25 side.pending_cont_.h = {};
HITCBC 133   25 side.pending_ex = {}; 133   25 side.pending_ex = {};
HITCBC 134   25 side.pending_claimed = nullptr; 134   25 side.pending_claimed = nullptr;
135   } 135   }
136   136  
137   // Set closed and resume any suspended readers 137   // Set closed and resume any suspended readers
138   // with eof on both sides. 138   // with eof on both sides.
HITCBC 139   214 void close() 139   214 void close()
140   { 140   {
HITCBC 141   214 closed = true; 141   214 closed = true;
HITCBC 142   642 for(auto& side : sides) 142   642 for(auto& side : sides)
HITCBC 143   428 wake(side); 143   428 wake(side);
HITCBC 144   214 } 144   214 }
145   }; 145   };
146   146  
147   // Wraps the maybe_fail() call. If the guard is 147   // Wraps the maybe_fail() call. If the guard is
148   // not disarmed before destruction (fuse returned 148   // not disarmed before destruction (fuse returned
149   // an error, or threw an exception), closes both 149   // an error, or threw an exception), closes both
150   // ends so any suspended peer gets eof. 150   // ends so any suspended peer gets eof.
151   struct close_guard 151   struct close_guard
152   { 152   {
153   state* st; 153   state* st;
154   bool armed = true; 154   bool armed = true;
HITCBC 155   327 void disarm() noexcept { armed = false; } 155   327 void disarm() noexcept { armed = false; }
HITCBC 156   541 ~close_guard() noexcept(false) { if(armed) st->close(); } 156   541 ~close_guard() noexcept(false) { if(armed) st->close(); }
157   }; 157   };
158   158  
159   std::shared_ptr<state> state_; 159   std::shared_ptr<state> state_;
160   int index_; 160   int index_;
161   161  
HITCBC 162   630 stream( 162   630 stream(
163   std::shared_ptr<state> sp, 163   std::shared_ptr<state> sp,
164   int index) noexcept 164   int index) noexcept
HITCBC 165   630 : state_(std::move(sp)) 165   630 : state_(std::move(sp))
HITCBC 166   630 , index_(index) 166   630 , index_(index)
167   { 167   {
HITCBC 168   630 } 168   630 }
169   169  
170   friend std::pair<stream, stream> 170   friend std::pair<stream, stream>
171   make_stream_pair(fuse); 171   make_stream_pair(fuse);
172   172  
173   public: 173   public:
174   /** Copy construction is disabled; a stream end is move-only. 174   /** Copy construction is disabled; a stream end is move-only.
175   175  
176   @param other The stream end that would be copied. 176   @param other The stream end that would be copied.
177   */ 177   */
178   stream(stream const& other) = delete; 178   stream(stream const& other) = delete;
179   179  
180   /** Copy assignment is disabled; a stream end is move-only. 180   /** Copy assignment is disabled; a stream end is move-only.
181   181  
182   @param other The stream end that would be assigned from. 182   @param other The stream end that would be assigned from.
183   183  
184   @return A reference to `*this`. 184   @return A reference to `*this`.
185   */ 185   */
186   stream& operator=(stream const& other) = delete; 186   stream& operator=(stream const& other) = delete;
187   187  
188   /** Move constructor. 188   /** Move constructor.
189   189  
190   @param other The stream end to move from. 190   @param other The stream end to move from.
191   */ 191   */
HITCBC 192   732 stream(stream&& other) = default; 192   732 stream(stream&& other) = default;
193   193  
194   /** Move assignment. 194   /** Move assignment.
195   195  
196   @param other The stream end to move from. 196   @param other The stream end to move from.
197   197  
198   @return A reference to `*this`. 198   @return A reference to `*this`.
199   */ 199   */
200   stream& operator=(stream&& other) = default; 200   stream& operator=(stream&& other) = default;
201   201  
202   /** Signal end-of-stream to the peer. 202   /** Signal end-of-stream to the peer.
203   203  
204   Marks the peer's read direction as closed. 204   Marks the peer's read direction as closed.
205   If the peer is suspended in @ref read_some, 205   If the peer is suspended in @ref read_some,
206   it is resumed. The peer drains any buffered 206   it is resumed. The peer drains any buffered
207   data before receiving `error::eof`. Writes 207   data before receiving `error::eof`. Writes
208   from the peer are unaffected. 208   from the peer are unaffected.
209   */ 209   */
210   void 210   void
HITCBC 211   8 close() 211   8 close()
212   { 212   {
HITCBC 213   8 int peer = 1 - index_; 213   8 int peer = 1 - index_;
HITCBC 214   8 auto& side = state_->sides[peer]; 214   8 auto& side = state_->sides[peer];
HITCBC 215   8 side.eof = true; 215   8 side.eof = true;
HITCBC 216   8 state::wake(side); 216   8 state::wake(side);
HITCBC 217   8 } 217   8 }
218   218  
219   /** Set the maximum bytes returned per read. 219   /** Set the maximum bytes returned per read.
220   220  
221   Limits how many bytes @ref read_some returns in 221   Limits how many bytes @ref read_some returns in
222   a single call, simulating chunked network delivery. 222   a single call, simulating chunked network delivery.
223   The default is unlimited. 223   The default is unlimited.
224   224  
225   @param n Maximum bytes per read. 225   @param n Maximum bytes per read.
226   */ 226   */
227   void 227   void
HITCBC 228   55 set_max_read_size(std::size_t n) noexcept 228   55 set_max_read_size(std::size_t n) noexcept
229   { 229   {
HITCBC 230   55 state_->sides[index_].max_read_size = n; 230   55 state_->sides[index_].max_read_size = n;
HITCBC 231   55 } 231   55 }
232   232  
233   /** Asynchronously read data from the stream. 233   /** Asynchronously read data from the stream.
234   234  
235   Transfers up to `buffer_size(buffers)` bytes from 235   Transfers up to `buffer_size(buffers)` bytes from
236   data written by the peer. If no data is available, 236   data written by the peer. If no data is available,
237   the calling coroutine suspends until the peer calls 237   the calling coroutine suspends until the peer calls
238   @ref write_some. Before every read, the attached 238   @ref write_some. Before every read, the attached
239   @ref fuse is consulted to possibly inject an error. 239   @ref fuse is consulted to possibly inject an error.
240   If the fuse fires, the pair is automatically closed. 240   If the fuse fires, the pair is automatically closed.
241   If the stream is closed, returns `error::eof`. 241   If the stream is closed, returns `error::eof`.
242   The returned `std::size_t` is the number of bytes 242   The returned `std::size_t` is the number of bytes
243   transferred. 243   transferred.
244   244  
245   @param buffers The mutable buffer sequence to receive data. 245   @param buffers The mutable buffer sequence to receive data.
246   246  
247   @return An awaitable that await-returns `(error_code,std::size_t)`. 247   @return An awaitable that await-returns `(error_code,std::size_t)`.
248   248  
249   @par Cancellation 249   @par Cancellation
250   Cancellation applies only to a read that would otherwise suspend. 250   Cancellation applies only to a read that would otherwise suspend.
251   If no data is available and the environment's stop token is 251   If no data is available and the environment's stop token is
252   requested, before or during the wait, the read resumes with 252   requested, before or during the wait, the read resumes with
253   `error::canceled`. A read that can complete immediately from 253   `error::canceled`. A read that can complete immediately from
254   buffered data is unaffected by the stop token. 254   buffered data is unaffected by the stop token.
255   255  
256   @see fuse, close 256   @see fuse, close
257   */ 257   */
258   template<MutableBufferSequence MB> 258   template<MutableBufferSequence MB>
259   auto 259   auto
HITCBC 260   302 read_some(MB buffers) 260   302 read_some(MB buffers)
261   { 261   {
262   // The read suspends when no data is available, parking its 262   // The read suspends when no data is available, parking its
263   // continuation on the side until the peer writes/closes. To 263   // continuation on the side until the peer writes/closes. To
264   // support cancellation it follows the same pattern as 264   // support cancellation it follows the same pattern as
265   // async_waker::wait_awaiter: a stop callback claims the resume 265   // async_waker::wait_awaiter: a stop callback claims the resume
266   // (racing the peer wake via an atomic) and posts the continuation 266   // (racing the peer wake via an atomic) and posts the continuation
267   // through the executor. Because it owns a std::atomic and a 267   // through the executor. Because it owns a std::atomic and a
268   // std::stop_callback, the awaitable needs explicit move and 268   // std::stop_callback, the awaitable needs explicit move and
269   // destruction (the task promise moves it into its 269   // destruction (the task promise moves it into its
270   // transform_awaiter before awaiting). 270   // transform_awaiter before awaiting).
271   struct awaitable 271   struct awaitable
272   { 272   {
273   stream* self_; 273   stream* self_;
274   MB buffers_; 274   MB buffers_;
275   275  
276   // Declared before stop_cb_buf_: the stop callback reads 276   // Declared before stop_cb_buf_: the stop callback reads
277   // these, so they must outlive a blocking stop_cb_ destructor. 277   // these, so they must outlive a blocking stop_cb_ destructor.
278   continuation cont_; 278   continuation cont_;
279   executor_ref ex_; 279   executor_ref ex_;
280   half* side_ = nullptr; 280   half* side_ = nullptr;
281   std::atomic<bool> claimed_{false}; 281   std::atomic<bool> claimed_{false};
282   bool canceled_ = false; 282   bool canceled_ = false;
283   bool stop_cb_active_ = false; 283   bool stop_cb_active_ = false;
284   284  
285   struct cancel_fn 285   struct cancel_fn
286   { 286   {
287   awaitable* self_; 287   awaitable* self_;
288   288  
HITCBC 289   15 void operator()() const noexcept 289   15 void operator()() const noexcept
290   { 290   {
HITCBC 291   15 if(! self_->claimed_.exchange( 291   15 if(! self_->claimed_.exchange(
292   true, std::memory_order_acq_rel)) 292   true, std::memory_order_acq_rel))
293   { 293   {
HITCBC 294   3 self_->canceled_ = true; 294   3 self_->canceled_ = true;
HITCBC 295   3 self_->ex_.post(self_->cont_); 295   3 self_->ex_.post(self_->cont_);
296   } 296   }
HITCBC 297   15 } 297   15 }
298   }; 298   };
299   299  
300   using stop_cb_t = std::stop_callback<cancel_fn>; 300   using stop_cb_t = std::stop_callback<cancel_fn>;
301   301  
302   // Declared last: its destructor may block while the callback 302   // Declared last: its destructor may block while the callback
303   // accesses the members above. A union gives correct alignment 303   // accesses the members above. A union gives correct alignment
304   // for stop_cb_t without an alignas specifier, which avoids 304   // for stop_cb_t without an alignas specifier, which avoids
305   // MSVC's C4324 padding warning on this function-local class 305   // MSVC's C4324 padding warning on this function-local class
306   // (the member-level pragma used by async_waker::wait_awaiter 306   // (the member-level pragma used by async_waker::wait_awaiter
307   // does not suppress it here). Lifetime is managed manually: 307   // does not suppress it here). Lifetime is managed manually:
308   // placement new in await_suspend, explicit destruction once done. 308   // placement new in await_suspend, explicit destruction once done.
309   union { stop_cb_t stop_cb_; }; 309   union { stop_cb_t stop_cb_; };
310   310  
HITCBC 311   302 awaitable(stream* self, MB buffers) noexcept 311   302 awaitable(stream* self, MB buffers) noexcept
HITCBC 312   302 : self_(self) 312   302 : self_(self)
HITCBC 313   302 , buffers_(buffers) 313   302 , buffers_(buffers)
314   { 314   {
HITCBC 315   302 } 315   302 }
316   316  
317   /// @pre Not yet awaited (no active stop callback). 317   /// @pre Not yet awaited (no active stop callback).
HITCBC 318   292 awaitable(awaitable&& o) noexcept 318   292 awaitable(awaitable&& o) noexcept
HITCBC 319   292 : self_(o.self_) 319   292 : self_(o.self_)
HITCBC 320   292 , buffers_(o.buffers_) 320   292 , buffers_(o.buffers_)
HITCBC 321   292 , cont_(o.cont_) 321   292 , cont_(o.cont_)
HITCBC 322   292 , ex_(o.ex_) 322   292 , ex_(o.ex_)
HITCBC 323   292 , side_(o.side_) 323   292 , side_(o.side_)
HITCBC 324   292 , claimed_(o.claimed_.load(std::memory_order_relaxed)) 324   292 , claimed_(o.claimed_.load(std::memory_order_relaxed))
HITCBC 325   292 , canceled_(o.canceled_) 325   292 , canceled_(o.canceled_)
HITCBC 326   292 , stop_cb_active_(std::exchange(o.stop_cb_active_, false)) 326   292 , stop_cb_active_(std::exchange(o.stop_cb_active_, false))
327   { 327   {
HITCBC 328   292 } 328   292 }
329   329  
HITCBC 330   594 ~awaitable() 330   594 ~awaitable()
331   { 331   {
HITCBC 332   594 if(stop_cb_active_) 332   594 if(stop_cb_active_)
HITCBC 333   1 stop_cb_.~stop_cb_t(); 333   1 stop_cb_.~stop_cb_t();
334   // Unlink from the side if still parked (e.g. the 334   // Unlink from the side if still parked (e.g. the
335   // coroutine was destroyed while suspended), so a later 335   // coroutine was destroyed while suspended), so a later
336   // peer wake does not dereference a freed claim flag. 336   // peer wake does not dereference a freed claim flag.
HITCBC 337   594 if(side_ && side_->pending_claimed == &claimed_) 337   594 if(side_ && side_->pending_claimed == &claimed_)
338   { 338   {
HITCBC 339   1 side_->pending_cont_.h = {}; 339   1 side_->pending_cont_.h = {};
HITCBC 340   1 side_->pending_ex = {}; 340   1 side_->pending_ex = {};
HITCBC 341   1 side_->pending_claimed = nullptr; 341   1 side_->pending_claimed = nullptr;
342   } 342   }
HITCBC 343   594 } 343   594 }
344   344  
345   awaitable(awaitable const&) = delete; 345   awaitable(awaitable const&) = delete;
346   awaitable& operator=(awaitable const&) = delete; 346   awaitable& operator=(awaitable const&) = delete;
347   awaitable& operator=(awaitable&&) = delete; 347   awaitable& operator=(awaitable&&) = delete;
348   348  
HITCBC 349   302 bool await_ready() const noexcept 349   302 bool await_ready() const noexcept
350   { 350   {
HITCBC 351   302 if(buffer_empty(buffers_)) 351   302 if(buffer_empty(buffers_))
HITCBC 352   8 return true; 352   8 return true;
HITCBC 353   294 auto* st = self_->state_.get(); 353   294 auto* st = self_->state_.get();
HITCBC 354   294 auto& side = st->sides[self_->index_]; 354   294 auto& side = st->sides[self_->index_];
HITCBC 355   576 return st->closed || side.eof || 355   576 return st->closed || side.eof ||
HITCBC 356   576 !side.buf.empty(); 356   576 !side.buf.empty();
357   } 357   }
358   358  
HITCBC 359   29 std::coroutine_handle<> await_suspend( 359   29 std::coroutine_handle<> await_suspend(
360   std::coroutine_handle<> h, 360   std::coroutine_handle<> h,
361   io_env const* env) noexcept 361   io_env const* env) noexcept
362   { 362   {
363   // Park the continuation, then register the stop callback. 363   // Park the continuation, then register the stop callback.
364   // If stop is already requested, the callback fires inline 364   // If stop is already requested, the callback fires inline
365   // during construction: it claims the resume and posts the 365   // during construction: it claims the resume and posts the
366   // continuation through the executor (never a symmetric 366   // continuation through the executor (never a symmetric
367   // self-transfer, which would leak this frame under 367   // self-transfer, which would leak this frame under
368   // run_async). The parked read is then resumed with 368   // run_async). The parked read is then resumed with
369   // error::canceled by the run loop. 369   // error::canceled by the run loop.
HITCBC 370   29 auto& side = self_->state_->sides[ 370   29 auto& side = self_->state_->sides[
HITCBC 371   29 self_->index_]; 371   29 self_->index_];
HITCBC 372   29 cont_.h = h; 372   29 cont_.h = h;
HITCBC 373   29 ex_ = env->executor; 373   29 ex_ = env->executor;
HITCBC 374   29 side_ = &side; 374   29 side_ = &side;
HITCBC 375   29 side.pending_cont_.h = h; 375   29 side.pending_cont_.h = h;
HITCBC 376   29 side.pending_ex = env->executor; 376   29 side.pending_ex = env->executor;
HITCBC 377   29 side.pending_claimed = &claimed_; 377   29 side.pending_claimed = &claimed_;
378   378  
HITCBC 379   29 ::new(static_cast<void*>(&stop_cb_)) stop_cb_t( 379   29 ::new(static_cast<void*>(&stop_cb_)) stop_cb_t(
HITCBC 380   29 env->stop_token, cancel_fn{this}); 380   29 env->stop_token, cancel_fn{this});
HITCBC 381   29 stop_cb_active_ = true; 381   29 stop_cb_active_ = true;
382   382  
HITCBC 383   29 return std::noop_coroutine(); 383   29 return std::noop_coroutine();
384   } 384   }
385   385  
386   [[nodiscard]] io_result<std::size_t> 386   [[nodiscard]] io_result<std::size_t>
HITCBC 387   301 await_resume() 387   301 await_resume()
388   { 388   {
HITCBC 389   301 if(stop_cb_active_) 389   301 if(stop_cb_active_)
390   { 390   {
HITCBC 391   28 stop_cb_.~stop_cb_t(); 391   28 stop_cb_.~stop_cb_t();
HITCBC 392   28 stop_cb_active_ = false; 392   28 stop_cb_active_ = false;
393   } 393   }
394   394  
HITCBC 395   301 if(buffer_empty(buffers_)) 395   301 if(buffer_empty(buffers_))
HITCBC 396   8 return {std::error_code(), 0}; 396   8 return {std::error_code(), 0};
397   397  
HITCBC 398   293 if(canceled_) 398   293 if(canceled_)
399   { 399   {
400   // The stop callback posted us but left the side 400   // The stop callback posted us but left the side
401   // untouched; unlink if a peer wake has not already. 401   // untouched; unlink if a peer wake has not already.
HITCBC 402   3 if(side_ && side_->pending_claimed == &claimed_) 402   3 if(side_ && side_->pending_claimed == &claimed_)
403   { 403   {
HITCBC 404   3 side_->pending_cont_.h = {}; 404   3 side_->pending_cont_.h = {};
HITCBC 405   3 side_->pending_ex = {}; 405   3 side_->pending_ex = {};
HITCBC 406   3 side_->pending_claimed = nullptr; 406   3 side_->pending_claimed = nullptr;
407   } 407   }
HITCBC 408   3 return {error::canceled, 0}; 408   3 return {error::canceled, 0};
409   } 409   }
410   410  
HITCBC 411   290 auto* st = self_->state_.get(); 411   290 auto* st = self_->state_.get();
HITCBC 412   290 auto& side = st->sides[ 412   290 auto& side = st->sides[
HITCBC 413   290 self_->index_]; 413   290 self_->index_];
414   414  
HITCBC 415   290 if(st->closed) 415   290 if(st->closed)
HITCBC 416   12 return {error::eof, 0}; 416   12 return {error::eof, 0};
417   417  
HITCBC 418   278 if(side.eof && side.buf.empty()) 418   278 if(side.eof && side.buf.empty())
HITCBC 419   8 return {error::eof, 0}; 419   8 return {error::eof, 0};
420   420  
HITCBC 421   270 if(!side.eof) 421   270 if(!side.eof)
422   { 422   {
HITCBC 423   265 close_guard g{st}; 423   265 close_guard g{st};
HITCBC 424   265 auto ec = st->f.maybe_fail(); 424   265 auto ec = st->f.maybe_fail();
HITCBC 425   211 if(ec) 425   211 if(ec)
HITCBC 426   54 return {ec, 0}; 426   54 return {ec, 0};
HITCBC 427   157 g.disarm(); 427   157 g.disarm();
HITCBC 428   265 } 428   265 }
429   429  
HITCBC 430   486 std::size_t const n = buffer_copy( 430   486 std::size_t const n = buffer_copy(
HITCBC 431   162 buffers_, make_buffer(side.buf), 431   162 buffers_, make_buffer(side.buf),
432   side.max_read_size); 432   side.max_read_size);
HITCBC 433   162 side.buf.erase(0, n); 433   162 side.buf.erase(0, n);
HITCBC 434   162 return {std::error_code(), n}; 434   162 return {std::error_code(), n};
435   } 435   }
436   }; 436   };
HITCBC 437   302 return awaitable{this, buffers}; 437   302 return awaitable{this, buffers};
438   } 438   }
439   439  
440   /** Asynchronously write data to the stream. 440   /** Asynchronously write data to the stream.
441   441  
442   Transfers up to `buffer_size(buffers)` bytes to the 442   Transfers up to `buffer_size(buffers)` bytes to the
443   peer's incoming buffer. If the peer is suspended in 443   peer's incoming buffer. If the peer is suspended in
444   @ref read_some, it is resumed. Before every write, 444   @ref read_some, it is resumed. Before every write,
445   the attached @ref fuse is consulted to possibly inject 445   the attached @ref fuse is consulted to possibly inject
446   an error. If the fuse fires, the pair is automatically 446   an error. If the fuse fires, the pair is automatically
447   closed. If the stream is closed, returns `error::eof`. 447   closed. If the stream is closed, returns `error::eof`.
448   The returned `std::size_t` is the number of bytes 448   The returned `std::size_t` is the number of bytes
449   transferred. 449   transferred.
450   450  
451   @param buffers The const buffer sequence containing 451   @param buffers The const buffer sequence containing
452   data to write. 452   data to write.
453   453  
454   @return An awaitable that await-returns `(error_code,std::size_t)`. 454   @return An awaitable that await-returns `(error_code,std::size_t)`.
455   455  
456   @par Cancellation 456   @par Cancellation
457   If the environment's stop token is requested, the write 457   If the environment's stop token is requested, the write
458   completes immediately with `error::canceled` and transfers no 458   completes immediately with `error::canceled` and transfers no
459   data. An empty buffer sequence is a no-op that completes 459   data. An empty buffer sequence is a no-op that completes
460   successfully regardless of the stop token. 460   successfully regardless of the stop token.
461   461  
462   @see fuse, close 462   @see fuse, close
463   */ 463   */
464   template<ConstBufferSequence CB> 464   template<ConstBufferSequence CB>
465   auto 465   auto
HITCBC 466   281 write_some(CB buffers) 466   281 write_some(CB buffers)
467   { 467   {
468   struct awaitable 468   struct awaitable
469   { 469   {
470   stream* self_; 470   stream* self_;
471   CB buffers_; 471   CB buffers_;
472   bool canceled_ = false; 472   bool canceled_ = false;
473   473  
HITCBC 474   281 bool await_ready() const noexcept { return false; } 474   281 bool await_ready() const noexcept { return false; }
475   475  
476   // The write completes synchronously; await_suspend is only 476   // The write completes synchronously; await_suspend is only
477   // used to observe the environment's stop token. Returning 477   // used to observe the environment's stop token. Returning
478   // false means the coroutine does not actually suspend. 478   // false means the coroutine does not actually suspend.
479   bool 479   bool
HITCBC 480   281 await_suspend( 480   281 await_suspend(
481   std::coroutine_handle<>, 481   std::coroutine_handle<>,
482   io_env const* env) noexcept 482   io_env const* env) noexcept
483   { 483   {
HITCBC 484   281 canceled_ = env->stop_token.stop_requested(); 484   281 canceled_ = env->stop_token.stop_requested();
HITCBC 485   281 return false; 485   281 return false;
486   } 486   }
487   487  
488   [[nodiscard]] io_result<std::size_t> 488   [[nodiscard]] io_result<std::size_t>
HITCBC 489   281 await_resume() 489   281 await_resume()
490   { 490   {
HITCBC 491   281 std::size_t n = buffer_size(buffers_); 491   281 std::size_t n = buffer_size(buffers_);
HITCBC 492   281 if(n == 0) 492   281 if(n == 0)
HITCBC 493   4 return {std::error_code(), 0}; 493   4 return {std::error_code(), 0};
494   494  
HITCBC 495   277 if(canceled_) 495   277 if(canceled_)
HITCBC 496   1 return {error::canceled, 0}; 496   1 return {error::canceled, 0};
497   497  
HITCBC 498   276 auto* st = self_->state_.get(); 498   276 auto* st = self_->state_.get();
499   499  
HITCBC 500   276 if(st->closed) 500   276 if(st->closed)
MISUBC 501   return {error::eof, 0}; 501   return {error::eof, 0};
502   502  
HITCBC 503   276 close_guard g{st}; 503   276 close_guard g{st};
HITCBC 504   276 auto ec = st->f.maybe_fail(); 504   276 auto ec = st->f.maybe_fail();
HITCBC 505   223 if(ec) 505   223 if(ec)
HITCBC 506   53 return {ec, 0}; 506   53 return {ec, 0};
HITCBC 507   170 g.disarm(); 507   170 g.disarm();
508   508  
HITCBC 509   170 int peer = 1 - self_->index_; 509   170 int peer = 1 - self_->index_;
HITCBC 510   170 auto& side = st->sides[peer]; 510   170 auto& side = st->sides[peer];
511   511  
HITCBC 512   170 std::size_t const old_size = side.buf.size(); 512   170 std::size_t const old_size = side.buf.size();
HITCBC 513   170 side.buf.resize(old_size + n); 513   170 side.buf.resize(old_size + n);
HITCBC 514   170 buffer_copy(make_buffer( 514   170 buffer_copy(make_buffer(
HITCBC 515   170 side.buf.data() + old_size, n), 515   170 side.buf.data() + old_size, n),
HITCBC 516   170 buffers_, n); 516   170 buffers_, n);
517   517  
HITCBC 518   170 state::wake(side); 518   170 state::wake(side);
519   519  
HITCBC 520   170 return {std::error_code(), n}; 520   170 return {std::error_code(), n};
HITCBC 521   276 } 521   276 }
522   }; 522   };
HITCBC 523   281 return awaitable{this, buffers}; 523   281 return awaitable{this, buffers};
524   } 524   }
525   525  
526   /** Inject data into this stream's peer for reading. 526   /** Inject data into this stream's peer for reading.
527   527  
528   Appends data directly to the peer's incoming buffer, 528   Appends data directly to the peer's incoming buffer,
529   bypassing the fuse. If the peer is suspended in 529   bypassing the fuse. If the peer is suspended in
530   @ref read_some, it is resumed. This is test setup, 530   @ref read_some, it is resumed. This is test setup,
531   not an operation under test. 531   not an operation under test.
532   532  
533   @param sv The data to inject. 533   @param sv The data to inject.
534   534  
535   @see make_stream_pair 535   @see make_stream_pair
536   */ 536   */
537   void 537   void
HITCBC 538   98 provide(std::string_view sv) 538   98 provide(std::string_view sv)
539   { 539   {
HITCBC 540   98 int peer = 1 - index_; 540   98 int peer = 1 - index_;
HITCBC 541   98 auto& side = state_->sides[peer]; 541   98 auto& side = state_->sides[peer];
HITCBC 542   98 side.buf.append(sv); 542   98 side.buf.append(sv);
HITCBC 543   98 state::wake(side); 543   98 state::wake(side);
HITCBC 544   98 } 544   98 }
545   545  
546   /** Read from this stream and verify the content. 546   /** Read from this stream and verify the content.
547   547  
548   Reads exactly `expected.size()` bytes from the stream 548   Reads exactly `expected.size()` bytes from the stream
549   and compares against the expected string. The read goes 549   and compares against the expected string. The read goes
550   through the normal path including the fuse. 550   through the normal path including the fuse.
551   551  
552   @param expected The expected content. 552   @param expected The expected content.
553   553  
554   @return A pair of `(error_code, bool)`. The error_code 554   @return A pair of `(error_code, bool)`. The error_code
555   is set if a read error occurs (e.g. fuse injection). 555   is set if a read error occurs (e.g. fuse injection).
556   The bool is true if the data matches. 556   The bool is true if the data matches.
557   557  
558   @see provide 558   @see provide
559   */ 559   */
560   std::pair<std::error_code, bool> 560   std::pair<std::error_code, bool>
HITCBC 561   38 expect(std::string_view expected) 561   38 expect(std::string_view expected)
562   { 562   {
HITCBC 563   38 std::error_code result; 563   38 std::error_code result;
HITCBC 564   38 bool match = false; 564   38 bool match = false;
HITCBC 565   141 run_blocking()([]( 565   141 run_blocking()([](
566   stream& self, 566   stream& self,
567   std::string_view expected, 567   std::string_view expected,
568   std::error_code& result, 568   std::error_code& result,
569   bool& match) -> task<> 569   bool& match) -> task<>
570   { 570   {
571   std::string buf(expected.size(), '\0'); 571   std::string buf(expected.size(), '\0');
572   auto [ec, n] = co_await read( 572   auto [ec, n] = co_await read(
573   self, mutable_buffer( 573   self, mutable_buffer(
574   buf.data(), buf.size())); 574   buf.data(), buf.size()));
575   if(ec) 575   if(ec)
576   { 576   {
577   result = ec; 577   result = ec;
578   co_return; 578   co_return;
579   } 579   }
580   match = (std::string_view( 580   match = (std::string_view(
581   buf.data(), n) == expected); 581   buf.data(), n) == expected);
HITCBC 582   161 }(*this, expected, result, match)); 582   161 }(*this, expected, result, match));
HITCBC 583   58 return {result, match}; 583   58 return {result, match};
584   } 584   }
585   585  
586   /** Return the stream's pending read data. 586   /** Return the stream's pending read data.
587   587  
588   Returns a view of the data waiting to be read 588   Returns a view of the data waiting to be read
589   from this stream. This is a direct peek at the 589   from this stream. This is a direct peek at the
590   internal buffer, bypassing the fuse. 590   internal buffer, bypassing the fuse.
591   591  
592   @return A view of the pending data. 592   @return A view of the pending data.
593   593  
594   @see provide, expect 594   @see provide, expect
595   */ 595   */
596   std::string_view 596   std::string_view
HITCBC 597   9 data() const noexcept 597   9 data() const noexcept
598   { 598   {
HITCBC 599   9 return state_->sides[index_].buf; 599   9 return state_->sides[index_].buf;
600   } 600   }
601   }; 601   };
602   602  
603   /** Create a connected pair of test streams. 603   /** Create a connected pair of test streams.
604   604  
605   Data written to one stream becomes readable on the other. 605   Data written to one stream becomes readable on the other.
606   If a coroutine calls @ref stream::read_some when no data 606   If a coroutine calls @ref stream::read_some when no data
607   is available, it suspends until the peer writes. Before 607   is available, it suspends until the peer writes. Before
608   every read or write, the @ref fuse is consulted to 608   every read or write, the @ref fuse is consulted to
609   possibly inject an error for testing fault scenarios. 609   possibly inject an error for testing fault scenarios.
610   When the fuse fires, the pair is automatically closed. 610   When the fuse fires, the pair is automatically closed.
611   611  
612   @param f The fuse used to inject errors during operations. 612   @param f The fuse used to inject errors during operations.
613   613  
614   @return A pair of connected streams. 614   @return A pair of connected streams.
615   615  
616   @see stream, fuse 616   @see stream, fuse
617   */ 617   */
618   inline std::pair<stream, stream> 618   inline std::pair<stream, stream>
HITCBC 619   315 make_stream_pair(fuse f = {}) 619   315 make_stream_pair(fuse f = {})
620   { 620   {
HITCBC 621   315 auto sp = std::make_shared<stream::state>(std::move(f)); 621   315 auto sp = std::make_shared<stream::state>(std::move(f));
HITCBC 622   630 return {stream(sp, 0), stream(sp, 1)}; 622   630 return {stream(sp, 0), stream(sp, 1)};
HITCBC 623   315 } 623   315 }
624   624  
625   } // test 625   } // test
626   } // capy 626   } // capy
627   } // boost 627   } // boost
628   628  
629   #endif 629   #endif