100.00% Lines (36/36) 100.00% Functions (8/8)
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_READ_STREAM_HPP 11   #ifndef BOOST_CAPY_TEST_READ_STREAM_HPP
12   #define BOOST_CAPY_TEST_READ_STREAM_HPP 12   #define BOOST_CAPY_TEST_READ_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/cond.hpp> 18   #include <boost/capy/cond.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/test/fuse.hpp> 22   #include <boost/capy/test/fuse.hpp>
23   23  
24   #include <string> 24   #include <string>
25   #include <string_view> 25   #include <string_view>
26   26  
27   namespace boost { 27   namespace boost {
28   namespace capy { 28   namespace capy {
29   namespace test { 29   namespace test {
30   30  
31   /** Buffers data supplied via `provide`, then hands it out through `read_some`. 31   /** Buffers data supplied via `provide`, then hands it out through `read_some`.
32   32  
33   Use this to verify code that performs reads without needing 33   Use this to verify code that performs reads without needing
34   real I/O. Call @ref provide to supply data, then @ref read_some 34   real I/O. Call @ref provide to supply data, then @ref read_some
35   to consume it. The associated @ref fuse enables error injection 35   to consume it. The associated @ref fuse enables error injection
36   at controlled points. An optional `max_read_size` constructor 36   at controlled points. An optional `max_read_size` constructor
37   parameter limits bytes per read to simulate chunked delivery. 37   parameter limits bytes per read to simulate chunked delivery.
38   38  
39   This class satisfies the @ref ReadStream concept. 39   This class satisfies the @ref ReadStream concept.
40   40  
41   @par Thread Safety 41   @par Thread Safety
42   Not thread-safe. 42   Not thread-safe.
43   43  
44   @par Example 44   @par Example
45   @code 45   @code
46   fuse f; 46   fuse f;
47   read_stream rs( f ); 47   read_stream rs( f );
48   rs.provide( "Hello, " ); 48   rs.provide( "Hello, " );
49   rs.provide( "World!" ); 49   rs.provide( "World!" );
50   50  
51   auto r = f.armed( [&]( fuse& ) -> task<void> { 51   auto r = f.armed( [&]( fuse& ) -> task<void> {
52   char buf[32]; 52   char buf[32];
53   auto [ec, n] = co_await rs.read_some( 53   auto [ec, n] = co_await rs.read_some(
54   mutable_buffer( buf, sizeof( buf ) ) ); 54   mutable_buffer( buf, sizeof( buf ) ) );
55   if( ec ) 55   if( ec )
56   co_return; 56   co_return;
57   // buf contains "Hello, World!" 57   // buf contains "Hello, World!"
58   } ); 58   } );
59   @endcode 59   @endcode
60   60  
61   @see fuse, ReadStream 61   @see fuse, ReadStream
62   */ 62   */
63   class read_stream 63   class read_stream
64   { 64   {
65   fuse f_; 65   fuse f_;
66   std::string data_; 66   std::string data_;
67   std::size_t pos_ = 0; 67   std::size_t pos_ = 0;
68   std::size_t max_read_size_; 68   std::size_t max_read_size_;
69   69  
70   public: 70   public:
71   /** Construct a read stream. 71   /** Construct a read stream.
72   72  
73   @param f The fuse used to inject errors during reads. 73   @param f The fuse used to inject errors during reads.
74   74  
75   @param max_read_size Maximum bytes returned per read. 75   @param max_read_size Maximum bytes returned per read.
76   Use to simulate chunked network delivery. 76   Use to simulate chunked network delivery.
77   */ 77   */
HITCBC 78   305 explicit read_stream( 78   305 explicit read_stream(
79   fuse f = {}, 79   fuse f = {},
80   std::size_t max_read_size = std::size_t(-1)) noexcept 80   std::size_t max_read_size = std::size_t(-1)) noexcept
HITCBC 81   305 : f_(std::move(f)) 81   305 : f_(std::move(f))
HITCBC 82   305 , max_read_size_(max_read_size) 82   305 , max_read_size_(max_read_size)
83   { 83   {
HITCBC 84   305 } 84   305 }
85   85  
86   /** Append data to be returned by subsequent reads. 86   /** Append data to be returned by subsequent reads.
87   87  
88   Multiple calls accumulate data that @ref read_some returns. 88   Multiple calls accumulate data that @ref read_some returns.
89   89  
90   @param sv The data to append. 90   @param sv The data to append.
91   */ 91   */
92   void 92   void
HITCBC 93   307 provide(std::string_view sv) 93   307 provide(std::string_view sv)
94   { 94   {
HITCBC 95   307 data_.append(sv); 95   307 data_.append(sv);
HITCBC 96   307 } 96   307 }
97   97  
98   /// Clear all data and reset the read position. 98   /// Clear all data and reset the read position.
99   void 99   void
HITCBC 100   6 clear() noexcept 100   6 clear() noexcept
101   { 101   {
HITCBC 102   6 data_.clear(); 102   6 data_.clear();
HITCBC 103   6 pos_ = 0; 103   6 pos_ = 0;
HITCBC 104   6 } 104   6 }
105   105  
106   /** Return the number of bytes available for reading. 106   /** Return the number of bytes available for reading.
107   107  
108   @return The number of provided bytes not yet consumed. 108   @return The number of provided bytes not yet consumed.
109   */ 109   */
110   std::size_t 110   std::size_t
HITCBC 111   24 available() const noexcept 111   24 available() const noexcept
112   { 112   {
HITCBC 113   24 return data_.size() - pos_; 113   24 return data_.size() - pos_;
114   } 114   }
115   115  
116   /** Asynchronously read data from the stream. 116   /** Asynchronously read data from the stream.
117   117  
118   Transfers up to `buffer_size( buffers )` bytes from the internal 118   Transfers up to `buffer_size( buffers )` bytes from the internal
119   buffer to the provided mutable buffer sequence. If no data remains, 119   buffer to the provided mutable buffer sequence. If no data remains,
120   returns `error::eof`. Before every read, the attached @ref fuse is 120   returns `error::eof`. Before every read, the attached @ref fuse is
121   consulted to possibly inject an error for testing fault scenarios. 121   consulted to possibly inject an error for testing fault scenarios.
122   The returned `std::size_t` is the number of bytes transferred. 122   The returned `std::size_t` is the number of bytes transferred.
123   123  
124   @par Effects 124   @par Effects
125   On success, advances the internal read position by the number of 125   On success, advances the internal read position by the number of
126   bytes copied. If an error is injected by the fuse, the read position 126   bytes copied. If an error is injected by the fuse, the read position
127   remains unchanged. 127   remains unchanged.
128   128  
129   @par Exception Safety 129   @par Exception Safety
130   Injected I/O conditions are reported via the `error_code` 130   Injected I/O conditions are reported via the `error_code`
131   component of the result. Throws `std::system_error` only when 131   component of the result. Throws `std::system_error` only when
132   the attached @ref fuse is in exception mode and reaches its 132   the attached @ref fuse is in exception mode and reaches its
133   failure point; no-throw otherwise. 133   failure point; no-throw otherwise.
134   134  
135   @par Cancellation 135   @par Cancellation
136   If the environment's stop token is requested, the read 136   If the environment's stop token is requested, the read
137   completes immediately with `error::canceled` and transfers no 137   completes immediately with `error::canceled` and transfers no
138   data. This lets code under test exercise its cancellation paths. 138   data. This lets code under test exercise its cancellation paths.
139   An empty buffer sequence is a no-op that completes successfully 139   An empty buffer sequence is a no-op that completes successfully
140   regardless of the stop token. 140   regardless of the stop token.
141   141  
142   @param buffers The mutable buffer sequence to receive data. 142   @param buffers The mutable buffer sequence to receive data.
143   143  
144   @return An awaitable that await-returns `(error_code,std::size_t)`. 144   @return An awaitable that await-returns `(error_code,std::size_t)`.
145   145  
146   @throws std::system_error When the attached @ref fuse is in 146   @throws std::system_error When the attached @ref fuse is in
147   exception mode and reaches its failure point. 147   exception mode and reaches its failure point.
148   148  
149   @see fuse 149   @see fuse
150   */ 150   */
151   template<MutableBufferSequence MB> 151   template<MutableBufferSequence MB>
152   auto 152   auto
HITCBC 153   430 read_some(MB buffers) 153   430 read_some(MB buffers)
154   { 154   {
155   struct awaitable 155   struct awaitable
156   { 156   {
157   read_stream* self_; 157   read_stream* self_;
158   MB buffers_; 158   MB buffers_;
159   bool canceled_ = false; 159   bool canceled_ = false;
160   160  
HITCBC 161   430 bool await_ready() const noexcept { return false; } 161   430 bool await_ready() const noexcept { return false; }
162   162  
163   // The operation completes synchronously, but await_suspend 163   // The operation completes synchronously, but await_suspend
164   // is the only place io_env is delivered (the promise's 164   // is the only place io_env is delivered (the promise's
165   // transform_awaiter forwards it here). Returning false means 165   // transform_awaiter forwards it here). Returning false means
166   // the coroutine does not actually suspend — it resumes 166   // the coroutine does not actually suspend — it resumes
167   // immediately — so the read still completes synchronously 167   // immediately — so the read still completes synchronously
168   // while having observed the stop token. See io_env, IoAwaitable. 168   // while having observed the stop token. See io_env, IoAwaitable.
169   bool 169   bool
HITCBC 170   430 await_suspend( 170   430 await_suspend(
171   std::coroutine_handle<>, 171   std::coroutine_handle<>,
172   io_env const* env) noexcept 172   io_env const* env) noexcept
173   { 173   {
HITCBC 174   430 canceled_ = env->stop_token.stop_requested(); 174   430 canceled_ = env->stop_token.stop_requested();
HITCBC 175   430 return false; 175   430 return false;
176   } 176   }
177   177  
178   [[nodiscard]] io_result<std::size_t> 178   [[nodiscard]] io_result<std::size_t>
HITCBC 179   430 await_resume() 179   430 await_resume()
180   { 180   {
181   // Empty buffer is a no-op regardless of 181   // Empty buffer is a no-op regardless of
182   // stream state, stop token, or fuse. 182   // stream state, stop token, or fuse.
HITCBC 183   430 if(buffer_empty(buffers_)) 183   430 if(buffer_empty(buffers_))
HITCBC 184   7 return {std::error_code(), 0}; 184   7 return {std::error_code(), 0};
185   185  
HITCBC 186   423 if(canceled_) 186   423 if(canceled_)
HITCBC 187   2 return {error::canceled, 0}; 187   2 return {error::canceled, 0};
188   188  
HITCBC 189   421 auto ec = self_->f_.maybe_fail(); 189   421 auto ec = self_->f_.maybe_fail();
HITCBC 190   330 if(ec) 190   330 if(ec)
HITCBC 191   91 return {ec, 0}; 191   91 return {ec, 0};
192   192  
HITCBC 193   239 if(self_->pos_ >= self_->data_.size()) 193   239 if(self_->pos_ >= self_->data_.size())
HITCBC 194   37 return {error::eof, 0}; 194   37 return {error::eof, 0};
195   195  
HITCBC 196   202 std::size_t avail = self_->data_.size() - self_->pos_; 196   202 std::size_t avail = self_->data_.size() - self_->pos_;
HITCBC 197   202 if(avail > self_->max_read_size_) 197   202 if(avail > self_->max_read_size_)
HITCBC 198   24 avail = self_->max_read_size_; 198   24 avail = self_->max_read_size_;
HITCBC 199   202 auto src = make_buffer(self_->data_.data() + self_->pos_, avail); 199   202 auto src = make_buffer(self_->data_.data() + self_->pos_, avail);
HITCBC 200   202 std::size_t const n = buffer_copy(buffers_, src); 200   202 std::size_t const n = buffer_copy(buffers_, src);
HITCBC 201   202 self_->pos_ += n; 201   202 self_->pos_ += n;
HITCBC 202   202 return {std::error_code(), n}; 202   202 return {std::error_code(), n};
203   } 203   }
204   }; 204   };
HITCBC 205   430 return awaitable{this, buffers}; 205   430 return awaitable{this, buffers};
206   } 206   }
207   }; 207   };
208   208  
209   } // test 209   } // test
210   } // capy 210   } // capy
211   } // boost 211   } // boost
212   212  
213   #endif 213   #endif