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