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: 78.9 % 308 243 65
Test Date: 2026-09-28 20:06:38 Functions: 97.1 % 70 68 2

           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_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_op.hpp>
      26                 : #include <boost/corosio/native/detail/reactor/reactor_op_complete.hpp>
      27                 : #include <boost/capy/buffers.hpp>
      28                 : 
      29                 : #include <coroutine>
      30                 : #include <memory>
      31                 : #include <mutex>
      32                 : #include <utility>
      33                 : 
      34                 : #include <errno.h>
      35                 : #include <sys/uio.h>
      36                 : #include <unistd.h>
      37                 : 
      38                 : /* Reactor-backed implementation of posix_descriptor.
      39                 : 
      40                 :    Deliberately does not derive from reactor_basic_socket: that base
      41                 :    is rooted in native_socket_base, which overrides local_endpoint(),
      42                 :    set_option() and get_option() on its ImplBase. posix_descriptor has
      43                 :    no socket verbs, so the shared logic (init_and_register, register_op,
      44                 :    the cancel/close/release op sweeps) is carried here against five op
      45                 :    slots instead of eight.
      46                 : 
      47                 :    The one behavior that is genuinely new: O_NONBLOCK is armed lazily,
      48                 :    on the first read_some/write_some and never from assign() or wait().
      49                 :    The flag lives on the shared open file description, so arming it is
      50                 :    visible to every other holder of that description -- which is why a
      51                 :    wait()-only user must never trigger it.
      52                 : 
      53                 :    PARALLEL COPY: init_and_register, register_op, cancel_single_op and
      54                 :    the cancel/close/release op sweeps here mirror the socket versions in
      55                 :    reactor_basic_socket.hpp (init_and_register, register_op,
      56                 :    cancel_single_op, do_cancel, do_close_socket, do_release_socket).
      57                 :    The two are separate code because that base carries native_socket_base
      58                 :    and its socket verbs; they are not separate protocols. A fix to the
      59                 :    cancel/park protocol -- slot claiming under desc_state_.mutex, the
      60                 :    cached-edge replay in register_op, the impl_ref_ pinning during
      61                 :    teardown -- belongs in both files.
      62                 : */
      63                 : 
      64                 : namespace boost::corosio::detail {
      65                 : 
      66                 : // ============================================================
      67                 : // Op types
      68                 : // ============================================================
      69                 : 
      70                 : /* Descriptor op family.
      71                 : 
      72                 :    Mirrors reactor_stream_ops.hpp. The Acceptor parameter is a
      73                 :    placeholder: reactor_op is parameterized on both a socket and an
      74                 :    acceptor impl type, and the shared completion helpers name
      75                 :    acceptor_impl_ in a branch that a descriptor op never takes but the
      76                 :    compiler still instantiates. Passing the backend's acceptor type (as
      77                 :    reactor_dgram_socket_impl already does) keeps those helpers shared
      78                 :    rather than duplicated here.
      79                 : 
      80                 :    @tparam Traits     Backend traits (epoll_traits, kqueue_traits, ...).
      81                 :    @tparam Descriptor The concrete descriptor type (forward-declared).
      82                 :    @tparam Acceptor   Placeholder acceptor type for the op base.
      83                 : */
      84                 : 
      85                 : template<class Traits, class Descriptor, class Acceptor>
      86                 : struct reactor_descriptor_base_op : reactor_op<Descriptor, Acceptor>
      87                 : {
      88                 :     void operator()() override;
      89                 :     void cancel() noexcept override;
      90                 : };
      91                 : 
      92                 : template<class Traits, class Descriptor, class Acceptor>
      93                 : struct reactor_descriptor_read_op final
      94                 :     : reactor_read_op<reactor_descriptor_base_op<Traits, Descriptor, Acceptor>>
      95                 : {};
      96                 : 
      97                 : template<class Traits, class Descriptor, class Acceptor>
      98                 : struct reactor_descriptor_write_op final
      99                 :     : reactor_write_op<
     100                 :           reactor_descriptor_base_op<Traits, Descriptor, Acceptor>,
     101                 :           typename Traits::descriptor_write_policy>
     102                 : {};
     103                 : 
     104                 : template<class Traits, class Descriptor, class Acceptor>
     105                 : struct reactor_descriptor_wait_op final
     106                 :     : reactor_wait_op<reactor_descriptor_base_op<Traits, Descriptor, Acceptor>>
     107                 : {
     108                 :     void operator()() override;
     109                 : };
     110                 : 
     111                 : // --- Deferred implementations (instantiated when Descriptor is complete) ---
     112                 : 
     113                 : template<class Traits, class Descriptor, class Acceptor>
     114                 : void
     115 HIT          25 : reactor_descriptor_base_op<Traits, Descriptor, Acceptor>::operator()()
     116                 : {
     117              25 :     complete_io_op(*this);
     118              25 : }
     119                 : 
     120                 : template<class Traits, class Descriptor, class Acceptor>
     121                 : void
     122               2 : reactor_descriptor_base_op<Traits, Descriptor, Acceptor>::cancel() noexcept
     123                 : {
     124                 :     // A descriptor op is only ever started against a descriptor impl, so
     125                 :     // the acceptor arm of the stream op's cancel() has no counterpart.
     126               2 :     if (this->socket_impl_)
     127               2 :         this->socket_impl_->cancel_single_op(*this);
     128                 :     else
     129 MIS           0 :         this->request_cancel();
     130 HIT           2 : }
     131                 : 
     132                 : template<class Traits, class Descriptor, class Acceptor>
     133                 : void
     134              16 : reactor_descriptor_wait_op<Traits, Descriptor, Acceptor>::operator()()
     135                 : {
     136              16 :     complete_wait_op(*this);
     137              16 : }
     138                 : 
     139                 : // ============================================================
     140                 : // Descriptor implementation
     141                 : // ============================================================
     142                 : 
     143                 : /** CRTP base for reactor-backed posix_descriptor implementations.
     144                 : 
     145                 :     Holds the adopted descriptor, its reactor registration state, and
     146                 :     the five op slots (read, write, and one wait per direction).
     147                 : 
     148                 :     @tparam Derived  The named final class (CRTP self).
     149                 :     @tparam Traits   Backend traits (epoll_traits, kqueue_traits, ...).
     150                 :     @tparam Service  The backend's descriptor service type.
     151                 :     @tparam Acceptor Placeholder acceptor type for the op base.
     152                 : */
     153                 : template<class Derived, class Traits, class Service, class Acceptor>
     154                 : class reactor_descriptor
     155                 :     : public posix_descriptor::implementation
     156                 :     , public std::enable_shared_from_this<Derived>
     157                 :     , public intrusive_list<Derived>::node
     158                 : {
     159                 :     using base_op  = reactor_descriptor_base_op<Traits, Derived, Acceptor>;
     160                 :     using read_op  = reactor_descriptor_read_op<Traits, Derived, Acceptor>;
     161                 :     using write_op = reactor_descriptor_write_op<Traits, Derived, Acceptor>;
     162                 :     using wait_op  = reactor_descriptor_wait_op<Traits, Derived, Acceptor>;
     163                 : 
     164                 : protected:
     165                 :     // NOLINTNEXTLINE(bugprone-crtp-constructor-accessibility)
     166              59 :     explicit reactor_descriptor(Service& svc) noexcept : svc_(svc) {}
     167                 : 
     168                 : public:
     169              59 :     ~reactor_descriptor() override = default;
     170                 : 
     171                 :     /// Per-descriptor state for persistent reactor registration.
     172                 :     typename Traits::desc_state_type desc_state_;
     173                 : 
     174                 :     // --- Virtual method overrides ---
     175                 : 
     176              23 :     std::coroutine_handle<> read_some(
     177                 :         std::coroutine_handle<> h,
     178                 :         capy::executor_ref ex,
     179                 :         buffer_param param,
     180                 :         std::stop_token token,
     181                 :         std::error_code* ec,
     182                 :         std::size_t* bytes_out) override
     183                 :     {
     184              23 :         return do_read_some(h, ex, param, token, ec, bytes_out);
     185                 :     }
     186                 : 
     187               6 :     std::coroutine_handle<> write_some(
     188                 :         std::coroutine_handle<> h,
     189                 :         capy::executor_ref ex,
     190                 :         buffer_param param,
     191                 :         std::stop_token token,
     192                 :         std::error_code* ec,
     193                 :         std::size_t* bytes_out) override
     194                 :     {
     195               6 :         return do_write_some(h, ex, param, token, ec, bytes_out);
     196                 :     }
     197                 : 
     198              16 :     std::coroutine_handle<> wait(
     199                 :         std::coroutine_handle<> h,
     200                 :         capy::executor_ref ex,
     201                 :         wait_type w,
     202                 :         std::stop_token token,
     203                 :         std::error_code* ec) override
     204                 :     {
     205              16 :         return do_wait(h, ex, w, token, ec);
     206                 :     }
     207                 : 
     208             139 :     native_handle_type native_handle() const noexcept override
     209                 :     {
     210             139 :         return fd_;
     211                 :     }
     212                 : 
     213                 :     native_handle_type release_descriptor() noexcept override;
     214                 : 
     215               3 :     void cancel() noexcept override
     216                 :     {
     217               3 :         do_cancel();
     218               3 :     }
     219                 : 
     220                 :     // --- Service-facing (non-virtual) ---
     221                 : 
     222                 :     /** Adopt the fd, initialize descriptor state, and register it.
     223                 : 
     224                 :         @param fd The descriptor to adopt.
     225                 : 
     226                 :         @return The error if the reactor rejects the descriptor, in
     227                 :         which case the implementation is left closed and the caller
     228                 :         retains ownership of @a fd; otherwise a default constructed
     229                 :         error code.
     230                 :     */
     231                 :     std::error_code init_and_register(int fd) noexcept;
     232                 : 
     233                 :     /// Close the descriptor and cancel pending operations.
     234                 :     void close_descriptor() noexcept;
     235                 : 
     236                 :     /// Cancel a single pending operation, claiming it from its slot.
     237                 :     template<class Op>
     238                 :     void cancel_single_op(Op& op) noexcept;
     239                 : 
     240                 : private:
     241                 :     /** Arm O_NONBLOCK, once, before the first speculative syscall.
     242                 : 
     243                 :         Reports an errno rather than an error_code because that is what
     244                 :         the op result model records; the round trip is lossless here
     245                 :         because fcntl only fails with codes make_err passes through.
     246                 :     */
     247              29 :     int arm_nonblocking() noexcept
     248                 :     {
     249              29 :         if (nonblocking_)
     250 MIS           0 :             return 0;
     251 HIT          29 :         if (auto ec = ensure_nonblocking(fd_))
     252 MIS           0 :             return ec.value();
     253 HIT          29 :         nonblocking_ = true;
     254              29 :         return 0;
     255                 :     }
     256                 : 
     257                 :     std::coroutine_handle<> do_read_some(
     258                 :         std::coroutine_handle<>,
     259                 :         capy::executor_ref,
     260                 :         buffer_param,
     261                 :         std::stop_token const&,
     262                 :         std::error_code*,
     263                 :         std::size_t*);
     264                 : 
     265                 :     std::coroutine_handle<> do_write_some(
     266                 :         std::coroutine_handle<>,
     267                 :         capy::executor_ref,
     268                 :         buffer_param,
     269                 :         std::stop_token const&,
     270                 :         std::error_code*,
     271                 :         std::size_t*);
     272                 : 
     273                 :     std::coroutine_handle<> do_wait(
     274                 :         std::coroutine_handle<>,
     275                 :         capy::executor_ref,
     276                 :         wait_type,
     277                 :         std::stop_token const&,
     278                 :         std::error_code*);
     279                 : 
     280                 :     void do_cancel() noexcept;
     281                 : 
     282                 :     /// Register an op with the reactor, handling cached edge events.
     283                 :     template<class Op>
     284                 :     void register_op(
     285                 :         Op& op,
     286                 :         reactor_op_base*& desc_slot,
     287                 :         bool& ready_flag,
     288                 :         bool is_write_direction = false) noexcept;
     289                 : 
     290                 :     /// Apply @a fn to each of the five op slots.
     291                 :     template<class Fn>
     292             213 :     void for_each_op(Fn fn) noexcept
     293                 :     {
     294             213 :         fn(rd_);
     295             213 :         fn(wr_);
     296             213 :         fn(wait_rd_);
     297             213 :         fn(wait_wr_);
     298             213 :         fn(wait_er_);
     299             213 :     }
     300                 : 
     301                 :     /** Claim every parked op out of its descriptor_state slot.
     302                 : 
     303                 :         @param claimed Receives the claimed ops; must hold five.
     304                 :         @param teardown Also clear the cached edge flags and, if the
     305                 :         state is queued in the scheduler, pin the impl alive.
     306                 :         @param self Keepalive used by @a teardown.
     307                 :         @return The number of ops claimed.
     308                 :     */
     309                 :     int claim_parked_ops(
     310                 :         reactor_op_base** claimed,
     311                 :         bool teardown,
     312                 :         std::shared_ptr<Derived> const& self) noexcept;
     313                 : 
     314                 :     /// Post claimed ops to the scheduler, keeping the impl alive.
     315                 :     void post_claimed_ops(
     316                 :         reactor_op_base** claimed,
     317                 :         int count,
     318                 :         std::shared_ptr<Derived> const& self) noexcept;
     319                 : 
     320                 :     /// Sweep every op slot, then drop the reactor registration.
     321                 :     void quiesce() noexcept;
     322                 : 
     323                 :     reactor_op_base** op_to_desc_slot(base_op& op) noexcept;
     324                 : 
     325                 :     Service& svc_;
     326                 :     int fd_           = -1;
     327                 :     bool nonblocking_ = false;
     328                 : 
     329                 :     read_op rd_;
     330                 :     write_op wr_;
     331                 :     wait_op wait_rd_;
     332                 :     wait_op wait_wr_;
     333                 :     wait_op wait_er_;
     334                 : };
     335                 : 
     336                 : // ============================================================
     337                 : // Registration and teardown
     338                 : // ============================================================
     339                 : 
     340                 : template<class Derived, class Traits, class Service, class Acceptor>
     341                 : std::error_code
     342              47 : reactor_descriptor<Derived, Traits, Service, Acceptor>::init_and_register(
     343                 :     int fd) noexcept
     344                 : {
     345              47 :     fd_            = fd;
     346              47 :     desc_state_.fd = fd;
     347                 :     {
     348                 :         // Every slot this type owns; connect_op is deliberately absent,
     349                 :         // a descriptor has no connect operation to park there.
     350              47 :         std::lock_guard lock(desc_state_.mutex);
     351              47 :         desc_state_.read_op       = nullptr;
     352              47 :         desc_state_.write_op      = nullptr;
     353              47 :         desc_state_.wait_read_op  = nullptr;
     354              47 :         desc_state_.wait_write_op = nullptr;
     355              47 :         desc_state_.wait_error_op = nullptr;
     356              47 :     }
     357              47 :     if (auto ec = svc_.scheduler().register_descriptor(fd, &desc_state_))
     358                 :     {
     359                 :         // Undo the partial state so a failed adopt is
     360                 :         // indistinguishable from a closed implementation.
     361 MIS           0 :         fd_                           = -1;
     362               0 :         desc_state_.fd                = -1;
     363               0 :         desc_state_.registered_events = 0;
     364               0 :         return ec;
     365                 :     }
     366 HIT          47 :     return {};
     367                 : }
     368                 : 
     369                 : template<class Derived, class Traits, class Service, class Acceptor>
     370                 : int
     371             213 : reactor_descriptor<Derived, Traits, Service, Acceptor>::claim_parked_ops(
     372                 :     reactor_op_base** claimed,
     373                 :     bool teardown,
     374                 :     std::shared_ptr<Derived> const& self) noexcept
     375                 : {
     376             213 :     int count = 0;
     377             213 :     std::lock_guard lock(desc_state_.mutex);
     378            1278 :     for (auto** slot :
     379             213 :          {&desc_state_.read_op, &desc_state_.write_op,
     380             213 :           &desc_state_.wait_read_op, &desc_state_.wait_write_op,
     381             213 :           &desc_state_.wait_error_op})
     382                 :     {
     383            1065 :         if (auto* c = std::exchange(*slot, nullptr))
     384               5 :             claimed[count++] = c;
     385                 :     }
     386             213 :     if (teardown)
     387                 :     {
     388             210 :         desc_state_.read_ready  = false;
     389             210 :         desc_state_.write_ready = false;
     390                 : 
     391                 :         // Must be set under the same lock that invoke_deferred_io clears
     392                 :         // is_enqueued_ under, or the impl could be destroyed while the
     393                 :         // scheduler still holds the queued descriptor_state.
     394             210 :         if (desc_state_.is_enqueued_.load(std::memory_order_acquire))
     395              33 :             desc_state_.impl_ref_ = self;
     396                 :     }
     397             426 :     return count;
     398             213 : }
     399                 : 
     400                 : template<class Derived, class Traits, class Service, class Acceptor>
     401                 : void
     402             213 : reactor_descriptor<Derived, Traits, Service, Acceptor>::post_claimed_ops(
     403                 :     reactor_op_base** claimed,
     404                 :     int count,
     405                 :     std::shared_ptr<Derived> const& self) noexcept
     406                 : {
     407             218 :     for (int i = 0; i < count; ++i)
     408                 :     {
     409               5 :         claimed[i]->impl_ptr = self;
     410               5 :         svc_.post(claimed[i]);
     411               5 :         svc_.work_finished();
     412                 :     }
     413             213 : }
     414                 : 
     415                 : template<class Derived, class Traits, class Service, class Acceptor>
     416                 : void
     417               3 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_cancel() noexcept
     418                 : {
     419               3 :     auto self = this->weak_from_this().lock();
     420               3 :     if (!self)
     421 MIS           0 :         return;
     422                 : 
     423 HIT          18 :     for_each_op([](auto& op) { op.request_cancel(); });
     424                 : 
     425                 :     reactor_op_base* claimed[5];
     426               3 :     int const count = claim_parked_ops(claimed, /*teardown=*/false, self);
     427               3 :     post_claimed_ops(claimed, count, self);
     428               3 : }
     429                 : 
     430                 : template<class Derived, class Traits, class Service, class Acceptor>
     431                 : void
     432             210 : reactor_descriptor<Derived, Traits, Service, Acceptor>::quiesce() noexcept
     433                 : {
     434             210 :     auto self = this->weak_from_this().lock();
     435             210 :     if (self)
     436                 :     {
     437            1260 :         for_each_op([](auto& op) { op.request_cancel(); });
     438                 : 
     439                 :         reactor_op_base* claimed[5];
     440             210 :         int const count = claim_parked_ops(claimed, /*teardown=*/true, self);
     441             210 :         post_claimed_ops(claimed, count, self);
     442                 :     }
     443                 : 
     444             210 :     if (fd_ >= 0 && desc_state_.registered_events != 0)
     445              47 :         svc_.scheduler().deregister_descriptor(fd_);
     446                 : 
     447             210 :     desc_state_.registered_events = 0;
     448                 :     // The next adopted fd starts from an unknown flag state.
     449             210 :     nonblocking_ = false;
     450             210 : }
     451                 : 
     452                 : template<class Derived, class Traits, class Service, class Acceptor>
     453                 : void
     454             208 : reactor_descriptor<Derived, Traits, Service, Acceptor>::
     455                 :     close_descriptor() noexcept
     456                 : {
     457             208 :     quiesce();
     458                 : 
     459             208 :     if (fd_ >= 0)
     460                 :     {
     461              45 :         ::close(fd_);
     462              45 :         fd_ = -1;
     463                 :     }
     464             208 :     desc_state_.fd = -1;
     465             208 : }
     466                 : 
     467                 : template<class Derived, class Traits, class Service, class Acceptor>
     468                 : native_handle_type
     469               2 : reactor_descriptor<Derived, Traits, Service, Acceptor>::
     470                 :     release_descriptor() noexcept
     471                 : {
     472               2 :     quiesce();
     473                 : 
     474                 :     // Do NOT close -- the caller takes ownership.
     475               2 :     native_handle_type released = fd_;
     476               2 :     fd_                         = -1;
     477               2 :     desc_state_.fd              = -1;
     478               2 :     return released;
     479                 : }
     480                 : 
     481                 : // ============================================================
     482                 : // Op registration and per-op cancellation
     483                 : // ============================================================
     484                 : 
     485                 : template<class Derived, class Traits, class Service, class Acceptor>
     486                 : template<class Op>
     487                 : void
     488              16 : reactor_descriptor<Derived, Traits, Service, Acceptor>::register_op(
     489                 :     Op& op,
     490                 :     reactor_op_base*& desc_slot,
     491                 :     bool& ready_flag,
     492                 :     bool is_write_direction) noexcept
     493                 : {
     494              16 :     svc_.work_started();
     495                 : 
     496              16 :     std::lock_guard lock(desc_state_.mutex);
     497              16 :     bool io_done = false;
     498              16 :     if (ready_flag)
     499                 :     {
     500               8 :         ready_flag = false;
     501               8 :         op.perform_io();
     502               8 :         io_done = (op.errn != EAGAIN && op.errn != EWOULDBLOCK);
     503               8 :         if (!io_done)
     504               8 :             op.errn = 0;
     505                 :     }
     506                 : 
     507              16 :     if (io_done || op.cancelled.load(std::memory_order_acquire))
     508                 :     {
     509 MIS           0 :         svc_.post(&op);
     510               0 :         svc_.work_finished();
     511                 :     }
     512                 :     else
     513                 :     {
     514 HIT          16 :         desc_slot = &op;
     515                 : 
     516                 :         // Select must rebuild its fd_sets when a write-direction op
     517                 :         // is parked, so select() watches for writability. Compiled
     518                 :         // away to nothing for epoll and kqueue.
     519                 :         if constexpr (Service::needs_write_notification)
     520                 :         {
     521               8 :             if (is_write_direction)
     522               2 :                 svc_.scheduler().notify_reactor();
     523                 :         }
     524                 :     }
     525              16 : }
     526                 : 
     527                 : template<class Derived, class Traits, class Service, class Acceptor>
     528                 : reactor_op_base**
     529               2 : reactor_descriptor<Derived, Traits, Service, Acceptor>::op_to_desc_slot(
     530                 :     base_op& op) noexcept
     531                 : {
     532               2 :     if (&op == static_cast<void*>(&rd_))
     533               2 :         return &desc_state_.read_op;
     534 MIS           0 :     if (&op == static_cast<void*>(&wr_))
     535               0 :         return &desc_state_.write_op;
     536               0 :     if (&op == static_cast<void*>(&wait_rd_))
     537               0 :         return &desc_state_.wait_read_op;
     538               0 :     if (&op == static_cast<void*>(&wait_wr_))
     539               0 :         return &desc_state_.wait_write_op;
     540               0 :     if (&op == static_cast<void*>(&wait_er_))
     541               0 :         return &desc_state_.wait_error_op;
     542               0 :     return nullptr;
     543                 : }
     544                 : 
     545                 : template<class Derived, class Traits, class Service, class Acceptor>
     546                 : template<class Op>
     547                 : void
     548 HIT           2 : reactor_descriptor<Derived, Traits, Service, Acceptor>::cancel_single_op(
     549                 :     Op& op) noexcept
     550                 : {
     551               2 :     auto self = this->weak_from_this().lock();
     552               2 :     if (!self)
     553 MIS           0 :         return;
     554                 : 
     555 HIT           2 :     op.request_cancel();
     556                 : 
     557               2 :     reactor_op_base** desc_op_ptr = op_to_desc_slot(op);
     558               2 :     if (!desc_op_ptr)
     559 MIS           0 :         return;
     560                 : 
     561 HIT           2 :     reactor_op_base* claimed = nullptr;
     562                 :     {
     563               2 :         std::lock_guard lock(desc_state_.mutex);
     564               2 :         if (*desc_op_ptr == &op)
     565 MIS           0 :             claimed = std::exchange(*desc_op_ptr, nullptr);
     566                 :         // Not in the slot: request_cancel() above already set
     567                 :         // op.cancelled, which register_op consults before parking
     568                 :         // and the completion decode consults on delivery. Latching
     569                 :         // a descriptor flag here instead would outlive this op and
     570                 :         // cancel the next wait in the same direction.
     571 HIT           2 :     }
     572               2 :     if (claimed)
     573                 :     {
     574 MIS           0 :         op.impl_ptr = self;
     575               0 :         svc_.post(&op);
     576               0 :         svc_.work_finished();
     577                 :     }
     578 HIT           2 : }
     579                 : 
     580                 : // ============================================================
     581                 : // I/O dispatch
     582                 : // ============================================================
     583                 : 
     584                 : template<class Derived, class Traits, class Service, class Acceptor>
     585                 : std::coroutine_handle<>
     586              23 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_read_some(
     587                 :     std::coroutine_handle<> h,
     588                 :     capy::executor_ref ex,
     589                 :     buffer_param param,
     590                 :     std::stop_token const& token,
     591                 :     std::error_code* ec,
     592                 :     std::size_t* bytes_out)
     593                 : {
     594              23 :     auto& op = rd_;
     595              23 :     op.reset();
     596              23 :     op.h         = h;
     597              23 :     op.ex        = ex;
     598              23 :     op.ec_out    = ec;
     599              23 :     op.bytes_out = bytes_out;
     600                 : 
     601                 :     // Closed-object contract: complete with bad_file_descriptor without
     602                 :     // touching the kernel or the unregistered descriptor state.
     603              23 :     if (fd_ < 0)
     604                 :     {
     605 MIS           0 :         op.start(token, static_cast<Derived*>(this));
     606               0 :         op.impl_ptr = this->shared_from_this();
     607               0 :         op.complete(EBADF, 0);
     608               0 :         svc_.post(&op);
     609               0 :         return std::noop_coroutine();
     610                 :     }
     611                 : 
     612 HIT          23 :     capy::mutable_buffer bufs[read_op::max_buffers];
     613              23 :     op.iovec_count =
     614              23 :         static_cast<int>(param.copy_to(bufs, read_op::max_buffers));
     615                 : 
     616              23 :     if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
     617                 :     {
     618 MIS           0 :         op.empty_buffer_read = true;
     619               0 :         op.start(token, static_cast<Derived*>(this));
     620               0 :         op.impl_ptr = this->shared_from_this();
     621               0 :         op.complete(0, 0);
     622               0 :         svc_.post(&op);
     623               0 :         return std::noop_coroutine();
     624                 :     }
     625                 : 
     626                 :     // The first transferring operation is what arms O_NONBLOCK; assign()
     627                 :     // and wait() never do.
     628 HIT          23 :     if (int const nerr = arm_nonblocking())
     629                 :     {
     630 MIS           0 :         op.start(token, static_cast<Derived*>(this));
     631               0 :         op.impl_ptr = this->shared_from_this();
     632               0 :         op.complete(nerr, 0);
     633               0 :         svc_.post(&op);
     634               0 :         return std::noop_coroutine();
     635                 :     }
     636                 : 
     637 HIT          48 :     for (int i = 0; i < op.iovec_count; ++i)
     638                 :     {
     639              25 :         op.iovecs[i].iov_base = bufs[i].data();
     640              25 :         op.iovecs[i].iov_len  = bufs[i].size();
     641                 :     }
     642                 : 
     643                 :     // Speculative read; the single-buffer case uses read() so the kernel
     644                 :     // skips the readv iov_iter setup.
     645                 :     ssize_t n;
     646              23 :     if (op.iovec_count == 1)
     647                 :     {
     648                 :         do
     649                 :         {
     650              21 :             n = ::read(fd_, bufs[0].data(), bufs[0].size());
     651                 :         }
     652              21 :         while (n < 0 && errno == EINTR);
     653                 :     }
     654                 :     else
     655                 :     {
     656                 :         do
     657                 :         {
     658               2 :             n = ::readv(fd_, op.iovecs, op.iovec_count);
     659                 :         }
     660               2 :         while (n < 0 && errno == EINTR);
     661                 :     }
     662                 : 
     663              23 :     if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
     664                 :     {
     665              17 :         int err    = (n < 0) ? errno : 0;
     666              17 :         auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
     667                 : 
     668              17 :         if (svc_.scheduler().try_consume_inline_budget())
     669                 :         {
     670               4 :             if (err)
     671 MIS           0 :                 *ec = make_err(err);
     672 HIT           4 :             else if (n == 0)
     673 MIS           0 :                 *ec = capy::error::eof;
     674                 :             else
     675 HIT           4 :                 *ec = {};
     676               4 :             *bytes_out = bytes;
     677               4 :             op.cont.h  = h;
     678               4 :             return dispatch_coro(ex, op.cont);
     679                 :         }
     680              13 :         op.start(token, static_cast<Derived*>(this));
     681              13 :         op.impl_ptr = this->shared_from_this();
     682              13 :         op.complete(err, bytes);
     683              13 :         svc_.post(&op);
     684              13 :         return std::noop_coroutine();
     685                 :     }
     686                 : 
     687                 :     // EAGAIN — register with reactor
     688               6 :     op.fd = fd_;
     689               6 :     op.start(token, static_cast<Derived*>(this));
     690               6 :     op.impl_ptr = this->shared_from_this();
     691                 : 
     692               6 :     register_op(op, desc_state_.read_op, desc_state_.read_ready);
     693               6 :     return std::noop_coroutine();
     694                 : }
     695                 : 
     696                 : template<class Derived, class Traits, class Service, class Acceptor>
     697                 : std::coroutine_handle<>
     698               6 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_write_some(
     699                 :     std::coroutine_handle<> h,
     700                 :     capy::executor_ref ex,
     701                 :     buffer_param param,
     702                 :     std::stop_token const& token,
     703                 :     std::error_code* ec,
     704                 :     std::size_t* bytes_out)
     705                 : {
     706               6 :     auto& op = wr_;
     707               6 :     op.reset();
     708               6 :     op.h         = h;
     709               6 :     op.ex        = ex;
     710               6 :     op.ec_out    = ec;
     711               6 :     op.bytes_out = bytes_out;
     712                 : 
     713               6 :     if (fd_ < 0)
     714                 :     {
     715 MIS           0 :         op.start(token, static_cast<Derived*>(this));
     716               0 :         op.impl_ptr = this->shared_from_this();
     717               0 :         op.complete(EBADF, 0);
     718               0 :         svc_.post(&op);
     719               0 :         return std::noop_coroutine();
     720                 :     }
     721                 : 
     722 HIT           6 :     capy::mutable_buffer bufs[write_op::max_buffers];
     723               6 :     op.iovec_count =
     724               6 :         static_cast<int>(param.copy_to(bufs, write_op::max_buffers));
     725                 : 
     726               6 :     if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
     727                 :     {
     728 MIS           0 :         op.start(token, static_cast<Derived*>(this));
     729               0 :         op.impl_ptr = this->shared_from_this();
     730               0 :         op.complete(0, 0);
     731               0 :         svc_.post(&op);
     732               0 :         return std::noop_coroutine();
     733                 :     }
     734                 : 
     735 HIT           6 :     if (int const nerr = arm_nonblocking())
     736                 :     {
     737 MIS           0 :         op.start(token, static_cast<Derived*>(this));
     738               0 :         op.impl_ptr = this->shared_from_this();
     739               0 :         op.complete(nerr, 0);
     740               0 :         svc_.post(&op);
     741               0 :         return std::noop_coroutine();
     742                 :     }
     743                 : 
     744 HIT          14 :     for (int i = 0; i < op.iovec_count; ++i)
     745                 :     {
     746               8 :         op.iovecs[i].iov_base = bufs[i].data();
     747               8 :         op.iovecs[i].iov_len  = bufs[i].size();
     748                 :     }
     749                 : 
     750                 :     // Speculative write; the single-buffer case skips the iov_iter setup.
     751                 :     ssize_t n;
     752               6 :     if (op.iovec_count == 1)
     753                 :     {
     754               8 :         n = write_op::write_policy::write_one(
     755               4 :             fd_, bufs[0].data(), bufs[0].size());
     756                 :     }
     757                 :     else
     758                 :     {
     759               2 :         n = write_op::write_policy::write(fd_, op.iovecs, op.iovec_count);
     760                 :     }
     761                 : 
     762               6 :     if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
     763                 :     {
     764               4 :         int err    = (n < 0) ? errno : 0;
     765               4 :         auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
     766                 : 
     767               4 :         if (svc_.scheduler().try_consume_inline_budget())
     768                 :         {
     769 MIS           0 :             *ec        = err ? make_err(err) : std::error_code{};
     770               0 :             *bytes_out = bytes;
     771               0 :             op.cont.h  = h;
     772               0 :             return dispatch_coro(ex, op.cont);
     773                 :         }
     774 HIT           4 :         op.start(token, static_cast<Derived*>(this));
     775               4 :         op.impl_ptr = this->shared_from_this();
     776               4 :         op.complete(err, bytes);
     777               4 :         svc_.post(&op);
     778               4 :         return std::noop_coroutine();
     779                 :     }
     780                 : 
     781                 :     // EAGAIN — register with reactor
     782               2 :     op.fd = fd_;
     783               2 :     op.start(token, static_cast<Derived*>(this));
     784               2 :     op.impl_ptr = this->shared_from_this();
     785                 : 
     786               2 :     register_op(op, desc_state_.write_op, desc_state_.write_ready, true);
     787               2 :     return std::noop_coroutine();
     788                 : }
     789                 : 
     790                 : template<class Derived, class Traits, class Service, class Acceptor>
     791                 : std::coroutine_handle<>
     792              16 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_wait(
     793                 :     std::coroutine_handle<> h,
     794                 :     capy::executor_ref ex,
     795                 :     wait_type w,
     796                 :     std::stop_token const& token,
     797                 :     std::error_code* ec)
     798                 : {
     799                 :     // Pick refs up-front to avoid duplicating the register_op call.
     800                 :     wait_op* op_ptr;
     801                 :     reactor_op_base** desc_slot_ptr;
     802                 :     std::uint32_t event;
     803                 : 
     804              16 :     if (w == wait_type::read)
     805                 :     {
     806              12 :         op_ptr        = &wait_rd_;
     807              12 :         desc_slot_ptr = &desc_state_.wait_read_op;
     808              12 :         event         = reactor_event_read;
     809                 :     }
     810               4 :     else if (w == wait_type::write)
     811                 :     {
     812               2 :         op_ptr        = &wait_wr_;
     813               2 :         desc_slot_ptr = &desc_state_.wait_write_op;
     814               2 :         event         = reactor_event_write;
     815                 :     }
     816                 :     else // wait_type::error
     817                 :     {
     818               2 :         op_ptr        = &wait_er_;
     819               2 :         desc_slot_ptr = &desc_state_.wait_error_op;
     820               2 :         event         = reactor_event_error;
     821                 :     }
     822                 : 
     823              16 :     auto& op = *op_ptr;
     824                 : 
     825                 :     // Speculative probe: an edge-triggered reactor cannot report a
     826                 :     // condition that already holds, so a wait initiated on an already
     827                 :     // ready descriptor would otherwise park forever. No syscall here
     828                 :     // modifies the descriptor -- in particular O_NONBLOCK is untouched.
     829              16 :     int perr = 0;
     830              16 :     if (wait_op::probe(fd_, event, perr))
     831                 :     {
     832               8 :         if (svc_.scheduler().try_consume_inline_budget())
     833                 :         {
     834 MIS           0 :             *ec       = perr ? make_err(perr) : std::error_code{};
     835               0 :             op.cont.h = h;
     836               0 :             return dispatch_coro(ex, op.cont);
     837                 :         }
     838 HIT           8 :         op.reset();
     839               8 :         op.wait_event = event;
     840               8 :         op.h          = h;
     841               8 :         op.ex         = ex;
     842               8 :         op.ec_out     = ec;
     843               8 :         op.fd         = fd_;
     844               8 :         op.start(token, static_cast<Derived*>(this));
     845               8 :         op.impl_ptr = this->shared_from_this();
     846               8 :         op.complete(perr, 0);
     847               8 :         svc_.post(&op);
     848               8 :         return std::noop_coroutine();
     849                 :     }
     850                 : 
     851               8 :     op.reset();
     852               8 :     op.wait_event = event;
     853               8 :     op.h          = h;
     854               8 :     op.ex         = ex;
     855               8 :     op.ec_out     = ec;
     856               8 :     op.fd         = fd_;
     857               8 :     op.start(token, static_cast<Derived*>(this));
     858               8 :     op.impl_ptr = this->shared_from_this();
     859                 : 
     860                 :     // Force register_op's ready path so the wait op re-probes under the
     861                 :     // descriptor mutex before parking. A stale write_ready latched at
     862                 :     // registration would otherwise report a full pipe as writable.
     863               8 :     bool force_probe = true;
     864               8 :     register_op(op, *desc_slot_ptr, force_probe, event == reactor_event_write);
     865               8 :     return std::noop_coroutine();
     866                 : }
     867                 : 
     868                 : } // namespace boost::corosio::detail
     869                 : 
     870                 : #endif // BOOST_COROSIO_POSIX
     871                 : 
     872                 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_HPP
        

Generated by: LCOV version 2.3