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_STREAM_FILE_HPP
11 : #define BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_STREAM_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/stream_file.hpp>
19 : #include <boost/corosio/file_base.hpp>
20 : #include <boost/corosio/detail/intrusive.hpp>
21 : #include <boost/corosio/detail/dispatch_coro.hpp>
22 : #include <boost/corosio/detail/scheduler_op.hpp>
23 : #include <boost/corosio/detail/thread_pool.hpp>
24 : #include <boost/corosio/detail/scheduler.hpp>
25 : #include <boost/corosio/detail/buffer_param.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 Stream File Implementation
55 : =================================
56 :
57 : Regular files cannot be monitored by epoll/kqueue/select — the kernel
58 : always reports them as ready. Blocking I/O (pread/pwrite) is dispatched
59 : to a shared thread pool, with completion posted back to the scheduler.
60 :
61 : This follows the same pattern as posix_resolver: pool_work_item for
62 : dispatch, scheduler_op for completion, shared_from_this for lifetime.
63 :
64 : Completion Flow
65 : ---------------
66 : 1. read_some() sets up file_read_op, posts to thread pool
67 : 2. Pool thread runs preadv() (blocking)
68 : 3. Pool thread stores results, posts scheduler_op to scheduler
69 : 4. Scheduler invokes op() which resumes the coroutine
70 :
71 : Single-Inflight Constraint
72 : --------------------------
73 : Only one asynchronous operation may be in flight at a time on a
74 : given file object. Concurrent read and write is not supported
75 : because both share offset_ without synchronization.
76 : */
77 :
78 : namespace boost::corosio::detail {
79 :
80 : struct scheduler;
81 : class posix_stream_file_service;
82 :
83 : /** Stream file implementation for POSIX backends.
84 :
85 : Each instance contains embedded operation objects (read_op_, write_op_)
86 : that are reused across calls. This avoids per-operation heap allocation.
87 : */
88 : class posix_stream_file final
89 : : public stream_file::implementation
90 : , public std::enable_shared_from_this<posix_stream_file>
91 : , public intrusive_list<posix_stream_file>::node
92 : {
93 : friend class posix_stream_file_service;
94 :
95 : public:
96 : static constexpr std::size_t max_buffers = 16;
97 :
98 : /** Operation state for a single file read or write.
99 :
100 : The coroutine, cancellation and keepalive machinery is inherited
101 : from `coro_op`; only the pool-path result state lives here.
102 : */
103 : struct file_op : coro_op
104 : {
105 : // Buffer data (copied from buffer_param at submission time)
106 : iovec iovecs[max_buffers];
107 : int iovec_count = 0;
108 :
109 : // Snapshotted at submission: assign() may replace fd_ and
110 : // offset_ while the worker runs.
111 : int fd = -1;
112 : std::uint64_t offset = 0;
113 : std::uint64_t generation = 0;
114 :
115 : // Result storage (populated by worker thread)
116 : int errn = 0;
117 : std::size_t bytes_transferred = 0;
118 :
119 HIT 1354 : file_op() = default;
120 :
121 650 : void reset() noexcept
122 : {
123 650 : iovec_count = 0;
124 650 : errn = 0;
125 650 : bytes_transferred = 0;
126 650 : is_read = false;
127 650 : cancelled.store(false, std::memory_order_relaxed);
128 650 : stop_cb.reset();
129 650 : impl_ptr.reset();
130 650 : ec_out = nullptr;
131 650 : bytes_out = nullptr;
132 650 : }
133 :
134 : void operator()() override;
135 : void destroy() override;
136 : };
137 :
138 : /** Pool work item for thread pool dispatch. */
139 : struct pool_op : pool_work_item
140 : {
141 : posix_stream_file* file_ = nullptr;
142 : std::shared_ptr<posix_stream_file> ref_;
143 : };
144 :
145 : explicit posix_stream_file(posix_stream_file_service& svc) noexcept;
146 :
147 : // -- io_stream::implementation --
148 :
149 : std::coroutine_handle<> read_some(
150 : std::coroutine_handle<>,
151 : capy::executor_ref,
152 : buffer_param,
153 : std::stop_token,
154 : std::error_code*,
155 : std::size_t*) override;
156 :
157 : std::coroutine_handle<> write_some(
158 : std::coroutine_handle<>,
159 : capy::executor_ref,
160 : buffer_param,
161 : std::stop_token,
162 : std::error_code*,
163 : std::size_t*) override;
164 :
165 : // -- stream_file::implementation --
166 :
167 2843 : native_handle_type native_handle() const noexcept override
168 : {
169 2843 : return fd_;
170 : }
171 :
172 2401 : void cancel() noexcept override
173 : {
174 2401 : read_op_.request_cancel();
175 2401 : write_op_.request_cancel();
176 2401 : }
177 :
178 : std::uint64_t size() const override;
179 : std::error_code resize(std::uint64_t new_size) noexcept override;
180 : std::error_code sync_data() noexcept override;
181 : std::error_code sync_all() noexcept override;
182 : native_handle_type release() override;
183 : std::error_code assign(native_handle_type handle) noexcept override;
184 : capy::io_result<std::uint64_t>
185 : seek(std::int64_t offset, file_base::seek_basis origin) noexcept override;
186 :
187 : // -- Internal --
188 :
189 : /** Open the file and store the fd. */
190 : std::error_code
191 : open_file(std::filesystem::path const& path, file_base::flags mode);
192 :
193 : /** Close the file descriptor. */
194 : void close_file() noexcept;
195 :
196 : private:
197 : posix_stream_file_service& svc_;
198 : int fd_ = -1;
199 : std::uint64_t offset_ = 0;
200 :
201 : // Bumped whenever the held file changes. A worker advances offset_
202 : // only if this still matches its op's snapshot; the mutex makes the
203 : // check and the advance one step against a concurrent assign().
204 : std::mutex offset_mutex_;
205 : std::uint64_t generation_ = 0;
206 :
207 : file_op read_op_;
208 : file_op write_op_;
209 : pool_op read_pool_op_;
210 : pool_op write_pool_op_;
211 :
212 : void bump_generation() noexcept;
213 :
214 : static void do_read_work(pool_work_item*) noexcept;
215 : static void do_write_work(pool_work_item*) noexcept;
216 : };
217 :
218 : // ---------------------------------------------------------------------------
219 : // Inline implementation
220 : // ---------------------------------------------------------------------------
221 :
222 677 : inline posix_stream_file::posix_stream_file(
223 677 : posix_stream_file_service& svc) noexcept
224 677 : : svc_(svc)
225 : {
226 677 : }
227 :
228 : inline std::error_code
229 656 : posix_stream_file::open_file(
230 : std::filesystem::path const& path, file_base::flags mode)
231 : {
232 656 : close_file();
233 :
234 656 : int oflags = 0;
235 :
236 : // Access mode
237 656 : unsigned access = static_cast<unsigned>(mode) & 3u;
238 656 : if (access == static_cast<unsigned>(file_base::read_write))
239 23 : oflags |= O_RDWR;
240 633 : else if (access == static_cast<unsigned>(file_base::write_only))
241 81 : oflags |= O_WRONLY;
242 : else
243 552 : oflags |= O_RDONLY;
244 :
245 : // Creation flags
246 656 : if ((mode & file_base::create) != file_base::flags(0))
247 40 : oflags |= O_CREAT;
248 656 : if ((mode & file_base::exclusive) != file_base::flags(0))
249 2 : oflags |= O_EXCL;
250 656 : if ((mode & file_base::truncate) != file_base::flags(0))
251 17 : oflags |= O_TRUNC;
252 656 : if ((mode & file_base::append) != file_base::flags(0))
253 8 : oflags |= O_APPEND;
254 656 : if ((mode & file_base::sync_all_on_write) != file_base::flags(0))
255 2 : oflags |= O_SYNC;
256 :
257 656 : int fd = ::open(path.c_str(), oflags | large_file_open_flag, 0666);
258 656 : if (fd < 0)
259 9 : return make_err(errno);
260 :
261 647 : fd_ = fd;
262 647 : offset_ = 0;
263 :
264 : // Append mode: position at end-of-file (preadv/pwritev use
265 : // explicit offsets, so O_APPEND alone is not sufficient).
266 647 : if ((mode & file_base::append) != file_base::flags(0))
267 : {
268 : file_stat_t st;
269 8 : if (file_fstat(fd, &st) < 0)
270 : {
271 5 : int err = errno;
272 5 : ::close(fd);
273 5 : fd_ = -1;
274 5 : return make_err(err);
275 : }
276 3 : offset_ = static_cast<std::uint64_t>(st.st_size);
277 : }
278 :
279 : #ifdef POSIX_FADV_SEQUENTIAL
280 642 : ::posix_fadvise(fd_, 0, 0, POSIX_FADV_SEQUENTIAL);
281 : #endif
282 :
283 642 : return {};
284 : }
285 :
286 : inline void
287 3055 : posix_stream_file::bump_generation() noexcept
288 : {
289 3055 : std::lock_guard<std::mutex> lock(offset_mutex_);
290 3055 : ++generation_;
291 3055 : }
292 :
293 : inline void
294 3051 : posix_stream_file::close_file() noexcept
295 : {
296 3051 : bump_generation();
297 3051 : if (fd_ >= 0)
298 : {
299 1045 : ::close(fd_);
300 1045 : fd_ = -1;
301 : }
302 3051 : }
303 :
304 : inline std::uint64_t
305 15 : posix_stream_file::size() const
306 : {
307 : file_stat_t st;
308 15 : if (file_fstat(fd_, &st) < 0)
309 5 : throw_system_error(make_err(errno), "stream_file::size");
310 10 : return static_cast<std::uint64_t>(st.st_size);
311 : }
312 :
313 : inline std::error_code
314 12 : posix_stream_file::resize(std::uint64_t new_size) noexcept
315 : {
316 12 : if (new_size >
317 12 : static_cast<std::uint64_t>((std::numeric_limits<file_off_t>::max)()))
318 2 : return make_err(EOVERFLOW);
319 10 : if (file_ftruncate(fd_, static_cast<file_off_t>(new_size)) < 0)
320 7 : return make_err(errno);
321 3 : return {};
322 : }
323 :
324 : inline std::error_code
325 8 : posix_stream_file::sync_data() noexcept
326 : {
327 : #if BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
328 8 : if (::fdatasync(fd_) < 0)
329 : #else // BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
330 : if (::fsync(fd_) < 0)
331 : #endif // BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
332 5 : return make_err(errno);
333 3 : return {};
334 : }
335 :
336 : inline std::error_code
337 8 : posix_stream_file::sync_all() noexcept
338 : {
339 8 : if (::fsync(fd_) < 0)
340 5 : return make_err(errno);
341 3 : return {};
342 : }
343 :
344 : inline native_handle_type
345 4 : posix_stream_file::release()
346 : {
347 : // A queued op has already copied the fd number; it must not run
348 : // after the caller closes it and the number is recycled.
349 4 : cancel();
350 4 : bump_generation();
351 4 : int fd = fd_;
352 4 : fd_ = -1;
353 4 : offset_ = 0;
354 4 : return fd;
355 : }
356 :
357 : inline std::error_code
358 411 : posix_stream_file::assign(native_handle_type handle) noexcept
359 : {
360 : // The public assign() guarantees the object is closed.
361 411 : if (auto ec = validate_file_fd(handle))
362 4 : return ec;
363 :
364 407 : fd_ = handle;
365 407 : offset_ = 0;
366 407 : return {};
367 : }
368 :
369 : inline capy::io_result<std::uint64_t>
370 435 : posix_stream_file::seek(
371 : std::int64_t offset, file_base::seek_basis origin) noexcept
372 : {
373 : // We track offset_ ourselves (not the kernel fd offset)
374 : // because preadv/pwritev use explicit offsets.
375 : std::int64_t new_pos;
376 :
377 435 : if (origin == file_base::seek_set)
378 : {
379 14 : new_pos = offset;
380 : }
381 421 : else if (origin == file_base::seek_cur)
382 : {
383 410 : new_pos = static_cast<std::int64_t>(offset_) + offset;
384 : }
385 : else
386 : {
387 : file_stat_t st;
388 11 : if (file_fstat(fd_, &st) < 0)
389 5 : return {make_err(errno), 0};
390 6 : new_pos = st.st_size + offset;
391 : }
392 :
393 430 : if (new_pos < 0)
394 6 : return {make_err(EINVAL), 0};
395 424 : if (new_pos >
396 424 : static_cast<std::int64_t>((std::numeric_limits<file_off_t>::max)()))
397 MIS 0 : return {make_err(EOVERFLOW), 0};
398 :
399 HIT 424 : offset_ = static_cast<std::uint64_t>(new_pos);
400 :
401 424 : return {std::error_code{}, offset_};
402 : }
403 :
404 : // -- file_op completion handler --
405 : // (read_some, write_some, do_read_work, do_write_work are
406 : // defined in posix_stream_file_service.hpp after the service)
407 :
408 : inline void
409 624 : posix_stream_file::file_op::operator()()
410 : {
411 624 : stop_cb.reset();
412 :
413 : // Empty buffers never reach the pool (diverted at initiation), so
414 : // empty_buffer stays false and a 0-byte read is a genuine EOF.
415 1234 : decode_io_result(
416 624 : ec_out, bytes_out, cancelled.load(std::memory_order_acquire),
417 624 : errn != 0 ? make_err(errn) : std::error_code{}, is_read,
418 : bytes_transferred, /*empty_buffer=*/false);
419 :
420 : // Move impl_ptr to a local so members remain valid through
421 : // dispatch — impl_ptr may be the last shared_ptr keeping
422 : // the parent posix_stream_file (which embeds this file_op) alive.
423 624 : auto prevent_destroy = std::move(impl_ptr);
424 624 : ex.on_work_finished();
425 624 : cont.h = h;
426 624 : dispatch_coro(ex, cont).resume();
427 624 : }
428 :
429 : inline void
430 2 : posix_stream_file::file_op::destroy()
431 : {
432 2 : stop_cb.reset();
433 2 : auto local_ex = ex;
434 2 : impl_ptr.reset();
435 2 : local_ex.on_work_finished();
436 2 : }
437 :
438 : } // namespace boost::corosio::detail
439 :
440 : #endif // BOOST_COROSIO_POSIX
441 :
442 : #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_STREAM_FILE_HPP
|