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

100.0% Lines (127 / 127) 100.0% Functions (17 / 17)
uring_stream_file.hpp
f(x) Functions (17)
Function Calls Lines Blocks
boost::corosio::detail::uring_stream_file::uring_stream_file(boost::corosio::detail::uring_scheduler&) :82 103x 100.0% 100.0% boost::corosio::detail::uring_stream_file::~uring_stream_file() :86 103x 100.0% 100.0% boost::corosio::detail::uring_stream_file::native_handle() const :111 325x 100.0% 100.0% boost::corosio::detail::uring_stream_file::cancel() :116 1x 100.0% 100.0% boost::corosio::detail::uring_stream_file::size() const :122 6x 100.0% 100.0% boost::corosio::detail::uring_stream_file::resize(unsigned long) :130 5x 100.0% 100.0% boost::corosio::detail::uring_stream_file::sync_data() :140 3x 100.0% 100.0% boost::corosio::detail::uring_stream_file::sync_all() :151 3x 100.0% 100.0% boost::corosio::detail::uring_stream_file::release() :158 1x 100.0% 100.0% boost::corosio::detail::uring_stream_file::assign(int) :170 3x 100.0% 100.0% boost::corosio::detail::uring_stream_file::seek(long, boost::corosio::file_base::seek_basis) :181 11x 100.0% 100.0% boost::corosio::detail::uring_stream_file::open_file(std::filesystem::__cxx11::path const&, boost::corosio::file_base::flags) :199 91x 100.0% 100.0% boost::corosio::detail::uring_stream_file::close_file() :241 485x 100.0% 100.0% boost::corosio::detail::uring_stream_file::read_some(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :257 38x 100.0% 92.0% boost::corosio::detail::uring_stream_file::write_some(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :292 40x 100.0% 92.0% boost::corosio::detail::uring_stream_file_service::uring_stream_file_service(boost::capy::execution_context&) :344 75x 100.0% 100.0% boost::corosio::detail::uring_stream_file_service::open_file(boost::corosio::stream_file::implementation&, std::filesystem::__cxx11::path const&, boost::corosio::file_base::flags) :352 91x 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_STREAM_FILE_HPP
12 #define BOOST_COROSIO_NATIVE_DETAIL_URING_URING_STREAM_FILE_HPP
13
14 #include <boost/corosio/detail/platform.hpp>
15
16 #if BOOST_COROSIO_HAS_URING
17
18 #include <boost/corosio/detail/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/stream_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_stream_file_service;
44
45 /** Native io_uring stream-file implementation.
46
47 Async `read_some` / `write_some` submit `IORING_OP_READV` /
48 `IORING_OP_WRITEV` with `offset == -1` (kernel f_pos). All
49 metadata operations (open, size, resize, sync, seek, close)
50 are synchronous syscalls.
51
52 @par Thread Safety
53 Concurrent `read_some` / `write_some` calls on the same file
54 interleave at the kernel level (matches POSIX `read(2)` /
55 `write(2)` semantics on a shared positional fd).
56
57 @note On `O_APPEND` open this backend relies on the kernel's
58 `f_pos` rather than tracking the offset in user space. Writes
59 still go to EOF atomically per `O_APPEND` semantics, but
60 `seek(0, seek_cur)` immediately after an append-mode open
61 returns `0` (the current f_pos), not the file size — observably
62 different from the POSIX backend, which seeds an internal offset
63 to size-at-open. Both behaviours are valid; documented for
64 cross-backend symmetry.
65 */
66 class BOOST_COROSIO_DECL uring_stream_file final
67 : public stream_file::implementation
68 , public std::enable_shared_from_this<uring_stream_file>
69 , public intrusive_list<uring_stream_file>::node
70 {
71 friend class uring_stream_file_service;
72
73 int fd_ = -1;
74 uring_scheduler* sched_ = nullptr;
75
76 // Per-fd op slots — embedded to eliminate per-call heap allocation.
77 // Single-pending invariant per slot.
78 uring_file_read_op rd_;
79 uring_file_write_op wr_;
80
81 public:
82 103x explicit uring_stream_file(uring_scheduler& sched) noexcept : sched_(&sched)
83 {
84 103x }
85
86 103x ~uring_stream_file() override
87 103x {
88 103x close_file();
89 103x }
90
91 // -- io_stream::implementation --
92
93 std::coroutine_handle<> read_some(
94 std::coroutine_handle<>,
95 capy::executor_ref,
96 buffer_param,
97 std::stop_token,
98 std::error_code*,
99 std::size_t*) override;
100
101 std::coroutine_handle<> write_some(
102 std::coroutine_handle<>,
103 capy::executor_ref,
104 buffer_param,
105 std::stop_token,
106 std::error_code*,
107 std::size_t*) override;
108
109 // -- stream_file::implementation --
110
111 325x native_handle_type native_handle() const noexcept override
112 {
113 325x return fd_;
114 }
115
116 1x void cancel() noexcept override
117 {
118 1x if (fd_ >= 0)
119 1x sched_->submit_cancel_by_fd(fd_);
120 1x }
121
122 6x std::uint64_t size() const override
123 {
124 file_stat_t st;
125 6x if (file_fstat(fd_, &st) < 0)
126 2x throw_system_error(make_err(errno), "stream_file::size");
127 4x return static_cast<std::uint64_t>(st.st_size);
128 }
129
130 5x std::error_code resize(std::uint64_t new_size) noexcept override
131 {
132 5x if (new_size > static_cast<std::uint64_t>(
133 5x (std::numeric_limits<file_off_t>::max)()))
134 1x return make_err(EOVERFLOW);
135 4x if (file_ftruncate(fd_, static_cast<file_off_t>(new_size)) < 0)
136 3x return make_err(errno);
137 1x return {};
138 }
139
140 3x std::error_code sync_data() noexcept override
141 {
142 #if BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
143 3x if (::fdatasync(fd_) < 0)
144 #else
145 if (::fsync(fd_) < 0)
146 #endif
147 2x return make_err(errno);
148 1x return {};
149 }
150
151 3x std::error_code sync_all() noexcept override
152 {
153 3x if (::fsync(fd_) < 0)
154 2x return make_err(errno);
155 1x return {};
156 }
157
158 1x native_handle_type release() override
159 {
160 // Flush the cancel while the fd is still open, so the kernel
161 // resolves it before the caller can close and recycle the
162 // number.
163 1x if (fd_ >= 0)
164 1x sched_->cancel_and_flush(fd_);
165 1x int fd = fd_;
166 1x fd_ = -1;
167 1x return fd;
168 }
169
170 3x std::error_code assign(native_handle_type handle) noexcept override
171 {
172 // The public assign() guarantees the object is closed.
173 3x if (auto ec = validate_file_fd(handle))
174 2x return ec;
175
176 1x fd_ = handle;
177 1x return {};
178 }
179
180 capy::io_result<std::uint64_t>
181 11x seek(std::int64_t offset, file_base::seek_basis origin) noexcept override
182 {
183 11x int whence = SEEK_SET;
184 11x if (origin == file_base::seek_cur)
185 2x whence = SEEK_CUR;
186 9x else if (origin == file_base::seek_end)
187 4x whence = SEEK_END;
188
189 11x file_off_t r = file_lseek(fd_, static_cast<file_off_t>(offset), whence);
190 11x if (r == static_cast<file_off_t>(-1))
191 5x return {make_err(errno), 0};
192 6x return {std::error_code{}, static_cast<std::uint64_t>(r)};
193 }
194
195 // -- Internal --
196
197 /// Open the file. Synchronous; sets `fd_`. Caller is the service.
198 std::error_code
199 91x open_file(std::filesystem::path const& path, file_base::flags mode)
200 {
201 91x close_file();
202
203 91x int oflags = 0;
204 91x unsigned access = static_cast<unsigned>(mode) & 3u;
205 91x if (access == static_cast<unsigned>(file_base::read_write))
206 7x oflags |= O_RDWR;
207 84x else if (access == static_cast<unsigned>(file_base::write_only))
208 37x oflags |= O_WRONLY;
209 else
210 47x oflags |= O_RDONLY;
211
212 91x if ((mode & file_base::create) != file_base::flags(0))
213 14x oflags |= O_CREAT;
214 91x if ((mode & file_base::exclusive) != file_base::flags(0))
215 1x oflags |= O_EXCL;
216 91x if ((mode & file_base::truncate) != file_base::flags(0))
217 8x oflags |= O_TRUNC;
218 91x if ((mode & file_base::append) != file_base::flags(0))
219 1x oflags |= O_APPEND;
220 91x if ((mode & file_base::sync_all_on_write) != file_base::flags(0))
221 1x oflags |= O_SYNC;
222
223 91x oflags |= O_CLOEXEC;
224
225 91x int fd = ::open(path.c_str(), oflags | large_file_open_flag, 0666);
226 91x if (fd < 0)
227 4x return make_err(errno);
228
229 87x fd_ = fd;
230
231 #ifdef POSIX_FADV_SEQUENTIAL
232 // Hint the page cache about the access pattern; matches the
233 // POSIX backend.
234 87x ::posix_fadvise(fd_, 0, 0, POSIX_FADV_SEQUENTIAL);
235 #endif
236
237 87x return {};
238 }
239
240 /// Cancel any in-flight ops and close the fd. Idempotent.
241 485x void close_file() noexcept
242 {
243 485x if (fd_ >= 0)
244 {
245 // The kernel may run a queued pipe write as task work at
246 // either kernel entry below; with the reader already gone
247 // that raises SIGPIPE.
248 87x scoped_sigpipe_block no_sigpipe;
249 87x sched_->cancel_and_flush(fd_);
250 87x ::close(fd_);
251 87x fd_ = -1;
252 87x }
253 485x }
254 };
255
256 inline std::coroutine_handle<>
257 38x uring_stream_file::read_some(
258 std::coroutine_handle<> h,
259 capy::executor_ref ex,
260 buffer_param buffers,
261 std::stop_token token,
262 std::error_code* ec,
263 std::size_t* bytes)
264 {
265 38x rd_.prepare(
266 76x h, ex, ec, bytes, fd_, /*file_offset=*/-1, sched_, shared_from_this(),
267 buffers, token);
268 38x sched_->work_started();
269
270 // Closed-object contract outranks the zero-length no-op.
271 38x if (fd_ < 0)
272 {
273 3x rd_.empty_buffer = false;
274 3x rd_.res = -EBADF;
275 3x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
276 3x sched_->push_completed_locked(&rd_);
277 3x return std::noop_coroutine();
278 3x }
279
280 35x if (rd_.empty_buffer || rd_.cancelled.load(std::memory_order_acquire))
281 {
282 1x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
283 1x sched_->push_completed_locked(&rd_);
284 1x return std::noop_coroutine();
285 1x }
286
287 34x uring_submit_op(*sched_, &rd_);
288 34x return std::noop_coroutine();
289 }
290
291 inline std::coroutine_handle<>
292 40x uring_stream_file::write_some(
293 std::coroutine_handle<> h,
294 capy::executor_ref ex,
295 buffer_param buffers,
296 std::stop_token token,
297 std::error_code* ec,
298 std::size_t* bytes)
299 {
300 40x wr_.prepare(
301 80x h, ex, ec, bytes, fd_, /*file_offset=*/-1, sched_, shared_from_this(),
302 buffers, token);
303 40x sched_->work_started();
304
305 // Closed-object contract outranks the zero-length no-op.
306 40x if (fd_ < 0)
307 {
308 3x wr_.empty_buffer = false;
309 3x wr_.res = -EBADF;
310 3x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
311 3x sched_->push_completed_locked(&wr_);
312 3x return std::noop_coroutine();
313 3x }
314
315 37x if (wr_.empty_buffer || wr_.cancelled.load(std::memory_order_acquire))
316 {
317 1x uring_scheduler::lock_type lock(sched_->dispatch_mutex());
318 1x sched_->push_completed_locked(&wr_);
319 1x return std::noop_coroutine();
320 1x }
321
322 36x uring_submit_op(*sched_, &wr_);
323 36x return std::noop_coroutine();
324 }
325
326 /** Native io_uring stream-file service.
327
328 Owns all `uring_stream_file` impls. Replaces
329 `posix_stream_file_service` for the io_uring backend; registered
330 under the abstract `file_service` key by `uring_t::construct`.
331 */
332 class BOOST_COROSIO_DECL uring_stream_file_service final
333 : public uring_file_service_base<
334 uring_stream_file_service,
335 file_service,
336 uring_stream_file>
337 {
338 using base_service = uring_file_service_base<
339 uring_stream_file_service,
340 file_service,
341 uring_stream_file>;
342
343 public:
344 75x explicit uring_stream_file_service(capy::execution_context& ctx)
345 75x : base_service(ctx.use_service<uring_scheduler>())
346 {
347 75x }
348
349 // construct / destroy / close / shutdown / scheduler() are inherited
350 // from uring_file_service_base.
351
352 91x std::error_code open_file(
353 stream_file::implementation& impl,
354 std::filesystem::path const& path,
355 file_base::flags mode) override
356 {
357 91x return static_cast<uring_stream_file&>(impl).open_file(path, mode);
358 }
359 };
360
361 } // namespace boost::corosio::detail
362
363 #endif // BOOST_COROSIO_HAS_URING
364
365 #endif // BOOST_COROSIO_NATIVE_DETAIL_URING_URING_STREAM_FILE_HPP
366