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

100.0% Lines (125 / 125) 100.0% Functions (16 / 16)
uring_random_access_file.hpp
f(x) Functions (16)
Function Calls Lines Blocks
boost::corosio::detail::uring_random_access_file::uring_random_access_file(boost::corosio::detail::uring_scheduler&) :74 111x 100.0% 100.0% boost::corosio::detail::uring_random_access_file::~uring_random_access_file() :79 110x 100.0% 100.0% boost::corosio::detail::uring_random_access_file::native_handle() const :104 542x 100.0% 100.0% boost::corosio::detail::uring_random_access_file::cancel() :109 2x 100.0% 100.0% boost::corosio::detail::uring_random_access_file::size() const :115 5x 100.0% 100.0% boost::corosio::detail::uring_random_access_file::resize(unsigned long) :123 6x 100.0% 100.0% boost::corosio::detail::uring_random_access_file::sync_data() :133 3x 100.0% 100.0% boost::corosio::detail::uring_random_access_file::sync_all() :144 3x 100.0% 100.0% boost::corosio::detail::uring_random_access_file::release() :151 1x 100.0% 100.0% boost::corosio::detail::uring_random_access_file::assign(int) :163 3x 100.0% 100.0% boost::corosio::detail::uring_random_access_file::open_file(std::filesystem::__cxx11::path const&, boost::corosio::file_base::flags) :177 102x 100.0% 100.0% boost::corosio::detail::uring_random_access_file::close_file() :217 522x 100.0% 100.0% boost::corosio::detail::uring_random_access_file::read_some_at(unsigned long, boost::capy::continuation&, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :233 169x 100.0% 71.0% boost::corosio::detail::uring_random_access_file::write_some_at(unsigned long, boost::capy::continuation&, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :272 45x 100.0% 71.0% boost::corosio::detail::uring_random_access_file_service::uring_random_access_file_service(boost::capy::execution_context&) :329 82x 100.0% 100.0% boost::corosio::detail::uring_random_access_file_service::open_file(boost::corosio::random_access_file::implementation&, std::filesystem::__cxx11::path const&, boost::corosio::file_base::flags) :338 102x 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_RANDOM_ACCESS_FILE_HPP
12 #define BOOST_COROSIO_NATIVE_DETAIL_URING_URING_RANDOM_ACCESS_FILE_HPP
13
14 #include <boost/corosio/detail/platform.hpp>
15
16 #if BOOST_COROSIO_HAS_URING
17
18 #include <boost/corosio/detail/random_access_file_service.hpp>
19 #include <boost/corosio/detail/intrusive.hpp>
20 #include <boost/corosio/native/detail/uring/uring_file_ops.hpp>
21 #include <boost/corosio/native/detail/uring/uring_file_service_base.hpp>
22 #include <boost/corosio/native/detail/uring/uring_scheduler.hpp>
23 #include <boost/corosio/native/detail/make_err.hpp>
24 #include <boost/corosio/native/detail/posix/large_file.hpp>
25 #include <boost/corosio/native/detail/validate_fd.hpp>
26 #include <boost/corosio/random_access_file.hpp>
27
28 #include <cstdint>
29 #include <filesystem>
30 #include <limits>
31 #include <memory>
32 #include <mutex>
33 #include <system_error>
34 #include <unordered_map>
35
36 #include <fcntl.h>
37 #include <sys/stat.h>
38 #include <sys/types.h>
39 #include <unistd.h>
40
41 namespace boost::corosio::detail {
42
43 class uring_random_access_file_service;
44
45 /** Native io_uring random-access-file implementation.
46
47 Async `read_some_at` / `write_some_at` submit `IORING_OP_READV`
48 / `IORING_OP_WRITEV` with the caller-supplied offset. Metadata
49 operations (open, size, resize, sync, close) are synchronous
50 syscalls.
51
52 @par Thread Safety
53 Concurrent `read_some_at` / `write_some_at` calls on the same
54 file at distinct offsets are safe; ordering between two
55 submissions at the same offset is unspecified at the kernel
56 level (matches POSIX `pread(2)` / `pwrite(2)` semantics).
57 */
58 class BOOST_COROSIO_DECL uring_random_access_file final
59 : public random_access_file::implementation
60 , public std::enable_shared_from_this<uring_random_access_file>
61 , public intrusive_list<uring_random_access_file>::node
62 {
63 friend class uring_random_access_file_service;
64
65 int fd_ = -1;
66 uring_scheduler* sched_ = nullptr;
67
68 // Random-access files legitimately support concurrent ops at
69 // different offsets on the same fd (e.g. parallel reads in
70 // testConcurrentReads). Embedding a single slot would smash
71 // state across calls; ops are heap-allocated per submission.
72
73 public:
74 111x explicit uring_random_access_file(uring_scheduler& sched) noexcept
75 111x : sched_(&sched)
76 {
77 111x }
78
79 110x ~uring_random_access_file() override
80 110x {
81 110x close_file();
82 110x }
83
84 // -- random_access_file::implementation --
85
86 std::coroutine_handle<> read_some_at(
87 std::uint64_t,
88 capy::continuation&,
89 capy::executor_ref,
90 buffer_param,
91 std::stop_token,
92 std::error_code*,
93 std::size_t*) override;
94
95 std::coroutine_handle<> write_some_at(
96 std::uint64_t,
97 capy::continuation&,
98 capy::executor_ref,
99 buffer_param,
100 std::stop_token,
101 std::error_code*,
102 std::size_t*) override;
103
104 542x native_handle_type native_handle() const noexcept override
105 {
106 542x return fd_;
107 }
108
109 2x void cancel() noexcept override
110 {
111 2x if (fd_ >= 0)
112 2x sched_->submit_cancel_by_fd(fd_);
113 2x }
114
115 5x std::uint64_t size() const override
116 {
117 file_stat_t st;
118 5x if (file_fstat(fd_, &st) < 0)
119 2x throw_system_error(make_err(errno), "random_access_file::size");
120 3x return static_cast<std::uint64_t>(st.st_size);
121 }
122
123 6x std::error_code resize(std::uint64_t new_size) noexcept override
124 {
125 6x if (new_size > static_cast<std::uint64_t>(
126 6x (std::numeric_limits<file_off_t>::max)()))
127 1x return make_err(EOVERFLOW);
128 5x if (file_ftruncate(fd_, static_cast<file_off_t>(new_size)) < 0)
129 3x return make_err(errno);
130 2x return {};
131 }
132
133 3x std::error_code sync_data() noexcept override
134 {
135 #if BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
136 3x if (::fdatasync(fd_) < 0)
137 #else
138 if (::fsync(fd_) < 0)
139 #endif
140 2x return make_err(errno);
141 1x return {};
142 }
143
144 3x std::error_code sync_all() noexcept override
145 {
146 3x if (::fsync(fd_) < 0)
147 2x return make_err(errno);
148 1x return {};
149 }
150
151 1x native_handle_type release() override
152 {
153 // Flush the cancel while the fd is still open, so the kernel
154 // resolves it before the caller can close and recycle the
155 // number.
156 1x if (fd_ >= 0)
157 1x sched_->cancel_and_flush(fd_);
158 1x int fd = fd_;
159 1x fd_ = -1;
160 1x return fd;
161 }
162
163 3x std::error_code assign(native_handle_type handle) noexcept override
164 {
165 // The public assign() guarantees the object is closed.
166 3x if (auto ec = validate_file_fd(handle))
167 2x return ec;
168
169 1x fd_ = handle;
170 1x return {};
171 }
172
173 // -- Internal --
174
175 /// Open the file. Synchronous; sets `fd_`. Caller is the service.
176 std::error_code
177 102x open_file(std::filesystem::path const& path, file_base::flags mode)
178 {
179 102x close_file();
180
181 102x int oflags = 0;
182 102x unsigned access = static_cast<unsigned>(mode) & 3u;
183 102x if (access == static_cast<unsigned>(file_base::read_write))
184 19x oflags |= O_RDWR;
185 83x else if (access == static_cast<unsigned>(file_base::write_only))
186 32x oflags |= O_WRONLY;
187 else
188 51x oflags |= O_RDONLY;
189
190 102x if ((mode & file_base::create) != file_base::flags(0))
191 14x oflags |= O_CREAT;
192 102x if ((mode & file_base::exclusive) != file_base::flags(0))
193 2x oflags |= O_EXCL;
194 102x if ((mode & file_base::truncate) != file_base::flags(0))
195 7x oflags |= O_TRUNC;
196 102x if ((mode & file_base::sync_all_on_write) != file_base::flags(0))
197 1x oflags |= O_SYNC;
198
199 102x oflags |= O_CLOEXEC;
200
201 102x int fd = ::open(path.c_str(), oflags | large_file_open_flag, 0666);
202 102x if (fd < 0)
203 4x return make_err(errno);
204
205 98x fd_ = fd;
206
207 #ifdef POSIX_FADV_RANDOM
208 // Hint the page cache that access will be random; matches
209 // the POSIX backend.
210 98x ::posix_fadvise(fd_, 0, 0, POSIX_FADV_RANDOM);
211 #endif
212
213 98x return {};
214 }
215
216 /// Cancel any in-flight ops and close the fd. Idempotent.
217 522x void close_file() noexcept
218 {
219 522x if (fd_ >= 0)
220 {
221 // The kernel may run a queued pipe write as task work at
222 // either kernel entry below; with the reader already gone
223 // that raises SIGPIPE.
224 98x scoped_sigpipe_block no_sigpipe;
225 98x sched_->cancel_and_flush(fd_);
226 98x ::close(fd_);
227 98x fd_ = -1;
228 98x }
229 522x }
230 };
231
232 inline std::coroutine_handle<>
233 169x uring_random_access_file::read_some_at(
234 std::uint64_t user_offset,
235 capy::continuation& cont,
236 capy::executor_ref ex,
237 buffer_param buffers,
238 std::stop_token token,
239 std::error_code* ec,
240 std::size_t* bytes)
241 {
242 169x auto op_guard = std::make_unique<uring_random_access_read_op>();
243 338x op_guard->prepare(
244 cont.h, ex, ec, bytes, fd_, static_cast<std::int64_t>(user_offset),
245 338x sched_, shared_from_this(), buffers, token);
246 169x op_guard->awaiting = &cont;
247 169x sched_->work_started();
248
249 // Closed-object contract outranks the zero-length no-op.
250 169x if (fd_ < 0)
251 {
252 2x op_guard->empty_buffer = false;
253 2x op_guard->res = -EBADF;
254 2x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
255 2x sched_->push_completed_locked(op_guard.release());
256 2x return std::noop_coroutine();
257 2x }
258
259 333x if (op_guard->empty_buffer ||
260 166x op_guard->cancelled.load(std::memory_order_acquire))
261 {
262 1x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
263 1x sched_->push_completed_locked(op_guard.release());
264 1x return std::noop_coroutine();
265 1x }
266
267 166x uring_submit_op(*sched_, op_guard.release());
268 166x return std::noop_coroutine();
269 169x }
270
271 inline std::coroutine_handle<>
272 45x uring_random_access_file::write_some_at(
273 std::uint64_t user_offset,
274 capy::continuation& cont,
275 capy::executor_ref ex,
276 buffer_param buffers,
277 std::stop_token token,
278 std::error_code* ec,
279 std::size_t* bytes)
280 {
281 45x auto op_guard = std::make_unique<uring_random_access_write_op>();
282 90x op_guard->prepare(
283 cont.h, ex, ec, bytes, fd_, static_cast<std::int64_t>(user_offset),
284 90x sched_, shared_from_this(), buffers, token);
285 45x op_guard->awaiting = &cont;
286 45x sched_->work_started();
287
288 // Closed-object contract outranks the zero-length no-op.
289 45x if (fd_ < 0)
290 {
291 1x op_guard->empty_buffer = false;
292 1x op_guard->res = -EBADF;
293 1x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
294 1x sched_->push_completed_locked(op_guard.release());
295 1x return std::noop_coroutine();
296 1x }
297
298 87x if (op_guard->empty_buffer ||
299 43x op_guard->cancelled.load(std::memory_order_acquire))
300 {
301 1x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
302 1x sched_->push_completed_locked(op_guard.release());
303 1x return std::noop_coroutine();
304 1x }
305
306 43x uring_submit_op(*sched_, op_guard.release());
307 43x return std::noop_coroutine();
308 45x }
309
310 /** Native io_uring random-access-file service.
311
312 Owns all `uring_random_access_file` impls. Replaces
313 `posix_random_access_file_service` for the io_uring backend;
314 registered under the abstract `random_access_file_service` key
315 by `uring_t::construct`.
316 */
317 class BOOST_COROSIO_DECL uring_random_access_file_service final
318 : public uring_file_service_base<
319 uring_random_access_file_service,
320 random_access_file_service,
321 uring_random_access_file>
322 {
323 using base_service = uring_file_service_base<
324 uring_random_access_file_service,
325 random_access_file_service,
326 uring_random_access_file>;
327
328 public:
329 82x explicit uring_random_access_file_service(
330 capy::execution_context& ctx)
331 82x : base_service(ctx.use_service<uring_scheduler>())
332 {
333 82x }
334
335 // construct / destroy / close / shutdown / scheduler() are inherited
336 // from uring_file_service_base.
337
338 102x std::error_code open_file(
339 random_access_file::implementation& impl,
340 std::filesystem::path const& path,
341 file_base::flags mode) override
342 {
343 102x return static_cast<uring_random_access_file&>(impl).open_file(
344 102x path, mode);
345 }
346 };
347
348 } // namespace boost::corosio::detail
349
350 #endif // BOOST_COROSIO_HAS_URING
351
352 #endif // BOOST_COROSIO_NATIVE_DETAIL_URING_URING_RANDOM_ACCESS_FILE_HPP
353