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