Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
55 changes: 54 additions & 1 deletion web-transport-moq/src/error.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
use std::sync::Arc;
use std::sync::{Arc, OnceLock};

use thiserror::Error;

Expand Down Expand Up @@ -68,6 +68,59 @@ impl From<noq::ConnectionError> 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<SessionError>,
}

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 {
Expand Down
25 changes: 10 additions & 15 deletions web-transport-moq/src/recv.rs
Original file line number Diff line number Diff line change
@@ -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<OnceLock<SessionError>>,
close: Arc<CloseReason>,
}

impl RecvStream {
pub(crate) fn new(stream: noq::RecvStream, error: Arc<OnceLock<SessionError>>) -> Self {
pub(crate) fn new(stream: noq::RecvStream, close: Arc<CloseReason>) -> 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>) -> 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
Expand Down Expand Up @@ -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"),
}
}
Expand Down
23 changes: 10 additions & 13 deletions web-transport-moq/src/send.rs
Original file line number Diff line number Diff line change
@@ -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`].
///
Expand All @@ -16,23 +16,20 @@ use crate::{ClosedStream, SessionError, WriteError};
#[derive(Debug)]
pub struct SendStream {
stream: noq::SendStream,
error: Arc<OnceLock<SessionError>>,
close: Arc<CloseReason>,
}

impl SendStream {
pub(crate) fn new(stream: noq::SendStream, error: Arc<OnceLock<SessionError>>) -> Self {
Self { stream, error }
pub(crate) fn new(stream: noq::SendStream, close: Arc<CloseReason>) -> 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>) -> 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`].
Expand All @@ -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"),
}
Expand Down
Loading
Loading