Skip to main content

tokio_quiche/http3/driver/
mod.rs

1// Copyright (C) 2025, Cloudflare, Inc.
2// All rights reserved.
3//
4// Redistribution and use in source and binary forms, with or without
5// modification, are permitted provided that the following conditions are
6// met:
7//
8//     * Redistributions of source code must retain the above copyright notice,
9//       this list of conditions and the following disclaimer.
10//
11//     * Redistributions in binary form must reproduce the above copyright
12//       notice, this list of conditions and the following disclaimer in the
13//       documentation and/or other materials provided with the distribution.
14//
15// THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS
16// IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO,
17// THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR
18// PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR
19// CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL,
20// EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO,
21// PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR
22// PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF
23// LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING
24// NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
25// SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
26
27mod client;
28/// Wrapper for running HTTP/3 connections.
29pub mod connection;
30mod datagram;
31// `DriverHooks` must stay private to prevent users from creating their own
32// H3Drivers.
33mod hooks;
34mod server;
35mod streams;
36#[cfg(test)]
37pub mod test_utils;
38#[cfg(test)]
39mod tests;
40mod waiting_streams;
41
42use std::collections::BTreeMap;
43use std::error::Error;
44use std::fmt;
45use std::marker::PhantomData;
46use std::sync::Arc;
47use std::time::Instant;
48
49use bytes::BufMut as _;
50use bytes::Bytes;
51use bytes::BytesMut;
52use datagram_socket::DgramBuffer;
53use datagram_socket::StreamClosureKind;
54use foundations::telemetry::log;
55use futures::FutureExt;
56use quiche::h3;
57use quiche::h3::WireErrorCode;
58use tokio::select;
59use tokio::sync::mpsc;
60use tokio::sync::mpsc::error::TryRecvError;
61use tokio::sync::mpsc::error::TrySendError;
62use tokio::sync::mpsc::UnboundedReceiver;
63use tokio::sync::mpsc::UnboundedSender;
64use tokio_stream::StreamExt;
65use tokio_util::sync::PollSender;
66
67use self::hooks::DriverHooks;
68use self::hooks::InboundHeaders;
69use self::streams::FlowCtx;
70use self::streams::HaveUpstreamCapacity;
71use self::streams::ReceivedDownstreamData;
72use self::streams::StreamCtx;
73use self::streams::StreamReady;
74use self::waiting_streams::WaitingStreams;
75use crate::http3::settings::Http3Settings;
76use crate::http3::H3AuditStats;
77use crate::metrics::Metrics;
78use crate::quic::HandshakeError;
79use crate::quic::HandshakeInfo;
80use crate::quic::QuicCommand;
81use crate::quic::QuicheConnection;
82use crate::ApplicationOverQuic;
83use crate::QuicResult;
84
85pub use self::client::ClientEventStream;
86pub use self::client::ClientH3Command;
87pub use self::client::ClientH3Controller;
88pub use self::client::ClientH3Driver;
89pub use self::client::ClientH3Event;
90pub use self::client::ClientRequestSender;
91pub use self::client::NewClientRequest;
92pub use self::server::IsInEarlyData;
93pub use self::server::RawPriorityValue;
94pub use self::server::ServerEventStream;
95pub use self::server::ServerH3Command;
96pub use self::server::ServerH3Controller;
97pub use self::server::ServerH3Driver;
98pub use self::server::ServerH3Event;
99
100// The default priority for HTTP/3 responses if the application didn't provide
101// one.
102const DEFAULT_PRIO: h3::Priority = h3::Priority::new(3, true);
103
104// For a stream use a channel with 16 entries, which works out to 16 * 64KB =
105// 1MB of max buffered data.
106#[cfg(not(any(test, debug_assertions)))]
107const STREAM_CAPACITY: usize = 16;
108// A single entry stresses `write_pending` in test and debug builds.
109#[cfg(any(test, debug_assertions))]
110const STREAM_CAPACITY: usize = 1;
111
112// For *all* flows use a shared channel with 2048 entries, which works out
113// to 3MB of max buffered data at 1500 bytes per datagram.
114const FLOW_CAPACITY: usize = 2048;
115
116// Floor for the lazily-allocated body receive buffer. The buffer is sized to
117// the amount currently readable on the stream (see [`process_h3_data`]), but we
118// never allocate below this floor, for two reasons:
119//
120// - A `Limit<BytesMut>` with a zero limit reports no remaining capacity and
121//   would make `recv_body_buf` a no-op.
122// - Sizing strictly to the readable length defeats allocation amortization: a
123//   body that trickles in a few bytes at a time (e.g. one byte per read) would
124//   reallocate the buffer on every read. Allocating at least this many bytes
125//   lets a single allocation absorb many small reads (each `split()` off)
126//   before it is exhausted and reallocated.
127//
128// The floor is a soft hint, not a hard minimum: it never overrides the
129// configured cap, so a cap smaller than the floor still bounds the allocation
130// (see [`body_recv_buf_size`]).
131const MIN_BODY_RECV_BUF_SIZE: usize = 1024;
132
133// Default cap for the body receive buffer when
134// [`Http3Settings::max_recv_body_buf_size`] is unset or zero. Chosen well below
135// `BufFactory::MAX_BUF_SIZE` (64 KiB) so a request that carries a body doesn't
136// pin a large allocation for the life of the stream; a larger streamed body
137// simply reallocates once per driver read-cycle.
138const DEFAULT_MAX_BODY_RECV_BUF_SIZE: usize = 16 * 1024;
139
140/// Computes the capacity to use for the body receive buffer given the number of
141/// bytes currently readable on the stream and the configured maximum.
142///
143/// The result is clamped to `[floor, max]`, where `floor` is
144/// [`MIN_BODY_RECV_BUF_SIZE`] capped by `max`. Reads below the floor still
145/// allocate the floor (so a trickle of tiny reads reuses one allocation instead
146/// of reallocating each time), while a single (potentially adversarial) read
147/// never allocates more than `max`. `max` is derived from
148/// [`Http3Settings::max_recv_body_buf_size`], defaulting to
149/// [`DEFAULT_MAX_BODY_RECV_BUF_SIZE`] when the setting is unset or zero.
150///
151/// `max` is the hard upper bound: when it is smaller than the floor it wins, so
152/// the range passed to `clamp` is always valid (`floor <= max`). The
153/// constructor guarantees `max >= 1`, so the result is never zero.
154fn body_recv_buf_size(readable: usize, max: usize) -> usize {
155    let floor = MIN_BODY_RECV_BUF_SIZE.min(max);
156    readable.clamp(floor, max)
157}
158
159/// Used by a local task to send [`OutboundFrame`]s to a peer on the
160/// stream or flow associated with this channel.
161pub type OutboundFrameSender = PollSender<OutboundFrame>;
162
163/// Used internally to receive [`OutboundFrame`]s which should be sent to a peer
164/// on the stream or flow associated with this channel.
165type OutboundFrameStream = mpsc::Receiver<OutboundFrame>;
166
167/// Used internally to send [`InboundFrame`]s (data) from the peer to a local
168/// task on the stream or flow associated with this channel.
169type InboundFrameSender = PollSender<InboundFrame>;
170
171/// Used by a local task to receive [`InboundFrame`]s (data) on the stream or
172/// flow associated with this channel.
173pub type InboundFrameStream = mpsc::Receiver<InboundFrame>;
174
175/// The error type used internally in [H3Driver].
176///
177/// Note that [`ApplicationOverQuic`] errors are not exposed to users at this
178/// time. The type is public to document the failure modes in [H3Driver].
179#[derive(Debug, PartialEq, Eq)]
180#[non_exhaustive]
181pub enum H3ConnectionError {
182    /// The controller task was shut down and is no longer listening.
183    ControllerWentAway,
184    /// Other error at the connection, but not stream level.
185    H3(h3::Error),
186    /// Received data for a stream that was closed or never opened.
187    NonexistentStream,
188    /// The server's post-accept timeout was hit.
189    /// The timeout can be configured in [`Http3Settings`].
190    PostAcceptTimeout,
191}
192
193impl From<h3::Error> for H3ConnectionError {
194    fn from(err: h3::Error) -> Self {
195        H3ConnectionError::H3(err)
196    }
197}
198
199impl From<quiche::Error> for H3ConnectionError {
200    fn from(err: quiche::Error) -> Self {
201        H3ConnectionError::H3(h3::Error::TransportError(err))
202    }
203}
204
205impl Error for H3ConnectionError {}
206
207impl fmt::Display for H3ConnectionError {
208    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
209        let s: &dyn fmt::Display = match self {
210            Self::ControllerWentAway => &"controller went away",
211            Self::H3(e) => e,
212            Self::NonexistentStream => &"nonexistent stream",
213            Self::PostAcceptTimeout => &"post accept timeout hit",
214        };
215
216        write!(f, "H3ConnectionError: {s}")
217    }
218}
219
220type H3ConnectionResult<T> = Result<T, H3ConnectionError>;
221
222/// HTTP/3 headers that were received on a stream.
223///
224/// `recv` is used to read the message body, while `send` is used to transmit
225/// data back to the peer.
226pub struct IncomingH3Headers {
227    /// Stream ID of the frame.
228    pub stream_id: u64,
229    /// The actual [`h3::Header`]s which were received.
230    pub headers: Vec<h3::Header>,
231    /// An [`OutboundFrameSender`] for streaming body data to the peer. For
232    /// [ClientH3Driver], note that the request body can also be passed a
233    /// cloned sender via [`NewClientRequest`].
234    pub send: OutboundFrameSender,
235    /// An [`InboundFrameStream`] of body data received from the peer.
236    pub recv: InboundFrameStream,
237    /// Whether there is a body associated with the incoming headers.
238    pub read_fin: bool,
239    /// Handle to the [`H3AuditStats`] for the message's stream.
240    pub h3_audit_stats: Arc<H3AuditStats>,
241}
242
243impl fmt::Debug for IncomingH3Headers {
244    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
245        f.debug_struct("IncomingH3Headers")
246            .field("stream_id", &self.stream_id)
247            .field("headers", &self.headers)
248            .field("read_fin", &self.read_fin)
249            .field("h3_audit_stats", &self.h3_audit_stats)
250            .finish()
251    }
252}
253
254/// [`H3Event`]s are produced by an [H3Driver] to describe HTTP/3 state updates.
255///
256/// Both [ServerH3Driver] and [ClientH3Driver] may extend this enum with
257/// endpoint-specific variants. The events must be consumed by users of the
258/// drivers, like a higher-level `Server` or `Client` controller.
259#[derive(Debug)]
260pub enum H3Event {
261    /// A SETTINGS frame was received.
262    IncomingSettings {
263        /// Raw HTTP/3 setting pairs, in the order received from the peer.
264        settings: Vec<(u64, u64)>,
265    },
266
267    /// A HEADERS frame was received on the given stream. This is either a
268    /// request or a response depending on the perspective of the [`H3Event`]
269    /// receiver.
270    IncomingHeaders(IncomingH3Headers),
271
272    /// A DATAGRAM flow was created and associated with the given `flow_id`.
273    /// This event is fired before a HEADERS event for CONNECT[-UDP] requests.
274    NewFlow {
275        /// Flow ID of the new flow.
276        flow_id: u64,
277        /// An [`OutboundFrameSender`] for transmitting datagrams to the peer.
278        send: OutboundFrameSender,
279        /// An [`InboundFrameStream`] for receiving datagrams from the peer.
280        recv: InboundFrameStream,
281    },
282    /// A RST_STREAM frame was seen on the given `stream_id`. The user of the
283    /// driver should clean up any state allocated for this stream.
284    ResetStream { stream_id: u64 },
285    /// The connection has irrecoverably errored and is shutting down.
286    ConnectionError(h3::Error),
287    /// The connection has been shutdown, optionally due to an
288    /// [`H3ConnectionError`].
289    ConnectionShutdown(Option<H3ConnectionError>),
290    /// Body data has been received over a stream.
291    BodyBytesReceived {
292        /// Stream ID of the body data.
293        stream_id: u64,
294        /// Number of bytes received.
295        num_bytes: u64,
296        /// Whether the stream is finished and won't yield any more data.
297        fin: bool,
298    },
299    /// The stream has been closed. This is used to signal stream closures that
300    /// don't result from RST_STREAM frames, unlike the
301    /// [`H3Event::ResetStream`] variant.
302    StreamClosed { stream_id: u64 },
303    /// A GOAWAY frame was received from the peer containing `id`,
304    /// as described in https://datatracker.ietf.org/doc/html/rfc9114#section-5.2.
305    GoAway { id: u64 },
306}
307
308impl H3Event {
309    /// Generates an event from an applicable [`H3ConnectionError`].
310    fn from_error(err: &H3ConnectionError) -> Option<Self> {
311        Some(match err {
312            H3ConnectionError::H3(e) => Self::ConnectionError(*e),
313            H3ConnectionError::PostAcceptTimeout => Self::ConnectionShutdown(
314                Some(H3ConnectionError::PostAcceptTimeout),
315            ),
316            _ => return None,
317        })
318    }
319}
320
321/// An [`OutboundFrame`] is a data frame that should be sent from a local task
322/// to a peer over a [`quiche::h3::Connection`].
323///
324/// This is used, for example, to send response body data to a peer, or proxied
325/// UDP datagrams.
326#[derive(Debug)]
327pub enum OutboundFrame {
328    /// Response headers to be sent to the peer, with optional priority.
329    Headers(Vec<h3::Header>, Option<quiche::h3::Priority>),
330    /// Response body/CONNECT downstream data plus FIN flag.
331    Body(Bytes, bool),
332    /// CONNECT-UDP (DATAGRAM) downstream data plus flow ID.
333    Datagram(DgramBuffer, u64),
334    /// Close the stream with a trailers, with optional priority.
335    Trailers(Vec<h3::Header>, Option<quiche::h3::Priority>),
336    /// An error encountered when serving the request. Stream should be closed.
337    PeerStreamError,
338    /// DATAGRAM flow explicitly closed.
339    FlowShutdown { flow_id: u64, stream_id: u64 },
340}
341
342/// An [`InboundFrame`] is a data frame that was received from the peer over a
343/// [`quiche::h3::Connection`]. This is used by peers to send body or datagrams
344/// to the local task.
345#[derive(Debug)]
346pub enum InboundFrame {
347    /// Request body/CONNECT upstream data plus FIN flag.
348    Body(BytesMut, bool),
349    /// CONNECT-UDP (DATAGRAM) upstream data.
350    Datagram(DgramBuffer),
351}
352
353/// A ready-made [`ApplicationOverQuic`] which can handle HTTP/3 and MASQUE.
354/// Depending on the `DriverHooks` in use, it powers either a client or a
355/// server.
356///
357/// Use the [ClientH3Driver] and [ServerH3Driver] aliases to access the
358/// respective driver types. The driver is passed into an I/O loop and
359/// communicates with the driver's user (e.g., an HTTP client or a server) via
360/// its associated [H3Controller]. The controller allows the application to both
361/// listen for [`H3Event`]s of note and send [`H3Command`]s into the I/O loop.
362pub struct H3Driver<H: DriverHooks> {
363    /// Configuration used to initialize `conn`. Created from [`Http3Settings`]
364    /// in the constructor.
365    h3_config: h3::Config,
366    /// The underlying HTTP/3 connection. Initialized in
367    /// `ApplicationOverQuic::on_conn_established`.
368    conn: Option<h3::Connection>,
369    /// State required by the client/server hooks.
370    hooks: H,
371    /// Sends [`H3Event`]s to the [H3Controller] paired with this driver.
372    h3_event_sender: mpsc::UnboundedSender<H::Event>,
373    /// Receives [`H3Command`]s from the [H3Controller] paired with this driver.
374    cmd_recv: mpsc::UnboundedReceiver<H::Command>,
375    /// A sender that feeds back into `cmd_recv`. Used by hooks that need to
376    /// re-queue commands (e.g. retrying blocked requests) without access to
377    /// the [H3Controller]'s copy of the sender.
378    cmd_sender: mpsc::UnboundedSender<H::Command>,
379
380    /// A map of stream IDs to their [StreamCtx]. This is mainly used to
381    /// retrieve the internal Tokio channels associated with the stream.
382    stream_map: BTreeMap<u64, StreamCtx>,
383    /// A map of flow IDs to their [FlowCtx]. This is mainly used to retrieve
384    /// the internal Tokio channels associated with the flow.
385    flow_map: BTreeMap<u64, FlowCtx>,
386    /// Set of [`WaitForStream`] futures. A stream is added to this set if
387    /// we need to send to it and its channel is at capacity, or if we need
388    /// data from its channel and the channel is empty.
389    waiting_streams: WaitingStreams,
390
391    /// Receives [`OutboundFrame`]s from all datagram flows on the connection.
392    dgram_recv: OutboundFrameStream,
393    /// Keeps the datagram channel open such that datagram flows can be created.
394    dgram_send: OutboundFrameSender,
395    /// A buffer to receive H3 body data from quiche. Lazily allocated on the
396    /// first body read and released once no streams remain, so idle
397    /// connections hold no receive buffer. We `split()` off filled parts until
398    /// we need to reallocate.
399    body_recv_buf: Option<bytes::buf::Limit<BytesMut>>,
400    /// Upper bound on the capacity of `body_recv_buf`, from
401    /// [`Http3Settings::max_recv_body_buf_size`]. Unset or `0` defaults to
402    /// [`DEFAULT_MAX_BODY_RECV_BUF_SIZE`], so this is always non-zero.
403    max_recv_body_buf_size: usize,
404
405    /// The maximum HTTP/3 stream ID seen on this connection.
406    max_stream_seen: u64,
407
408    /// Tracks whether we have forwarded the HTTP/3 SETTINGS frame
409    /// to the [H3Controller] once.
410    settings_received_and_forwarded: bool,
411    /// Tracks whether the H3 event receiver has been dropped.
412    /// Used to avoid busy-looping on `h3_event_sender.closed()`.
413    h3_event_receiver_dropped: bool,
414}
415
416impl<H: DriverHooks> H3Driver<H> {
417    /// Builds a new [H3Driver] and an associated [H3Controller].
418    ///
419    /// The driver should then be passed to
420    /// [`InitialQuicConnection`](crate::InitialQuicConnection)'s `start`
421    /// method.
422    pub fn new(http3_settings: Http3Settings) -> (Self, H3Controller<H>) {
423        let (dgram_send, dgram_recv) = mpsc::channel(FLOW_CAPACITY);
424        let (cmd_sender, cmd_recv) = mpsc::unbounded_channel();
425        let (h3_event_sender, h3_event_recv) = mpsc::unbounded_channel();
426
427        (
428            H3Driver {
429                h3_config: (&http3_settings).into(),
430                conn: None,
431                hooks: H::new(&http3_settings),
432                h3_event_sender,
433                cmd_recv,
434                cmd_sender: cmd_sender.clone(),
435
436                stream_map: BTreeMap::new(),
437                flow_map: BTreeMap::new(),
438
439                dgram_recv,
440                dgram_send: PollSender::new(dgram_send),
441                max_stream_seen: 0,
442                body_recv_buf: None,
443                // `Some(0)` would make `recv_body_buf` a no-op, so treat it as
444                // unset.
445                max_recv_body_buf_size: http3_settings
446                    .max_recv_body_buf_size
447                    .filter(|&size| size > 0)
448                    .unwrap_or(DEFAULT_MAX_BODY_RECV_BUF_SIZE),
449
450                waiting_streams: WaitingStreams::new(),
451
452                settings_received_and_forwarded: false,
453                h3_event_receiver_dropped: false,
454            },
455            H3Controller {
456                cmd_sender,
457                h3_event_recv: Some(h3_event_recv),
458            },
459        )
460    }
461
462    /// Returns a sender that feeds back into this driver's own `cmd_recv`.
463    ///
464    /// Hooks that need to re-queue commands (e.g. retrying a request that
465    /// was temporarily blocked) can use this sender without needing access
466    /// to the paired [H3Controller].
467    pub(crate) fn self_cmd_sender(&self) -> &mpsc::UnboundedSender<H::Command> {
468        &self.cmd_sender
469    }
470
471    /// Retrieve the [FlowCtx] associated with the given `flow_id`. If no
472    /// context is found, a new one will be created.
473    fn get_or_insert_flow(
474        &mut self, flow_id: u64,
475    ) -> H3ConnectionResult<&mut FlowCtx> {
476        use std::collections::btree_map::Entry;
477        Ok(match self.flow_map.entry(flow_id) {
478            Entry::Vacant(e) => {
479                // This is a datagram for a new flow we haven't seen before
480                let (flow, recv) = FlowCtx::new(FLOW_CAPACITY);
481                let flow_req = H3Event::NewFlow {
482                    flow_id,
483                    recv,
484                    send: self.dgram_send.clone(),
485                };
486                self.h3_event_sender
487                    .send(flow_req.into())
488                    .map_err(|_| H3ConnectionError::ControllerWentAway)?;
489                e.insert(flow)
490            },
491            Entry::Occupied(e) => e.into_mut(),
492        })
493    }
494
495    /// Adds a [StreamCtx] to the stream map with the given `stream_id`.
496    fn insert_stream(&mut self, stream_id: u64, ctx: StreamCtx) {
497        self.stream_map.insert(stream_id, ctx);
498        self.max_stream_seen = self.max_stream_seen.max(stream_id);
499    }
500
501    /// Fetches body chunks from the [`quiche::h3::Connection`] and forwards
502    /// them to the stream's associated [`InboundFrameStream`].
503    fn process_h3_data(
504        &mut self, qconn: &mut QuicheConnection, stream_id: u64,
505    ) -> H3ConnectionResult<()> {
506        // Split self borrow between conn and stream_map
507        let conn = self.conn.as_mut().ok_or(Self::connection_not_present())?;
508        let ctx = self
509            .stream_map
510            .get_mut(&stream_id)
511            .ok_or(H3ConnectionError::NonexistentStream)?;
512
513        enum StreamStatus {
514            Done { close: bool },
515            Reset { wire_err_code: u64 },
516            Blocked,
517        }
518
519        let status = loop {
520            let Some(sender) = ctx.send.as_ref().and_then(PollSender::get_ref)
521            else {
522                // already waiting for capacity
523                break StreamStatus::Done { close: false };
524            };
525
526            let try_reserve_result = sender.try_reserve();
527            let permit = match try_reserve_result {
528                Ok(permit) => permit,
529                Err(TrySendError::Closed(())) => {
530                    // The channel has closed before we delivered a fin or reset
531                    // to the application.
532                    if !ctx.fin_or_reset_recv &&
533                        ctx.associated_dgram_flow_id.is_none()
534                    // The channel might be closed if the stream was used to
535                    // initiate a datagram exchange.
536                    // TODO: ideally, the application would still shut down the
537                    // stream properly. Once applications code
538                    // is fixed, we can remove this check.
539                    {
540                        let err = h3::WireErrorCode::RequestCancelled as u64;
541                        let _ = qconn.stream_shutdown(
542                            stream_id,
543                            quiche::Shutdown::Read,
544                            err,
545                        );
546                        // Release the borrow of `ctx` before updating it.
547                        drop(try_reserve_result);
548                        ctx.handle_sent_stop_sending(err);
549                        // TODO: should we send an H3Event event to
550                        // h3_event_sender? We can only get here if the app
551                        // actively closed or dropped
552                        // the channel so any event we send would be more for
553                        // logging or auditing
554                    }
555                    break StreamStatus::Done {
556                        close: ctx.both_directions_done(),
557                    };
558                },
559                Err(TrySendError::Full(())) => {
560                    if ctx.fin_or_reset_recv || qconn.stream_readable(stream_id) {
561                        break StreamStatus::Blocked;
562                    }
563                    break StreamStatus::Done { close: false };
564                },
565            };
566
567            if ctx.fin_or_reset_recv {
568                // Signal end-of-body to upstream
569                permit.send(InboundFrame::Body(Default::default(), true));
570                break StreamStatus::Done {
571                    close: ctx.fin_or_reset_sent,
572                };
573            }
574
575            // Size the buffer for currently readable, in-order data, capped at
576            // the configured maximum. The readable length includes H3 framing,
577            // so it can drain all readable body bytes. The minimum keeps the
578            // allocation nonzero and amortizes a trickle of small reads.
579            let want = body_recv_buf_size(
580                qconn.stream_readable_len(stream_id, self.max_recv_body_buf_size),
581                self.max_recv_body_buf_size,
582            );
583            // Lazily allocate the receive buffer on first use; idle
584            // connections never receive body bytes and never allocate it.
585            let body_recv_buf = self
586                .body_recv_buf
587                .get_or_insert_with(|| BytesMut::with_capacity(want).limit(want));
588            // `Limit<BytesMut>` reports only capacity below its limit. Plain
589            // `BytesMut` can reallocate and always reports available space.
590            //
591            // Reallocate when the remaining room cannot hold this read. This
592            // handles exhausted buffers and reads larger than the previous
593            // allocation. Keep capacity equal to the limit to preserve the
594            // `split()` invariant checked below.
595            if body_recv_buf.remaining_mut() < want {
596                *body_recv_buf = BytesMut::with_capacity(want).limit(want);
597            }
598            match conn.recv_body_buf(qconn, stream_id, &mut *body_recv_buf) {
599                Ok(n) => {
600                    ctx.audit_stats.add_downstream_bytes_recvd(n as u64);
601                    let event = H3Event::BodyBytesReceived {
602                        stream_id,
603                        num_bytes: n as u64,
604                        fin: false,
605                    };
606                    let _ = self.h3_event_sender.send(event.into());
607                    // Take the filled part, leave the remaining capacity
608                    let filled_body = body_recv_buf.get_mut().split();
609                    // Sanity check: the remaining spare capacity should equal
610                    // the limit.
611                    debug_assert_eq!(
612                        body_recv_buf.get_mut().spare_capacity_mut().len(),
613                        body_recv_buf.remaining_mut()
614                    );
615                    // A full split leaves only an empty shared handle, so let
616                    // the forwarded frame own the allocation.
617                    if !body_recv_buf.has_remaining_mut() {
618                        self.body_recv_buf = None;
619                    }
620                    permit.send(InboundFrame::Body(filled_body, false));
621                },
622                Err(h3::Error::Done) =>
623                    break StreamStatus::Done { close: false },
624                Err(h3::Error::TransportError(quiche::Error::StreamReset(
625                    code,
626                ))) => {
627                    break StreamStatus::Reset {
628                        wire_err_code: code,
629                    };
630                },
631                Err(_) => break StreamStatus::Done { close: true },
632            }
633        };
634
635        match status {
636            StreamStatus::Done { close } => {
637                if close {
638                    return self.cleanup_stream(qconn, stream_id);
639                }
640
641                // The QUIC stream is finished, manually invoke `process_h3_fin`
642                // in case `h3::poll()` is never called again.
643                //
644                // Note that this case will not conflict with StreamStatus::Done
645                // being returned due to the body channel being
646                // blocked. qconn.stream_finished() will guarantee
647                // that we've fully parsed the body as it only returns true
648                // if we've seen a Fin for the read half of the stream.
649                if !ctx.fin_or_reset_recv && qconn.stream_finished(stream_id) {
650                    return self.process_h3_fin(qconn, stream_id);
651                }
652            },
653            StreamStatus::Reset { wire_err_code } => {
654                debug_assert!(ctx.send.is_some());
655                ctx.handle_recvd_reset(wire_err_code);
656                let cleanup = ctx.both_directions_done();
657                H::stream_recv_closed(self, stream_id);
658                self.h3_event_sender
659                    .send(H3Event::ResetStream { stream_id }.into())
660                    .map_err(|_| H3ConnectionError::ControllerWentAway)?;
661                if cleanup {
662                    return self.cleanup_stream(qconn, stream_id);
663                }
664            },
665            StreamStatus::Blocked => {
666                self.waiting_streams.push(ctx.wait_for_send(stream_id))?;
667            },
668        }
669
670        Ok(())
671    }
672
673    /// Processes an end-of-stream event from the [`quiche::h3::Connection`].
674    fn process_h3_fin(
675        &mut self, qconn: &mut QuicheConnection, stream_id: u64,
676    ) -> H3ConnectionResult<()> {
677        let Some(ctx) = self.stream_map.get_mut(&stream_id) else {
678            // Report a pending stop because no context remains to handle it.
679            //
680            // This removes its writable notification and permits transport
681            // cleanup.
682            let _ = qconn.stream_capacity(stream_id);
683            return Ok(());
684        };
685
686        if ctx.fin_or_reset_recv {
687            return Ok(());
688        }
689
690        ctx.fin_or_reset_recv = true;
691        ctx.audit_stats
692            .set_recvd_stream_fin(StreamClosureKind::Explicit);
693        H::stream_recv_closed(self, stream_id);
694
695        // It's important to send this H3Event before process_h3_data so that
696        // a server can (potentially) generate the control response before the
697        // corresponding receiver drops.
698        let event = H3Event::BodyBytesReceived {
699            stream_id,
700            num_bytes: 0,
701            fin: true,
702        };
703        let _ = self.h3_event_sender.send(event.into());
704
705        // Communicate fin to upstream. Since `ctx.fin_recv` is true now,
706        // there can't be a recursive loop.
707        self.process_h3_data(qconn, stream_id)
708    }
709
710    /// Processes a single [`quiche::h3::Event`] received from the underlying
711    /// [`quiche::h3::Connection`]. Some events are dispatched to helper
712    /// methods.
713    fn process_read_event(
714        &mut self, qconn: &mut QuicheConnection, stream_id: u64, event: h3::Event,
715    ) -> H3ConnectionResult<()> {
716        self.forward_settings()?;
717
718        match event {
719            // Requests/responses are exclusively handled by hooks.
720            h3::Event::Headers { list, more_frames } =>
721                H::headers_received(self, qconn, InboundHeaders {
722                    stream_id,
723                    headers: list,
724                    has_body: more_frames,
725                }),
726
727            h3::Event::Data => self.process_h3_data(qconn, stream_id),
728            h3::Event::Finished => self.process_h3_fin(qconn, stream_id),
729
730            h3::Event::Reset(code) => {
731                if let Some(ctx) = self.stream_map.get_mut(&stream_id) {
732                    ctx.handle_recvd_reset(code);
733                    self.waiting_streams.cancel_upstream(stream_id);
734
735                    let cleanup = ctx.both_directions_done();
736                    H::stream_recv_closed(self, stream_id);
737                    self.h3_event_sender
738                        .send(H3Event::ResetStream { stream_id }.into())
739                        .map_err(|_| H3ConnectionError::ControllerWentAway)?;
740                    if cleanup {
741                        return self.cleanup_stream(qconn, stream_id);
742                    }
743                } else {
744                    // A reset can finish a stream that never had a context.
745                    let _ = qconn.stream_capacity(stream_id);
746                }
747
748                // TODO: if we don't have the stream in our map: should we
749                // send the H3Event::ResetStream?
750                Ok(())
751            },
752
753            h3::Event::PriorityUpdate => Ok(()),
754            h3::Event::GoAway => {
755                self.h3_event_sender
756                    .send(H3Event::GoAway { id: stream_id }.into())
757                    .map_err(|_| H3ConnectionError::ControllerWentAway)?;
758                Ok(())
759            },
760        }
761    }
762
763    /// The SETTINGS frame can be received at any point, so we
764    /// need to check `peer_settings_raw` to decide if we've received it.
765    ///
766    /// Settings should only be sent once, so we generate a single event
767    /// when `peer_settings_raw` transitions from None to Some.
768    fn forward_settings(&mut self) -> H3ConnectionResult<()> {
769        if self.settings_received_and_forwarded {
770            return Ok(());
771        }
772
773        // capture the peer settings and forward it
774        if let Some(settings) = self.conn_mut()?.peer_settings_raw() {
775            let incoming_settings = H3Event::IncomingSettings {
776                settings: settings.to_vec(),
777            };
778
779            self.h3_event_sender
780                .send(incoming_settings.into())
781                .map_err(|_| H3ConnectionError::ControllerWentAway)?;
782
783            self.settings_received_and_forwarded = true;
784        }
785        Ok(())
786    }
787
788    /// Send an individual frame to the underlying [`quiche::h3::Connection`] to
789    /// be flushed at a later time.
790    ///
791    /// `Self::process_writes` will iterate over all writable streams and call
792    /// this method in a loop for each stream to send all writable packets.
793    fn process_write_frame(
794        conn: &mut h3::Connection, qconn: &mut QuicheConnection,
795        ctx: &mut StreamCtx,
796    ) -> h3::Result<()> {
797        let audit_stats = &ctx.audit_stats;
798        let stream_id = audit_stats.stream_id();
799
800        match &mut ctx.queued_frame {
801            Some(OutboundFrame::Headers(headers, priority)) => {
802                let prio = priority.as_ref().unwrap_or(&DEFAULT_PRIO);
803
804                let res = if ctx.initial_headers_sent {
805                    // Initial headers were already sent, send additional
806                    // headers now.
807                    conn.send_additional_headers_with_priority(
808                        qconn, stream_id, headers, prio, false, false,
809                    )
810                } else {
811                    // Send initial headers.
812                    conn.send_response_with_priority(
813                        qconn, stream_id, headers, prio, false,
814                    )
815                    .inspect(|_| ctx.initial_headers_sent = true)
816                };
817
818                if let Err(h3::Error::StreamBlocked) = res {
819                    ctx.first_full_headers_flush_fail_time
820                        .get_or_insert(Instant::now());
821                }
822
823                if res.is_ok() {
824                    if let Some(first) =
825                        ctx.first_full_headers_flush_fail_time.take()
826                    {
827                        ctx.audit_stats.add_header_flush_duration(
828                            Instant::now().duration_since(first),
829                        );
830                    }
831                }
832
833                return res;
834            },
835
836            Some(OutboundFrame::Body(body, fin)) if !body.is_empty() || *fin => {
837                let len = body.len();
838
839                if *fin {
840                    // Drop the receiver on FIN to prevent more frames.
841                    // We cannot use `mpsc::Receiver::close()` because of a
842                    // Tokio inconsistency when reading a closed channel.
843                    // See https://github.com/tokio-rs/tokio/issues/7631.
844                    ctx.recv = None;
845                }
846
847                let n = conn.send_body_zc(qconn, stream_id, body, *fin)?;
848
849                audit_stats.add_downstream_bytes_sent(n as _);
850
851                return if n != len {
852                    // Couldn't write the entire body, `send_body_zc` will
853                    // have trimmed `body` accordingly. The driver keeps
854                    // the remainder of the body to send in the future.
855                    debug_assert_eq!(
856                        n + body.len(),
857                        len,
858                        "send_body_zc() should have trimmed body but did not"
859                    );
860                    Err(h3::Error::StreamBlocked)
861                } else {
862                    if *fin {
863                        Self::on_fin_sent(ctx)?;
864                    }
865
866                    Ok(())
867                };
868            },
869
870            Some(OutboundFrame::Body(..)) => {
871                // An empty non-FIN body falls through to the STOP probe below.
872            },
873
874            Some(OutboundFrame::Trailers(headers, priority)) => {
875                let prio = priority.as_ref().unwrap_or(&DEFAULT_PRIO);
876
877                // trailers always set fin=true
878                let res = conn.send_additional_headers_with_priority(
879                    qconn, stream_id, headers, prio, true, true,
880                );
881
882                if res.is_ok() {
883                    Self::on_fin_sent(ctx)?;
884                }
885
886                return res;
887            },
888
889            Some(OutboundFrame::PeerStreamError) => {
890                return Err(h3::Error::MessageError);
891            },
892
893            Some(OutboundFrame::FlowShutdown { .. }) => {
894                unreachable!("Only flows send shutdowns")
895            },
896
897            Some(OutboundFrame::Datagram(..)) => {
898                unreachable!("Only flows send datagrams")
899            },
900
901            None => {
902                // With no queued frame, fall through to the STOP probe below.
903            },
904        }
905
906        // H3 send methods normally report STOP_SENDING as a transport error,
907        // but neither an idle stream nor an empty non-FIN body invokes them.
908        //
909        // Check for cancellation here so the caller can close the outbound
910        // channel and record the stop code without waiting for another frame.
911        if let Err(error @ quiche::Error::StreamStopped(_)) =
912            qconn.stream_capacity(stream_id)
913        {
914            return Err(h3::Error::TransportError(error));
915        }
916
917        Ok(())
918    }
919
920    fn on_fin_sent(ctx: &mut StreamCtx) -> h3::Result<()> {
921        ctx.recv = None;
922        ctx.fin_or_reset_sent = true;
923        ctx.audit_stats
924            .set_sent_stream_fin(StreamClosureKind::Explicit);
925        if ctx.fin_or_reset_recv {
926            // Return a TransportError to trigger stream cleanup
927            // instead of h3::Error::Done
928            Err(h3::Error::TransportError(quiche::Error::Done))
929        } else {
930            Ok(())
931        }
932    }
933
934    /// Resumes reads or writes to the connection when a stream channel becomes
935    /// unblocked.
936    ///
937    /// If we were waiting for more data from a channel, we resume writing to
938    /// the connection. Otherwise, we were blocked on channel capacity and
939    /// continue reading from the connection. `Upstream` in this context is
940    /// the consumer of the stream.
941    fn upstream_ready(
942        &mut self, qconn: &mut QuicheConnection, ready: StreamReady,
943    ) -> H3ConnectionResult<()> {
944        match ready {
945            StreamReady::Downstream(r) => self.upstream_read_ready(qconn, r),
946            StreamReady::Upstream(r) => self.upstream_write_ready(qconn, r),
947        }
948    }
949
950    fn upstream_read_ready(
951        &mut self, qconn: &mut QuicheConnection,
952        read_ready: ReceivedDownstreamData,
953    ) -> H3ConnectionResult<()> {
954        let ReceivedDownstreamData {
955            stream_id,
956            chan,
957            data,
958        } = read_ready;
959
960        let Some(stream) = self.stream_map.get_mut(&stream_id) else {
961            return Ok(());
962        };
963
964        // After a stop, the closed parked receiver can still return frames
965        // queued before it. The write side is done, so drop the receiver and
966        // any such frame, except a peer stream error, which must still shut
967        // down the read side.
968        if !stream.fin_or_reset_sent {
969            stream.recv = Some(chan);
970            stream.queued_frame = data;
971        } else if let Some(OutboundFrame::PeerStreamError) = data {
972            stream.queued_frame = data;
973        } else {
974            return Ok(());
975        }
976
977        self.process_writable_stream(qconn, stream_id)
978    }
979
980    fn upstream_write_ready(
981        &mut self, qconn: &mut QuicheConnection,
982        write_ready: HaveUpstreamCapacity,
983    ) -> H3ConnectionResult<()> {
984        let HaveUpstreamCapacity {
985            stream_id,
986            mut chan,
987        } = write_ready;
988
989        match self.stream_map.get_mut(&stream_id) {
990            None => Ok(()),
991            Some(stream) => {
992                // Release the associated permit before retaining the channel.
993                chan.abort_send();
994                stream.send = Some(chan);
995                self.process_h3_data(qconn, stream_id)
996            },
997        }
998    }
999
1000    /// Processes all queued outbound datagrams from the `dgram_recv` channel.
1001    fn dgram_ready(
1002        &mut self, qconn: &mut QuicheConnection, frame: OutboundFrame,
1003    ) -> H3ConnectionResult<()> {
1004        let mut frame = Ok(frame);
1005
1006        loop {
1007            match frame {
1008                Ok(OutboundFrame::Datagram(dgram, flow_id)) => {
1009                    // Drop datagrams if there is no capacity
1010                    let _ = datagram::send_h3_dgram(qconn, flow_id, dgram);
1011                },
1012                Ok(OutboundFrame::FlowShutdown { flow_id, stream_id }) => {
1013                    self.shutdown_stream(
1014                        qconn,
1015                        stream_id,
1016                        StreamShutdown::Both {
1017                            read_error_code: WireErrorCode::NoError as u64,
1018                            write_error_code: WireErrorCode::NoError as u64,
1019                        },
1020                    )?;
1021                    self.flow_map.remove(&flow_id);
1022                    self.close_if_idle(qconn);
1023                    break;
1024                },
1025                Ok(_) => unreachable!("Flows can't send frame of other types"),
1026                Err(TryRecvError::Empty) => break,
1027                Err(TryRecvError::Disconnected) =>
1028                    return Err(H3ConnectionError::ControllerWentAway),
1029            }
1030
1031            frame = self.dgram_recv.try_recv();
1032        }
1033
1034        Ok(())
1035    }
1036
1037    /// Return a mutable reference to the driver's HTTP/3 connection.
1038    ///
1039    /// If the connection doesn't exist yet, this function returns
1040    /// a `Self::connection_not_present()` error.
1041    fn conn_mut(&mut self) -> H3ConnectionResult<&mut h3::Connection> {
1042        self.conn.as_mut().ok_or(Self::connection_not_present())
1043    }
1044
1045    /// Alias for [`quiche::Error::TlsFail`], which is used in the case where
1046    /// this driver doesn't have an established HTTP/3 connection attached
1047    /// to it yet.
1048    const fn connection_not_present() -> H3ConnectionError {
1049        H3ConnectionError::H3(h3::Error::TransportError(quiche::Error::TlsFail))
1050    }
1051
1052    /// Cleans up internal state for the indicated HTTP/3 stream.
1053    ///
1054    /// This function removes the stream from the stream map, closes any pending
1055    /// futures, removes associated DATAGRAM flows, and sends a
1056    /// [`H3Event::StreamClosed`] event (for servers).
1057    fn cleanup_stream(
1058        &mut self, qconn: &mut QuicheConnection, stream_id: u64,
1059    ) -> H3ConnectionResult<()> {
1060        let Some(mut stream_ctx) = self.stream_map.remove(&stream_id) else {
1061            return Ok(());
1062        };
1063
1064        // Receive-side cleanup can precede the writable stop notification.
1065        if let Err(quiche::Error::StreamStopped(code)) =
1066            qconn.stream_capacity(stream_id)
1067        {
1068            stream_ctx.handle_recvd_stop_sending(code);
1069        }
1070
1071        self.waiting_streams.cancel_upstream(stream_id);
1072        self.waiting_streams.cancel_downstream(stream_id);
1073
1074        // Close any DATAGRAM-proxying channels associated with the stream.
1075        if let Some(mapped_flow_id) = stream_ctx.associated_dgram_flow_id {
1076            self.flow_map.remove(&mapped_flow_id);
1077        }
1078
1079        H::stream_closed(self, stream_id);
1080
1081        self.close_if_idle(qconn);
1082
1083        Ok(())
1084    }
1085
1086    /// Handles connection cleanup once no streams or flows remain.
1087    ///
1088    /// Releases the body receive buffer (it is reallocated lazily on the next
1089    /// body read; body bytes only flow on active streams, so an empty stream
1090    /// map means it is unused) and closes the connection with `NoError` if the
1091    /// H3 event receiver has been dropped.
1092    fn close_if_idle(&mut self, qconn: &mut QuicheConnection) {
1093        if self.stream_map.is_empty() && self.flow_map.is_empty() {
1094            self.body_recv_buf = None;
1095
1096            if self.h3_event_receiver_dropped {
1097                let _ = qconn.close(
1098                    true,
1099                    quiche::h3::WireErrorCode::NoError as u64,
1100                    &[],
1101                );
1102            }
1103        }
1104    }
1105
1106    /// Shuts down the indicated HTTP/3 stream by sending frames and cleaning
1107    /// up then cleans up internal state by calling
1108    /// [`Self::cleanup_stream`].
1109    fn shutdown_stream(
1110        &mut self, qconn: &mut QuicheConnection, stream_id: u64,
1111        shutdown: StreamShutdown,
1112    ) -> H3ConnectionResult<()> {
1113        let Some(stream_ctx) = self.stream_map.get(&stream_id) else {
1114            return Ok(());
1115        };
1116
1117        let audit_stats = &stream_ctx.audit_stats;
1118
1119        match shutdown {
1120            StreamShutdown::Read { error_code } => {
1121                audit_stats.set_sent_stop_sending_error_code(error_code as _);
1122                let _ = qconn.stream_shutdown(
1123                    stream_id,
1124                    quiche::Shutdown::Read,
1125                    error_code,
1126                );
1127            },
1128            StreamShutdown::Write { error_code } => {
1129                audit_stats.set_sent_reset_stream_error_code(error_code as _);
1130                let _ = qconn.stream_shutdown(
1131                    stream_id,
1132                    quiche::Shutdown::Write,
1133                    error_code,
1134                );
1135            },
1136            StreamShutdown::Both {
1137                read_error_code,
1138                write_error_code,
1139            } => {
1140                audit_stats
1141                    .set_sent_stop_sending_error_code(read_error_code as _);
1142                let _ = qconn.stream_shutdown(
1143                    stream_id,
1144                    quiche::Shutdown::Read,
1145                    read_error_code,
1146                );
1147                audit_stats
1148                    .set_sent_reset_stream_error_code(write_error_code as _);
1149                let _ = qconn.stream_shutdown(
1150                    stream_id,
1151                    quiche::Shutdown::Write,
1152                    write_error_code,
1153                );
1154            },
1155        }
1156
1157        self.cleanup_stream(qconn, stream_id)
1158    }
1159
1160    /// Handles a regular [`H3Command`]. May be called internally by
1161    /// [DriverHooks] for non-endpoint-specific [`H3Command`]s.
1162    fn handle_core_command(
1163        &mut self, qconn: &mut QuicheConnection, cmd: H3Command,
1164    ) -> H3ConnectionResult<()> {
1165        match cmd {
1166            H3Command::QuicCmd(cmd) => cmd.execute(qconn),
1167            H3Command::GoAway => {
1168                let max_id = self.max_stream_seen;
1169                self.conn_mut()
1170                    .expect("connection should be established")
1171                    .send_goaway(qconn, max_id)?;
1172            },
1173            H3Command::ShutdownStream {
1174                stream_id,
1175                shutdown,
1176            } => {
1177                self.shutdown_stream(qconn, stream_id, shutdown)?;
1178            },
1179        }
1180        Ok(())
1181    }
1182}
1183
1184impl<H: DriverHooks> H3Driver<H> {
1185    /// Reads all buffered datagrams out of `qconn` and distributes them to
1186    /// their flow channels.
1187    fn process_available_dgrams(
1188        &mut self, qconn: &mut QuicheConnection,
1189    ) -> H3ConnectionResult<()> {
1190        loop {
1191            match datagram::receive_h3_dgram(qconn) {
1192                Ok((flow_id, dgram))
1193                    if !qconn.is_server() ||
1194                        self.hooks.extended_connect_enabled() =>
1195                {
1196                    self.get_or_insert_flow(flow_id)?.send_best_effort(dgram);
1197                },
1198                Ok(_) => {},
1199                Err(quiche::Error::Done) => return Ok(()),
1200                Err(err) => return Err(H3ConnectionError::from(err)),
1201            }
1202        }
1203    }
1204
1205    /// Flushes any queued-up frames for `stream_id` into `qconn` until either
1206    /// there is no more capacity in `qconn` or no more frames to send.
1207    fn process_writable_stream(
1208        &mut self, qconn: &mut QuicheConnection, stream_id: u64,
1209    ) -> H3ConnectionResult<()> {
1210        // Split self borrow between conn and stream_map
1211        let conn = self.conn.as_mut().ok_or(Self::connection_not_present())?;
1212        let Some(ctx) = self.stream_map.get_mut(&stream_id) else {
1213            // Keep the stop pending while headers can still create a context.
1214            if qconn.stream_finished(stream_id) {
1215                let _ = qconn.stream_capacity(stream_id);
1216            }
1217
1218            return Ok(());
1219        };
1220
1221        loop {
1222            // Process each writable frame, queue the next frame for processing
1223            // and shut down any errored streams.
1224            match Self::process_write_frame(conn, qconn, ctx) {
1225                Ok(()) => ctx.queued_frame = None,
1226                Err(h3::Error::StreamBlocked | h3::Error::Done) => break,
1227                Err(h3::Error::MessageError) => {
1228                    return self.shutdown_stream(
1229                        qconn,
1230                        stream_id,
1231                        StreamShutdown::Both {
1232                            read_error_code: WireErrorCode::MessageError as u64,
1233                            write_error_code: WireErrorCode::MessageError as u64,
1234                        },
1235                    );
1236                },
1237                Err(h3::Error::TransportError(quiche::Error::StreamStopped(
1238                    e,
1239                ))) => {
1240                    ctx.handle_recvd_stop_sending(e);
1241
1242                    // An idle stream's outbound receiver can be owned by a
1243                    // waiting future instead of ctx.recv. Cancel it so the
1244                    // application cannot send another frame.
1245                    self.waiting_streams.cancel_downstream(stream_id);
1246
1247                    if ctx.both_directions_done() {
1248                        return self.cleanup_stream(qconn, stream_id);
1249                    } else {
1250                        return Ok(());
1251                    }
1252                },
1253                Err(h3::Error::TransportError(
1254                    quiche::Error::InvalidStreamState(stream),
1255                )) => {
1256                    return self.cleanup_stream(qconn, stream);
1257                },
1258                Err(_) => {
1259                    return self.cleanup_stream(qconn, stream_id);
1260                },
1261            }
1262
1263            let Some(recv) = ctx.recv.as_mut() else {
1264                // This stream is already waiting for data or we wrote a fin and
1265                // closed the channel.
1266                debug_assert!(
1267                    ctx.queued_frame.is_none(),
1268                    "We MUST NOT have a queued frame if we are already waiting on 
1269                    more data from the channel"
1270                );
1271                return Ok(());
1272            };
1273
1274            // Attempt to queue the next frame for processing. The corresponding
1275            // sender is created at the same time as the `StreamCtx`
1276            // and ultimately ends up in an `H3Body`. The body then
1277            // determines which frames to send to the peer via
1278            // this processing loop.
1279            match recv.try_recv() {
1280                Ok(frame) => ctx.queued_frame = Some(frame),
1281                Err(TryRecvError::Disconnected) => {
1282                    if !ctx.fin_or_reset_sent &&
1283                        ctx.associated_dgram_flow_id.is_none()
1284                    // The channel might be closed if the stream was used to
1285                    // initiate a datagram exchange.
1286                    // TODO: ideally, the application would still shut down the
1287                    // stream properly. Once applications code
1288                    // is fixed, we can remove this check.
1289                    {
1290                        // No FIN was written before closure. Send RESET_STREAM
1291                        // because no more data will follow.
1292                        let err = h3::WireErrorCode::RequestCancelled as u64;
1293                        let _ = qconn.stream_shutdown(
1294                            stream_id,
1295                            quiche::Shutdown::Write,
1296                            err,
1297                        );
1298                        ctx.handle_sent_reset(err);
1299                        if ctx.both_directions_done() {
1300                            return self.cleanup_stream(qconn, stream_id);
1301                        }
1302                    }
1303                    break;
1304                },
1305                Err(TryRecvError::Empty) => {
1306                    self.waiting_streams.push(ctx.wait_for_recv(stream_id))?;
1307                    break;
1308                },
1309            }
1310        }
1311
1312        Ok(())
1313    }
1314
1315    /// Tests `qconn` for either a local or peer error and increments
1316    /// the associated HTTP/3 or QUIC error counter.
1317    fn record_quiche_error(qconn: &mut QuicheConnection, metrics: &impl Metrics) {
1318        // split metrics between local/peer and QUIC/HTTP/3 level errors
1319        if let Some(err) = qconn.local_error() {
1320            if err.is_app {
1321                metrics.local_h3_conn_close_error_count(err.error_code.into())
1322            } else {
1323                metrics.local_quic_conn_close_error_count(err.error_code.into())
1324            }
1325            .inc();
1326        } else if let Some(err) = qconn.peer_error() {
1327            if err.is_app {
1328                metrics.peer_h3_conn_close_error_count(err.error_code.into())
1329            } else {
1330                metrics.peer_quic_conn_close_error_count(err.error_code.into())
1331            }
1332            .inc();
1333        }
1334    }
1335}
1336
1337impl<H: DriverHooks> ApplicationOverQuic for H3Driver<H> {
1338    fn on_conn_established(
1339        &mut self, quiche_conn: &mut QuicheConnection,
1340        handshake_info: &HandshakeInfo,
1341    ) -> QuicResult<()> {
1342        let conn = h3::Connection::with_transport(quiche_conn, &self.h3_config)?;
1343        self.conn = Some(conn);
1344
1345        H::conn_established(self, quiche_conn, handshake_info)?;
1346        Ok(())
1347    }
1348
1349    #[inline]
1350    fn should_act(&self) -> bool {
1351        self.conn.is_some()
1352    }
1353
1354    /// Poll the underlying [`quiche::h3::Connection`] for
1355    /// [`quiche::h3::Event`]s and DATAGRAMs, delegating processing to
1356    /// `Self::process_read_event`.
1357    ///
1358    /// If a DATAGRAM is found, it is sent to the receiver on its channel.
1359    fn process_reads(&mut self, qconn: &mut QuicheConnection) -> QuicResult<()> {
1360        loop {
1361            match self.conn_mut()?.poll(qconn) {
1362                Ok((stream_id, event)) =>
1363                    self.process_read_event(qconn, stream_id, event)?,
1364                Err(h3::Error::Done) => break,
1365                Err(err) => {
1366                    // Don't bubble error up, instead keep the worker loop going
1367                    // until quiche reports the connection is
1368                    // closed.
1369                    log::debug!("connection closed due to h3 protocol error"; "error"=>?err);
1370                    return Ok(());
1371                },
1372            };
1373        }
1374
1375        self.process_available_dgrams(qconn)?;
1376        Ok(())
1377    }
1378
1379    /// Write as much data as possible into the [`quiche::h3::Connection`] from
1380    /// all sources. This will attempt to write any queued frames into their
1381    /// respective streams, if writable.
1382    fn process_writes(&mut self, qconn: &mut QuicheConnection) -> QuicResult<()> {
1383        while let Some(stream_id) = qconn.stream_writable_next() {
1384            self.process_writable_stream(qconn, stream_id)?;
1385        }
1386
1387        // Also optimistically check for any ready streams
1388        while let Some(Some(ready)) = self.waiting_streams.next().now_or_never() {
1389            self.upstream_ready(qconn, ready)?;
1390        }
1391
1392        Ok(())
1393    }
1394
1395    /// Reports connection-level error metrics and forwards
1396    /// IOWorker errors to the associated [H3Controller].
1397    fn on_conn_close<M: Metrics>(
1398        &mut self, quiche_conn: &mut QuicheConnection, metrics: &M,
1399        work_loop_result: &QuicResult<()>,
1400    ) {
1401        let max_stream_seen = self.max_stream_seen;
1402        metrics
1403            .maximum_writable_streams()
1404            .observe(max_stream_seen as f64);
1405
1406        Self::record_quiche_error(quiche_conn, metrics);
1407
1408        let Err(work_loop_error) = work_loop_result else {
1409            return;
1410        };
1411
1412        let Some(h3_err) = work_loop_error.downcast_ref::<H3ConnectionError>()
1413        else {
1414            // Errors returned by the IoWorker rather than the driver are
1415            // expected and already counted in metrics:
1416            // - `HandshakeError`: with 0-RTT early data, the handshake can
1417            //   still fail (e.g. time out) after the application was started.
1418            // - `quiche::Error`: the connection was closed by quiche, typically
1419            //   because of a peer protocol violation.
1420            if work_loop_error.is::<HandshakeError>() ||
1421                work_loop_error.is::<quiche::Error>()
1422            {
1423                log::debug!("connection closed by IoWorker"; "error" => %work_loop_error);
1424            } else {
1425                log::error!("Found non-H3ConnectionError"; "error" => %work_loop_error);
1426            }
1427            return;
1428        };
1429
1430        if matches!(h3_err, H3ConnectionError::ControllerWentAway) {
1431            // Inform client that we won't (can't) respond anymore
1432            let _ = quiche_conn.close(true, WireErrorCode::NoError as u64, &[]);
1433            return;
1434        }
1435
1436        if let Some(ev) = H3Event::from_error(h3_err) {
1437            let _ = self.h3_event_sender.send(ev.into());
1438            #[expect(clippy::needless_return)]
1439            return; // avoid accidental fallthrough in the future
1440        }
1441    }
1442
1443    /// Wait for incoming data from the [H3Controller]. The next iteration of
1444    /// the I/O loop commences when one of the `select!`ed futures triggers.
1445    #[inline]
1446    async fn wait_for_data(
1447        &mut self, qconn: &mut QuicheConnection,
1448    ) -> QuicResult<()> {
1449        select! {
1450            biased;
1451            Some(ready) = self.waiting_streams.next() => self.upstream_ready(qconn, ready),
1452            Some(dgram) = self.dgram_recv.recv() => self.dgram_ready(qconn, dgram),
1453            Some(cmd) = self.cmd_recv.recv() => H::conn_command(self, qconn, cmd),
1454            r = self.hooks.wait_for_action(qconn), if H::has_wait_action(self) => r,
1455            _ = self.h3_event_sender.closed(), if !self.h3_event_receiver_dropped => {
1456                self.h3_event_receiver_dropped = true;
1457                self.close_if_idle(qconn);
1458                Ok(())
1459            }
1460        }?;
1461
1462        // Make sure controller is not starved, but also not prioritized in the
1463        // biased select. So poll it last, however also perform a try_recv
1464        // each iteration.
1465        if let Ok(cmd) = self.cmd_recv.try_recv() {
1466            H::conn_command(self, qconn, cmd)?;
1467        }
1468
1469        Ok(())
1470    }
1471}
1472
1473impl<H: DriverHooks> Drop for H3Driver<H> {
1474    fn drop(&mut self) {
1475        for stream in self.stream_map.values() {
1476            stream
1477                .audit_stats
1478                .set_recvd_stream_fin(StreamClosureKind::Implicit);
1479        }
1480    }
1481}
1482
1483/// [`H3Command`]s are sent by the [H3Controller] to alter the [H3Driver]'s
1484/// state.
1485///
1486/// Both [ServerH3Driver] and [ClientH3Driver] may extend this enum with
1487/// endpoint-specific variants.
1488#[derive(Debug)]
1489pub enum H3Command {
1490    /// A connection-level command that executes directly on the
1491    /// [`quiche::Connection`].
1492    QuicCmd(QuicCommand),
1493    /// Send a GOAWAY frame to the peer to initiate a graceful connection
1494    /// shutdown.
1495    GoAway,
1496    /// Shuts down a stream in the specified direction(s) and removes it from
1497    /// local state.
1498    ///
1499    /// This removes the stream from local state and sends a `RESET_STREAM`
1500    /// frame (for write direction) and/or a `STOP_SENDING` frame (for read
1501    /// direction) to the peer. See [`quiche::Connection::stream_shutdown`]
1502    /// for details.
1503    ShutdownStream {
1504        stream_id: u64,
1505        shutdown: StreamShutdown,
1506    },
1507}
1508
1509/// Specifies which direction(s) of a stream to shut down.
1510///
1511/// Used with [`H3Controller::shutdown_stream`] and the internal
1512/// `shutdown_stream` function to control whether to send a `STOP_SENDING` frame
1513/// (read direction), and/or a `RESET_STREAM` frame (write direction)
1514///
1515/// Note: Despite its name, "shutdown" here refers to signaling the peer about
1516/// stream termination, not sending a FIN flag. `STOP_SENDING` asks the peer to
1517/// stop sending data, while `RESET_STREAM` abruptly terminates the write side.
1518#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1519pub enum StreamShutdown {
1520    /// Shut down only the read direction (sends `STOP_SENDING` frame with the
1521    /// given error code).
1522    Read { error_code: u64 },
1523    /// Shut down only the write direction (sends `RESET_STREAM` frame with the
1524    /// given error code).
1525    Write { error_code: u64 },
1526    /// Shut down both directions (sends both `STOP_SENDING` and `RESET_STREAM`
1527    /// frames).
1528    Both {
1529        read_error_code: u64,
1530        write_error_code: u64,
1531    },
1532}
1533
1534/// Sends [`H3Command`]s to an [H3Driver]. The sender is typed and internally
1535/// wraps instances of `T` in the appropriate `H3Command` variant.
1536pub struct RequestSender<C, T> {
1537    sender: UnboundedSender<C>,
1538    // Required to work around dangling type parameter
1539    _r: PhantomData<fn() -> T>,
1540}
1541
1542impl<C, T: Into<C>> RequestSender<C, T> {
1543    /// Send a request to the [H3Driver]. This can only fail if the driver is
1544    /// gone.
1545    #[inline(always)]
1546    pub fn send(&self, v: T) -> Result<(), mpsc::error::SendError<C>> {
1547        self.sender.send(v.into())
1548    }
1549}
1550
1551impl<C, T> Clone for RequestSender<C, T> {
1552    fn clone(&self) -> Self {
1553        Self {
1554            sender: self.sender.clone(),
1555            _r: Default::default(),
1556        }
1557    }
1558}
1559
1560/// Interface to communicate with a paired [H3Driver].
1561///
1562/// An [H3Controller] receives [`H3Event`]s from its driver, which must be
1563/// consumed by the application built on top of the driver to react to incoming
1564/// events. The controller also allows the application to send ad-hoc
1565/// [`H3Command`]s to the driver, which will be processed when the driver waits
1566/// for incoming data.
1567pub struct H3Controller<H: DriverHooks> {
1568    /// Sends [`H3Command`]s to the [H3Driver], like [`QuicCommand`]s or
1569    /// outbound HTTP requests.
1570    cmd_sender: UnboundedSender<H::Command>,
1571    /// Receives [`H3Event`]s from the [H3Driver]. Can be extracted and
1572    /// used independently of the [H3Controller].
1573    h3_event_recv: Option<UnboundedReceiver<H::Event>>,
1574}
1575
1576impl<H: DriverHooks> H3Controller<H> {
1577    /// Gets a mut reference to the [`H3Event`] receiver for the paired
1578    /// [H3Driver].
1579    pub fn event_receiver_mut(&mut self) -> &mut UnboundedReceiver<H::Event> {
1580        self.h3_event_recv
1581            .as_mut()
1582            .expect("No event receiver on H3Controller")
1583    }
1584
1585    /// Takes the [`H3Event`] receiver for the paired [H3Driver].
1586    pub fn take_event_receiver(&mut self) -> UnboundedReceiver<H::Event> {
1587        self.h3_event_recv
1588            .take()
1589            .expect("No event receiver on H3Controller")
1590    }
1591
1592    /// Creates a [`QuicCommand`] sender for the paired [H3Driver].
1593    pub fn cmd_sender(&self) -> RequestSender<H::Command, QuicCommand> {
1594        RequestSender {
1595            sender: self.cmd_sender.clone(),
1596            _r: Default::default(),
1597        }
1598    }
1599
1600    /// Sends a GOAWAY frame to initiate a graceful connection shutdown.
1601    pub fn send_goaway(&self) {
1602        let _ = self.cmd_sender.send(H3Command::GoAway.into());
1603    }
1604
1605    /// Creates an [`H3Command`] sender for the paired [H3Driver].
1606    pub fn h3_cmd_sender(&self) -> RequestSender<H::Command, H3Command> {
1607        RequestSender {
1608            sender: self.cmd_sender.clone(),
1609            _r: Default::default(),
1610        }
1611    }
1612
1613    /// Shuts down a stream in the specified direction(s) and removes it from
1614    /// local state.
1615    ///
1616    /// This removes the stream from local state and sends a `RESET_STREAM`
1617    /// frame (for write direction) and/or a `STOP_SENDING` frame (for read
1618    /// direction) to the peer, depending on the [`StreamShutdown`] variant.
1619    pub fn shutdown_stream(&self, stream_id: u64, shutdown: StreamShutdown) {
1620        let _ = self.cmd_sender.send(
1621            H3Command::ShutdownStream {
1622                stream_id,
1623                shutdown,
1624            }
1625            .into(),
1626        );
1627    }
1628}