src/corosio/src/tls/detail/engine_driver.hpp
100.0% Lines (38 / 38)
100.0% Functions (13 / 13)
Functions (13)
Function
Calls
Lines
Blocks
boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::flush_output()
:198
51911x
100.0%
44.0%
boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::best_effort_flush()
:252
20x
100.0%
44.0%
boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::fill_input(unsigned long)
:260
3949x
100.0%
44.0%
boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::take_pending_flush_ec()
:316
90572x
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
2543x
100.0%
57.0%
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<16ul, false>)
:380
45177x
100.0%
42.0%
boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::do_write_some(boost::capy::detail::buffer_array<16ul, true>)
:471
45248x
100.0%
42.0%
boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::do_handshake(boost::corosio::tls_role)
:566
1861x
100.0%
44.0%
boost::corosio::detail::engine_driver<boost::corosio::detail::openssl::engine>::do_shutdown()
:641
147x
100.0%
44.0%
| Line | 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 | 51911x | 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 | 103822x | } | |
| 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 | 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 | 3949x | 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 | 7898x | } | |
| 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 | 90572x | std::error_code take_pending_flush_ec() noexcept | |
| 317 | { | ||
| 318 | 90572x | std::error_code ec = pending_flush_ec_; | |
| 319 | 90572x | pending_flush_ec_ = {}; | |
| 320 | 90572x | 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 | 2543x | engine_driver(capy::any_stream& s, tls_context ctx) | |
| 332 | 2543x | : s_(&s) | |
| 333 | 2543x | , ctx_(std::move(ctx)) | |
| 334 | { | ||
| 335 | 2543x | out_buf_.resize(buffer_size_); | |
| 336 | 2543x | } | |
| 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 | 45177x | 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 | 90354x | } | |
| 470 | |||
| 471 | 45248x | 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 | 90496x | } | |
| 565 | |||
| 566 | 1861x | 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 | 3722x | } | |
| 640 | |||
| 641 | 147x | 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 | 294x | } | |
| 695 | }; | ||
| 696 | |||
| 697 | } // namespace detail | ||
| 698 | |||
| 699 | } // namespace boost::corosio | ||
| 700 | |||
| 701 | #endif | ||
| 702 |