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

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