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

99.2% Lines (124 / 125) 100.0% Functions (17 / 17)
posix_stream_file.hpp
f(x) Functions (17)
Function Calls Lines Blocks
boost::corosio::detail::posix_stream_file::file_op::file_op() :110 550x 100.0% 100.0% boost::corosio::detail::posix_stream_file::file_op::reset() :112 284x 100.0% 100.0% boost::corosio::detail::posix_stream_file::reuse() :151 64x 100.0% 67.0% boost::corosio::detail::posix_stream_file::native_handle() const :178 941x 100.0% 100.0% boost::corosio::detail::posix_stream_file::cancel() :183 988x 100.0% 100.0% boost::corosio::detail::posix_stream_file::posix_stream_file(boost::corosio::detail::posix_stream_file_service&) :225 275x 100.0% 100.0% boost::corosio::detail::posix_stream_file::open_file(std::filesystem::__cxx11::path const&, boost::corosio::file_base::flags) :232 320x 100.0% 100.0% boost::corosio::detail::posix_stream_file::close_file() :290 1312x 100.0% 100.0% boost::corosio::detail::posix_stream_file::size() const :300 17x 100.0% 100.0% boost::corosio::detail::posix_stream_file::resize(unsigned long) :309 12x 100.0% 100.0% boost::corosio::detail::posix_stream_file::sync_data() :320 10x 100.0% 100.0% boost::corosio::detail::posix_stream_file::sync_all() :332 10x 100.0% 100.0% boost::corosio::detail::posix_stream_file::release() :340 2x 100.0% 100.0% boost::corosio::detail::posix_stream_file::assign(int) :349 6x 100.0% 100.0% boost::corosio::detail::posix_stream_file::seek(long, boost::corosio::file_base::seek_basis) :358 30x 93.3% 87.0% boost::corosio::detail::posix_stream_file::file_op::operator()() :397 258x 100.0% 82.0% boost::corosio::detail::posix_stream_file::file_op::destroy() :418 2x 100.0% 100.0%
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/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 <optional>
41 #include <stop_token>
42 #include <system_error>
43
44 #include <errno.h>
45 #include <fcntl.h>
46 #include <sys/stat.h>
47 #include <sys/uio.h>
48 #include <unistd.h>
49
50 /*
51 POSIX Stream File Implementation
52 =================================
53
54 Regular files cannot be monitored by epoll/kqueue/select — the kernel
55 always reports them as ready. Blocking I/O (pread/pwrite) is dispatched
56 to a shared thread pool, with completion posted back to the scheduler.
57
58 This follows the same pattern as posix_resolver: pool_work_item for
59 dispatch, scheduler_op for completion, an object_ref keepalive on
60 each op for lifetime while its pool work is in flight.
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 intrusive_list<posix_stream_file>::node
89 {
90 friend class posix_stream_file_service;
91
92 public:
93 static constexpr std::size_t max_buffers = 16;
94
95 /** Operation state for a single file read or write.
96
97 The coroutine, cancellation and keepalive machinery is inherited
98 from `coro_op`; only the pool-path result state lives here.
99 */
100 struct file_op : coro_op
101 {
102 // Buffer data (copied from buffer_param at submission time)
103 iovec iovecs[max_buffers];
104 int iovec_count = 0;
105
106 // Result storage (populated by worker thread)
107 int errn = 0;
108 std::size_t bytes_transferred = 0;
109
110 550x file_op() = default;
111
112 284x void reset() noexcept
113 {
114 284x iovec_count = 0;
115 284x errn = 0;
116 284x bytes_transferred = 0;
117 284x is_read = false;
118 284x cancelled.store(false, std::memory_order_relaxed);
119 284x stop_cb.reset();
120 284x object_ref_.reset();
121 284x ec_out = nullptr;
122 284x bytes_out = nullptr;
123 284x }
124
125 void operator()() override;
126 void destroy() override;
127 };
128
129 /** Pool work item for thread pool dispatch. */
130 struct pool_op : pool_work_item
131 {
132 posix_stream_file* file_ = nullptr;
133 detail::object_ref ref_;
134 };
135
136 explicit posix_stream_file(posix_stream_file_service& svc) noexcept;
137
138 /// Recycle into the owning service's pool. Defined out-of-line
139 /// after posix_stream_file_service for its complete type.
140 void retire() noexcept override;
141
142 /** Reset op slots for recycling.
143
144 `close_file()` already drove fd_ to its closed value before
145 the refcount reached zero; each op's own `reset()` runs again
146 before its next use, so only the `stop_cb` invariant is worth
147 asserting here.
148
149 @pre refs_ == 0, fd closed, no op in flight.
150 */
151 64x void reuse() noexcept
152 {
153 64x BOOST_COROSIO_ASSERT(fd_ == -1);
154 64x BOOST_COROSIO_ASSERT(!read_op_.stop_cb);
155 64x BOOST_COROSIO_ASSERT(!write_op_.stop_cb);
156 64x }
157
158 // -- io_stream::implementation --
159
160 std::coroutine_handle<> read_some(
161 std::coroutine_handle<>,
162 capy::executor_ref,
163 buffer_param,
164 std::stop_token,
165 std::error_code*,
166 std::size_t*) override;
167
168 std::coroutine_handle<> write_some(
169 std::coroutine_handle<>,
170 capy::executor_ref,
171 buffer_param,
172 std::stop_token,
173 std::error_code*,
174 std::size_t*) override;
175
176 // -- stream_file::implementation --
177
178 941x native_handle_type native_handle() const noexcept override
179 {
180 941x return fd_;
181 }
182
183 988x void cancel() noexcept override
184 {
185 988x read_op_.request_cancel();
186 988x write_op_.request_cancel();
187 988x }
188
189 std::uint64_t size() const override;
190 std::error_code resize(std::uint64_t new_size) noexcept override;
191 std::error_code sync_data() noexcept override;
192 std::error_code sync_all() noexcept override;
193 native_handle_type release() override;
194 std::error_code assign(native_handle_type handle) noexcept override;
195 capy::io_result<std::uint64_t>
196 seek(std::int64_t offset, file_base::seek_basis origin) noexcept override;
197
198 // -- Internal --
199
200 /** Open the file and store the fd. */
201 std::error_code
202 open_file(std::filesystem::path const& path, file_base::flags mode);
203
204 /** Close the file descriptor. */
205 void close_file() noexcept;
206
207 private:
208 posix_stream_file_service& svc_;
209 int fd_ = -1;
210 std::uint64_t offset_ = 0;
211
212 file_op read_op_;
213 file_op write_op_;
214 pool_op read_pool_op_;
215 pool_op write_pool_op_;
216
217 static void do_read_work(pool_work_item*) noexcept;
218 static void do_write_work(pool_work_item*) noexcept;
219 };
220
221 // ---------------------------------------------------------------------------
222 // Inline implementation
223 // ---------------------------------------------------------------------------
224
225 275x inline posix_stream_file::posix_stream_file(
226 275x posix_stream_file_service& svc) noexcept
227 275x : svc_(svc)
228 {
229 275x }
230
231 inline std::error_code
232 320x posix_stream_file::open_file(
233 std::filesystem::path const& path, file_base::flags mode)
234 {
235 320x close_file();
236
237 320x int oflags = 0;
238
239 // Access mode
240 320x unsigned access = static_cast<unsigned>(mode) & 3u;
241 320x if (access == static_cast<unsigned>(file_base::read_write))
242 87x oflags |= O_RDWR;
243 233x else if (access == static_cast<unsigned>(file_base::write_only))
244 81x oflags |= O_WRONLY;
245 else
246 152x oflags |= O_RDONLY;
247
248 // Creation flags
249 320x if ((mode & file_base::create) != file_base::flags(0))
250 106x oflags |= O_CREAT;
251 320x if ((mode & file_base::exclusive) != file_base::flags(0))
252 2x oflags |= O_EXCL;
253 320x if ((mode & file_base::truncate) != file_base::flags(0))
254 83x oflags |= O_TRUNC;
255 320x if ((mode & file_base::append) != file_base::flags(0))
256 8x oflags |= O_APPEND;
257 320x if ((mode & file_base::sync_all_on_write) != file_base::flags(0))
258 2x oflags |= O_SYNC;
259
260 320x int fd = ::open(path.c_str(), oflags, 0666);
261 320x if (fd < 0)
262 9x return make_err(errno);
263
264 311x fd_ = fd;
265 311x offset_ = 0;
266
267 // Append mode: position at end-of-file (preadv/pwritev use
268 // explicit offsets, so O_APPEND alone is not sufficient).
269 311x if ((mode & file_base::append) != file_base::flags(0))
270 {
271 struct stat st;
272 8x if (::fstat(fd, &st) < 0)
273 {
274 5x int err = errno;
275 5x ::close(fd);
276 5x fd_ = -1;
277 5x return make_err(err);
278 }
279 3x offset_ = static_cast<std::uint64_t>(st.st_size);
280 }
281
282 #ifdef POSIX_FADV_SEQUENTIAL
283 306x ::posix_fadvise(fd_, 0, 0, POSIX_FADV_SEQUENTIAL);
284 #endif
285
286 306x return {};
287 }
288
289 inline void
290 1312x posix_stream_file::close_file() noexcept
291 {
292 1312x if (fd_ >= 0)
293 {
294 310x ::close(fd_);
295 310x fd_ = -1;
296 }
297 1312x }
298
299 inline std::uint64_t
300 17x posix_stream_file::size() const
301 {
302 struct stat st;
303 17x if (::fstat(fd_, &st) < 0)
304 5x throw_system_error(make_err(errno), "stream_file::size");
305 12x return static_cast<std::uint64_t>(st.st_size);
306 }
307
308 inline std::error_code
309 12x posix_stream_file::resize(std::uint64_t new_size) noexcept
310 {
311 12x if (new_size >
312 12x static_cast<std::uint64_t>((std::numeric_limits<off_t>::max)()))
313 2x return make_err(EOVERFLOW);
314 10x if (::ftruncate(fd_, static_cast<off_t>(new_size)) < 0)
315 7x return make_err(errno);
316 3x return {};
317 }
318
319 inline std::error_code
320 10x posix_stream_file::sync_data() noexcept
321 {
322 #if BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
323 10x if (::fdatasync(fd_) < 0)
324 #else // BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
325 if (::fsync(fd_) < 0)
326 #endif // BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
327 7x return make_err(errno);
328 3x return {};
329 }
330
331 inline std::error_code
332 10x posix_stream_file::sync_all() noexcept
333 {
334 10x if (::fsync(fd_) < 0)
335 7x return make_err(errno);
336 3x return {};
337 }
338
339 inline native_handle_type
340 2x posix_stream_file::release()
341 {
342 2x int fd = fd_;
343 2x fd_ = -1;
344 2x offset_ = 0;
345 2x return fd;
346 }
347
348 inline std::error_code
349 6x posix_stream_file::assign(native_handle_type handle) noexcept
350 {
351 6x close_file();
352 6x fd_ = handle;
353 6x offset_ = 0;
354 6x return {};
355 }
356
357 inline capy::io_result<std::uint64_t>
358 30x posix_stream_file::seek(
359 std::int64_t offset, file_base::seek_basis origin) noexcept
360 {
361 // We track offset_ ourselves (not the kernel fd offset)
362 // because preadv/pwritev use explicit offsets.
363 std::int64_t new_pos;
364
365 30x if (origin == file_base::seek_set)
366 {
367 14x new_pos = offset;
368 }
369 16x else if (origin == file_base::seek_cur)
370 {
371 5x new_pos = static_cast<std::int64_t>(offset_) + offset;
372 }
373 else
374 {
375 struct stat st;
376 11x if (::fstat(fd_, &st) < 0)
377 5x return {make_err(errno), 0};
378 6x new_pos = st.st_size + offset;
379 }
380
381 25x if (new_pos < 0)
382 6x return {make_err(EINVAL), 0};
383 19x if (new_pos >
384 19x static_cast<std::int64_t>((std::numeric_limits<off_t>::max)()))
385 ✗ return {make_err(EOVERFLOW), 0};
386
387 19x offset_ = static_cast<std::uint64_t>(new_pos);
388
389 19x return {std::error_code{}, offset_};
390 }
391
392 // -- file_op completion handler --
393 // (read_some, write_some, do_read_work, do_write_work are
394 // defined in posix_stream_file_service.hpp after the service)
395
396 inline void
397 258x posix_stream_file::file_op::operator()()
398 {
399 258x stop_cb.reset();
400
401 // Empty buffers never reach the pool (diverted at initiation), so
402 // empty_buffer stays false and a 0-byte read is a genuine EOF.
403 502x decode_io_result(
404 258x ec_out, bytes_out, cancelled.load(std::memory_order_acquire),
405 258x errn != 0 ? make_err(errn) : std::error_code{}, is_read,
406 bytes_transferred, /*empty_buffer=*/false);
407
408 // Move object_ref_ to a local so members remain valid through
409 // dispatch — this may be the last reference keeping the parent
410 // posix_stream_file (which embeds this file_op) alive.
411 258x auto prevent_destroy = std::move(object_ref_);
412 258x ex.on_work_finished();
413 258x cont.h = h;
414 258x dispatch_coro(ex, cont).resume();
415 258x }
416
417 inline void
418 2x posix_stream_file::file_op::destroy()
419 {
420 2x stop_cb.reset();
421 2x auto local_ex = ex;
422 2x object_ref_.reset();
423 2x local_ex.on_work_finished();
424 2x }
425
426 } // namespace boost::corosio::detail
427
428 #endif // BOOST_COROSIO_POSIX
429
430 #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_STREAM_FILE_HPP
431