LCOV - code coverage report
Current view: top level - corosio/native/detail/reactor - reactor_descriptor.hpp (source / functions) Coverage Total Hit Missed
Test: coverage_remapped.info Lines: 82.3 % 260 214 46
Test Date: 2026-10-08 17:58:26 Functions: 92.9 % 56 52 4

           TLA  Line data    Source code
       1                 : //
       2                 : // Copyright (c) 2026 Michael Vandeberg
       3                 : //
       4                 : // Distributed under the Boost Software License, Version 1.0. (See accompanying
       5                 : // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
       6                 : //
       7                 : // Official repository: https://github.com/cppalliance/corosio
       8                 : //
       9                 : 
      10                 : #ifndef BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_HPP
      11                 : #define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_HPP
      12                 : 
      13                 : #include <boost/corosio/detail/platform.hpp>
      14                 : 
      15                 : #if BOOST_COROSIO_POSIX
      16                 : 
      17                 : #include <boost/corosio/posix_stream_descriptor.hpp>
      18                 : #include <boost/corosio/wait_type.hpp>
      19                 : #include <boost/corosio/detail/dispatch_coro.hpp>
      20                 : #include <boost/corosio/detail/intrusive.hpp>
      21                 : #include <boost/corosio/detail/native_handle.hpp>
      22                 : #include <boost/corosio/native/detail/make_err.hpp>
      23                 : #include <boost/corosio/native/detail/validate_fd.hpp>
      24                 : #include <boost/corosio/native/detail/reactor/reactor_descriptor_state.hpp>
      25                 : #include <boost/corosio/native/detail/reactor/reactor_io_core.hpp>
      26                 : #include <boost/corosio/native/detail/reactor/reactor_op.hpp>
      27                 : #include <boost/corosio/native/detail/reactor/reactor_op_complete.hpp>
      28                 : #include <boost/capy/buffers.hpp>
      29                 : 
      30                 : #include <atomic>
      31                 : #include <coroutine>
      32                 : #include <memory>
      33                 : #include <mutex>
      34                 : #include <utility>
      35                 : 
      36                 : #include <errno.h>
      37                 : #include <sys/stat.h>
      38                 : #include <sys/uio.h>
      39                 : #include <unistd.h>
      40                 : 
      41                 : /* Reactor-backed implementation of posix_stream_descriptor.
      42                 : 
      43                 :    The register/park/cancel/teardown protocol lives in reactor_io_core.
      44                 : 
      45                 :    The one behavior that is genuinely new: O_NONBLOCK is armed lazily,
      46                 :    on the first read_some/write_some and never from assign() or wait().
      47                 :    The flag lives on the shared open file description, so arming it is
      48                 :    visible to every other holder of that description -- which is why a
      49                 :    wait()-only user must never trigger it.
      50                 : */
      51                 : 
      52                 : namespace boost::corosio::detail {
      53                 : 
      54                 : // ============================================================
      55                 : // Op types
      56                 : // ============================================================
      57                 : 
      58                 : /* Descriptor op family.
      59                 : 
      60                 :    Mirrors reactor_stream_ops.hpp. The Acceptor parameter is a
      61                 :    placeholder: reactor_op is parameterized on both a socket and an
      62                 :    acceptor impl type, and the shared completion helpers name
      63                 :    acceptor_impl_ in a branch that a descriptor op never takes but the
      64                 :    compiler still instantiates. Passing the backend's acceptor type (as
      65                 :    reactor_dgram_socket_impl already does) keeps those helpers shared
      66                 :    rather than duplicated here.
      67                 : 
      68                 :    @tparam Traits     Backend traits (epoll_traits, kqueue_traits, ...).
      69                 :    @tparam Descriptor The concrete descriptor type (forward-declared).
      70                 :    @tparam Acceptor   Placeholder acceptor type for the op base.
      71                 : */
      72                 : 
      73                 : template<class Traits, class Descriptor, class Acceptor>
      74                 : struct reactor_descriptor_base_op : reactor_op<Descriptor, Acceptor>
      75                 : {
      76                 :     void operator()() override;
      77                 :     void cancel() noexcept override;
      78                 : };
      79                 : 
      80                 : template<class Traits, class Descriptor, class Acceptor>
      81                 : struct reactor_descriptor_read_op final
      82                 :     : reactor_read_op<reactor_descriptor_base_op<Traits, Descriptor, Acceptor>>
      83                 : {};
      84                 : 
      85                 : template<class Traits, class Descriptor, class Acceptor>
      86                 : struct reactor_descriptor_write_op final
      87                 :     : reactor_write_op<
      88                 :           reactor_descriptor_base_op<Traits, Descriptor, Acceptor>,
      89                 :           typename Traits::descriptor_write_policy>
      90                 : {};
      91                 : 
      92                 : template<class Traits, class Descriptor, class Acceptor>
      93                 : struct reactor_descriptor_wait_op final
      94                 :     : reactor_wait_op<reactor_descriptor_base_op<Traits, Descriptor, Acceptor>>
      95                 : {
      96                 :     using base_type = reactor_wait_op<
      97                 :         reactor_descriptor_base_op<Traits, Descriptor, Acceptor>>;
      98                 : 
      99                 :     /** Probe like reactor_wait_op::probe, but name a code for an error
     100                 :         wait that sees POLLERR or POLLHUP. A non-socket has no SO_ERROR
     101                 :         to consult, and success would claim no error condition exists.
     102                 :     */
     103 HIT          72 :     static bool probe(int fd, std::uint32_t event, int& err) noexcept
     104                 :     {
     105              72 :         if (event != reactor_event_error || fd < 0)
     106              48 :             return base_type::probe(fd, event, err);
     107                 : 
     108              24 :         pollfd pfd{};
     109              24 :         pfd.fd     = fd;
     110              24 :         pfd.events = POLLPRI;
     111                 :         int r;
     112                 :         do
     113                 :         {
     114              24 :             r = ::poll(&pfd, 1, 0);
     115                 :         }
     116              24 :         while (r < 0 && errno == EINTR);
     117              24 :         if (r < 0)
     118                 :         {
     119 MIS           0 :             err = (errno == EAGAIN || errno == EWOULDBLOCK) ? ENOMEM : errno;
     120               0 :             return true;
     121                 :         }
     122 HIT          24 :         if (r == 0)
     123              22 :             return false;
     124               2 :         if (pfd.revents & POLLNVAL)
     125 MIS           0 :             err = EBADF;
     126 HIT           2 :         else if (pfd.revents & (POLLERR | POLLHUP))
     127               2 :             err = EIO;
     128 MIS           0 :         else if (is_fifo(fd))
     129                 :             // Darwin reports POLLPRI on a pipe that merely holds data;
     130                 :             // a FIFO has no exceptional condition to signal.
     131               0 :             return false;
     132 HIT           2 :         return true;
     133                 :     }
     134                 : 
     135 MIS           0 :     static bool is_fifo(int fd) noexcept
     136                 :     {
     137                 :         struct stat st;
     138               0 :         return ::fstat(fd, &st) == 0 && S_ISFIFO(st.st_mode);
     139                 :     }
     140                 : 
     141 HIT          35 :     void perform_io() noexcept override
     142                 :     {
     143              35 :         int err = 0;
     144              35 :         if (probe(this->fd, this->wait_event, err))
     145              10 :             this->complete(err, 0);
     146                 :         else
     147              25 :             this->complete(EAGAIN, 0);
     148              35 :     }
     149                 : 
     150                 :     void operator()() override;
     151                 : };
     152                 : 
     153                 : // --- Deferred implementations (instantiated when Descriptor is complete) ---
     154                 : 
     155                 : template<class Traits, class Descriptor, class Acceptor>
     156                 : void
     157              55 : reactor_descriptor_base_op<Traits, Descriptor, Acceptor>::operator()()
     158                 : {
     159              55 :     complete_io_op(*this);
     160              55 : }
     161                 : 
     162                 : template<class Traits, class Descriptor, class Acceptor>
     163                 : void
     164               6 : reactor_descriptor_base_op<Traits, Descriptor, Acceptor>::cancel() noexcept
     165                 : {
     166                 :     // A descriptor op is only ever started against a descriptor impl, so
     167                 :     // the acceptor arm of the stream op's cancel() has no counterpart.
     168               6 :     if (this->socket_impl_)
     169               6 :         this->socket_impl_->cancel_single_op(*this);
     170                 :     else
     171 MIS           0 :         this->request_cancel();
     172 HIT           6 : }
     173                 : 
     174                 : template<class Traits, class Descriptor, class Acceptor>
     175                 : void
     176              27 : reactor_descriptor_wait_op<Traits, Descriptor, Acceptor>::operator()()
     177                 : {
     178              27 :     complete_wait_op(*this);
     179              27 : }
     180                 : 
     181                 : // ============================================================
     182                 : // Descriptor implementation
     183                 : // ============================================================
     184                 : 
     185                 : /** CRTP base for reactor-backed posix_stream_descriptor implementations.
     186                 : 
     187                 :     Holds the adopted descriptor, its reactor registration state, and
     188                 :     the five op slots (read, write, and one wait per direction).
     189                 : 
     190                 :     @tparam Derived  The named final class (CRTP self).
     191                 :     @tparam Traits   Backend traits (epoll_traits, kqueue_traits, ...).
     192                 :     @tparam Service  The backend's descriptor service type.
     193                 :     @tparam Acceptor Placeholder acceptor type for the op base.
     194                 : */
     195                 : template<class Derived, class Traits, class Service, class Acceptor>
     196                 : class reactor_descriptor
     197                 :     : public posix_stream_descriptor::implementation
     198                 :     , public std::enable_shared_from_this<Derived>
     199                 :     , public reactor_io_core<Derived, Service, typename Traits::desc_state_type>
     200                 :     , public intrusive_list<Derived>::node
     201                 : {
     202                 :     using base_op  = reactor_descriptor_base_op<Traits, Derived, Acceptor>;
     203                 :     using read_op  = reactor_descriptor_read_op<Traits, Derived, Acceptor>;
     204                 :     using write_op = reactor_descriptor_write_op<Traits, Derived, Acceptor>;
     205                 :     using wait_op  = reactor_descriptor_wait_op<Traits, Derived, Acceptor>;
     206                 : 
     207                 :     using core_type =
     208                 :         reactor_io_core<Derived, Service, typename Traits::desc_state_type>;
     209                 :     friend core_type;
     210                 : 
     211                 : protected:
     212                 :     // NOLINTNEXTLINE(bugprone-crtp-constructor-accessibility)
     213             121 :     explicit reactor_descriptor(Service& svc) noexcept : core_type(svc) {}
     214                 : 
     215                 :     using core_type::svc_;
     216                 : 
     217                 : public:
     218             121 :     ~reactor_descriptor() override = default;
     219                 : 
     220                 :     using core_type::desc_state_;
     221                 : 
     222                 :     // --- Virtual method overrides ---
     223                 : 
     224              49 :     std::coroutine_handle<> read_some(
     225                 :         std::coroutine_handle<> h,
     226                 :         capy::executor_ref ex,
     227                 :         buffer_param param,
     228                 :         std::stop_token token,
     229                 :         std::error_code* ec,
     230                 :         std::size_t* bytes_out) override
     231                 :     {
     232              49 :         return do_read_some(h, ex, param, token, ec, bytes_out);
     233                 :     }
     234                 : 
     235              20 :     std::coroutine_handle<> write_some(
     236                 :         std::coroutine_handle<> h,
     237                 :         capy::executor_ref ex,
     238                 :         buffer_param param,
     239                 :         std::stop_token token,
     240                 :         std::error_code* ec,
     241                 :         std::size_t* bytes_out) override
     242                 :     {
     243              20 :         return do_write_some(h, ex, param, token, ec, bytes_out);
     244                 :     }
     245                 : 
     246              37 :     std::coroutine_handle<> wait(
     247                 :         std::coroutine_handle<> h,
     248                 :         capy::executor_ref ex,
     249                 :         wait_type w,
     250                 :         std::stop_token token,
     251                 :         std::error_code* ec) override
     252                 :     {
     253              37 :         return do_wait(h, ex, w, token, ec);
     254                 :     }
     255                 : 
     256             262 :     native_handle_type native_handle() const noexcept override
     257                 :     {
     258             262 :         return fd_;
     259                 :     }
     260                 : 
     261                 :     native_handle_type release_descriptor() noexcept override;
     262                 : 
     263               7 :     void cancel() noexcept override
     264                 :     {
     265               7 :         this->cancel_all();
     266               7 :     }
     267                 : 
     268                 :     // --- Service-facing (non-virtual) ---
     269                 : 
     270                 :     /** Adopt the fd, initialize descriptor state, and register it.
     271                 : 
     272                 :         @param fd The descriptor to adopt.
     273                 : 
     274                 :         @return The error if the reactor rejects the descriptor, in
     275                 :         which case the implementation is left closed and the caller
     276                 :         retains ownership of @a fd; otherwise a default constructed
     277                 :         error code.
     278                 :     */
     279                 :     std::error_code init_and_register(int fd) noexcept;
     280                 : 
     281                 :     /// Close the descriptor and cancel pending operations.
     282                 :     void close_descriptor() noexcept;
     283                 : 
     284                 : private:
     285                 :     /** Arm O_NONBLOCK, once, before the first speculative syscall.
     286                 : 
     287                 :         Reports an errno rather than an error_code because that is what
     288                 :         the op result model records; the round trip is lossless here
     289                 :         because fcntl only fails with codes make_err passes through.
     290                 :     */
     291              67 :     int arm_nonblocking() noexcept
     292                 :     {
     293                 :         // Relaxed: ensure_nonblocking is idempotent, so a duplicate
     294                 :         // fcntl from a racing first read and write is harmless.
     295              67 :         if (nonblocking_.load(std::memory_order_relaxed))
     296               6 :             return 0;
     297              61 :         if (auto ec = ensure_nonblocking(fd_))
     298 MIS           0 :             return ec.value();
     299 HIT          61 :         nonblocking_.store(true, std::memory_order_relaxed);
     300              61 :         return 0;
     301                 :     }
     302                 : 
     303                 :     std::coroutine_handle<> do_read_some(
     304                 :         std::coroutine_handle<>,
     305                 :         capy::executor_ref,
     306                 :         buffer_param,
     307                 :         std::stop_token const&,
     308                 :         std::error_code*,
     309                 :         std::size_t*);
     310                 : 
     311                 :     std::coroutine_handle<> do_write_some(
     312                 :         std::coroutine_handle<>,
     313                 :         capy::executor_ref,
     314                 :         buffer_param,
     315                 :         std::stop_token const&,
     316                 :         std::error_code*,
     317                 :         std::size_t*);
     318                 : 
     319                 :     std::coroutine_handle<> do_wait(
     320                 :         std::coroutine_handle<>,
     321                 :         capy::executor_ref,
     322                 :         wait_type,
     323                 :         std::stop_token const&,
     324                 :         std::error_code*);
     325                 : 
     326                 :     /// Apply @a fn to each of the five op slots.
     327                 :     template<class Fn>
     328             327 :     void for_each_op(Fn fn) noexcept
     329                 :     {
     330             327 :         fn(rd_);
     331             327 :         fn(wr_);
     332             327 :         fn(wait_rd_);
     333             327 :         fn(wait_wr_);
     334             327 :         fn(wait_er_);
     335             327 :     }
     336                 : 
     337                 :     /// Apply @a fn to each op and the descriptor_state slot it parks in.
     338                 :     template<class Fn>
     339             433 :     void for_each_desc_entry(Fn fn) noexcept
     340                 :     {
     341             433 :         fn(rd_, desc_state_.read_op);
     342             433 :         fn(wr_, desc_state_.write_op);
     343             433 :         fn(wait_rd_, desc_state_.wait_read_op);
     344             433 :         fn(wait_wr_, desc_state_.wait_write_op);
     345             433 :         fn(wait_er_, desc_state_.wait_error_op);
     346             433 :     }
     347                 : 
     348                 :     /// Sweep every op slot, then drop the reactor registration.
     349                 :     void quiesce() noexcept;
     350                 : 
     351                 :     reactor_op_base** op_to_desc_slot(base_op& op) noexcept;
     352                 : 
     353                 :     int fd_           = -1;
     354                 :     std::atomic<bool> nonblocking_{false};
     355                 : 
     356                 :     read_op rd_;
     357                 :     write_op wr_;
     358                 :     wait_op wait_rd_;
     359                 :     wait_op wait_wr_;
     360                 :     wait_op wait_er_;
     361                 : };
     362                 : 
     363                 : // ============================================================
     364                 : // Registration and teardown
     365                 : // ============================================================
     366                 : 
     367                 : template<class Derived, class Traits, class Service, class Acceptor>
     368                 : std::error_code
     369             106 : reactor_descriptor<Derived, Traits, Service, Acceptor>::init_and_register(
     370                 :     int fd) noexcept
     371                 : {
     372             106 :     fd_ = fd;
     373             106 :     if (auto ec = this->register_fd(fd))
     374                 :     {
     375 MIS           0 :         fd_ = -1;
     376               0 :         return ec;
     377                 :     }
     378 HIT         106 :     return {};
     379                 : }
     380                 : 
     381                 : template<class Derived, class Traits, class Service, class Acceptor>
     382                 : void
     383             320 : reactor_descriptor<Derived, Traits, Service, Acceptor>::quiesce() noexcept
     384                 : {
     385             320 :     this->abandon_all();
     386             320 :     this->unregister_fd(fd_);
     387                 :     // The next adopted fd starts from an unknown flag state.
     388             320 :     nonblocking_.store(false, std::memory_order_relaxed);
     389             320 : }
     390                 : 
     391                 : template<class Derived, class Traits, class Service, class Acceptor>
     392                 : void
     393             316 : reactor_descriptor<Derived, Traits, Service, Acceptor>::
     394                 :     close_descriptor() noexcept
     395                 : {
     396             316 :     quiesce();
     397                 : 
     398             316 :     if (fd_ >= 0)
     399                 :     {
     400             102 :         ::close(fd_);
     401             102 :         fd_ = -1;
     402                 :     }
     403             316 : }
     404                 : 
     405                 : template<class Derived, class Traits, class Service, class Acceptor>
     406                 : native_handle_type
     407               4 : reactor_descriptor<Derived, Traits, Service, Acceptor>::
     408                 :     release_descriptor() noexcept
     409                 : {
     410               4 :     quiesce();
     411                 : 
     412                 :     // Do NOT close -- the caller takes ownership.
     413               4 :     native_handle_type released = fd_;
     414               4 :     fd_                         = -1;
     415               4 :     return released;
     416                 : }
     417                 : 
     418                 : // ============================================================
     419                 : // Op slot lookup
     420                 : // ============================================================
     421                 : 
     422                 : template<class Derived, class Traits, class Service, class Acceptor>
     423                 : reactor_op_base**
     424               6 : reactor_descriptor<Derived, Traits, Service, Acceptor>::op_to_desc_slot(
     425                 :     base_op& op) noexcept
     426                 : {
     427               6 :     if (&op == static_cast<void*>(&rd_))
     428               6 :         return &desc_state_.read_op;
     429 MIS           0 :     if (&op == static_cast<void*>(&wr_))
     430               0 :         return &desc_state_.write_op;
     431               0 :     if (&op == static_cast<void*>(&wait_rd_))
     432               0 :         return &desc_state_.wait_read_op;
     433               0 :     if (&op == static_cast<void*>(&wait_wr_))
     434               0 :         return &desc_state_.wait_write_op;
     435               0 :     if (&op == static_cast<void*>(&wait_er_))
     436               0 :         return &desc_state_.wait_error_op;
     437               0 :     return nullptr;
     438                 : }
     439                 : 
     440                 : // ============================================================
     441                 : // I/O dispatch
     442                 : // ============================================================
     443                 : 
     444                 : template<class Derived, class Traits, class Service, class Acceptor>
     445                 : std::coroutine_handle<>
     446 HIT          49 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_read_some(
     447                 :     std::coroutine_handle<> h,
     448                 :     capy::executor_ref ex,
     449                 :     buffer_param param,
     450                 :     std::stop_token const& token,
     451                 :     std::error_code* ec,
     452                 :     std::size_t* bytes_out)
     453                 : {
     454              49 :     auto& op = rd_;
     455              49 :     op.reset();
     456              49 :     op.h         = h;
     457              49 :     op.ex        = ex;
     458              49 :     op.ec_out    = ec;
     459              49 :     op.bytes_out = bytes_out;
     460                 : 
     461                 :     // Closed-object contract: complete with bad_file_descriptor without
     462                 :     // touching the kernel or the unregistered descriptor state.
     463              49 :     if (fd_ < 0)
     464                 :     {
     465 MIS           0 :         op.start(token, static_cast<Derived*>(this));
     466               0 :         op.impl_ptr = this->shared_from_this();
     467               0 :         op.complete(EBADF, 0);
     468               0 :         svc_.post(&op);
     469               0 :         return std::noop_coroutine();
     470                 :     }
     471                 : 
     472 HIT          49 :     capy::mutable_buffer bufs[read_op::max_buffers];
     473              49 :     op.iovec_count =
     474              49 :         static_cast<int>(param.copy_to(bufs, read_op::max_buffers));
     475                 : 
     476              49 :     if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
     477                 :     {
     478               2 :         op.empty_buffer_read = true;
     479               2 :         op.start(token, static_cast<Derived*>(this));
     480               2 :         op.impl_ptr = this->shared_from_this();
     481               2 :         op.complete(0, 0);
     482               2 :         svc_.post(&op);
     483               2 :         return std::noop_coroutine();
     484                 :     }
     485                 : 
     486                 :     // The first transferring operation is what arms O_NONBLOCK; assign()
     487                 :     // and wait() never do.
     488              47 :     if (int const nerr = arm_nonblocking())
     489                 :     {
     490 MIS           0 :         op.start(token, static_cast<Derived*>(this));
     491               0 :         op.impl_ptr = this->shared_from_this();
     492               0 :         op.complete(nerr, 0);
     493               0 :         svc_.post(&op);
     494               0 :         return std::noop_coroutine();
     495                 :     }
     496                 : 
     497 HIT          96 :     for (int i = 0; i < op.iovec_count; ++i)
     498                 :     {
     499              49 :         op.iovecs[i].iov_base = bufs[i].data();
     500              49 :         op.iovecs[i].iov_len  = bufs[i].size();
     501                 :     }
     502                 : 
     503                 :     // Speculative read; the single-buffer case uses read() so the kernel
     504                 :     // skips the readv iov_iter setup.
     505                 :     ssize_t n;
     506              47 :     if (op.iovec_count == 1)
     507                 :     {
     508                 :         do
     509                 :         {
     510              45 :             n = ::read(fd_, bufs[0].data(), bufs[0].size());
     511                 :         }
     512              45 :         while (n < 0 && errno == EINTR);
     513                 :     }
     514                 :     else
     515                 :     {
     516                 :         do
     517                 :         {
     518               2 :             n = ::readv(fd_, op.iovecs, op.iovec_count);
     519                 :         }
     520               2 :         while (n < 0 && errno == EINTR);
     521                 :     }
     522                 : 
     523              47 :     if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
     524                 :     {
     525              23 :         int err    = (n < 0) ? errno : 0;
     526              23 :         auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
     527                 : 
     528              23 :         if (svc_.scheduler().try_consume_inline_budget())
     529                 :         {
     530               6 :             if (err)
     531 MIS           0 :                 *ec = make_err(err);
     532 HIT           6 :             else if (n == 0)
     533               2 :                 *ec = capy::error::eof;
     534                 :             else
     535               4 :                 *ec = {};
     536               6 :             *bytes_out = bytes;
     537               6 :             op.cont.h  = h;
     538               6 :             return dispatch_coro(ex, op.cont);
     539                 :         }
     540              17 :         op.start(token, static_cast<Derived*>(this));
     541              17 :         op.impl_ptr = this->shared_from_this();
     542              17 :         op.complete(err, bytes);
     543              17 :         svc_.post(&op);
     544              17 :         return std::noop_coroutine();
     545                 :     }
     546                 : 
     547                 :     // EAGAIN — register with reactor
     548              24 :     op.fd = fd_;
     549              24 :     op.start(token, static_cast<Derived*>(this));
     550              24 :     op.impl_ptr = this->shared_from_this();
     551                 : 
     552              24 :     this->register_op(op, desc_state_.read_op, desc_state_.read_ready);
     553              24 :     return std::noop_coroutine();
     554                 : }
     555                 : 
     556                 : template<class Derived, class Traits, class Service, class Acceptor>
     557                 : std::coroutine_handle<>
     558              20 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_write_some(
     559                 :     std::coroutine_handle<> h,
     560                 :     capy::executor_ref ex,
     561                 :     buffer_param param,
     562                 :     std::stop_token const& token,
     563                 :     std::error_code* ec,
     564                 :     std::size_t* bytes_out)
     565                 : {
     566              20 :     auto& op = wr_;
     567              20 :     op.reset();
     568              20 :     op.h         = h;
     569              20 :     op.ex        = ex;
     570              20 :     op.ec_out    = ec;
     571              20 :     op.bytes_out = bytes_out;
     572                 : 
     573              20 :     if (fd_ < 0)
     574                 :     {
     575 MIS           0 :         op.start(token, static_cast<Derived*>(this));
     576               0 :         op.impl_ptr = this->shared_from_this();
     577               0 :         op.complete(EBADF, 0);
     578               0 :         svc_.post(&op);
     579               0 :         return std::noop_coroutine();
     580                 :     }
     581                 : 
     582 HIT          20 :     capy::mutable_buffer bufs[write_op::max_buffers];
     583              20 :     op.iovec_count =
     584              20 :         static_cast<int>(param.copy_to(bufs, write_op::max_buffers));
     585                 : 
     586              20 :     if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
     587                 :     {
     588 MIS           0 :         op.start(token, static_cast<Derived*>(this));
     589               0 :         op.impl_ptr = this->shared_from_this();
     590               0 :         op.complete(0, 0);
     591               0 :         svc_.post(&op);
     592               0 :         return std::noop_coroutine();
     593                 :     }
     594                 : 
     595 HIT          20 :     if (int const nerr = arm_nonblocking())
     596                 :     {
     597 MIS           0 :         op.start(token, static_cast<Derived*>(this));
     598               0 :         op.impl_ptr = this->shared_from_this();
     599               0 :         op.complete(nerr, 0);
     600               0 :         svc_.post(&op);
     601               0 :         return std::noop_coroutine();
     602                 :     }
     603                 : 
     604 HIT          42 :     for (int i = 0; i < op.iovec_count; ++i)
     605                 :     {
     606              22 :         op.iovecs[i].iov_base = bufs[i].data();
     607              22 :         op.iovecs[i].iov_len  = bufs[i].size();
     608                 :     }
     609                 : 
     610                 :     // Speculative write; the single-buffer case skips the iov_iter setup.
     611                 :     ssize_t n;
     612              20 :     if (op.iovec_count == 1)
     613                 :     {
     614              36 :         n = write_op::write_policy::write_one(
     615              18 :             fd_, bufs[0].data(), bufs[0].size());
     616                 :     }
     617                 :     else
     618                 :     {
     619               2 :         n = write_op::write_policy::write(fd_, op.iovecs, op.iovec_count);
     620                 :     }
     621                 : 
     622              20 :     if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
     623                 :     {
     624              14 :         int err    = (n < 0) ? errno : 0;
     625              14 :         auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
     626                 : 
     627              14 :         if (svc_.scheduler().try_consume_inline_budget())
     628                 :         {
     629               2 :             *ec        = err ? make_err(err) : std::error_code{};
     630               2 :             *bytes_out = bytes;
     631               2 :             op.cont.h  = h;
     632               2 :             return dispatch_coro(ex, op.cont);
     633                 :         }
     634              12 :         op.start(token, static_cast<Derived*>(this));
     635              12 :         op.impl_ptr = this->shared_from_this();
     636              12 :         op.complete(err, bytes);
     637              12 :         svc_.post(&op);
     638              12 :         return std::noop_coroutine();
     639                 :     }
     640                 : 
     641                 :     // EAGAIN — register with reactor
     642               6 :     op.fd = fd_;
     643               6 :     op.start(token, static_cast<Derived*>(this));
     644               6 :     op.impl_ptr = this->shared_from_this();
     645                 : 
     646               6 :     this->register_op(op, desc_state_.write_op, desc_state_.write_ready, true);
     647               6 :     return std::noop_coroutine();
     648                 : }
     649                 : 
     650                 : template<class Derived, class Traits, class Service, class Acceptor>
     651                 : std::coroutine_handle<>
     652              37 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_wait(
     653                 :     std::coroutine_handle<> h,
     654                 :     capy::executor_ref ex,
     655                 :     wait_type w,
     656                 :     std::stop_token const& token,
     657                 :     std::error_code* ec)
     658                 : {
     659                 :     // Pick refs up-front to avoid duplicating the register_op call.
     660                 :     wait_op* op_ptr;
     661                 :     reactor_op_base** desc_slot_ptr;
     662                 :     std::uint32_t event;
     663                 : 
     664              37 :     if (w == wait_type::read)
     665                 :     {
     666              16 :         op_ptr        = &wait_rd_;
     667              16 :         desc_slot_ptr = &desc_state_.wait_read_op;
     668              16 :         event         = reactor_event_read;
     669                 :     }
     670              21 :     else if (w == wait_type::write)
     671                 :     {
     672               8 :         op_ptr        = &wait_wr_;
     673               8 :         desc_slot_ptr = &desc_state_.wait_write_op;
     674               8 :         event         = reactor_event_write;
     675                 :     }
     676                 :     else // wait_type::error
     677                 :     {
     678              13 :         op_ptr        = &wait_er_;
     679              13 :         desc_slot_ptr = &desc_state_.wait_error_op;
     680              13 :         event         = reactor_event_error;
     681                 :     }
     682                 : 
     683              37 :     auto& op = *op_ptr;
     684                 : 
     685                 :     // Speculative probe: an edge-triggered reactor cannot report a
     686                 :     // condition that already holds, so a wait initiated on an already
     687                 :     // ready descriptor would otherwise park forever. No syscall here
     688                 :     // modifies the descriptor -- in particular O_NONBLOCK is untouched.
     689              37 :     int perr = 0;
     690              37 :     if (wait_op::probe(fd_, event, perr))
     691                 :     {
     692              12 :         if (svc_.scheduler().try_consume_inline_budget())
     693                 :         {
     694               2 :             *ec       = perr ? make_err(perr) : std::error_code{};
     695               2 :             op.cont.h = h;
     696               2 :             return dispatch_coro(ex, op.cont);
     697                 :         }
     698              10 :         op.reset();
     699              10 :         op.wait_event = event;
     700              10 :         op.h          = h;
     701              10 :         op.ex         = ex;
     702              10 :         op.ec_out     = ec;
     703              10 :         op.fd         = fd_;
     704              10 :         op.start(token, static_cast<Derived*>(this));
     705              10 :         op.impl_ptr = this->shared_from_this();
     706              10 :         op.complete(perr, 0);
     707              10 :         svc_.post(&op);
     708              10 :         return std::noop_coroutine();
     709                 :     }
     710                 : 
     711              25 :     op.reset();
     712              25 :     op.wait_event = event;
     713              25 :     op.h          = h;
     714              25 :     op.ex         = ex;
     715              25 :     op.ec_out     = ec;
     716              25 :     op.fd         = fd_;
     717              25 :     op.start(token, static_cast<Derived*>(this));
     718              25 :     op.impl_ptr = this->shared_from_this();
     719                 : 
     720                 :     // Force register_op's ready path so the wait op re-probes under the
     721                 :     // descriptor mutex before parking. A stale write_ready latched at
     722                 :     // registration would otherwise report a full pipe as writable.
     723              25 :     bool force_probe = true;
     724              25 :     this->register_op(
     725                 :         op, *desc_slot_ptr, force_probe, event == reactor_event_write);
     726              25 :     return std::noop_coroutine();
     727                 : }
     728                 : 
     729                 : } // namespace boost::corosio::detail
     730                 : 
     731                 : #endif // BOOST_COROSIO_POSIX
     732                 : 
     733                 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_HPP
        

Generated by: LCOV version 2.3