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

100.0% Lines (184/184) 100.0% List of functions (14/14) 76.6% Branches (49/64)
posix_stream_file_service.hpp
f(x) Functions (14)
Function Calls Lines Branches Blocks
boost::corosio::detail::posix_stream_file_service::posix_stream_file_service(boost::capy::execution_context&) :35 1158x 100.0% 50.0% 100.0% boost::corosio::detail::posix_stream_file_service::~posix_stream_file_service() :41 1737x 100.0% – 100.0% boost::corosio::detail::posix_stream_file_service::construct() :47 635x 100.0% 50.0% 42.0% boost::corosio::detail::posix_stream_file_service::destroy(boost::corosio::io_object::implementation*) :61 633x 100.0% – 100.0% boost::corosio::detail::posix_stream_file_service::close(boost::corosio::io_object::handle&) :69 1634x 100.0% 50.0% 100.0% boost::corosio::detail::posix_stream_file_service::open_file(boost::corosio::stream_file::implementation&, std::__1::__fs::filesystem::path const&, boost::corosio::file_base::flags) :79 613x 100.0% 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::shutdown() :91 579x 100.0% 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::destroy_impl(boost::corosio::detail::posix_stream_file&) :103 633x 100.0% 50.0% 50.0% boost::corosio::detail::posix_stream_file_service::post(boost::corosio::detail::scheduler_op*) :110 611x 100.0% – 100.0% boost::corosio::detail::posix_stream_file_service::pool() :138 620x 100.0% – 100.0% boost::corosio::detail::posix_stream_file::read_some(std::__1::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::__1::stop_token, std::__1::error_code*, unsigned long*) :157 551x 100.0% 100.0% 93.0% boost::corosio::detail::posix_stream_file::do_read_work(boost::corosio::detail::pool_work_item*) :228 536x 100.0% 72.2% 94.0% boost::corosio::detail::posix_stream_file::write_some(std::__1::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::__1::stop_token, std::__1::error_code*, unsigned long*) :268 85x 100.0% 100.0% 93.0% boost::corosio::detail::posix_stream_file::do_write_work(boost::corosio::detail::pool_work_item*) :339 75x 100.0% 68.8% 94.0%
Line Branch 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 1158x explicit posix_stream_file_service(capy::execution_context& ctx)
36
1/2
✓ Branch 0 taken 579 times.
✗ Branch 1 not taken.
579x : sched_(&get_scheduler(ctx))
37 579x , pool_(ctx)
38 1158x {
39 1158x }
40
41 1737x ~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 635x io_object::implementation* construct() override
48 {
49 635x auto ptr = std::make_shared<posix_stream_file>(*this);
50 635x auto* impl = ptr.get();
51
52 {
53
1/2
✓ Branch 0 taken 635 times.
✗ Branch 1 not taken.
635x std::lock_guard<std::mutex> lock(mutex_);
54 635x file_list_.push_back(impl);
55
1/2
✓ Branch 0 taken 635 times.
✗ Branch 1 not taken.
635x file_ptrs_[impl] = std::move(ptr);
56 635x }
57
58 635x return impl;
59 635x }
60
61 633x void destroy(io_object::implementation* p) override
62 {
63 633x auto& impl = static_cast<posix_stream_file&>(*p);
64 633x impl.cancel();
65 633x impl.close_file();
66 633x destroy_impl(impl);
67 633x }
68
69 1634x void close(io_object::handle& h) override
70 {
71
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 1634 times.
1634x if (h.get())
72 {
73 1634x auto& impl = static_cast<posix_stream_file&>(*h.get());
74 1634x impl.cancel();
75 1634x impl.close_file();
76 1634x }
77 1634x }
78
79 613x 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
2/2
✓ Branch 0 taken 2 times.
✓ Branch 1 taken 611 times.
613x if (sched_->scheduler_locking_disabled())
87 2x return std::make_error_code(std::errc::operation_not_supported);
88 611x return static_cast<posix_stream_file&>(impl).open_file(path, mode);
89 613x }
90
91 579x void shutdown() override
92 {
93 579x std::lock_guard<std::mutex> lock(mutex_);
94
2/2
✓ Branch 0 taken 2 times.
✓ Branch 1 taken 579 times.
581x 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 2x }
100 579x file_ptrs_.clear();
101 579x }
102
103 633x void destroy_impl(posix_stream_file& impl)
104 {
105 633x std::lock_guard<std::mutex> lock(mutex_);
106 633x file_list_.remove(&impl);
107
1/2
✓ Branch 0 taken 633 times.
✗ Branch 1 not taken.
633x file_ptrs_.erase(&impl);
108 633x }
109
110 611x void post(scheduler_op* op)
111 {
112 611x sched_->post(op);
113 611x }
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 620x thread_pool& pool()
139 {
140 620x 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 551x 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 551x auto& op = read_op_;
166 551x op.reset();
167 551x op.is_read = true;
168
169 // Closed-object contract outranks the zero-length no-op.
170
2/2
✓ Branch 0 taken 545 times.
✓ Branch 1 taken 6 times.
551x 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 545x capy::mutable_buffer bufs[max_buffers];
179 545x op.iovec_count = static_cast<int>(param.copy_to(bufs, max_buffers));
180
181
2/2
✓ Branch 0 taken 543 times.
✓ Branch 1 taken 2 times.
545x 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
2/2
✓ Branch 0 taken 543 times.
✓ Branch 1 taken 543 times.
1086x for (int i = 0; i < op.iovec_count; ++i)
190 {
191 543x op.iovecs[i].iov_base = bufs[i].data();
192 543x op.iovecs[i].iov_len = bufs[i].size();
193 543x }
194
195 543x op.h = h;
196 543x op.ex = ex;
197 543x op.ec_out = ec;
198 543x op.bytes_out = bytes_out;
199 543x op.start(token);
200
201 543x op.fd = fd_;
202 543x op.offset = offset_;
203 543x op.generation = generation_;
204
205 543x op.ex.on_work_started();
206
207 543x read_pool_op_.file_ = this;
208 543x read_pool_op_.ref_ = this->shared_from_this();
209 543x read_pool_op_.func_ = &posix_stream_file::do_read_work;
210
2/2
✓ Branch 0 taken 7 times.
✓ Branch 1 taken 536 times.
543x 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 7x read_pool_op_.ref_.reset();
217 7x op.stop_cb.reset();
218 7x op.ex.on_work_finished();
219 7x *ec = pec;
220 7x *bytes_out = 0;
221 7x op.cont.h = h;
222 7x return dispatch_coro(ex, op.cont);
223 }
224 536x return std::noop_coroutine();
225 551x }
226
227 inline void
228 536x posix_stream_file::do_read_work(pool_work_item* w) noexcept
229 {
230 536x auto* pw = static_cast<pool_op*>(w);
231 536x auto* self = pw->file_;
232 536x auto& op = self->read_op_;
233
234
2/2
✓ Branch 0 taken 327 times.
✓ Branch 1 taken 209 times.
536x if (!op.cancelled.load(std::memory_order_acquire))
235 {
236 ssize_t n;
237 209x do
238 {
239 209x n = file_preadv(
240 209x op.fd, op.iovecs, op.iovec_count,
241 209x static_cast<file_off_t>(op.offset));
242 418x }
243
3/4
✓ Branch 0 taken 197 times.
✓ Branch 1 taken 12 times.
✓ Branch 2 taken 12 times.
✗ Branch 3 not taken.
209x while (n < 0 && errno == EINTR);
244
245
2/2
✓ Branch 0 taken 197 times.
✓ Branch 1 taken 12 times.
209x if (n >= 0)
246 {
247 197x op.errn = 0;
248 197x 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
1/2
✓ Branch 0 taken 197 times.
✗ Branch 1 not taken.
197x std::lock_guard<std::mutex> lock(self->offset_mutex_);
253
2/2
✓ Branch 0 taken 60 times.
✓ Branch 1 taken 137 times.
197x if (self->generation_ == op.generation)
254 137x self->offset_ += static_cast<std::uint64_t>(n);
255 197x }
256 else
257 {
258
1/2
✓ Branch 0 taken 12 times.
✗ Branch 1 not taken.
12x op.errn = errno;
259 12x op.bytes_transferred = 0;
260 }
261 209x }
262
263
1/2
✓ Branch 0 taken 536 times.
✗ Branch 1 not taken.
536x op.impl_ptr = std::move(pw->ref_);
264
1/2
✓ Branch 0 taken 536 times.
✗ Branch 1 not taken.
536x self->svc_.post(&op);
265 536x }
266
267 inline std::coroutine_handle<>
268 85x 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 85x auto& op = write_op_;
277 85x op.reset();
278 85x op.is_read = false;
279
280 // Closed-object contract outranks the zero-length no-op.
281
2/2
✓ Branch 0 taken 79 times.
✓ Branch 1 taken 6 times.
85x 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 79x capy::mutable_buffer bufs[max_buffers];
290 79x op.iovec_count = static_cast<int>(param.copy_to(bufs, max_buffers));
291
292
2/2
✓ Branch 0 taken 77 times.
✓ Branch 1 taken 2 times.
79x 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
2/2
✓ Branch 0 taken 77 times.
✓ Branch 1 taken 77 times.
154x for (int i = 0; i < op.iovec_count; ++i)
301 {
302 77x op.iovecs[i].iov_base = bufs[i].data();
303 77x op.iovecs[i].iov_len = bufs[i].size();
304 77x }
305
306 77x op.h = h;
307 77x op.ex = ex;
308 77x op.ec_out = ec;
309 77x op.bytes_out = bytes_out;
310 77x op.start(token);
311
312 77x op.fd = fd_;
313 77x op.offset = offset_;
314 77x op.generation = generation_;
315
316 77x op.ex.on_work_started();
317
318 77x write_pool_op_.file_ = this;
319 77x write_pool_op_.ref_ = this->shared_from_this();
320 77x write_pool_op_.func_ = &posix_stream_file::do_write_work;
321
2/2
✓ Branch 0 taken 2 times.
✓ Branch 1 taken 75 times.
77x 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 75x return std::noop_coroutine();
336 85x }
337
338 inline void
339 75x posix_stream_file::do_write_work(pool_work_item* w) noexcept
340 {
341 75x auto* pw = static_cast<pool_op*>(w);
342 75x auto* self = pw->file_;
343 75x auto& op = self->write_op_;
344
345
2/2
✓ Branch 0 taken 52 times.
✓ Branch 1 taken 23 times.
75x if (!op.cancelled.load(std::memory_order_acquire))
346 {
347 ssize_t n;
348 23x do
349 {
350 23x n = file_pwritev(
351 23x op.fd, op.iovecs, op.iovec_count,
352 23x static_cast<file_off_t>(op.offset));
353 46x }
354
3/4
✓ Branch 0 taken 19 times.
✓ Branch 1 taken 4 times.
✓ Branch 2 taken 4 times.
✗ Branch 3 not taken.
23x while (n < 0 && errno == EINTR);
355
356
2/2
✓ Branch 0 taken 19 times.
✓ Branch 1 taken 4 times.
23x if (n >= 0)
357 {
358 19x op.errn = 0;
359 19x 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
1/2
✓ Branch 0 taken 19 times.
✗ Branch 1 not taken.
19x std::lock_guard<std::mutex> lock(self->offset_mutex_);
364
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 19 times.
19x if (self->generation_ == op.generation)
365 19x self->offset_ += static_cast<std::uint64_t>(n);
366 19x }
367 else
368 {
369
1/2
✓ Branch 0 taken 4 times.
✗ Branch 1 not taken.
4x op.errn = errno;
370 4x op.bytes_transferred = 0;
371 }
372 23x }
373
374 75x op.impl_ptr = std::move(pw->ref_);
375
1/2
✓ Branch 0 taken 75 times.
✗ Branch 1 not taken.
75x self->svc_.post(&op);
376 75x }
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