include/boost/corosio/native/detail/posix/posix_stream_file_service.hpp

100.0% Lines (153 / 153) 100.0% Functions (15 / 15)
posix_stream_file_service.hpp
f(x) Functions (15)
Function Calls Lines Blocks
boost::corosio::detail::posix_stream_file_service::posix_stream_file_service(boost::capy::execution_context&) :38 219x 100.0% 86.0% boost::corosio::detail::posix_stream_file_service::~posix_stream_file_service() :44 438x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::construct() :50 339x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::destroy(boost::corosio::io_object::implementation*) :55 337x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::close(boost::corosio::io_object::handle&) :63 645x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::open_file(boost::corosio::stream_file::implementation&, std::filesystem::__cxx11::path const&, boost::corosio::file_base::flags) :73 322x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::shutdown() :85 219x 100.0% 81.0% boost::corosio::detail::posix_stream_file_service::shutdown()::{lambda(boost::corosio::detail::posix_stream_file*)#1}::operator()(boost::corosio::detail::posix_stream_file*) const :94 4x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::post(boost::corosio::detail::scheduler_op*) :107 260x 100.0% 100.0% boost::corosio::detail::posix_stream_file_service::pool() :135 268x 100.0% 100.0% boost::corosio::detail::posix_stream_file::read_some(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :151 193x 100.0% 97.0% boost::corosio::detail::posix_stream_file::do_read_work(boost::corosio::detail::pool_work_item*) :218 179x 100.0% 100.0% boost::corosio::detail::posix_stream_file::write_some(std::__n4861::coroutine_handle<void>, boost::capy::executor_ref, boost::corosio::buffer_param, std::stop_token, std::error_code*, unsigned long*) :253 91x 100.0% 97.0% boost::corosio::detail::posix_stream_file::do_write_work(boost::corosio::detail::pool_work_item*) :320 81x 100.0% 100.0% boost::corosio::detail::posix_stream_file::retire() :355 337x 100.0% 100.0%
Line TLA Hits 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 219x explicit posix_stream_file_service(capy::execution_context& ctx)
39 657x : sched_(&get_scheduler(ctx))
40 219x , pool_(ctx)
41 {
42 219x }
43
44 438x ~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 339x io_object::implementation* construct() override
51 {
52 339x return object_pool_.acquire(*this);
53 }
54
55 337x void destroy(io_object::implementation* p) override
56 {
57 337x auto& impl = static_cast<posix_stream_file&>(*p);
58 337x impl.cancel();
59 337x impl.close_file();
60 337x release(&impl);
61 337x }
62
63 645x void close(io_object::handle& h) override
64 {
65 645x if (h.get())
66 {
67 645x auto& impl = static_cast<posix_stream_file&>(*h.get());
68 645x impl.cancel();
69 645x impl.close_file();
70 }
71 645x }
72
73 322x 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 322x if (sched_->scheduler_locking_disabled())
81 2x return std::make_error_code(std::errc::operation_not_supported);
82 320x return static_cast<posix_stream_file&>(impl).open_file(path, mode);
83 }
84
85 219x 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 219x std::vector<posix_stream_file*> live;
93 219x object_pool_.shutdown(
94 219x [&](posix_stream_file* f)
95 {
96 4x acquire(f);
97 4x live.push_back(f);
98 4x });
99 223x for (auto* f : live)
100 {
101 4x f->cancel();
102 4x f->close_file();
103 4x release(f);
104 }
105 219x }
106
107 260x void post(scheduler_op* op)
108 {
109 260x sched_->post(op);
110 260x }
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 268x thread_pool& pool()
136 {
137 268x 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 193x 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 193x auto& op = read_op_;
160 193x op.reset();
161 193x op.is_read = true;
162
163 // Closed-object contract outranks the zero-length no-op.
164 193x if (fd_ < 0)
165 {
166 6x *ec = make_error_code(std::errc::bad_file_descriptor);
167 6x *bytes_out = 0;
168 6x op.cont.h = h;
169 6x return dispatch_coro(ex, op.cont);
170 }
171
172 187x capy::mutable_buffer bufs[max_buffers];
173 187x op.iovec_count = static_cast<int>(param.copy_to(bufs, max_buffers));
174
175 187x if (op.iovec_count == 0)
176 {
177 2x *ec = {};
178 2x *bytes_out = 0;
179 2x op.cont.h = h;
180 2x return dispatch_coro(ex, op.cont);
181 }
182
183 370x for (int i = 0; i < op.iovec_count; ++i)
184 {
185 185x op.iovecs[i].iov_base = bufs[i].data();
186 185x op.iovecs[i].iov_len = bufs[i].size();
187 }
188
189 185x op.h = h;
190 185x op.ex = ex;
191 185x op.ec_out = ec;
192 185x op.bytes_out = bytes_out;
193 185x op.start(token);
194
195 185x op.ex.on_work_started();
196
197 185x read_pool_op_.file_ = this;
198 185x read_pool_op_.ref_ = detail::object_ref(this);
199 185x read_pool_op_.func_ = &posix_stream_file::do_read_work;
200 185x 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 6x read_pool_op_.ref_.reset();
207 6x op.stop_cb.reset();
208 6x op.ex.on_work_finished();
209 6x *ec = pec;
210 6x *bytes_out = 0;
211 6x op.cont.h = h;
212 6x return dispatch_coro(ex, op.cont);
213 }
214 179x return std::noop_coroutine();
215 }
216
217 inline void
218 179x posix_stream_file::do_read_work(pool_work_item* w) noexcept
219 {
220 179x auto* pw = static_cast<pool_op*>(w);
221 179x auto* self = pw->file_;
222 179x auto& op = self->read_op_;
223
224 179x if (!op.cancelled.load(std::memory_order_acquire))
225 {
226 ssize_t n;
227 do
228 {
229 254x n = ::preadv(
230 127x self->fd_, op.iovecs, op.iovec_count,
231 127x static_cast<off_t>(self->offset_));
232 }
233 127x while (n < 0 && errno == EINTR);
234
235 127x if (n >= 0)
236 {
237 120x op.errn = 0;
238 120x op.bytes_transferred = static_cast<std::size_t>(n);
239 120x self->offset_ += static_cast<std::uint64_t>(n);
240 }
241 else
242 {
243 7x op.errn = errno;
244 7x op.bytes_transferred = 0;
245 }
246 }
247
248 179x op.object_ref_ = std::move(pw->ref_);
249 179x self->svc_.post(&op);
250 179x }
251
252 inline std::coroutine_handle<>
253 91x 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 91x auto& op = write_op_;
262 91x op.reset();
263 91x op.is_read = false;
264
265 // Closed-object contract outranks the zero-length no-op.
266 91x if (fd_ < 0)
267 {
268 6x *ec = make_error_code(std::errc::bad_file_descriptor);
269 6x *bytes_out = 0;
270 6x op.cont.h = h;
271 6x return dispatch_coro(ex, op.cont);
272 }
273
274 85x capy::mutable_buffer bufs[max_buffers];
275 85x op.iovec_count = static_cast<int>(param.copy_to(bufs, max_buffers));
276
277 85x if (op.iovec_count == 0)
278 {
279 2x *ec = {};
280 2x *bytes_out = 0;
281 2x op.cont.h = h;
282 2x return dispatch_coro(ex, op.cont);
283 }
284
285 166x for (int i = 0; i < op.iovec_count; ++i)
286 {
287 83x op.iovecs[i].iov_base = bufs[i].data();
288 83x op.iovecs[i].iov_len = bufs[i].size();
289 }
290
291 83x op.h = h;
292 83x op.ex = ex;
293 83x op.ec_out = ec;
294 83x op.bytes_out = bytes_out;
295 83x op.start(token);
296
297 83x op.ex.on_work_started();
298
299 83x write_pool_op_.file_ = this;
300 83x write_pool_op_.ref_ = detail::object_ref(this);
301 83x write_pool_op_.func_ = &posix_stream_file::do_write_work;
302 83x 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 2x write_pool_op_.ref_.reset();
309 2x op.stop_cb.reset();
310 2x op.ex.on_work_finished();
311 2x *ec = pec;
312 2x *bytes_out = 0;
313 2x op.cont.h = h;
314 2x return dispatch_coro(ex, op.cont);
315 }
316 81x return std::noop_coroutine();
317 }
318
319 inline void
320 81x posix_stream_file::do_write_work(pool_work_item* w) noexcept
321 {
322 81x auto* pw = static_cast<pool_op*>(w);
323 81x auto* self = pw->file_;
324 81x auto& op = self->write_op_;
325
326 81x if (!op.cancelled.load(std::memory_order_acquire))
327 {
328 ssize_t n;
329 do
330 {
331 62x n = ::pwritev(
332 31x self->fd_, op.iovecs, op.iovec_count,
333 31x static_cast<off_t>(self->offset_));
334 }
335 31x while (n < 0 && errno == EINTR);
336
337 31x if (n >= 0)
338 {
339 24x op.errn = 0;
340 24x op.bytes_transferred = static_cast<std::size_t>(n);
341 24x self->offset_ += static_cast<std::uint64_t>(n);
342 }
343 else
344 {
345 7x op.errn = errno;
346 7x op.bytes_transferred = 0;
347 }
348 }
349
350 81x op.object_ref_ = std::move(pw->ref_);
351 81x self->svc_.post(&op);
352 81x }
353
354 inline void
355 337x posix_stream_file::retire() noexcept
356 {
357 337x svc_.object_pool_.recycle(this);
358 337x }
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
365