From d26fdedfbd7a854df5b268e38ae8e006887a46be Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 10:14:18 -0700 Subject: [PATCH 1/2] fix(web-transport-moq): report a raw QUIC peer's close code A raw QUIC session closes with the application code as is, but the peer decoded ApplicationClosed only through the HTTP/3 code space, so session_error() was None. The session and its streams now share a CloseReason that knows the code space, and close() after the connection already closed keeps that reason instead of recording LocallyClosed. Co-Authored-By: Claude Opus 5.5 --- web-transport-moq/src/error.rs | 55 +++++++++++++++++- web-transport-moq/src/recv.rs | 25 ++++---- web-transport-moq/src/send.rs | 23 ++++---- web-transport-moq/src/session.rs | 83 +++++++++++---------------- web-transport-moq/tests/raw_close.rs | 85 ++++++++++++++++++++++++++++ 5 files changed, 193 insertions(+), 78 deletions(-) create mode 100644 web-transport-moq/tests/raw_close.rs diff --git a/web-transport-moq/src/error.rs b/web-transport-moq/src/error.rs index bb1c9ca86..50476cc8a 100644 --- a/web-transport-moq/src/error.rs +++ b/web-transport-moq/src/error.rs @@ -1,4 +1,4 @@ -use std::sync::Arc; +use std::sync::{Arc, OnceLock}; use thiserror::Error; @@ -68,6 +68,59 @@ impl From for SessionError { } } +/// Why a session closed, shared by the session and its streams. The first close wins. +#[derive(Debug)] +pub(crate) struct CloseReason { + // Raw QUIC carries the application code as is; HTTP/3 maps it into its own code space. + raw: bool, + reason: OnceLock, +} + +impl CloseReason { + pub(crate) fn new(raw: bool) -> Self { + Self { + raw, + reason: OnceLock::new(), + } + } + + /// Record the close reason, returning false if one was already recorded. + pub(crate) fn set(&self, err: SessionError) -> bool { + self.reason.set(err).is_ok() + } + + /// Replace an error caused by the connection closing with the recorded reason, or + /// decode the peer's raw QUIC close code when none was recorded. + pub(crate) fn map(&self, err: SessionError) -> SessionError { + let conn = match &err { + SessionError::ConnectionError(conn) + | SessionError::SendDatagramError(noq::SendDatagramError::ConnectionLost(conn)) => { + Some(conn) + } + SessionError::WebTransportError(WebTransportError::Closed(..)) => None, + _ => return err, + }; + + if let Some(reason) = self.reason.get() { + return reason.clone(); + } + + match conn { + Some(noq::ConnectionError::ApplicationClosed(close)) if self.raw => { + match u32::try_from(close.error_code.into_inner()) { + Ok(code) => WebTransportError::Closed( + code, + String::from_utf8_lossy(&close.reason).into_owned(), + ) + .into(), + Err(_) => err, + } + } + _ => err, + } + } +} + /// An error that can occur when reading/writing the WebTransport stream header. #[derive(Clone, Error, Debug)] pub enum WebTransportError { diff --git a/web-transport-moq/src/recv.rs b/web-transport-moq/src/recv.rs index fa0da2104..25e1059e4 100644 --- a/web-transport-moq/src/recv.rs +++ b/web-transport-moq/src/recv.rs @@ -1,38 +1,35 @@ use std::{ io, pin::Pin, - sync::{Arc, OnceLock}, + sync::Arc, task::{Context, Poll}, }; use bytes::Bytes; -use crate::{ClosedStream, ReadError, ReadExactError, ReadToEndError, SessionError}; +use crate::{CloseReason, ClosedStream, ReadError, ReadExactError, ReadToEndError, SessionError}; /// A stream that can be used to receive bytes. See [`noq::RecvStream`]. #[derive(Debug)] pub struct RecvStream { inner: noq::RecvStream, - error: Arc>, + close: Arc, } impl RecvStream { - pub(crate) fn new(stream: noq::RecvStream, error: Arc>) -> Self { + pub(crate) fn new(stream: noq::RecvStream, close: Arc) -> Self { Self { inner: stream, - error, + close, } } - /// Replace connection-level errors with the stored session error if available. + /// Report a session error as the session's close reason. fn map_error(&self, e: impl Into) -> ReadError { - let e = e.into(); - if let Some(err) = self.error.get() { - if matches!(&e, ReadError::SessionError(_)) { - return ReadError::SessionError(err.clone()); - } + match e.into() { + ReadError::SessionError(e) => ReadError::SessionError(self.close.map(e)), + e => e, } - e } /// Tell the other end to stop sending data with the given error code. See @@ -97,9 +94,7 @@ impl RecvStream { match self.inner.received_reset().await { Ok(None) => Ok(None), Ok(Some(code)) => Ok(web_transport_proto::error_from_http3(code.into_inner())), - Err(noq::ResetError::ConnectionLost(conn_err)) => { - Err(self.error.get().cloned().unwrap_or_else(|| conn_err.into())) - } + Err(noq::ResetError::ConnectionLost(conn_err)) => Err(self.close.map(conn_err.into())), Err(noq::ResetError::ZeroRttRejected) => unreachable!("0-RTT not supported"), } } diff --git a/web-transport-moq/src/send.rs b/web-transport-moq/src/send.rs index 0e67d64cf..26f0a7213 100644 --- a/web-transport-moq/src/send.rs +++ b/web-transport-moq/src/send.rs @@ -1,13 +1,13 @@ use std::{ io, pin::Pin, - sync::{Arc, OnceLock}, + sync::Arc, task::{Context, Poll}, }; use bytes::{Buf, Bytes}; -use crate::{ClosedStream, SessionError, WriteError}; +use crate::{CloseReason, ClosedStream, SessionError, WriteError}; /// A stream that can be used to send bytes. See [`noq::SendStream`]. /// @@ -16,23 +16,20 @@ use crate::{ClosedStream, SessionError, WriteError}; #[derive(Debug)] pub struct SendStream { stream: noq::SendStream, - error: Arc>, + close: Arc, } impl SendStream { - pub(crate) fn new(stream: noq::SendStream, error: Arc>) -> Self { - Self { stream, error } + pub(crate) fn new(stream: noq::SendStream, close: Arc) -> Self { + Self { stream, close } } - /// Replace connection-level errors with the stored session error if available. + /// Report a session error as the session's close reason. fn map_error(&self, e: impl Into) -> WriteError { - let e = e.into(); - if let Some(err) = self.error.get() { - if matches!(&e, WriteError::SessionError(_)) { - return WriteError::SessionError(err.clone()); - } + match e.into() { + WriteError::SessionError(e) => WriteError::SessionError(self.close.map(e)), + e => e, } - e } /// Abruptly reset the stream with the provided error code. See [`noq::SendStream::reset`]. @@ -54,7 +51,7 @@ impl SendStream { Ok(Some(code)) => Ok(web_transport_proto::error_from_http3(code.into_inner())), Ok(None) => Ok(None), Err(noq::StoppedError::ConnectionLost(conn_err)) => { - Err(self.error.get().cloned().unwrap_or_else(|| conn_err.into())) + Err(self.close.map(conn_err.into())) } Err(noq::StoppedError::ZeroRttRejected) => unreachable!("0-RTT not supported"), } diff --git a/web-transport-moq/src/session.rs b/web-transport-moq/src/session.rs index 65314070d..67558fe7a 100644 --- a/web-transport-moq/src/session.rs +++ b/web-transport-moq/src/session.rs @@ -15,7 +15,8 @@ use kio::{Fan, Waiter}; use crate::{ proto::{ConnectRequest, ConnectResponse, Frame, StreamUni, VarInt}, - ClientError, Connected, RecvStream, SendStream, SessionError, Settings, WebTransportError, + ClientError, CloseReason, Connected, RecvStream, SendStream, SessionError, Settings, + WebTransportError, }; /// The ALPN the QUIC handshake negotiated, or `None` if there was none, the handshake @@ -77,10 +78,9 @@ pub struct Session { // Wrapped in Arc>> so close() can take it exactly once. connect_send: Arc>>, - // Session error, set once by either local close() or the background task + // Why the session closed: set once by a local close() or the background task // when a remote CloseWebTransportSession capsule is received. - // Uses OnceLock for set-once, first-writer-wins semantics with lock-free reads. - error: Arc>, + close: Arc, // The request sent by the client, or None for a raw QUIC session. request: Option, @@ -118,10 +118,10 @@ impl Session { let mut header_datagram = Vec::new(); session_id.encode(&mut header_datagram); - let error: Arc> = Arc::new(OnceLock::new()); + let close = Arc::new(CloseReason::new(false)); // Accept logic is stateful, so use an Arc to share it. - let accept = SessionAccept::new(conn.clone(), session_id, error.clone()); + let accept = SessionAccept::new(conn.clone(), session_id, close.clone()); let this = Self { conn, @@ -132,7 +132,7 @@ impl Session { header_datagram, settings: Some(Arc::new(settings)), connect_send: Arc::new(Mutex::new(Some(connect.send))), - error: error.clone(), + close: close.clone(), request: Some(connect.request.clone()), response: Some(connect.response.clone()), alpn: Default::default(), @@ -140,18 +140,14 @@ impl Session { // Run a background task to read capsules from the CONNECT recv stream. let conn2 = this.conn.clone(); - tokio::spawn(Self::run_recv(conn2, connect.recv, error)); + tokio::spawn(Self::run_recv(conn2, connect.recv, close)); this } // Read capsules from the CONNECT recv stream until it's closed, // then record the close error and tear down the connection. - async fn run_recv( - conn: noq::Connection, - recv: noq::RecvStream, - error: Arc>, - ) { + async fn run_recv(conn: noq::Connection, recv: noq::RecvStream, close: Arc) { let close_info = Self::read_capsules(recv).await; let code = close_info.as_ref().map_or(0, |(c, _)| *c); @@ -164,14 +160,14 @@ impl Session { match close_info { Some((code, reason)) => { let err = WebTransportError::Closed(code, reason.clone()); - if error.set(err.into()).is_err() { + if !close.set(err.into()) { return; } conn.close(http3_code, reason.as_bytes()); } None => { let err = noq::ConnectionError::LocallyClosed.into(); - if error.set(err).is_err() { + if !close.set(err) { return; } conn.close(http3_code, b""); @@ -239,7 +235,7 @@ impl Session { .accept_uni() .await .map_err(|e| self.map_error(e))?; - Ok(RecvStream::new(recv, self.error.clone())) + Ok(RecvStream::new(recv, self.close.clone())) } } @@ -252,8 +248,8 @@ impl Session { } else { let (send, recv) = self.conn.accept_bi().await.map_err(|e| self.map_error(e))?; Ok(( - SendStream::new(send, self.error.clone()), - RecvStream::new(recv, self.error.clone()), + SendStream::new(send, self.close.clone()), + RecvStream::new(recv, self.close.clone()), )) } } @@ -273,7 +269,7 @@ impl Session { // Reset the stream priority back to the default of 0. send.set_priority(0).ok(); - Ok(SendStream::new(send, self.error.clone())) + Ok(SendStream::new(send, self.close.clone())) } /// Open a new bidirectional stream. See [`noq::Connection::open_bi`]. @@ -292,8 +288,8 @@ impl Session { // Reset the stream priority back to the default of 0. send.set_priority(0).ok(); Ok(( - SendStream::new(send, self.error.clone()), - RecvStream::new(recv, self.error.clone()), + SendStream::new(send, self.close.clone()), + RecvStream::new(recv, self.close.clone()), )) } @@ -401,11 +397,14 @@ impl Session { /// Callers should `await` [`Session::closed()`] to ensure the capsule has been /// delivered. Session operations will fail once the QUIC connection is closed. pub fn close(&self, code: u32, reason: &[u8]) { - // Record the local close error. First writer wins — if the background - // task already set a remote close error, or close() was already called, - // this is a no-op. + // First close wins: a no-op if the connection already closed, the background + // task recorded a remote close, or close() was already called. + if let Some(err) = self.conn.close_reason() { + self.close.set(self.map_error(err)); + return; + } let err = SessionError::ConnectionError(noq::ConnectionError::LocallyClosed); - if self.error.set(err).is_err() { + if !self.close.set(err) { return; } @@ -503,19 +502,9 @@ impl Session { self.conn.close_reason().map(|e| self.map_error(e)) } - /// Replace connection-level errors with the stored session error if available. + /// Report an error caused by the connection closing as the session's close reason. fn map_error(&self, e: impl Into) -> SessionError { - let e = e.into(); - if let Some(err) = self.error.get() { - if matches!( - &e, - SessionError::ConnectionError(_) - | SessionError::SendDatagramError(noq::SendDatagramError::ConnectionLost(_)) - ) { - return err.clone(); - } - } - e + self.close.map(e.into()) } async fn write_full(send: &mut noq::SendStream, buf: &[u8]) -> Result<(), SessionError> { @@ -545,7 +534,7 @@ impl Session { accept: None, settings: None, connect_send: Arc::new(Mutex::new(None)), - error: Arc::new(OnceLock::new()), + close: Arc::new(CloseReason::new(true)), request: None, response: None, alpn: Default::default(), @@ -695,8 +684,8 @@ type PendingBi = pub struct SessionAccept { session_id: VarInt, - // Shared session error for propagation to accepted streams. - error: Arc>, + // Shared close reason for propagation to accepted streams. + close: Arc, // We also need to keep a reference to the qpack streams if the endpoint (incorrectly) creates // them. Again, this is just so they don't get closed until we drop the session. @@ -726,11 +715,7 @@ pub struct SessionAccept { } impl SessionAccept { - pub(crate) fn new( - conn: noq::Connection, - session_id: VarInt, - error: Arc>, - ) -> Self { + pub(crate) fn new(conn: noq::Connection, session_id: VarInt, close: Arc) -> Self { // Create a stream that just outputs new streams, so it's easy to call from poll. let accept_uni = Box::pin(futures::stream::unfold(conn.clone(), |conn| async { Some((conn.accept_uni().await, conn)) @@ -747,7 +732,7 @@ impl SessionAccept { Self { session_id, - error, + close, qpack_decoder: None, qpack_encoder: None, @@ -818,7 +803,7 @@ impl SessionAccept { // Decide if we keep looping based on the type. match typ { StreamUni::WEBTRANSPORT => { - let recv = RecvStream::new(recv, self.error.clone()); + let recv = RecvStream::new(recv, self.close.clone()); return Poll::Ready(Ok(recv)); } StreamUni::QPACK_DECODER => { @@ -905,8 +890,8 @@ impl SessionAccept { if let Some((send, recv)) = res { // Wrap the streams in our own types for correct error codes. - let send = SendStream::new(send, self.error.clone()); - let recv = RecvStream::new(recv, self.error.clone()); + let send = SendStream::new(send, self.close.clone()); + let recv = RecvStream::new(recv, self.close.clone()); return Poll::Ready(Ok((send, recv))); } diff --git a/web-transport-moq/tests/raw_close.rs b/web-transport-moq/tests/raw_close.rs new file mode 100644 index 000000000..8766de032 --- /dev/null +++ b/web-transport-moq/tests/raw_close.rs @@ -0,0 +1,85 @@ +//! A raw QUIC session carries the application close code directly, with no HTTP/3 +//! mapping, so the peer must decode it the same way and keep the first close. + +use std::{net::Ipv4Addr, sync::Arc, time::Duration}; + +use anyhow::{Context as _, Result}; +use rcgen::{CertifiedKey, KeyPair}; +use rustls::pki_types::{PrivateKeyDer, PrivatePkcs8KeyDer}; +use tokio::time::timeout; +use web_transport_moq::{generic, noq, Session}; + +/// With both crypto features enabled, as `--all-features` does, the crate cannot pick a +/// provider for us, so the test installs one. +fn install_crypto_provider() { + #[cfg(all(feature = "aws-lc-rs", feature = "ring"))] + { + let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); + } +} + +/// A connected raw QUIC client and server session, plus the endpoints that drive them. +async fn raw_pair() -> Result<(Session, Session, [noq::Endpoint; 2])> { + install_crypto_provider(); + + let CertifiedKey { cert, signing_key } = + rcgen::generate_simple_self_signed(vec!["localhost".into()])?; + let key = PrivateKeyDer::Pkcs8(PrivatePkcs8KeyDer::from(KeyPair::serialize_der( + &signing_key, + ))); + + let server_config = noq::ServerConfig::with_single_cert(vec![cert.der().clone()], key)?; + let server = noq::Endpoint::server(server_config, (Ipv4Addr::LOCALHOST, 0).into())?; + + let mut roots = rustls::RootCertStore::empty(); + roots.add(cert.der().clone())?; + let client_config = noq::ClientConfig::with_root_certificates(Arc::new(roots))?; + let client = noq::Endpoint::client((Ipv4Addr::LOCALHOST, 0).into())?; + + let connecting = client.connect_with(client_config, server.local_addr()?, "localhost")?; + let accept = async { anyhow::Ok(server.accept().await.context("no connection")?.await?) }; + let (client_conn, server_conn) = + tokio::try_join!(async { anyhow::Ok(connecting.await?) }, accept)?; + + Ok(( + Session::raw(client_conn), + Session::raw(server_conn), + [client, server], + )) +} + +fn assert_code(err: &impl generic::Error, code: u32) { + assert_eq!( + err.session_error(), + Some((code, "kicked".to_string())), + "{err}" + ); +} + +/// The peer's raw close code reaches `closed`, the accept and stream paths, and +/// survives a later local `close`. +#[tokio::test] +async fn peer_close_code() -> Result<()> { + let (client, server, _endpoints) = raw_pair().await?; + + let mut send = server.open_uni().await?; + send.write_all(b"x").await?; + let mut recv = timeout(Duration::from_secs(5), client.accept_uni()).await??; + let mut buf = [0u8; 1]; + recv.read_exact(&mut buf).await?; + + server.close(4075, b"kicked"); + + let err = timeout(Duration::from_secs(5), client.closed()).await?; + assert_code(&err, 4075); + assert_code(&client.accept_uni().await.err().unwrap(), 4075); + assert_code(&client.open_uni().await.err().unwrap(), 4075); + assert_code(&recv.read(&mut buf).await.err().unwrap(), 4075); + + // Closing an already closed session changes nothing. + client.close(1, b""); + assert_code(&client.closed().await, 4075); + assert_code(&client.close_reason().unwrap(), 4075); + + Ok(()) +} From e82e02a4ffb9e7f421c34ba4ae0411aedd57da09 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 12:24:24 -0700 Subject: [PATCH 2/2] test(web-transport-moq): use web-transport-trait directly Co-Authored-By: Claude Opus 5.5 --- web-transport-moq/tests/raw_close.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/web-transport-moq/tests/raw_close.rs b/web-transport-moq/tests/raw_close.rs index 8766de032..d5993c807 100644 --- a/web-transport-moq/tests/raw_close.rs +++ b/web-transport-moq/tests/raw_close.rs @@ -7,7 +7,7 @@ use anyhow::{Context as _, Result}; use rcgen::{CertifiedKey, KeyPair}; use rustls::pki_types::{PrivateKeyDer, PrivatePkcs8KeyDer}; use tokio::time::timeout; -use web_transport_moq::{generic, noq, Session}; +use web_transport_moq::{noq, Session}; /// With both crypto features enabled, as `--all-features` does, the crate cannot pick a /// provider for us, so the test installs one. @@ -48,7 +48,7 @@ async fn raw_pair() -> Result<(Session, Session, [noq::Endpoint; 2])> { )) } -fn assert_code(err: &impl generic::Error, code: u32) { +fn assert_code(err: &impl web_transport_trait::Error, code: u32) { assert_eq!( err.session_error(), Some((code, "kicked".to_string())),