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