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 172x 100.0% 89.0% boost::corosio::detail::posix_random_access_file_service::~posix_random_access_file_service() :39 344x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::construct() :46 227x 100.0% 71.0% boost::corosio::detail::posix_random_access_file_service::destroy(boost::corosio::io_object::implementation*) :60 225x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::close(boost::corosio::io_object::handle&) :68 421x 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 211x 80.0% 83.0% boost::corosio::detail::posix_random_access_file_service::shutdown() :91 172x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::destroy_impl(boost::corosio::detail::posix_random_access_file&) :103 225x 100.0% 67.0% boost::corosio::detail::posix_random_access_file_service::post(boost::corosio::detail::scheduler_op*) :110 440x 100.0% 100.0% boost::corosio::detail::posix_random_access_file_service::pool() :138 444x 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 353x 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 101x 100.0% 84.0% boost::corosio::detail::posix_random_access_file::raf_op::do_work(boost::corosio::detail::pool_work_item*) :307 440x 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 172x explicit posix_random_access_file_service(capy::execution_context& ctx)
34 516x : sched_(&get_scheduler(ctx))
35 172x , pool_(ctx)
36 {
37 172x }
38
39 344x ~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 227x io_object::implementation* construct() override
47 {
48 227x auto ptr = std::make_shared<posix_random_access_file>(*this);
49 227x auto* impl = ptr.get();
50
51 {
52 227x std::lock_guard<std::mutex> lock(mutex_);
53 227x file_list_.push_back(impl);
54 227x file_ptrs_[impl] = std::move(ptr);
55 227x }
56
57 227x return impl;
58 227x }
59
60 225x void destroy(io_object::implementation* p) override
61 {
62 225x auto& impl = static_cast<posix_random_access_file&>(*p);
63 225x impl.cancel();
64 225x impl.close_file();
65 225x destroy_impl(impl);
66 225x }
67
68 421x void close(io_object::handle& h) override
69 {
70 421x if (h.get())
71 {
72 421x auto& impl = static_cast<posix_random_access_file&>(*h.get());
73 421x impl.cancel();
74 421x impl.close_file();
75 }
76 421x }
77
78 211x 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 211x if (sched_->scheduler_locking_disabled())
86 ✗ return std::make_error_code(std::errc::operation_not_supported);
87 211x return static_cast<posix_random_access_file&>(impl).open_file(
88 211x path, mode);
89 }
90
91 172x void shutdown() override
92 {
93 172x std::lock_guard<std::mutex> lock(mutex_);
94 174x 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 172x file_ptrs_.clear();
101 172x }
102
103 225x void destroy_impl(posix_random_access_file& impl)
104 {
105 225x std::lock_guard<std::mutex> lock(mutex_);
106 225x file_list_.remove(&impl);
107 225x file_ptrs_.erase(&impl);
108 225x }
109
110 440x void post(scheduler_op* op)
111 {
112 440x sched_->post(op);
113 440x }
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 444x thread_pool& pool()
139 {
140 444x 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 353x 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 353x 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 349x capy::mutable_buffer bufs[max_buffers];
177 349x auto count = param.copy_to(bufs, max_buffers);
178
179 349x if (count == 0)
180 {
181 2x *ec = {};
182 2x *bytes_out = 0;
183 2x return cont.h;
184 }
185
186 347x auto* op = new raf_op();
187 347x op->is_read = true;
188 347x op->offset = offset;
189 347x op->fd = fd_;
190
191 347x op->iovec_count = static_cast<int>(count);
192 694x for (int i = 0; i < op->iovec_count; ++i)
193 {
194 347x op->iovecs[i].iov_base = bufs[i].data();
195 347x op->iovecs[i].iov_len = bufs[i].size();
196 }
197
198 347x op->h = cont.h;
199 347x op->awaiting = &cont;
200 347x op->ex = ex;
201 347x op->ec_out = ec;
202 347x op->bytes_out = bytes_out;
203 347x op->file_ = this;
204 347x op->impl_ptr = this->shared_from_this();
205 347x op->start(token);
206
207 347x op->ex.on_work_started();
208
209 {
210 347x std::lock_guard<std::mutex> lock(ops_mutex_);
211 347x outstanding_ops_.push_back(op);
212 347x }
213
214 347x static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
215 347x 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 345x return std::noop_coroutine();
229 }
230
231 inline std::coroutine_handle<>
232 101x 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 101x 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 99x capy::mutable_buffer bufs[max_buffers];
250 99x auto count = param.copy_to(bufs, max_buffers);
251
252 99x if (count == 0)
253 {
254 2x *ec = {};
255 2x *bytes_out = 0;
256 2x return cont.h;
257 }
258
259 97x auto* op = new raf_op();
260 97x op->is_read = false;
261 97x op->offset = offset;
262 97x op->fd = fd_;
263
264 97x op->iovec_count = static_cast<int>(count);
265 194x for (int i = 0; i < op->iovec_count; ++i)
266 {
267 97x op->iovecs[i].iov_base = bufs[i].data();
268 97x op->iovecs[i].iov_len = bufs[i].size();
269 }
270
271 97x op->h = cont.h;
272 97x op->awaiting = &cont;
273 97x op->ex = ex;
274 97x op->ec_out = ec;
275 97x op->bytes_out = bytes_out;
276 97x op->file_ = this;
277 97x op->impl_ptr = this->shared_from_this();
278 97x op->start(token);
279
280 97x op->ex.on_work_started();
281
282 {
283 97x std::lock_guard<std::mutex> lock(ops_mutex_);
284 97x outstanding_ops_.push_back(op);
285 97x }
286
287 97x static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
288 97x 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 95x return std::noop_coroutine();
302 }
303
304 // -- raf_op thread-pool work function --
305
306 inline void
307 440x posix_random_access_file::raf_op::do_work(pool_work_item* w) noexcept
308 {
309 440x auto* op = static_cast<raf_op*>(w);
310 440x auto* self = op->file_;
311
312 440x if (op->cancelled.load(std::memory_order_acquire))
313 {
314 67x op->errn = ECANCELED;
315 67x op->bytes_transferred = 0;
316 }
317 373x else if (
318 746x op->offset >
319 373x 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 371x if (op->is_read)
328 {
329 do
330 {
331 291x n = file_preadv(
332 291x op->fd, op->iovecs, op->iovec_count,
333 291x static_cast<file_off_t>(op->offset));
334 }
335 291x while (n < 0 && errno == EINTR);
336 }
337 else
338 {
339 do
340 {
341 80x n = file_pwritev(
342 80x op->fd, op->iovecs, op->iovec_count,
343 80x static_cast<file_off_t>(op->offset));
344 }
345 80x while (n < 0 && errno == EINTR);
346 }
347
348 371x if (n >= 0)
349 {
350 353x op->errn = 0;
351 353x op->bytes_transferred = static_cast<std::size_t>(n);
352 }
353 else
354 {
355 18x op->errn = errno;
356 18x op->bytes_transferred = 0;
357 }
358 }
359
360 440x self->svc_.post(static_cast<scheduler_op*>(op));
361 440x }
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