TLA Line data Source code
1 : //
2 : // Copyright (c) 2025 Vinnie Falco (vinnie.falco@gmail.com)
3 : // Copyright (c) 2026 Michael Vandeberg
4 : //
5 : // Distributed under the Boost Software License, Version 1.0. (See accompanying
6 : // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7 : //
8 : // Official repository: https://github.com/cppalliance/capy
9 : //
10 :
11 : #ifndef BOOST_CAPY_IO_ANY_READ_STREAM_HPP
12 : #define BOOST_CAPY_IO_ANY_READ_STREAM_HPP
13 :
14 : #include <boost/capy/detail/config.hpp>
15 : #include <boost/capy/detail/await_suspend_helper.hpp>
16 : #include <boost/capy/buffers.hpp>
17 : #include <boost/capy/detail/buffer_array.hpp>
18 : #include <boost/capy/concept/io_awaitable.hpp>
19 : #include <boost/capy/concept/read_stream.hpp>
20 : #include <boost/capy/ex/io_env.hpp>
21 : #include <boost/capy/io_result.hpp>
22 :
23 : #include <concepts>
24 : #include <coroutine>
25 : #include <cstddef>
26 : #include <exception>
27 : #include <new>
28 : #include <span>
29 : #include <stop_token>
30 : #include <system_error>
31 : #include <utility>
32 :
33 : namespace boost {
34 : namespace capy {
35 :
36 : /** Dispatches `read_some` through a type-erased vtable, using preallocated awaitable storage.
37 :
38 : This class provides type erasure for any type satisfying the
39 : @ref ReadStream concept, enabling runtime polymorphism for
40 : read operations. It uses cached awaitable storage to achieve
41 : zero steady-state allocation after construction.
42 :
43 : The wrapper supports two construction modes:
44 : - **Owning**: Pass by value to transfer ownership. The wrapper
45 : allocates storage and owns the stream.
46 : - **Reference**: Pass a pointer to wrap without ownership. The
47 : pointed-to stream must outlive this wrapper.
48 :
49 : @par Awaitable Preallocation
50 : The constructor preallocates storage for the type-erased awaitable.
51 : This reserves all virtual address space at server startup
52 : so memory usage can be measured up front, rather than
53 : allocating piecemeal as traffic arrives.
54 :
55 : @par Immediate Completion
56 : When the underlying stream's awaitable reports ready immediately
57 : (e.g. buffered data already available), the wrapper skips
58 : coroutine suspension entirely and returns the result inline.
59 :
60 : @par Thread Safety
61 : Not thread-safe. Concurrent operations on the same wrapper
62 : are undefined behavior.
63 :
64 : @par Example
65 : @par !example example
66 :
67 :
68 : @see any_write_stream, any_stream, ReadStream
69 : */
70 : class any_read_stream
71 : {
72 : struct vtable;
73 :
74 : template<ReadStream S>
75 : struct vtable_for_impl;
76 :
77 : // ordered for cache line coherence
78 : void* stream_ = nullptr;
79 : vtable const* vt_ = nullptr;
80 : void* cached_awaitable_ = nullptr;
81 : void* storage_ = nullptr;
82 : bool awaitable_active_ = false;
83 :
84 : public:
85 : /** Destructor.
86 :
87 : Destroys the owned stream (if any) and releases the cached
88 : awaitable storage.
89 : */
90 : ~any_read_stream();
91 :
92 : /** Construct a default instance.
93 :
94 : Constructs an empty wrapper. @ref has_value and `operator bool`
95 : report the empty state; calling @ref read_some before the
96 : wrapper holds a stream is undefined behavior.
97 : */
98 HIT 4 : any_read_stream() = default;
99 :
100 : /** Non-copyable.
101 :
102 : The awaitable cache is per-instance and cannot be shared.
103 :
104 : @param other The wrapper that would be copied.
105 : */
106 : any_read_stream(any_read_stream const& other) = delete;
107 :
108 : /** Copy assignment is disabled.
109 :
110 : The awaitable cache is per-instance and cannot be shared.
111 :
112 : @param other The wrapper that would be assigned from.
113 :
114 : @return A reference to `*this`.
115 : */
116 : any_read_stream& operator=(any_read_stream const& other) = delete;
117 :
118 : /** Construct by moving.
119 :
120 : Transfers ownership of the wrapped stream (if owned) and
121 : cached awaitable storage from `other`. After the move, `other` is
122 : in a default-constructed state.
123 :
124 : @param other The wrapper to move from.
125 : */
126 4 : any_read_stream(any_read_stream&& other) noexcept
127 4 : : stream_(std::exchange(other.stream_, nullptr))
128 4 : , vt_(std::exchange(other.vt_, nullptr))
129 4 : , cached_awaitable_(std::exchange(other.cached_awaitable_, nullptr))
130 4 : , storage_(std::exchange(other.storage_, nullptr))
131 4 : , awaitable_active_(std::exchange(other.awaitable_active_, false))
132 : {
133 4 : }
134 :
135 : /** Assign by moving.
136 :
137 : Destroys any owned stream and releases existing resources,
138 : then transfers ownership from `other`.
139 :
140 : @param other The wrapper to move from.
141 : @return Reference to this wrapper.
142 : */
143 : any_read_stream&
144 : operator=(any_read_stream&& other) noexcept;
145 :
146 : /** Construct by taking ownership of a ReadStream.
147 :
148 : Allocates storage and moves the stream into this wrapper.
149 : The wrapper owns the stream and destroys it.
150 :
151 : @param s The stream to take ownership of.
152 : */
153 : template<ReadStream S>
154 : requires (!std::same_as<std::decay_t<S>, any_read_stream>)
155 : any_read_stream(S s);
156 :
157 : /** Construct by wrapping a ReadStream without ownership.
158 :
159 : Wraps the given stream by pointer. The stream must remain
160 : valid for the lifetime of this wrapper.
161 :
162 : @param s Pointer to the stream to wrap.
163 : */
164 : template<ReadStream S>
165 : any_read_stream(S* s);
166 :
167 : /** Check if the wrapper contains a valid stream.
168 :
169 : @return `true` if wrapping a stream, `false` if default-constructed
170 : or moved-from.
171 : */
172 : bool
173 31 : has_value() const noexcept
174 : {
175 31 : return stream_ != nullptr;
176 : }
177 :
178 : /** Check if the wrapper contains a valid stream.
179 :
180 : @return `true` if wrapping a stream, `false` if default-constructed
181 : or moved-from.
182 : */
183 : explicit
184 3 : operator bool() const noexcept
185 : {
186 3 : return has_value();
187 : }
188 :
189 : /** Initiate an asynchronous read operation.
190 :
191 : Reads data into the provided buffer sequence. The operation
192 : completes when at least one byte is read, or an error
193 : occurs.
194 :
195 : @param buffers The buffer sequence to read into. Passed by
196 : value to ensure the sequence lives in the coroutine frame
197 : across suspension points.
198 :
199 : @return An awaitable that await-returns `(error_code,std::size_t)`.
200 :
201 : @par Immediate Completion
202 : The operation completes immediately without suspending
203 : the calling coroutine when the underlying stream's
204 : awaitable reports immediate readiness via `await_ready`.
205 :
206 : @note This is a partial operation and may not process the
207 : entire buffer sequence. Use the composed @ref read algorithm
208 : for guaranteed complete transfer.
209 :
210 : @par Preconditions
211 : The wrapper must contain a valid stream (`has_value() == true`).
212 :
213 : @par After an Error
214 : A subsequent call is permitted. The wrapper forwards directly
215 : to the underlying stream, imposing no stricter rule than
216 : @ref ReadStream.
217 : */
218 : template<MutableBufferSequence MB>
219 : auto
220 : read_some(MB buffers);
221 :
222 : protected:
223 : /** Rebind to a new stream after move.
224 :
225 : Updates the internal pointer to reference a new stream object.
226 : Used by owning wrappers after move assignment when the owned
227 : object has moved to a new location.
228 :
229 : @param new_stream The new stream to bind to. Must be the same
230 : type as the original stream.
231 :
232 : @note Terminates if called with a stream of different type
233 : than the original.
234 : */
235 : template<ReadStream S>
236 : void
237 : rebind(S& new_stream) noexcept
238 : {
239 : if(vt_ != &vtable_for_impl<S>::value)
240 : std::terminate();
241 : stream_ = &new_stream;
242 : }
243 : };
244 :
245 : struct any_read_stream::vtable
246 : {
247 : // ordered by call frequency for cache line coherence
248 : void (*construct_awaitable)(
249 : void* stream,
250 : void* storage,
251 : std::span<mutable_buffer const> buffers);
252 : bool (*await_ready)(void*);
253 : std::coroutine_handle<> (*await_suspend)(void*, std::coroutine_handle<>, io_env const*);
254 : io_result<std::size_t> (*await_resume)(void*);
255 : void (*destroy_awaitable)(void*) noexcept;
256 : std::size_t awaitable_size;
257 : std::size_t awaitable_align;
258 : void (*destroy)(void*) noexcept;
259 : };
260 :
261 : template<ReadStream S>
262 : struct any_read_stream::vtable_for_impl
263 : {
264 : using Awaitable = decltype(std::declval<S&>().read_some(
265 : std::span<mutable_buffer const>{}));
266 :
267 : static void
268 4 : do_destroy_impl(void* stream) noexcept
269 : {
270 4 : static_cast<S*>(stream)->~S();
271 4 : }
272 :
273 : static void
274 103 : construct_awaitable_impl(
275 : void* stream,
276 : void* storage,
277 : std::span<mutable_buffer const> buffers)
278 : {
279 103 : auto& s = *static_cast<S*>(stream);
280 103 : ::new(storage) Awaitable(s.read_some(buffers));
281 103 : }
282 :
283 : static constexpr vtable value = {
284 : &construct_awaitable_impl,
285 103 : +[](void* p) {
286 103 : return static_cast<Awaitable*>(p)->await_ready();
287 : },
288 77 : +[](void* p, std::coroutine_handle<> h, io_env const* env) {
289 77 : return detail::call_await_suspend(
290 77 : static_cast<Awaitable*>(p), h, env);
291 : },
292 101 : +[](void* p) {
293 101 : return static_cast<Awaitable*>(p)->await_resume();
294 : },
295 115 : +[](void* p) noexcept {
296 26 : static_cast<Awaitable*>(p)->~Awaitable();
297 : },
298 : sizeof(Awaitable),
299 : alignof(Awaitable),
300 : &do_destroy_impl
301 : };
302 : };
303 :
304 : inline
305 123 : any_read_stream::~any_read_stream()
306 : {
307 123 : if(storage_)
308 : {
309 3 : vt_->destroy(stream_);
310 3 : ::operator delete(storage_);
311 : }
312 123 : if(cached_awaitable_)
313 : {
314 106 : if(awaitable_active_)
315 1 : vt_->destroy_awaitable(cached_awaitable_);
316 106 : ::operator delete(cached_awaitable_);
317 : }
318 123 : }
319 :
320 : inline any_read_stream&
321 10 : any_read_stream::operator=(any_read_stream&& other) noexcept
322 : {
323 10 : if(this != &other)
324 : {
325 10 : if(storage_)
326 : {
327 1 : vt_->destroy(stream_);
328 1 : ::operator delete(storage_);
329 : }
330 10 : if(cached_awaitable_)
331 : {
332 4 : if(awaitable_active_)
333 1 : vt_->destroy_awaitable(cached_awaitable_);
334 4 : ::operator delete(cached_awaitable_);
335 : }
336 10 : stream_ = std::exchange(other.stream_, nullptr);
337 10 : vt_ = std::exchange(other.vt_, nullptr);
338 10 : cached_awaitable_ = std::exchange(other.cached_awaitable_, nullptr);
339 10 : storage_ = std::exchange(other.storage_, nullptr);
340 10 : awaitable_active_ = std::exchange(other.awaitable_active_, false);
341 : }
342 10 : return *this;
343 : }
344 :
345 : template<ReadStream S>
346 : requires (!std::same_as<std::decay_t<S>, any_read_stream>)
347 5 : any_read_stream::any_read_stream(S s)
348 5 : : vt_(&vtable_for_impl<S>::value)
349 : {
350 : struct guard {
351 : any_read_stream* self;
352 : bool committed = false;
353 5 : ~guard() {
354 5 : if(!committed && self->storage_) {
355 1 : if(self->stream_)
356 : self->vt_->destroy(self->stream_); // LCOV_EXCL_LINE OOM rollback: only when the cached-awaitable allocation throws
357 1 : ::operator delete(self->storage_);
358 1 : self->storage_ = nullptr;
359 1 : self->stream_ = nullptr;
360 : }
361 5 : }
362 5 : } g{this};
363 :
364 5 : storage_ = ::operator new(sizeof(S));
365 5 : stream_ = ::new(storage_) S(std::move(s));
366 :
367 : // Preallocate the awaitable storage
368 4 : cached_awaitable_ = ::operator new(vt_->awaitable_size);
369 :
370 4 : g.committed = true;
371 5 : }
372 :
373 : template<ReadStream S>
374 106 : any_read_stream::any_read_stream(S* s)
375 106 : : stream_(s)
376 106 : , vt_(&vtable_for_impl<S>::value)
377 : {
378 : // Preallocate the awaitable storage
379 106 : cached_awaitable_ = ::operator new(vt_->awaitable_size);
380 106 : }
381 :
382 : template<MutableBufferSequence MB>
383 : auto
384 103 : any_read_stream::read_some(MB buffers)
385 : {
386 : // VFALCO in theory, we could use if constexpr to detect a
387 : // span and then pass that through to read_some without the array
388 : // LCOV_EXCL_START read_some awaitable: exercised by tests, but the
389 : // coverage tooling reports its templated body uncovered per-instantiation
390 : struct awaitable
391 : {
392 : any_read_stream* self_;
393 : detail::mutable_buffer_array<detail::max_iovec_> ba_;
394 :
395 : bool
396 : await_ready()
397 : {
398 : self_->vt_->construct_awaitable(
399 : self_->stream_,
400 : self_->cached_awaitable_,
401 : ba_.to_span());
402 : self_->awaitable_active_ = true;
403 :
404 : return self_->vt_->await_ready(
405 : self_->cached_awaitable_);
406 : }
407 :
408 : std::coroutine_handle<>
409 : await_suspend(std::coroutine_handle<> h, io_env const* env)
410 : {
411 : return self_->vt_->await_suspend(
412 : self_->cached_awaitable_, h, env);
413 : }
414 :
415 : [[nodiscard]] io_result<std::size_t>
416 : await_resume()
417 : {
418 : struct guard {
419 : any_read_stream* self;
420 : ~guard() {
421 : self->vt_->destroy_awaitable(self->cached_awaitable_);
422 : self->awaitable_active_ = false;
423 : }
424 : } g{self_};
425 : return self_->vt_->await_resume(
426 : self_->cached_awaitable_);
427 : }
428 : };
429 : // LCOV_EXCL_STOP
430 : return awaitable{this,
431 103 : detail::mutable_buffer_array<detail::max_iovec_>(buffers)};
432 103 : }
433 :
434 : } // namespace capy
435 : } // namespace boost
436 :
437 : #endif
|