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

99.3% Lines (147 / 148) 100.0% Functions (13 / 13)
posix_random_access_file_service.hpp
f(x) Functions (13)
Function Calls Lines Blocks
boost::corosio::detail::posix_random_access_file_service::posix_random_access_file_service(boost::capy::execution_context&) :33 160x 100.0% 89.0% boost::corosio::detail::posix_random_access_file_service::~posix_random_access_file_service() :39 320x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::construct() :46 215x 100.0% 71.0% boost::corosio::detail::posix_random_access_file_service::destroy(boost::corosio::io_object::implementation*) :60 213x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::close(boost::corosio::io_object::handle&) :68 405x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::open_file(boost::corosio::random_access_file::implementation&, std::filesystem::__cxx11::path const&, boost::corosio::file_base::flags) :78 199x 80.0% 83.0% boost::corosio::detail::posix_random_access_file_service::shutdown() :91 160x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::destroy_impl(boost::corosio::detail::posix_random_access_file&) :103 213x 100.0% 67.0% boost::corosio::detail::posix_random_access_file_service::post(boost::corosio::detail::scheduler_op*) :110 422x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::pool() :138 426x 100.0% 100.0% boost::corosio::detail::posix_random_access_file::read_some_at(unsigned long, std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :159 343x 100.0% 84.0% boost::corosio::detail::posix_random_access_file::write_some_at(unsigned long, std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :230 93x 100.0% 84.0% boost::corosio::detail::posix_random_access_file::raf_op::do_work(boost::corosio::detail::pool_work_item*) :303 422x 100.0% 95.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_RANDOM_ACCESS_FILE_SERVICE_HPP
11 #define BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RANDOM_ACCESS_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_random_access_file.hpp>
18 #include <boost/corosio/native/detail/reactor/reactor_scheduler.hpp>
19 #include <boost/corosio/detail/random_access_file_service.hpp>
20 #include <boost/corosio/detail/thread_pool.hpp>
21
22 #include <limits>
23 #include <mutex>
24 #include <unordered_map>
25
26 namespace boost::corosio::detail {
27
28 /** Random-access file service for POSIX backends. */
29 class BOOST_COROSIO_DECL posix_random_access_file_service final
30 : public random_access_file_service
31 {
32 public:
33 160x explicit posix_random_access_file_service(capy::execution_context& ctx)
34 480x : sched_(&get_scheduler(ctx))
35 160x , pool_(ctx)
36 {
37 160x }
38
39 320x ~posix_random_access_file_service() override = default;
40
41 posix_random_access_file_service(posix_random_access_file_service const&) =
42 delete;
43 posix_random_access_file_service&
44 operator=(posix_random_access_file_service const&) = delete;
45
46 215x io_object::implementation* construct() override
47 {
48 215x auto ptr = std::make_shared<posix_random_access_file>(*this);
49 215x auto* impl = ptr.get();
50
51 {
52 215x std::lock_guard<std::mutex> lock(mutex_);
53 215x file_list_.push_back(impl);
54 215x file_ptrs_[impl] = std::move(ptr);
55 215x }
56
57 215x return impl;
58 215x }
59
60 213x void destroy(io_object::implementation* p) override
61 {
62 213x auto& impl = static_cast<posix_random_access_file&>(*p);
63 213x impl.cancel();
64 213x impl.close_file();
65 213x destroy_impl(impl);
66 213x }
67
68 405x void close(io_object::handle& h) override
69 {
70 405x if (h.get())
71 {
72 405x auto& impl = static_cast<posix_random_access_file&>(*h.get());
73 405x impl.cancel();
74 405x impl.close_file();
75 }
76 405x }
77
78 199x std::error_code open_file(
79 random_access_file::implementation& impl,
80 std::filesystem::path const& path,
81 file_base::flags mode) override
82 {
83 // Unavailable in the unsafe tier: the file thread pool completes
84 // cross-thread, which the lockless scheduler cannot accept.
85 199x if (sched_->scheduler_locking_disabled())
86 ✗ return std::make_error_code(std::errc::operation_not_supported);
87 199x return static_cast<posix_random_access_file&>(impl).open_file(
88 199x path, mode);
89 }
90
91 160x void shutdown() override
92 {
93 160x std::lock_guard<std::mutex> lock(mutex_);
94 162x 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 160x file_ptrs_.clear();
101 160x }
102
103 213x void destroy_impl(posix_random_access_file& impl)
104 {
105 213x std::lock_guard<std::mutex> lock(mutex_);
106 213x file_list_.remove(&impl);
107 213x file_ptrs_.erase(&impl);
108 213x }
109
110 422x void post(scheduler_op* op)
111 {
112 422x sched_->post(op);
113 422x }
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 426x thread_pool& pool()
139 {
140 426x return pool_.get();
141 }
142
143 private:
144 scheduler* sched_;
145 thread_pool_ref pool_;
146 std::mutex mutex_;
147 intrusive_list<posix_random_access_file> file_list_;
148 std::unordered_map<
149 posix_random_access_file*,
150 std::shared_ptr<posix_random_access_file>>
151 file_ptrs_;
152 };
153
154 // ---------------------------------------------------------------------------
155 // posix_random_access_file inline implementations (require complete service)
156 // ---------------------------------------------------------------------------
157
158 inline std::coroutine_handle<>
159 343x posix_random_access_file::read_some_at(
160 std::uint64_t offset,
161 std::coroutine_handle<> h,
162 capy::executor_ref ex,
163 buffer_param param,
164 std::stop_token token,
165 std::error_code* ec,
166 std::size_t* bytes_out)
167 {
168 // Closed-object contract outranks the zero-length no-op.
169 343x if (fd_ < 0)
170 {
171 4x *ec = make_error_code(std::errc::bad_file_descriptor);
172 4x *bytes_out = 0;
173 4x return h;
174 }
175
176 339x capy::mutable_buffer bufs[max_buffers];
177 339x auto count = param.copy_to(bufs, max_buffers);
178
179 339x if (count == 0)
180 {
181 2x *ec = {};
182 2x *bytes_out = 0;
183 2x return h;
184 }
185
186 337x auto* op = new raf_op();
187 337x op->is_read = true;
188 337x op->offset = offset;
189
190 337x op->iovec_count = static_cast<int>(count);
191 674x for (int i = 0; i < op->iovec_count; ++i)
192 {
193 337x op->iovecs[i].iov_base = bufs[i].data();
194 337x op->iovecs[i].iov_len = bufs[i].size();
195 }
196
197 337x op->h = h;
198 337x op->ex = ex;
199 337x op->ec_out = ec;
200 337x op->bytes_out = bytes_out;
201 337x op->file_ = this;
202 337x op->impl_ptr = this->shared_from_this();
203 337x op->start(token);
204
205 337x op->ex.on_work_started();
206
207 {
208 337x std::lock_guard<std::mutex> lock(ops_mutex_);
209 337x outstanding_ops_.push_back(op);
210 337x }
211
212 337x static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
213 337x if (auto pec = svc_.pool().post(static_cast<pool_work_item*>(op)))
214 {
215 // The pool is shutting down, or the system refused it a thread.
216 // Nothing of this read went cross-thread, so it answers here
217 // like the closed-descriptor and zero-length exits above rather
218 // than through a completion the scheduler has to carry back.
219 // destroy() is the discard the op never reaching the queue
220 // needs: it unlinks, unwinds the work count and frees.
221 2x op->destroy();
222 2x *ec = pec;
223 2x *bytes_out = 0;
224 2x return h;
225 }
226 335x return std::noop_coroutine();
227 }
228
229 inline std::coroutine_handle<>
230 93x posix_random_access_file::write_some_at(
231 std::uint64_t offset,
232 std::coroutine_handle<> h,
233 capy::executor_ref ex,
234 buffer_param param,
235 std::stop_token token,
236 std::error_code* ec,
237 std::size_t* bytes_out)
238 {
239 // Closed-object contract outranks the zero-length no-op.
240 93x if (fd_ < 0)
241 {
242 2x *ec = make_error_code(std::errc::bad_file_descriptor);
243 2x *bytes_out = 0;
244 2x return h;
245 }
246
247 91x capy::mutable_buffer bufs[max_buffers];
248 91x auto count = param.copy_to(bufs, max_buffers);
249
250 91x if (count == 0)
251 {
252 2x *ec = {};
253 2x *bytes_out = 0;
254 2x return h;
255 }
256
257 89x auto* op = new raf_op();
258 89x op->is_read = false;
259 89x op->offset = offset;
260
261 89x op->iovec_count = static_cast<int>(count);
262 178x for (int i = 0; i < op->iovec_count; ++i)
263 {
264 89x op->iovecs[i].iov_base = bufs[i].data();
265 89x op->iovecs[i].iov_len = bufs[i].size();
266 }
267
268 89x op->h = h;
269 89x op->ex = ex;
270 89x op->ec_out = ec;
271 89x op->bytes_out = bytes_out;
272 89x op->file_ = this;
273 89x op->impl_ptr = this->shared_from_this();
274 89x op->start(token);
275
276 89x op->ex.on_work_started();
277
278 {
279 89x std::lock_guard<std::mutex> lock(ops_mutex_);
280 89x outstanding_ops_.push_back(op);
281 89x }
282
283 89x static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
284 89x if (auto pec = svc_.pool().post(static_cast<pool_work_item*>(op)))
285 {
286 // The pool is shutting down, or the system refused it a thread.
287 // Nothing of this write went cross-thread, so it answers here
288 // like the closed-descriptor and zero-length exits above rather
289 // than through a completion the scheduler has to carry back.
290 // destroy() is the discard the op never reaching the queue
291 // needs: it unlinks, unwinds the work count and frees.
292 2x op->destroy();
293 2x *ec = pec;
294 2x *bytes_out = 0;
295 2x return h;
296 }
297 87x return std::noop_coroutine();
298 }
299
300 // -- raf_op thread-pool work function --
301
302 inline void
303 422x posix_random_access_file::raf_op::do_work(pool_work_item* w) noexcept
304 {
305 422x auto* op = static_cast<raf_op*>(w);
306 422x auto* self = op->file_;
307
308 422x if (op->cancelled.load(std::memory_order_acquire))
309 {
310 57x op->errn = ECANCELED;
311 57x op->bytes_transferred = 0;
312 }
313 365x else if (
314 730x op->offset >
315 365x static_cast<std::uint64_t>(std::numeric_limits<off_t>::max()))
316 {
317 2x op->errn = EOVERFLOW;
318 2x op->bytes_transferred = 0;
319 }
320 else
321 {
322 ssize_t n;
323 363x if (op->is_read)
324 {
325 do
326 {
327 562x n = ::preadv(
328 281x self->fd_, op->iovecs, op->iovec_count,
329 281x static_cast<off_t>(op->offset));
330 }
331 281x while (n < 0 && errno == EINTR);
332 }
333 else
334 {
335 do
336 {
337 164x n = ::pwritev(
338 82x self->fd_, op->iovecs, op->iovec_count,
339 82x static_cast<off_t>(op->offset));
340 }
341 82x while (n < 0 && errno == EINTR);
342 }
343
344 363x if (n >= 0)
345 {
346 349x op->errn = 0;
347 349x op->bytes_transferred = static_cast<std::size_t>(n);
348 }
349 else
350 {
351 14x op->errn = errno;
352 14x op->bytes_transferred = 0;
353 }
354 }
355
356 422x self->svc_.post(static_cast<scheduler_op*>(op));
357 422x }
358
359 } // namespace boost::corosio::detail
360
361 #endif // BOOST_COROSIO_POSIX
362
363 #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RANDOM_ACCESS_FILE_SERVICE_HPP
364