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
|