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