include/boost/corosio/native/detail/io_uring/io_uring_scheduler.hpp

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