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

98.4% Lines (123 / 125) 100.0% Functions (18 / 18)
uring_file_ops.hpp
f(x) Functions (18)
Function Calls Lines Blocks
boost::corosio::detail::uring_file_read_op_base::uring_file_read_op_base(void (*)(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int)) :53 355x 100.0% 100.0% boost::corosio::detail::uring_file_read_op_base::prepare(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::error_code*, unsigned long*, int, long, boost::corosio::detail::uring_scheduler*, std::shared_ptr<void>, boost::corosio::buffer_param, std::stop_token const&) :66 234x 100.0% 100.0% boost::corosio::detail::uring_file_read_op_base::do_prep(boost::corosio::detail::uring_op*, io_uring_sqe*) :93 226x 100.0% 100.0% boost::corosio::detail::uring_file_read_op_base::do_cqe(boost::corosio::detail::uring_op*, int, unsigned int, boost::corosio::detail::ready_queue&) :102 226x 100.0% 100.0% boost::corosio::detail::uring_file_read_op_base::finish(boost::corosio::detail::uring_file_read_op_base*) :111 167x 100.0% 100.0% boost::corosio::detail::uring_file_read_op::uring_file_read_op() :125 103x 100.0% 100.0% boost::corosio::detail::uring_file_read_op::do_handler(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int) :127 38x 90.9% 90.0% boost::corosio::detail::uring_random_access_read_op::uring_random_access_read_op() :157 169x 100.0% 100.0% boost::corosio::detail::uring_random_access_read_op::do_handler(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int) :162 169x 100.0% 100.0% boost::corosio::detail::uring_file_write_op_base::uring_file_write_op_base(void (*)(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int)) :204 222x 100.0% 100.0% boost::corosio::detail::uring_file_write_op_base::prepare(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, std::error_code*, unsigned long*, int, long, boost::corosio::detail::uring_scheduler*, std::shared_ptr<void>, boost::corosio::buffer_param, std::stop_token const&) :214 96x 100.0% 100.0% boost::corosio::detail::uring_file_write_op_base::do_prep(boost::corosio::detail::uring_op*, io_uring_sqe*) :241 90x 100.0% 100.0% boost::corosio::detail::uring_file_write_op_base::do_cqe(boost::corosio::detail::uring_op*, int, unsigned int, boost::corosio::detail::ready_queue&) :250 87x 100.0% 100.0% boost::corosio::detail::uring_file_write_op_base::finish(boost::corosio::detail::uring_file_write_op_base*) :259 43x 100.0% 100.0% boost::corosio::detail::uring_file_write_op::uring_file_write_op() :271 103x 100.0% 100.0% boost::corosio::detail::uring_file_write_op::do_handler(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int) :273 40x 90.9% 90.0% boost::corosio::detail::uring_random_access_write_op::uring_random_access_write_op() :301 45x 100.0% 100.0% boost::corosio::detail::uring_random_access_write_op::do_handler(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int) :306 44x 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_FILE_OPS_HPP
12 #define BOOST_COROSIO_NATIVE_DETAIL_URING_URING_FILE_OPS_HPP
13
14 #include <boost/corosio/detail/platform.hpp>
15
16 #if BOOST_COROSIO_HAS_URING
17
18 #include <boost/corosio/native/detail/uring/uring_op.hpp>
19 #include <boost/corosio/native/detail/uring/uring_socket_ops.hpp>
20 #include <boost/corosio/native/detail/coro_op_complete.hpp>
21 #include <boost/corosio/detail/dispatch_coro.hpp>
22
23 #include <cstdint>
24 #include <sys/uio.h>
25
26 namespace boost::corosio::detail {
27
28 /** Scatter-gather file read via `IORING_OP_READV`.
29
30 Stream files pass `offset == -1` so the kernel uses (and updates)
31 the fd's `f_pos`, matching POSIX `read(2)` semantics. Random-
32 access files pass an explicit caller-supplied offset.
33
34 @par Handler dispatch
35 `do_cqe` captures `res`/`cqe_flags` and queues self into `local`;
36 `do_handler` runs from the scheduler queue and resumes the
37 coroutine.
38 */
39 /// Shared state and submission logic for file read ops. Concrete
40 /// subclasses pick a `do_handler` that matches their storage model:
41 /// `uring_file_read_op` for embedded slots (stream_file), and
42 /// `uring_random_access_read_op` for heap-allocated per-call ops
43 /// (random_access_file, where concurrent reads at different offsets
44 /// are legitimate).
45 struct uring_file_read_op_base : uring_op
46 {
47 iovec iovecs[uring_max_iov];
48 int iovec_count = 0;
49 int fd = -1;
50 std::int64_t offset = -1; // -1 means kernel f_pos
51
52 protected:
53 355x explicit uring_file_read_op_base(func_type handler) noexcept
54 355x : uring_op(handler, &do_cqe, &do_prep)
55 {
56 355x is_read = true;
57 355x }
58
59 public:
60 /** Reset and initialize for a new submission.
61
62 @param file_offset -1 selects the kernel's `f_pos` (POSIX
63 `read(2)` semantics for stream files); otherwise the explicit
64 offset for random-access files.
65 */
66 234x void prepare(
67 std::coroutine_handle<> handle,
68 capy::executor_ref executor,
69 std::error_code* ec,
70 std::size_t* bytes,
71 int file_descriptor,
72 std::int64_t file_offset,
73 uring_scheduler* scheduler,
74 std::shared_ptr<void> impl,
75 buffer_param buffers,
76 std::stop_token const& token) noexcept
77 {
78 234x h = handle;
79 234x ex = executor;
80 234x ec_out = ec;
81 234x bytes_out = bytes;
82 234x fd = file_descriptor;
83 234x offset = file_offset;
84 234x sched_ = scheduler;
85 234x impl_ptr = std::move(impl);
86 234x res = 0;
87 234x cqe_flags = 0;
88 234x iovec_count = copy_to_iovec(buffers, iovecs);
89 234x empty_buffer = (iovec_count == 0);
90 234x start(token);
91 234x }
92
93 226x static void do_prep(uring_op* base, ::io_uring_sqe* sqe) noexcept
94 {
95 226x auto* self = static_cast<uring_file_read_op_base*>(base);
96 226x ::io_uring_prep_readv(
97 226x sqe, self->fd, self->iovecs, self->iovec_count,
98 226x static_cast<__u64>(self->offset));
99 226x }
100
101 static void
102 226x do_cqe(uring_op* base, int res, unsigned flags, ready_queue& local) noexcept
103 {
104 226x auto* self = static_cast<uring_file_read_op_base*>(base);
105 226x self->res = res;
106 226x self->cqe_flags = flags;
107 226x local.push(self);
108 226x }
109
110 /// Fill ec_out and bytes_out from the completion.
111 167x static void finish(uring_file_read_op_base* self) noexcept
112 {
113 167x uring_set_result(self, /*is_read=*/true, self->empty_buffer);
114 167x if (self->bytes_out)
115 167x *self->bytes_out =
116 167x self->res >= 0 ? static_cast<std::size_t>(self->res) : 0u;
117 167x }
118 };
119
120 /// Scatter-gather file read embedded as a member of stream_file
121 /// (single-pending per fd). Handler uses the suicide-move pattern;
122 /// the impl owns this slot.
123 struct uring_file_read_op : uring_file_read_op_base
124 {
125 103x uring_file_read_op() noexcept : uring_file_read_op_base(&do_handler) {}
126
127 38x static void do_handler(
128 void* owner,
129 scheduler_op* base,
130 std::uint32_t /*bytes*/,
131 std::uint32_t /*error*/) noexcept
132 {
133 38x auto* self = static_cast<uring_file_read_op*>(base);
134 38x if (coro_drain_if_shutdown(owner, self))
135 ✗ return;
136
137 38x if (self->sched_)
138 38x self->sched_->reset_inline_budget();
139
140 38x uring_set_result(self, /*is_read=*/true, self->empty_buffer);
141 38x if (self->bytes_out)
142 38x *self->bytes_out =
143 38x self->res >= 0 ? static_cast<std::size_t>(self->res) : 0u;
144 38x coro_resume(self);
145 }
146 };
147
148 /// Heap-allocated scatter-gather file read for random_access_file —
149 /// each `read_some_at` call allocates a fresh op so multiple reads
150 /// at different offsets on the same fd can be in flight concurrently.
151 struct uring_random_access_read_op : uring_file_read_op_base
152 {
153 /// The awaitable's continuation, not the embedded `cont`: this op is
154 /// freed before the coroutine resumes, possibly on another thread.
155 capy::continuation* awaiting = nullptr;
156
157 169x uring_random_access_read_op() noexcept
158 169x : uring_file_read_op_base(&do_handler)
159 {
160 169x }
161
162 169x static void do_handler(
163 void* owner,
164 scheduler_op* base,
165 std::uint32_t /*bytes*/,
166 std::uint32_t /*error*/) noexcept
167 {
168 169x auto* self = static_cast<uring_random_access_read_op*>(base);
169 169x self->stop_cb.reset();
170
171 169x if (owner == nullptr)
172 {
173 2x delete self;
174 2x return;
175 }
176
177 167x finish(self);
178 167x auto* c = self->awaiting;
179 167x auto ex = self->ex;
180 167x delete self;
181 167x dispatch_coro(ex, *c).resume();
182 }
183 };
184
185 /** Scatter-gather file write via `IORING_OP_WRITEV`.
186
187 Stream files pass `offset == -1` (kernel f_pos); random-access
188 files pass an explicit caller-supplied offset. Unlike socket
189 writes there is no `MSG_NOSIGNAL` equivalent: a write to a pipe
190 or FIFO whose reader has closed raises SIGPIPE when the kernel
191 executes it, so the teardown paths that flush queued writes hold
192 the signal blocked (see `scoped_sigpipe_block`).
193 */
194 /// Shared state and submission logic for file write ops. Concrete
195 /// subclasses pick a `do_handler` matching their storage model.
196 struct uring_file_write_op_base : uring_op
197 {
198 iovec iovecs[uring_max_iov];
199 int iovec_count = 0;
200 int fd = -1;
201 std::int64_t offset = -1;
202
203 protected:
204 222x explicit uring_file_write_op_base(func_type handler) noexcept
205 222x : uring_op(handler, &do_cqe, &do_prep)
206 {
207 222x }
208
209 public:
210 /** Reset and initialize for a new submission.
211
212 See uring_file_read_op_base::prepare for the offset convention.
213 */
214 96x void prepare(
215 std::coroutine_handle<> handle,
216 capy::executor_ref executor,
217 std::error_code* ec,
218 std::size_t* bytes,
219 int file_descriptor,
220 std::int64_t file_offset,
221 uring_scheduler* scheduler,
222 std::shared_ptr<void> impl,
223 buffer_param buffers,
224 std::stop_token const& token) noexcept
225 {
226 96x h = handle;
227 96x ex = executor;
228 96x ec_out = ec;
229 96x bytes_out = bytes;
230 96x fd = file_descriptor;
231 96x offset = file_offset;
232 96x sched_ = scheduler;
233 96x impl_ptr = std::move(impl);
234 96x res = 0;
235 96x cqe_flags = 0;
236 96x iovec_count = copy_to_iovec(buffers, iovecs);
237 96x empty_buffer = (iovec_count == 0);
238 96x start(token);
239 96x }
240
241 90x static void do_prep(uring_op* base, ::io_uring_sqe* sqe) noexcept
242 {
243 90x auto* self = static_cast<uring_file_write_op_base*>(base);
244 90x ::io_uring_prep_writev(
245 90x sqe, self->fd, self->iovecs, self->iovec_count,
246 90x static_cast<__u64>(self->offset));
247 90x }
248
249 static void
250 87x do_cqe(uring_op* base, int res, unsigned flags, ready_queue& local) noexcept
251 {
252 87x auto* self = static_cast<uring_file_write_op_base*>(base);
253 87x self->res = res;
254 87x self->cqe_flags = flags;
255 87x local.push(self);
256 87x }
257
258 /// Fill ec_out and bytes_out from the completion.
259 43x static void finish(uring_file_write_op_base* self) noexcept
260 {
261 43x uring_set_result(self, /*is_read=*/false, self->empty_buffer);
262 43x if (self->bytes_out)
263 43x *self->bytes_out =
264 43x self->res >= 0 ? static_cast<std::size_t>(self->res) : 0u;
265 43x }
266 };
267
268 /// Embedded file write op for stream_file.
269 struct uring_file_write_op : uring_file_write_op_base
270 {
271 103x uring_file_write_op() noexcept : uring_file_write_op_base(&do_handler) {}
272
273 40x static void do_handler(
274 void* owner,
275 scheduler_op* base,
276 std::uint32_t /*bytes*/,
277 std::uint32_t /*error*/) noexcept
278 {
279 40x auto* self = static_cast<uring_file_write_op*>(base);
280 40x if (coro_drain_if_shutdown(owner, self))
281 ✗ return;
282
283 40x if (self->sched_)
284 40x self->sched_->reset_inline_budget();
285
286 40x uring_set_result(self, /*is_read=*/false, self->empty_buffer);
287 40x if (self->bytes_out)
288 40x *self->bytes_out =
289 40x self->res >= 0 ? static_cast<std::size_t>(self->res) : 0u;
290 40x coro_resume(self);
291 }
292 };
293
294 /// Heap-allocated file write op for random_access_file.
295 struct uring_random_access_write_op : uring_file_write_op_base
296 {
297 /// The awaitable's continuation, not the embedded `cont`: this op is
298 /// freed before the coroutine resumes, possibly on another thread.
299 capy::continuation* awaiting = nullptr;
300
301 45x uring_random_access_write_op() noexcept
302 45x : uring_file_write_op_base(&do_handler)
303 {
304 45x }
305
306 44x static void do_handler(
307 void* owner,
308 scheduler_op* base,
309 std::uint32_t /*bytes*/,
310 std::uint32_t /*error*/) noexcept
311 {
312 44x auto* self = static_cast<uring_random_access_write_op*>(base);
313 44x self->stop_cb.reset();
314
315 44x if (owner == nullptr)
316 {
317 1x delete self;
318 1x return;
319 }
320
321 43x finish(self);
322 43x auto* c = self->awaiting;
323 43x auto ex = self->ex;
324 43x delete self;
325 43x dispatch_coro(ex, *c).resume();
326 }
327 };
328
329 } // namespace boost::corosio::detail
330
331 #endif // BOOST_COROSIO_HAS_URING
332
333 #endif // BOOST_COROSIO_NATIVE_DETAIL_URING_URING_FILE_OPS_HPP
334