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