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