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

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