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

99.3% Lines (266/0/268) 100.0% List of functions (30/0/30)
uring_socket_ops.hpp
f(x) Functions (30)
Function Calls Lines Blocks
boost::corosio::detail::copy_to_iovec(boost::corosio::buffer_param const&, iovec (&) [16]) :67 396517x 100.0% 100.0% boost::corosio::detail::uring_set_result(boost::corosio::detail::uring_op*, bool, bool) :90 25697x 100.0% 100.0% boost::corosio::detail::uring_read_op::uring_read_op() :111 10001x 100.0% 100.0% boost::corosio::detail::uring_read_op::prepare(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::error_code*, unsigned long*, int, boost::corosio::detail::uring_scheduler*, std::shared_ptr<void>, boost::corosio::detail::speculative_state*, boost::corosio::buffer_param, std::stop_token const&) :124 11211x 100.0% 100.0% boost::corosio::detail::uring_read_op::do_prep(boost::corosio::detail::uring_op*, io_uring_sqe*) :151 245x 100.0% 100.0% boost::corosio::detail::uring_read_op::do_cqe(boost::corosio::detail::uring_op*, int, unsigned int, boost::corosio::detail::ready_queue&) :172 243x 100.0% 100.0% boost::corosio::detail::uring_read_op::do_handler(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int) :180 11209x 100.0% 100.0% boost::corosio::detail::uring_write_op::uring_write_op() :223 10001x 100.0% 100.0% boost::corosio::detail::uring_write_op::prepare(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::error_code*, unsigned long*, int, boost::corosio::detail::uring_scheduler*, std::shared_ptr<void>, boost::corosio::detail::speculative_state*, boost::corosio::buffer_param, std::stop_token const&) :226 10973x 100.0% 100.0% boost::corosio::detail::uring_write_op::do_prep(boost::corosio::detail::uring_op*, io_uring_sqe*) :259 5x 100.0% 100.0% boost::corosio::detail::uring_write_op::do_cqe(boost::corosio::detail::uring_op*, int, unsigned int, boost::corosio::detail::ready_queue&) :278 4x 100.0% 100.0% boost::corosio::detail::uring_write_op::do_handler(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int) :286 10972x 100.0% 100.0% boost::corosio::detail::uring_connect_op::uring_connect_op() :330 9974x 100.0% 100.0% boost::corosio::detail::uring_connect_op::prepare(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::error_code*, int, boost::corosio::detail::uring_scheduler*, std::shared_ptr<void>, boost::corosio::endpoint, boost::corosio::endpoint*, boost::corosio::endpoint*, std::stop_token const&) :340 3286x 100.0% 100.0% boost::corosio::detail::uring_connect_op::do_prep(boost::corosio::detail::uring_op*, io_uring_sqe*) :368 3284x 100.0% 100.0% boost::corosio::detail::uring_connect_op::do_cqe(boost::corosio::detail::uring_op*, int, unsigned int, boost::corosio::detail::ready_queue&) :377 3283x 100.0% 100.0% boost::corosio::detail::uring_connect_op::do_handler(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int) :385 3285x 100.0% 100.0% boost::corosio::detail::uring_do_submit_op(boost::corosio::detail::uring_scheduler&, boost::corosio::detail::uring_op*, bool) :440 4111x 100.0% 100.0% boost::corosio::detail::uring_submit_op(boost::corosio::detail::uring_scheduler&, boost::corosio::detail::uring_op*) :535 3889x 100.0% 100.0% boost::corosio::detail::uring_try_submit_op(boost::corosio::detail::uring_scheduler&, boost::corosio::detail::uring_op*) :566 222x 100.0% 100.0% boost::corosio::detail::uring_wait_op::uring_wait_op() :587 10531x 100.0% 100.0% boost::corosio::detail::uring_wait_op::prepare(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::error_code*, int, boost::corosio::detail::uring_scheduler*, std::shared_ptr<void>, int, std::stop_token const&) :590 50x 100.0% 100.0% boost::corosio::detail::uring_wait_op::do_prep(boost::corosio::detail::uring_op*, io_uring_sqe*) :613 40x 100.0% 100.0% boost::corosio::detail::uring_wait_op::do_cqe(boost::corosio::detail::uring_op*, int, unsigned int, boost::corosio::detail::ready_queue&) :620 40x 100.0% 100.0% boost::corosio::detail::uring_wait_op::do_handler(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int) :628 50x 90.0% 90.0% boost::corosio::detail::uring_local_connect_op::uring_local_connect_op() :686 285x 100.0% 100.0% boost::corosio::detail::uring_local_connect_op::prepare(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::error_code*, int, boost::corosio::detail::uring_scheduler*, std::shared_ptr<void>, boost::corosio::local_endpoint, boost::corosio::local_endpoint*, boost::corosio::local_endpoint*, std::stop_token const&) :694 44x 100.0% 100.0% boost::corosio::detail::uring_local_connect_op::do_prep(boost::corosio::detail::uring_op*, io_uring_sqe*) :721 42x 100.0% 100.0% boost::corosio::detail::uring_local_connect_op::do_cqe(boost::corosio::detail::uring_op*, int, unsigned int, boost::corosio::detail::ready_queue&) :730 41x 100.0% 100.0% boost::corosio::detail::uring_local_connect_op::do_handler(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int) :738 43x 100.0% 100.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_SOCKET_OPS_HPP
12 #define BOOST_COROSIO_NATIVE_DETAIL_URING_URING_SOCKET_OPS_HPP
13
14 #include <boost/corosio/detail/platform.hpp>
15
16 #if BOOST_COROSIO_HAS_URING
17
18 #include <liburing.h>
19
20 #include <boost/capy/buffers.hpp>
21 #include <boost/capy/error.hpp>
22 #include <boost/corosio/detail/buffer_param.hpp>
23 #include <boost/corosio/detail/dispatch_coro.hpp>
24 #include <boost/corosio/local_endpoint.hpp>
25 #include <boost/corosio/native/detail/uring/uring_buffer.hpp>
26 #include <boost/corosio/native/detail/uring/uring_op.hpp>
27 #include <boost/corosio/native/detail/uring/uring_scheduler.hpp>
28 #include <boost/corosio/native/detail/coro_op_complete.hpp>
29 #include <boost/corosio/native/detail/make_err.hpp>
30 #include <boost/corosio/native/detail/speculative_state.hpp>
31
32 #include <system_error>
33
34 #include <errno.h>
35 #include <netinet/in.h>
36 #include <poll.h>
37 #include <sys/socket.h>
38 #include <sys/uio.h>
39
40 namespace boost::corosio::detail {
41
42 /// Maximum scatter/gather segments per read/write/dgram op.
43 ///
44 /// Bounded well below `IOV_MAX` (1024 on Linux) so each op's
45 /// `iovec[uring_max_iov]` lives inside the uring_op object on
46 /// the same allocation as the rest of its state. Plan 4's registered-
47 /// buffer work will revisit; until then 16 covers typical scatter use
48 /// cases (fragmented buffers from buffer_sequence) without bloating
49 /// per-op memory.
50 inline constexpr std::size_t uring_max_iov = 16;
51
52 /** Fill an `iovec` array from a buffer sequence without type-punning.
53
54 `buffer_param::copy_to` writes `capy::mutable_buffer` objects. Those
55 cannot legally alias `iovec` storage: no `mutable_buffer` lives at
56 that address, and `mutable_buffer` is not an implicit-lifetime type,
57 so its lifetime cannot be started in place (which also rules out
58 `std::start_lifetime_as_array`). We copy into a real `mutable_buffer`
59 scratch array, then translate field by field into the caller's
60 `iovec` array. Zero-size buffers are already skipped by `copy_to`.
61
62 @param buffers The source buffer sequence.
63 @param iovecs Destination array, filled with the non-zero buffers.
64 @return The number of `iovec` entries written.
65 */
66 inline int
67 396517x copy_to_iovec(
68 buffer_param const& buffers, iovec (&iovecs)[uring_max_iov]) noexcept
69 {
70 396517x capy::mutable_buffer bufs[uring_max_iov];
71 396517x std::size_t const n = buffers.copy_to(bufs, uring_max_iov);
72 793025x for (std::size_t i = 0; i < n; ++i)
73 {
74 396508x iovecs[i].iov_base = bufs[i].data();
75 396508x iovecs[i].iov_len = bufs[i].size();
76 }
77 396517x return static_cast<int>(n);
78 }
79
80 /** Resolve ec_out/bytes_out from a CQE result for a completed I/O op.
81
82 Shared by read, write, and connect handlers. For reads, `res == 0`
83 with a non-empty buffer means the peer closed the connection (EOF).
84
85 @param self The completed op.
86 @param is_read True if this is a receive/read operation.
87 @param empty_buf True if the submitted buffer was zero-length.
88 */
89 inline void
90 25697x uring_set_result(uring_op* self, bool is_read, bool empty_buf) noexcept
91 {
92 51241x decode_io_result(
93 25697x self->ec_out, self->cancelled.load(std::memory_order_acquire),
94 153x self->res < 0 ? make_err(-self->res) : std::error_code{}, is_read,
95 25697x self->res >= 0 ? static_cast<std::size_t>(self->res) : 0u, empty_buf);
96 25697x }
97
98 /** Scatter-gather read via `IORING_OP_READV`.
99
100 @par Handler dispatch
101 do_cqe captures `res`/`cqe_flags` and queues self into `local`;
102 do_handler runs from the scheduler queue and resumes the coroutine.
103 */
104 struct uring_read_op : uring_op
105 {
106 iovec iovecs[uring_max_iov];
107 int iovec_count = 0;
108 int fd = -1;
109 detail::speculative_state* spec_state = nullptr;
110
111 10001x uring_read_op() noexcept : uring_op(&do_handler, &do_cqe, &do_prep)
112 {
113 10001x is_read = true;
114 10001x }
115
116 /** Reset and initialize for a new submission.
117
118 Embedded ops are reused across calls; every mutable field the
119 handler may read must be re-initialized here. `start(token)`
120 also resets `cancelled`, `sqe_set`, and `stop_cb`.
121
122 @pre This slot has no in-flight op (its prior op completed).
123 */
124 11211x void prepare(
125 std::coroutine_handle<> handle,
126 capy::executor_ref executor,
127 std::error_code* ec,
128 std::size_t* bytes,
129 int file_descriptor,
130 uring_scheduler* scheduler,
131 std::shared_ptr<void> impl,
132 detail::speculative_state* spec,
133 buffer_param buffers,
134 std::stop_token const& token) noexcept
135 {
136 11211x h = handle;
137 11211x ex = executor;
138 11211x ec_out = ec;
139 11211x bytes_out = bytes;
140 11211x fd = file_descriptor;
141 11211x sched_ = scheduler;
142 11211x impl_ptr = std::move(impl);
143 11211x spec_state = spec;
144 11211x res = 0;
145 11211x cqe_flags = 0;
146 11211x iovec_count = copy_to_iovec(buffers, iovecs);
147 11211x empty_buffer = (iovec_count == 0);
148 11211x start(token);
149 11211x }
150
151 245x static void do_prep(uring_op* base, ::io_uring_sqe* sqe) noexcept
152 {
153 245x auto* self = static_cast<uring_read_op*>(base);
154 // Single-buffer fast path: IORING_OP_RECV with a flat
155 // (buffer, length) skips the iovec-array indirection that
156 // IORING_OP_READV pays. For multi-iovec scatter reads, fall
157 // back to readv.
158 245x if (self->iovec_count == 1)
159 {
160 244x ::io_uring_prep_recv(
161 sqe, self->fd, self->iovecs[0].iov_base,
162 self->iovecs[0].iov_len, 0);
163 }
164 else
165 {
166 1x ::io_uring_prep_readv(
167 1x sqe, self->fd, self->iovecs, self->iovec_count, 0);
168 }
169 245x }
170
171 static void
172 243x do_cqe(uring_op* base, int res, unsigned flags, ready_queue& local) noexcept
173 {
174 243x auto* self = static_cast<uring_read_op*>(base);
175 243x self->res = res;
176 243x self->cqe_flags = flags;
177 243x local.push(self);
178 243x }
179
180 11209x static void do_handler(
181 void* owner,
182 scheduler_op* base,
183 std::uint32_t /*bytes*/,
184 std::uint32_t /*error*/) noexcept
185 {
186 11209x auto* self = static_cast<uring_read_op*>(base);
187 11209x if (coro_drain_if_shutdown(owner, self))
188 2x return;
189
190 11207x if (self->sched_)
191 11207x self->sched_->reset_inline_budget();
192
193 11207x uring_set_result(self, true, self->empty_buffer);
194
195 11207x if (self->res > 0 && self->spec_state)
196 {
197 // Kernel signalled readiness — restore speculation.
198 11083x self->spec_state->on_async_read_ready();
199 }
200
201 11207x if (self->bytes_out)
202 11207x *self->bytes_out =
203 11207x self->res >= 0 ? static_cast<std::size_t>(self->res) : 0u;
204
205 11207x coro_resume(self);
206 // suicide drops here; may destroy impl + self.
207 }
208 };
209
210 /** Scatter-gather write via `IORING_OP_SENDMSG` with `MSG_NOSIGNAL`.
211
212 `MSG_NOSIGNAL` prevents `SIGPIPE` when the peer has closed the
213 connection; the error is surfaced as `EPIPE` instead.
214 */
215 struct uring_write_op : uring_op
216 {
217 iovec iovecs[uring_max_iov];
218 int iovec_count = 0;
219 int fd = -1;
220 msghdr msg{};
221 detail::speculative_state* spec_state = nullptr;
222
223 10001x uring_write_op() noexcept : uring_op(&do_handler, &do_cqe, &do_prep) {}
224
225 /** Reset and initialize for a new submission. See uring_read_op::prepare. */
226 10973x void prepare(
227 std::coroutine_handle<> handle,
228 capy::executor_ref executor,
229 std::error_code* ec,
230 std::size_t* bytes,
231 int file_descriptor,
232 uring_scheduler* scheduler,
233 std::shared_ptr<void> impl,
234 detail::speculative_state* spec,
235 buffer_param buffers,
236 std::stop_token const& token) noexcept
237 {
238 10973x h = handle;
239 10973x ex = executor;
240 10973x ec_out = ec;
241 10973x bytes_out = bytes;
242 10973x fd = file_descriptor;
243 10973x sched_ = scheduler;
244 10973x impl_ptr = std::move(impl);
245 10973x spec_state = spec;
246 10973x res = 0;
247 10973x cqe_flags = 0;
248 10973x iovec_count = copy_to_iovec(buffers, iovecs);
249 10973x empty_buffer = (iovec_count == 0);
250 10973x if (!empty_buffer)
251 {
252 10973x msg = {};
253 10973x msg.msg_iov = iovecs;
254 10973x msg.msg_iovlen = static_cast<decltype(msg.msg_iovlen)>(iovec_count);
255 }
256 10973x start(token);
257 10973x }
258
259 5x static void do_prep(uring_op* base, ::io_uring_sqe* sqe) noexcept
260 {
261 5x auto* self = static_cast<uring_write_op*>(base);
262 // Single-buffer fast path: IORING_OP_SEND with MSG_NOSIGNAL
263 // skips the msghdr indirection that IORING_OP_SENDMSG pays.
264 // For multi-iovec scatter writes, fall back to sendmsg.
265 5x if (self->iovec_count == 1)
266 {
267 4x ::io_uring_prep_send(
268 4x sqe, self->fd, self->iovecs[0].iov_base,
269 self->iovecs[0].iov_len, MSG_NOSIGNAL);
270 }
271 else
272 {
273 1x ::io_uring_prep_sendmsg(sqe, self->fd, &self->msg, MSG_NOSIGNAL);
274 }
275 5x }
276
277 static void
278 4x do_cqe(uring_op* base, int res, unsigned flags, ready_queue& local) noexcept
279 {
280 4x auto* self = static_cast<uring_write_op*>(base);
281 4x self->res = res;
282 4x self->cqe_flags = flags;
283 4x local.push(self);
284 4x }
285
286 10972x static void do_handler(
287 void* owner,
288 scheduler_op* base,
289 std::uint32_t /*bytes*/,
290 std::uint32_t /*error*/) noexcept
291 {
292 10972x auto* self = static_cast<uring_write_op*>(base);
293 10972x if (coro_drain_if_shutdown(owner, self))
294 1x return;
295
296 10971x if (self->sched_)
297 10971x self->sched_->reset_inline_budget();
298
299 10971x uring_set_result(self, false, self->empty_buffer);
300
301 10971x if (self->res > 0 && self->spec_state)
302 {
303 // Kernel signalled readiness — restore speculation.
304 10966x self->spec_state->on_async_write_ready();
305 }
306
307 10971x if (self->bytes_out)
308 10971x *self->bytes_out =
309 10971x self->res >= 0 ? static_cast<std::size_t>(self->res) : 0u;
310
311 10971x coro_resume(self);
312 }
313 };
314
315 /** Non-blocking connect via `IORING_OP_CONNECT`.
316
317 Negative `res` is the connect error; zero means success.
318 `remote_endpoint_out` is written only on success so a failed
319 connect does not corrupt the socket's cached remote endpoint.
320 */
321 struct uring_connect_op : uring_op
322 {
323 sockaddr_storage addr{};
324 socklen_t addrlen = 0;
325 int fd = -1;
326 endpoint target_endpoint{};
327 endpoint* remote_endpoint_out = nullptr;
328 endpoint* local_endpoint_out = nullptr;
329
330 9974x uring_connect_op() noexcept : uring_op(&do_handler, &do_cqe, &do_prep) {}
331
332 /** Reset and initialize for a new submission.
333
334 The caller must fill `addr` and `addrlen` before calling this
335 (typically via `to_sockaddr(ep, family, conn_.addr)` which
336 returns the addrlen) — `to_sockaddr` is the family-aware
337 helper and requires the socket family which is known to the
338 caller, not the op.
339 */
340 3286x void prepare(
341 std::coroutine_handle<> handle,
342 capy::executor_ref executor,
343 std::error_code* ec,
344 int file_descriptor,
345 uring_scheduler* scheduler,
346 std::shared_ptr<void> impl,
347 endpoint target,
348 endpoint* remote_out,
349 endpoint* local_out,
350 std::stop_token const& token) noexcept
351 {
352 3286x h = handle;
353 3286x ex = executor;
354 3286x ec_out = ec;
355 3286x bytes_out = nullptr;
356 3286x fd = file_descriptor;
357 3286x sched_ = scheduler;
358 3286x impl_ptr = std::move(impl);
359 3286x res = 0;
360 3286x cqe_flags = 0;
361 3286x target_endpoint = target;
362 3286x remote_endpoint_out = remote_out;
363 3286x local_endpoint_out = local_out;
364 // addr / addrlen are pre-filled by the caller.
365 3286x start(token);
366 3286x }
367
368 3284x static void do_prep(uring_op* base, ::io_uring_sqe* sqe) noexcept
369 {
370 3284x auto* self = static_cast<uring_connect_op*>(base);
371 3284x ::io_uring_prep_connect(
372 3284x sqe, self->fd, reinterpret_cast<sockaddr const*>(&self->addr),
373 self->addrlen);
374 3284x }
375
376 static void
377 3283x do_cqe(uring_op* base, int res, unsigned flags, ready_queue& local) noexcept
378 {
379 3283x auto* self = static_cast<uring_connect_op*>(base);
380 3283x self->res = res;
381 3283x self->cqe_flags = flags;
382 3283x local.push(self);
383 3283x }
384
385 3285x static void do_handler(
386 void* owner,
387 scheduler_op* base,
388 std::uint32_t /*bytes*/,
389 std::uint32_t /*error*/) noexcept
390 {
391 3285x auto* self = static_cast<uring_connect_op*>(base);
392 3285x if (coro_drain_if_shutdown(owner, self))
393 1x return;
394
395 3284x if (self->sched_)
396 3284x self->sched_->reset_inline_budget();
397
398 3284x uring_set_result(self, false, false);
399
400 // Write endpoints only on success.
401 3284x if (self->res >= 0)
402 {
403 3273x if (self->remote_endpoint_out)
404 3273x *self->remote_endpoint_out = self->target_endpoint;
405 3273x if (self->local_endpoint_out && self->fd >= 0)
406 {
407 3273x sockaddr_storage local{};
408 3273x socklen_t len = sizeof(local);
409 3273x if (::getsockname(
410 3273x self->fd, reinterpret_cast<sockaddr*>(&local), &len) ==
411 0)
412 3273x *self->local_endpoint_out = sockaddr_to_endpoint(local);
413 }
414 }
415
416 3284x coro_resume(self);
417 }
418 };
419
420 /** Submit an `uring_op`, reporting an SQ that stayed full.
421
422 The body behind @ref uring_submit_op and
423 @ref uring_try_submit_op; `counted` selects which of the two
424 answers an exhausted SQ ring gets.
425
426 @pre `op->prep_func != nullptr`.
427
428 @par Exception Safety
429 Nothrow.
430
431 @param sched The scheduler owning the ring.
432 @param op The operation to submit.
433 @param counted True when a `work_started()` backs this op, so its
434 completion may ride the scheduler's queue.
435
436 @return `false` when the SQ stayed full and the op was left for the
437 caller to report; `true` otherwise.
438 */
439 inline bool
440 4111x uring_do_submit_op(uring_scheduler& sched, uring_op* op, bool counted) noexcept
441 {
442 4111x sched.lazy_init_ring();
443
444 4111x bool need_post = false;
445 {
446 4111x typename uring_scheduler::lock_type ring_lock(sched.ring_mutex());
447
448 4111x ::io_uring_sqe* sqe = ::io_uring_get_sqe(sched.ring());
449 4111x if (!sqe)
450 {
451 // SQ ring full — flush to kernel and retry once.
452 8x ::io_uring_submit(sched.ring());
453 8x sqe = ::io_uring_get_sqe(sched.ring());
454 }
455
456 4111x if (!sqe)
457 {
458 // SQ stayed full after one flush — synchronous failure path.
459 // Report EAGAIN the way a CQE would, because every handler
460 // re-derives its result from `res`: a code written straight
461 // to ec_out is overwritten on the way out, and a res left at
462 // zero reads as end-of-file, a zero-byte write, or a
463 // successful connect that never happened. Queue the op as
464 // completed so do_one dispatches the handler; the caller's
465 // work_started() pays for the work_finished() do_one spends
466 // on it. (CAS path is not entered here.)
467 8x op->res = -EAGAIN;
468 8x if (!counted)
469 {
470 // Nothing counted this op, so queueing it would spend a
471 // work_finished() the context never owed and drive
472 // outstanding_work_ below what is really outstanding.
473 // The owner reports the failure instead.
474 5x return false;
475 }
476 3x typename uring_scheduler::lock_type lock(sched.dispatch_mutex());
477 3x sched.push_completed_locked(op);
478 3x return true;
479 3x }
480
481 4103x op->prep_func(op, sqe);
482 4103x ::io_uring_sqe_set_data(sqe, op);
483 // Count this op against the in-flight gate in do_one: it
484 // expects exactly one F_MORE-less CQE per submitted SQE
485 // (multishot ops decrement only on the terminal CQE).
486 4103x sched.inflight_inc();
487 // Release pairs with the acquire in uring_op::request_cancel:
488 // a stop_token firing after we release the mutex will see
489 // sqe_set==true and submit a cancel-by-user_data SQE.
490 4103x op->sqe_set.store(true, std::memory_order_release);
491
492 // First submitter in a batch wins the CAS and will post
493 // submit_sqes_op; others piggyback on the same flush.
494 4103x if (!sched.submit_op_posted_exchange(true))
495 314x need_post = true;
496 4111x }
497
498 4103x if (need_post)
499 {
500 // Flush is deferred to submit_sqes_op; post() owns the wake.
501 314x sched.post(&sched.submit_op_ref());
502 }
503 4103x return true;
504 }
505
506 /** Submit an `uring_op` a `work_started()` already paid for.
507
508 Acquires the ring mutex, prepares the SQE, and (under the same
509 mutex) CAS-sets `submit_op_posted_`. The first submitter of a
510 batch wins the CAS and posts the scheduler's `submit_sqes_op`,
511 which later flushes all queued SQEs in a single
512 `io_uring_submit_and_get_events` call and drains any ready CQEs.
513 Subsequent submitters in the same batch piggyback — their SQEs
514 sit in the user-space SQ ring until that op dispatches.
515
516 On SQ-ring exhaustion (after one flush retry), completes the op
517 with `EAGAIN` and queues it so its handler dispatches on the next
518 `do_one` cycle, exactly as if the kernel had returned that error.
519 That is why this spelling reports nothing: the failure reaches the
520 caller as the operation's own `EAGAIN` completion.
521
522 @pre `op->prep_func != nullptr`.
523 @pre A `work_started()` backs this op, so the `work_finished()` the
524 scheduler spends on everything it dispatches is owed.
525
526 @par Exception Safety
527 Nothrow.
528
529 @param sched The scheduler owning the ring.
530 @param op The operation to submit.
531
532 @see uring_try_submit_op
533 */
534 inline void
535 3889x uring_submit_op(uring_scheduler& sched, uring_op* op) noexcept
536 {
537 // The result is true by construction: a counted op's SQ-full path
538 // queues the op and answers through its own completion.
539 3889x uring_do_submit_op(sched, op, true);
540 3889x }
541
542 /** Submit an `uring_op` nothing counted, reporting a full SQ.
543
544 The scheduler spends a `work_finished()` on everything it
545 dispatches out of `completed_ops_`, so an op no `work_started()`
546 backs cannot be completed through that queue: doing so drives
547 `outstanding_work_` below what is really outstanding. This spelling
548 hands an SQ that stayed full back to the owner instead, which
549 reports it through its own channel — the multishot accept arm
550 latches it and answers the accepts it can no longer serve.
551
552 @pre `op->prep_func != nullptr`.
553
554 @par Exception Safety
555 Nothrow.
556
557 @param sched The scheduler owning the ring.
558 @param op The operation to submit.
559
560 @return True when the SQE was prepared; false when the SQ stayed
561 full after one flush and the caller owns the failure.
562
563 @see uring_submit_op
564 */
565 [[nodiscard]] inline bool
566 222x uring_try_submit_op(uring_scheduler& sched, uring_op* op) noexcept
567 {
568 222x return uring_do_submit_op(sched, op, false);
569 }
570
571 /** Readiness wait via `IORING_OP_POLL_ADD`.
572
573 Used to implement the `wait()` virtual for socket and acceptor
574 implementations. The op submits a one-shot poll on `fd` for the
575 requested set of poll flags (POLLIN / POLLOUT / POLLPRI|POLLERR|
576 POLLHUP) and reports completion without transferring any data.
577
578 The CQE's `res` carries the actual revents, but we surface only
579 success/cancel/error on `*ec_out` — callers of `wait()` just need
580 a readiness signal, not the specific event mask.
581 */
582 struct uring_wait_op : uring_op
583 {
584 int fd = -1;
585 int poll_flags = 0;
586
587 10531x uring_wait_op() noexcept : uring_op(&do_handler, &do_cqe, &do_prep) {}
588
589 /** Reset and initialize for a new submission. */
590 50x void prepare(
591 std::coroutine_handle<> handle,
592 capy::executor_ref executor,
593 std::error_code* ec,
594 int file_descriptor,
595 uring_scheduler* scheduler,
596 std::shared_ptr<void> impl,
597 int flags,
598 std::stop_token const& token) noexcept
599 {
600 50x h = handle;
601 50x ex = executor;
602 50x ec_out = ec;
603 50x bytes_out = nullptr;
604 50x fd = file_descriptor;
605 50x sched_ = scheduler;
606 50x impl_ptr = std::move(impl);
607 50x poll_flags = flags;
608 50x res = 0;
609 50x cqe_flags = 0;
610 50x start(token);
611 50x }
612
613 40x static void do_prep(uring_op* base, ::io_uring_sqe* sqe) noexcept
614 {
615 40x auto* self = static_cast<uring_wait_op*>(base);
616 40x ::io_uring_prep_poll_add(sqe, self->fd, self->poll_flags);
617 40x }
618
619 static void
620 40x do_cqe(uring_op* base, int res, unsigned flags, ready_queue& local) noexcept
621 {
622 40x auto* self = static_cast<uring_wait_op*>(base);
623 40x self->res = res;
624 40x self->cqe_flags = flags;
625 40x local.push(self);
626 40x }
627
628 50x static void do_handler(
629 void* owner,
630 scheduler_op* base,
631 std::uint32_t /*bytes*/,
632 std::uint32_t /*error*/) noexcept
633 {
634 50x auto* self = static_cast<uring_wait_op*>(base);
635 50x if (coro_drain_if_shutdown(owner, self))
636 3x return;
637
638 47x if (self->sched_)
639 47x self->sched_->reset_inline_budget();
640
641 // A POLL_ADD completion carries the error band in its revents
642 // (res), not as a negative res, so name the reason the reactor
643 // way — SO_ERROR, or EIO when the kernel has none — instead of
644 // completing wait(error) with an empty, benign-looking code.
645 // OOB (POLLPRI) is a readiness signal, not an error.
646 47x std::error_code ec{};
647 47x if (self->res < 0)
648 {
649 20x ec = make_err(-self->res);
650 }
651 27x else if (self->res & (POLLERR | POLLHUP | POLLNVAL))
652 {
653 2x int so_err = 0;
654 2x socklen_t len = sizeof(so_err);
655 2x if (::getsockopt(self->fd, SOL_SOCKET, SO_ERROR, &so_err, &len) < 0)
656 so_err = errno;
657 2x if (so_err == 0)
658 so_err = EIO;
659 2x ec = make_err(so_err);
660 }
661
662 // Wait reports only success/cancel/error — no bytes, no EOF.
663 47x decode_io_result(
664 47x self->ec_out, self->cancelled.load(std::memory_order_acquire), ec,
665 /*is_read=*/false, /*bytes=*/0, /*empty_buffer=*/false);
666
667 47x coro_resume(self);
668 }
669 };
670
671 /** Non-blocking connect for Unix domain sockets via `IORING_OP_CONNECT`.
672
673 Like `uring_connect_op` but stores `local_endpoint` for the target
674 and out-pointers, since `sockaddr_to_local_endpoint` returns
675 `local_endpoint`, not `endpoint`.
676 */
677 struct uring_local_connect_op : uring_op
678 {
679 sockaddr_storage addr{};
680 socklen_t addrlen = 0;
681 int fd = -1;
682 corosio::local_endpoint target_endpoint{};
683 corosio::local_endpoint* remote_endpoint_out = nullptr;
684 corosio::local_endpoint* local_endpoint_out = nullptr;
685
686 285x uring_local_connect_op() noexcept : uring_op(&do_handler, &do_cqe, &do_prep)
687 {
688 285x }
689
690 /** Reset and initialize for a new submission.
691
692 Caller pre-fills `addr` and `addrlen` (see uring_connect_op::prepare).
693 */
694 44x void prepare(
695 std::coroutine_handle<> handle,
696 capy::executor_ref executor,
697 std::error_code* ec,
698 int file_descriptor,
699 uring_scheduler* scheduler,
700 std::shared_ptr<void> impl,
701 corosio::local_endpoint target,
702 corosio::local_endpoint* remote_out,
703 corosio::local_endpoint* local_out,
704 std::stop_token const& token) noexcept
705 {
706 44x h = handle;
707 44x ex = executor;
708 44x ec_out = ec;
709 44x bytes_out = nullptr;
710 44x fd = file_descriptor;
711 44x sched_ = scheduler;
712 44x impl_ptr = std::move(impl);
713 44x res = 0;
714 44x cqe_flags = 0;
715 44x target_endpoint = target;
716 44x remote_endpoint_out = remote_out;
717 44x local_endpoint_out = local_out;
718 44x start(token);
719 44x }
720
721 42x static void do_prep(uring_op* base, ::io_uring_sqe* sqe) noexcept
722 {
723 42x auto* self = static_cast<uring_local_connect_op*>(base);
724 42x ::io_uring_prep_connect(
725 42x sqe, self->fd, reinterpret_cast<sockaddr const*>(&self->addr),
726 self->addrlen);
727 42x }
728
729 static void
730 41x do_cqe(uring_op* base, int res, unsigned flags, ready_queue& local) noexcept
731 {
732 41x auto* self = static_cast<uring_local_connect_op*>(base);
733 41x self->res = res;
734 41x self->cqe_flags = flags;
735 41x local.push(self);
736 41x }
737
738 43x static void do_handler(
739 void* owner,
740 scheduler_op* base,
741 std::uint32_t /*bytes*/,
742 std::uint32_t /*error*/) noexcept
743 {
744 43x auto* self = static_cast<uring_local_connect_op*>(base);
745 43x if (coro_drain_if_shutdown(owner, self))
746 1x return;
747
748 42x if (self->sched_)
749 42x self->sched_->reset_inline_budget();
750
751 42x uring_set_result(self, false, false);
752
753 // Write endpoints only on success.
754 42x if (self->res >= 0)
755 {
756 39x if (self->remote_endpoint_out)
757 39x *self->remote_endpoint_out = self->target_endpoint;
758 39x if (self->local_endpoint_out && self->fd >= 0)
759 {
760 39x sockaddr_storage local{};
761 39x socklen_t len = sizeof(local);
762 39x if (::getsockname(
763 39x self->fd, reinterpret_cast<sockaddr*>(&local), &len) ==
764 0)
765 39x *self->local_endpoint_out =
766 39x sockaddr_to_local_endpoint(local, len);
767 }
768 }
769
770 42x coro_resume(self);
771 }
772 };
773
774 } // namespace boost::corosio::detail
775
776 #endif // BOOST_COROSIO_HAS_URING
777
778 #endif // BOOST_COROSIO_NATIVE_DETAIL_URING_URING_SOCKET_OPS_HPP
779