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/capy/ex/executor_ref.hpp>
29 : #include <boost/capy/error.hpp>
30 : #include <boost/capy/buffers.hpp>
31 :
32 : #include <atomic>
33 : #include <coroutine>
34 : #include <cstddef>
35 : #include <cstdint>
36 : #include <filesystem>
37 : #include <limits>
38 : #include <mutex>
39 : #include <optional>
40 : #include <stop_token>
41 : #include <system_error>
42 :
43 : #include <errno.h>
44 : #include <fcntl.h>
45 : #include <sys/stat.h>
46 : #include <sys/uio.h>
47 : #include <unistd.h>
48 :
49 : /*
50 : POSIX Random-Access File Implementation
51 : ========================================
52 :
53 : Each async read/write acquires an raf_op that serves as both the
54 : thread-pool work item and the scheduler completion op. This
55 : allows unlimited concurrent operations on the same file object.
56 :
57 : Ops recycle through the file's free_ops_ list on completion or
58 : shutdown, so steady-state reads and writes allocate nothing; the
59 : file destructor frees the recycled storage.
60 : */
61 :
62 : namespace boost::corosio::detail {
63 :
64 : struct scheduler;
65 : class posix_random_access_file_service;
66 :
67 : /** Random-access file implementation for POSIX backends. */
68 : class posix_random_access_file final
69 : : public random_access_file::implementation
70 : , public intrusive_list<posix_random_access_file>::node
71 : {
72 : friend class posix_random_access_file_service;
73 :
74 : public:
75 : static constexpr std::size_t max_buffers = 16;
76 :
77 : /** Per-operation state, acquired from the file's free list.
78 :
79 : Inherits from `coro_op` (for scheduler completion plus the shared
80 : coroutine, cancellation and keepalive machinery) and
81 : `pool_work_item` (for thread-pool dispatch). The intrusive hook
82 : links it into `outstanding_ops_` while in flight and `free_ops_`
83 : once recycled — membership is strictly sequential, so one hook
84 : serves both under `ops_mutex_`. `coro_op` leads the base list so
85 : 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; `object_ref_` 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 : /// Recycle into the owning service's pool. Defined out-of-line
113 : /// after posix_random_access_file_service for its complete type.
114 : void retire() noexcept override;
115 :
116 : /** Reset for recycling.
117 :
118 : `close_file()` already drove fd_ to its closed value before
119 : the refcount reached zero. Every `raf_op` holds its own
120 : `object_ref` while in flight, so the refcount cannot reach
121 : zero while one is outstanding — asserting `outstanding_ops_`
122 : is empty is therefore a precondition check, not a defensive
123 : one. `free_ops_` deliberately survives recycling: the storage
124 : belongs to this impl and only the destructor frees it.
125 :
126 : @pre refs_ == 0, fd closed, no op in flight.
127 : */
128 MIS 0 : void reuse() noexcept
129 : {
130 0 : BOOST_COROSIO_ASSERT(fd_ == -1);
131 0 : BOOST_COROSIO_ASSERT(outstanding_ops_.empty());
132 0 : }
133 :
134 HIT 438 : ~posix_random_access_file() override
135 219 : {
136 621 : while (auto* op = free_ops_.pop_front())
137 402 : delete op;
138 438 : }
139 :
140 : // -- random_access_file::implementation --
141 :
142 : std::coroutine_handle<> read_some_at(
143 : std::uint64_t offset,
144 : std::coroutine_handle<>,
145 : capy::executor_ref,
146 : buffer_param,
147 : std::stop_token,
148 : std::error_code*,
149 : std::size_t*) override;
150 :
151 : std::coroutine_handle<> write_some_at(
152 : std::uint64_t offset,
153 : std::coroutine_handle<>,
154 : capy::executor_ref,
155 : buffer_param,
156 : std::stop_token,
157 : std::error_code*,
158 : std::size_t*) override;
159 :
160 1160 : native_handle_type native_handle() const noexcept override
161 : {
162 1160 : return fd_;
163 : }
164 :
165 638 : void cancel() noexcept override
166 : {
167 638 : std::lock_guard<std::mutex> lock(ops_mutex_);
168 638 : outstanding_ops_.for_each([](raf_op* op) {
169 8 : op->cancelled.store(true, std::memory_order_release);
170 8 : });
171 638 : }
172 :
173 : std::uint64_t size() const override;
174 : std::error_code resize(std::uint64_t new_size) noexcept override;
175 : std::error_code sync_data() noexcept override;
176 : std::error_code sync_all() noexcept override;
177 : native_handle_type release() override;
178 : std::error_code assign(native_handle_type handle) noexcept override;
179 :
180 : std::error_code
181 : open_file(std::filesystem::path const& path, file_base::flags mode);
182 : void close_file() noexcept;
183 :
184 : private:
185 : /** Pop a recycled op, or allocate on the cold path.
186 :
187 : Per-use fields are filled by the caller; the fields `prepare`-
188 : style reuse must not inherit from the previous run (`errn`,
189 : `bytes_transferred`) are reset here. `start()` resets the
190 : cancellation machinery.
191 : */
192 494 : raf_op* acquire_op()
193 : {
194 : raf_op* op;
195 : {
196 494 : std::lock_guard<std::mutex> lock(ops_mutex_);
197 494 : op = free_ops_.pop_front();
198 494 : }
199 494 : if (!op)
200 402 : op = new raf_op();
201 494 : op->errn = 0;
202 494 : op->bytes_transferred = 0;
203 494 : return op;
204 : }
205 :
206 : posix_random_access_file_service& svc_;
207 : int fd_ = -1;
208 : std::mutex ops_mutex_;
209 : intrusive_list<raf_op> outstanding_ops_;
210 : intrusive_list<raf_op> free_ops_;
211 : };
212 :
213 : // ---------------------------------------------------------------------------
214 : // Inline implementation
215 : // ---------------------------------------------------------------------------
216 :
217 219 : inline posix_random_access_file::posix_random_access_file(
218 219 : posix_random_access_file_service& svc) noexcept
219 219 : : svc_(svc)
220 : {
221 219 : }
222 :
223 : inline std::error_code
224 203 : posix_random_access_file::open_file(
225 : std::filesystem::path const& path, file_base::flags mode)
226 : {
227 203 : close_file();
228 :
229 203 : int oflags = 0;
230 :
231 203 : unsigned access = static_cast<unsigned>(mode) & 3u;
232 203 : if (access == static_cast<unsigned>(file_base::read_write))
233 33 : oflags |= O_RDWR;
234 170 : else if (access == static_cast<unsigned>(file_base::write_only))
235 64 : oflags |= O_WRONLY;
236 : else
237 106 : oflags |= O_RDONLY;
238 :
239 203 : if ((mode & file_base::create) != file_base::flags(0))
240 28 : oflags |= O_CREAT;
241 203 : if ((mode & file_base::exclusive) != file_base::flags(0))
242 4 : oflags |= O_EXCL;
243 203 : if ((mode & file_base::truncate) != file_base::flags(0))
244 14 : oflags |= O_TRUNC;
245 203 : if ((mode & file_base::sync_all_on_write) != file_base::flags(0))
246 2 : oflags |= O_SYNC;
247 : // Note: no O_APPEND for random access files
248 :
249 203 : int fd = ::open(path.c_str(), oflags, 0666);
250 203 : if (fd < 0)
251 9 : return make_err(errno);
252 :
253 194 : fd_ = fd;
254 :
255 : #ifdef POSIX_FADV_RANDOM
256 194 : ::posix_fadvise(fd_, 0, 0, POSIX_FADV_RANDOM);
257 : #endif
258 :
259 194 : return {};
260 : }
261 :
262 : inline void
263 844 : posix_random_access_file::close_file() noexcept
264 : {
265 844 : if (fd_ >= 0)
266 : {
267 198 : ::close(fd_);
268 198 : fd_ = -1;
269 : }
270 844 : }
271 :
272 : inline std::uint64_t
273 13 : posix_random_access_file::size() const
274 : {
275 : struct stat st;
276 13 : if (::fstat(fd_, &st) < 0)
277 5 : throw_system_error(make_err(errno), "random_access_file::size");
278 8 : return static_cast<std::uint64_t>(st.st_size);
279 : }
280 :
281 : inline std::error_code
282 13 : posix_random_access_file::resize(std::uint64_t new_size) noexcept
283 : {
284 13 : if (new_size >
285 13 : static_cast<std::uint64_t>((std::numeric_limits<off_t>::max)()))
286 2 : return make_err(EOVERFLOW);
287 11 : if (::ftruncate(fd_, static_cast<off_t>(new_size)) < 0)
288 7 : return make_err(errno);
289 4 : return {};
290 : }
291 :
292 : inline std::error_code
293 9 : posix_random_access_file::sync_data() noexcept
294 : {
295 : #if BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
296 9 : if (::fdatasync(fd_) < 0)
297 : #else // BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
298 : if (::fsync(fd_) < 0)
299 : #endif // BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
300 7 : return make_err(errno);
301 2 : return {};
302 : }
303 :
304 : inline std::error_code
305 9 : posix_random_access_file::sync_all() noexcept
306 : {
307 9 : if (::fsync(fd_) < 0)
308 7 : return make_err(errno);
309 2 : return {};
310 : }
311 :
312 : inline native_handle_type
313 3 : posix_random_access_file::release()
314 : {
315 3 : int fd = fd_;
316 3 : fd_ = -1;
317 3 : return fd;
318 : }
319 :
320 : inline std::error_code
321 7 : posix_random_access_file::assign(native_handle_type handle) noexcept
322 : {
323 7 : close_file();
324 7 : fd_ = handle;
325 7 : return {};
326 : }
327 :
328 : // read_some_at, write_some_at are defined in
329 : // posix_random_access_file_service.hpp after the service.
330 :
331 : // -- raf_op completion handler (scheduler thread) --
332 :
333 : inline void
334 488 : posix_random_access_file::raf_op::operator()()
335 : {
336 488 : stop_cb.reset();
337 :
338 : // Empty buffers never reach the pool (diverted at initiation), so
339 : // empty_buffer stays false and a 0-byte read is a genuine EOF.
340 901 : decode_io_result(
341 488 : ec_out, bytes_out, cancelled.load(std::memory_order_acquire),
342 488 : errn != 0 ? make_err(errn) : std::error_code{}, is_read,
343 : bytes_transferred, /*empty_buffer=*/false);
344 :
345 : // Copy out everything needed after recycling: once this op is on
346 : // free_ops_ a concurrent initiation may pop and refill it. The
347 : // keepalive drops after the push so a final release never runs
348 : // under ops_mutex_.
349 488 : auto keep = std::move(object_ref_);
350 488 : auto coro = h;
351 488 : auto exec = ex;
352 : {
353 488 : std::lock_guard<std::mutex> lock(file_->ops_mutex_);
354 488 : file_->outstanding_ops_.remove(this);
355 488 : file_->free_ops_.push_front(this);
356 488 : }
357 488 : keep.reset();
358 488 : exec.on_work_finished();
359 488 : coro.resume();
360 488 : }
361 :
362 : // -- raf_op shutdown cleanup --
363 :
364 : inline void
365 6 : posix_random_access_file::raf_op::destroy()
366 : {
367 6 : stop_cb.reset();
368 6 : auto keep = std::move(object_ref_);
369 6 : auto exec = ex;
370 : {
371 6 : std::lock_guard<std::mutex> lock(file_->ops_mutex_);
372 6 : file_->outstanding_ops_.remove(this);
373 6 : file_->free_ops_.push_front(this);
374 6 : }
375 6 : keep.reset();
376 6 : exec.on_work_finished();
377 6 : }
378 :
379 : } // namespace boost::corosio::detail
380 :
381 : #endif // BOOST_COROSIO_POSIX
382 :
383 : #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RANDOM_ACCESS_FILE_HPP
|