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/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 HIT 550 : file_op() = default;
111 :
112 284 : void reset() noexcept
113 : {
114 284 : iovec_count = 0;
115 284 : errn = 0;
116 284 : bytes_transferred = 0;
117 284 : is_read = false;
118 284 : cancelled.store(false, std::memory_order_relaxed);
119 284 : stop_cb.reset();
120 284 : object_ref_.reset();
121 284 : ec_out = nullptr;
122 284 : bytes_out = nullptr;
123 284 : }
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 64 : void reuse() noexcept
152 : {
153 64 : BOOST_COROSIO_ASSERT(fd_ == -1);
154 64 : BOOST_COROSIO_ASSERT(!read_op_.stop_cb);
155 64 : BOOST_COROSIO_ASSERT(!write_op_.stop_cb);
156 64 : }
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 941 : native_handle_type native_handle() const noexcept override
179 : {
180 941 : return fd_;
181 : }
182 :
183 988 : void cancel() noexcept override
184 : {
185 988 : read_op_.request_cancel();
186 988 : write_op_.request_cancel();
187 988 : }
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 275 : inline posix_stream_file::posix_stream_file(
226 275 : posix_stream_file_service& svc) noexcept
227 275 : : svc_(svc)
228 : {
229 275 : }
230 :
231 : inline std::error_code
232 320 : posix_stream_file::open_file(
233 : std::filesystem::path const& path, file_base::flags mode)
234 : {
235 320 : close_file();
236 :
237 320 : int oflags = 0;
238 :
239 : // Access mode
240 320 : unsigned access = static_cast<unsigned>(mode) & 3u;
241 320 : if (access == static_cast<unsigned>(file_base::read_write))
242 87 : oflags |= O_RDWR;
243 233 : else if (access == static_cast<unsigned>(file_base::write_only))
244 81 : oflags |= O_WRONLY;
245 : else
246 152 : oflags |= O_RDONLY;
247 :
248 : // Creation flags
249 320 : if ((mode & file_base::create) != file_base::flags(0))
250 106 : oflags |= O_CREAT;
251 320 : if ((mode & file_base::exclusive) != file_base::flags(0))
252 2 : oflags |= O_EXCL;
253 320 : if ((mode & file_base::truncate) != file_base::flags(0))
254 83 : oflags |= O_TRUNC;
255 320 : if ((mode & file_base::append) != file_base::flags(0))
256 8 : oflags |= O_APPEND;
257 320 : if ((mode & file_base::sync_all_on_write) != file_base::flags(0))
258 2 : oflags |= O_SYNC;
259 :
260 320 : int fd = ::open(path.c_str(), oflags, 0666);
261 320 : if (fd < 0)
262 9 : return make_err(errno);
263 :
264 311 : fd_ = fd;
265 311 : 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 311 : if ((mode & file_base::append) != file_base::flags(0))
270 : {
271 : struct stat st;
272 8 : if (::fstat(fd, &st) < 0)
273 : {
274 5 : int err = errno;
275 5 : ::close(fd);
276 5 : fd_ = -1;
277 5 : return make_err(err);
278 : }
279 3 : offset_ = static_cast<std::uint64_t>(st.st_size);
280 : }
281 :
282 : #ifdef POSIX_FADV_SEQUENTIAL
283 306 : ::posix_fadvise(fd_, 0, 0, POSIX_FADV_SEQUENTIAL);
284 : #endif
285 :
286 306 : return {};
287 : }
288 :
289 : inline void
290 1312 : posix_stream_file::close_file() noexcept
291 : {
292 1312 : if (fd_ >= 0)
293 : {
294 310 : ::close(fd_);
295 310 : fd_ = -1;
296 : }
297 1312 : }
298 :
299 : inline std::uint64_t
300 17 : posix_stream_file::size() const
301 : {
302 : struct stat st;
303 17 : if (::fstat(fd_, &st) < 0)
304 5 : throw_system_error(make_err(errno), "stream_file::size");
305 12 : return static_cast<std::uint64_t>(st.st_size);
306 : }
307 :
308 : inline std::error_code
309 12 : posix_stream_file::resize(std::uint64_t new_size) noexcept
310 : {
311 12 : if (new_size >
312 12 : static_cast<std::uint64_t>((std::numeric_limits<off_t>::max)()))
313 2 : return make_err(EOVERFLOW);
314 10 : if (::ftruncate(fd_, static_cast<off_t>(new_size)) < 0)
315 7 : return make_err(errno);
316 3 : return {};
317 : }
318 :
319 : inline std::error_code
320 10 : posix_stream_file::sync_data() noexcept
321 : {
322 : #if BOOST_COROSIO_HAS_POSIX_SYNCHRONIZED_IO
323 10 : 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 7 : return make_err(errno);
328 3 : return {};
329 : }
330 :
331 : inline std::error_code
332 10 : posix_stream_file::sync_all() noexcept
333 : {
334 10 : if (::fsync(fd_) < 0)
335 7 : return make_err(errno);
336 3 : return {};
337 : }
338 :
339 : inline native_handle_type
340 2 : posix_stream_file::release()
341 : {
342 2 : int fd = fd_;
343 2 : fd_ = -1;
344 2 : offset_ = 0;
345 2 : return fd;
346 : }
347 :
348 : inline std::error_code
349 6 : posix_stream_file::assign(native_handle_type handle) noexcept
350 : {
351 6 : close_file();
352 6 : fd_ = handle;
353 6 : offset_ = 0;
354 6 : return {};
355 : }
356 :
357 : inline capy::io_result<std::uint64_t>
358 30 : 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 30 : if (origin == file_base::seek_set)
366 : {
367 14 : new_pos = offset;
368 : }
369 16 : else if (origin == file_base::seek_cur)
370 : {
371 5 : new_pos = static_cast<std::int64_t>(offset_) + offset;
372 : }
373 : else
374 : {
375 : struct stat st;
376 11 : if (::fstat(fd_, &st) < 0)
377 5 : return {make_err(errno), 0};
378 6 : new_pos = st.st_size + offset;
379 : }
380 :
381 25 : if (new_pos < 0)
382 6 : return {make_err(EINVAL), 0};
383 19 : if (new_pos >
384 19 : static_cast<std::int64_t>((std::numeric_limits<off_t>::max)()))
385 MIS 0 : return {make_err(EOVERFLOW), 0};
386 :
387 HIT 19 : offset_ = static_cast<std::uint64_t>(new_pos);
388 :
389 19 : 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 258 : posix_stream_file::file_op::operator()()
398 : {
399 258 : 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 502 : decode_io_result(
404 258 : ec_out, bytes_out, cancelled.load(std::memory_order_acquire),
405 258 : 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 258 : auto prevent_destroy = std::move(object_ref_);
412 258 : ex.on_work_finished();
413 258 : cont.h = h;
414 258 : dispatch_coro(ex, cont).resume();
415 258 : }
416 :
417 : inline void
418 2 : posix_stream_file::file_op::destroy()
419 : {
420 2 : stop_cb.reset();
421 2 : auto local_ex = ex;
422 2 : object_ref_.reset();
423 2 : local_ex.on_work_finished();
424 2 : }
425 :
426 : } // namespace boost::corosio::detail
427 :
428 : #endif // BOOST_COROSIO_POSIX
429 :
430 : #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_STREAM_FILE_HPP
|