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/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 HIT 544 : file_op() = default;
112 :
113 241 : void reset() noexcept
114 : {
115 241 : iovec_count = 0;
116 241 : errn = 0;
117 241 : bytes_transferred = 0;
118 241 : is_read = false;
119 241 : cancelled.store(false, std::memory_order_relaxed);
120 241 : stop_cb.reset();
121 241 : impl_ptr.reset();
122 241 : ec_out = nullptr;
123 241 : bytes_out = nullptr;
124 241 : }
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 815 : native_handle_type native_handle() const noexcept override
160 : {
161 815 : return fd_;
162 : }
163 :
164 787 : void cancel() noexcept override
165 : {
166 787 : read_op_.request_cancel();
167 787 : write_op_.request_cancel();
168 787 : }
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 272 : inline posix_stream_file::posix_stream_file(
207 272 : posix_stream_file_service& svc) noexcept
208 272 : : svc_(svc)
209 : {
210 272 : }
211 :
212 : inline std::error_code
213 253 : posix_stream_file::open_file(
214 : std::filesystem::path const& path, file_base::flags mode)
215 : {
216 253 : close_file();
217 :
218 253 : int oflags = 0;
219 :
220 : // Access mode
221 253 : unsigned access = static_cast<unsigned>(mode) & 3u;
222 253 : if (access == static_cast<unsigned>(file_base::read_write))
223 25 : oflags |= O_RDWR;
224 228 : else if (access == static_cast<unsigned>(file_base::write_only))
225 81 : oflags |= O_WRONLY;
226 : else
227 147 : oflags |= O_RDONLY;
228 :
229 : // Creation flags
230 253 : if ((mode & file_base::create) != file_base::flags(0))
231 44 : oflags |= O_CREAT;
232 253 : if ((mode & file_base::exclusive) != file_base::flags(0))
233 2 : oflags |= O_EXCL;
234 253 : if ((mode & file_base::truncate) != file_base::flags(0))
235 17 : oflags |= O_TRUNC;
236 253 : if ((mode & file_base::append) != file_base::flags(0))
237 8 : oflags |= O_APPEND;
238 253 : if ((mode & file_base::sync_all_on_write) != file_base::flags(0))
239 2 : oflags |= O_SYNC;
240 :
241 253 : int fd = ::open(path.c_str(), oflags, 0666);
242 253 : if (fd < 0)
243 9 : return make_err(errno);
244 :
245 244 : fd_ = fd;
246 244 : 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 244 : if ((mode & file_base::append) != file_base::flags(0))
251 : {
252 : struct stat st;
253 8 : if (::fstat(fd, &st) < 0)
254 : {
255 5 : int err = errno;
256 5 : ::close(fd);
257 5 : fd_ = -1;
258 5 : return make_err(err);
259 : }
260 3 : offset_ = static_cast<std::uint64_t>(st.st_size);
261 : }
262 :
263 : #ifdef POSIX_FADV_SEQUENTIAL
264 239 : ::posix_fadvise(fd_, 0, 0, POSIX_FADV_SEQUENTIAL);
265 : #endif
266 :
267 239 : return {};
268 : }
269 :
270 : inline void
271 1038 : posix_stream_file::close_file() noexcept
272 : {
273 1038 : if (fd_ >= 0)
274 : {
275 243 : ::close(fd_);
276 243 : fd_ = -1;
277 : }
278 1038 : }
279 :
280 : inline std::uint64_t
281 17 : posix_stream_file::size() const
282 : {
283 : struct stat st;
284 17 : if (::fstat(fd_, &st) < 0)
285 5 : throw_system_error(make_err(errno), "stream_file::size");
286 12 : return static_cast<std::uint64_t>(st.st_size);
287 : }
288 :
289 : inline std::error_code
290 12 : posix_stream_file::resize(std::uint64_t new_size) noexcept
291 : {
292 12 : if (new_size >
293 12 : static_cast<std::uint64_t>((std::numeric_limits<off_t>::max)()))
294 2 : return make_err(EOVERFLOW);
295 10 : if (::ftruncate(fd_, static_cast<off_t>(new_size)) < 0)
296 7 : return make_err(errno);
297 3 : return {};
298 : }
299 :
300 : inline std::error_code
301 8 : posix_stream_file::sync_data() noexcept
302 : {
303 : #if BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
304 8 : 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 5 : return make_err(errno);
309 3 : return {};
310 : }
311 :
312 : inline std::error_code
313 8 : posix_stream_file::sync_all() noexcept
314 : {
315 8 : if (::fsync(fd_) < 0)
316 5 : return make_err(errno);
317 3 : return {};
318 : }
319 :
320 : inline native_handle_type
321 2 : posix_stream_file::release()
322 : {
323 2 : int fd = fd_;
324 2 : fd_ = -1;
325 2 : offset_ = 0;
326 2 : return fd;
327 : }
328 :
329 : inline std::error_code
330 12 : 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 12 : if (handle >= 0 && handle == fd_)
335 2 : 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 10 : if (auto ec = validate_file_fd(handle))
340 4 : 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 6 : cancel();
349 6 : close_file();
350 6 : fd_ = handle;
351 6 : offset_ = 0;
352 6 : return {};
353 : }
354 :
355 : inline capy::io_result<std::uint64_t>
356 30 : 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 30 : if (origin == file_base::seek_set)
364 : {
365 14 : new_pos = offset;
366 : }
367 16 : else if (origin == file_base::seek_cur)
368 : {
369 5 : new_pos = static_cast<std::int64_t>(offset_) + offset;
370 : }
371 : else
372 : {
373 : struct stat st;
374 11 : if (::fstat(fd_, &st) < 0)
375 5 : return {make_err(errno), 0};
376 6 : new_pos = st.st_size + offset;
377 : }
378 :
379 25 : if (new_pos < 0)
380 6 : return {make_err(EINVAL), 0};
381 19 : if (new_pos >
382 19 : static_cast<std::int64_t>((std::numeric_limits<off_t>::max)()))
383 MIS 0 : return {make_err(EOVERFLOW), 0};
384 :
385 HIT 19 : offset_ = static_cast<std::uint64_t>(new_pos);
386 :
387 19 : 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 215 : posix_stream_file::file_op::operator()()
396 : {
397 215 : 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 416 : decode_io_result(
402 215 : ec_out, bytes_out, cancelled.load(std::memory_order_acquire),
403 215 : 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 215 : auto prevent_destroy = std::move(impl_ptr);
410 215 : ex.on_work_finished();
411 215 : cont.h = h;
412 215 : dispatch_coro(ex, cont).resume();
413 215 : }
414 :
415 : inline void
416 2 : posix_stream_file::file_op::destroy()
417 : {
418 2 : stop_cb.reset();
419 2 : auto local_ex = ex;
420 2 : impl_ptr.reset();
421 2 : local_ex.on_work_finished();
422 2 : }
423 :
424 : } // namespace boost::corosio::detail
425 :
426 : #endif // BOOST_COROSIO_POSIX
427 :
428 : #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_STREAM_FILE_HPP
|