LCOV - code coverage report
Current view: top level - corosio/native/detail/reactor - reactor_stream_socket.hpp (source / functions) Coverage Total Hit Missed
Test: coverage_remapped.info Lines: 96.8 % 282 273 9
Test Date: 2026-10-08 17:58:26 Functions: 95.8 % 96 92 4

           TLA  Line data    Source code
       1                 : //
       2                 : // Copyright (c) 2026 Steve Gerbino
       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/corosio
       9                 : //
      10                 : 
      11                 : #ifndef BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_STREAM_SOCKET_HPP
      12                 : #define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_STREAM_SOCKET_HPP
      13                 : 
      14                 : #include <boost/corosio/tcp_socket.hpp>
      15                 : #include <boost/corosio/shutdown_type.hpp>
      16                 : #include <boost/corosio/wait_type.hpp>
      17                 : #include <boost/corosio/native/detail/reactor/reactor_basic_socket.hpp>
      18                 : #include <boost/corosio/native/detail/reactor/reactor_descriptor_state.hpp>
      19                 : #include <boost/corosio/detail/dispatch_coro.hpp>
      20                 : #include <boost/capy/buffers.hpp>
      21                 : 
      22                 : #include <coroutine>
      23                 : 
      24                 : #include <errno.h>
      25                 : #include <sys/socket.h>
      26                 : #include <sys/uio.h>
      27                 : 
      28                 : namespace boost::corosio::detail {
      29                 : 
      30                 : /** CRTP base for reactor-backed stream socket implementations.
      31                 : 
      32                 :     Inherits shared data members and cancel/close/register logic
      33                 :     from reactor_basic_socket. Adds the stream-specific remote
      34                 :     endpoint, shutdown, and I/O dispatch (connect, read, write, wait).
      35                 : 
      36                 :     @tparam Derived   The concrete socket type (CRTP).
      37                 :     @tparam Service   The backend's socket service type.
      38                 :     @tparam ConnOp    The backend's connect op type.
      39                 :     @tparam ReadOp    The backend's read op type.
      40                 :     @tparam WriteOp   The backend's write op type.
      41                 :     @tparam WaitOp    The backend's wait op type.
      42                 :     @tparam DescState The backend's descriptor_state type.
      43                 :     @tparam ImplBase  The public vtable base
      44                 :                       (tcp_socket::implementation or
      45                 :                        local_stream_socket::implementation).
      46                 :     @tparam Endpoint  The endpoint type (endpoint or local_endpoint).
      47                 : */
      48                 : template<
      49                 :     class Derived,
      50                 :     class Service,
      51                 :     class ConnOp,
      52                 :     class ReadOp,
      53                 :     class WriteOp,
      54                 :     class WaitOp,
      55                 :     class DescState,
      56                 :     class ImplBase = tcp_socket::implementation,
      57                 :     class Endpoint = endpoint>
      58                 : class reactor_stream_socket
      59                 :     : public reactor_basic_socket<
      60                 :           Derived,
      61                 :           ImplBase,
      62                 :           Service,
      63                 :           DescState,
      64                 :           Endpoint>
      65                 : {
      66                 :     using base_type =
      67                 :         reactor_basic_socket<Derived, ImplBase, Service, DescState, Endpoint>;
      68                 :     using self_type = reactor_stream_socket<
      69                 :         Derived,
      70                 :         Service,
      71                 :         ConnOp,
      72                 :         ReadOp,
      73                 :         WriteOp,
      74                 :         WaitOp,
      75                 :         DescState,
      76                 :         ImplBase,
      77                 :         Endpoint>;
      78                 :     friend base_type;
      79                 :     friend reactor_io_core<Derived, Service, DescState>;
      80                 :     friend Derived;
      81                 : 
      82                 : protected:
      83                 :     // NOLINTNEXTLINE(bugprone-crtp-constructor-accessibility)
      84 HIT       14466 :     explicit reactor_stream_socket(Service& svc) noexcept : base_type(svc) {}
      85                 : 
      86                 : protected:
      87                 :     Endpoint remote_endpoint_;
      88                 : 
      89                 : public:
      90                 :     /// Pending connect operation slot.
      91                 :     ConnOp conn_;
      92                 : 
      93                 :     /// Pending read operation slot.
      94                 :     ReadOp rd_;
      95                 : 
      96                 :     /// Pending write operation slot.
      97                 :     WriteOp wr_;
      98                 : 
      99                 :     /// Pending wait-for-read operation slot.
     100                 :     WaitOp wait_rd_;
     101                 : 
     102                 :     /// Pending wait-for-write operation slot.
     103                 :     WaitOp wait_wr_;
     104                 : 
     105                 :     /// Pending wait-for-error operation slot.
     106                 :     WaitOp wait_er_;
     107                 : 
     108           14466 :     ~reactor_stream_socket() override = default;
     109                 : 
     110                 :     /// Return the cached remote endpoint.
     111              62 :     Endpoint remote_endpoint() const noexcept override
     112                 :     {
     113              62 :         return remote_endpoint_;
     114                 :     }
     115                 : 
     116                 :     // --- Virtual method overrides (satisfy ImplBase pure virtuals) ---
     117                 : 
     118            4667 :     std::coroutine_handle<> connect(
     119                 :         std::coroutine_handle<> h,
     120                 :         capy::executor_ref ex,
     121                 :         Endpoint ep,
     122                 :         std::stop_token token,
     123                 :         std::error_code* ec) override
     124                 :     {
     125            4667 :         return do_connect(h, ex, ep, token, ec);
     126                 :     }
     127                 : 
     128          210443 :     std::coroutine_handle<> read_some(
     129                 :         std::coroutine_handle<> h,
     130                 :         capy::executor_ref ex,
     131                 :         buffer_param param,
     132                 :         std::stop_token token,
     133                 :         std::error_code* ec,
     134                 :         std::size_t* bytes_out) override
     135                 :     {
     136          210443 :         return do_read_some(h, ex, param, token, ec, bytes_out);
     137                 :     }
     138                 : 
     139          209742 :     std::coroutine_handle<> write_some(
     140                 :         std::coroutine_handle<> h,
     141                 :         capy::executor_ref ex,
     142                 :         buffer_param param,
     143                 :         std::stop_token token,
     144                 :         std::error_code* ec,
     145                 :         std::size_t* bytes_out) override
     146                 :     {
     147          209742 :         return do_write_some(h, ex, param, token, ec, bytes_out);
     148                 :     }
     149                 : 
     150              96 :     std::coroutine_handle<> wait(
     151                 :         std::coroutine_handle<> h,
     152                 :         capy::executor_ref ex,
     153                 :         wait_type w,
     154                 :         std::stop_token token,
     155                 :         std::error_code* ec) override
     156                 :     {
     157              96 :         return do_wait(h, ex, w, token, ec);
     158                 :     }
     159                 : 
     160              25 :     std::error_code shutdown(corosio::shutdown_type what) noexcept override
     161                 :     {
     162              25 :         return do_shutdown(static_cast<int>(what));
     163                 :     }
     164                 : 
     165             226 :     void cancel() noexcept override
     166                 :     {
     167             226 :         this->do_cancel();
     168             226 :     }
     169                 : 
     170                 :     // --- End virtual overrides ---
     171                 : 
     172                 :     /// Close the socket (non-virtual, called by the service).
     173                 :     void close_socket() noexcept
     174                 :     {
     175                 :         this->do_close_socket();
     176                 :     }
     177                 : 
     178                 :     /** Shut down part or all of the full-duplex connection.
     179                 : 
     180                 :         @param what 0 = receive, 1 = send, 2 = both.
     181                 :     */
     182              25 :     std::error_code do_shutdown(int what) noexcept
     183                 :     {
     184                 :         int how;
     185              25 :         switch (what)
     186                 :         {
     187               4 :         case 0: // shutdown_receive
     188               4 :             how = SHUT_RD;
     189               4 :             break;
     190              17 :         case 1: // shutdown_send
     191              17 :             how = SHUT_WR;
     192              17 :             break;
     193               4 :         case 2: // shutdown_both
     194               4 :             how = SHUT_RDWR;
     195               4 :             break;
     196 MIS           0 :         default:
     197               0 :             return make_err(EINVAL);
     198                 :         }
     199 HIT          25 :         if (::shutdown(this->fd_, how) != 0)
     200               2 :             return make_err(errno);
     201              23 :         return {};
     202                 :     }
     203                 : 
     204                 :     /// Cache local and remote endpoints.
     205            9392 :     void set_endpoints(Endpoint local, Endpoint remote) noexcept
     206                 :     {
     207            9392 :         this->local_endpoint_ = std::move(local);
     208            9392 :         remote_endpoint_      = std::move(remote);
     209            9392 :     }
     210                 : 
     211                 :     /** Shared connect dispatch.
     212                 : 
     213                 :         Tries the connect syscall speculatively. On synchronous
     214                 :         completion, returns via inline budget or posts through queue.
     215                 :         On EINPROGRESS, registers with the reactor.
     216                 :     */
     217                 :     std::coroutine_handle<> do_connect(
     218                 :         std::coroutine_handle<>,
     219                 :         capy::executor_ref,
     220                 :         Endpoint const&,
     221                 :         std::stop_token const&,
     222                 :         std::error_code*);
     223                 : 
     224                 :     /** Shared scatter-read dispatch.
     225                 : 
     226                 :         Tries readv() speculatively. On success or hard error,
     227                 :         returns via inline budget or posts through queue.
     228                 :         On EAGAIN, registers with the reactor.
     229                 :     */
     230                 :     std::coroutine_handle<> do_read_some(
     231                 :         std::coroutine_handle<>,
     232                 :         capy::executor_ref,
     233                 :         buffer_param,
     234                 :         std::stop_token const&,
     235                 :         std::error_code*,
     236                 :         std::size_t*);
     237                 : 
     238                 :     /** Shared gather-write dispatch.
     239                 : 
     240                 :         Tries the write via WriteOp::write_policy speculatively.
     241                 :         On success or hard error, returns via inline budget or
     242                 :         posts through queue. On EAGAIN, registers with the reactor.
     243                 :     */
     244                 :     std::coroutine_handle<> do_write_some(
     245                 :         std::coroutine_handle<>,
     246                 :         capy::executor_ref,
     247                 :         buffer_param,
     248                 :         std::stop_token const&,
     249                 :         std::error_code*,
     250                 :         std::size_t*);
     251                 : 
     252                 :     /** Shared readiness-wait dispatch.
     253                 : 
     254                 :         Every wait type probes the descriptor with a zero-timeout
     255                 :         `poll()` and completes at once if the condition already
     256                 :         holds; otherwise the op re-probes under the descriptor mutex
     257                 :         and parks, completing when a reactor event arrives and a
     258                 :         fresh probe confirms the condition. A write wait therefore
     259                 :         completes only while a non-blocking write can make progress.
     260                 :     */
     261                 :     std::coroutine_handle<> do_wait(
     262                 :         std::coroutine_handle<>,
     263                 :         capy::executor_ref,
     264                 :         wait_type,
     265                 :         std::stop_token const&,
     266                 :         std::error_code*);
     267                 : 
     268                 :     /** Close the socket and cancel pending operations.
     269                 : 
     270                 :         Extends the base do_close_socket() to also reset
     271                 :         the remote endpoint.
     272                 :     */
     273           43218 :     void do_close_socket() noexcept
     274                 :     {
     275           43218 :         base_type::do_close_socket();
     276           43218 :         remote_endpoint_ = Endpoint{};
     277           43218 :     }
     278                 : 
     279                 :     /// Release ownership of the descriptor and drop the cached peer.
     280               8 :     native_handle_type do_release_socket() noexcept
     281                 :     {
     282               8 :         auto fd          = base_type::do_release_socket();
     283               8 :         remote_endpoint_ = Endpoint{};
     284               8 :         return fd;
     285                 :     }
     286                 : 
     287                 : private:
     288                 :     // CRTP callbacks for reactor_io_core cancel/close
     289                 : 
     290                 :     template<class Op>
     291             219 :     reactor_op_base** op_to_desc_slot(Op& op) noexcept
     292                 :     {
     293             219 :         if (&op == static_cast<void*>(&conn_))
     294 MIS           0 :             return &this->desc_state_.connect_op;
     295 HIT         219 :         if (&op == static_cast<void*>(&rd_))
     296             208 :             return &this->desc_state_.read_op;
     297              11 :         if (&op == static_cast<void*>(&wr_))
     298               6 :             return &this->desc_state_.write_op;
     299               5 :         if (&op == static_cast<void*>(&wait_rd_))
     300               3 :             return &this->desc_state_.wait_read_op;
     301               2 :         if (&op == static_cast<void*>(&wait_wr_))
     302 MIS           0 :             return &this->desc_state_.wait_write_op;
     303 HIT           2 :         if (&op == static_cast<void*>(&wait_er_))
     304               2 :             return &this->desc_state_.wait_error_op;
     305 MIS           0 :         return nullptr;
     306                 :     }
     307                 : 
     308                 :     template<class Fn>
     309 HIT       43452 :     void for_each_op(Fn fn) noexcept
     310                 :     {
     311           43452 :         fn(conn_);
     312           43452 :         fn(rd_);
     313           43452 :         fn(wr_);
     314           43452 :         fn(wait_rd_);
     315           43452 :         fn(wait_wr_);
     316           43452 :         fn(wait_er_);
     317           43452 :     }
     318                 : 
     319                 :     template<class Fn>
     320           48374 :     void for_each_desc_entry(Fn fn) noexcept
     321                 :     {
     322           48374 :         fn(conn_, this->desc_state_.connect_op);
     323           48374 :         fn(rd_, this->desc_state_.read_op);
     324           48374 :         fn(wr_, this->desc_state_.write_op);
     325           48374 :         fn(wait_rd_, this->desc_state_.wait_read_op);
     326           48374 :         fn(wait_wr_, this->desc_state_.wait_write_op);
     327           48374 :         fn(wait_er_, this->desc_state_.wait_error_op);
     328           48374 :     }
     329                 : };
     330                 : 
     331                 : template<
     332                 :     class Derived,
     333                 :     class Service,
     334                 :     class ConnOp,
     335                 :     class ReadOp,
     336                 :     class WriteOp,
     337                 :     class WaitOp,
     338                 :     class DescState,
     339                 :     class ImplBase,
     340                 :     class Endpoint>
     341                 : std::coroutine_handle<>
     342            4667 : reactor_stream_socket<
     343                 :     Derived,
     344                 :     Service,
     345                 :     ConnOp,
     346                 :     ReadOp,
     347                 :     WriteOp,
     348                 :     WaitOp,
     349                 :     DescState,
     350                 :     ImplBase,
     351                 :     Endpoint>::
     352                 :     do_connect(
     353                 :         std::coroutine_handle<> h,
     354                 :         capy::executor_ref ex,
     355                 :         Endpoint const& ep,
     356                 :         std::stop_token const& token,
     357                 :         std::error_code* ec)
     358                 : {
     359            4667 :     auto& op = conn_;
     360                 : 
     361            4667 :     sockaddr_storage storage{};
     362            4667 :     socklen_t addrlen = to_sockaddr(ep, socket_family(this->fd_), storage);
     363                 :     int result =
     364            4667 :         ::connect(this->fd_, reinterpret_cast<sockaddr*>(&storage), addrlen);
     365                 : 
     366            4667 :     if (result == 0)
     367                 :     {
     368              29 :         sockaddr_storage local_storage{};
     369              29 :         socklen_t local_len = sizeof(local_storage);
     370              29 :         if (::getsockname(
     371                 :                 this->fd_, reinterpret_cast<sockaddr*>(&local_storage),
     372              29 :                 &local_len) == 0)
     373 MIS           0 :             this->local_endpoint_ =
     374 HIT          29 :                 from_sockaddr_as(local_storage, local_len, Endpoint{});
     375              29 :         remote_endpoint_ = ep;
     376                 :     }
     377                 : 
     378            4667 :     if (result == 0 || errno != EINPROGRESS)
     379                 :     {
     380              37 :         int err = (result < 0) ? errno : 0;
     381              37 :         if (this->svc_.scheduler().try_consume_inline_budget())
     382                 :         {
     383 MIS           0 :             *ec       = err ? make_err(err) : std::error_code{};
     384               0 :             op.cont.h = h;
     385               0 :             return dispatch_coro(ex, op.cont);
     386                 :         }
     387 HIT          37 :         op.reset();
     388              37 :         op.h               = h;
     389              37 :         op.ex              = ex;
     390              37 :         op.ec_out          = ec;
     391              37 :         op.fd              = this->fd_;
     392              37 :         op.target_endpoint = ep;
     393              37 :         op.start(token, static_cast<Derived*>(this));
     394              37 :         op.impl_ptr = this->shared_from_this();
     395              37 :         op.complete(err, 0);
     396              37 :         this->svc_.post(&op);
     397              37 :         return std::noop_coroutine();
     398                 :     }
     399                 : 
     400                 :     // EINPROGRESS — register with reactor
     401            4630 :     op.reset();
     402            4630 :     op.h               = h;
     403            4630 :     op.ex              = ex;
     404            4630 :     op.ec_out          = ec;
     405            4630 :     op.fd              = this->fd_;
     406            4630 :     op.target_endpoint = ep;
     407            4630 :     op.start(token, static_cast<Derived*>(this));
     408            4630 :     op.impl_ptr = this->shared_from_this();
     409                 : 
     410            4630 :     this->register_op(
     411            4630 :         op, this->desc_state_.connect_op, this->desc_state_.write_ready, true);
     412            4630 :     return std::noop_coroutine();
     413                 : }
     414                 : 
     415                 : template<
     416                 :     class Derived,
     417                 :     class Service,
     418                 :     class ConnOp,
     419                 :     class ReadOp,
     420                 :     class WriteOp,
     421                 :     class WaitOp,
     422                 :     class DescState,
     423                 :     class ImplBase,
     424                 :     class Endpoint>
     425                 : std::coroutine_handle<>
     426          210443 : reactor_stream_socket<
     427                 :     Derived,
     428                 :     Service,
     429                 :     ConnOp,
     430                 :     ReadOp,
     431                 :     WriteOp,
     432                 :     WaitOp,
     433                 :     DescState,
     434                 :     ImplBase,
     435                 :     Endpoint>::
     436                 :     do_read_some(
     437                 :         std::coroutine_handle<> h,
     438                 :         capy::executor_ref ex,
     439                 :         buffer_param param,
     440                 :         std::stop_token const& token,
     441                 :         std::error_code* ec,
     442                 :         std::size_t* bytes_out)
     443                 : {
     444          210443 :     auto& op = rd_;
     445          210443 :     op.reset();
     446                 : 
     447                 :     // Closed-object contract: complete with bad_file_descriptor without
     448                 :     // touching the kernel or the unregistered descriptor state.
     449          210443 :     if (this->fd_ < 0)
     450                 :     {
     451               8 :         op.h         = h;
     452               8 :         op.ex        = ex;
     453               8 :         op.ec_out    = ec;
     454               8 :         op.bytes_out = bytes_out;
     455               8 :         op.start(token, static_cast<Derived*>(this));
     456               8 :         op.impl_ptr = this->shared_from_this();
     457               8 :         op.complete(EBADF, 0);
     458               8 :         this->svc_.post(&op);
     459               8 :         return std::noop_coroutine();
     460                 :     }
     461                 : 
     462          210435 :     capy::mutable_buffer bufs[ReadOp::max_buffers];
     463          210435 :     op.iovec_count = static_cast<int>(param.copy_to(bufs, ReadOp::max_buffers));
     464                 : 
     465          210435 :     if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
     466                 :     {
     467               4 :         op.empty_buffer_read = true;
     468               4 :         op.h                 = h;
     469               4 :         op.ex                = ex;
     470               4 :         op.ec_out            = ec;
     471               4 :         op.bytes_out         = bytes_out;
     472               4 :         op.start(token, static_cast<Derived*>(this));
     473               4 :         op.impl_ptr = this->shared_from_this();
     474               4 :         op.complete(0, 0);
     475               4 :         this->svc_.post(&op);
     476               4 :         return std::noop_coroutine();
     477                 :     }
     478                 : 
     479          420880 :     for (int i = 0; i < op.iovec_count; ++i)
     480                 :     {
     481          210449 :         op.iovecs[i].iov_base = bufs[i].data();
     482          210449 :         op.iovecs[i].iov_len  = bufs[i].size();
     483                 :     }
     484                 : 
     485                 :     // Speculative read; for the single-buffer case use recv() so the
     486                 :     // kernel skips the readv iov_iter setup.
     487                 :     ssize_t n;
     488          210431 :     if (op.iovec_count == 1)
     489                 :     {
     490                 :         do
     491                 :         {
     492          210419 :             n = ::recv(this->fd_, bufs[0].data(), bufs[0].size(), 0);
     493                 :         }
     494          210419 :         while (n < 0 && errno == EINTR);
     495                 :     }
     496                 :     else
     497                 :     {
     498                 :         do
     499                 :         {
     500              16 :             n = ::readv(this->fd_, op.iovecs, op.iovec_count);
     501                 :         }
     502              16 :         while (n < 0 && errno == EINTR);
     503                 :     }
     504                 : 
     505          210431 :     if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
     506                 :     {
     507          209650 :         int err    = (n < 0) ? errno : 0;
     508          209650 :         auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
     509                 : 
     510          209650 :         if (this->svc_.scheduler().try_consume_inline_budget())
     511                 :         {
     512          167752 :             if (err)
     513               8 :                 *ec = make_err(err);
     514          167744 :             else if (n == 0)
     515              17 :                 *ec = capy::error::eof;
     516                 :             else
     517          167727 :                 *ec = {};
     518          167752 :             *bytes_out = bytes;
     519          167752 :             op.cont.h  = h;
     520          167752 :             return dispatch_coro(ex, op.cont);
     521                 :         }
     522           41898 :         op.h         = h;
     523           41898 :         op.ex        = ex;
     524           41898 :         op.ec_out    = ec;
     525           41898 :         op.bytes_out = bytes_out;
     526           41898 :         op.start(token, static_cast<Derived*>(this));
     527           41898 :         op.impl_ptr = this->shared_from_this();
     528           41898 :         op.complete(err, bytes);
     529           41898 :         this->svc_.post(&op);
     530           41898 :         return std::noop_coroutine();
     531                 :     }
     532                 : 
     533                 :     // EAGAIN — register with reactor
     534             781 :     op.h         = h;
     535             781 :     op.ex        = ex;
     536             781 :     op.ec_out    = ec;
     537             781 :     op.bytes_out = bytes_out;
     538             781 :     op.fd        = this->fd_;
     539             781 :     op.start(token, static_cast<Derived*>(this));
     540             781 :     op.impl_ptr = this->shared_from_this();
     541                 : 
     542             781 :     this->register_op(
     543             781 :         op, this->desc_state_.read_op, this->desc_state_.read_ready);
     544             781 :     return std::noop_coroutine();
     545                 : }
     546                 : 
     547                 : template<
     548                 :     class Derived,
     549                 :     class Service,
     550                 :     class ConnOp,
     551                 :     class ReadOp,
     552                 :     class WriteOp,
     553                 :     class WaitOp,
     554                 :     class DescState,
     555                 :     class ImplBase,
     556                 :     class Endpoint>
     557                 : std::coroutine_handle<>
     558          209742 : reactor_stream_socket<
     559                 :     Derived,
     560                 :     Service,
     561                 :     ConnOp,
     562                 :     ReadOp,
     563                 :     WriteOp,
     564                 :     WaitOp,
     565                 :     DescState,
     566                 :     ImplBase,
     567                 :     Endpoint>::
     568                 :     do_write_some(
     569                 :         std::coroutine_handle<> h,
     570                 :         capy::executor_ref ex,
     571                 :         buffer_param param,
     572                 :         std::stop_token const& token,
     573                 :         std::error_code* ec,
     574                 :         std::size_t* bytes_out)
     575                 : {
     576          209742 :     auto& op = wr_;
     577          209742 :     op.reset();
     578                 : 
     579                 :     // Closed-object contract: complete with bad_file_descriptor without
     580                 :     // touching the kernel or the unregistered descriptor state.
     581          209742 :     if (this->fd_ < 0)
     582                 :     {
     583              10 :         op.h         = h;
     584              10 :         op.ex        = ex;
     585              10 :         op.ec_out    = ec;
     586              10 :         op.bytes_out = bytes_out;
     587              10 :         op.start(token, static_cast<Derived*>(this));
     588              10 :         op.impl_ptr = this->shared_from_this();
     589              10 :         op.complete(EBADF, 0);
     590              10 :         this->svc_.post(&op);
     591              10 :         return std::noop_coroutine();
     592                 :     }
     593                 : 
     594          209732 :     capy::mutable_buffer bufs[WriteOp::max_buffers];
     595          209732 :     op.iovec_count =
     596          209732 :         static_cast<int>(param.copy_to(bufs, WriteOp::max_buffers));
     597                 : 
     598          209732 :     if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
     599                 :     {
     600               4 :         op.h         = h;
     601               4 :         op.ex        = ex;
     602               4 :         op.ec_out    = ec;
     603               4 :         op.bytes_out = bytes_out;
     604               4 :         op.start(token, static_cast<Derived*>(this));
     605               4 :         op.impl_ptr = this->shared_from_this();
     606               4 :         op.complete(0, 0);
     607               4 :         this->svc_.post(&op);
     608               4 :         return std::noop_coroutine();
     609                 :     }
     610                 : 
     611          419472 :     for (int i = 0; i < op.iovec_count; ++i)
     612                 :     {
     613          209744 :         op.iovecs[i].iov_base = bufs[i].data();
     614          209744 :         op.iovecs[i].iov_len  = bufs[i].size();
     615                 :     }
     616                 : 
     617                 :     // Speculative write; the single-buffer case dispatches to a
     618                 :     // backend-specific fast path so the kernel skips msghdr/iov_iter
     619                 :     // setup (and so each backend can pick the right SIGPIPE strategy).
     620                 :     ssize_t n;
     621          209728 :     if (op.iovec_count == 1)
     622                 :     {
     623          419432 :         n = WriteOp::write_policy::write_one(
     624          209716 :             this->fd_, bufs[0].data(), bufs[0].size());
     625                 :     }
     626                 :     else
     627                 :     {
     628              12 :         n = WriteOp::write_policy::write(this->fd_, op.iovecs, op.iovec_count);
     629                 :     }
     630                 : 
     631          209728 :     if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
     632                 :     {
     633          209584 :         int err    = (n < 0) ? errno : 0;
     634          209584 :         auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
     635                 : 
     636          209584 :         if (this->svc_.scheduler().try_consume_inline_budget())
     637                 :         {
     638          167611 :             *ec        = err ? make_err(err) : std::error_code{};
     639          167611 :             *bytes_out = bytes;
     640          167611 :             op.cont.h  = h;
     641          167611 :             return dispatch_coro(ex, op.cont);
     642                 :         }
     643           41973 :         op.h         = h;
     644           41973 :         op.ex        = ex;
     645           41973 :         op.ec_out    = ec;
     646           41973 :         op.bytes_out = bytes_out;
     647           41973 :         op.start(token, static_cast<Derived*>(this));
     648           41973 :         op.impl_ptr = this->shared_from_this();
     649           41973 :         op.complete(err, bytes);
     650           41973 :         this->svc_.post(&op);
     651           41973 :         return std::noop_coroutine();
     652                 :     }
     653                 : 
     654                 :     // EAGAIN — register with reactor
     655             144 :     op.h         = h;
     656             144 :     op.ex        = ex;
     657             144 :     op.ec_out    = ec;
     658             144 :     op.bytes_out = bytes_out;
     659             144 :     op.fd        = this->fd_;
     660             144 :     op.start(token, static_cast<Derived*>(this));
     661             144 :     op.impl_ptr = this->shared_from_this();
     662                 : 
     663             144 :     this->register_op(
     664             144 :         op, this->desc_state_.write_op, this->desc_state_.write_ready, true);
     665             144 :     return std::noop_coroutine();
     666                 : }
     667                 : 
     668                 : template<
     669                 :     class Derived,
     670                 :     class Service,
     671                 :     class ConnOp,
     672                 :     class ReadOp,
     673                 :     class WriteOp,
     674                 :     class WaitOp,
     675                 :     class DescState,
     676                 :     class ImplBase,
     677                 :     class Endpoint>
     678                 : std::coroutine_handle<>
     679              96 : reactor_stream_socket<
     680                 :     Derived,
     681                 :     Service,
     682                 :     ConnOp,
     683                 :     ReadOp,
     684                 :     WriteOp,
     685                 :     WaitOp,
     686                 :     DescState,
     687                 :     ImplBase,
     688                 :     Endpoint>::
     689                 :     do_wait(
     690                 :         std::coroutine_handle<> h,
     691                 :         capy::executor_ref ex,
     692                 :         wait_type w,
     693                 :         std::stop_token const& token,
     694                 :         std::error_code* ec)
     695                 : {
     696                 :     // Pick refs up-front to avoid duplicating the register_op call.
     697                 :     WaitOp* op_ptr;
     698                 :     reactor_op_base** desc_slot_ptr;
     699                 :     std::uint32_t event;
     700                 : 
     701              96 :     if (w == wait_type::read)
     702                 :     {
     703              55 :         op_ptr        = &wait_rd_;
     704              55 :         desc_slot_ptr = &this->desc_state_.wait_read_op;
     705              55 :         event         = reactor_event_read;
     706                 :     }
     707              41 :     else if (w == wait_type::write)
     708                 :     {
     709              23 :         op_ptr        = &wait_wr_;
     710              23 :         desc_slot_ptr = &this->desc_state_.wait_write_op;
     711              23 :         event         = reactor_event_write;
     712                 :     }
     713                 :     else // wait_type::error
     714                 :     {
     715              18 :         op_ptr        = &wait_er_;
     716              18 :         desc_slot_ptr = &this->desc_state_.wait_error_op;
     717              18 :         event         = reactor_event_error;
     718                 :     }
     719                 : 
     720              96 :     auto& op = *op_ptr;
     721                 : 
     722                 :     // Speculative probe, mirroring the speculative read: an
     723                 :     // edge-triggered reactor cannot report a condition that already
     724                 :     // holds, so a wait initiated on an already-ready socket would
     725                 :     // otherwise park forever.
     726              96 :     int perr = 0;
     727              96 :     if (WaitOp::probe(this->fd_, event, perr))
     728                 :     {
     729              40 :         if (this->svc_.scheduler().try_consume_inline_budget())
     730                 :         {
     731               8 :             *ec       = perr ? make_err(perr) : std::error_code{};
     732               8 :             op.cont.h = h;
     733               8 :             return dispatch_coro(ex, op.cont);
     734                 :         }
     735              32 :         op.reset();
     736              32 :         op.wait_event = event;
     737              32 :         op.h          = h;
     738              32 :         op.ex         = ex;
     739              32 :         op.ec_out     = ec;
     740              32 :         op.fd         = this->fd_;
     741              32 :         op.start(token, static_cast<Derived*>(this));
     742              32 :         op.impl_ptr = this->shared_from_this();
     743              32 :         op.complete(perr, 0);
     744              32 :         this->svc_.post(&op);
     745              32 :         return std::noop_coroutine();
     746                 :     }
     747                 : 
     748              56 :     op.reset();
     749              56 :     op.wait_event = event;
     750              56 :     op.h          = h;
     751              56 :     op.ex         = ex;
     752              56 :     op.ec_out     = ec;
     753              56 :     op.fd         = this->fd_;
     754              56 :     op.start(token, static_cast<Derived*>(this));
     755              56 :     op.impl_ptr = this->shared_from_this();
     756                 : 
     757                 :     // Force register_op's ready path so the wait op re-probes under
     758                 :     // the descriptor mutex before parking. An edge consumed between
     759                 :     // the speculative probe above and the park (a concurrent short
     760                 :     // read, or an error event dispatched to an empty slot) would
     761                 :     // otherwise leave the wait parked on a ready socket.
     762              56 :     bool force_probe = true;
     763              56 :     this->register_op(
     764                 :         op, *desc_slot_ptr, force_probe, event == reactor_event_write);
     765              56 :     return std::noop_coroutine();
     766                 : }
     767                 : 
     768                 : } // namespace boost::corosio::detail
     769                 : 
     770                 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_STREAM_SOCKET_HPP
        

Generated by: LCOV version 2.3