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