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

99.3% Lines (151 / 152) 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 170x 100.0% 89.0% boost::corosio::detail::posix_random_access_file_service::~posix_random_access_file_service() :39 340x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::construct() :46 225x 100.0% 71.0% boost::corosio::detail::posix_random_access_file_service::destroy(boost::corosio::io_object::implementation*) :60 223x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::close(boost::corosio::io_object::handle&) :68 417x 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 207x 80.0% 83.0% boost::corosio::detail::posix_random_access_file_service::shutdown() :91 170x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::destroy_impl(boost::corosio::detail::posix_random_access_file&) :103 223x 100.0% 67.0% boost::corosio::detail::posix_random_access_file_service::post(boost::corosio::detail::scheduler_op*) :110 432x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::pool() :138 436x 100.0% 100.0% boost::corosio::detail::posix_random_access_file::read_some_at(unsigned long, boost::capy::continuation&, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :159 349x 100.0% 84.0% boost::corosio::detail::posix_random_access_file::write_some_at(unsigned long, boost::capy::continuation&, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :232 97x 100.0% 84.0% boost::corosio::detail::posix_random_access_file::raf_op::do_work(boost::corosio::detail::pool_work_item*) :307 432x 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 170x explicit posix_random_access_file_service(capy::execution_context& ctx)
34 510x : sched_(&get_scheduler(ctx))
35 170x , pool_(ctx)
36 {
37 170x }
38
39 340x ~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 225x io_object::implementation* construct() override
47 {
48 225x auto ptr = std::make_shared<posix_random_access_file>(*this);
49 225x auto* impl = ptr.get();
50
51 {
52 225x std::lock_guard<std::mutex> lock(mutex_);
53 225x file_list_.push_back(impl);
54 225x file_ptrs_[impl] = std::move(ptr);
55 225x }
56
57 225x return impl;
58 225x }
59
60 223x void destroy(io_object::implementation* p) override
61 {
62 223x auto& impl = static_cast<posix_random_access_file&>(*p);
63 223x impl.cancel();
64 223x impl.close_file();
65 223x destroy_impl(impl);
66 223x }
67
68 417x void close(io_object::handle& h) override
69 {
70 417x if (h.get())
71 {
72 417x auto& impl = static_cast<posix_random_access_file&>(*h.get());
73 417x impl.cancel();
74 417x impl.close_file();
75 }
76 417x }
77
78 207x 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 207x if (sched_->scheduler_locking_disabled())
86 ✗ return std::make_error_code(std::errc::operation_not_supported);
87 207x return static_cast<posix_random_access_file&>(impl).open_file(
88 207x path, mode);
89 }
90
91 170x void shutdown() override
92 {
93 170x std::lock_guard<std::mutex> lock(mutex_);
94 172x 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 170x file_ptrs_.clear();
101 170x }
102
103 223x void destroy_impl(posix_random_access_file& impl)
104 {
105 223x std::lock_guard<std::mutex> lock(mutex_);
106 223x file_list_.remove(&impl);
107 223x file_ptrs_.erase(&impl);
108 223x }
109
110 432x void post(scheduler_op* op)
111 {
112 432x sched_->post(op);
113 432x }
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 436x thread_pool& pool()
139 {
140 436x 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 349x posix_random_access_file::read_some_at(
160 std::uint64_t offset,
161 capy::continuation& cont,
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 349x if (fd_ < 0)
170 {
171 4x *ec = make_error_code(std::errc::bad_file_descriptor);
172 4x *bytes_out = 0;
173 4x return cont.h;
174 }
175
176 345x capy::mutable_buffer bufs[max_buffers];
177 345x auto count = param.copy_to(bufs, max_buffers);
178
179 345x if (count == 0)
180 {
181 2x *ec = {};
182 2x *bytes_out = 0;
183 2x return cont.h;
184 }
185
186 343x auto* op = new raf_op();
187 343x op->is_read = true;
188 343x op->offset = offset;
189 343x op->fd = fd_;
190
191 343x op->iovec_count = static_cast<int>(count);
192 686x for (int i = 0; i < op->iovec_count; ++i)
193 {
194 343x op->iovecs[i].iov_base = bufs[i].data();
195 343x op->iovecs[i].iov_len = bufs[i].size();
196 }
197
198 343x op->h = cont.h;
199 343x op->awaiting = &cont;
200 343x op->ex = ex;
201 343x op->ec_out = ec;
202 343x op->bytes_out = bytes_out;
203 343x op->file_ = this;
204 343x op->impl_ptr = this->shared_from_this();
205 343x op->start(token);
206
207 343x op->ex.on_work_started();
208
209 {
210 343x std::lock_guard<std::mutex> lock(ops_mutex_);
211 343x outstanding_ops_.push_back(op);
212 343x }
213
214 343x static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
215 343x if (auto pec = svc_.pool().post(static_cast<pool_work_item*>(op)))
216 {
217 // The pool is shutting down, or the system refused it a thread.
218 // Nothing of this read went cross-thread, so it answers here
219 // like the closed-descriptor and zero-length exits above rather
220 // than through a completion the scheduler has to carry back.
221 // destroy() is the discard the op never reaching the queue
222 // needs: it unlinks, unwinds the work count and frees.
223 2x op->destroy();
224 2x *ec = pec;
225 2x *bytes_out = 0;
226 2x return cont.h;
227 }
228 341x return std::noop_coroutine();
229 }
230
231 inline std::coroutine_handle<>
232 97x posix_random_access_file::write_some_at(
233 std::uint64_t offset,
234 capy::continuation& cont,
235 capy::executor_ref ex,
236 buffer_param param,
237 std::stop_token token,
238 std::error_code* ec,
239 std::size_t* bytes_out)
240 {
241 // Closed-object contract outranks the zero-length no-op.
242 97x if (fd_ < 0)
243 {
244 2x *ec = make_error_code(std::errc::bad_file_descriptor);
245 2x *bytes_out = 0;
246 2x return cont.h;
247 }
248
249 95x capy::mutable_buffer bufs[max_buffers];
250 95x auto count = param.copy_to(bufs, max_buffers);
251
252 95x if (count == 0)
253 {
254 2x *ec = {};
255 2x *bytes_out = 0;
256 2x return cont.h;
257 }
258
259 93x auto* op = new raf_op();
260 93x op->is_read = false;
261 93x op->offset = offset;
262 93x op->fd = fd_;
263
264 93x op->iovec_count = static_cast<int>(count);
265 186x for (int i = 0; i < op->iovec_count; ++i)
266 {
267 93x op->iovecs[i].iov_base = bufs[i].data();
268 93x op->iovecs[i].iov_len = bufs[i].size();
269 }
270
271 93x op->h = cont.h;
272 93x op->awaiting = &cont;
273 93x op->ex = ex;
274 93x op->ec_out = ec;
275 93x op->bytes_out = bytes_out;
276 93x op->file_ = this;
277 93x op->impl_ptr = this->shared_from_this();
278 93x op->start(token);
279
280 93x op->ex.on_work_started();
281
282 {
283 93x std::lock_guard<std::mutex> lock(ops_mutex_);
284 93x outstanding_ops_.push_back(op);
285 93x }
286
287 93x static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
288 93x if (auto pec = svc_.pool().post(static_cast<pool_work_item*>(op)))
289 {
290 // The pool is shutting down, or the system refused it a thread.
291 // Nothing of this write went cross-thread, so it answers here
292 // like the closed-descriptor and zero-length exits above rather
293 // than through a completion the scheduler has to carry back.
294 // destroy() is the discard the op never reaching the queue
295 // needs: it unlinks, unwinds the work count and frees.
296 2x op->destroy();
297 2x *ec = pec;
298 2x *bytes_out = 0;
299 2x return cont.h;
300 }
301 91x return std::noop_coroutine();
302 }
303
304 // -- raf_op thread-pool work function --
305
306 inline void
307 432x posix_random_access_file::raf_op::do_work(pool_work_item* w) noexcept
308 {
309 432x auto* op = static_cast<raf_op*>(w);
310 432x auto* self = op->file_;
311
312 432x if (op->cancelled.load(std::memory_order_acquire))
313 {
314 102x op->errn = ECANCELED;
315 102x op->bytes_transferred = 0;
316 }
317 330x else if (
318 660x op->offset >
319 330x static_cast<std::uint64_t>(std::numeric_limits<file_off_t>::max()))
320 {
321 2x op->errn = EOVERFLOW;
322 2x op->bytes_transferred = 0;
323 }
324 else
325 {
326 ssize_t n;
327 328x if (op->is_read)
328 {
329 do
330 {
331 287x n = file_preadv(
332 287x op->fd, op->iovecs, op->iovec_count,
333 287x static_cast<file_off_t>(op->offset));
334 }
335 287x while (n < 0 && errno == EINTR);
336 }
337 else
338 {
339 do
340 {
341 41x n = file_pwritev(
342 41x op->fd, op->iovecs, op->iovec_count,
343 41x static_cast<file_off_t>(op->offset));
344 }
345 41x while (n < 0 && errno == EINTR);
346 }
347
348 328x if (n >= 0)
349 {
350 314x op->errn = 0;
351 314x op->bytes_transferred = static_cast<std::size_t>(n);
352 }
353 else
354 {
355 14x op->errn = errno;
356 14x op->bytes_transferred = 0;
357 }
358 }
359
360 432x self->svc_.post(static_cast<scheduler_op*>(op));
361 432x }
362
363 } // namespace boost::corosio::detail
364
365 #endif // BOOST_COROSIO_POSIX
366
367 #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RANDOM_ACCESS_FILE_SERVICE_HPP
368