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: 92.3 % 156 144 12
Test Date: 2026-10-08 17:58:26 Functions: 96.6 % 116 112 4

           TLA  Line data    Source code
       1                 : //
       2                 : // Copyright (c) 2026 Steve Gerbino
       3                 : // Copyright (c) 2026 Michael Vandeberg
       4                 : //
       5                 : // Distributed under the Boost Software License, Version 1.0. (See accompanying
       6                 : // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
       7                 : //
       8                 : // Official repository: https://github.com/cppalliance/corosio
       9                 : //
      10                 : 
      11                 : #ifndef BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_ACCEPTOR_HPP
      12                 : #define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_ACCEPTOR_HPP
      13                 : 
      14                 : #include <boost/corosio/tcp_acceptor.hpp>
      15                 : #include <boost/corosio/wait_type.hpp>
      16                 : #include <boost/corosio/detail/intrusive.hpp>
      17                 : #include <boost/corosio/native/detail/reactor/reactor_io_core.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 <memory>
      23                 : #include <mutex>
      24                 : #include <utility>
      25                 : 
      26                 : #include <errno.h>
      27                 : #include <netinet/in.h>
      28                 : #include <sys/socket.h>
      29                 : #include <unistd.h>
      30                 : 
      31                 : namespace boost::corosio::detail {
      32                 : 
      33                 : /** CRTP base for reactor-backed acceptor implementations.
      34                 : 
      35                 :     Provides shared data members, trivial virtual overrides, and
      36                 :     non-virtual helper methods for cancellation and close. Concrete
      37                 :     backends inherit and add `cancel()`, `close_socket()`, and
      38                 :     `accept()` overrides that delegate to the `do_*` helpers.
      39                 : 
      40                 :     @tparam Derived   The concrete acceptor type (CRTP).
      41                 :     @tparam Service   The backend's acceptor service type.
      42                 :     @tparam Op        The backend's base op type.
      43                 :     @tparam AcceptOp  The backend's accept op type.
      44                 :     @tparam WaitOp    The backend's wait op type.
      45                 :     @tparam DescState The backend's descriptor_state type.
      46                 :     @tparam ImplBase  The public vtable base
      47                 :                       (tcp_acceptor::implementation or
      48                 :                        local_stream_acceptor::implementation).
      49                 :     @tparam Endpoint  The endpoint type (endpoint or local_endpoint).
      50                 : */
      51                 : template<
      52                 :     class Derived,
      53                 :     class Service,
      54                 :     class Op,
      55                 :     class AcceptOp,
      56                 :     class WaitOp,
      57                 :     class DescState,
      58                 :     class ImplBase = tcp_acceptor::implementation,
      59                 :     class Endpoint = endpoint>
      60                 : class reactor_acceptor
      61                 :     : public ImplBase
      62                 :     , public std::enable_shared_from_this<Derived>
      63                 :     , public reactor_io_core<Derived, Service, DescState>
      64                 :     , public intrusive_list<Derived>::node
      65                 : {
      66                 :     friend Derived;
      67                 : 
      68                 :     using core_type = reactor_io_core<Derived, Service, DescState>;
      69                 :     friend core_type;
      70                 : 
      71                 : protected:
      72                 :     // NOLINTNEXTLINE(bugprone-crtp-constructor-accessibility)
      73 HIT         788 :     explicit reactor_acceptor(Service& svc) noexcept : core_type(svc) {}
      74                 : 
      75                 : protected:
      76                 :     using core_type::svc_;
      77                 :     int fd_ = -1;
      78                 :     Endpoint local_endpoint_;
      79                 : 
      80                 : public:
      81                 :     /// Pending accept operation slot.
      82                 :     AcceptOp acc_;
      83                 : 
      84                 :     /// Pending wait-for-read operation slot.
      85                 :     WaitOp wait_rd_;
      86                 : 
      87                 :     /// Pending wait-for-write operation slot.
      88                 :     WaitOp wait_wr_;
      89                 : 
      90                 :     /// Pending wait-for-error operation slot.
      91                 :     WaitOp wait_er_;
      92                 : 
      93                 :     using core_type::desc_state_;
      94                 : 
      95             788 :     ~reactor_acceptor() override = default;
      96                 : 
      97                 :     /// Return the underlying file descriptor.
      98              44 :     native_handle_type native_handle() const noexcept override
      99                 :     {
     100              44 :         return fd_;
     101                 :     }
     102                 : 
     103             642 :     corosio::family family() const noexcept override
     104                 :     {
     105             642 :         return to_family(socket_family(fd_));
     106                 :     }
     107                 : 
     108                 :     /// Release and return the native handle without closing it.
     109              20 :     native_handle_type release_socket() noexcept override
     110                 :     {
     111              20 :         return do_release_socket();
     112                 :     }
     113                 : 
     114                 :     /// Return the cached local endpoint.
     115            5169 :     Endpoint local_endpoint() const noexcept override
     116                 :     {
     117            5169 :         return local_endpoint_;
     118                 :     }
     119                 : 
     120                 :     /// Return true if the acceptor has an open file descriptor.
     121            9558 :     bool is_open() const noexcept override
     122                 :     {
     123            9558 :         return fd_ >= 0;
     124                 :     }
     125                 : 
     126                 :     /// Set a socket option.
     127             617 :     std::error_code set_option(
     128                 :         int level,
     129                 :         int optname,
     130                 :         void const* data,
     131                 :         std::size_t size) noexcept override
     132                 :     {
     133             617 :         if (::setsockopt(
     134             617 :                 fd_, level, optname, data, static_cast<socklen_t>(size)) != 0)
     135              10 :             return make_err(errno);
     136             607 :         return {};
     137                 :     }
     138                 : 
     139                 :     /// Get a socket option.
     140                 :     std::error_code
     141              25 :     get_option(int level, int optname, void* data, std::size_t* size)
     142                 :         const noexcept override
     143                 :     {
     144              25 :         socklen_t len = static_cast<socklen_t>(*size);
     145              25 :         if (::getsockopt(fd_, level, optname, data, &len) != 0)
     146              10 :             return make_err(errno);
     147              15 :         *size = static_cast<std::size_t>(len);
     148              15 :         return {};
     149                 :     }
     150                 : 
     151                 :     /// Cache the local endpoint.
     152             689 :     void set_local_endpoint(Endpoint ep) noexcept
     153                 :     {
     154             689 :         local_endpoint_ = std::move(ep);
     155             689 :     }
     156                 : 
     157                 :     /// Assign the fd and initialize descriptor state for the acceptor.
     158             720 :     void init_acceptor_fd(int fd) noexcept
     159                 :     {
     160             720 :         fd_            = fd;
     161             720 :         desc_state_.fd = fd;
     162                 :         {
     163             720 :             std::lock_guard lock(desc_state_.mutex);
     164             720 :             desc_state_.read_op       = nullptr;
     165             720 :             desc_state_.wait_read_op  = nullptr;
     166             720 :             desc_state_.wait_write_op = nullptr;
     167             720 :             desc_state_.wait_error_op = nullptr;
     168             720 :         }
     169             720 :     }
     170                 : 
     171                 :     /** Assign the fd, initialize descriptor state, and register with
     172                 :         the reactor.
     173                 : 
     174                 :         Adoption skips `do_listen`, so the registration it performs
     175                 :         has to happen here instead.
     176                 : 
     177                 :         @param fd The already-listening descriptor to adopt.
     178                 : 
     179                 :         @return The error if the reactor rejects the descriptor, in
     180                 :         which case the implementation is left closed and the caller
     181                 :         retains ownership of @a fd; otherwise a default constructed
     182                 :         error code.
     183                 :     */
     184              14 :     std::error_code init_and_register(int fd) noexcept
     185                 :     {
     186              14 :         fd_ = fd;
     187              14 :         if (auto ec = this->register_fd(fd))
     188                 :         {
     189               1 :             fd_ = -1;
     190               1 :             return ec;
     191                 :         }
     192              13 :         return {};
     193                 :     }
     194                 : 
     195                 :     /// Return a reference to the owning service.
     196            4603 :     Service& service() noexcept
     197                 :     {
     198            4603 :         return svc_;
     199                 :     }
     200                 : 
     201              21 :     void cancel() noexcept override
     202                 :     {
     203              21 :         do_cancel();
     204              21 :     }
     205                 : 
     206                 :     /// Close the acceptor (non-virtual, called by the service).
     207            3012 :     void close_socket() noexcept
     208                 :     {
     209            3012 :         do_close_socket();
     210            3012 :     }
     211                 : 
     212              41 :     std::coroutine_handle<> wait(
     213                 :         std::coroutine_handle<> h,
     214                 :         capy::executor_ref ex,
     215                 :         wait_type w,
     216                 :         std::stop_token token,
     217                 :         std::error_code* ec) override
     218                 :     {
     219              41 :         return do_wait(h, ex, w, token, ec);
     220                 :     }
     221                 : 
     222                 :     /** Wait for readiness on the listen socket.
     223                 : 
     224                 :         For `wait_type::read`, completion signals that an incoming
     225                 :         connection is pending and a subsequent accept succeeds
     226                 :         without blocking; a connection already queued when the wait
     227                 :         begins completes it immediately via an initiation probe.
     228                 : 
     229                 :         `wait_type::write` fails with `operation_not_supported` on
     230                 :         every backend: writability carries no meaning for a
     231                 :         listening socket.
     232                 :     */
     233                 :     std::coroutine_handle<> do_wait(
     234                 :         std::coroutine_handle<>,
     235                 :         capy::executor_ref,
     236                 :         wait_type,
     237                 :         std::stop_token const&,
     238                 :         std::error_code*);
     239                 : 
     240                 :     /** Cancel the pending accept operation. */
     241              21 :     void do_cancel() noexcept
     242                 :     {
     243              21 :         this->cancel_all();
     244              21 :     }
     245                 : 
     246                 :     /** Close the acceptor and cancel pending operations.
     247                 : 
     248                 :         Invoked by the derived class's close_socket(). The
     249                 :         derived class may add backend-specific cleanup after
     250                 :         calling this method.
     251                 :     */
     252            3012 :     void do_close_socket() noexcept
     253                 :     {
     254            3012 :         this->abandon_all();
     255            3012 :         this->unregister_fd(fd_);
     256            3012 :         if (fd_ >= 0)
     257                 :         {
     258             713 :             ::close(fd_);
     259             713 :             fd_ = -1;
     260                 :         }
     261            3012 :         local_endpoint_ = Endpoint{};
     262            3012 :     }
     263                 : 
     264                 :     /** Release the acceptor without closing the fd. */
     265              20 :     native_handle_type do_release_socket() noexcept
     266                 :     {
     267              20 :         this->abandon_all();
     268              20 :         native_handle_type released = fd_;
     269              20 :         this->unregister_fd(fd_);
     270              20 :         fd_             = -1;
     271              20 :         local_endpoint_ = Endpoint{};
     272              20 :         return released;
     273                 :     }
     274                 : 
     275                 :     /** Bind the acceptor socket to an endpoint.
     276                 : 
     277                 :         Caches the resolved local endpoint (including ephemeral
     278                 :         port) after a successful bind.
     279                 : 
     280                 :         @param ep The endpoint to bind to.
     281                 :         @return The error code from bind(), or success.
     282                 :     */
     283                 :     std::error_code do_bind(Endpoint const& ep);
     284                 : 
     285                 :     /** Start listening on the acceptor socket.
     286                 : 
     287                 :         Registers the file descriptor with the reactor after
     288                 :         a successful listen() call.
     289                 : 
     290                 :         @param backlog The listen backlog.
     291                 :         @return The error code from listen() or from reactor
     292                 :         registration, or success.
     293                 :     */
     294                 :     std::error_code do_listen(int backlog);
     295                 : 
     296                 : private:
     297                 :     // CRTP callbacks for reactor_io_core cancel/close
     298                 : 
     299                 :     template<class AnyOp>
     300              70 :     reactor_op_base** op_to_desc_slot(AnyOp& op) noexcept
     301                 :     {
     302              70 :         if (&op == static_cast<void*>(&acc_))
     303              70 :             return &desc_state_.read_op;
     304 MIS           0 :         if (&op == static_cast<void*>(&wait_rd_))
     305               0 :             return &desc_state_.wait_read_op;
     306               0 :         if (&op == static_cast<void*>(&wait_wr_))
     307               0 :             return &desc_state_.wait_write_op;
     308               0 :         if (&op == static_cast<void*>(&wait_er_))
     309               0 :             return &desc_state_.wait_error_op;
     310               0 :         return nullptr;
     311                 :     }
     312                 : 
     313                 :     template<class Fn>
     314 HIT        3053 :     void for_each_op(Fn fn) noexcept
     315                 :     {
     316            3053 :         fn(acc_);
     317            3053 :         fn(wait_rd_);
     318            3053 :         fn(wait_wr_);
     319            3053 :         fn(wait_er_);
     320            3053 :     }
     321                 : 
     322                 :     template<class Fn>
     323            3067 :     void for_each_desc_entry(Fn fn) noexcept
     324                 :     {
     325            3067 :         fn(acc_, desc_state_.read_op);
     326            3067 :         fn(wait_rd_, desc_state_.wait_read_op);
     327            3067 :         fn(wait_wr_, desc_state_.wait_write_op);
     328            3067 :         fn(wait_er_, desc_state_.wait_error_op);
     329            3067 :     }
     330                 : };
     331                 : 
     332                 : template<
     333                 :     class Derived,
     334                 :     class Service,
     335                 :     class Op,
     336                 :     class AcceptOp,
     337                 :     class WaitOp,
     338                 :     class DescState,
     339                 :     class ImplBase,
     340                 :     class Endpoint>
     341                 : std::error_code
     342             692 : reactor_acceptor<
     343                 :     Derived,
     344                 :     Service,
     345                 :     Op,
     346                 :     AcceptOp,
     347                 :     WaitOp,
     348                 :     DescState,
     349                 :     ImplBase,
     350                 :     Endpoint>::do_bind(Endpoint const& ep)
     351                 : {
     352             692 :     sockaddr_storage storage{};
     353             692 :     socklen_t addrlen = to_sockaddr(ep, storage);
     354             692 :     if (::bind(fd_, reinterpret_cast<sockaddr*>(&storage), addrlen) < 0)
     355              16 :         return make_err(errno);
     356                 : 
     357                 :     // Cache local endpoint (resolves ephemeral port / path)
     358             676 :     sockaddr_storage local{};
     359             676 :     socklen_t local_len = sizeof(local);
     360             676 :     if (::getsockname(fd_, reinterpret_cast<sockaddr*>(&local), &local_len) ==
     361                 :         0)
     362             676 :         set_local_endpoint(from_sockaddr_as(local, local_len, Endpoint{}));
     363                 : 
     364             676 :     return {};
     365                 : }
     366                 : 
     367                 : template<
     368                 :     class Derived,
     369                 :     class Service,
     370                 :     class Op,
     371                 :     class AcceptOp,
     372                 :     class WaitOp,
     373                 :     class DescState,
     374                 :     class ImplBase,
     375                 :     class Endpoint>
     376                 : std::error_code
     377             648 : reactor_acceptor<
     378                 :     Derived,
     379                 :     Service,
     380                 :     Op,
     381                 :     AcceptOp,
     382                 :     WaitOp,
     383                 :     DescState,
     384                 :     ImplBase,
     385                 :     Endpoint>::do_listen(int backlog)
     386                 : {
     387             648 :     if (::listen(fd_, backlog) < 0)
     388              12 :         return make_err(errno);
     389                 : 
     390                 :     // A re-listen only changes the backlog; the descriptor is already
     391                 :     // registered and re-adding it would fail on epoll.
     392             636 :     if (desc_state_.registered_events != 0)
     393               2 :         return {};
     394                 : 
     395             634 :     return svc_.scheduler().register_descriptor(fd_, &desc_state_);
     396                 : }
     397                 : 
     398                 : template<
     399                 :     class Derived,
     400                 :     class Service,
     401                 :     class Op,
     402                 :     class AcceptOp,
     403                 :     class WaitOp,
     404                 :     class DescState,
     405                 :     class ImplBase,
     406                 :     class Endpoint>
     407                 : std::coroutine_handle<>
     408              41 : reactor_acceptor<
     409                 :     Derived,
     410                 :     Service,
     411                 :     Op,
     412                 :     AcceptOp,
     413                 :     WaitOp,
     414                 :     DescState,
     415                 :     ImplBase,
     416                 :     Endpoint>::
     417                 :     do_wait(
     418                 :         std::coroutine_handle<> h,
     419                 :         capy::executor_ref ex,
     420                 :         wait_type w,
     421                 :         std::stop_token const& token,
     422                 :         std::error_code* ec)
     423                 : {
     424                 :     // Writability carries no meaning for a listening socket; some
     425                 :     // backends could only lie about it and others could never report
     426                 :     // it, so the wait fails the same way everywhere instead.
     427              41 :     if (w == wait_type::write)
     428                 :     {
     429               6 :         auto& op = wait_wr_;
     430               6 :         op.reset();
     431               6 :         op.wait_event = reactor_event_write;
     432               6 :         op.h          = h;
     433               6 :         op.ex         = ex;
     434               6 :         op.ec_out     = ec;
     435               6 :         op.fd         = this->fd_;
     436               6 :         op.start(token, static_cast<Derived*>(this));
     437               6 :         op.impl_ptr = this->shared_from_this();
     438               6 :         op.complete(ENOTSUP, 0);
     439               6 :         svc_.post(&op);
     440               6 :         return std::noop_coroutine();
     441                 :     }
     442                 : 
     443                 :     WaitOp* op_ptr;
     444                 :     reactor_op_base** desc_slot_ptr;
     445                 :     std::uint32_t event;
     446                 : 
     447              35 :     if (w == wait_type::read)
     448                 :     {
     449              29 :         op_ptr        = &wait_rd_;
     450              29 :         desc_slot_ptr = &desc_state_.wait_read_op;
     451              29 :         event         = reactor_event_read;
     452                 :     }
     453                 :     else // wait_type::error
     454                 :     {
     455               6 :         op_ptr        = &wait_er_;
     456               6 :         desc_slot_ptr = &desc_state_.wait_error_op;
     457               6 :         event         = reactor_event_error;
     458                 :     }
     459                 : 
     460              35 :     auto& op = *op_ptr;
     461              35 :     op.reset();
     462              35 :     op.wait_event = event;
     463              35 :     op.h          = h;
     464              35 :     op.ex         = ex;
     465              35 :     op.ec_out     = ec;
     466              35 :     op.fd         = this->fd_;
     467              35 :     op.start(token, static_cast<Derived*>(this));
     468              35 :     op.impl_ptr = this->shared_from_this();
     469                 : 
     470                 :     // A listener's readiness can predate the wait: an adopted or
     471                 :     // shared descriptor has history the reactor never saw, and an
     472                 :     // edge already dispatched will not be re-announced. Probe before
     473                 :     // parking.
     474              35 :     int perr = 0;
     475              35 :     if (WaitOp::probe(this->fd_, event, perr))
     476                 :     {
     477              12 :         op.complete(perr, 0);
     478              12 :         svc_.post(&op);
     479              12 :         return std::noop_coroutine();
     480                 :     }
     481                 : 
     482              23 :     svc_.work_started();
     483                 : 
     484              23 :     std::lock_guard lock(desc_state_.mutex);
     485              23 :     if (op.cancelled.load(std::memory_order_acquire))
     486                 :     {
     487 MIS           0 :         svc_.post(&op);
     488               0 :         svc_.work_finished();
     489                 :     }
     490 HIT          23 :     else if (WaitOp::probe(this->fd_, event, perr))
     491                 :     {
     492                 :         // Close the probe-to-park window: an edge that landed after
     493                 :         // the first probe was consumed, so re-check under the mutex
     494                 :         // the dispatch path holds.
     495 MIS           0 :         op.complete(perr, 0);
     496               0 :         svc_.post(&op);
     497               0 :         svc_.work_finished();
     498                 :     }
     499                 :     else
     500                 :     {
     501 HIT          23 :         *desc_slot_ptr = &op;
     502                 : 
     503                 :         // Select watches an fd only while an op is parked; see
     504                 :         // register_op.
     505                 :         if constexpr (Service::needs_park_notification)
     506              11 :             svc_.scheduler().notify_reactor();
     507                 :     }
     508              23 :     return std::noop_coroutine();
     509              23 : }
     510                 : 
     511                 : } // namespace boost::corosio::detail
     512                 : 
     513                 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_ACCEPTOR_HPP
        

Generated by: LCOV version 2.3