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_STREAM_FILE_SERVICE_HPP
11 : #define BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_STREAM_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_stream_file.hpp>
18 : #include <boost/corosio/native/detail/reactor/reactor_scheduler.hpp>
19 : #include <boost/corosio/detail/file_service.hpp>
20 : #include <boost/corosio/detail/thread_pool.hpp>
21 :
22 : #include <mutex>
23 : #include <unordered_map>
24 :
25 : namespace boost::corosio::detail {
26 :
27 : /** Stream file service for POSIX backends.
28 :
29 : Owns all posix_stream_file instances. Thread lifecycle is
30 : managed by the thread_pool service (shared with resolver).
31 : */
32 : class BOOST_COROSIO_DECL posix_stream_file_service final : public file_service
33 : {
34 : public:
35 HIT 621 : explicit posix_stream_file_service(capy::execution_context& ctx)
36 1863 : : sched_(&get_scheduler(ctx))
37 621 : , pool_(ctx)
38 : {
39 621 : }
40 :
41 1242 : ~posix_stream_file_service() override = default;
42 :
43 : posix_stream_file_service(posix_stream_file_service const&) = delete;
44 : posix_stream_file_service&
45 : operator=(posix_stream_file_service const&) = delete;
46 :
47 677 : io_object::implementation* construct() override
48 : {
49 677 : auto ptr = std::make_shared<posix_stream_file>(*this);
50 677 : auto* impl = ptr.get();
51 :
52 : {
53 677 : std::lock_guard<std::mutex> lock(mutex_);
54 677 : file_list_.push_back(impl);
55 677 : file_ptrs_[impl] = std::move(ptr);
56 677 : }
57 :
58 677 : return impl;
59 677 : }
60 :
61 675 : void destroy(io_object::implementation* p) override
62 : {
63 675 : auto& impl = static_cast<posix_stream_file&>(*p);
64 675 : impl.cancel();
65 675 : impl.close_file();
66 675 : destroy_impl(impl);
67 675 : }
68 :
69 1718 : void close(io_object::handle& h) override
70 : {
71 1718 : if (h.get())
72 : {
73 1718 : auto& impl = static_cast<posix_stream_file&>(*h.get());
74 1718 : impl.cancel();
75 1718 : impl.close_file();
76 : }
77 1718 : }
78 :
79 658 : std::error_code open_file(
80 : stream_file::implementation& impl,
81 : std::filesystem::path const& path,
82 : file_base::flags mode) override
83 : {
84 : // Unavailable in the unsafe tier: the file thread pool completes
85 : // cross-thread, which the lockless scheduler cannot accept.
86 658 : if (sched_->scheduler_locking_disabled())
87 2 : return std::make_error_code(std::errc::operation_not_supported);
88 656 : return static_cast<posix_stream_file&>(impl).open_file(path, mode);
89 : }
90 :
91 621 : void shutdown() override
92 : {
93 621 : std::lock_guard<std::mutex> lock(mutex_);
94 623 : 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 621 : file_ptrs_.clear();
101 621 : }
102 :
103 675 : void destroy_impl(posix_stream_file& impl)
104 : {
105 675 : std::lock_guard<std::mutex> lock(mutex_);
106 675 : file_list_.remove(&impl);
107 675 : file_ptrs_.erase(&impl);
108 675 : }
109 :
110 626 : void post(scheduler_op* op)
111 : {
112 626 : sched_->post(op);
113 626 : }
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 634 : thread_pool& pool()
139 : {
140 634 : return pool_.get();
141 : }
142 :
143 : private:
144 : scheduler* sched_;
145 : thread_pool_ref pool_;
146 : std::mutex mutex_;
147 : intrusive_list<posix_stream_file> file_list_;
148 : std::unordered_map<posix_stream_file*, std::shared_ptr<posix_stream_file>>
149 : file_ptrs_;
150 : };
151 :
152 : // ---------------------------------------------------------------------------
153 : // posix_stream_file inline implementations (require complete service type)
154 : // ---------------------------------------------------------------------------
155 :
156 : inline std::coroutine_handle<>
157 559 : posix_stream_file::read_some(
158 : std::coroutine_handle<> h,
159 : capy::executor_ref ex,
160 : buffer_param param,
161 : std::stop_token token,
162 : std::error_code* ec,
163 : std::size_t* bytes_out)
164 : {
165 559 : auto& op = read_op_;
166 559 : op.reset();
167 559 : op.is_read = true;
168 :
169 : // Closed-object contract outranks the zero-length no-op.
170 559 : if (fd_ < 0)
171 : {
172 6 : *ec = make_error_code(std::errc::bad_file_descriptor);
173 6 : *bytes_out = 0;
174 6 : op.cont.h = h;
175 6 : return dispatch_coro(ex, op.cont);
176 : }
177 :
178 553 : capy::mutable_buffer bufs[max_buffers];
179 553 : op.iovec_count = static_cast<int>(param.copy_to(bufs, max_buffers));
180 :
181 553 : if (op.iovec_count == 0)
182 : {
183 2 : *ec = {};
184 2 : *bytes_out = 0;
185 2 : op.cont.h = h;
186 2 : return dispatch_coro(ex, op.cont);
187 : }
188 :
189 1102 : for (int i = 0; i < op.iovec_count; ++i)
190 : {
191 551 : op.iovecs[i].iov_base = bufs[i].data();
192 551 : op.iovecs[i].iov_len = bufs[i].size();
193 : }
194 :
195 551 : op.h = h;
196 551 : op.ex = ex;
197 551 : op.ec_out = ec;
198 551 : op.bytes_out = bytes_out;
199 551 : op.start(token);
200 :
201 551 : op.fd = fd_;
202 551 : op.offset = offset_;
203 551 : op.generation = generation_;
204 :
205 551 : op.ex.on_work_started();
206 :
207 551 : read_pool_op_.file_ = this;
208 551 : read_pool_op_.ref_ = this->shared_from_this();
209 551 : read_pool_op_.func_ = &posix_stream_file::do_read_work;
210 551 : if (auto pec = svc_.pool().post(&read_pool_op_))
211 : {
212 : // The pool is shutting down, or the system refused it a thread.
213 : // Nothing of this read went cross-thread, so it answers here
214 : // like the closed-descriptor and zero-length exits above rather
215 : // than through a completion the scheduler has to carry back.
216 6 : read_pool_op_.ref_.reset();
217 6 : op.stop_cb.reset();
218 6 : op.ex.on_work_finished();
219 6 : *ec = pec;
220 6 : *bytes_out = 0;
221 6 : op.cont.h = h;
222 6 : return dispatch_coro(ex, op.cont);
223 : }
224 545 : return std::noop_coroutine();
225 : }
226 :
227 : inline void
228 545 : posix_stream_file::do_read_work(pool_work_item* w) noexcept
229 : {
230 545 : auto* pw = static_cast<pool_op*>(w);
231 545 : auto* self = pw->file_;
232 545 : auto& op = self->read_op_;
233 :
234 545 : if (!op.cancelled.load(std::memory_order_acquire))
235 : {
236 : ssize_t n;
237 : do
238 : {
239 127 : n = file_preadv(
240 127 : op.fd, op.iovecs, op.iovec_count,
241 127 : static_cast<file_off_t>(op.offset));
242 : }
243 127 : while (n < 0 && errno == EINTR);
244 :
245 127 : if (n >= 0)
246 : {
247 120 : op.errn = 0;
248 120 : op.bytes_transferred = static_cast<std::size_t>(n);
249 :
250 : // The file may have been replaced while this ran; its
251 : // position belongs to whatever it now holds.
252 120 : std::lock_guard<std::mutex> lock(self->offset_mutex_);
253 120 : if (self->generation_ == op.generation)
254 115 : self->offset_ += static_cast<std::uint64_t>(n);
255 120 : }
256 : else
257 : {
258 7 : op.errn = errno;
259 7 : op.bytes_transferred = 0;
260 : }
261 : }
262 :
263 545 : op.impl_ptr = std::move(pw->ref_);
264 545 : self->svc_.post(&op);
265 545 : }
266 :
267 : inline std::coroutine_handle<>
268 91 : posix_stream_file::write_some(
269 : std::coroutine_handle<> h,
270 : capy::executor_ref ex,
271 : buffer_param param,
272 : std::stop_token token,
273 : std::error_code* ec,
274 : std::size_t* bytes_out)
275 : {
276 91 : auto& op = write_op_;
277 91 : op.reset();
278 91 : op.is_read = false;
279 :
280 : // Closed-object contract outranks the zero-length no-op.
281 91 : if (fd_ < 0)
282 : {
283 6 : *ec = make_error_code(std::errc::bad_file_descriptor);
284 6 : *bytes_out = 0;
285 6 : op.cont.h = h;
286 6 : return dispatch_coro(ex, op.cont);
287 : }
288 :
289 85 : capy::mutable_buffer bufs[max_buffers];
290 85 : op.iovec_count = static_cast<int>(param.copy_to(bufs, max_buffers));
291 :
292 85 : if (op.iovec_count == 0)
293 : {
294 2 : *ec = {};
295 2 : *bytes_out = 0;
296 2 : op.cont.h = h;
297 2 : return dispatch_coro(ex, op.cont);
298 : }
299 :
300 166 : for (int i = 0; i < op.iovec_count; ++i)
301 : {
302 83 : op.iovecs[i].iov_base = bufs[i].data();
303 83 : op.iovecs[i].iov_len = bufs[i].size();
304 : }
305 :
306 83 : op.h = h;
307 83 : op.ex = ex;
308 83 : op.ec_out = ec;
309 83 : op.bytes_out = bytes_out;
310 83 : op.start(token);
311 :
312 83 : op.fd = fd_;
313 83 : op.offset = offset_;
314 83 : op.generation = generation_;
315 :
316 83 : op.ex.on_work_started();
317 :
318 83 : write_pool_op_.file_ = this;
319 83 : write_pool_op_.ref_ = this->shared_from_this();
320 83 : write_pool_op_.func_ = &posix_stream_file::do_write_work;
321 83 : if (auto pec = svc_.pool().post(&write_pool_op_))
322 : {
323 : // The pool is shutting down, or the system refused it a thread.
324 : // Nothing of this write went cross-thread, so it answers here
325 : // like the closed-descriptor and zero-length exits above rather
326 : // than through a completion the scheduler has to carry back.
327 2 : write_pool_op_.ref_.reset();
328 2 : op.stop_cb.reset();
329 2 : op.ex.on_work_finished();
330 2 : *ec = pec;
331 2 : *bytes_out = 0;
332 2 : op.cont.h = h;
333 2 : return dispatch_coro(ex, op.cont);
334 : }
335 81 : return std::noop_coroutine();
336 : }
337 :
338 : inline void
339 81 : posix_stream_file::do_write_work(pool_work_item* w) noexcept
340 : {
341 81 : auto* pw = static_cast<pool_op*>(w);
342 81 : auto* self = pw->file_;
343 81 : auto& op = self->write_op_;
344 :
345 81 : if (!op.cancelled.load(std::memory_order_acquire))
346 : {
347 : ssize_t n;
348 : do
349 : {
350 29 : n = file_pwritev(
351 29 : op.fd, op.iovecs, op.iovec_count,
352 29 : static_cast<file_off_t>(op.offset));
353 : }
354 29 : while (n < 0 && errno == EINTR);
355 :
356 29 : if (n >= 0)
357 : {
358 22 : op.errn = 0;
359 22 : op.bytes_transferred = static_cast<std::size_t>(n);
360 :
361 : // The file may have been replaced while this ran; its
362 : // position belongs to whatever it now holds.
363 22 : std::lock_guard<std::mutex> lock(self->offset_mutex_);
364 22 : if (self->generation_ == op.generation)
365 22 : self->offset_ += static_cast<std::uint64_t>(n);
366 22 : }
367 : else
368 : {
369 7 : op.errn = errno;
370 7 : op.bytes_transferred = 0;
371 : }
372 : }
373 :
374 81 : op.impl_ptr = std::move(pw->ref_);
375 81 : self->svc_.post(&op);
376 81 : }
377 :
378 : } // namespace boost::corosio::detail
379 :
380 : #endif // BOOST_COROSIO_POSIX
381 :
382 : #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_STREAM_FILE_SERVICE_HPP
|