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