src/corosio/src/tls/detail/engine_driver.hpp

100.0% Lines (38/38) 100.0% List of functions (13/13) 100.0% Branches (8/8)
engine_driver.hpp
f(x) Functions (13)
Function Calls Lines Branches Blocks
boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::flush_output() :198 70218x 100.0% 100.0% 43.8% boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::best_effort_flush() :252 20x 100.0% 100.0% 43.8% boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::fill_input(unsigned long long) :260 34912x 100.0% 100.0% 43.8% boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::take_pending_flush_ec() :316 64328x 100.0% – 100.0% boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::engine_driver(boost::capy::any_stream&, boost::corosio::tls_context) :331 3063x 100.0% 100.0% 57.1% boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::rebind_stream(boost::capy::any_stream&) :351 3x 100.0% – 100.0% boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::set_hostname(std::basic_string_view<char, std::char_traits<char> >) :357 13x 100.0% – 100.0% boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::alpn_protocol() const :363 3x 100.0% – 100.0% boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::reset() :369 65x 100.0% – 100.0% boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::do_read_some(boost::capy::detail::buffer_array<16ull, false>) :380 32137x 100.0% 100.0% 42.1% boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::do_write_some(boost::capy::detail::buffer_array<16ull, true>) :471 32038x 100.0% 100.0% 42.1% boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::do_handshake(boost::corosio::tls_role) :566 2205x 100.0% 100.0% 43.8% boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::do_shutdown() :641 153x 100.0% 100.0% 43.8%
Line Branch TLA Hits Source Code
1 //
2 // Copyright (c) 2026 Steve Gerbino
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 SRC_TLS_DETAIL_ENGINE_DRIVER_HPP
11 #define SRC_TLS_DETAIL_ENGINE_DRIVER_HPP
12
13 #include <boost/corosio/tls_context.hpp>
14 #include <boost/corosio/tls_stream.hpp>
15
16 #include "src/tls/detail/engine_types.hpp"
17
18 #include <boost/capy/buffers.hpp>
19 #include <boost/capy/detail/buffer_array.hpp>
20 #include <boost/capy/ex/async_mutex.hpp>
21 #include <boost/capy/io/any_stream.hpp>
22 #include <boost/capy/io_task.hpp>
23 #include <boost/capy/task.hpp>
24 #include <boost/capy/write.hpp>
25
26 #include <concepts>
27 #include <cstddef>
28 #include <cstdint>
29 #include <cstring>
30 #include <limits>
31 #include <string>
32 #include <string_view>
33 #include <system_error>
34 #include <tuple>
35 #include <utility>
36 #include <vector>
37
38 // This header is instantiated in both backends' TUs, whose vendor
39 // headers cannot coexist (WolfSSL's OpenSSL-compat layer clashes
40 // with genuine OpenSSL declarations on some builds); it must stay
41 // vendor-free.
42 #ifdef OPENSSL_VERSION_NUMBER
43 #error engine_driver.hpp must stay vendor-header-free
44 #endif
45 #ifdef WOLFSSL_VERSION
46 #error engine_driver.hpp must stay vendor-header-free
47 #endif
48
49 /*
50 TLS Driver Architecture
51
52 TLS layer wrapping an underlying stream (via any_stream). One
53 read_some and one write_some may be in flight concurrently.
54
55 Engine / Driver Split. The engine (each backend's detail/engine.hpp)
56 owns the TLS session, its byte interface, and every library-result-
57 to-error-code decision; it is synchronous and transport-free. This
58 driver owns the transport, the per-direction claims, and all
59 buffering, and loops on engine verdicts:
60
61 1. Call eng_.perform(op, data, len)
62 2. output_then_retry: flush pending ciphertext, then retry
63 3. input: feed transport bytes into the engine, then retry
64 4. done / output_then_done: flush if pending, report the
65 (already mapped) ec
66
67 Renegotiation causes cross-direction I/O: a read may need to
68 write handshake data, a write may need to read. Each operation
69 services whatever direction the engine requests.
70
71 Full-Duplex Concurrency. rd_cm_ and wr_cm_ are separate transport
72 claims: a read parked on the network cannot stall a concurrent
73 write's flush, and vice versa. A cross-direction request (e.g. a
74 read needing to flush queued handshake output) takes the other
75 direction's claim itself.
76
77 read_gen_ counts successful deposits into the engine. Each engine
78 call snapshots it before running; if fill_input later finds the
79 generation has moved (a concurrent operation's flush suspension let
80 a deposit land in the meantime) it retries the engine instead of
81 issuing a redundant transport read.
82
83 The transport reads ciphertext straight into the engine's input
84 staging via input_area(), so no driver-side buffer or deposit copy
85 sits on the read path. The write path stays copy-based on purpose:
86 ciphertext drained from the engine into out_buf_[0, out_len_) is
87 decoupled from the engine's shared output staging, so a concurrent
88 reader can still emit its own records (an alert, a key update) while
89 a write is in flight. out_len_ holds only the abandoned tail a failed
90 or canceled write kept for retry; the length being written is a local,
91 never published, so a concurrent flush probe cannot mistake in-flight
92 bytes for a tail and park on the write claim behind a healthy write.
93
94 A read or write that fully satisfies the caller's buffer still owes
95 a trailing flush (queued handshake or session-ticket output must
96 reach the peer before the call can report success). If that flush
97 fails, the transfer already succeeded and the stream contract
98 forbids reporting both a full count and an error in the same
99 result, so the error is stashed in pending_flush_ec_ and surfaced
100 by the next read, write, or shutdown call instead.
101 */
102
103 namespace boost::corosio {
104
105 namespace detail {
106
107 /** The engine surface `engine_driver` drives.
108
109 A synchronous, transport-free TLS record engine: `perform`
110 advances one operation and returns a fully mapped verdict, while
111 `put_input` / `get_output` shuttle wire bytes. The handshake
112 hooks (`check_context`, `check_session`, `prepare`) let each
113 backend gate and configure a handshake at the protocol-mandated
114 points without the driver knowing the backend's session model.
115 */
116 template<class Engine>
117 concept tls_engine = requires(
118 Engine& e,
119 Engine const& ce,
120 engine_op op,
121 void* buf,
122 unsigned char const* in,
123 unsigned char* out,
124 std::size_t n,
125 tls_context const& ctx,
126 tls_role role,
127 std::string const& hostname,
128 std::string& alpn) {
129 { e.perform(op, buf, n) } -> std::same_as<engine_result>;
130 { e.put_input(in, n) } -> std::same_as<std::size_t>;
131 { e.input_area() } -> std::same_as<std::pair<unsigned char*, std::size_t>>;
132 { e.input_committed(n) };
133 { ce.pending_output() } -> std::same_as<std::size_t>;
134 { e.get_output(out, n) } -> std::same_as<std::size_t>;
135 { ce.received_shutdown() } -> std::same_as<bool>;
136 { ce.capture_alpn(alpn) };
137 { e.reset() };
138 { ce.check_context() } -> std::same_as<std::error_code>;
139 { ce.check_session() } -> std::same_as<std::error_code>;
140 { e.prepare(ctx, role, hostname) } -> std::same_as<std::error_code>;
141 };
142
143 /** Coroutine driver shared by every TLS backend.
144
145 Hosts the transport, the per-direction claims, the buffering, and
146 the four operation loops; the engine type supplies all TLS
147 mechanics. Each backend's stream impl instantiates this template
148 with its own engine so the driver logic exists exactly once.
149
150 @par Thread Safety
151 Distinct objects: Safe.@n
152 Shared objects: Unsafe.
153 */
154 template<tls_engine Engine>
155 class engine_driver
156 {
157 // Large enough to hold the largest possible TLS record.
158 static constexpr std::size_t buffer_size_ = std::size_t{17} * 1024;
159
160 capy::any_stream* s_;
161 tls_context ctx_;
162 Engine eng_;
163
164 // A handshake was attempted (successfully or not); the stream
165 // must be reset before the next handshake.
166 bool used_ = false;
167
168 // Per-stream SNI/verification hostname, set via set_hostname().
169 std::string hostname_;
170
171 // ALPN protocol negotiated during the handshake (empty if none).
172 std::string alpn_selected_;
173
174 std::vector<char> out_buf_;
175
176 // One transport claim per direction so a parked reader cannot block
177 // a writer's flush. read_gen_ counts deposits into the engine: each
178 // operation captures it at its engine call and passes it to
179 // fill_input, which retries the engine instead of issuing a stale
180 // transport read when a deposit landed anywhere after the engine's
181 // verdict (including during an intervening flush suspension).
182 capy::async_mutex rd_cm_;
183 capy::async_mutex wr_cm_;
184 std::uint64_t read_gen_ = 0;
185
186 // Engine ciphertext drained but not yet written: the abandoned
187 // tail a failed or canceled transport write kept for retry.
188 // Zero while a write is in flight (that length is a local).
189 std::size_t out_len_ = 0;
190
191 // A trailing flush that fails after the engine already accepted a
192 // full payload cannot be reported alongside n == size (the stream
193 // contract forbids error + full transfer, and capy::write's
194 // composed loop discards ec exactly in that case); the error is
195 // held here and surfaced by the next operation instead.
196 std::error_code pending_flush_ec_;
197
198
1/1
✓ Branch 5 → 6 taken 70218 times.
70218x capy::task<std::error_code> flush_output()
199 {
200 if (eng_.pending_output() == 0 && out_len_ == 0)
201 co_return std::error_code{};
202
203 auto [lec] = co_await wr_cm_.lock();
204 if (lec)
205 co_return lec;
206 capy::async_mutex::lock_guard wr_guard(&wr_cm_);
207
208 // Drain under the claim so records reach the wire in engine
209 // order even when both directions produce output.
210 while (eng_.pending_output() > 0 || out_len_ > 0)
211 {
212 // Take the abandoned tail (from an earlier failed write in
213 // this loop or a prior flush) into a local and append fresh
214 // engine ciphertext behind it, preserving record order.
215 // out_len_ drops to zero for the duration of the write: the
216 // in-flight length lives only in `n`, so a concurrent flush
217 // probe never sees these bytes as a tail and never parks on
218 // the write claim behind a healthy write.
219 std::size_t n = out_len_;
220 out_len_ = 0;
221 if (n < out_buf_.size())
222 n += eng_.get_output(
223 reinterpret_cast<unsigned char*>(out_buf_.data()) + n,
224 out_buf_.size() - n);
225 // The loop guard just confirmed pending bytes exist, so a
226 // drain failure here is unreachable in practice; fail loudly
227 // rather than silently drop already-accepted ciphertext.
228 if (n == 0) // LCOV_EXCL_LINE unreachable: pending bytes confirmed
229 co_return make_error_code(
230 std::errc::
231 no_buffer_space); // LCOV_EXCL_LINE unreachable: transport returned 0 with no error unreachable: pending bytes confirmed
232 auto [ec, wn] = co_await capy::write(
233 *s_, capy::const_buffer(out_buf_.data(), n));
234 if (ec)
235 {
236 // wn bytes already reached the peer; keep only the unsent
237 // remainder so a post-cancellation flush retry resends
238 // neither the delivered prefix nor loses the rest.
239 std::memmove(out_buf_.data(), out_buf_.data() + wn, n - wn);
240 out_len_ = n - wn;
241 co_return ec;
242 }
243 }
244 co_return std::error_code{};
245 140436x }
246
247 // Flushes on a fatal/alert exit path: the engine may have queued a
248 // close_notify or alert that should still reach the peer, but a
249 // caller already reporting the triggering error cannot also report
250 // this flush's own failure. flush_output() itself no-ops when
251 // nothing is pending.
252
1/1
✓ Branch 5 → 6 taken 20 times.
20x capy::task<void> best_effort_flush()
253 {
254 std::ignore = co_await flush_output();
255 40x }
256
257 // gen is the caller's read_gen_ snapshot from its engine call; a
258 // later snapshot would miss deposits landing while the caller's
259 // flush_output was suspended.
260
1/1
✓ Branch 5 → 6 taken 34912 times.
34912x capy::task<std::error_code> fill_input(std::uint64_t gen)
261 {
262 // Input already arrived since the engine's verdict; retry the
263 // engine rather than queue behind a re-parked reader.
264 if (read_gen_ != gen)
265 co_return std::error_code{};
266
267 auto [lec] = co_await rd_cm_.lock();
268 if (lec)
269 co_return lec;
270 capy::async_mutex::lock_guard rd_guard(&rd_cm_);
271
272 // Input arrived while we queued for the claim; the caller's
273 // engine retry consumes it.
274 if (read_gen_ != gen)
275 co_return std::error_code{};
276
277 // Read the transport straight into the engine's input staging:
278 // input_area() hands back the contiguous writable run, so no
279 // staging buffer or deposit copy sits between the socket and the
280 // engine.
281 auto [dst, cap] = eng_.input_area();
282 // The staging is full while the caller still wants input: its own
283 // engine retry must decrypt the staged record to free space, so
284 // report success to re-run it rather than issue a zero-length read.
285 if (cap == 0)
286 co_return std::error_code{};
287
288 auto [ec, n] = co_await s_->read_some(capy::mutable_buffer(dst, cap));
289
290 // ReadStream permits n>0 alongside ec (IOCP forwards
291 // bytes_transferred on failed completions; a canceled read can
292 // likewise deliver bytes); commit whatever arrived before
293 // surfacing ec, or the record stream loses wire bytes the
294 // transport already handed over.
295 if (n > 0)
296 {
297 eng_.input_committed(n);
298 ++read_gen_;
299 co_return ec;
300 }
301
302 if (ec)
303 co_return ec;
304
305 // The transport delivered nothing without an error, so it cannot
306 // make progress: fail loudly rather than spin the engine's input
307 // retry against a staging that will never fill.
308 co_return make_error_code(
309 std::errc::
310 no_buffer_space); // LCOV_EXCL_LINE unreachable: staging cannot stay empty
311 69824x }
312
313 // A prior read/write already reported its full transfer as success;
314 // the trailing flush error it deferred surfaces on the next call
315 // instead. Each entry point takes it exactly once.
316 64328x std::error_code take_pending_flush_ec() noexcept
317 {
318 64328x std::error_code ec = pending_flush_ec_;
319 64328x pending_flush_ec_ = {};
320 64328x return ec;
321 }
322
323 public:
324 /** Construct a driver over a transport stream.
325
326 @param s The transport; the caller keeps it alive and repoints
327 it on move via `rebind_stream`.
328 @param ctx The TLS context handed to the engine's handshake
329 preparation.
330 */
331 3063x engine_driver(capy::any_stream& s, tls_context ctx)
332 3063x : s_(&s)
333 3063x , ctx_(std::move(ctx))
334 {
335
1/1
✓ Branch 11 → 12 taken 3063 times.
3063x out_buf_.resize(buffer_size_);
336 3063x }
337
338 /// Return the engine for backend-specific setup.
339 Engine& engine() noexcept
340 {
341 return eng_;
342 }
343
344 /// Return the TLS context this driver was constructed with.
345 tls_context const& context() const noexcept
346 {
347 return ctx_;
348 }
349
350 /// Point the driver at the transport's post-move location.
351 3x void rebind_stream(capy::any_stream& s) noexcept
352 {
353 3x s_ = &s;
354 3x }
355
356 /// Set the hostname applied to the next client handshake.
357 13x void set_hostname(std::string_view hostname)
358 {
359 13x hostname_ = hostname;
360 13x }
361
362 /// Return the ALPN protocol negotiated by the last handshake.
363 3x std::string_view alpn_protocol() const noexcept
364 {
365 3x return alpn_selected_;
366 }
367
368 /// Reset the engine and every per-connection driver state.
369 65x void reset()
370 {
371 65x eng_.reset();
372
373 65x out_len_ = 0;
374
375 65x alpn_selected_.clear();
376 65x pending_flush_ec_ = {};
377 65x used_ = false;
378 65x }
379
380
1/1
✓ Branch 6 → 7 taken 32137 times.
32137x capy::io_task<std::size_t> do_read_some(
381 capy::detail::mutable_buffer_array<capy::detail::max_iovec_> buffers)
382 {
383 if (auto ec = take_pending_flush_ec())
384 co_return {ec, 0};
385
386 std::error_code ec;
387 std::size_t total_read = 0;
388 std::size_t const bufs_size = capy::buffer_size(buffers);
389
390 for (auto& buf : buffers)
391 {
392 char* const dest = static_cast<char*>(buf.data());
393 // Clamp for the engines' int-sized transfer API: a >=2 GiB
394 // buffer becomes a partial read, never a truncated length.
395 int const remaining = static_cast<int>(
396 (std::min)(buf.size(),
397 std::size_t((std::numeric_limits<int>::max)())));
398 if (remaining == 0)
399 continue;
400
401 // Exits by co_return: success after the first transferred
402 // chunk, or any terminal error.
403 while (true)
404 {
405 auto const gen = read_gen_;
406 auto r = eng_.perform(
407 engine_op::read, dest, static_cast<std::size_t>(remaining));
408
409 if (r.ec)
410 {
411 // Terminal, already mapped by the engine. eof (a
412 // received close_notify) arrives as a plain done:
413 // it queues no output, so no flush is needed.
414 if (r.want == engine_want::output_then_done)
415 co_await best_effort_flush();
416 co_return {r.ec, total_read};
417 }
418
419 if (r.want == engine_want::done ||
420 r.want == engine_want::output_then_done)
421 {
422 total_read += r.bytes;
423
424 // r.bytes > 0 already satisfies ReadStream's "at
425 // least one byte transferred" success condition;
426 // report now rather than loop for more (another
427 // engine call could park on input). out_len_ > 0
428 // covers a retained tail the engine cannot see.
429 if (r.want == engine_want::output_then_done || out_len_ > 0)
430 ec = co_await flush_output();
431 if (ec && total_read == bufs_size)
432 {
433 // First failure wins: concurrent directions share
434 // one transport failure domain; the earliest
435 // error is the meaningful one.
436 if (!pending_flush_ec_)
437 pending_flush_ec_ = ec;
438 ec = {};
439 }
440 co_return {ec, total_read};
441 }
442
443 if (r.want == engine_want::output_then_retry)
444 {
445 ec = co_await flush_output();
446 if (ec)
447 co_return {ec, total_read};
448 continue;
449 }
450
451 // want == input. Flush first: a retained tail (an
452 // earlier failed write) may hold bytes the peer needs
453 // before it will send more; no-op when nothing pends.
454 ec = co_await flush_output();
455 if (ec)
456 co_return {ec, total_read};
457
458 ec = co_await fill_input(gen);
459 if (ec)
460 {
461 ec = map_fill_error(
462 engine_op::read, ec, eng_.received_shutdown());
463 co_return {ec, total_read};
464 }
465 }
466 }
467
468 co_return {std::error_code{}, total_read};
469 64274x }
470
471
1/1
✓ Branch 6 → 7 taken 32038 times.
32038x capy::io_task<std::size_t> do_write_some(
472 capy::detail::const_buffer_array<capy::detail::max_iovec_> buffers)
473 {
474 if (auto ec = take_pending_flush_ec())
475 co_return {ec, 0};
476
477 std::error_code ec;
478 std::size_t total_written = 0;
479 std::size_t const bufs_size = capy::buffer_size(buffers);
480
481 for (auto const& buf : buffers)
482 {
483 // The engine only reads through this pointer for a write
484 // op; the cast satisfies perform's single transfer
485 // signature.
486 void* const src = const_cast<void*>(buf.data());
487 // Clamp for the engines' int-sized transfer API: a >=2 GiB
488 // buffer becomes a partial write, never a truncated length.
489 int const remaining = static_cast<int>(
490 (std::min)(buf.size(),
491 std::size_t((std::numeric_limits<int>::max)())));
492 if (remaining == 0)
493 continue;
494
495 // Exits by co_return: success after the first transferred
496 // chunk, or any terminal error.
497 while (true)
498 {
499 auto const gen = read_gen_;
500 auto r = eng_.perform(
501 engine_op::write, src, static_cast<std::size_t>(remaining));
502
503 if (r.ec)
504 {
505 // Terminal, already mapped. eof means the peer's
506 // close_notify WAS received (an announced close);
507 // it arrives as a plain done and needs no flush.
508 if (r.want == engine_want::output_then_done)
509 co_await best_effort_flush();
510 co_return {r.ec, total_written};
511 }
512
513 if (r.want == engine_want::done ||
514 r.want == engine_want::output_then_done)
515 {
516 total_written += r.bytes;
517
518 // r.bytes > 0 already satisfies WriteStream's "at
519 // least one byte transferred" success condition;
520 // report now rather than loop for more. out_len_ > 0
521 // covers a retained tail the engine cannot see.
522 if (r.want == engine_want::output_then_done || out_len_ > 0)
523 ec = co_await flush_output();
524 if (ec && total_written == bufs_size)
525 {
526 // First failure wins: concurrent directions share
527 // one transport failure domain; the earliest
528 // error is the meaningful one.
529 if (!pending_flush_ec_)
530 pending_flush_ec_ = ec;
531 ec = {};
532 }
533 co_return {ec, total_written};
534 }
535
536 if (r.want == engine_want::output_then_retry)
537 {
538 ec = co_await flush_output();
539 if (ec)
540 co_return {ec, total_written};
541 continue;
542 }
543
544 // want == input: a write can need a read mid-rekey; the
545 // transport eof this surfaces means the same thing it
546 // means on the read path, so map it the same way. Flush
547 // first for the same retained-tail reason as the read
548 // path.
549 ec = co_await flush_output();
550 if (ec)
551 co_return {ec, total_written};
552
553 ec = co_await fill_input(gen);
554 if (ec)
555 {
556 ec = map_fill_error(
557 engine_op::write, ec, eng_.received_shutdown());
558 co_return {ec, total_written};
559 }
560 }
561 }
562
563 co_return {std::error_code{}, total_written};
564 64076x }
565
566
1/1
✓ Branch 5 → 6 taken 2205 times.
2205x capy::io_task<> do_handshake(tls_role role)
567 {
568 // Refuse the handshake while the engine's configuration is
569 // unusable, before consuming any per-connection state.
570 if (auto cec = eng_.check_context())
571 co_return {cec};
572
573 if (used_)
574 reset();
575
576 // reset() may not have restored a usable session; refuse to
577 // hand out a handshake on it.
578 if (auto sec = eng_.check_session())
579 co_return {sec};
580
581 // A failed attempt leaves the session in a dead state; any
582 // attempt, not just a completed handshake, consumes the stream
583 // so the next handshake() starts fresh.
584 used_ = true;
585
586 if (auto pec = eng_.prepare(ctx_, role, hostname_))
587 co_return {pec};
588
589 auto const op = role == tls_role::client ? engine_op::handshake_client
590 : engine_op::handshake_server;
591
592 std::error_code ec;
593
594 while (true)
595 {
596 auto const gen = read_gen_;
597 auto r = eng_.perform(op, nullptr, 0);
598
599 if (r.ec)
600 {
601 if (r.want == engine_want::output_then_done)
602 co_await best_effort_flush();
603 co_return {r.ec};
604 }
605
606 if (r.want == engine_want::done ||
607 r.want == engine_want::output_then_done)
608 {
609 eng_.capture_alpn(alpn_selected_);
610 ec = co_await flush_output();
611 co_return {ec};
612 }
613
614 if (r.want == engine_want::output_then_retry)
615 {
616 // Must flush (e.g. ClientHello) before reading the
617 // peer's reply.
618 ec = co_await flush_output();
619 if (ec)
620 co_return {ec};
621 continue;
622 }
623
624 // want == input. The flush is a no-op unless a retained
625 // tail pends; the fill error passes through map_fill_error
626 // unchanged (no close is clean before the session is
627 // established).
628 ec = co_await flush_output();
629 if (ec)
630 co_return {ec};
631
632 ec = co_await fill_input(gen);
633 if (ec)
634 {
635 ec = map_fill_error(op, ec, eng_.received_shutdown());
636 co_return {ec};
637 }
638 }
639 4410x }
640
641
1/1
✓ Branch 5 → 6 taken 153 times.
153x capy::io_task<> do_shutdown()
642 {
643 if (auto ec = take_pending_flush_ec())
644 co_return {ec};
645
646 std::error_code ec;
647
648 while (true)
649 {
650 auto const gen = read_gen_;
651 auto r = eng_.perform(engine_op::shutdown, nullptr, 0);
652
653 if (r.ec)
654 {
655 if (r.want == engine_want::output_then_done)
656 co_await best_effort_flush();
657 co_return {r.ec};
658 }
659
660 if (r.want == engine_want::done ||
661 r.want == engine_want::output_then_done)
662 {
663 // Covers both the bidirectional-complete result and the
664 // engine's shutdown-after-close success mapping.
665 ec = co_await flush_output();
666 co_return {ec};
667 }
668
669 if (r.want == engine_want::output_then_retry)
670 {
671 // Sends our close_notify before parking for the peer's.
672 ec = co_await flush_output();
673 if (ec)
674 co_return {ec};
675 continue;
676 }
677
678 // want == input: awaiting the peer's close_notify. It may
679 // already have been deposited (and consumed by a concurrent
680 // reader) during a flush; fill_input then returns without
681 // reading and the loop retries the engine instead of parking.
682 ec = co_await flush_output();
683 if (ec)
684 co_return {ec};
685
686 ec = co_await fill_input(gen);
687 if (ec)
688 {
689 ec = map_fill_error(
690 engine_op::shutdown, ec, eng_.received_shutdown());
691 co_return {ec};
692 }
693 }
694 306x }
695 };
696
697 } // namespace detail
698
699 } // namespace boost::corosio
700
701 #endif
702