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