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

100.0% Lines (170 / 170) 100.0% Functions (14 / 14)
posix_stream_file_service.hpp
f(x) Functions (14)
Function Calls Lines Blocks
boost::corosio::detail::posix_stream_file_service::posix_stream_file_service(boost::capy::execution_context&) :35 621x 100.0% 89.0% boost::corosio::detail::posix_stream_file_service::~posix_stream_file_service() :41 1242x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::construct() :47 677x 100.0% 71.0% boost::corosio::detail::posix_stream_file_service::destroy(boost::corosio::io_object::implementation*) :61 675x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::close(boost::corosio::io_object::handle&) :69 1718x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::open_file(boost::corosio::stream_file::implementation&, std::filesystem::__cxx11::path const&, boost::corosio::file_base::flags) :79 658x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::shutdown() :91 621x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::destroy_impl(boost::corosio::detail::posix_stream_file&) :103 675x 100.0% 67.0% boost::corosio::detail::posix_stream_file_service::post(boost::corosio::detail::scheduler_op*) :110 626x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::pool() :138 634x 100.0% 100.0% boost::corosio::detail::posix_stream_file::read_some(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :157 559x 100.0% 95.0% boost::corosio::detail::posix_stream_file::do_read_work(boost::corosio::detail::pool_work_item*) :228 545x 100.0% 100.0% boost::corosio::detail::posix_stream_file::write_some(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :268 91x 100.0% 95.0% boost::corosio::detail::posix_stream_file::do_write_work(boost::corosio::detail::pool_work_item*) :339 81x 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_SERVICE_HPP
11 #define BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_STREAM_FILE_SERVICE_HPP
12
13 #include <boost/corosio/detail/platform.hpp>
14
15 #if BOOST_COROSIO_POSIX
16
17 #include <boost/corosio/native/detail/posix/posix_stream_file.hpp>
18 #include <boost/corosio/native/detail/reactor/reactor_scheduler.hpp>
19 #include <boost/corosio/detail/file_service.hpp>
20 #include <boost/corosio/detail/thread_pool.hpp>
21
22 #include <mutex>
23 #include <unordered_map>
24
25 namespace boost::corosio::detail {
26
27 /** Stream file service for POSIX backends.
28
29 Owns all posix_stream_file instances. Thread lifecycle is
30 managed by the thread_pool service (shared with resolver).
31 */
32 class BOOST_COROSIO_DECL posix_stream_file_service final : public file_service
33 {
34 public:
35 621x explicit posix_stream_file_service(capy::execution_context& ctx)
36 1863x : sched_(&get_scheduler(ctx))
37 621x , pool_(ctx)
38 {
39 621x }
40
41 1242x ~posix_stream_file_service() override = default;
42
43 posix_stream_file_service(posix_stream_file_service const&) = delete;
44 posix_stream_file_service&
45 operator=(posix_stream_file_service const&) = delete;
46
47 677x io_object::implementation* construct() override
48 {
49 677x auto ptr = std::make_shared<posix_stream_file>(*this);
50 677x auto* impl = ptr.get();
51
52 {
53 677x std::lock_guard<std::mutex> lock(mutex_);
54 677x file_list_.push_back(impl);
55 677x file_ptrs_[impl] = std::move(ptr);
56 677x }
57
58 677x return impl;
59 677x }
60
61 675x void destroy(io_object::implementation* p) override
62 {
63 675x auto& impl = static_cast<posix_stream_file&>(*p);
64 675x impl.cancel();
65 675x impl.close_file();
66 675x destroy_impl(impl);
67 675x }
68
69 1718x void close(io_object::handle& h) override
70 {
71 1718x if (h.get())
72 {
73 1718x auto& impl = static_cast<posix_stream_file&>(*h.get());
74 1718x impl.cancel();
75 1718x impl.close_file();
76 }
77 1718x }
78
79 658x std::error_code open_file(
80 stream_file::implementation& impl,
81 std::filesystem::path const& path,
82 file_base::flags mode) override
83 {
84 // Unavailable in the unsafe tier: the file thread pool completes
85 // cross-thread, which the lockless scheduler cannot accept.
86 658x if (sched_->scheduler_locking_disabled())
87 2x return std::make_error_code(std::errc::operation_not_supported);
88 656x return static_cast<posix_stream_file&>(impl).open_file(path, mode);
89 }
90
91 621x void shutdown() override
92 {
93 621x std::lock_guard<std::mutex> lock(mutex_);
94 623x for (auto* impl = file_list_.pop_front(); impl != nullptr;
95 2x impl = file_list_.pop_front())
96 {
97 2x impl->cancel();
98 2x impl->close_file();
99 }
100 621x file_ptrs_.clear();
101 621x }
102
103 675x void destroy_impl(posix_stream_file& impl)
104 {
105 675x std::lock_guard<std::mutex> lock(mutex_);
106 675x file_list_.remove(&impl);
107 675x file_ptrs_.erase(&impl);
108 675x }
109
110 626x void post(scheduler_op* op)
111 {
112 626x sched_->post(op);
113 626x }
114
115 void work_started() noexcept
116 {
117 sched_->work_started();
118 }
119
120 void work_finished() noexcept
121 {
122 sched_->work_finished();
123 }
124
125 /** Return the thread pool that runs this service's file work.
126
127 The pool's service is created on first use, so this can fail
128 where a plain accessor could not. Its workers start later, on
129 the first post, and a thread the system refuses there is
130 reported by that post rather than thrown here.
131
132 @throws std::bad_alloc If the service cannot be allocated.
133
134 @return The context's shared blocking-I/O pool.
135
136 @see thread_pool_ref::get
137 */
138 634x thread_pool& pool()
139 {
140 634x return pool_.get();
141 }
142
143 private:
144 scheduler* sched_;
145 thread_pool_ref pool_;
146 std::mutex mutex_;
147 intrusive_list<posix_stream_file> file_list_;
148 std::unordered_map<posix_stream_file*, std::shared_ptr<posix_stream_file>>
149 file_ptrs_;
150 };
151
152 // ---------------------------------------------------------------------------
153 // posix_stream_file inline implementations (require complete service type)
154 // ---------------------------------------------------------------------------
155
156 inline std::coroutine_handle<>
157 559x posix_stream_file::read_some(
158 std::coroutine_handle<> h,
159 capy::executor_ref ex,
160 buffer_param param,
161 std::stop_token token,
162 std::error_code* ec,
163 std::size_t* bytes_out)
164 {
165 559x auto& op = read_op_;
166 559x op.reset();
167 559x op.is_read = true;
168
169 // Closed-object contract outranks the zero-length no-op.
170 559x if (fd_ < 0)
171 {
172 6x *ec = make_error_code(std::errc::bad_file_descriptor);
173 6x *bytes_out = 0;
174 6x op.cont.h = h;
175 6x return dispatch_coro(ex, op.cont);
176 }
177
178 553x capy::mutable_buffer bufs[max_buffers];
179 553x op.iovec_count = static_cast<int>(param.copy_to(bufs, max_buffers));
180
181 553x if (op.iovec_count == 0)
182 {
183 2x *ec = {};
184 2x *bytes_out = 0;
185 2x op.cont.h = h;
186 2x return dispatch_coro(ex, op.cont);
187 }
188
189 1102x for (int i = 0; i < op.iovec_count; ++i)
190 {
191 551x op.iovecs[i].iov_base = bufs[i].data();
192 551x op.iovecs[i].iov_len = bufs[i].size();
193 }
194
195 551x op.h = h;
196 551x op.ex = ex;
197 551x op.ec_out = ec;
198 551x op.bytes_out = bytes_out;
199 551x op.start(token);
200
201 551x op.fd = fd_;
202 551x op.offset = offset_;
203 551x op.generation = generation_;
204
205 551x op.ex.on_work_started();
206
207 551x read_pool_op_.file_ = this;
208 551x read_pool_op_.ref_ = this->shared_from_this();
209 551x read_pool_op_.func_ = &posix_stream_file::do_read_work;
210 551x if (auto pec = svc_.pool().post(&read_pool_op_))
211 {
212 // The pool is shutting down, or the system refused it a thread.
213 // Nothing of this read went cross-thread, so it answers here
214 // like the closed-descriptor and zero-length exits above rather
215 // than through a completion the scheduler has to carry back.
216 6x read_pool_op_.ref_.reset();
217 6x op.stop_cb.reset();
218 6x op.ex.on_work_finished();
219 6x *ec = pec;
220 6x *bytes_out = 0;
221 6x op.cont.h = h;
222 6x return dispatch_coro(ex, op.cont);
223 }
224 545x return std::noop_coroutine();
225 }
226
227 inline void
228 545x posix_stream_file::do_read_work(pool_work_item* w) noexcept
229 {
230 545x auto* pw = static_cast<pool_op*>(w);
231 545x auto* self = pw->file_;
232 545x auto& op = self->read_op_;
233
234 545x if (!op.cancelled.load(std::memory_order_acquire))
235 {
236 ssize_t n;
237 do
238 {
239 127x n = file_preadv(
240 127x op.fd, op.iovecs, op.iovec_count,
241 127x static_cast<file_off_t>(op.offset));
242 }
243 127x while (n < 0 && errno == EINTR);
244
245 127x if (n >= 0)
246 {
247 120x op.errn = 0;
248 120x op.bytes_transferred = static_cast<std::size_t>(n);
249
250 // The file may have been replaced while this ran; its
251 // position belongs to whatever it now holds.
252 120x std::lock_guard<std::mutex> lock(self->offset_mutex_);
253 120x if (self->generation_ == op.generation)
254 115x self->offset_ += static_cast<std::uint64_t>(n);
255 120x }
256 else
257 {
258 7x op.errn = errno;
259 7x op.bytes_transferred = 0;
260 }
261 }
262
263 545x op.impl_ptr = std::move(pw->ref_);
264 545x self->svc_.post(&op);
265 545x }
266
267 inline std::coroutine_handle<>
268 91x posix_stream_file::write_some(
269 std::coroutine_handle<> h,
270 capy::executor_ref ex,
271 buffer_param param,
272 std::stop_token token,
273 std::error_code* ec,
274 std::size_t* bytes_out)
275 {
276 91x auto& op = write_op_;
277 91x op.reset();
278 91x op.is_read = false;
279
280 // Closed-object contract outranks the zero-length no-op.
281 91x if (fd_ < 0)
282 {
283 6x *ec = make_error_code(std::errc::bad_file_descriptor);
284 6x *bytes_out = 0;
285 6x op.cont.h = h;
286 6x return dispatch_coro(ex, op.cont);
287 }
288
289 85x capy::mutable_buffer bufs[max_buffers];
290 85x op.iovec_count = static_cast<int>(param.copy_to(bufs, max_buffers));
291
292 85x if (op.iovec_count == 0)
293 {
294 2x *ec = {};
295 2x *bytes_out = 0;
296 2x op.cont.h = h;
297 2x return dispatch_coro(ex, op.cont);
298 }
299
300 166x for (int i = 0; i < op.iovec_count; ++i)
301 {
302 83x op.iovecs[i].iov_base = bufs[i].data();
303 83x op.iovecs[i].iov_len = bufs[i].size();
304 }
305
306 83x op.h = h;
307 83x op.ex = ex;
308 83x op.ec_out = ec;
309 83x op.bytes_out = bytes_out;
310 83x op.start(token);
311
312 83x op.fd = fd_;
313 83x op.offset = offset_;
314 83x op.generation = generation_;
315
316 83x op.ex.on_work_started();
317
318 83x write_pool_op_.file_ = this;
319 83x write_pool_op_.ref_ = this->shared_from_this();
320 83x write_pool_op_.func_ = &posix_stream_file::do_write_work;
321 83x if (auto pec = svc_.pool().post(&write_pool_op_))
322 {
323 // The pool is shutting down, or the system refused it a thread.
324 // Nothing of this write went cross-thread, so it answers here
325 // like the closed-descriptor and zero-length exits above rather
326 // than through a completion the scheduler has to carry back.
327 2x write_pool_op_.ref_.reset();
328 2x op.stop_cb.reset();
329 2x op.ex.on_work_finished();
330 2x *ec = pec;
331 2x *bytes_out = 0;
332 2x op.cont.h = h;
333 2x return dispatch_coro(ex, op.cont);
334 }
335 81x return std::noop_coroutine();
336 }
337
338 inline void
339 81x posix_stream_file::do_write_work(pool_work_item* w) noexcept
340 {
341 81x auto* pw = static_cast<pool_op*>(w);
342 81x auto* self = pw->file_;
343 81x auto& op = self->write_op_;
344
345 81x if (!op.cancelled.load(std::memory_order_acquire))
346 {
347 ssize_t n;
348 do
349 {
350 29x n = file_pwritev(
351 29x op.fd, op.iovecs, op.iovec_count,
352 29x static_cast<file_off_t>(op.offset));
353 }
354 29x while (n < 0 && errno == EINTR);
355
356 29x if (n >= 0)
357 {
358 22x op.errn = 0;
359 22x op.bytes_transferred = static_cast<std::size_t>(n);
360
361 // The file may have been replaced while this ran; its
362 // position belongs to whatever it now holds.
363 22x std::lock_guard<std::mutex> lock(self->offset_mutex_);
364 22x if (self->generation_ == op.generation)
365 22x self->offset_ += static_cast<std::uint64_t>(n);
366 22x }
367 else
368 {
369 7x op.errn = errno;
370 7x op.bytes_transferred = 0;
371 }
372 }
373
374 81x op.impl_ptr = std::move(pw->ref_);
375 81x self->svc_.post(&op);
376 81x }
377
378 } // namespace boost::corosio::detail
379
380 #endif // BOOST_COROSIO_POSIX
381
382 #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_STREAM_FILE_SERVICE_HPP
383