LCOV - code coverage report
Current view: top level - corosio/native/detail/reactor - reactor_acceptor.hpp (source / functions) Coverage Total Hit Missed
Test: coverage_remapped.info Lines: 98.0 % 246 241 5
Test Date: 2026-10-08 18:13:32 Functions: 97.1 % 136 132 4

           TLA  Line data    Source code
       1                 : //
       2                 : // Copyright (c) 2026 Steve Gerbino
       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_ACCEPTOR_HPP
      11                 : #define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_ACCEPTOR_HPP
      12                 : 
      13                 : #include <boost/corosio/tcp_acceptor.hpp>
      14                 : #include <boost/corosio/wait_type.hpp>
      15                 : #include <boost/corosio/detail/config.hpp>
      16                 : #include <boost/corosio/detail/intrusive.hpp>
      17                 : #include <boost/corosio/native/detail/reactor/reactor_op_base.hpp>
      18                 : #include <boost/corosio/native/detail/reactor/reactor_descriptor_state.hpp>
      19                 : #include <boost/corosio/native/detail/make_err.hpp>
      20                 : #include <boost/corosio/native/detail/endpoint_convert.hpp>
      21                 : 
      22                 : #include <cstring>
      23                 : #include <memory>
      24                 : #include <mutex>
      25                 : #include <utility>
      26                 : 
      27                 : #include <errno.h>
      28                 : #include <netinet/in.h>
      29                 : #include <sys/socket.h>
      30                 : #include <unistd.h>
      31                 : 
      32                 : namespace boost::corosio::detail {
      33                 : 
      34                 : /** CRTP base for reactor-backed acceptor implementations.
      35                 : 
      36                 :     Provides shared data members, trivial virtual overrides, and
      37                 :     non-virtual helper methods for cancellation and close. Concrete
      38                 :     backends inherit and add `cancel()`, `close_socket()`, and
      39                 :     `accept()` overrides that delegate to the `do_*` helpers.
      40                 : 
      41                 :     @tparam Derived   The concrete acceptor type (CRTP).
      42                 :     @tparam Service   The backend's acceptor service type.
      43                 :     @tparam Op        The backend's base op type.
      44                 :     @tparam AcceptOp  The backend's accept op type.
      45                 :     @tparam WaitOp    The backend's wait op type.
      46                 :     @tparam DescState The backend's descriptor_state type.
      47                 :     @tparam ImplBase  The public vtable base
      48                 :                       (tcp_acceptor::implementation or
      49                 :                        local_stream_acceptor::implementation).
      50                 :     @tparam Endpoint  The endpoint type (endpoint or local_endpoint).
      51                 : */
      52                 : template<
      53                 :     class Derived,
      54                 :     class Service,
      55                 :     class Op,
      56                 :     class AcceptOp,
      57                 :     class WaitOp,
      58                 :     class DescState,
      59                 :     class ImplBase = tcp_acceptor::implementation,
      60                 :     class Endpoint = endpoint>
      61                 : class reactor_acceptor
      62                 :     : public ImplBase
      63                 :     , public intrusive_list<Derived>::node
      64                 : {
      65                 :     friend Derived;
      66                 : 
      67                 : protected:
      68                 :     // NOLINTNEXTLINE(bugprone-crtp-constructor-accessibility)
      69 HIT         780 :     explicit reactor_acceptor(Service& svc) noexcept : svc_(svc) {}
      70                 : 
      71                 : protected:
      72                 :     Service& svc_;
      73                 :     int fd_ = -1;
      74                 :     Endpoint local_endpoint_;
      75                 : 
      76                 : public:
      77                 :     /// Pending accept operation slot.
      78                 :     AcceptOp acc_;
      79                 : 
      80                 :     /// Pending wait-for-read operation slot.
      81                 :     WaitOp wait_rd_;
      82                 : 
      83                 :     /// Pending wait-for-write operation slot.
      84                 :     WaitOp wait_wr_;
      85                 : 
      86                 :     /// Pending wait-for-error operation slot.
      87                 :     WaitOp wait_er_;
      88                 : 
      89                 :     /// Per-descriptor state for persistent reactor registration.
      90                 :     DescState desc_state_;
      91                 : 
      92             780 :     ~reactor_acceptor() override = default;
      93                 : 
      94                 :     /// Return the underlying file descriptor.
      95              59 :     native_handle_type native_handle() const noexcept override
      96                 :     {
      97              59 :         return fd_;
      98                 :     }
      99                 : 
     100             673 :     corosio::family family() const noexcept override
     101                 :     {
     102             673 :         return to_family(socket_family(fd_));
     103                 :     }
     104                 : 
     105                 :     /// Release and return the native handle without closing it.
     106              22 :     native_handle_type release_socket() noexcept override
     107                 :     {
     108              22 :         return do_release_socket();
     109                 :     }
     110                 : 
     111                 :     /// Return the cached local endpoint.
     112            5448 :     Endpoint local_endpoint() const noexcept override
     113                 :     {
     114            5448 :         return local_endpoint_;
     115                 :     }
     116                 : 
     117                 :     /// Return true if the acceptor has an open file descriptor.
     118           10085 :     bool is_open() const noexcept override
     119                 :     {
     120           10085 :         return fd_ >= 0;
     121                 :     }
     122                 : 
     123                 :     /// Set a socket option.
     124             648 :     std::error_code set_option(
     125                 :         int level,
     126                 :         int optname,
     127                 :         void const* data,
     128                 :         std::size_t size) noexcept override
     129                 :     {
     130             648 :         if (::setsockopt(
     131             648 :                 fd_, level, optname, data, static_cast<socklen_t>(size)) != 0)
     132              10 :             return make_err(errno);
     133             638 :         return {};
     134                 :     }
     135                 : 
     136                 :     /// Get a socket option.
     137                 :     std::error_code
     138              25 :     get_option(int level, int optname, void* data, std::size_t* size)
     139                 :         const noexcept override
     140                 :     {
     141              25 :         socklen_t len = static_cast<socklen_t>(*size);
     142              25 :         if (::getsockopt(fd_, level, optname, data, &len) != 0)
     143              10 :             return make_err(errno);
     144              15 :         *size = static_cast<std::size_t>(len);
     145              15 :         return {};
     146                 :     }
     147                 : 
     148                 :     /// Cache the local endpoint.
     149             724 :     void set_local_endpoint(Endpoint ep) noexcept
     150                 :     {
     151             724 :         local_endpoint_ = std::move(ep);
     152             724 :     }
     153                 : 
     154                 :     /// Assign the fd and initialize descriptor state for the acceptor.
     155             765 :     void init_acceptor_fd(int fd) noexcept
     156                 :     {
     157             765 :         fd_            = fd;
     158             765 :         desc_state_.fd = fd;
     159                 :         {
     160             765 :             std::lock_guard lock(desc_state_.mutex);
     161             765 :             desc_state_.read_op       = nullptr;
     162             765 :             desc_state_.wait_read_op  = nullptr;
     163             765 :             desc_state_.wait_write_op = nullptr;
     164             765 :             desc_state_.wait_error_op = nullptr;
     165             765 :         }
     166             765 :     }
     167                 : 
     168                 :     /** Assign the fd, initialize descriptor state, and register with
     169                 :         the reactor.
     170                 : 
     171                 :         Adoption skips `do_listen`, so the registration it performs
     172                 :         has to happen here instead.
     173                 : 
     174                 :         @param fd The already-listening descriptor to adopt.
     175                 : 
     176                 :         @return The error if the reactor rejects the descriptor, in
     177                 :         which case the implementation is left closed and the caller
     178                 :         retains ownership of @a fd; otherwise a default constructed
     179                 :         error code.
     180                 :     */
     181              18 :     std::error_code init_and_register(int fd) noexcept
     182                 :     {
     183              18 :         init_acceptor_fd(fd);
     184              18 :         if (auto ec = svc_.scheduler().register_descriptor(fd, &desc_state_))
     185                 :         {
     186               1 :             fd_                           = -1;
     187               1 :             desc_state_.fd                = -1;
     188               1 :             desc_state_.registered_events = 0;
     189               1 :             return ec;
     190                 :         }
     191              17 :         return {};
     192                 :     }
     193                 : 
     194                 :     /// Return a reference to the owning service.
     195            4853 :     Service& service() noexcept
     196                 :     {
     197            4853 :         return svc_;
     198                 :     }
     199                 : 
     200              19 :     void cancel() noexcept override
     201                 :     {
     202              19 :         do_cancel();
     203              19 :     }
     204                 : 
     205                 :     /// Close the acceptor (non-virtual, called by the service).
     206            3410 :     void close_socket() noexcept
     207                 :     {
     208            3410 :         do_close_socket();
     209            3410 :     }
     210                 : 
     211              38 :     std::coroutine_handle<> wait(
     212                 :         std::coroutine_handle<> h,
     213                 :         capy::executor_ref ex,
     214                 :         wait_type w,
     215                 :         std::stop_token token,
     216                 :         std::error_code* ec) override
     217                 :     {
     218              38 :         return do_wait(h, ex, w, token, ec);
     219                 :     }
     220                 : 
     221                 :     /** Wait for readiness on the listen socket.
     222                 : 
     223                 :         For `wait_type::read`, completion signals that an incoming
     224                 :         connection is pending and a subsequent accept succeeds
     225                 :         without blocking; a connection already queued when the wait
     226                 :         begins completes it immediately via an initiation probe.
     227                 : 
     228                 :         `wait_type::write` fails with `operation_not_supported` on
     229                 :         every backend: writability carries no meaning for a
     230                 :         listening socket.
     231                 :     */
     232                 :     std::coroutine_handle<> do_wait(
     233                 :         std::coroutine_handle<>,
     234                 :         capy::executor_ref,
     235                 :         wait_type,
     236                 :         std::stop_token const&,
     237                 :         std::error_code*);
     238                 : 
     239                 :     /** Cancel a single pending operation.
     240                 : 
     241                 :         Claims the operation from the read_op descriptor slot
     242                 :         under the mutex and posts it to the scheduler as cancelled.
     243                 : 
     244                 :         @param op The operation to cancel.
     245                 :     */
     246                 :     void cancel_single_op(Op& op) noexcept;
     247                 : 
     248                 :     /** Cancel the pending accept operation. */
     249                 :     void do_cancel() noexcept;
     250                 : 
     251                 :     /** Close the acceptor and cancel pending operations.
     252                 : 
     253                 :         Invoked by the derived class's close_socket(). The
     254                 :         derived class may add backend-specific cleanup after
     255                 :         calling this method.
     256                 :     */
     257                 :     void do_close_socket() noexcept;
     258                 : 
     259                 :     /** Release the acceptor without closing the fd. */
     260                 :     native_handle_type do_release_socket() noexcept;
     261                 : 
     262                 :     /** Reset descriptor state for recycling.
     263                 : 
     264                 :         Actively re-initializes every field it owns rather than
     265                 :         trusting `close_socket()`'s prior writes to still be there —
     266                 :         see reactor_basic_socket::reuse() for the same rationale.
     267                 :         `object_ref_` and `is_enqueued_` are asserted instead: `poison()`
     268                 :         never touches them, so they must already hold their
     269                 :         zero-action-complete state.
     270                 : 
     271                 :         @pre refs_ == 0, fd closed and deregistered, no op in flight.
     272                 :     */
     273             163 :     void reuse() noexcept
     274                 :     {
     275             163 :         fd_                           = -1;
     276             163 :         local_endpoint_               = Endpoint{};
     277             163 :         desc_state_.fd                = -1;
     278             163 :         desc_state_.registered_events = 0;
     279             163 :         desc_state_.read_op           = nullptr;
     280             163 :         desc_state_.wait_read_op      = nullptr;
     281             163 :         desc_state_.wait_write_op     = nullptr;
     282             163 :         desc_state_.wait_error_op     = nullptr;
     283             163 :         desc_state_.read_ready        = false;
     284             163 :         desc_state_.write_ready       = false;
     285             163 :         BOOST_COROSIO_ASSERT(!desc_state_.object_ref_);
     286             163 :         BOOST_COROSIO_ASSERT(!desc_state_.is_enqueued_.load(
     287                 :             std::memory_order_relaxed));
     288             163 :         BOOST_COROSIO_ASSERT(!acc_.stop_cb);
     289             163 :         BOOST_COROSIO_ASSERT(!wait_rd_.stop_cb);
     290             163 :         BOOST_COROSIO_ASSERT(!wait_wr_.stop_cb);
     291             163 :         BOOST_COROSIO_ASSERT(!wait_er_.stop_cb);
     292             163 :     }
     293                 : 
     294                 : #if !defined(NDEBUG)
     295                 :     /** Poison the fields `reuse()` re-initializes, before the pool
     296                 :         parks this impl on the free list. See
     297                 :         reactor_basic_socket::poison() for the full rationale and the
     298                 :         list of fields deliberately left untouched.
     299                 : 
     300                 :         @pre refs_ == 0, fd closed and deregistered, no op in flight.
     301                 :     */
     302             938 :     void poison() noexcept
     303                 :     {
     304            9380 :         auto smash = [](auto& field) {
     305            9380 :             std::memset(static_cast<void*>(&field), 0xDB, sizeof(field));
     306                 :         };
     307             938 :         smash(fd_);
     308             938 :         smash(local_endpoint_);
     309             938 :         smash(desc_state_.fd);
     310             938 :         smash(desc_state_.registered_events);
     311             938 :         smash(desc_state_.read_op);
     312             938 :         smash(desc_state_.wait_read_op);
     313             938 :         smash(desc_state_.wait_write_op);
     314             938 :         smash(desc_state_.wait_error_op);
     315             938 :         smash(desc_state_.read_ready);
     316             938 :         smash(desc_state_.write_ready);
     317             938 :     }
     318                 : #endif
     319                 : 
     320                 :     /** Bind the acceptor socket to an endpoint.
     321                 : 
     322                 :         Caches the resolved local endpoint (including ephemeral
     323                 :         port) after a successful bind.
     324                 : 
     325                 :         @param ep The endpoint to bind to.
     326                 :         @return The error code from bind(), or success.
     327                 :     */
     328                 :     std::error_code do_bind(Endpoint const& ep);
     329                 : 
     330                 :     /** Start listening on the acceptor socket.
     331                 : 
     332                 :         Registers the file descriptor with the reactor after
     333                 :         a successful listen() call.
     334                 : 
     335                 :         @param backlog The listen backlog.
     336                 :         @return The error code from listen() or from reactor
     337                 :         registration, or success.
     338                 :     */
     339                 :     std::error_code do_listen(int backlog);
     340                 : };
     341                 : 
     342                 : template<
     343                 :     class Derived,
     344                 :     class Service,
     345                 :     class Op,
     346                 :     class AcceptOp,
     347                 :     class WaitOp,
     348                 :     class DescState,
     349                 :     class ImplBase,
     350                 :     class Endpoint>
     351                 : void
     352             162 : reactor_acceptor<
     353                 :     Derived,
     354                 :     Service,
     355                 :     Op,
     356                 :     AcceptOp,
     357                 :     WaitOp,
     358                 :     DescState,
     359                 :     ImplBase,
     360                 :     Endpoint>::cancel_single_op(Op& op) noexcept
     361                 : {
     362             162 :     op.request_cancel();
     363                 : 
     364             162 :     reactor_op_base* claimed = nullptr;
     365                 :     {
     366             162 :         std::lock_guard lock(desc_state_.mutex);
     367            1458 :         auto try_claim = [&](reactor_op_base*& slot) {
     368             648 :             if (!claimed && slot == &op)
     369             103 :                 claimed = std::exchange(slot, nullptr);
     370                 :         };
     371             162 :         try_claim(desc_state_.read_op);
     372             162 :         try_claim(desc_state_.wait_read_op);
     373             162 :         try_claim(desc_state_.wait_write_op);
     374             162 :         try_claim(desc_state_.wait_error_op);
     375             162 :     }
     376             162 :     if (claimed)
     377                 :     {
     378             103 :         op.object_ref_ = detail::object_ref(this);
     379             103 :         svc_.post(&op);
     380             103 :         svc_.work_finished();
     381                 :     }
     382             162 : }
     383                 : 
     384                 : template<
     385                 :     class Derived,
     386                 :     class Service,
     387                 :     class Op,
     388                 :     class AcceptOp,
     389                 :     class WaitOp,
     390                 :     class DescState,
     391                 :     class ImplBase,
     392                 :     class Endpoint>
     393                 : void
     394              19 : reactor_acceptor<
     395                 :     Derived,
     396                 :     Service,
     397                 :     Op,
     398                 :     AcceptOp,
     399                 :     WaitOp,
     400                 :     DescState,
     401                 :     ImplBase,
     402                 :     Endpoint>::do_cancel() noexcept
     403                 : {
     404              19 :     cancel_single_op(acc_);
     405              19 :     cancel_single_op(wait_rd_);
     406              19 :     cancel_single_op(wait_wr_);
     407              19 :     cancel_single_op(wait_er_);
     408              19 : }
     409                 : 
     410                 : template<
     411                 :     class Derived,
     412                 :     class Service,
     413                 :     class Op,
     414                 :     class AcceptOp,
     415                 :     class WaitOp,
     416                 :     class DescState,
     417                 :     class ImplBase,
     418                 :     class Endpoint>
     419                 : void
     420            3410 : reactor_acceptor<
     421                 :     Derived,
     422                 :     Service,
     423                 :     Op,
     424                 :     AcceptOp,
     425                 :     WaitOp,
     426                 :     DescState,
     427                 :     ImplBase,
     428                 :     Endpoint>::do_close_socket() noexcept
     429                 : {
     430            3410 :     acc_.request_cancel();
     431            3410 :     wait_rd_.request_cancel();
     432            3410 :     wait_wr_.request_cancel();
     433            3410 :     wait_er_.request_cancel();
     434                 : 
     435            3410 :     reactor_op_base* claimed_acc = nullptr;
     436            3410 :     reactor_op_base* claimed_wr  = nullptr;
     437            3410 :     reactor_op_base* claimed_ww  = nullptr;
     438            3410 :     reactor_op_base* claimed_we  = nullptr;
     439                 :     {
     440            3410 :         std::lock_guard lock(desc_state_.mutex);
     441            3410 :         claimed_acc = std::exchange(desc_state_.read_op, nullptr);
     442            3410 :         claimed_wr  = std::exchange(desc_state_.wait_read_op, nullptr);
     443            3410 :         claimed_ww  = std::exchange(desc_state_.wait_write_op, nullptr);
     444            3410 :         claimed_we  = std::exchange(desc_state_.wait_error_op, nullptr);
     445            3410 :         desc_state_.read_ready  = false;
     446            3410 :         desc_state_.write_ready = false;
     447                 : 
     448            3410 :         if (desc_state_.is_enqueued_.load(std::memory_order_acquire))
     449              48 :             desc_state_.object_ref_ = detail::object_ref(this);
     450            3410 :     }
     451                 : 
     452           30690 :     auto repost = [&](reactor_op_base* claimed, reactor_op_base& op) {
     453           13640 :         if (claimed)
     454                 :         {
     455              19 :             op.object_ref_ = detail::object_ref(this);
     456              19 :             svc_.post(&op);
     457              19 :             svc_.work_finished();
     458                 :         }
     459                 :     };
     460            3410 :     repost(claimed_acc, acc_);
     461            3410 :     repost(claimed_wr, wait_rd_);
     462            3410 :     repost(claimed_ww, wait_wr_);
     463            3410 :     repost(claimed_we, wait_er_);
     464                 : 
     465            3410 :     if (fd_ >= 0)
     466                 :     {
     467             742 :         if (desc_state_.registered_events != 0)
     468             659 :             svc_.scheduler().deregister_descriptor(fd_);
     469             742 :         ::close(fd_);
     470             742 :         fd_ = -1;
     471                 :     }
     472                 : 
     473            3410 :     desc_state_.fd                = -1;
     474            3410 :     desc_state_.registered_events = 0;
     475                 : 
     476            3410 :     local_endpoint_ = Endpoint{};
     477            3410 : }
     478                 : 
     479                 : template<
     480                 :     class Derived,
     481                 :     class Service,
     482                 :     class Op,
     483                 :     class AcceptOp,
     484                 :     class WaitOp,
     485                 :     class DescState,
     486                 :     class ImplBase,
     487                 :     class Endpoint>
     488                 : native_handle_type
     489              22 : reactor_acceptor<
     490                 :     Derived,
     491                 :     Service,
     492                 :     Op,
     493                 :     AcceptOp,
     494                 :     WaitOp,
     495                 :     DescState,
     496                 :     ImplBase,
     497                 :     Endpoint>::do_release_socket() noexcept
     498                 : {
     499              22 :     acc_.request_cancel();
     500              22 :     wait_rd_.request_cancel();
     501              22 :     wait_wr_.request_cancel();
     502              22 :     wait_er_.request_cancel();
     503                 : 
     504              22 :     reactor_op_base* claimed_acc = nullptr;
     505              22 :     reactor_op_base* claimed_wr  = nullptr;
     506              22 :     reactor_op_base* claimed_ww  = nullptr;
     507              22 :     reactor_op_base* claimed_we  = nullptr;
     508                 :     {
     509              22 :         std::lock_guard lock(desc_state_.mutex);
     510              22 :         claimed_acc = std::exchange(desc_state_.read_op, nullptr);
     511              22 :         claimed_wr  = std::exchange(desc_state_.wait_read_op, nullptr);
     512              22 :         claimed_ww  = std::exchange(desc_state_.wait_write_op, nullptr);
     513              22 :         claimed_we  = std::exchange(desc_state_.wait_error_op, nullptr);
     514              22 :         desc_state_.read_ready  = false;
     515              22 :         desc_state_.write_ready = false;
     516                 : 
     517              22 :         if (desc_state_.is_enqueued_.load(std::memory_order_acquire))
     518               1 :             desc_state_.object_ref_ = detail::object_ref(this);
     519              22 :     }
     520                 : 
     521             198 :     auto repost = [&](reactor_op_base* claimed, reactor_op_base& op) {
     522              88 :         if (claimed)
     523                 :         {
     524               5 :             op.object_ref_ = detail::object_ref(this);
     525               5 :             svc_.post(&op);
     526               5 :             svc_.work_finished();
     527                 :         }
     528                 :     };
     529              22 :     repost(claimed_acc, acc_);
     530              22 :     repost(claimed_wr, wait_rd_);
     531              22 :     repost(claimed_ww, wait_wr_);
     532              22 :     repost(claimed_we, wait_er_);
     533                 : 
     534              22 :     native_handle_type released = fd_;
     535                 : 
     536              22 :     if (fd_ >= 0)
     537                 :     {
     538              22 :         if (desc_state_.registered_events != 0)
     539              22 :             svc_.scheduler().deregister_descriptor(fd_);
     540              22 :         fd_ = -1;
     541                 :     }
     542                 : 
     543              22 :     desc_state_.fd                = -1;
     544              22 :     desc_state_.registered_events = 0;
     545                 : 
     546              22 :     local_endpoint_ = Endpoint{};
     547                 : 
     548              22 :     return released;
     549                 : }
     550                 : 
     551                 : template<
     552                 :     class Derived,
     553                 :     class Service,
     554                 :     class Op,
     555                 :     class AcceptOp,
     556                 :     class WaitOp,
     557                 :     class DescState,
     558                 :     class ImplBase,
     559                 :     class Endpoint>
     560                 : std::error_code
     561             723 : reactor_acceptor<
     562                 :     Derived,
     563                 :     Service,
     564                 :     Op,
     565                 :     AcceptOp,
     566                 :     WaitOp,
     567                 :     DescState,
     568                 :     ImplBase,
     569                 :     Endpoint>::do_bind(Endpoint const& ep)
     570                 : {
     571             723 :     sockaddr_storage storage{};
     572             723 :     socklen_t addrlen = to_sockaddr(ep, storage);
     573             723 :     if (::bind(fd_, reinterpret_cast<sockaddr*>(&storage), addrlen) < 0)
     574              16 :         return make_err(errno);
     575                 : 
     576                 :     // Cache local endpoint (resolves ephemeral port / path)
     577             707 :     sockaddr_storage local{};
     578             707 :     socklen_t local_len = sizeof(local);
     579             707 :     if (::getsockname(fd_, reinterpret_cast<sockaddr*>(&local), &local_len) ==
     580                 :         0)
     581             707 :         set_local_endpoint(from_sockaddr_as(local, local_len, Endpoint{}));
     582                 : 
     583             707 :     return {};
     584                 : }
     585                 : 
     586                 : template<
     587                 :     class Derived,
     588                 :     class Service,
     589                 :     class Op,
     590                 :     class AcceptOp,
     591                 :     class WaitOp,
     592                 :     class DescState,
     593                 :     class ImplBase,
     594                 :     class Endpoint>
     595                 : std::error_code
     596             679 : reactor_acceptor<
     597                 :     Derived,
     598                 :     Service,
     599                 :     Op,
     600                 :     AcceptOp,
     601                 :     WaitOp,
     602                 :     DescState,
     603                 :     ImplBase,
     604                 :     Endpoint>::do_listen(int backlog)
     605                 : {
     606             679 :     if (::listen(fd_, backlog) < 0)
     607              12 :         return make_err(errno);
     608                 : 
     609                 :     // A re-listen only changes the backlog; the descriptor is already
     610                 :     // registered and re-adding it would fail on epoll.
     611             667 :     if (desc_state_.registered_events != 0)
     612               2 :         return {};
     613                 : 
     614             665 :     return svc_.scheduler().register_descriptor(fd_, &desc_state_);
     615                 : }
     616                 : 
     617                 : template<
     618                 :     class Derived,
     619                 :     class Service,
     620                 :     class Op,
     621                 :     class AcceptOp,
     622                 :     class WaitOp,
     623                 :     class DescState,
     624                 :     class ImplBase,
     625                 :     class Endpoint>
     626                 : std::coroutine_handle<>
     627              38 : reactor_acceptor<
     628                 :     Derived,
     629                 :     Service,
     630                 :     Op,
     631                 :     AcceptOp,
     632                 :     WaitOp,
     633                 :     DescState,
     634                 :     ImplBase,
     635                 :     Endpoint>::
     636                 :     do_wait(
     637                 :         std::coroutine_handle<> h,
     638                 :         capy::executor_ref ex,
     639                 :         wait_type w,
     640                 :         std::stop_token const& token,
     641                 :         std::error_code* ec)
     642                 : {
     643                 :     // Writability carries no meaning for a listening socket; some
     644                 :     // backends could only lie about it and others could never report
     645                 :     // it, so the wait fails the same way everywhere instead.
     646              38 :     if (w == wait_type::write)
     647                 :     {
     648               6 :         auto& op = wait_wr_;
     649               6 :         op.reset();
     650               6 :         op.wait_event = reactor_event_write;
     651               6 :         op.h          = h;
     652               6 :         op.ex         = ex;
     653               6 :         op.ec_out     = ec;
     654               6 :         op.fd         = this->fd_;
     655               6 :         op.start(token, static_cast<Derived*>(this));
     656               6 :         op.object_ref_ = detail::object_ref(this);
     657               6 :         op.complete(ENOTSUP, 0);
     658               6 :         svc_.post(&op);
     659               6 :         return std::noop_coroutine();
     660                 :     }
     661                 : 
     662                 :     WaitOp* op_ptr;
     663                 :     reactor_op_base** desc_slot_ptr;
     664                 :     std::uint32_t event;
     665                 : 
     666              32 :     if (w == wait_type::read)
     667                 :     {
     668              26 :         op_ptr        = &wait_rd_;
     669              26 :         desc_slot_ptr = &desc_state_.wait_read_op;
     670              26 :         event         = reactor_event_read;
     671                 :     }
     672                 :     else // wait_type::error
     673                 :     {
     674               6 :         op_ptr        = &wait_er_;
     675               6 :         desc_slot_ptr = &desc_state_.wait_error_op;
     676               6 :         event         = reactor_event_error;
     677                 :     }
     678                 : 
     679              32 :     auto& op = *op_ptr;
     680              32 :     op.reset();
     681              32 :     op.wait_event = event;
     682              32 :     op.h          = h;
     683              32 :     op.ex         = ex;
     684              32 :     op.ec_out     = ec;
     685              32 :     op.fd         = this->fd_;
     686              32 :     op.start(token, static_cast<Derived*>(this));
     687              32 :     op.object_ref_ = detail::object_ref(this);
     688                 : 
     689                 :     // A listener's readiness can predate the wait: an adopted or
     690                 :     // shared descriptor has history the reactor never saw, and an
     691                 :     // edge already dispatched will not be re-announced. Probe before
     692                 :     // parking.
     693              32 :     int perr = 0;
     694              32 :     if (WaitOp::probe(this->fd_, event, perr))
     695                 :     {
     696              11 :         op.complete(perr, 0);
     697              11 :         svc_.post(&op);
     698              11 :         return std::noop_coroutine();
     699                 :     }
     700                 : 
     701              21 :     svc_.work_started();
     702                 : 
     703              21 :     std::lock_guard lock(desc_state_.mutex);
     704              21 :     if (op.cancelled.load(std::memory_order_acquire))
     705                 :     {
     706 MIS           0 :         svc_.post(&op);
     707               0 :         svc_.work_finished();
     708                 :     }
     709 HIT          21 :     else if (WaitOp::probe(this->fd_, event, perr))
     710                 :     {
     711                 :         // Close the probe-to-park window: an edge that landed after
     712                 :         // the first probe was consumed, so re-check under the mutex
     713                 :         // the dispatch path holds.
     714 MIS           0 :         op.complete(perr, 0);
     715               0 :         svc_.post(&op);
     716               0 :         svc_.work_finished();
     717                 :     }
     718                 :     else
     719                 :     {
     720 HIT          21 :         *desc_slot_ptr = &op;
     721                 :     }
     722              21 :     return std::noop_coroutine();
     723              21 : }
     724                 : 
     725                 : } // namespace boost::corosio::detail
     726                 : 
     727                 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_ACCEPTOR_HPP
        

Generated by: LCOV version 2.3