100.00% Lines (75/75)
100.00% Functions (16/16)
| 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_IO_ANY_READ_STREAM_HPP | 11 | #ifndef BOOST_CAPY_IO_ANY_READ_STREAM_HPP | |||||
| 12 | #define BOOST_CAPY_IO_ANY_READ_STREAM_HPP | 12 | #define BOOST_CAPY_IO_ANY_READ_STREAM_HPP | |||||
| 13 | 13 | |||||||
| 14 | #include <boost/capy/detail/config.hpp> | 14 | #include <boost/capy/detail/config.hpp> | |||||
| 15 | #include <boost/capy/detail/await_suspend_helper.hpp> | 15 | #include <boost/capy/detail/await_suspend_helper.hpp> | |||||
| 16 | #include <boost/capy/buffers.hpp> | 16 | #include <boost/capy/buffers.hpp> | |||||
| 17 | #include <boost/capy/detail/buffer_array.hpp> | 17 | #include <boost/capy/detail/buffer_array.hpp> | |||||
| 18 | #include <boost/capy/concept/io_awaitable.hpp> | 18 | #include <boost/capy/concept/io_awaitable.hpp> | |||||
| 19 | #include <boost/capy/concept/read_stream.hpp> | 19 | #include <boost/capy/concept/read_stream.hpp> | |||||
| 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 | 22 | |||||||
| 23 | #include <concepts> | 23 | #include <concepts> | |||||
| 24 | #include <coroutine> | 24 | #include <coroutine> | |||||
| 25 | #include <cstddef> | 25 | #include <cstddef> | |||||
| 26 | #include <exception> | 26 | #include <exception> | |||||
| 27 | #include <new> | 27 | #include <new> | |||||
| 28 | #include <span> | 28 | #include <span> | |||||
| 29 | #include <stop_token> | 29 | #include <stop_token> | |||||
| 30 | #include <system_error> | 30 | #include <system_error> | |||||
| 31 | #include <utility> | 31 | #include <utility> | |||||
| 32 | 32 | |||||||
| 33 | namespace boost { | 33 | namespace boost { | |||||
| 34 | namespace capy { | 34 | namespace capy { | |||||
| 35 | 35 | |||||||
| 36 | /** Dispatches `read_some` through a type-erased vtable, using preallocated awaitable storage. | 36 | /** Dispatches `read_some` through a type-erased vtable, using preallocated awaitable storage. | |||||
| 37 | 37 | |||||||
| 38 | This class provides type erasure for any type satisfying the | 38 | This class provides type erasure for any type satisfying the | |||||
| 39 | @ref ReadStream concept, enabling runtime polymorphism for | 39 | @ref ReadStream concept, enabling runtime polymorphism for | |||||
| 40 | read operations. It uses cached awaitable storage to achieve | 40 | read operations. It uses cached awaitable storage to achieve | |||||
| 41 | zero steady-state allocation after construction. | 41 | zero steady-state allocation after construction. | |||||
| 42 | 42 | |||||||
| 43 | The wrapper supports two construction modes: | 43 | The wrapper supports two construction modes: | |||||
| 44 | - **Owning**: Pass by value to transfer ownership. The wrapper | 44 | - **Owning**: Pass by value to transfer ownership. The wrapper | |||||
| 45 | allocates storage and owns the stream. | 45 | allocates storage and owns the stream. | |||||
| 46 | - **Reference**: Pass a pointer to wrap without ownership. The | 46 | - **Reference**: Pass a pointer to wrap without ownership. The | |||||
| 47 | pointed-to stream must outlive this wrapper. | 47 | pointed-to stream must outlive this wrapper. | |||||
| 48 | 48 | |||||||
| 49 | @par Awaitable Preallocation | 49 | @par Awaitable Preallocation | |||||
| 50 | The constructor preallocates storage for the type-erased awaitable. | 50 | The constructor preallocates storage for the type-erased awaitable. | |||||
| 51 | This reserves all virtual address space at server startup | 51 | This reserves all virtual address space at server startup | |||||
| 52 | so memory usage can be measured up front, rather than | 52 | so memory usage can be measured up front, rather than | |||||
| 53 | allocating piecemeal as traffic arrives. | 53 | allocating piecemeal as traffic arrives. | |||||
| 54 | 54 | |||||||
| 55 | @par Immediate Completion | 55 | @par Immediate Completion | |||||
| 56 | When the underlying stream's awaitable reports ready immediately | 56 | When the underlying stream's awaitable reports ready immediately | |||||
| 57 | (e.g. buffered data already available), the wrapper skips | 57 | (e.g. buffered data already available), the wrapper skips | |||||
| 58 | coroutine suspension entirely and returns the result inline. | 58 | coroutine suspension entirely and returns the result inline. | |||||
| 59 | 59 | |||||||
| 60 | @par Thread Safety | 60 | @par Thread Safety | |||||
| 61 | Not thread-safe. Concurrent operations on the same wrapper | 61 | Not thread-safe. Concurrent operations on the same wrapper | |||||
| 62 | are undefined behavior. | 62 | are undefined behavior. | |||||
| 63 | 63 | |||||||
| 64 | @par Example | 64 | @par Example | |||||
| 65 | - | @code | 65 | + | @par !example example | |||
| 66 | - | // Owning - takes ownership of the stream | ||||||
| 67 | - | any_read_stream owning_stream(socket{ioc}); | ||||||
| 68 | - | |||||||
| 69 | - | // Reference - wraps without ownership | ||||||
| 70 | - | socket sock(ioc); | ||||||
| 71 | - | any_read_stream ref_stream(&sock); | ||||||
| 72 | - | char data[1024]; | ||||||
| 73 | - | mutable_buffer buf(data, sizeof(data)); | ||||||
| 74 | - | auto [ec, n] = co_await owning_stream.read_some(buf); | ||||||
| 75 | - | @endcode | ||||||
| 76 | 66 | |||||||
| 77 | 67 | |||||||
| 78 | @see any_write_stream, any_stream, ReadStream | 68 | @see any_write_stream, any_stream, ReadStream | |||||
| 79 | */ | 69 | */ | |||||
| 80 | class any_read_stream | 70 | class any_read_stream | |||||
| 81 | { | 71 | { | |||||
| 82 | struct vtable; | 72 | struct vtable; | |||||
| 83 | 73 | |||||||
| 84 | template<ReadStream S> | 74 | template<ReadStream S> | |||||
| 85 | struct vtable_for_impl; | 75 | struct vtable_for_impl; | |||||
| 86 | 76 | |||||||
| 87 | // ordered for cache line coherence | 77 | // ordered for cache line coherence | |||||
| 88 | void* stream_ = nullptr; | 78 | void* stream_ = nullptr; | |||||
| 89 | vtable const* vt_ = nullptr; | 79 | vtable const* vt_ = nullptr; | |||||
| 90 | void* cached_awaitable_ = nullptr; | 80 | void* cached_awaitable_ = nullptr; | |||||
| 91 | void* storage_ = nullptr; | 81 | void* storage_ = nullptr; | |||||
| 92 | bool awaitable_active_ = false; | 82 | bool awaitable_active_ = false; | |||||
| 93 | 83 | |||||||
| 94 | public: | 84 | public: | |||||
| 95 | /** Destructor. | 85 | /** Destructor. | |||||
| 96 | 86 | |||||||
| 97 | Destroys the owned stream (if any) and releases the cached | 87 | Destroys the owned stream (if any) and releases the cached | |||||
| 98 | awaitable storage. | 88 | awaitable storage. | |||||
| 99 | */ | 89 | */ | |||||
| 100 | ~any_read_stream(); | 90 | ~any_read_stream(); | |||||
| 101 | 91 | |||||||
| 102 | /** Construct a default instance. | 92 | /** Construct a default instance. | |||||
| 103 | 93 | |||||||
| 104 | Constructs an empty wrapper. @ref has_value and `operator bool` | 94 | Constructs an empty wrapper. @ref has_value and `operator bool` | |||||
| 105 | report the empty state; calling @ref read_some before the | 95 | report the empty state; calling @ref read_some before the | |||||
| 106 | wrapper holds a stream is undefined behavior. | 96 | wrapper holds a stream is undefined behavior. | |||||
| 107 | */ | 97 | */ | |||||
| HITCBC | 108 | 4 | any_read_stream() = default; | 98 | 4 | any_read_stream() = default; | ||
| 109 | 99 | |||||||
| 110 | /** Non-copyable. | 100 | /** Non-copyable. | |||||
| 111 | 101 | |||||||
| 112 | The awaitable cache is per-instance and cannot be shared. | 102 | The awaitable cache is per-instance and cannot be shared. | |||||
| 113 | 103 | |||||||
| 114 | @param other The wrapper that would be copied. | 104 | @param other The wrapper that would be copied. | |||||
| 115 | */ | 105 | */ | |||||
| 116 | any_read_stream(any_read_stream const& other) = delete; | 106 | any_read_stream(any_read_stream const& other) = delete; | |||||
| 117 | 107 | |||||||
| 118 | /** Copy assignment is disabled. | 108 | /** Copy assignment is disabled. | |||||
| 119 | 109 | |||||||
| 120 | The awaitable cache is per-instance and cannot be shared. | 110 | The awaitable cache is per-instance and cannot be shared. | |||||
| 121 | 111 | |||||||
| 122 | @param other The wrapper that would be assigned from. | 112 | @param other The wrapper that would be assigned from. | |||||
| 123 | 113 | |||||||
| 124 | @return A reference to `*this`. | 114 | @return A reference to `*this`. | |||||
| 125 | */ | 115 | */ | |||||
| 126 | any_read_stream& operator=(any_read_stream const& other) = delete; | 116 | any_read_stream& operator=(any_read_stream const& other) = delete; | |||||
| 127 | 117 | |||||||
| 128 | /** Construct by moving. | 118 | /** Construct by moving. | |||||
| 129 | 119 | |||||||
| 130 | Transfers ownership of the wrapped stream (if owned) and | 120 | Transfers ownership of the wrapped stream (if owned) and | |||||
| 131 | cached awaitable storage from `other`. After the move, `other` is | 121 | cached awaitable storage from `other`. After the move, `other` is | |||||
| 132 | in a default-constructed state. | 122 | in a default-constructed state. | |||||
| 133 | 123 | |||||||
| 134 | @param other The wrapper to move from. | 124 | @param other The wrapper to move from. | |||||
| 135 | */ | 125 | */ | |||||
| HITCBC | 136 | 4 | any_read_stream(any_read_stream&& other) noexcept | 126 | 4 | any_read_stream(any_read_stream&& other) noexcept | ||
| HITCBC | 137 | 4 | : stream_(std::exchange(other.stream_, nullptr)) | 127 | 4 | : stream_(std::exchange(other.stream_, nullptr)) | ||
| HITCBC | 138 | 4 | , vt_(std::exchange(other.vt_, nullptr)) | 128 | 4 | , vt_(std::exchange(other.vt_, nullptr)) | ||
| HITCBC | 139 | 4 | , cached_awaitable_(std::exchange(other.cached_awaitable_, nullptr)) | 129 | 4 | , cached_awaitable_(std::exchange(other.cached_awaitable_, nullptr)) | ||
| HITCBC | 140 | 4 | , storage_(std::exchange(other.storage_, nullptr)) | 130 | 4 | , storage_(std::exchange(other.storage_, nullptr)) | ||
| HITCBC | 141 | 4 | , awaitable_active_(std::exchange(other.awaitable_active_, false)) | 131 | 4 | , awaitable_active_(std::exchange(other.awaitable_active_, false)) | ||
| 142 | { | 132 | { | |||||
| HITCBC | 143 | 4 | } | 133 | 4 | } | ||
| 144 | 134 | |||||||
| 145 | /** Assign by moving. | 135 | /** Assign by moving. | |||||
| 146 | 136 | |||||||
| 147 | Destroys any owned stream and releases existing resources, | 137 | Destroys any owned stream and releases existing resources, | |||||
| 148 | then transfers ownership from `other`. | 138 | then transfers ownership from `other`. | |||||
| 149 | 139 | |||||||
| 150 | @param other The wrapper to move from. | 140 | @param other The wrapper to move from. | |||||
| 151 | @return Reference to this wrapper. | 141 | @return Reference to this wrapper. | |||||
| 152 | */ | 142 | */ | |||||
| 153 | any_read_stream& | 143 | any_read_stream& | |||||
| 154 | operator=(any_read_stream&& other) noexcept; | 144 | operator=(any_read_stream&& other) noexcept; | |||||
| 155 | 145 | |||||||
| 156 | /** Construct by taking ownership of a ReadStream. | 146 | /** Construct by taking ownership of a ReadStream. | |||||
| 157 | 147 | |||||||
| 158 | Allocates storage and moves the stream into this wrapper. | 148 | Allocates storage and moves the stream into this wrapper. | |||||
| 159 | The wrapper owns the stream and destroys it. | 149 | The wrapper owns the stream and destroys it. | |||||
| 160 | 150 | |||||||
| 161 | @param s The stream to take ownership of. | 151 | @param s The stream to take ownership of. | |||||
| 162 | */ | 152 | */ | |||||
| 163 | template<ReadStream S> | 153 | template<ReadStream S> | |||||
| 164 | requires (!std::same_as<std::decay_t<S>, any_read_stream>) | 154 | requires (!std::same_as<std::decay_t<S>, any_read_stream>) | |||||
| 165 | any_read_stream(S s); | 155 | any_read_stream(S s); | |||||
| 166 | 156 | |||||||
| 167 | /** Construct by wrapping a ReadStream without ownership. | 157 | /** Construct by wrapping a ReadStream without ownership. | |||||
| 168 | 158 | |||||||
| 169 | Wraps the given stream by pointer. The stream must remain | 159 | Wraps the given stream by pointer. The stream must remain | |||||
| 170 | valid for the lifetime of this wrapper. | 160 | valid for the lifetime of this wrapper. | |||||
| 171 | 161 | |||||||
| 172 | @param s Pointer to the stream to wrap. | 162 | @param s Pointer to the stream to wrap. | |||||
| 173 | */ | 163 | */ | |||||
| 174 | template<ReadStream S> | 164 | template<ReadStream S> | |||||
| 175 | any_read_stream(S* s); | 165 | any_read_stream(S* s); | |||||
| 176 | 166 | |||||||
| 177 | /** Check if the wrapper contains a valid stream. | 167 | /** Check if the wrapper contains a valid stream. | |||||
| 178 | 168 | |||||||
| 179 | @return `true` if wrapping a stream, `false` if default-constructed | 169 | @return `true` if wrapping a stream, `false` if default-constructed | |||||
| 180 | or moved-from. | 170 | or moved-from. | |||||
| 181 | */ | 171 | */ | |||||
| 182 | bool | 172 | bool | |||||
| HITCBC | 183 | 31 | has_value() const noexcept | 173 | 31 | has_value() const noexcept | ||
| 184 | { | 174 | { | |||||
| HITCBC | 185 | 31 | return stream_ != nullptr; | 175 | 31 | return stream_ != nullptr; | ||
| 186 | } | 176 | } | |||||
| 187 | 177 | |||||||
| 188 | /** Check if the wrapper contains a valid stream. | 178 | /** Check if the wrapper contains a valid stream. | |||||
| 189 | 179 | |||||||
| 190 | @return `true` if wrapping a stream, `false` if default-constructed | 180 | @return `true` if wrapping a stream, `false` if default-constructed | |||||
| 191 | or moved-from. | 181 | or moved-from. | |||||
| 192 | */ | 182 | */ | |||||
| 193 | explicit | 183 | explicit | |||||
| HITCBC | 194 | 3 | operator bool() const noexcept | 184 | 3 | operator bool() const noexcept | ||
| 195 | { | 185 | { | |||||
| HITCBC | 196 | 3 | return has_value(); | 186 | 3 | return has_value(); | ||
| 197 | } | 187 | } | |||||
| 198 | 188 | |||||||
| 199 | /** Initiate an asynchronous read operation. | 189 | /** Initiate an asynchronous read operation. | |||||
| 200 | 190 | |||||||
| 201 | Reads data into the provided buffer sequence. The operation | 191 | Reads data into the provided buffer sequence. The operation | |||||
| 202 | completes when at least one byte is read, or an error | 192 | completes when at least one byte is read, or an error | |||||
| 203 | occurs. | 193 | occurs. | |||||
| 204 | 194 | |||||||
| 205 | @param buffers The buffer sequence to read into. Passed by | 195 | @param buffers The buffer sequence to read into. Passed by | |||||
| 206 | value to ensure the sequence lives in the coroutine frame | 196 | value to ensure the sequence lives in the coroutine frame | |||||
| 207 | across suspension points. | 197 | across suspension points. | |||||
| 208 | 198 | |||||||
| 209 | @return An awaitable that await-returns `(error_code,std::size_t)`. | 199 | @return An awaitable that await-returns `(error_code,std::size_t)`. | |||||
| 210 | 200 | |||||||
| 211 | @par Immediate Completion | 201 | @par Immediate Completion | |||||
| 212 | The operation completes immediately without suspending | 202 | The operation completes immediately without suspending | |||||
| 213 | the calling coroutine when the underlying stream's | 203 | the calling coroutine when the underlying stream's | |||||
| 214 | awaitable reports immediate readiness via `await_ready`. | 204 | awaitable reports immediate readiness via `await_ready`. | |||||
| 215 | 205 | |||||||
| 216 | @note This is a partial operation and may not process the | 206 | @note This is a partial operation and may not process the | |||||
| 217 | entire buffer sequence. Use the composed @ref read algorithm | 207 | entire buffer sequence. Use the composed @ref read algorithm | |||||
| 218 | for guaranteed complete transfer. | 208 | for guaranteed complete transfer. | |||||
| 219 | 209 | |||||||
| 220 | @par Preconditions | 210 | @par Preconditions | |||||
| 221 | The wrapper must contain a valid stream (`has_value() == true`). | 211 | The wrapper must contain a valid stream (`has_value() == true`). | |||||
| 222 | 212 | |||||||
| 223 | @par After an Error | 213 | @par After an Error | |||||
| 224 | A subsequent call is permitted. The wrapper forwards directly | 214 | A subsequent call is permitted. The wrapper forwards directly | |||||
| 225 | to the underlying stream, imposing no stricter rule than | 215 | to the underlying stream, imposing no stricter rule than | |||||
| 226 | @ref ReadStream. | 216 | @ref ReadStream. | |||||
| 227 | */ | 217 | */ | |||||
| 228 | template<MutableBufferSequence MB> | 218 | template<MutableBufferSequence MB> | |||||
| 229 | auto | 219 | auto | |||||
| 230 | read_some(MB buffers); | 220 | read_some(MB buffers); | |||||
| 231 | 221 | |||||||
| 232 | protected: | 222 | protected: | |||||
| 233 | /** Rebind to a new stream after move. | 223 | /** Rebind to a new stream after move. | |||||
| 234 | 224 | |||||||
| 235 | Updates the internal pointer to reference a new stream object. | 225 | Updates the internal pointer to reference a new stream object. | |||||
| 236 | Used by owning wrappers after move assignment when the owned | 226 | Used by owning wrappers after move assignment when the owned | |||||
| 237 | object has moved to a new location. | 227 | object has moved to a new location. | |||||
| 238 | 228 | |||||||
| 239 | @param new_stream The new stream to bind to. Must be the same | 229 | @param new_stream The new stream to bind to. Must be the same | |||||
| 240 | type as the original stream. | 230 | type as the original stream. | |||||
| 241 | 231 | |||||||
| 242 | @note Terminates if called with a stream of different type | 232 | @note Terminates if called with a stream of different type | |||||
| 243 | than the original. | 233 | than the original. | |||||
| 244 | */ | 234 | */ | |||||
| 245 | template<ReadStream S> | 235 | template<ReadStream S> | |||||
| 246 | void | 236 | void | |||||
| 247 | rebind(S& new_stream) noexcept | 237 | rebind(S& new_stream) noexcept | |||||
| 248 | { | 238 | { | |||||
| 249 | if(vt_ != &vtable_for_impl<S>::value) | 239 | if(vt_ != &vtable_for_impl<S>::value) | |||||
| 250 | std::terminate(); | 240 | std::terminate(); | |||||
| 251 | stream_ = &new_stream; | 241 | stream_ = &new_stream; | |||||
| 252 | } | 242 | } | |||||
| 253 | }; | 243 | }; | |||||
| 254 | 244 | |||||||
| 255 | struct any_read_stream::vtable | 245 | struct any_read_stream::vtable | |||||
| 256 | { | 246 | { | |||||
| 257 | // ordered by call frequency for cache line coherence | 247 | // ordered by call frequency for cache line coherence | |||||
| 258 | void (*construct_awaitable)( | 248 | void (*construct_awaitable)( | |||||
| 259 | void* stream, | 249 | void* stream, | |||||
| 260 | void* storage, | 250 | void* storage, | |||||
| 261 | std::span<mutable_buffer const> buffers); | 251 | std::span<mutable_buffer const> buffers); | |||||
| 262 | bool (*await_ready)(void*); | 252 | bool (*await_ready)(void*); | |||||
| 263 | std::coroutine_handle<> (*await_suspend)(void*, std::coroutine_handle<>, io_env const*); | 253 | std::coroutine_handle<> (*await_suspend)(void*, std::coroutine_handle<>, io_env const*); | |||||
| 264 | io_result<std::size_t> (*await_resume)(void*); | 254 | io_result<std::size_t> (*await_resume)(void*); | |||||
| 265 | void (*destroy_awaitable)(void*) noexcept; | 255 | void (*destroy_awaitable)(void*) noexcept; | |||||
| 266 | std::size_t awaitable_size; | 256 | std::size_t awaitable_size; | |||||
| 267 | std::size_t awaitable_align; | 257 | std::size_t awaitable_align; | |||||
| 268 | void (*destroy)(void*) noexcept; | 258 | void (*destroy)(void*) noexcept; | |||||
| 269 | }; | 259 | }; | |||||
| 270 | 260 | |||||||
| 271 | template<ReadStream S> | 261 | template<ReadStream S> | |||||
| 272 | struct any_read_stream::vtable_for_impl | 262 | struct any_read_stream::vtable_for_impl | |||||
| 273 | { | 263 | { | |||||
| 274 | using Awaitable = decltype(std::declval<S&>().read_some( | 264 | using Awaitable = decltype(std::declval<S&>().read_some( | |||||
| 275 | std::span<mutable_buffer const>{})); | 265 | std::span<mutable_buffer const>{})); | |||||
| 276 | 266 | |||||||
| 277 | static void | 267 | static void | |||||
| HITCBC | 278 | 4 | do_destroy_impl(void* stream) noexcept | 268 | 4 | do_destroy_impl(void* stream) noexcept | ||
| 279 | { | 269 | { | |||||
| HITCBC | 280 | 4 | static_cast<S*>(stream)->~S(); | 270 | 4 | static_cast<S*>(stream)->~S(); | ||
| HITCBC | 281 | 4 | } | 271 | 4 | } | ||
| 282 | 272 | |||||||
| 283 | static void | 273 | static void | |||||
| HITCBC | 284 | 103 | construct_awaitable_impl( | 274 | 103 | construct_awaitable_impl( | ||
| 285 | void* stream, | 275 | void* stream, | |||||
| 286 | void* storage, | 276 | void* storage, | |||||
| 287 | std::span<mutable_buffer const> buffers) | 277 | std::span<mutable_buffer const> buffers) | |||||
| 288 | { | 278 | { | |||||
| HITCBC | 289 | 103 | auto& s = *static_cast<S*>(stream); | 279 | 103 | auto& s = *static_cast<S*>(stream); | ||
| HITCBC | 290 | 103 | ::new(storage) Awaitable(s.read_some(buffers)); | 280 | 103 | ::new(storage) Awaitable(s.read_some(buffers)); | ||
| HITCBC | 291 | 103 | } | 281 | 103 | } | ||
| 292 | 282 | |||||||
| 293 | static constexpr vtable value = { | 283 | static constexpr vtable value = { | |||||
| 294 | &construct_awaitable_impl, | 284 | &construct_awaitable_impl, | |||||
| HITCBC | 295 | 103 | +[](void* p) { | 285 | 103 | +[](void* p) { | ||
| HITCBC | 296 | 103 | return static_cast<Awaitable*>(p)->await_ready(); | 286 | 103 | return static_cast<Awaitable*>(p)->await_ready(); | ||
| 297 | }, | 287 | }, | |||||
| HITCBC | 298 | 77 | +[](void* p, std::coroutine_handle<> h, io_env const* env) { | 288 | 77 | +[](void* p, std::coroutine_handle<> h, io_env const* env) { | ||
| HITCBC | 299 | 77 | return detail::call_await_suspend( | 289 | 77 | return detail::call_await_suspend( | ||
| HITCBC | 300 | 77 | static_cast<Awaitable*>(p), h, env); | 290 | 77 | static_cast<Awaitable*>(p), h, env); | ||
| 301 | }, | 291 | }, | |||||
| HITCBC | 302 | 101 | +[](void* p) { | 292 | 101 | +[](void* p) { | ||
| HITCBC | 303 | 101 | return static_cast<Awaitable*>(p)->await_resume(); | 293 | 101 | return static_cast<Awaitable*>(p)->await_resume(); | ||
| 304 | }, | 294 | }, | |||||
| HITCBC | 305 | 115 | +[](void* p) noexcept { | 295 | 115 | +[](void* p) noexcept { | ||
| HITCBC | 306 | 26 | static_cast<Awaitable*>(p)->~Awaitable(); | 296 | 26 | static_cast<Awaitable*>(p)->~Awaitable(); | ||
| 307 | }, | 297 | }, | |||||
| 308 | sizeof(Awaitable), | 298 | sizeof(Awaitable), | |||||
| 309 | alignof(Awaitable), | 299 | alignof(Awaitable), | |||||
| 310 | &do_destroy_impl | 300 | &do_destroy_impl | |||||
| 311 | }; | 301 | }; | |||||
| 312 | }; | 302 | }; | |||||
| 313 | 303 | |||||||
| 314 | inline | 304 | inline | |||||
| HITCBC | 315 | 123 | any_read_stream::~any_read_stream() | 305 | 123 | any_read_stream::~any_read_stream() | ||
| 316 | { | 306 | { | |||||
| HITCBC | 317 | 123 | if(storage_) | 307 | 123 | if(storage_) | ||
| 318 | { | 308 | { | |||||
| HITCBC | 319 | 3 | vt_->destroy(stream_); | 309 | 3 | vt_->destroy(stream_); | ||
| HITCBC | 320 | 3 | ::operator delete(storage_); | 310 | 3 | ::operator delete(storage_); | ||
| 321 | } | 311 | } | |||||
| HITCBC | 322 | 123 | if(cached_awaitable_) | 312 | 123 | if(cached_awaitable_) | ||
| 323 | { | 313 | { | |||||
| HITCBC | 324 | 106 | if(awaitable_active_) | 314 | 106 | if(awaitable_active_) | ||
| HITCBC | 325 | 1 | vt_->destroy_awaitable(cached_awaitable_); | 315 | 1 | vt_->destroy_awaitable(cached_awaitable_); | ||
| HITCBC | 326 | 106 | ::operator delete(cached_awaitable_); | 316 | 106 | ::operator delete(cached_awaitable_); | ||
| 327 | } | 317 | } | |||||
| HITCBC | 328 | 123 | } | 318 | 123 | } | ||
| 329 | 319 | |||||||
| 330 | inline any_read_stream& | 320 | inline any_read_stream& | |||||
| HITCBC | 331 | 10 | any_read_stream::operator=(any_read_stream&& other) noexcept | 321 | 10 | any_read_stream::operator=(any_read_stream&& other) noexcept | ||
| 332 | { | 322 | { | |||||
| HITCBC | 333 | 10 | if(this != &other) | 323 | 10 | if(this != &other) | ||
| 334 | { | 324 | { | |||||
| HITCBC | 335 | 10 | if(storage_) | 325 | 10 | if(storage_) | ||
| 336 | { | 326 | { | |||||
| HITCBC | 337 | 1 | vt_->destroy(stream_); | 327 | 1 | vt_->destroy(stream_); | ||
| HITCBC | 338 | 1 | ::operator delete(storage_); | 328 | 1 | ::operator delete(storage_); | ||
| 339 | } | 329 | } | |||||
| HITCBC | 340 | 10 | if(cached_awaitable_) | 330 | 10 | if(cached_awaitable_) | ||
| 341 | { | 331 | { | |||||
| HITCBC | 342 | 4 | if(awaitable_active_) | 332 | 4 | if(awaitable_active_) | ||
| HITCBC | 343 | 1 | vt_->destroy_awaitable(cached_awaitable_); | 333 | 1 | vt_->destroy_awaitable(cached_awaitable_); | ||
| HITCBC | 344 | 4 | ::operator delete(cached_awaitable_); | 334 | 4 | ::operator delete(cached_awaitable_); | ||
| 345 | } | 335 | } | |||||
| HITCBC | 346 | 10 | stream_ = std::exchange(other.stream_, nullptr); | 336 | 10 | stream_ = std::exchange(other.stream_, nullptr); | ||
| HITCBC | 347 | 10 | vt_ = std::exchange(other.vt_, nullptr); | 337 | 10 | vt_ = std::exchange(other.vt_, nullptr); | ||
| HITCBC | 348 | 10 | cached_awaitable_ = std::exchange(other.cached_awaitable_, nullptr); | 338 | 10 | cached_awaitable_ = std::exchange(other.cached_awaitable_, nullptr); | ||
| HITCBC | 349 | 10 | storage_ = std::exchange(other.storage_, nullptr); | 339 | 10 | storage_ = std::exchange(other.storage_, nullptr); | ||
| HITCBC | 350 | 10 | awaitable_active_ = std::exchange(other.awaitable_active_, false); | 340 | 10 | awaitable_active_ = std::exchange(other.awaitable_active_, false); | ||
| 351 | } | 341 | } | |||||
| HITCBC | 352 | 10 | return *this; | 342 | 10 | return *this; | ||
| 353 | } | 343 | } | |||||
| 354 | 344 | |||||||
| 355 | template<ReadStream S> | 345 | template<ReadStream S> | |||||
| 356 | requires (!std::same_as<std::decay_t<S>, any_read_stream>) | 346 | requires (!std::same_as<std::decay_t<S>, any_read_stream>) | |||||
| HITCBC | 357 | 5 | any_read_stream::any_read_stream(S s) | 347 | 5 | any_read_stream::any_read_stream(S s) | ||
| HITCBC | 358 | 5 | : vt_(&vtable_for_impl<S>::value) | 348 | 5 | : vt_(&vtable_for_impl<S>::value) | ||
| 359 | { | 349 | { | |||||
| 360 | struct guard { | 350 | struct guard { | |||||
| 361 | any_read_stream* self; | 351 | any_read_stream* self; | |||||
| 362 | bool committed = false; | 352 | bool committed = false; | |||||
| HITCBC | 363 | 5 | ~guard() { | 353 | 5 | ~guard() { | ||
| HITCBC | 364 | 5 | if(!committed && self->storage_) { | 354 | 5 | if(!committed && self->storage_) { | ||
| HITCBC | 365 | 1 | if(self->stream_) | 355 | 1 | if(self->stream_) | ||
| 366 | self->vt_->destroy(self->stream_); // LCOV_EXCL_LINE OOM rollback: only when the cached-awaitable allocation throws | 356 | self->vt_->destroy(self->stream_); // LCOV_EXCL_LINE OOM rollback: only when the cached-awaitable allocation throws | |||||
| HITCBC | 367 | 1 | ::operator delete(self->storage_); | 357 | 1 | ::operator delete(self->storage_); | ||
| HITCBC | 368 | 1 | self->storage_ = nullptr; | 358 | 1 | self->storage_ = nullptr; | ||
| HITCBC | 369 | 1 | self->stream_ = nullptr; | 359 | 1 | self->stream_ = nullptr; | ||
| 370 | } | 360 | } | |||||
| HITCBC | 371 | 5 | } | 361 | 5 | } | ||
| HITCBC | 372 | 5 | } g{this}; | 362 | 5 | } g{this}; | ||
| 373 | 363 | |||||||
| HITCBC | 374 | 5 | storage_ = ::operator new(sizeof(S)); | 364 | 5 | storage_ = ::operator new(sizeof(S)); | ||
| HITCBC | 375 | 5 | stream_ = ::new(storage_) S(std::move(s)); | 365 | 5 | stream_ = ::new(storage_) S(std::move(s)); | ||
| 376 | 366 | |||||||
| 377 | // Preallocate the awaitable storage | 367 | // Preallocate the awaitable storage | |||||
| HITCBC | 378 | 4 | cached_awaitable_ = ::operator new(vt_->awaitable_size); | 368 | 4 | cached_awaitable_ = ::operator new(vt_->awaitable_size); | ||
| 379 | 369 | |||||||
| HITCBC | 380 | 4 | g.committed = true; | 370 | 4 | g.committed = true; | ||
| HITCBC | 381 | 5 | } | 371 | 5 | } | ||
| 382 | 372 | |||||||
| 383 | template<ReadStream S> | 373 | template<ReadStream S> | |||||
| HITCBC | 384 | 106 | any_read_stream::any_read_stream(S* s) | 374 | 106 | any_read_stream::any_read_stream(S* s) | ||
| HITCBC | 385 | 106 | : stream_(s) | 375 | 106 | : stream_(s) | ||
| HITCBC | 386 | 106 | , vt_(&vtable_for_impl<S>::value) | 376 | 106 | , vt_(&vtable_for_impl<S>::value) | ||
| 387 | { | 377 | { | |||||
| 388 | // Preallocate the awaitable storage | 378 | // Preallocate the awaitable storage | |||||
| HITCBC | 389 | 106 | cached_awaitable_ = ::operator new(vt_->awaitable_size); | 379 | 106 | cached_awaitable_ = ::operator new(vt_->awaitable_size); | ||
| HITCBC | 390 | 106 | } | 380 | 106 | } | ||
| 391 | 381 | |||||||
| 392 | template<MutableBufferSequence MB> | 382 | template<MutableBufferSequence MB> | |||||
| 393 | auto | 383 | auto | |||||
| HITCBC | 394 | 103 | any_read_stream::read_some(MB buffers) | 384 | 103 | any_read_stream::read_some(MB buffers) | ||
| 395 | { | 385 | { | |||||
| 396 | // VFALCO in theory, we could use if constexpr to detect a | 386 | // VFALCO in theory, we could use if constexpr to detect a | |||||
| 397 | // span and then pass that through to read_some without the array | 387 | // span and then pass that through to read_some without the array | |||||
| 398 | // LCOV_EXCL_START read_some awaitable: exercised by tests, but the | 388 | // LCOV_EXCL_START read_some awaitable: exercised by tests, but the | |||||
| 399 | // coverage tooling reports its templated body uncovered per-instantiation | 389 | // coverage tooling reports its templated body uncovered per-instantiation | |||||
| 400 | struct awaitable | 390 | struct awaitable | |||||
| 401 | { | 391 | { | |||||
| 402 | any_read_stream* self_; | 392 | any_read_stream* self_; | |||||
| 403 | detail::mutable_buffer_array<detail::max_iovec_> ba_; | 393 | detail::mutable_buffer_array<detail::max_iovec_> ba_; | |||||
| 404 | 394 | |||||||
| 405 | bool | 395 | bool | |||||
| 406 | await_ready() | 396 | await_ready() | |||||
| 407 | { | 397 | { | |||||
| 408 | self_->vt_->construct_awaitable( | 398 | self_->vt_->construct_awaitable( | |||||
| 409 | self_->stream_, | 399 | self_->stream_, | |||||
| 410 | self_->cached_awaitable_, | 400 | self_->cached_awaitable_, | |||||
| 411 | ba_.to_span()); | 401 | ba_.to_span()); | |||||
| 412 | self_->awaitable_active_ = true; | 402 | self_->awaitable_active_ = true; | |||||
| 413 | 403 | |||||||
| 414 | return self_->vt_->await_ready( | 404 | return self_->vt_->await_ready( | |||||
| 415 | self_->cached_awaitable_); | 405 | self_->cached_awaitable_); | |||||
| 416 | } | 406 | } | |||||
| 417 | 407 | |||||||
| 418 | std::coroutine_handle<> | 408 | std::coroutine_handle<> | |||||
| 419 | await_suspend(std::coroutine_handle<> h, io_env const* env) | 409 | await_suspend(std::coroutine_handle<> h, io_env const* env) | |||||
| 420 | { | 410 | { | |||||
| 421 | return self_->vt_->await_suspend( | 411 | return self_->vt_->await_suspend( | |||||
| 422 | self_->cached_awaitable_, h, env); | 412 | self_->cached_awaitable_, h, env); | |||||
| 423 | } | 413 | } | |||||
| 424 | 414 | |||||||
| 425 | [[nodiscard]] io_result<std::size_t> | 415 | [[nodiscard]] io_result<std::size_t> | |||||
| 426 | await_resume() | 416 | await_resume() | |||||
| 427 | { | 417 | { | |||||
| 428 | struct guard { | 418 | struct guard { | |||||
| 429 | any_read_stream* self; | 419 | any_read_stream* self; | |||||
| 430 | ~guard() { | 420 | ~guard() { | |||||
| 431 | self->vt_->destroy_awaitable(self->cached_awaitable_); | 421 | self->vt_->destroy_awaitable(self->cached_awaitable_); | |||||
| 432 | self->awaitable_active_ = false; | 422 | self->awaitable_active_ = false; | |||||
| 433 | } | 423 | } | |||||
| 434 | } g{self_}; | 424 | } g{self_}; | |||||
| 435 | return self_->vt_->await_resume( | 425 | return self_->vt_->await_resume( | |||||
| 436 | self_->cached_awaitable_); | 426 | self_->cached_awaitable_); | |||||
| 437 | } | 427 | } | |||||
| 438 | }; | 428 | }; | |||||
| 439 | // LCOV_EXCL_STOP | 429 | // LCOV_EXCL_STOP | |||||
| 440 | return awaitable{this, | 430 | return awaitable{this, | |||||
| HITCBC | 441 | 103 | detail::mutable_buffer_array<detail::max_iovec_>(buffers)}; | 431 | 103 | detail::mutable_buffer_array<detail::max_iovec_>(buffers)}; | ||
| HITCBC | 442 | 103 | } | 432 | 103 | } | ||
| 443 | 433 | |||||||
| 444 | } // namespace capy | 434 | } // namespace capy | |||||
| 445 | } // namespace boost | 435 | } // namespace boost | |||||
| 446 | 436 | |||||||
| 447 | #endif | 437 | #endif | |||||