LCOV - code coverage report
Current view: top level - capy/io - any_read_stream.hpp (source / functions) Coverage Total Hit Missed
Test: coverage_remapped.info Lines: 100.0 % 75 75
Test Date: 2026-08-21 22:12:46 Functions: 58.8 % 80 47 33

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

Generated by: LCOV version 2.3