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