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

99.3% Lines (146/0/147) 100.0% List of functions (14/0/14)
posix_random_access_file_service.hpp
f(x) Functions (14)
Function Calls Lines Blocks
boost::corosio::detail::posix_random_access_file_service::posix_random_access_file_service(boost::capy::execution_context&, boost::corosio::detail::scheduler&) :33 2095x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::~posix_random_access_file_service() :40 4190x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::construct() :47 111x 100.0% 71.0% boost::corosio::detail::posix_random_access_file_service::destroy(boost::corosio::io_object::implementation*) :61 109x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::close(boost::corosio::io_object::handle&) :69 199x 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) :79 97x 80.0% 83.0% boost::corosio::detail::posix_random_access_file_service::shutdown() :92 2095x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::destroy_impl(boost::corosio::detail::posix_random_access_file&) :104 109x 100.0% 67.0% boost::corosio::detail::posix_random_access_file_service::post(boost::corosio::detail::scheduler_op*) :111 328x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::pool() :139 332x 100.0% 100.0% boost::corosio::detail::get_random_access_file_service(boost::capy::execution_context&, boost::corosio::detail::scheduler&) :157 2095x 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*) :167 297x 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*) :238 45x 100.0% 84.0% boost::corosio::detail::posix_random_access_file::raf_op::do_work(boost::corosio::detail::pool_work_item*) :311 328x 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 2095x posix_random_access_file_service(
34 capy::execution_context& ctx, scheduler& sched)
35 4190x : sched_(&sched)
36 2095x , pool_(ctx)
37 {
38 2095x }
39
40 4190x ~posix_random_access_file_service() override = default;
41
42 posix_random_access_file_service(
43 posix_random_access_file_service const&) = delete;
44 posix_random_access_file_service& operator=(
45 posix_random_access_file_service const&) = delete;
46
47 111x io_object::implementation* construct() override
48 {
49 111x auto ptr = std::make_shared<posix_random_access_file>(*this);
50 111x auto* impl = ptr.get();
51
52 {
53 111x std::lock_guard<std::mutex> lock(mutex_);
54 111x file_list_.push_back(impl);
55 111x file_ptrs_[impl] = std::move(ptr);
56 111x }
57
58 111x return impl;
59 111x }
60
61 109x void destroy(io_object::implementation* p) override
62 {
63 109x auto& impl = static_cast<posix_random_access_file&>(*p);
64 109x impl.cancel();
65 109x impl.close_file();
66 109x destroy_impl(impl);
67 109x }
68
69 199x void close(io_object::handle& h) override
70 {
71 199x if (h.get())
72 {
73 199x auto& impl = static_cast<posix_random_access_file&>(*h.get());
74 199x impl.cancel();
75 199x impl.close_file();
76 }
77 199x }
78
79 97x std::error_code open_file(
80 random_access_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 97x if (sched_->scheduler_locking_disabled())
87 return std::make_error_code(std::errc::operation_not_supported);
88 97x return static_cast<posix_random_access_file&>(impl).open_file(
89 97x path, mode);
90 }
91
92 2095x void shutdown() override
93 {
94 2095x std::lock_guard<std::mutex> lock(mutex_);
95 2097x for (auto* impl = file_list_.pop_front(); impl != nullptr;
96 2x impl = file_list_.pop_front())
97 {
98 2x impl->cancel();
99 2x impl->close_file();
100 }
101 2095x file_ptrs_.clear();
102 2095x }
103
104 109x void destroy_impl(posix_random_access_file& impl)
105 {
106 109x std::lock_guard<std::mutex> lock(mutex_);
107 109x file_list_.remove(&impl);
108 109x file_ptrs_.erase(&impl);
109 109x }
110
111 328x void post(scheduler_op* op)
112 {
113 328x sched_->post(op);
114 328x }
115
116 void work_started() noexcept
117 {
118 sched_->work_started();
119 }
120
121 void work_finished() noexcept
122 {
123 sched_->work_finished();
124 }
125
126 /** Return the thread pool that runs this service's file work.
127
128 The pool's service is created on first use, so this can fail
129 where a plain accessor could not. Its workers start later, on
130 the first post, and a thread the system refuses there is
131 reported by that post rather than thrown here.
132
133 @throws std::bad_alloc If the service cannot be allocated.
134
135 @return The context's shared blocking-I/O pool.
136
137 @see thread_pool_ref::get
138 */
139 332x thread_pool& pool()
140 {
141 332x return pool_.get();
142 }
143
144 private:
145 scheduler* sched_;
146 thread_pool_ref pool_;
147 std::mutex mutex_;
148 intrusive_list<posix_random_access_file> file_list_;
149 std::unordered_map<
150 posix_random_access_file*,
151 std::shared_ptr<posix_random_access_file>>
152 file_ptrs_;
153 };
154
155 /** Get or create the random-access file service for the given context. */
156 inline posix_random_access_file_service&
157 2095x get_random_access_file_service(capy::execution_context& ctx, scheduler& sched)
158 {
159 2095x return ctx.make_service<posix_random_access_file_service>(sched);
160 }
161
162 // ---------------------------------------------------------------------------
163 // posix_random_access_file inline implementations (require complete service)
164 // ---------------------------------------------------------------------------
165
166 inline std::coroutine_handle<>
167 297x posix_random_access_file::read_some_at(
168 std::uint64_t offset,
169 std::coroutine_handle<> h,
170 capy::executor_ref ex,
171 buffer_param param,
172 std::stop_token token,
173 std::error_code* ec,
174 std::size_t* bytes_out)
175 {
176 // Closed-object contract outranks the zero-length no-op.
177 297x if (fd_ < 0)
178 {
179 4x *ec = make_error_code(std::errc::bad_file_descriptor);
180 4x *bytes_out = 0;
181 4x return h;
182 }
183
184 293x capy::mutable_buffer bufs[max_buffers];
185 293x auto count = param.copy_to(bufs, max_buffers);
186
187 293x if (count == 0)
188 {
189 2x *ec = {};
190 2x *bytes_out = 0;
191 2x return h;
192 }
193
194 291x auto* op = new raf_op();
195 291x op->is_read = true;
196 291x op->offset = offset;
197
198 291x op->iovec_count = static_cast<int>(count);
199 582x for (int i = 0; i < op->iovec_count; ++i)
200 {
201 291x op->iovecs[i].iov_base = bufs[i].data();
202 291x op->iovecs[i].iov_len = bufs[i].size();
203 }
204
205 291x op->h = h;
206 291x op->ex = ex;
207 291x op->ec_out = ec;
208 291x op->bytes_out = bytes_out;
209 291x op->file_ = this;
210 291x op->impl_ptr = this->shared_from_this();
211 291x op->start(token);
212
213 291x op->ex.on_work_started();
214
215 {
216 291x std::lock_guard<std::mutex> lock(ops_mutex_);
217 291x outstanding_ops_.push_back(op);
218 291x }
219
220 291x static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
221 291x if (auto pec = svc_.pool().post(static_cast<pool_work_item*>(op)))
222 {
223 // The pool is shutting down, or the system refused it a thread.
224 // Nothing of this read went cross-thread, so it answers here
225 // like the closed-descriptor and zero-length exits above rather
226 // than through a completion the scheduler has to carry back.
227 // destroy() is the discard the op never reaching the queue
228 // needs: it unlinks, unwinds the work count and frees.
229 2x op->destroy();
230 2x *ec = pec;
231 2x *bytes_out = 0;
232 2x return h;
233 }
234 289x return std::noop_coroutine();
235 }
236
237 inline std::coroutine_handle<>
238 45x posix_random_access_file::write_some_at(
239 std::uint64_t offset,
240 std::coroutine_handle<> h,
241 capy::executor_ref ex,
242 buffer_param param,
243 std::stop_token token,
244 std::error_code* ec,
245 std::size_t* bytes_out)
246 {
247 // Closed-object contract outranks the zero-length no-op.
248 45x if (fd_ < 0)
249 {
250 2x *ec = make_error_code(std::errc::bad_file_descriptor);
251 2x *bytes_out = 0;
252 2x return h;
253 }
254
255 43x capy::mutable_buffer bufs[max_buffers];
256 43x auto count = param.copy_to(bufs, max_buffers);
257
258 43x if (count == 0)
259 {
260 2x *ec = {};
261 2x *bytes_out = 0;
262 2x return h;
263 }
264
265 41x auto* op = new raf_op();
266 41x op->is_read = false;
267 41x op->offset = offset;
268
269 41x op->iovec_count = static_cast<int>(count);
270 82x for (int i = 0; i < op->iovec_count; ++i)
271 {
272 41x op->iovecs[i].iov_base = bufs[i].data();
273 41x op->iovecs[i].iov_len = bufs[i].size();
274 }
275
276 41x op->h = h;
277 41x op->ex = ex;
278 41x op->ec_out = ec;
279 41x op->bytes_out = bytes_out;
280 41x op->file_ = this;
281 41x op->impl_ptr = this->shared_from_this();
282 41x op->start(token);
283
284 41x op->ex.on_work_started();
285
286 {
287 41x std::lock_guard<std::mutex> lock(ops_mutex_);
288 41x outstanding_ops_.push_back(op);
289 41x }
290
291 41x static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
292 41x if (auto pec = svc_.pool().post(static_cast<pool_work_item*>(op)))
293 {
294 // The pool is shutting down, or the system refused it a thread.
295 // Nothing of this write went cross-thread, so it answers here
296 // like the closed-descriptor and zero-length exits above rather
297 // than through a completion the scheduler has to carry back.
298 // destroy() is the discard the op never reaching the queue
299 // needs: it unlinks, unwinds the work count and frees.
300 2x op->destroy();
301 2x *ec = pec;
302 2x *bytes_out = 0;
303 2x return h;
304 }
305 39x return std::noop_coroutine();
306 }
307
308 // -- raf_op thread-pool work function --
309
310 inline void
311 328x posix_random_access_file::raf_op::do_work(pool_work_item* w) noexcept
312 {
313 328x auto* op = static_cast<raf_op*>(w);
314 328x auto* self = op->file_;
315
316 328x if (op->cancelled.load(std::memory_order_acquire))
317 {
318 8x op->errn = ECANCELED;
319 8x op->bytes_transferred = 0;
320 }
321 640x else if (op->offset >
322 320x static_cast<std::uint64_t>(std::numeric_limits<off_t>::max()))
323 {
324 2x op->errn = EOVERFLOW;
325 2x op->bytes_transferred = 0;
326 }
327 else
328 {
329 ssize_t n;
330 318x if (op->is_read)
331 {
332 do
333 {
334 562x n = ::preadv(self->fd_, op->iovecs, op->iovec_count,
335 281x static_cast<off_t>(op->offset));
336 }
337 281x while (n < 0 && errno == EINTR);
338 }
339 else
340 {
341 do
342 {
343 74x n = ::pwritev(self->fd_, op->iovecs, op->iovec_count,
344 37x static_cast<off_t>(op->offset));
345 }
346 37x while (n < 0 && errno == EINTR);
347 }
348
349 318x if (n >= 0)
350 {
351 304x op->errn = 0;
352 304x op->bytes_transferred = static_cast<std::size_t>(n);
353 }
354 else
355 {
356 14x op->errn = errno;
357 14x op->bytes_transferred = 0;
358 }
359 }
360
361 328x self->svc_.post(static_cast<scheduler_op*>(op));
362 328x }
363
364 } // namespace boost::corosio::detail
365
366 #endif // BOOST_COROSIO_POSIX
367
368 #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RANDOM_ACCESS_FILE_SERVICE_HPP
369