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

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