include/boost/corosio/native/detail/posix/posix_stream_file.hpp

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