TLA Line data Source code
1 : //
2 : // Copyright (c) 2026 Michael Vandeberg
3 : //
4 : // Distributed under the Boost Software License, Version 1.0. (See accompanying
5 : // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
6 : //
7 : // Official repository: https://github.com/cppalliance/corosio
8 : //
9 :
10 : #ifndef BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RANDOM_ACCESS_FILE_HPP
11 : #define BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RANDOM_ACCESS_FILE_HPP
12 :
13 : #include <boost/corosio/detail/platform.hpp>
14 :
15 : #if BOOST_COROSIO_POSIX
16 :
17 : #include <boost/corosio/detail/config.hpp>
18 : #include <boost/corosio/random_access_file.hpp>
19 : #include <boost/corosio/file_base.hpp>
20 : #include <boost/corosio/detail/intrusive.hpp>
21 : #include <boost/corosio/detail/scheduler_op.hpp>
22 : #include <boost/corosio/detail/thread_pool.hpp>
23 : #include <boost/corosio/detail/scheduler.hpp>
24 : #include <boost/corosio/detail/buffer_param.hpp>
25 : #include <boost/corosio/native/detail/coro_op.hpp>
26 : #include <boost/corosio/native/detail/coro_op_complete.hpp>
27 : #include <boost/corosio/native/detail/make_err.hpp>
28 : #include <boost/corosio/native/detail/validate_fd.hpp>
29 : #include <boost/capy/ex/executor_ref.hpp>
30 : #include <boost/capy/error.hpp>
31 : #include <boost/capy/buffers.hpp>
32 :
33 : #include <atomic>
34 : #include <coroutine>
35 : #include <cstddef>
36 : #include <cstdint>
37 : #include <filesystem>
38 : #include <limits>
39 : #include <memory>
40 : #include <mutex>
41 : #include <optional>
42 : #include <stop_token>
43 : #include <system_error>
44 :
45 : #include <errno.h>
46 : #include <fcntl.h>
47 : #include <sys/stat.h>
48 : #include <sys/uio.h>
49 : #include <unistd.h>
50 :
51 : /*
52 : POSIX Random-Access File Implementation
53 : ========================================
54 :
55 : Each async read/write heap-allocates an raf_op that serves
56 : as both the thread-pool work item and the scheduler completion
57 : op. This allows unlimited concurrent operations on the same
58 : file object, matching Asio's per-op allocation model.
59 :
60 : The raf_op self-deletes on completion or shutdown.
61 : */
62 :
63 : namespace boost::corosio::detail {
64 :
65 : struct scheduler;
66 : class posix_random_access_file_service;
67 :
68 : /** Random-access file implementation for POSIX backends. */
69 : class posix_random_access_file final
70 : : public random_access_file::implementation
71 : , public std::enable_shared_from_this<posix_random_access_file>
72 : , public intrusive_list<posix_random_access_file>::node
73 : {
74 : friend class posix_random_access_file_service;
75 :
76 : public:
77 : static constexpr std::size_t max_buffers = 16;
78 :
79 : /** Per-operation state, heap-allocated for each async call.
80 :
81 : Inherits from `coro_op` (for scheduler completion plus the shared
82 : coroutine, cancellation and keepalive machinery) and
83 : `pool_work_item` (for thread-pool dispatch). Linked into the
84 : file's outstanding_ops_ list for cancellation tracking. `coro_op`
85 : leads the base list so a `scheduler_op*` round-trips.
86 : */
87 : struct raf_op final
88 : : coro_op
89 : , pool_work_item
90 : , intrusive_list<raf_op>::node
91 : {
92 : iovec iovecs[max_buffers];
93 : int iovec_count = 0;
94 : std::uint64_t offset = 0;
95 :
96 : int errn = 0;
97 : std::size_t bytes_transferred = 0;
98 :
99 : // Raw back-pointer for the typed work; `impl_ptr` is the keepalive.
100 : posix_random_access_file* file_ = nullptr;
101 :
102 : void operator()() override;
103 : void destroy() override;
104 :
105 : /// Thread-pool work function: executes preadv/pwritev.
106 : static void do_work(pool_work_item*) noexcept;
107 : };
108 :
109 : explicit posix_random_access_file(
110 : posix_random_access_file_service& svc) noexcept;
111 :
112 : // -- random_access_file::implementation --
113 :
114 : std::coroutine_handle<> read_some_at(
115 : std::uint64_t offset,
116 : std::coroutine_handle<>,
117 : capy::executor_ref,
118 : buffer_param,
119 : std::stop_token,
120 : std::error_code*,
121 : std::size_t*) override;
122 :
123 : std::coroutine_handle<> write_some_at(
124 : std::uint64_t offset,
125 : std::coroutine_handle<>,
126 : capy::executor_ref,
127 : buffer_param,
128 : std::stop_token,
129 : std::error_code*,
130 : std::size_t*) override;
131 :
132 HIT 1105 : native_handle_type native_handle() const noexcept override
133 : {
134 1105 : return fd_;
135 : }
136 :
137 645 : void cancel() noexcept override
138 : {
139 645 : std::lock_guard<std::mutex> lock(ops_mutex_);
140 645 : outstanding_ops_.for_each([](raf_op* op) {
141 8 : op->cancelled.store(true, std::memory_order_release);
142 8 : });
143 645 : }
144 :
145 : std::uint64_t size() const override;
146 : std::error_code resize(std::uint64_t new_size) noexcept override;
147 : std::error_code sync_data() noexcept override;
148 : std::error_code sync_all() noexcept override;
149 : native_handle_type release() override;
150 : std::error_code assign(native_handle_type handle) noexcept override;
151 :
152 : std::error_code
153 : open_file(std::filesystem::path const& path, file_base::flags mode);
154 : void close_file() noexcept;
155 :
156 : private:
157 : posix_random_access_file_service& svc_;
158 : int fd_ = -1;
159 : std::mutex ops_mutex_;
160 : intrusive_list<raf_op> outstanding_ops_;
161 : };
162 :
163 : // ---------------------------------------------------------------------------
164 : // Inline implementation
165 : // ---------------------------------------------------------------------------
166 :
167 221 : inline posix_random_access_file::posix_random_access_file(
168 221 : posix_random_access_file_service& svc) noexcept
169 221 : : svc_(svc)
170 : {
171 221 : }
172 :
173 : inline std::error_code
174 205 : posix_random_access_file::open_file(
175 : std::filesystem::path const& path, file_base::flags mode)
176 : {
177 205 : close_file();
178 :
179 205 : int oflags = 0;
180 :
181 205 : unsigned access = static_cast<unsigned>(mode) & 3u;
182 205 : if (access == static_cast<unsigned>(file_base::read_write))
183 35 : oflags |= O_RDWR;
184 170 : else if (access == static_cast<unsigned>(file_base::write_only))
185 64 : oflags |= O_WRONLY;
186 : else
187 106 : oflags |= O_RDONLY;
188 :
189 205 : if ((mode & file_base::create) != file_base::flags(0))
190 32 : oflags |= O_CREAT;
191 205 : if ((mode & file_base::exclusive) != file_base::flags(0))
192 4 : oflags |= O_EXCL;
193 205 : if ((mode & file_base::truncate) != file_base::flags(0))
194 14 : oflags |= O_TRUNC;
195 205 : if ((mode & file_base::sync_all_on_write) != file_base::flags(0))
196 2 : oflags |= O_SYNC;
197 : // Note: no O_APPEND for random access files
198 :
199 205 : int fd = ::open(path.c_str(), oflags, 0666);
200 205 : if (fd < 0)
201 9 : return make_err(errno);
202 :
203 196 : fd_ = fd;
204 :
205 : #ifdef POSIX_FADV_RANDOM
206 196 : ::posix_fadvise(fd_, 0, 0, POSIX_FADV_RANDOM);
207 : #endif
208 :
209 196 : return {};
210 : }
211 :
212 : inline void
213 846 : posix_random_access_file::close_file() noexcept
214 : {
215 846 : if (fd_ >= 0)
216 : {
217 200 : ::close(fd_);
218 200 : fd_ = -1;
219 : }
220 846 : }
221 :
222 : inline std::uint64_t
223 13 : posix_random_access_file::size() const
224 : {
225 : struct stat st;
226 13 : if (::fstat(fd_, &st) < 0)
227 5 : throw_system_error(make_err(errno), "random_access_file::size");
228 8 : return static_cast<std::uint64_t>(st.st_size);
229 : }
230 :
231 : inline std::error_code
232 13 : posix_random_access_file::resize(std::uint64_t new_size) noexcept
233 : {
234 13 : if (new_size >
235 13 : static_cast<std::uint64_t>((std::numeric_limits<off_t>::max)()))
236 2 : return make_err(EOVERFLOW);
237 11 : if (::ftruncate(fd_, static_cast<off_t>(new_size)) < 0)
238 7 : return make_err(errno);
239 4 : return {};
240 : }
241 :
242 : inline std::error_code
243 7 : posix_random_access_file::sync_data() noexcept
244 : {
245 : #if BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
246 7 : if (::fdatasync(fd_) < 0)
247 : #else // BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
248 : if (::fsync(fd_) < 0)
249 : #endif // BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
250 5 : return make_err(errno);
251 2 : return {};
252 : }
253 :
254 : inline std::error_code
255 7 : posix_random_access_file::sync_all() noexcept
256 : {
257 7 : if (::fsync(fd_) < 0)
258 5 : return make_err(errno);
259 2 : return {};
260 : }
261 :
262 : inline native_handle_type
263 3 : posix_random_access_file::release()
264 : {
265 3 : int fd = fd_;
266 3 : fd_ = -1;
267 3 : return fd;
268 : }
269 :
270 : inline std::error_code
271 13 : posix_random_access_file::assign(native_handle_type handle) noexcept
272 : {
273 : // handle >= 0 guard: an unset impl reports native_handle() == -1, and a
274 : // caller-supplied -1 must fail as a bad fd, not a self-assign.
275 13 : if (handle >= 0 && handle == fd_)
276 2 : return std::make_error_code(std::errc::invalid_argument);
277 :
278 : // Validate before touching the held fd: a failed assign must leave
279 : // this object unchanged and the caller still owning handle.
280 11 : if (auto ec = validate_file_fd(handle))
281 4 : return ec;
282 :
283 : // cancel() first: an in-flight read_at/write_at's pool-thread
284 : // completion reads fd_ at execution time, not at post time, so
285 : // without this a pending op silently completes against the newly
286 : // adopted file instead of being cancelled. The service's
287 : // close(handle) / destroy() normally pair cancel()+close_file();
288 : // assign() bypasses that path and must do the same pairing itself.
289 7 : cancel();
290 7 : close_file();
291 7 : fd_ = handle;
292 7 : return {};
293 : }
294 :
295 : // read_some_at, write_some_at are defined in
296 : // posix_random_access_file_service.hpp after the service.
297 :
298 : // -- raf_op completion handler (scheduler thread) --
299 :
300 : inline void
301 422 : posix_random_access_file::raf_op::operator()()
302 : {
303 422 : stop_cb.reset();
304 :
305 : // Empty buffers never reach the pool (diverted at initiation), so
306 : // empty_buffer stays false and a 0-byte read is a genuine EOF.
307 773 : decode_io_result(
308 422 : ec_out, bytes_out, cancelled.load(std::memory_order_acquire),
309 422 : errn != 0 ? make_err(errn) : std::error_code{}, is_read,
310 : bytes_transferred, /*empty_buffer=*/false);
311 :
312 : {
313 422 : std::lock_guard<std::mutex> lock(file_->ops_mutex_);
314 422 : file_->outstanding_ops_.remove(this);
315 422 : }
316 :
317 422 : impl_ptr.reset();
318 :
319 422 : auto coro = h;
320 422 : ex.on_work_finished();
321 422 : delete this;
322 422 : coro.resume();
323 422 : }
324 :
325 : // -- raf_op shutdown cleanup --
326 :
327 : inline void
328 6 : posix_random_access_file::raf_op::destroy()
329 : {
330 6 : stop_cb.reset();
331 : {
332 6 : std::lock_guard<std::mutex> lock(file_->ops_mutex_);
333 6 : file_->outstanding_ops_.remove(this);
334 6 : }
335 6 : impl_ptr.reset();
336 6 : ex.on_work_finished();
337 6 : delete this;
338 6 : }
339 :
340 : } // namespace boost::corosio::detail
341 :
342 : #endif // BOOST_COROSIO_POSIX
343 :
344 : #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RANDOM_ACCESS_FILE_HPP
|