include/boost/corosio/native/detail/uring/uring_scheduler.hpp

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