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)
io_uring_scheduler.hpp
f(x) 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_, &params);
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