include/boost/corosio/native/detail/uring/uring_scheduler.hpp
95.2% Lines (536 / 563, 9 excl)
96.7% Functions (58 / 60)
Functions (60)
Function
Calls
Lines
Blocks
boost::corosio::detail::scoped_sigpipe_block::~scoped_sigpipe_block()
:73
6027x
90.9%
93.0%
boost::corosio::detail::scoped_sigpipe_block::scoped_sigpipe_block()
:93
6027x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::ring()
:161
3044x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::dispatch_mutex() const
:168
20646x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::ring_mutex() const
:174
3028x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::submit_op_posted_exchange(bool) const
:199
3020x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::submit_op_ref() const
:207
413x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::inflight_inc() const
:216
8393x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::inflight() const
:233
4x
100.0%
68.0%
boost::corosio::detail::uring_scheduler::init_ring() const
:256
938x
100.0%
100.0%
void boost::corosio::detail::uring_scheduler::retire_op<boost::corosio::detail::uring_multi_accept_op, boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::retire_multishot()::{lambda(boost::corosio::detail::uring_multi_accept_op&)#1}>(std::unique_ptr<boost::corosio::detail::uring_multi_accept_op, std::default_delete<boost::corosio::detail::uring_multi_accept_op> >&, boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_local_stream_acceptor, boost::corosio::local_stream_acceptor::implementation, boost::corosio::local_endpoint, boost::corosio::detail::uring_local_stream_service>::retire_multishot()::{lambda(boost::corosio::detail::uring_multi_accept_op&)#1})
:350
30x
100.0%
100.0%
void boost::corosio::detail::uring_scheduler::retire_op<boost::corosio::detail::uring_multi_accept_op, boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::retire_multishot()::{lambda(boost::corosio::detail::uring_multi_accept_op&)#1}>(std::unique_ptr<boost::corosio::detail::uring_multi_accept_op, std::default_delete<boost::corosio::detail::uring_multi_accept_op> >&, boost::corosio::detail::uring_multishot_acceptor_base<boost::corosio::detail::uring_tcp_acceptor, boost::corosio::tcp_acceptor::implementation, boost::corosio::endpoint, boost::corosio::detail::uring_tcp_service>::retire_multishot()::{lambda(boost::corosio::detail::uring_multi_accept_op&)#1})
:350
193x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::push_completed_locked(boost::corosio::detail::scheduler_op*) const
:398
20646x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::configure_threading(boost::corosio::detail::scheduler::threading_config)
:403
938x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::configure_sqpoll(bool, unsigned int, int)
:432
0
0.0%
0.0%
boost::corosio::detail::uring_scheduler::scheduler_locking_disabled() const
:440
5x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::signal_drain_op::operator()()
:526
159x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::signal_drain_op::destroy()
:534
0
0.0%
0.0%
boost::corosio::detail::uring_scheduler::submit_sqes_op::submit_sqes_op()
:546
938x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::uring_scheduler(boost::capy::execution_context&, int)
:582
938x
100.0%
68.0%
boost::corosio::detail::uring_scheduler::uring_scheduler(boost::capy::execution_context&, int)::{lambda(void*)#1}::operator()(void*) const
:593
2677x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::~uring_scheduler()
:605
1876x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::lazy_init_ring() const
:619
30675x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::lazy_init_ring() const::{lambda()#1}::operator()() const
:621
938x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::lazy_init_ring_unlocked() const
:625
938x
84.6%
89.0%
boost::corosio::detail::uring_scheduler::shutdown()
:732
938x
100.0%
88.0%
boost::corosio::detail::uring_scheduler::stop()
:770
1466x
100.0%
90.0%
boost::corosio::detail::uring_scheduler::stopped() const
:790
452x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::restart()
:796
643x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::work_started()
:802
33748x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::work_finished()
:808
49666x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::interrupt_reactor() const
:815
10863x
71.4%
67.0%
boost::corosio::detail::uring_scheduler::drain_wakeup_eventfd() const
:843
10702x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::prep_multishot_poll(int, void*)
:854
71x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::register_signal_reader(int)
:876
66x
100.0%
88.0%
boost::corosio::detail::uring_scheduler::post(std::__n4861::coroutine_handle<void>) const
:901
1843x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::post_handler(std::__n4861::coroutine_handle<void>)
:906
1843x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::operator()()
:908
1837x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::destroy()
:915
6x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::post(boost::corosio::detail::scheduler_op*) const
:940
5831x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::post(boost::capy::continuation&) const
:957
8360x
100.0%
100.0%
boost::corosio::detail::uring_run_guard::uring_run_guard(boost::corosio::detail::uring_scheduler const*)
:1003
1868x
100.0%
100.0%
boost::corosio::detail::uring_run_guard::~uring_run_guard()
:1011
1868x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::running_in_this_thread() const
:1018
4468x
100.0%
86.0%
boost::corosio::detail::uring_scheduler::reset_inline_budget() const
:1029
62765x
100.0%
83.0%
boost::corosio::detail::uring_scheduler::try_consume_inline_budget() const
:1042
351157x
87.5%
78.0%
boost::corosio::detail::uring_scheduler::run()
:1060
778x
100.0%
78.0%
boost::corosio::detail::uring_scheduler::run_one()
:1090
71x
100.0%
73.0%
boost::corosio::detail::uring_scheduler::wait_one(long)
:1103
1289x
100.0%
73.0%
boost::corosio::detail::uring_scheduler::poll()
:1116
36x
100.0%
77.0%
boost::corosio::detail::uring_scheduler::poll_one()
:1135
4x
100.0%
73.0%
boost::corosio::detail::uring_scheduler::do_one(long)
:1148
40499x
98.7%
77.0%
boost::corosio::detail::uring_scheduler::process_completions()
:1370
13296x
100.0%
98.0%
boost::corosio::detail::uring_scheduler::submit_sqes_op::do_handler(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int)
:1476
29x
30.0%
38.0%
boost::corosio::detail::uring_scheduler::submit_cancel_by_user_data(boost::corosio::detail::uring_op*)
:1496
218x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::submit_cancel_by_fd(int)
:1521
97x
100.0%
100.0%
boost::corosio::detail::uring_op::on_cancel()
:1541
218x
100.0%
89.0%
boost::corosio::detail::uring_scheduler::cancel_and_flush(int)
:1552
4857x
100.0%
100.0%
boost::corosio::detail::uring_scheduler::release_retired_op(boost::corosio::detail::uring_op*)
:1581
4x
100.0%
94.0%
boost::corosio::detail::uring_scheduler::drain_cqes_for(boost::corosio::detail::uring_op*)
:1604
215x
94.0%
94.0%
| Line | TLA | Hits | Source Code |
|---|---|---|---|
| 1 | // | ||
| 2 | // Copyright (c) 2026 Steve Gerbino | ||
| 3 | // Copyright (c) 2026 Michael Vandeberg | ||
| 4 | // | ||
| 5 | // Distributed under the Boost Software License, Version 1.0. (See accompanying | ||
| 6 | // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt) | ||
| 7 | // | ||
| 8 | // Official repository: https://github.com/cppalliance/corosio | ||
| 9 | // | ||
| 10 | |||
| 11 | #ifndef BOOST_COROSIO_NATIVE_DETAIL_URING_URING_SCHEDULER_HPP | ||
| 12 | #define BOOST_COROSIO_NATIVE_DETAIL_URING_URING_SCHEDULER_HPP | ||
| 13 | |||
| 14 | #include <boost/corosio/detail/platform.hpp> | ||
| 15 | |||
| 16 | #if BOOST_COROSIO_HAS_URING | ||
| 17 | |||
| 18 | // Include before any project headers open a namespace — prevents the | ||
| 19 | // boost::corosio::uring tag variable from shadowing struct ::io_uring. | ||
| 20 | #include <liburing.h> | ||
| 21 | |||
| 22 | #include <boost/corosio/detail/conditionally_enabled_event.hpp> | ||
| 23 | #include <boost/corosio/detail/conditionally_enabled_mutex.hpp> | ||
| 24 | #include <boost/corosio/detail/config.hpp> | ||
| 25 | #include <boost/corosio/detail/except.hpp> | ||
| 26 | #include <boost/corosio/detail/ready_queue.hpp> | ||
| 27 | #include <boost/corosio/detail/scheduler.hpp> | ||
| 28 | #include <boost/corosio/detail/scheduler_op.hpp> | ||
| 29 | #include <boost/corosio/detail/timer_service.hpp> | ||
| 30 | #include <boost/corosio/native/detail/uring/uring_op.hpp> | ||
| 31 | #include <boost/corosio/native/detail/make_err.hpp> | ||
| 32 | #include <boost/corosio/native/detail/posix/posix_signal_service.hpp> | ||
| 33 | #include <boost/capy/ex/execution_context.hpp> | ||
| 34 | |||
| 35 | #include <atomic> | ||
| 36 | #include <chrono> | ||
| 37 | #include <coroutine> | ||
| 38 | #include <cstddef> | ||
| 39 | #include <cstdint> | ||
| 40 | #include <limits> | ||
| 41 | #include <memory> | ||
| 42 | #include <mutex> | ||
| 43 | #include <tuple> | ||
| 44 | #include <vector> | ||
| 45 | |||
| 46 | #include <errno.h> | ||
| 47 | #include <poll.h> | ||
| 48 | #include <pthread.h> | ||
| 49 | #include <signal.h> | ||
| 50 | #include <sys/eventfd.h> | ||
| 51 | #include <time.h> | ||
| 52 | #include <unistd.h> | ||
| 53 | |||
| 54 | namespace boost::corosio::detail { | ||
| 55 | |||
| 56 | /** Block SIGPIPE on the calling thread for the current scope. | ||
| 57 | |||
| 58 | A queued pipe write can execute inline while this thread is in | ||
| 59 | the kernel submitting SQEs; if the pipe's reader is already gone | ||
| 60 | the kernel raises a thread-directed SIGPIPE, which kills any | ||
| 61 | process that has not ignored the signal. The write's CQE still | ||
| 62 | reports `EPIPE`, so the signal carries no information the | ||
| 63 | completion path does not already deliver. The destructor consumes | ||
| 64 | any SIGPIPE raised while blocked and restores the caller's mask. | ||
| 65 | */ | ||
| 66 | class scoped_sigpipe_block | ||
| 67 | { | ||
| 68 | sigset_t old_{}; | ||
| 69 | bool restore_ = false; | ||
| 70 | |||
| 71 | public: | ||
| 72 | /// Consume any SIGPIPE raised in scope and restore the mask. | ||
| 73 | 6027x | ~scoped_sigpipe_block() | |
| 74 | 6027x | { | |
| 75 | 6027x | if (!restore_) | |
| 76 | ✗ | return; | |
| 77 | 6027x | if (!sigismember(&old_, SIGPIPE)) | |
| 78 | { | ||
| 79 | // Only consume what this scope could have generated; a | ||
| 80 | // caller who blocked SIGPIPE keeps their pending state. | ||
| 81 | sigset_t set; | ||
| 82 | 5789x | sigemptyset(&set); | |
| 83 | 5789x | sigaddset(&set, SIGPIPE); | |
| 84 | 5789x | timespec zero{}; | |
| 85 | 5790x | while (::sigtimedwait(&set, nullptr, &zero) == SIGPIPE) | |
| 86 | { | ||
| 87 | } | ||
| 88 | } | ||
| 89 | 6027x | ::pthread_sigmask(SIG_SETMASK, &old_, nullptr); | |
| 90 | 6027x | } | |
| 91 | |||
| 92 | /// Construct and block SIGPIPE for the calling thread. | ||
| 93 | 6027x | scoped_sigpipe_block() noexcept | |
| 94 | 6027x | { | |
| 95 | sigset_t set; | ||
| 96 | 6027x | sigemptyset(&set); | |
| 97 | 6027x | sigaddset(&set, SIGPIPE); | |
| 98 | 6027x | restore_ = ::pthread_sigmask(SIG_BLOCK, &set, &old_) == 0; | |
| 99 | 6027x | } | |
| 100 | |||
| 101 | scoped_sigpipe_block(scoped_sigpipe_block const&) = delete; | ||
| 102 | scoped_sigpipe_block& operator=(scoped_sigpipe_block const&) = delete; | ||
| 103 | }; | ||
| 104 | |||
| 105 | // Forward-declared so the out-of-line inline definitions below the class | ||
| 106 | // can reference the frame stack without a circular dependency. | ||
| 107 | struct uring_scheduler_frame; | ||
| 108 | extern thread_local uring_scheduler_frame* tl_running_scheduler_frame_; | ||
| 109 | |||
| 110 | /** io_uring scheduler — proactor model on Linux 6.x+. | ||
| 111 | |||
| 112 | Owns one io_uring per io_context. Lazy batched submit; | ||
| 113 | cross-thread post wakes a registered eventfd via multishot | ||
| 114 | POLL_ADD. | ||
| 115 | |||
| 116 | @par Thread Safety | ||
| 117 | All public member functions are thread-safe. | ||
| 118 | */ | ||
| 119 | class BOOST_COROSIO_DECL uring_scheduler final | ||
| 120 | : public scheduler | ||
| 121 | { | ||
| 122 | public: | ||
| 123 | using mutex_type = conditionally_enabled_mutex; | ||
| 124 | using lock_type = mutex_type::scoped_lock; | ||
| 125 | using event_type = conditionally_enabled_event; | ||
| 126 | |||
| 127 | uring_scheduler(capy::execution_context& ctx, int concurrency_hint = -1); | ||
| 128 | ~uring_scheduler() override; | ||
| 129 | uring_scheduler(uring_scheduler const&) = delete; | ||
| 130 | uring_scheduler& operator=(uring_scheduler const&) = delete; | ||
| 131 | |||
| 132 | void shutdown() override; | ||
| 133 | |||
| 134 | void post(std::coroutine_handle<>) const override; | ||
| 135 | void post(scheduler_op*) const override; | ||
| 136 | void post(capy::continuation&) const override; | ||
| 137 | |||
| 138 | bool running_in_this_thread() const noexcept override; | ||
| 139 | void stop() override; | ||
| 140 | bool stopped() const noexcept override; | ||
| 141 | void restart() override; | ||
| 142 | std::size_t run() override; | ||
| 143 | std::size_t run_one() override; | ||
| 144 | std::size_t wait_one(long usec) override; | ||
| 145 | std::size_t poll() override; | ||
| 146 | std::size_t poll_one() override; | ||
| 147 | void work_started() noexcept override; | ||
| 148 | void work_finished() noexcept override; | ||
| 149 | |||
| 150 | /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp). | ||
| 151 | /// Submits a multishot POLL on @p read_fd; on readiness the drain+deliver | ||
| 152 | /// runs in dispatch context via signal_drain_op_. | ||
| 153 | [[nodiscard]] std::error_code register_signal_reader(int read_fd) override; | ||
| 154 | |||
| 155 | /** Return the underlying liburing ring. | ||
| 156 | |||
| 157 | Triggers lazy ring initialisation on first call. Used by | ||
| 158 | socket op submission helpers (e.g. `uring_submit_op`) and | ||
| 159 | any other code path that needs a live ring pointer. | ||
| 160 | */ | ||
| 161 | 3044x | struct ::io_uring* ring() noexcept | |
| 162 | { | ||
| 163 | 3044x | lazy_init_ring(); | |
| 164 | 3044x | return &ring_; | |
| 165 | } | ||
| 166 | |||
| 167 | /// Return the dispatch mutex (protects completed_ops_ / cond_). | ||
| 168 | 20646x | mutex_type& dispatch_mutex() const noexcept | |
| 169 | { | ||
| 170 | 20646x | return dispatch_mutex_; | |
| 171 | } | ||
| 172 | |||
| 173 | /// Return the ring mutex (serialises userspace SQ/CQ access). | ||
| 174 | 3028x | mutex_type& ring_mutex() const noexcept | |
| 175 | { | ||
| 176 | 3028x | return ring_mutex_; | |
| 177 | } | ||
| 178 | |||
| 179 | /** Reset the calling thread's inline-budget for this scheduler. | ||
| 180 | |||
| 181 | Called at the top of each dispatched op in `do_one` so each | ||
| 182 | op handler gets a fresh budget for inline speculative | ||
| 183 | completions. Walks the frame stack; no-op if this scheduler | ||
| 184 | isn't on the stack (i.e. called from a non-run thread). | ||
| 185 | */ | ||
| 186 | void reset_inline_budget() const noexcept; | ||
| 187 | |||
| 188 | /** Consume one unit of inline budget if available. | ||
| 189 | |||
| 190 | @return `true` if budget was available and consumed; `false` | ||
| 191 | if the budget is exhausted or this scheduler is not on | ||
| 192 | the calling thread's run stack. | ||
| 193 | */ | ||
| 194 | bool try_consume_inline_budget() const noexcept; | ||
| 195 | |||
| 196 | /// Exchange the submit-batch posted flag. Returns the prior value. | ||
| 197 | /// Caller MUST hold ring_mutex_ — the flag is plain bool, not atomic, | ||
| 198 | /// and the mutex provides the read-modify-write atomicity. | ||
| 199 | 3020x | bool submit_op_posted_exchange(bool desired) const noexcept | |
| 200 | { | ||
| 201 | 3020x | bool prev = submit_op_posted_; | |
| 202 | 3020x | submit_op_posted_ = desired; | |
| 203 | 3020x | return prev; | |
| 204 | } | ||
| 205 | |||
| 206 | /// Return a reference to the mutable embedded submit_sqes_op. | ||
| 207 | 413x | scheduler_op& submit_op_ref() const noexcept | |
| 208 | { | ||
| 209 | 413x | return submit_op_; | |
| 210 | } | ||
| 211 | |||
| 212 | /// Increment the io_uring in-flight counter. Callers prep an SQE | ||
| 213 | /// whose CQE will require IORING_ENTER_GETEVENTS to surface under | ||
| 214 | /// DEFER_TASKRUN. Excluded: the wakeup-eventfd multishot SQE, whose | ||
| 215 | /// progress doesn't depend on userspace getevents. | ||
| 216 | 8393x | void inflight_inc() const noexcept | |
| 217 | { | ||
| 218 | 8393x | uring_inflight_.fetch_add(1, std::memory_order_release); | |
| 219 | 8393x | } | |
| 220 | |||
| 221 | /** Return the current io_uring in-flight counter. | ||
| 222 | |||
| 223 | Test-only helper: `uring_inflight_` is an internal accounting | ||
| 224 | counter (it gates the `do_one` ring pump), with no bearing on the | ||
| 225 | public API. It is exposed solely so tests can assert the counter | ||
| 226 | stays balanced across op submission and teardown — in particular | ||
| 227 | that `drain_cqes_for` does not leak counts. Not thread-safe with | ||
| 228 | respect to a concurrently running scheduler; call it from a quiesced | ||
| 229 | context. | ||
| 230 | |||
| 231 | @return The number of SQEs currently counted as in flight. | ||
| 232 | */ | ||
| 233 | 4x | std::int64_t inflight() const noexcept | |
| 234 | { | ||
| 235 | 8x | return uring_inflight_.load(std::memory_order_acquire); | |
| 236 | } | ||
| 237 | |||
| 238 | /// Initialize the io_uring ring on first access. Idempotent. | ||
| 239 | void lazy_init_ring() const; | ||
| 240 | |||
| 241 | /** Create the io_uring ring. | ||
| 242 | |||
| 243 | Called once the owning `io_context` has applied every option | ||
| 244 | that feeds the ring's setup flags, so that a kernel that | ||
| 245 | refuses the ring is reported from the `io_context` | ||
| 246 | constructor rather than from the first operation submitted | ||
| 247 | through it. Idempotent. | ||
| 248 | |||
| 249 | @par Exception Safety | ||
| 250 | Strong guarantee. A ring that cannot be created leaves no | ||
| 251 | descriptor behind. | ||
| 252 | |||
| 253 | @throws std::system_error If the ring, its wakeup eventfd, or | ||
| 254 | the poll watching that eventfd could not be set up. | ||
| 255 | */ | ||
| 256 | 938x | void init_ring() const | |
| 257 | { | ||
| 258 | 938x | lazy_init_ring(); | |
| 259 | 932x | } | |
| 260 | |||
| 261 | /// Wake the leader if it's blocked in `submit_and_wait_timeout`. | ||
| 262 | /// Best-effort: the wakeup is suppressed if the leader has already | ||
| 263 | /// been signalled and not yet acked. | ||
| 264 | void interrupt_reactor() const noexcept; | ||
| 265 | |||
| 266 | /** Submit `IORING_OP_ASYNC_CANCEL` targeting an in-flight op by its | ||
| 267 | user_data pointer. | ||
| 268 | |||
| 269 | The kernel delivers `-ECANCELED` on the target's CQE if it was | ||
| 270 | still in flight; the op's completion handler then reports | ||
| 271 | `operation_aborted`. Best-effort: if the SQ is full after one | ||
| 272 | flush attempt the function returns without cancelling (the op | ||
| 273 | will complete normally on its own). | ||
| 274 | |||
| 275 | @param target The in-flight op to cancel. | ||
| 276 | */ | ||
| 277 | void submit_cancel_by_user_data(uring_op* target) noexcept; | ||
| 278 | |||
| 279 | /** Submit `IORING_OP_ASYNC_CANCEL` with `IORING_ASYNC_CANCEL_FD` | ||
| 280 | to cancel every in-flight op on the given fd in one SQE. | ||
| 281 | |||
| 282 | Best-effort: if the SQ is full after one flush attempt the | ||
| 283 | function returns without cancelling. | ||
| 284 | |||
| 285 | @param fd The file descriptor whose in-flight ops should be | ||
| 286 | cancelled. | ||
| 287 | */ | ||
| 288 | void submit_cancel_by_fd(int fd) noexcept; | ||
| 289 | |||
| 290 | /** Submit `IORING_OP_ASYNC_CANCEL` for `fd` and immediately flush | ||
| 291 | the submission ring to the kernel. | ||
| 292 | |||
| 293 | Must be called while `fd` is still open so the kernel can | ||
| 294 | resolve the file from the fd number before it is closed and | ||
| 295 | potentially recycled. | ||
| 296 | |||
| 297 | Best-effort: if the SQ is full the function still flushes any | ||
| 298 | earlier pending SQEs to the kernel. | ||
| 299 | |||
| 300 | @param fd The file descriptor whose in-flight ops should be | ||
| 301 | cancelled. | ||
| 302 | */ | ||
| 303 | void cancel_and_flush(int fd) noexcept; | ||
| 304 | |||
| 305 | /** Take ownership of an op the kernel may still complete. | ||
| 306 | |||
| 307 | For member-owned ops (e.g. a multishot accept arming) whose | ||
| 308 | owner goes away on a *non-terminal* path — the owning acceptor | ||
| 309 | adopting a different descriptor, say. Draining the op's CQEs | ||
| 310 | there is not an option: `drain_cqes_for` consumes without | ||
| 311 | dispatching, which is only tolerable while the whole context | ||
| 312 | is being torn down. | ||
| 313 | |||
| 314 | The op stays allocated, and its `user_data` therefore stays | ||
| 315 | reserved, until the kernel delivers its terminal CQE. That is | ||
| 316 | what keeps a freshly submitted op from aliasing the retired | ||
| 317 | one. `process_completions` routes a retired op's CQEs to its | ||
| 318 | `retire_func` and frees the op on the terminal CQE; anything | ||
| 319 | still outstanding is released when the scheduler is destroyed, | ||
| 320 | after `io_uring_queue_exit`. | ||
| 321 | |||
| 322 | @par Invariant | ||
| 323 | Retirement state — the `retired` flag, and every owner | ||
| 324 | back-pointer @p prepare clears — is written and read only | ||
| 325 | under `ring_mutex_`. `process_completions` runs the whole CQE | ||
| 326 | dispatch under that mutex, so taking it here is the only thing | ||
| 327 | that orders this publication against a leader mid-dispatch. | ||
| 328 | Everything is therefore done inside one critical section: | ||
| 329 | @p prepare runs, the flag is set, and the op moves into the | ||
| 330 | retired list before the mutex is released. | ||
| 331 | |||
| 332 | @par Thread Safety | ||
| 333 | Safe to call from any thread. Takes `ring_mutex_` then | ||
| 334 | `retired_mutex_`, the same order `release_retired_op` uses. | ||
| 335 | Callers must hold neither. | ||
| 336 | |||
| 337 | @pre After @p prepare runs, @p slot must hold no reference back | ||
| 338 | to its owner (`impl_ptr` cleared, back-pointers nulled). | ||
| 339 | `release_retired_op` deletes the op from inside the CQE | ||
| 340 | loop with `ring_mutex_` held, so its destructor must not | ||
| 341 | re-enter the scheduler or touch the owner. | ||
| 342 | |||
| 343 | @param slot The owner's pointer to the op. Emptied on return; | ||
| 344 | ignored if already empty. | ||
| 345 | @param prepare Invoked as `prepare(*slot)` under `ring_mutex_` | ||
| 346 | to clear back-pointers and install a `retire_func`. Must | ||
| 347 | not take `ring_mutex_` or post work. | ||
| 348 | */ | ||
| 349 | template<class Op, class PrepareFn> | ||
| 350 | 223x | void retire_op(std::unique_ptr<Op>& slot, PrepareFn prepare) noexcept | |
| 351 | { | ||
| 352 | // Interrupt first so a leader parked in the kernel drops | ||
| 353 | // ring_mutex_ promptly, as cancel_and_flush does. | ||
| 354 | 223x | interrupt_reactor(); | |
| 355 | 223x | lock_type lock(ring_mutex_); | |
| 356 | 223x | if (!slot) | |
| 357 | 218x | return; | |
| 358 | 5x | prepare(*slot); | |
| 359 | 5x | slot->retired = true; | |
| 360 | 5x | std::lock_guard<std::mutex> retired_lock(retired_mutex_); | |
| 361 | 5x | retired_ops_.push_back(std::move(slot)); | |
| 362 | 223x | } | |
| 363 | |||
| 364 | /** Drain pending CQEs for a specific op's `user_data`. | ||
| 365 | |||
| 366 | Submits an ASYNC_CANCEL by user_data to short-circuit any | ||
| 367 | in-flight op holding `target`, then iterates the CQ ring and | ||
| 368 | consumes every CQE matching `target` so its memory can be | ||
| 369 | freed safely. Used by member-owned ops (e.g. | ||
| 370 | `uring_multi_accept_op`) whose destructor cannot tolerate | ||
| 371 | outstanding CQEs. | ||
| 372 | |||
| 373 | @warning Teardown paths only. Every CQE this walks is consumed | ||
| 374 | and none are dispatched, so an unrelated op completing inside | ||
| 375 | the drain window loses its handler and its coroutine never | ||
| 376 | resumes. That is only tolerable when the owning object is | ||
| 377 | being destroyed. To retire a still-live op on a non-terminal | ||
| 378 | path, use @ref retire_op instead. | ||
| 379 | |||
| 380 | @par Thread Safety | ||
| 381 | Safe to call from any thread. Internally takes `ring_mutex_` | ||
| 382 | to serialise against the run-loop leader; calls | ||
| 383 | `interrupt_reactor()` first so the leader returns from its | ||
| 384 | kernel wait promptly. | ||
| 385 | |||
| 386 | @param target The op pointer used as user_data on the SQE. | ||
| 387 | */ | ||
| 388 | void drain_cqes_for(uring_op* target) noexcept; | ||
| 389 | |||
| 390 | /** Queue an already-counted op while the caller holds dispatch_mutex_. | ||
| 391 | |||
| 392 | Does NOT increment `outstanding_work_`. Use for synchronous | ||
| 393 | completion paths (e.g. SQE backpressure) where the caller called | ||
| 394 | `work_started()` and already holds the dispatch lock. | ||
| 395 | |||
| 396 | @pre `dispatch_mutex_` must be locked by the calling thread. | ||
| 397 | */ | ||
| 398 | 20646x | void push_completed_locked(scheduler_op* op) const noexcept | |
| 399 | { | ||
| 400 | 20646x | completed_ops_.push(op); | |
| 401 | 20646x | } | |
| 402 | |||
| 403 | 938x | void configure_threading(threading_config cfg) noexcept override | |
| 404 | { | ||
| 405 | 938x | scheduler_locking_disabled_ = !cfg.scheduler_locking; | |
| 406 | // reactor_io_locking off also drives SINGLE_ISSUER/DEFER_TASKRUN and | ||
| 407 | // the eventfd wake elision (see lazy_init_ring_unlocked and | ||
| 408 | // interrupt_reactor). one_thread is unused: the leader-follower wake | ||
| 409 | // model is not gated on it. | ||
| 410 | 938x | reactor_io_locking_ = cfg.reactor_io_locking; | |
| 411 | 938x | dispatch_mutex_.set_enabled(cfg.scheduler_locking); | |
| 412 | 938x | ring_mutex_.set_enabled(cfg.reactor_io_locking); | |
| 413 | 938x | cond_.set_enabled(cfg.scheduler_locking); | |
| 414 | 938x | } | |
| 415 | |||
| 416 | /** Configure SQPOLL parameters. | ||
| 417 | |||
| 418 | Must be called before the ring is created — the values are | ||
| 419 | cached and read when it is. No-op if `enable` is false (the | ||
| 420 | default). | ||
| 421 | |||
| 422 | @note When combined with single-threaded mode, | ||
| 423 | IORING_SETUP_DEFER_TASKRUN is suppressed — the kernel | ||
| 424 | rejects that combination. SINGLE_ISSUER still applies. | ||
| 425 | |||
| 426 | @param enable Set IORING_SETUP_SQPOLL on ring init. | ||
| 427 | @param idle_ms sq_thread_idle in milliseconds; 0 = kernel | ||
| 428 | default (1ms). | ||
| 429 | @param cpu Pin the polling thread to this CPU; -1 to | ||
| 430 | not pin. | ||
| 431 | */ | ||
| 432 | ✗ | void configure_sqpoll(bool enable, unsigned idle_ms, int cpu) noexcept | |
| 433 | { | ||
| 434 | ✗ | enable_sqpoll_ = enable; | |
| 435 | ✗ | sq_thread_idle_ms_ = idle_ms; | |
| 436 | ✗ | sq_thread_cpu_ = cpu; | |
| 437 | ✗ | } | |
| 438 | |||
| 439 | /// Return true when scheduler locking is disabled (fully-lockless tier). | ||
| 440 | 5x | bool scheduler_locking_disabled() const noexcept override | |
| 441 | { | ||
| 442 | 5x | return scheduler_locking_disabled_; | |
| 443 | } | ||
| 444 | |||
| 445 | private: | ||
| 446 | // ring_ + wakeup_eventfd_ are mutable so lazy_init_ring() (called | ||
| 447 | // from const contexts like post()) can populate them on first use. | ||
| 448 | mutable struct ::io_uring ring_{}; | ||
| 449 | mutable int wakeup_eventfd_ = -1; | ||
| 450 | timer_service* timer_svc_ = nullptr; | ||
| 451 | |||
| 452 | // dispatch_mutex_ protects completed_ops_, cond_, task_running_. | ||
| 453 | // ring_mutex_ protects every userspace touch of ring_ (SQ tail, | ||
| 454 | // CQ head): get_sqe / submit / submit_and_wait_timeout / | ||
| 455 | // for_each_cqe / cq_advance. | ||
| 456 | // | ||
| 457 | // process_completions runs under ring_mutex_ and briefly takes | ||
| 458 | // dispatch_mutex_ to splice into completed_ops_. The locks are | ||
| 459 | // never held simultaneously for the full duration of any other | ||
| 460 | // path's critical section, so no deadlock. | ||
| 461 | mutable mutex_type dispatch_mutex_{true}; | ||
| 462 | mutable mutex_type ring_mutex_{true}; | ||
| 463 | mutable event_type cond_{true}; | ||
| 464 | mutable ready_queue completed_ops_; | ||
| 465 | // outstanding_work_ and uring_inflight_ are both atomic | ||
| 466 | // counters updated at high frequency on different paths: | ||
| 467 | // - outstanding_work_ : every work_started / work_finished call, | ||
| 468 | // including timers, posts, and SQE submits. | ||
| 469 | // - uring_inflight_ : only SQE submit + non-F_MORE CQE consume. | ||
| 470 | // Under multi-thread workloads the threads tend to update these | ||
| 471 | // from different code paths; placing them on the same cache line | ||
| 472 | // would cause false sharing and unnecessary cache-line ping-pong. | ||
| 473 | // Hold each on its own line. | ||
| 474 | alignas(64) mutable std::atomic<std::int64_t> outstanding_work_{0}; | ||
| 475 | // Count of io_uring SQEs in flight whose completion requires user- | ||
| 476 | // space to enter the kernel via IORING_ENTER_GETEVENTS for task | ||
| 477 | // work to progress under IORING_SETUP_DEFER_TASKRUN. Excludes the | ||
| 478 | // wakeup-eventfd multishot poll (registered in lazy_init_ring), and | ||
| 479 | // is updated by uring_submit_op and by process_completions on | ||
| 480 | // each non-F_MORE, non-eventfd CQE. Used by do_one to skip the | ||
| 481 | // ring pump when there is no io_uring work pending. | ||
| 482 | alignas(64) mutable std::atomic<std::int64_t> uring_inflight_{0}; | ||
| 483 | std::atomic<bool> stopped_{false}; | ||
| 484 | // Leader-follower flag: true while a thread is blocked in | ||
| 485 | // io_uring_submit_and_wait_timeout. Protected by dispatch_mutex_. | ||
| 486 | mutable bool task_running_ = false; | ||
| 487 | bool scheduler_locking_disabled_ = false; | ||
| 488 | bool reactor_io_locking_ = true; | ||
| 489 | bool enable_sqpoll_ = false; | ||
| 490 | unsigned sq_thread_idle_ms_ = 0; | ||
| 491 | int sq_thread_cpu_ = -1; | ||
| 492 | |||
| 493 | int cancel_sentinel_ = 0; | ||
| 494 | |||
| 495 | // Ops adopted by retire_op, kept alive so the kernel never sees | ||
| 496 | // their user_data reused. Declared before ring_ is exited only in | ||
| 497 | // the sense that the destructor body runs io_uring_queue_exit | ||
| 498 | // first; the vector is then destroyed with the rest of the members, | ||
| 499 | // by which point the kernel can no longer reference these ops. | ||
| 500 | // A plain std::mutex (not the conditionally-enabled mutex_type): | ||
| 501 | // retirement happens on user threads regardless of the threading | ||
| 502 | // configuration. Leaf lock — nothing else is taken under it. | ||
| 503 | mutable std::mutex retired_mutex_; | ||
| 504 | std::vector<std::unique_ptr<uring_op>> retired_ops_; | ||
| 505 | |||
| 506 | /// Free a retired op once its terminal CQE has been consumed. | ||
| 507 | void release_retired_op(uring_op* op) noexcept; | ||
| 508 | |||
| 509 | // Signal self-pipe integration. The read end is watched via a multishot | ||
| 510 | // POLL SQE tagged with &signal_pipe_sentinel_ (distinct from nullptr = | ||
| 511 | // wakeup eventfd and &cancel_sentinel_). On its CQE we re-arm the poll if | ||
| 512 | // needed and enqueue signal_drain_op_ so the drain+deliver runs in | ||
| 513 | // dispatch context — never under ring_mutex_ — keeping deliver_signal's | ||
| 514 | // mutex locking off the ring critical section. | ||
| 515 | int signal_pipe_read_fd_ = -1; | ||
| 516 | int signal_pipe_sentinel_ = 0; | ||
| 517 | |||
| 518 | /// Dispatch-context op that drains the signal self-pipe and delivers each | ||
| 519 | /// pending signal. Enqueued (once at a time, guarded by queued_) from | ||
| 520 | /// process_completions when the poll CQE fires. Scheduler-owned; destroy() | ||
| 521 | /// is a no-op. | ||
| 522 | struct signal_drain_op final : scheduler_op | ||
| 523 | { | ||
| 524 | std::atomic<bool> queued_{false}; | ||
| 525 | |||
| 526 | 159x | void operator()() override | |
| 527 | { | ||
| 528 | // Clear before draining so a signal that arrives mid-drain re-arms | ||
| 529 | // the op via a fresh CQE rather than being lost. | ||
| 530 | 159x | queued_.store(false, std::memory_order_release); | |
| 531 | 159x | posix_signal_detail::drain_signal_pipe(); | |
| 532 | 159x | } | |
| 533 | |||
| 534 | ✗ | void destroy() override {} | |
| 535 | }; | ||
| 536 | mutable signal_drain_op signal_drain_op_; | ||
| 537 | |||
| 538 | /// Flushes the SQ ring and drains CQEs in one mutex-held pass. | ||
| 539 | /// One instance covers a whole batch; subsequent SQEs in the same | ||
| 540 | /// batch skip the post, amortising syscall cost across the batch. | ||
| 541 | /// Mirrors Asio's `submit_sqes_op` (`io_uring_service.ipp:730-742`). | ||
| 542 | struct submit_sqes_op final : scheduler_op | ||
| 543 | { | ||
| 544 | uring_scheduler* sched_ = nullptr; | ||
| 545 | |||
| 546 | 938x | submit_sqes_op() noexcept : scheduler_op(&do_handler) {} | |
| 547 | |||
| 548 | static void do_handler( | ||
| 549 | void* owner, | ||
| 550 | scheduler_op* base, | ||
| 551 | std::uint32_t /*bytes*/, | ||
| 552 | std::uint32_t /*error*/) noexcept; | ||
| 553 | }; | ||
| 554 | |||
| 555 | /// True between the first submitter of a batch posting `submit_op_` | ||
| 556 | /// and the dispatched op clearing the flag inside its handler. Read | ||
| 557 | /// and written only while holding `ring_mutex_`. | ||
| 558 | mutable bool submit_op_posted_ = false; | ||
| 559 | |||
| 560 | /// Single embedded `submit_sqes_op` instance, owned by the scheduler. | ||
| 561 | mutable submit_sqes_op submit_op_; | ||
| 562 | |||
| 563 | // drain_cqes_for tuning. The bound exists to avoid stalling a | ||
| 564 | // destructor if the kernel never returns a cancel completion (best- | ||
| 565 | // effort drain); 8 rounds * 1ms == 8ms worst case. | ||
| 566 | static constexpr int drain_cqes_max_rounds = 8; | ||
| 567 | static constexpr unsigned long drain_cqes_kick_ns = 1'000'000; | ||
| 568 | |||
| 569 | // ring_inited_ goes true once the ring exists. The init is deferred | ||
| 570 | // from the constructor so configure_threading() and configure_sqpoll() | ||
| 571 | // can take effect before io_uring_queue_init_params chooses flags. | ||
| 572 | mutable std::once_flag ring_init_once_; | ||
| 573 | mutable bool ring_inited_ = false; | ||
| 574 | |||
| 575 | std::size_t do_one(long timeout_us); | ||
| 576 | void process_completions(); | ||
| 577 | void drain_wakeup_eventfd() const noexcept; | ||
| 578 | bool prep_multishot_poll(int fd, void* data) noexcept; | ||
| 579 | void lazy_init_ring_unlocked() const; | ||
| 580 | }; | ||
| 581 | |||
| 582 | 938x | inline uring_scheduler::uring_scheduler( | |
| 583 | 938x | capy::execution_context& ctx, int /*concurrency_hint*/) | |
| 584 | { | ||
| 585 | // sched_ cannot be set in the member initialiser — `this` is not | ||
| 586 | // available there. | ||
| 587 | 938x | submit_op_.sched_ = this; | |
| 588 | |||
| 589 | // Wire timer service. on_earliest_changed wakes the run loop so it | ||
| 590 | // recomputes its wait timeout. | ||
| 591 | 938x | timer_svc_ = &get_timer_service(ctx, *this); | |
| 592 | 938x | timer_svc_->set_on_earliest_changed( | |
| 593 | 3615x | timer_service::callback(this, [](void* p) { | |
| 594 | 2677x | static_cast<uring_scheduler*>(p)->interrupt_reactor(); | |
| 595 | 2677x | })); | |
| 596 | |||
| 597 | // Ring init is deferred so the options the io_context applies | ||
| 598 | // after this constructor — the locking tier and SQPOLL — can feed | ||
| 599 | // the flags io_uring_queue_init_params is given. The io_context | ||
| 600 | // calls init_ring() once they are in place; the lazy_init_ring() | ||
| 601 | // calls scattered through the op paths only matter for a | ||
| 602 | // scheduler used without one. | ||
| 603 | 938x | } | |
| 604 | |||
| 605 | 1876x | inline uring_scheduler::~uring_scheduler() | |
| 606 | { | ||
| 607 | 938x | if (ring_inited_) | |
| 608 | { | ||
| 609 | // Ring teardown can still run a doomed pipe write as task | ||
| 610 | // work; absorb the SIGPIPE it would raise. | ||
| 611 | 932x | scoped_sigpipe_block no_sigpipe; | |
| 612 | 932x | if (wakeup_eventfd_ >= 0) | |
| 613 | 932x | ::close(wakeup_eventfd_); | |
| 614 | 932x | ::io_uring_queue_exit(&ring_); | |
| 615 | 932x | } | |
| 616 | 1876x | } | |
| 617 | |||
| 618 | inline void | ||
| 619 | 30675x | uring_scheduler::lazy_init_ring() const | |
| 620 | { | ||
| 621 | 31613x | std::call_once(ring_init_once_, [this] { lazy_init_ring_unlocked(); }); | |
| 622 | 30669x | } | |
| 623 | |||
| 624 | inline void | ||
| 625 | 938x | uring_scheduler::lazy_init_ring_unlocked() const | |
| 626 | { | ||
| 627 | 938x | io_uring_params params{}; | |
| 628 | // The unsafe_io and unsafe tiers guarantee a single ring submitter. | ||
| 629 | 938x | if (!reactor_io_locking_) | |
| 630 | { | ||
| 631 | // SINGLE_ISSUER promises the kernel one submitter thread, | ||
| 632 | // letting it skip internal SQ locking. DEFER_TASKRUN tells | ||
| 633 | // it to batch task_work delivery at io_uring_enter(GETEVENTS) | ||
| 634 | // boundaries instead of interrupting the run thread via | ||
| 635 | // TWA_SIGNAL — eliminates cache pollution from mid-flight | ||
| 636 | // task_work and gives a meaningful single-threaded | ||
| 637 | // throughput uplift. | ||
| 638 | // | ||
| 639 | // Plan 3 disabled DEFER_TASKRUN defensively over a misread | ||
| 640 | // of the GETEVENTS contract. Plan 4a re-enabled it: liburing's | ||
| 641 | // io_uring_submit_and_wait_timeout always sets | ||
| 642 | // IORING_ENTER_GETEVENTS when wait_nr > 0, regardless of | ||
| 643 | // ts. Our run loop's only kernel-wait call passes wait_nr=1. | ||
| 644 | // Submit-only paths (cancel_and_flush, etc.) leave their | ||
| 645 | // CQEs queued until the leader's next GETEVENTS-bearing | ||
| 646 | // wait — benign. | ||
| 647 | // | ||
| 648 | // Multi-thread mode never sets these flags: SINGLE_ISSUER | ||
| 649 | // would be unsafe with multiple submitter threads. | ||
| 650 | // | ||
| 651 | // DEFER_TASKRUN is suppressed when SQPOLL is also enabled | ||
| 652 | // — the kernel rejects that combination with -EINVAL. The | ||
| 653 | // SQPOLL polling thread already delivers completions | ||
| 654 | // without TWA_SIGNAL interruption, so DEFER_TASKRUN's | ||
| 655 | // benefit is moot in that mode. | ||
| 656 | 2x | params.flags = IORING_SETUP_SINGLE_ISSUER; | |
| 657 | 2x | if (!enable_sqpoll_) | |
| 658 | 2x | params.flags |= IORING_SETUP_DEFER_TASKRUN; | |
| 659 | } | ||
| 660 | |||
| 661 | 938x | if (enable_sqpoll_) | |
| 662 | { | ||
| 663 | // SQPOLL forks a kernel thread that busy-polls the SQ ring; | ||
| 664 | // submission becomes a userspace-only memory store. Combines | ||
| 665 | // with SINGLE_ISSUER (the kernel accepts that pair) but NOT | ||
| 666 | // with DEFER_TASKRUN (kernel returns -EINVAL); the | ||
| 667 | // reactor-I/O-lockless branch above suppresses DEFER_TASKRUN | ||
| 668 | // when SQPOLL is also set. Idle timeout 0 means kernel | ||
| 669 | // default (1ms); we only forward when explicitly set so | ||
| 670 | // the kernel default is preserved. | ||
| 671 | ✗ | params.flags |= IORING_SETUP_SQPOLL; | |
| 672 | ✗ | if (sq_thread_idle_ms_ != 0) | |
| 673 | ✗ | params.sq_thread_idle = sq_thread_idle_ms_; | |
| 674 | ✗ | if (sq_thread_cpu_ >= 0) | |
| 675 | { | ||
| 676 | ✗ | params.flags |= IORING_SETUP_SQ_AFF; | |
| 677 | ✗ | params.sq_thread_cpu = static_cast<__u32>(sq_thread_cpu_); | |
| 678 | } | ||
| 679 | } | ||
| 680 | |||
| 681 | 938x | int rc = ::io_uring_queue_init_params(256, &ring_, ¶ms); | |
| 682 | 938x | if (rc < 0) | |
| 683 | 1x | detail::throw_system_error(make_err(-rc), "io_uring_queue_init_params"); | |
| 684 | |||
| 685 | 937x | wakeup_eventfd_ = ::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC); | |
| 686 | 937x | if (wakeup_eventfd_ < 0) | |
| 687 | { | ||
| 688 | 1x | int errn = errno; | |
| 689 | 1x | ::io_uring_queue_exit(&ring_); | |
| 690 | 1x | detail::throw_system_error(make_err(errn), "eventfd"); | |
| 691 | } | ||
| 692 | |||
| 693 | // Register a one-shot poll on the wake eventfd. user_data nullptr | ||
| 694 | // is the sentinel recognized by process_completions, which calls | ||
| 695 | // drain_wakeup_eventfd() to consume the eventfd byte AND re-arm | ||
| 696 | // the poll. Plan 5a switched away from IORING_POLL_MULTISHOT | ||
| 697 | // because multishot ops can silently terminate (e.g. under CQ | ||
| 698 | // pressure), and we don't observe the termination — leaving the | ||
| 699 | // wake mechanism dead and the leader stuck in kernel wait. One- | ||
| 700 | // shot rearm-on-fire is fail-fast: every wake event is paired | ||
| 701 | // with an explicit rearm, so a missed rearm would manifest | ||
| 702 | // immediately as the next wake being lost (test-visible). | ||
| 703 | 936x | ::io_uring_sqe* sqe = ::io_uring_get_sqe(&ring_); | |
| 704 | 936x | if (!sqe) | |
| 705 | { | ||
| 706 | 3x | ::close(wakeup_eventfd_); | |
| 707 | 3x | wakeup_eventfd_ = -1; | |
| 708 | 3x | ::io_uring_queue_exit(&ring_); | |
| 709 | 3x | detail::throw_system_error( | |
| 710 | 6x | make_err(ENOSPC), "io_uring_get_sqe (wakeup)"); | |
| 711 | } | ||
| 712 | // Multishot poll: fires a CQE on each eventfd POLLIN without | ||
| 713 | // consuming the SQE. Avoids the re-arm hazard of one-shot poll | ||
| 714 | // (where drain_wakeup_eventfd's get_sqe could return null on a | ||
| 715 | // full SQ, leaving no SQE to detect future wakes). | ||
| 716 | 933x | ::io_uring_prep_poll_multishot(sqe, wakeup_eventfd_, POLLIN); | |
| 717 | 933x | ::io_uring_sqe_set_data(sqe, nullptr); | |
| 718 | 933x | int submit_rc = ::io_uring_submit(&ring_); | |
| 719 | 933x | if (submit_rc < 0) | |
| 720 | { | ||
| 721 | 1x | ::close(wakeup_eventfd_); | |
| 722 | 1x | wakeup_eventfd_ = -1; | |
| 723 | 1x | ::io_uring_queue_exit(&ring_); | |
| 724 | 1x | detail::throw_system_error( | |
| 725 | 2x | make_err(-submit_rc), "io_uring_submit (wakeup)"); | |
| 726 | } | ||
| 727 | |||
| 728 | 932x | ring_inited_ = true; | |
| 729 | 932x | } | |
| 730 | |||
| 731 | inline void | ||
| 732 | 938x | uring_scheduler::shutdown() | |
| 733 | { | ||
| 734 | 938x | stopped_.store(true, std::memory_order_release); | |
| 735 | |||
| 736 | // Drain posted ops, calling destroy() on each so embedded handles | ||
| 737 | // (coroutine frames, error_code outputs) get torn down rather | ||
| 738 | // than leaked. Mirrors reactor_scheduler::shutdown_drain. | ||
| 739 | // | ||
| 740 | // Service shutdown order (driven by capy::execution_context): | ||
| 741 | // each socket/acceptor service::shutdown() submits a cancel SQE | ||
| 742 | // for every live impl. The CQEs that result either land in | ||
| 743 | // completed_ops_ (drained here as op->destroy()) or stay in the | ||
| 744 | // kernel ring; ~scheduler's io_uring_queue_exit cleans the | ||
| 745 | // latter up at process teardown. Self-referential impl_ptr | ||
| 746 | // cycles (e.g. multishot acceptor's multi_op_->impl_ptr) are | ||
| 747 | // broken explicitly inside each service before the scheduler | ||
| 748 | // shutdown runs. | ||
| 749 | 938x | lock_type lock(dispatch_mutex_); | |
| 750 | 1003x | while (auto e = completed_ops_.pop()) | |
| 751 | { | ||
| 752 | 65x | if (ready_is_continuation(e)) | |
| 753 | { | ||
| 754 | 4x | lock.unlock(); | |
| 755 | 4x | if (auto h = ready_as_cont(e)->h) | |
| 756 | 4x | h.destroy(); | |
| 757 | 4x | lock.lock(); | |
| 758 | } | ||
| 759 | else | ||
| 760 | { | ||
| 761 | 61x | lock.unlock(); | |
| 762 | 61x | ready_as_op(e)->destroy(); | |
| 763 | 61x | lock.lock(); | |
| 764 | } | ||
| 765 | 65x | } | |
| 766 | 938x | cond_.notify_all(); | |
| 767 | 938x | } | |
| 768 | |||
| 769 | inline void | ||
| 770 | 1466x | uring_scheduler::stop() | |
| 771 | { | ||
| 772 | 1466x | stopped_.store(true, std::memory_order_release); | |
| 773 | { | ||
| 774 | 1466x | lock_type lock(dispatch_mutex_); | |
| 775 | 1466x | cond_.notify_all(); | |
| 776 | 1466x | } | |
| 777 | // Force-wake unconditionally — bypass interrupt_reactor's CAS | ||
| 778 | // coalescing. A dropped wake here leaves the leader blocked | ||
| 779 | // forever in submit_and_wait_timeout (no further CQE will | ||
| 780 | // arrive after stop()). With multishot poll on wakeup_eventfd_, | ||
| 781 | // this write reliably produces a CQE. | ||
| 782 | 1466x | if (ring_inited_) | |
| 783 | { | ||
| 784 | 1466x | std::uint64_t v = 1; | |
| 785 | 1466x | [[maybe_unused]] auto r = ::write(wakeup_eventfd_, &v, sizeof(v)); | |
| 786 | } | ||
| 787 | 1466x | } | |
| 788 | |||
| 789 | inline bool | ||
| 790 | 452x | uring_scheduler::stopped() const noexcept | |
| 791 | { | ||
| 792 | 452x | return stopped_.load(std::memory_order_acquire); | |
| 793 | } | ||
| 794 | |||
| 795 | inline void | ||
| 796 | 643x | uring_scheduler::restart() | |
| 797 | { | ||
| 798 | 643x | stopped_.store(false, std::memory_order_release); | |
| 799 | 643x | } | |
| 800 | |||
| 801 | inline void | ||
| 802 | 33748x | uring_scheduler::work_started() noexcept | |
| 803 | { | ||
| 804 | 33748x | outstanding_work_.fetch_add(1, std::memory_order_relaxed); | |
| 805 | 33748x | } | |
| 806 | |||
| 807 | inline void | ||
| 808 | 49666x | uring_scheduler::work_finished() noexcept | |
| 809 | { | ||
| 810 | 99332x | if (outstanding_work_.fetch_sub(1, std::memory_order_acq_rel) == 1) | |
| 811 | 1153x | stop(); | |
| 812 | 49666x | } | |
| 813 | |||
| 814 | inline void | ||
| 815 | 10863x | uring_scheduler::interrupt_reactor() const noexcept | |
| 816 | { | ||
| 817 | // Skip if the ring hasn't been initialised yet — there's no leader | ||
| 818 | // to wake and no eventfd to write. | ||
| 819 | 10863x | if (!ring_inited_) | |
| 820 | ✗ | return; | |
| 821 | |||
| 822 | // Lockless tiers (reactor-I/O locking off): cross-thread post() is | ||
| 823 | // forbidden, so interrupt_reactor is only ever reached from the leader | ||
| 824 | // thread's own coroutines — it is not in kernel wait, nothing to wake. | ||
| 825 | // Under the safe tier (including concurrency_hint == 1) cross-thread | ||
| 826 | // post() is allowed, so the eventfd write below must always fire. | ||
| 827 | 10863x | if (!reactor_io_locking_) | |
| 828 | ✗ | return; | |
| 829 | |||
| 830 | // Multi-thread: write the eventfd unconditionally. CAS-coalescing | ||
| 831 | // is unsafe here because the leader's Phase 2 in do_one waits | ||
| 832 | // indefinitely for a CQE; a dropped wake leaves the leader | ||
| 833 | // blocked forever when there is no other CQE-producing activity. | ||
| 834 | // Multishot poll on wakeup_eventfd_ delivers a CQE for every | ||
| 835 | // write, so multiple writes in flight produce multiple CQEs | ||
| 836 | // (drained together by drain_wakeup_eventfd's single read of | ||
| 837 | // the eventfd counter). | ||
| 838 | 10863x | std::uint64_t v = 1; | |
| 839 | 10863x | std::ignore = ::write(wakeup_eventfd_, &v, sizeof(v)); | |
| 840 | } | ||
| 841 | |||
| 842 | inline void | ||
| 843 | 10702x | uring_scheduler::drain_wakeup_eventfd() const noexcept | |
| 844 | { | ||
| 845 | std::uint64_t v; | ||
| 846 | 10702x | std::ignore = ::read(wakeup_eventfd_, &v, sizeof(v)); | |
| 847 | |||
| 848 | // Multishot poll never needs re-arming. The poll-add was queued | ||
| 849 | // once at lazy_init_ring with IORING_POLL_ADD_MULTI; each eventfd | ||
| 850 | // POLLIN produces a CQE without consuming the SQE. | ||
| 851 | 10702x | } | |
| 852 | |||
| 853 | inline bool | ||
| 854 | 71x | uring_scheduler::prep_multishot_poll(int fd, void* data) noexcept | |
| 855 | { | ||
| 856 | // Prepare a multishot POLLIN SQE on `fd` tagged with `data`. Caller holds | ||
| 857 | // ring_mutex_ and flushes separately (re-arm sites ride the batch submit; | ||
| 858 | // register/init submit explicitly). Returns false when no SQE could be | ||
| 859 | // had after one flush, so a caller that has someone to report to does | ||
| 860 | // not mistake the submit of an empty queue for an armed poll. Shared by | ||
| 861 | // the wakeup-eventfd and signal self-pipe multishot polls. | ||
| 862 | 71x | ::io_uring_sqe* sqe = ::io_uring_get_sqe(&ring_); | |
| 863 | 71x | if (!sqe) | |
| 864 | { | ||
| 865 | 2x | ::io_uring_submit(&ring_); | |
| 866 | 2x | sqe = ::io_uring_get_sqe(&ring_); | |
| 867 | } | ||
| 868 | 71x | if (!sqe) | |
| 869 | 1x | return false; | |
| 870 | 70x | ::io_uring_prep_poll_multishot(sqe, fd, POLLIN); | |
| 871 | 70x | ::io_uring_sqe_set_data(sqe, data); | |
| 872 | 70x | return true; | |
| 873 | } | ||
| 874 | |||
| 875 | inline std::error_code | ||
| 876 | 66x | uring_scheduler::register_signal_reader(int read_fd) | |
| 877 | { | ||
| 878 | // Called once per service from add_signal(), holding neither the | ||
| 879 | // signal_state mutex nor the service mutex (see the call site). Submit a | ||
| 880 | // multishot POLL on the pipe read end; its CQE (tagged with | ||
| 881 | // &signal_pipe_sentinel_) is recognised in process_completions. The actual | ||
| 882 | // drain+deliver is deferred to dispatch context, so the only lock taken | ||
| 883 | // here is ring_mutex_ with no outer lock held, hence no lock-order | ||
| 884 | // inversion with the reactor drain path. | ||
| 885 | 66x | signal_pipe_read_fd_ = read_fd; | |
| 886 | 66x | lazy_init_ring(); | |
| 887 | |||
| 888 | 66x | lock_type lock(ring_mutex_); | |
| 889 | // EAGAIN, not ENOBUFS: an exhausted submission queue reports the | ||
| 890 | // same retryable code wherever it is hit, and the next add() is | ||
| 891 | // the retry. | ||
| 892 | 66x | if (!prep_multishot_poll(read_fd, &signal_pipe_sentinel_)) | |
| 893 | 1x | return make_err(EAGAIN); | |
| 894 | 65x | int rc = ::io_uring_submit(&ring_); | |
| 895 | 65x | if (rc < 0) | |
| 896 | 1x | return make_err(-rc); | |
| 897 | 64x | return {}; | |
| 898 | 66x | } | |
| 899 | |||
| 900 | inline void | ||
| 901 | 1843x | uring_scheduler::post(std::coroutine_handle<> h) const | |
| 902 | { | ||
| 903 | struct post_handler final : scheduler_op | ||
| 904 | { | ||
| 905 | std::coroutine_handle<> h_; | ||
| 906 | 1843x | explicit post_handler(std::coroutine_handle<> h) noexcept : h_(h) {} | |
| 907 | |||
| 908 | 1837x | void operator()() override | |
| 909 | { | ||
| 910 | 1837x | auto saved = h_; | |
| 911 | 1837x | delete this; | |
| 912 | 1837x | saved.resume(); | |
| 913 | 1837x | } | |
| 914 | |||
| 915 | 6x | void destroy() override | |
| 916 | { | ||
| 917 | 6x | auto saved = h_; | |
| 918 | 6x | delete this; | |
| 919 | 6x | if (saved) | |
| 920 | 6x | saved.destroy(); | |
| 921 | 6x | } | |
| 922 | }; | ||
| 923 | |||
| 924 | 1843x | auto* op = new post_handler(h); | |
| 925 | 1843x | lazy_init_ring(); | |
| 926 | 1843x | outstanding_work_.fetch_add(1, std::memory_order_relaxed); | |
| 927 | bool wake_leader; | ||
| 928 | { | ||
| 929 | 1843x | lock_type lock(dispatch_mutex_); | |
| 930 | 1843x | completed_ops_.push(op); | |
| 931 | 1843x | wake_leader = task_running_; | |
| 932 | 1843x | if (!wake_leader) | |
| 933 | 1839x | cond_.notify_one(); | |
| 934 | 1843x | } | |
| 935 | 1843x | if (wake_leader) | |
| 936 | 4x | interrupt_reactor(); | |
| 937 | 1843x | } | |
| 938 | |||
| 939 | inline void | ||
| 940 | 5831x | uring_scheduler::post(scheduler_op* op) const | |
| 941 | { | ||
| 942 | 5831x | lazy_init_ring(); | |
| 943 | 5831x | outstanding_work_.fetch_add(1, std::memory_order_relaxed); | |
| 944 | bool wake_leader; | ||
| 945 | { | ||
| 946 | 5831x | lock_type lock(dispatch_mutex_); | |
| 947 | 5831x | completed_ops_.push(op); | |
| 948 | 5831x | wake_leader = task_running_; | |
| 949 | 5831x | if (!wake_leader) | |
| 950 | 3526x | cond_.notify_one(); | |
| 951 | 5831x | } | |
| 952 | 5831x | if (wake_leader) | |
| 953 | 2305x | interrupt_reactor(); | |
| 954 | 5831x | } | |
| 955 | |||
| 956 | inline void | ||
| 957 | 8360x | uring_scheduler::post(capy::continuation& c) const | |
| 958 | { | ||
| 959 | 8360x | lazy_init_ring(); | |
| 960 | 8360x | outstanding_work_.fetch_add(1, std::memory_order_relaxed); | |
| 961 | bool wake_leader; | ||
| 962 | { | ||
| 963 | 8360x | lock_type lock(dispatch_mutex_); | |
| 964 | 8360x | completed_ops_.push(c); | |
| 965 | 8360x | wake_leader = task_running_; | |
| 966 | 8360x | if (!wake_leader) | |
| 967 | 8093x | cond_.notify_one(); | |
| 968 | 8360x | } | |
| 969 | 8360x | if (wake_leader) | |
| 970 | 267x | interrupt_reactor(); | |
| 971 | 8360x | } | |
| 972 | |||
| 973 | // Thread-local stack of frames for io_uring schedulers being run on the | ||
| 974 | // current thread. Holds the running-scheduler pointer (for | ||
| 975 | // running_in_this_thread reporting) and the inline completion budget | ||
| 976 | // used by the speculative non-blocking I/O path (plan 5j). Nesting | ||
| 977 | // stacks frames via prev_ so each scheduler gets its own budget. | ||
| 978 | struct uring_scheduler_frame | ||
| 979 | { | ||
| 980 | uring_scheduler const* sched; | ||
| 981 | uring_scheduler_frame* prev; | ||
| 982 | int inline_budget; | ||
| 983 | int inline_budget_max; | ||
| 984 | }; | ||
| 985 | |||
| 986 | inline thread_local uring_scheduler_frame* tl_running_scheduler_frame_ = | ||
| 987 | nullptr; | ||
| 988 | |||
| 989 | // Default inline budget. Matches reactor's initial budget (2). Adaptive | ||
| 990 | // ramp-up to a max is intentionally NOT implemented yet — keep it simple | ||
| 991 | // for plan 5j and revisit if benches show fairness issues. | ||
| 992 | inline constexpr int uring_inline_budget_initial = 2; | ||
| 993 | inline constexpr int uring_inline_budget_max = 16; | ||
| 994 | |||
| 995 | /// RAII guard: pushes a frame onto the thread's running-scheduler stack | ||
| 996 | /// on construction, restores the previous on destruction. Used by | ||
| 997 | /// run/run_one/wait_one/poll/poll_one to mark the running thread and | ||
| 998 | /// hold a fresh inline budget for speculative completions. | ||
| 999 | struct uring_run_guard | ||
| 1000 | { | ||
| 1001 | uring_scheduler_frame frame_; | ||
| 1002 | |||
| 1003 | 1868x | explicit uring_run_guard(uring_scheduler const* self) noexcept | |
| 1004 | 1868x | : frame_{ | |
| 1005 | self, tl_running_scheduler_frame_, uring_inline_budget_initial, | ||
| 1006 | uring_inline_budget_max} | ||
| 1007 | { | ||
| 1008 | 1868x | tl_running_scheduler_frame_ = &frame_; | |
| 1009 | 1868x | } | |
| 1010 | |||
| 1011 | 1868x | ~uring_run_guard() noexcept | |
| 1012 | { | ||
| 1013 | 1868x | tl_running_scheduler_frame_ = frame_.prev; | |
| 1014 | 1868x | } | |
| 1015 | }; | ||
| 1016 | |||
| 1017 | inline bool | ||
| 1018 | 4468x | uring_scheduler::running_in_this_thread() const noexcept | |
| 1019 | { | ||
| 1020 | 4468x | for (auto* f = tl_running_scheduler_frame_; f != nullptr; f = f->prev) | |
| 1021 | { | ||
| 1022 | 265x | if (f->sched == this) | |
| 1023 | 265x | return true; | |
| 1024 | } | ||
| 1025 | 4203x | return false; | |
| 1026 | } | ||
| 1027 | |||
| 1028 | inline void | ||
| 1029 | 62765x | uring_scheduler::reset_inline_budget() const noexcept | |
| 1030 | { | ||
| 1031 | 62765x | for (auto* f = tl_running_scheduler_frame_; f != nullptr; f = f->prev) | |
| 1032 | { | ||
| 1033 | 62765x | if (f->sched == this) | |
| 1034 | { | ||
| 1035 | 62765x | f->inline_budget = f->inline_budget_max; | |
| 1036 | 62765x | return; | |
| 1037 | } | ||
| 1038 | } | ||
| 1039 | } | ||
| 1040 | |||
| 1041 | inline bool | ||
| 1042 | 351157x | uring_scheduler::try_consume_inline_budget() const noexcept | |
| 1043 | { | ||
| 1044 | 351157x | for (auto* f = tl_running_scheduler_frame_; f != nullptr; f = f->prev) | |
| 1045 | { | ||
| 1046 | 351157x | if (f->sched == this) | |
| 1047 | { | ||
| 1048 | 351157x | if (f->inline_budget > 0) | |
| 1049 | { | ||
| 1050 | 330530x | --f->inline_budget; | |
| 1051 | 330530x | return true; | |
| 1052 | } | ||
| 1053 | 20627x | return false; | |
| 1054 | } | ||
| 1055 | } | ||
| 1056 | ✗ | return false; | |
| 1057 | } | ||
| 1058 | |||
| 1059 | inline std::size_t | ||
| 1060 | 778x | uring_scheduler::run() | |
| 1061 | { | ||
| 1062 | 778x | lazy_init_ring(); | |
| 1063 | 1556x | if (outstanding_work_.load(std::memory_order_acquire) == 0) | |
| 1064 | { | ||
| 1065 | 39x | stop(); | |
| 1066 | 39x | return 0; | |
| 1067 | } | ||
| 1068 | |||
| 1069 | 739x | uring_run_guard guard(this); | |
| 1070 | 739x | std::size_t n = 0; | |
| 1071 | for (;;) | ||
| 1072 | { | ||
| 1073 | 39337x | std::size_t r = do_one(-1); | |
| 1074 | 39336x | if (r) | |
| 1075 | { | ||
| 1076 | 38598x | if (n != (std::numeric_limits<std::size_t>::max)()) | |
| 1077 | 38598x | ++n; | |
| 1078 | 38598x | continue; | |
| 1079 | } | ||
| 1080 | 1478x | if (outstanding_work_.load(std::memory_order_acquire) == 0 || | |
| 1081 | 2x | stopped_.load(std::memory_order_acquire)) | |
| 1082 | 738x | break; | |
| 1083 | // do_one returned 0 but work still outstanding (e.g. timer | ||
| 1084 | // expiry dispatched async work). Continue. | ||
| 1085 | 38598x | } | |
| 1086 | 738x | return n; | |
| 1087 | 739x | } | |
| 1088 | |||
| 1089 | inline std::size_t | ||
| 1090 | 71x | uring_scheduler::run_one() | |
| 1091 | { | ||
| 1092 | 71x | lazy_init_ring(); | |
| 1093 | 142x | if (outstanding_work_.load(std::memory_order_acquire) == 0) | |
| 1094 | { | ||
| 1095 | 2x | stop(); | |
| 1096 | 2x | return 0; | |
| 1097 | } | ||
| 1098 | 69x | uring_run_guard guard(this); | |
| 1099 | 69x | return do_one(-1); | |
| 1100 | 69x | } | |
| 1101 | |||
| 1102 | inline std::size_t | ||
| 1103 | 1289x | uring_scheduler::wait_one(long usec) | |
| 1104 | { | ||
| 1105 | 1289x | lazy_init_ring(); | |
| 1106 | 2578x | if (outstanding_work_.load(std::memory_order_acquire) == 0) | |
| 1107 | { | ||
| 1108 | 266x | stop(); | |
| 1109 | 266x | return 0; | |
| 1110 | } | ||
| 1111 | 1023x | uring_run_guard guard(this); | |
| 1112 | 1023x | return do_one(usec); | |
| 1113 | 1023x | } | |
| 1114 | |||
| 1115 | inline std::size_t | ||
| 1116 | 36x | uring_scheduler::poll() | |
| 1117 | { | ||
| 1118 | 36x | lazy_init_ring(); | |
| 1119 | 72x | if (outstanding_work_.load(std::memory_order_acquire) == 0) | |
| 1120 | { | ||
| 1121 | 1x | stop(); | |
| 1122 | 1x | return 0; | |
| 1123 | } | ||
| 1124 | 35x | uring_run_guard guard(this); | |
| 1125 | 35x | std::size_t n = 0; | |
| 1126 | 68x | while (do_one(0)) | |
| 1127 | { | ||
| 1128 | 33x | if (n != (std::numeric_limits<std::size_t>::max)()) | |
| 1129 | 33x | ++n; | |
| 1130 | } | ||
| 1131 | 35x | return n; | |
| 1132 | 35x | } | |
| 1133 | |||
| 1134 | inline std::size_t | ||
| 1135 | 4x | uring_scheduler::poll_one() | |
| 1136 | { | ||
| 1137 | 4x | lazy_init_ring(); | |
| 1138 | 8x | if (outstanding_work_.load(std::memory_order_acquire) == 0) | |
| 1139 | { | ||
| 1140 | 2x | stop(); | |
| 1141 | 2x | return 0; | |
| 1142 | } | ||
| 1143 | 2x | uring_run_guard guard(this); | |
| 1144 | 2x | return do_one(0); | |
| 1145 | 2x | } | |
| 1146 | |||
| 1147 | inline std::size_t | ||
| 1148 | 40499x | uring_scheduler::do_one(long timeout_us) | |
| 1149 | { | ||
| 1150 | // Leader-follower: only one thread at a time may call | ||
| 1151 | // io_uring_submit_and_wait_timeout on a shared ring (liburing's | ||
| 1152 | // userspace head/tail bookkeeping is not thread-safe). Other | ||
| 1153 | // threads either dispatch ready ops from completed_ops_ or wait | ||
| 1154 | // on cond_ until the leader returns from the kernel. | ||
| 1155 | 40499x | if (stopped_.load(std::memory_order_acquire)) | |
| 1156 | 734x | return 0; | |
| 1157 | |||
| 1158 | // submit_sqes_op only pumps the ring once per SQE batch. If the user | ||
| 1159 | // keeps a non-empty completed_ops_ (e.g. timer with 0ns expiry as a | ||
| 1160 | // yield primitive), the leader-phase kernel pass below never runs | ||
| 1161 | // and CQEs accumulate in the ring forever — sub_request's read CQE | ||
| 1162 | // never gets drained and the bench spins. submit_and_get_events | ||
| 1163 | // (not plain submit) is required because IORING_SETUP_DEFER_TASKRUN | ||
| 1164 | // gates task work on IORING_ENTER_GETEVENTS. | ||
| 1165 | // | ||
| 1166 | // Gate the kernel pump on there being io_uring-specific work. The | ||
| 1167 | // check is performed under ring_mutex_ so a concurrent cross-thread | ||
| 1168 | // submitter cannot prep an SQE that we then race past — both this | ||
| 1169 | // path and uring_submit_op acquire ring_mutex_ before touching | ||
| 1170 | // the ring. When all three sources are empty (no io_uring ops in | ||
| 1171 | // flight needing DEFER_TASKRUN GETEVENTS, no userspace-pending | ||
| 1172 | // SQEs, no kernel-ready CQEs) a kernel entry would have no work — | ||
| 1173 | // saves ~8 pp of cycles on the no-I/O microbenchmark | ||
| 1174 | // (io_context:single_threaded). We deliberately do NOT include | ||
| 1175 | // outstanding_work_ here, because that counter mixes coroutine | ||
| 1176 | // posts (in completed_ops_) with io_uring work — IOCTX has many | ||
| 1177 | // coroutine posts and no io_uring work, and the kernel pump there | ||
| 1178 | // is pure overhead. | ||
| 1179 | 39765x | if (ring_inited_) | |
| 1180 | { | ||
| 1181 | 39765x | lock_type ring_lock(ring_mutex_); | |
| 1182 | 39765x | if (uring_inflight_.load(std::memory_order_acquire) != 0 || | |
| 1183 | 69883x | ::io_uring_sq_ready(&ring_) != 0 || | |
| 1184 | 30118x | ::io_uring_cq_ready(&ring_) != 0) | |
| 1185 | { | ||
| 1186 | 10583x | ::io_uring_submit_and_get_events(&ring_); | |
| 1187 | 10583x | process_completions(); | |
| 1188 | } | ||
| 1189 | 39765x | } | |
| 1190 | |||
| 1191 | // Drain expired timers eagerly, for the same reason the kernel CQE | ||
| 1192 | // pump runs unconditionally above: when completed_ops_ stays non- | ||
| 1193 | // empty (e.g. continuous loopback I/O whose CQEs land in the top- | ||
| 1194 | // of-do_one process_completions call), the leader-wait branch | ||
| 1195 | // below — the only other place process_expired() runs — is never | ||
| 1196 | // reached. Without this, stopper-timer-based shutdowns (and any | ||
| 1197 | // other timer dependent on a busy I/O loop yielding) deadlock. | ||
| 1198 | // | ||
| 1199 | // empty() is a single relaxed-acquire atomic load on | ||
| 1200 | // timer_service::cached_nearest_ns_ (lock-free, no clock_gettime). | ||
| 1201 | // Skipping process_expired() when no timer is registered avoids the | ||
| 1202 | // mutex + clock_gettime hot-path cost that dominates IOCTX cycles | ||
| 1203 | // (~25 pp on io_context:single_threaded). When a timer IS | ||
| 1204 | // registered the call runs exactly as before, preserving the | ||
| 1205 | // deadlock fix this guard was originally written to address. | ||
| 1206 | 39765x | if (!timer_svc_->empty()) | |
| 1207 | 30440x | timer_svc_->process_expired(); | |
| 1208 | |||
| 1209 | 39765x | lock_type lock(dispatch_mutex_); | |
| 1210 | for (;;) | ||
| 1211 | { | ||
| 1212 | 42298x | if (stopped_.load(std::memory_order_acquire)) | |
| 1213 | 14x | return 0; | |
| 1214 | |||
| 1215 | 42284x | if (auto e = completed_ops_.pop()) | |
| 1216 | { | ||
| 1217 | // Hand off any remaining queued work to a follower so we | ||
| 1218 | // dispatch in parallel. | ||
| 1219 | 39562x | if (!completed_ops_.empty()) | |
| 1220 | 24366x | cond_.notify_one(); | |
| 1221 | 39562x | lock.unlock(); | |
| 1222 | // Speculative follow-ups in the handler share this budget. | ||
| 1223 | 39562x | reset_inline_budget(); | |
| 1224 | 39562x | if (ready_is_continuation(e)) | |
| 1225 | 8356x | ready_as_cont(e)->h.resume(); | |
| 1226 | else | ||
| 1227 | 31206x | (*ready_as_op(e))(); | |
| 1228 | 39562x | work_finished(); | |
| 1229 | 39562x | return 1; | |
| 1230 | } | ||
| 1231 | |||
| 1232 | 5444x | if (outstanding_work_.load(std::memory_order_acquire) == 0) | |
| 1233 | 1x | return 0; | |
| 1234 | |||
| 1235 | 2721x | if (task_running_) | |
| 1236 | { | ||
| 1237 | // LCOV_EXCL_START: follower path — entered only when another | ||
| 1238 | // thread holds ring leadership at the instant this thread | ||
| 1239 | // checks. That leader/follower timing overlap is a race no | ||
| 1240 | // deterministic test reliably forces (see the io_uring | ||
| 1241 | // mt_scheduler regression test, which exercises MT run but | ||
| 1242 | // does not pin this window). | ||
| 1243 | // Another thread holds leadership; either return (poll) | ||
| 1244 | // or wait for it to deliver work / release leadership. | ||
| 1245 | − | if (timeout_us == 0) | |
| 1246 | − | return 0; | |
| 1247 | − | if (timeout_us < 0) | |
| 1248 | − | cond_.wait(lock); | |
| 1249 | else | ||
| 1250 | { | ||
| 1251 | − | cond_.wait_for(lock, std::chrono::microseconds(timeout_us)); | |
| 1252 | // wait_one honoured its timeout; if nothing arrived, | ||
| 1253 | // return rather than re-arm. | ||
| 1254 | − | if (completed_ops_.empty() && | |
| 1255 | − | !stopped_.load(std::memory_order_acquire)) | |
| 1256 | − | return 0; | |
| 1257 | } | ||
| 1258 | − | continue; | |
| 1259 | // LCOV_EXCL_STOP | ||
| 1260 | } | ||
| 1261 | |||
| 1262 | // Become the leader: run the kernel poll. We drop the lock | ||
| 1263 | // for the blocking wait, then take it back to release | ||
| 1264 | // leadership and wake any follower that should pick up new | ||
| 1265 | // work. | ||
| 1266 | 2714x | __kernel_timespec ts{}; | |
| 1267 | 2714x | __kernel_timespec* ts_ptr = nullptr; | |
| 1268 | 2714x | auto next_expiry = timer_svc_->nearest_expiry(); | |
| 1269 | 2714x | auto now = std::chrono::steady_clock::now(); | |
| 1270 | |||
| 1271 | 2714x | if (timeout_us == 0) | |
| 1272 | { | ||
| 1273 | 24x | ts.tv_sec = 0; | |
| 1274 | 24x | ts.tv_nsec = 0; | |
| 1275 | 24x | ts_ptr = &ts; | |
| 1276 | } | ||
| 1277 | 2690x | else if (next_expiry != timer_service::time_point::max()) | |
| 1278 | { | ||
| 1279 | auto delta_ns = | ||
| 1280 | 2227x | std::chrono::duration_cast<std::chrono::nanoseconds>( | |
| 1281 | 2227x | next_expiry - now) | |
| 1282 | 2227x | .count(); | |
| 1283 | 2227x | if (delta_ns < 0) | |
| 1284 | ✗ | delta_ns = 0; | |
| 1285 | 2227x | ts.tv_sec = delta_ns / 1'000'000'000; | |
| 1286 | 2227x | ts.tv_nsec = delta_ns % 1'000'000'000; | |
| 1287 | 2227x | ts_ptr = &ts; | |
| 1288 | } | ||
| 1289 | 463x | else if (timeout_us > 0) | |
| 1290 | { | ||
| 1291 | 428x | ts.tv_sec = timeout_us / 1'000'000; | |
| 1292 | 428x | ts.tv_nsec = (timeout_us % 1'000'000) * 1000; | |
| 1293 | 428x | ts_ptr = &ts; | |
| 1294 | } | ||
| 1295 | else | ||
| 1296 | { | ||
| 1297 | // run() with no pending timers: cap the kernel wait at 1s | ||
| 1298 | // so the leader periodically re-checks state. Defense in | ||
| 1299 | // depth against a lost wakeup (e.g. multishot poll on the | ||
| 1300 | // wakeup eventfd terminates and the re-arm SQE doesn't | ||
| 1301 | // reach the kernel in time). Worst case: one extra | ||
| 1302 | // wake-up per io_context per second when truly idle. | ||
| 1303 | 35x | ts.tv_sec = 1; | |
| 1304 | 35x | ts.tv_nsec = 0; | |
| 1305 | 35x | ts_ptr = &ts; | |
| 1306 | } | ||
| 1307 | |||
| 1308 | 2714x | task_running_ = true; | |
| 1309 | 2714x | lock.unlock(); | |
| 1310 | |||
| 1311 | // Three-phase kernel wait, matching Boost.Asio's | ||
| 1312 | // io_uring_service::run pattern. ring_mutex_ is held briefly | ||
| 1313 | // to push pending SQEs and to drain CQEs, but NOT during | ||
| 1314 | // the blocking io_uring_wait_cqe_timeout. Cross-thread | ||
| 1315 | // submitters (uring_submit_op, cancel paths) can take | ||
| 1316 | // ring_mutex_ during the wait and prep new SQEs without | ||
| 1317 | // blocking on the leader; their wake eventfd write fires the | ||
| 1318 | // multishot poll and returns the leader from wait_cqe_timeout | ||
| 1319 | // promptly. | ||
| 1320 | // | ||
| 1321 | // Phase 1 — submit any pending SQEs to the kernel. | ||
| 1322 | { | ||
| 1323 | 2714x | lock_type ring_lock(ring_mutex_); | |
| 1324 | 2714x | ::io_uring_submit(&ring_); | |
| 1325 | 2714x | } | |
| 1326 | |||
| 1327 | // Phase 2 — wait for at least one CQE without holding the | ||
| 1328 | // mutex. Multi-thread `io_uring_enter` is permitted without | ||
| 1329 | // SINGLE_ISSUER. wait_cqe_timeout only peeks the CQ ring; | ||
| 1330 | // head advancement happens under the mutex in | ||
| 1331 | // process_completions below. | ||
| 1332 | 2714x | ::io_uring_cqe* cqe = nullptr; | |
| 1333 | 2714x | int rc = ::io_uring_wait_cqe_timeout(&ring_, &cqe, ts_ptr); | |
| 1334 | |||
| 1335 | // Phase 3 — drain CQEs under the mutex. | ||
| 1336 | { | ||
| 1337 | 2714x | lock_type ring_lock(ring_mutex_); | |
| 1338 | 2714x | if (rc == 0 || rc == -ETIME || rc == -EINTR) | |
| 1339 | 2713x | process_completions(); | |
| 1340 | 2714x | } | |
| 1341 | |||
| 1342 | 2714x | if (rc < 0 && rc != -ETIME && rc != -EINTR) | |
| 1343 | { | ||
| 1344 | // Restore state before propagating so followers don't | ||
| 1345 | // deadlock waiting for a leader that never returns. | ||
| 1346 | 1x | lock.lock(); | |
| 1347 | 1x | task_running_ = false; | |
| 1348 | 1x | cond_.notify_all(); | |
| 1349 | 1x | detail::throw_system_error( | |
| 1350 | 2x | make_err(-rc), "io_uring_wait_cqe_timeout"); | |
| 1351 | } | ||
| 1352 | |||
| 1353 | 2713x | if (!timer_svc_->empty()) | |
| 1354 | 2232x | timer_svc_->process_expired(); | |
| 1355 | |||
| 1356 | 2713x | lock.lock(); | |
| 1357 | 2713x | task_running_ = false; | |
| 1358 | 2713x | cond_.notify_all(); | |
| 1359 | |||
| 1360 | // For poll() / wait_one() we honour the timeout: one kernel | ||
| 1361 | // pass is the contract. If still nothing dispatchable, exit. | ||
| 1362 | // For run() (timeout < 0) keep looping until work arrives or | ||
| 1363 | // someone calls stop(). | ||
| 1364 | 2713x | if (timeout_us >= 0 && completed_ops_.empty()) | |
| 1365 | 187x | return 0; | |
| 1366 | 2533x | } | |
| 1367 | 39765x | } | |
| 1368 | |||
| 1369 | inline void | ||
| 1370 | 13296x | uring_scheduler::process_completions() | |
| 1371 | { | ||
| 1372 | unsigned head; | ||
| 1373 | ::io_uring_cqe* cqe; | ||
| 1374 | 13296x | unsigned consumed = 0; | |
| 1375 | |||
| 1376 | // Collect completed I/O ops locally; splice into completed_ops_ | ||
| 1377 | // after the loop so do_one dispatches them one at a time. | ||
| 1378 | 13296x | ready_queue local_ops; | |
| 1379 | |||
| 1380 | 13296x | std::int64_t inflight_dec = 0; | |
| 1381 | 32383x | io_uring_for_each_cqe(&ring_, head, cqe) | |
| 1382 | { | ||
| 1383 | 19087x | void* ud = io_uring_cqe_get_data(cqe); | |
| 1384 | 19087x | if (ud == nullptr) | |
| 1385 | { | ||
| 1386 | // Wakeup eventfd CQE: drain the eventfd byte. Not counted | ||
| 1387 | // by uring_inflight_; we never incremented for the | ||
| 1388 | // wakeup multishot SQE (its progress doesn't depend on | ||
| 1389 | // userspace getevents). | ||
| 1390 | 9971x | drain_wakeup_eventfd(); | |
| 1391 | // If multishot terminated (kernel dropped under memory | ||
| 1392 | // pressure or similar), re-arm. Each CQE except the last | ||
| 1393 | // sets IORING_CQE_F_MORE. | ||
| 1394 | // Best-effort: nothing here has a caller to report to, and | ||
| 1395 | // prep_multishot_poll already flushed once to make room. | ||
| 1396 | 9971x | if ((cqe->flags & IORING_CQE_F_MORE) == 0) | |
| 1397 | 3x | std::ignore = prep_multishot_poll(wakeup_eventfd_, nullptr); | |
| 1398 | } | ||
| 1399 | 9116x | else if (ud == &cancel_sentinel_) | |
| 1400 | { | ||
| 1401 | // CQE for an ASYNC_CANCEL op — ignore; the actual op's | ||
| 1402 | // CQE arrives separately and is dispatched via cqe_func. | ||
| 1403 | // Cancels are one-shot, no F_MORE, decrement inflight. | ||
| 1404 | 4138x | ++inflight_dec; | |
| 1405 | } | ||
| 1406 | 4978x | else if (ud == &signal_pipe_sentinel_) | |
| 1407 | { | ||
| 1408 | // Signal self-pipe readiness. Re-arm the multishot poll if it | ||
| 1409 | // terminated (F_MORE cleared), then enqueue signal_drain_op_ to | ||
| 1410 | // drain + deliver in dispatch context. Not counted in | ||
| 1411 | // uring_inflight_ (like the wakeup eventfd poll): its progress | ||
| 1412 | // does not gate DEFER_TASKRUN GETEVENTS. | ||
| 1413 | 161x | if ((cqe->flags & IORING_CQE_F_MORE) == 0) | |
| 1414 | 1x | std::ignore = prep_multishot_poll( | |
| 1415 | 1x | signal_pipe_read_fd_, &signal_pipe_sentinel_); | |
| 1416 | 161x | bool expected = false; | |
| 1417 | 161x | if (signal_drain_op_.queued_.compare_exchange_strong( | |
| 1418 | expected, true, std::memory_order_acq_rel, | ||
| 1419 | std::memory_order_relaxed)) | ||
| 1420 | { | ||
| 1421 | // Balance the work_finished() do_one runs after dispatching. | ||
| 1422 | 159x | work_started(); | |
| 1423 | 159x | local_ops.push(&signal_drain_op_); | |
| 1424 | } | ||
| 1425 | } | ||
| 1426 | else | ||
| 1427 | { | ||
| 1428 | 4817x | auto* iop = static_cast<uring_op*>(ud); | |
| 1429 | 4817x | if (iop->retired) | |
| 1430 | { | ||
| 1431 | // The owner handed this op to retire_op and is no | ||
| 1432 | // longer listening. Never dispatch cqe_func — the | ||
| 1433 | // handler's output pointers and coroutine are gone. | ||
| 1434 | // retire_func disposes of whatever the result owns | ||
| 1435 | // (an accepted descriptor, say), since only the op | ||
| 1436 | // type knows what `res` means. | ||
| 1437 | 4x | if (iop->retire_func) | |
| 1438 | 4x | (*iop->retire_func)(iop, cqe->res, cqe->flags); | |
| 1439 | 4x | if ((cqe->flags & IORING_CQE_F_MORE) == 0) | |
| 1440 | { | ||
| 1441 | 4x | ++inflight_dec; | |
| 1442 | 4x | release_retired_op(iop); | |
| 1443 | } | ||
| 1444 | } | ||
| 1445 | else | ||
| 1446 | { | ||
| 1447 | 4813x | (*iop->cqe_func)(iop, cqe->res, cqe->flags, local_ops); | |
| 1448 | // Decrement inflight on the terminal CQE only — multishot | ||
| 1449 | // ops (acceptor) hold the SQE alive across F_MORE CQEs and | ||
| 1450 | // free it only when F_MORE is cleared. | ||
| 1451 | 4813x | if ((cqe->flags & IORING_CQE_F_MORE) == 0) | |
| 1452 | 2803x | ++inflight_dec; | |
| 1453 | } | ||
| 1454 | } | ||
| 1455 | 19087x | ++consumed; | |
| 1456 | } | ||
| 1457 | 13296x | if (inflight_dec) | |
| 1458 | 6651x | uring_inflight_.fetch_sub(inflight_dec, std::memory_order_acq_rel); | |
| 1459 | |||
| 1460 | 13296x | if (consumed) | |
| 1461 | 10213x | io_uring_cq_advance(&ring_, consumed); | |
| 1462 | |||
| 1463 | // Caller holds ring_mutex_. Take dispatch_mutex_ briefly to | ||
| 1464 | // splice locally-collected ops onto the global queue (lock order | ||
| 1465 | // ring_mutex_ -> dispatch_mutex_). | ||
| 1466 | 13296x | if (!local_ops.empty()) | |
| 1467 | { | ||
| 1468 | 2931x | lock_type lock(dispatch_mutex_); | |
| 1469 | 2931x | completed_ops_.splice(local_ops); | |
| 1470 | // Wake any follower waiting on cond_; it'll pop and dispatch. | ||
| 1471 | 2931x | cond_.notify_one(); | |
| 1472 | 2931x | } | |
| 1473 | 13296x | } | |
| 1474 | |||
| 1475 | inline void | ||
| 1476 | 29x | uring_scheduler::submit_sqes_op::do_handler( | |
| 1477 | void* owner, | ||
| 1478 | scheduler_op* base, | ||
| 1479 | std::uint32_t /*bytes*/, | ||
| 1480 | std::uint32_t /*error*/) noexcept | ||
| 1481 | { | ||
| 1482 | 29x | if (owner == nullptr) | |
| 1483 | 29x | return; // shutdown drain — nothing to do; SQE storage is | |
| 1484 | // kernel-mapped and discarded by io_uring_queue_exit. | ||
| 1485 | |||
| 1486 | ✗ | auto* self = static_cast<submit_sqes_op*>(base); | |
| 1487 | ✗ | auto* sched = self->sched_; | |
| 1488 | |||
| 1489 | ✗ | uring_scheduler::lock_type ring_lock(sched->ring_mutex_); | |
| 1490 | ✗ | sched->submit_op_posted_ = false; | |
| 1491 | ✗ | ::io_uring_submit_and_get_events(&sched->ring_); | |
| 1492 | ✗ | sched->process_completions(); | |
| 1493 | ✗ | } | |
| 1494 | |||
| 1495 | inline void | ||
| 1496 | 218x | uring_scheduler::submit_cancel_by_user_data(uring_op* target) noexcept | |
| 1497 | { | ||
| 1498 | 218x | lazy_init_ring(); | |
| 1499 | // Wake the leader (if any) so its submit_and_wait_timeout returns | ||
| 1500 | // and releases ring_mutex_; otherwise we'd block here until the | ||
| 1501 | // next CQE arrives organically. Cancellation is best-effort if | ||
| 1502 | // the SQ stays full after one flush — the op completes on its | ||
| 1503 | // own and reports cancelled via the in-flight `cancelled` flag. | ||
| 1504 | 218x | interrupt_reactor(); | |
| 1505 | 218x | lock_type lock(ring_mutex_); | |
| 1506 | 218x | io_uring_sqe* sqe = io_uring_get_sqe(&ring_); | |
| 1507 | 218x | if (!sqe) | |
| 1508 | { | ||
| 1509 | 3x | io_uring_submit(&ring_); | |
| 1510 | 3x | sqe = io_uring_get_sqe(&ring_); | |
| 1511 | } | ||
| 1512 | 218x | if (!sqe) | |
| 1513 | 3x | return; | |
| 1514 | |||
| 1515 | 215x | io_uring_prep_cancel(sqe, target, 0); | |
| 1516 | 215x | io_uring_sqe_set_data(sqe, &cancel_sentinel_); | |
| 1517 | 215x | inflight_inc(); | |
| 1518 | 218x | } | |
| 1519 | |||
| 1520 | inline void | ||
| 1521 | 97x | uring_scheduler::submit_cancel_by_fd(int fd) noexcept | |
| 1522 | { | ||
| 1523 | 97x | lazy_init_ring(); | |
| 1524 | 97x | interrupt_reactor(); | |
| 1525 | 97x | lock_type lock(ring_mutex_); | |
| 1526 | 97x | io_uring_sqe* sqe = io_uring_get_sqe(&ring_); | |
| 1527 | 97x | if (!sqe) | |
| 1528 | { | ||
| 1529 | 1x | io_uring_submit(&ring_); | |
| 1530 | 1x | sqe = io_uring_get_sqe(&ring_); | |
| 1531 | } | ||
| 1532 | 97x | if (!sqe) | |
| 1533 | 1x | return; | |
| 1534 | |||
| 1535 | 96x | io_uring_prep_cancel_fd(sqe, fd, IORING_ASYNC_CANCEL_ALL); | |
| 1536 | 96x | io_uring_sqe_set_data(sqe, &cancel_sentinel_); | |
| 1537 | 96x | inflight_inc(); | |
| 1538 | 97x | } | |
| 1539 | |||
| 1540 | inline void | ||
| 1541 | 218x | uring_op::on_cancel() noexcept | |
| 1542 | { | ||
| 1543 | 218x | request_cancel(); // coro_op: records the cancellation (sets the flag) | |
| 1544 | // Skip the cancel SQE if we never linked an SQE to this op — the | ||
| 1545 | // bypass path in the caller will see cancelled=true and complete | ||
| 1546 | // synchronously without a kernel round-trip. | ||
| 1547 | 218x | if (sched_ && sqe_set.load(std::memory_order_acquire)) | |
| 1548 | 218x | sched_->submit_cancel_by_user_data(this); | |
| 1549 | 218x | } | |
| 1550 | |||
| 1551 | inline void | ||
| 1552 | 4857x | uring_scheduler::cancel_and_flush(int fd) noexcept | |
| 1553 | { | ||
| 1554 | // The flush can execute a queued write on `fd` inline; when the | ||
| 1555 | // fd is a pipe whose reader has already closed — service | ||
| 1556 | // shutdown closes impls one at a time, so teardown itself | ||
| 1557 | // creates that state — the kernel raises SIGPIPE. | ||
| 1558 | 4857x | scoped_sigpipe_block no_sigpipe; | |
| 1559 | |||
| 1560 | 4857x | lazy_init_ring(); | |
| 1561 | 4857x | interrupt_reactor(); | |
| 1562 | 4857x | lock_type lock(ring_mutex_); | |
| 1563 | 4857x | io_uring_sqe* sqe = io_uring_get_sqe(&ring_); | |
| 1564 | 4857x | if (!sqe) | |
| 1565 | { | ||
| 1566 | 9x | io_uring_submit(&ring_); | |
| 1567 | 9x | sqe = io_uring_get_sqe(&ring_); | |
| 1568 | } | ||
| 1569 | 4857x | if (sqe) | |
| 1570 | { | ||
| 1571 | 4851x | io_uring_prep_cancel_fd(sqe, fd, IORING_ASYNC_CANCEL_ALL); | |
| 1572 | 4851x | io_uring_sqe_set_data(sqe, &cancel_sentinel_); | |
| 1573 | 4851x | inflight_inc(); | |
| 1574 | } | ||
| 1575 | // Flush while fd is still open so the kernel resolves the file | ||
| 1576 | // from the fd number before the caller closes and recycles it. | ||
| 1577 | 4857x | io_uring_submit(&ring_); | |
| 1578 | 4857x | } | |
| 1579 | |||
| 1580 | inline void | ||
| 1581 | 4x | uring_scheduler::release_retired_op(uring_op* op) noexcept | |
| 1582 | { | ||
| 1583 | // Called from the CQE loop with ring_mutex_ held; takes | ||
| 1584 | // retired_mutex_ under it, matching retire_op's order. Deleting | ||
| 1585 | // here is safe only because retire_op requires the op to hold no | ||
| 1586 | // reference back to its owner, so ~op cannot re-enter the | ||
| 1587 | // scheduler or touch a destroyed acceptor. | ||
| 1588 | 4x | std::unique_ptr<uring_op> owned; | |
| 1589 | { | ||
| 1590 | 4x | std::lock_guard<std::mutex> lock(retired_mutex_); | |
| 1591 | 4x | for (auto it = retired_ops_.begin(); it != retired_ops_.end(); ++it) | |
| 1592 | { | ||
| 1593 | 4x | if (it->get() == op) | |
| 1594 | { | ||
| 1595 | 4x | owned = std::move(*it); | |
| 1596 | 4x | retired_ops_.erase(it); | |
| 1597 | 4x | break; | |
| 1598 | } | ||
| 1599 | } | ||
| 1600 | 4x | } | |
| 1601 | 4x | } | |
| 1602 | |||
| 1603 | inline void | ||
| 1604 | 215x | uring_scheduler::drain_cqes_for(uring_op* target) noexcept | |
| 1605 | { | ||
| 1606 | 215x | lazy_init_ring(); | |
| 1607 | // Submit a cancel by user_data so the kernel returns CQEs for | ||
| 1608 | // the target promptly, then iterate the CQ ring and consume | ||
| 1609 | // every CQE that matches `target`. ring_mutex_ serializes against | ||
| 1610 | // the leader's kernel wait and any concurrent cancel path; the | ||
| 1611 | // interrupt_reactor() ensures the leader returns promptly so we | ||
| 1612 | // can take the mutex. | ||
| 1613 | 215x | interrupt_reactor(); | |
| 1614 | { | ||
| 1615 | 215x | lock_type lock(ring_mutex_); | |
| 1616 | 215x | if (auto* sqe = io_uring_get_sqe(&ring_)) | |
| 1617 | { | ||
| 1618 | 211x | io_uring_prep_cancel(sqe, target, 0); | |
| 1619 | 211x | io_uring_sqe_set_data(sqe, &cancel_sentinel_); | |
| 1620 | 211x | inflight_inc(); | |
| 1621 | } | ||
| 1622 | 215x | io_uring_submit(&ring_); | |
| 1623 | 215x | } | |
| 1624 | |||
| 1625 | // Loop a few rounds: cancel SQE submission, then drain CQEs. | ||
| 1626 | // Bounded loop avoids stalls if the kernel never returns a | ||
| 1627 | // cancel completion — best-effort. | ||
| 1628 | 238x | for (int rounds = 0; rounds < drain_cqes_max_rounds; ++rounds) | |
| 1629 | { | ||
| 1630 | 238x | lock_type lock(ring_mutex_); | |
| 1631 | |||
| 1632 | unsigned head; | ||
| 1633 | ::io_uring_cqe* cqe; | ||
| 1634 | 238x | unsigned consumed = 0; | |
| 1635 | 238x | bool saw_target = false; | |
| 1636 | 238x | std::int64_t inflight_dec = 0; | |
| 1637 | |||
| 1638 | 1697x | io_uring_for_each_cqe(&ring_, head, cqe) | |
| 1639 | { | ||
| 1640 | // Mirror process_completions' uring_inflight_ accounting. | ||
| 1641 | // That counter gates the do_one ring pump, so every CQE we | ||
| 1642 | // advance past here must adjust it exactly as the normal | ||
| 1643 | // drain would — otherwise it drifts upward (each teardown | ||
| 1644 | // leaks the counts of the CQEs it swallows), defeating the | ||
| 1645 | // idle-skip optimisation for the lifetime of the io_context. | ||
| 1646 | // We do NOT dispatch real ops — the target is being | ||
| 1647 | // destructed and siblings may already be freed — but we still | ||
| 1648 | // account for and house-keep each CQE we consume. | ||
| 1649 | 1459x | void* ud = io_uring_cqe_get_data(cqe); | |
| 1650 | 1459x | if (ud == nullptr) | |
| 1651 | { | ||
| 1652 | // Wakeup eventfd CQE — our own interrupt_reactor() above | ||
| 1653 | // very likely produced one. Drain the byte and re-arm if | ||
| 1654 | // the multishot terminated, exactly as process_completions | ||
| 1655 | // does. Never incremented, so never decremented. | ||
| 1656 | 731x | drain_wakeup_eventfd(); | |
| 1657 | 731x | if ((cqe->flags & IORING_CQE_F_MORE) == 0) | |
| 1658 | 1x | std::ignore = prep_multishot_poll(wakeup_eventfd_, nullptr); | |
| 1659 | } | ||
| 1660 | 728x | else if (ud == &signal_pipe_sentinel_) | |
| 1661 | { | ||
| 1662 | // Signal self-pipe readiness. Re-arm if the multishot | ||
| 1663 | // terminated; the still-readable pipe re-fires on the next | ||
| 1664 | // kernel enter so process_completions delivers the signal — | ||
| 1665 | // we deliberately do NOT enqueue signal_drain_op_ from this | ||
| 1666 | // teardown path. Not counted by uring_inflight_ (the poll | ||
| 1667 | // was armed via prep_multishot_poll, which never increments), | ||
| 1668 | // so it must NOT be decremented. | ||
| 1669 | ✗ | if ((cqe->flags & IORING_CQE_F_MORE) == 0) | |
| 1670 | ✗ | std::ignore = prep_multishot_poll( | |
| 1671 | ✗ | signal_pipe_read_fd_, &signal_pipe_sentinel_); | |
| 1672 | } | ||
| 1673 | 728x | else if (ud == &cancel_sentinel_) | |
| 1674 | { | ||
| 1675 | // ASYNC_CANCEL CQE (one-shot, no F_MORE), including the | ||
| 1676 | // cancel SQE we submitted just above. Decrement inflight. | ||
| 1677 | 523x | ++inflight_dec; | |
| 1678 | } | ||
| 1679 | 205x | else if (ud == target) | |
| 1680 | { | ||
| 1681 | 201x | saw_target = true; | |
| 1682 | // Don't dispatch — caller is destructing target; just | ||
| 1683 | // consume so the CQE doesn't dangle. Decrement inflight on | ||
| 1684 | // the terminal CQE only: the target is a multishot op | ||
| 1685 | // whose intermediate CQEs carry F_MORE. | ||
| 1686 | 201x | if ((cqe->flags & IORING_CQE_F_MORE) == 0) | |
| 1687 | 199x | ++inflight_dec; | |
| 1688 | } | ||
| 1689 | else | ||
| 1690 | { | ||
| 1691 | // Some other op's CQE. Intentionally NOT dispatched: it | ||
| 1692 | // may belong to an op freed by a sibling teardown (other | ||
| 1693 | // acceptors / sockets), and dispatching would UAF. We | ||
| 1694 | // still account for its terminal CQE so inflight stays | ||
| 1695 | // balanced — the submit that produced it incremented the | ||
| 1696 | // counter. The io_context's destructor sequence runs | ||
| 1697 | // services' shutdowns before ~scheduler, so any still-live | ||
| 1698 | // ops drain through their own paths first. | ||
| 1699 | 4x | if ((cqe->flags & IORING_CQE_F_MORE) == 0) | |
| 1700 | 4x | ++inflight_dec; | |
| 1701 | } | ||
| 1702 | 1459x | ++consumed; | |
| 1703 | } | ||
| 1704 | 238x | if (consumed) | |
| 1705 | { | ||
| 1706 | 218x | io_uring_cq_advance(&ring_, consumed); | |
| 1707 | 218x | if (inflight_dec) | |
| 1708 | 213x | uring_inflight_.fetch_sub( | |
| 1709 | inflight_dec, std::memory_order_acq_rel); | ||
| 1710 | 218x | if (saw_target) | |
| 1711 | 199x | break; | |
| 1712 | 19x | continue; | |
| 1713 | } | ||
| 1714 | |||
| 1715 | // Nothing in the CQ — kick the kernel briefly. Hold | ||
| 1716 | // ring_mutex_ across the wait so we don't race with the | ||
| 1717 | // run-loop leader. | ||
| 1718 | 20x | __kernel_timespec ts{0, static_cast<long long>(drain_cqes_kick_ns)}; | |
| 1719 | 20x | ::io_uring_cqe* one = nullptr; | |
| 1720 | int rc = | ||
| 1721 | 20x | ::io_uring_submit_and_wait_timeout(&ring_, &one, 1, &ts, nullptr); | |
| 1722 | 20x | if (rc < 0 && rc != -ETIME && rc != -EINTR) | |
| 1723 | 1x | break; | |
| 1724 | 19x | if (rc == -ETIME) | |
| 1725 | 15x | break; | |
| 1726 | 238x | } | |
| 1727 | 215x | } | |
| 1728 | |||
| 1729 | } // namespace boost::corosio::detail | ||
| 1730 | |||
| 1731 | #endif // BOOST_COROSIO_HAS_URING | ||
| 1732 | |||
| 1733 | #endif // BOOST_COROSIO_NATIVE_DETAIL_URING_URING_SCHEDULER_HPP | ||
| 1734 |