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