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/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 HIT 164 : explicit posix_random_access_file_service(capy::execution_context& ctx)
37 492 : : sched_(&get_scheduler(ctx))
38 164 : , pool_(ctx)
39 : {
40 164 : }
41 :
42 328 : ~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 219 : io_object::implementation* construct() override
50 : {
51 219 : return object_pool_.acquire(*this);
52 : }
53 :
54 217 : void destroy(io_object::implementation* p) override
55 : {
56 217 : auto& impl = static_cast<posix_random_access_file&>(*p);
57 217 : impl.cancel();
58 217 : impl.close_file();
59 217 : release(&impl);
60 217 : }
61 :
62 413 : void close(io_object::handle& h) override
63 : {
64 413 : if (h.get())
65 : {
66 413 : auto& impl = static_cast<posix_random_access_file&>(*h.get());
67 413 : impl.cancel();
68 413 : impl.close_file();
69 : }
70 413 : }
71 :
72 203 : 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 203 : if (sched_->scheduler_locking_disabled())
80 MIS 0 : return std::make_error_code(std::errc::operation_not_supported);
81 HIT 203 : return static_cast<posix_random_access_file&>(impl).open_file(
82 203 : path, mode);
83 : }
84 :
85 164 : 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 164 : std::vector<posix_random_access_file*> live;
93 164 : object_pool_.shutdown(
94 164 : [&](posix_random_access_file* f)
95 : {
96 4 : acquire(f);
97 4 : live.push_back(f);
98 4 : });
99 168 : for (auto* f : live)
100 : {
101 4 : f->cancel();
102 4 : f->close_file();
103 4 : release(f);
104 : }
105 164 : }
106 :
107 490 : void post(scheduler_op* op)
108 : {
109 490 : sched_->post(op);
110 490 : }
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 494 : thread_pool& pool()
136 : {
137 494 : 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 379 : 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 379 : if (fd_ < 0)
162 : {
163 4 : *ec = make_error_code(std::errc::bad_file_descriptor);
164 4 : *bytes_out = 0;
165 4 : return h;
166 : }
167 :
168 375 : capy::mutable_buffer bufs[max_buffers];
169 375 : auto count = param.copy_to(bufs, max_buffers);
170 :
171 375 : if (count == 0)
172 : {
173 2 : *ec = {};
174 2 : *bytes_out = 0;
175 2 : return h;
176 : }
177 :
178 373 : auto* op = acquire_op();
179 373 : op->is_read = true;
180 373 : op->offset = offset;
181 :
182 373 : op->iovec_count = static_cast<int>(count);
183 746 : for (int i = 0; i < op->iovec_count; ++i)
184 : {
185 373 : op->iovecs[i].iov_base = bufs[i].data();
186 373 : op->iovecs[i].iov_len = bufs[i].size();
187 : }
188 :
189 373 : op->h = h;
190 373 : op->ex = ex;
191 373 : op->ec_out = ec;
192 373 : op->bytes_out = bytes_out;
193 373 : op->file_ = this;
194 373 : op->object_ref_ = detail::object_ref(this);
195 373 : op->start(token);
196 :
197 373 : op->ex.on_work_started();
198 :
199 : {
200 373 : std::lock_guard<std::mutex> lock(ops_mutex_);
201 373 : outstanding_ops_.push_back(op);
202 373 : }
203 :
204 373 : static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
205 373 : 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 2 : op->destroy();
214 2 : *ec = pec;
215 2 : *bytes_out = 0;
216 2 : return h;
217 : }
218 371 : return std::noop_coroutine();
219 : }
220 :
221 : inline std::coroutine_handle<>
222 125 : 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 125 : if (fd_ < 0)
233 : {
234 2 : *ec = make_error_code(std::errc::bad_file_descriptor);
235 2 : *bytes_out = 0;
236 2 : return h;
237 : }
238 :
239 123 : capy::mutable_buffer bufs[max_buffers];
240 123 : auto count = param.copy_to(bufs, max_buffers);
241 :
242 123 : if (count == 0)
243 : {
244 2 : *ec = {};
245 2 : *bytes_out = 0;
246 2 : return h;
247 : }
248 :
249 121 : auto* op = acquire_op();
250 121 : op->is_read = false;
251 121 : op->offset = offset;
252 :
253 121 : op->iovec_count = static_cast<int>(count);
254 242 : for (int i = 0; i < op->iovec_count; ++i)
255 : {
256 121 : op->iovecs[i].iov_base = bufs[i].data();
257 121 : op->iovecs[i].iov_len = bufs[i].size();
258 : }
259 :
260 121 : op->h = h;
261 121 : op->ex = ex;
262 121 : op->ec_out = ec;
263 121 : op->bytes_out = bytes_out;
264 121 : op->file_ = this;
265 121 : op->object_ref_ = detail::object_ref(this);
266 121 : op->start(token);
267 :
268 121 : op->ex.on_work_started();
269 :
270 : {
271 121 : std::lock_guard<std::mutex> lock(ops_mutex_);
272 121 : outstanding_ops_.push_back(op);
273 121 : }
274 :
275 121 : static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
276 121 : 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 2 : op->destroy();
285 2 : *ec = pec;
286 2 : *bytes_out = 0;
287 2 : return h;
288 : }
289 119 : return std::noop_coroutine();
290 : }
291 :
292 : // -- raf_op thread-pool work function --
293 :
294 : inline void
295 490 : posix_random_access_file::raf_op::do_work(pool_work_item* w) noexcept
296 : {
297 490 : auto* op = static_cast<raf_op*>(w);
298 490 : auto* self = op->file_;
299 :
300 490 : if (op->cancelled.load(std::memory_order_acquire))
301 : {
302 61 : op->errn = ECANCELED;
303 61 : op->bytes_transferred = 0;
304 : }
305 429 : else if (
306 858 : op->offset >
307 429 : static_cast<std::uint64_t>(std::numeric_limits<off_t>::max()))
308 : {
309 2 : op->errn = EOVERFLOW;
310 2 : op->bytes_transferred = 0;
311 : }
312 : else
313 : {
314 : ssize_t n;
315 427 : if (op->is_read)
316 : {
317 : do
318 : {
319 634 : n = ::preadv(
320 317 : self->fd_, op->iovecs, op->iovec_count,
321 317 : static_cast<off_t>(op->offset));
322 : }
323 317 : while (n < 0 && errno == EINTR);
324 : }
325 : else
326 : {
327 : do
328 : {
329 220 : n = ::pwritev(
330 110 : self->fd_, op->iovecs, op->iovec_count,
331 110 : static_cast<off_t>(op->offset));
332 : }
333 110 : while (n < 0 && errno == EINTR);
334 : }
335 :
336 427 : if (n >= 0)
337 : {
338 413 : op->errn = 0;
339 413 : op->bytes_transferred = static_cast<std::size_t>(n);
340 : }
341 : else
342 : {
343 14 : op->errn = errno;
344 14 : op->bytes_transferred = 0;
345 : }
346 : }
347 :
348 490 : self->svc_.post(static_cast<scheduler_op*>(op));
349 490 : }
350 :
351 : inline void
352 217 : posix_random_access_file::retire() noexcept
353 : {
354 217 : svc_.object_pool_.recycle(this);
355 217 : }
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
|