0.3.5-alpha.3
This commit is contained in:
@@ -1,9 +1,11 @@
|
||||
// file: crates/common/game-realtime-webtransport-lib/src/webtransport.rs
|
||||
// version: 1
|
||||
// version: 2
|
||||
|
||||
const CERTIFICATE_HASH_SIZE: usize = 32;
|
||||
const LOCAL_CERTIFICATE_CLOCK_SKEW: std::time::Duration = std::time::Duration::from_secs(60);
|
||||
const LOCAL_CERTIFICATE_VALIDITY: std::time::Duration = std::time::Duration::from_secs(7 * 24 * 60 * 60);
|
||||
const PRIMARY_FRAME_HEADER_SIZE: usize = 4;
|
||||
const PRIMARY_FRAME_MAX_PAYLOAD_SIZE: usize = 1024 * 1024;
|
||||
const TRACING_TARGET: &str = "games::realtime::webtransport";
|
||||
|
||||
/// SHA-256 fingerprint of one certificate accepted by the native WebTransport client.
|
||||
@@ -151,7 +153,7 @@ impl WebTransportServerConfig {
|
||||
}
|
||||
}
|
||||
|
||||
/// Established native WebTransport session before application-stream adaptation.
|
||||
/// Established native WebTransport session before or while the single primary application stream is selected.
|
||||
pub struct WebTransportSession {
|
||||
inner: web_transport_quinn::Session,
|
||||
}
|
||||
@@ -161,6 +163,37 @@ impl WebTransportSession {
|
||||
return Self { inner };
|
||||
}
|
||||
|
||||
/// Accepts the peer-created primary bidirectional stream and adapts it to the transport-neutral realtime contract.
|
||||
///
|
||||
/// The native WebTransport wrapper writes the required stream/session header while opening the stream, so the peer can
|
||||
/// accept it before the first application frame is sent.
|
||||
pub async fn accept_primary_connection(self) -> Result<WebTransportConnection, game_realtime_transport_lib::TransportError> {
|
||||
let (sender, receiver) = match self.inner.accept_bi().await {
|
||||
Ok(value) => value,
|
||||
Err(error) => {
|
||||
let mapped = transport_error(game_realtime_transport_lib::TransportErrorKind::Protocol, error.to_string());
|
||||
tracing::warn!(target: TRACING_TARGET, detail = mapped.detail(), "WebTransport primary bidirectional stream accept failed");
|
||||
return Err(mapped);
|
||||
},
|
||||
};
|
||||
tracing::debug!(target: TRACING_TARGET, peer = %self.inner.remote_address(), "WebTransport primary bidirectional stream accepted");
|
||||
return Ok(WebTransportConnection::new(self.inner, sender, receiver));
|
||||
}
|
||||
|
||||
/// Opens the single primary bidirectional stream and adapts it to the transport-neutral realtime contract.
|
||||
pub async fn open_primary_connection(self) -> Result<WebTransportConnection, game_realtime_transport_lib::TransportError> {
|
||||
let (sender, receiver) = match self.inner.open_bi().await {
|
||||
Ok(value) => value,
|
||||
Err(error) => {
|
||||
let mapped = transport_error(game_realtime_transport_lib::TransportErrorKind::Protocol, error.to_string());
|
||||
tracing::warn!(target: TRACING_TARGET, detail = mapped.detail(), "WebTransport primary bidirectional stream open failed");
|
||||
return Err(mapped);
|
||||
},
|
||||
};
|
||||
tracing::debug!(target: TRACING_TARGET, peer = %self.inner.remote_address(), "WebTransport primary bidirectional stream opened");
|
||||
return Ok(WebTransportConnection::new(self.inner, sender, receiver));
|
||||
}
|
||||
|
||||
/// Returns the remote UDP socket address backing the established QUIC connection.
|
||||
#[must_use]
|
||||
pub fn remote_addr(&self) -> std::net::SocketAddr {
|
||||
@@ -177,6 +210,121 @@ impl WebTransportSession {
|
||||
}
|
||||
}
|
||||
|
||||
/// Established WebTransport realtime connection carried by one primary reliable bidirectional stream.
|
||||
pub struct WebTransportConnection {
|
||||
receiver: web_transport_quinn::RecvStream,
|
||||
sender: web_transport_quinn::SendStream,
|
||||
session: web_transport_quinn::Session,
|
||||
}
|
||||
|
||||
impl WebTransportConnection {
|
||||
fn new(session: web_transport_quinn::Session, sender: web_transport_quinn::SendStream, receiver: web_transport_quinn::RecvStream) -> Self {
|
||||
return Self { receiver, sender, session };
|
||||
}
|
||||
}
|
||||
|
||||
impl game_realtime_transport_lib::RealtimeConnection for WebTransportConnection {
|
||||
type Receiver = crate::WebTransportReceiver;
|
||||
type Sender = crate::WebTransportSender;
|
||||
|
||||
fn split(self) -> (Self::Sender, Self::Receiver) {
|
||||
let receiver_session = self.session.clone();
|
||||
return (
|
||||
crate::WebTransportSender { inner: self.sender, _session: self.session },
|
||||
crate::WebTransportReceiver { inner: self.receiver, _session: receiver_session },
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// Receive half of the primary reliable WebTransport stream.
|
||||
pub struct WebTransportReceiver {
|
||||
inner: web_transport_quinn::RecvStream,
|
||||
_session: web_transport_quinn::Session,
|
||||
}
|
||||
|
||||
impl game_realtime_transport_lib::RealtimeReceiver for WebTransportReceiver {
|
||||
type ReceiveFuture<'a>
|
||||
= std::pin::Pin<
|
||||
Box<dyn core::future::Future<Output = Result<game_realtime_transport_lib::TransportReceive, game_realtime_transport_lib::TransportError>> + 'a>,
|
||||
>
|
||||
where
|
||||
Self: 'a;
|
||||
|
||||
fn receive(&mut self) -> Self::ReceiveFuture<'_> {
|
||||
return Box::pin(async move {
|
||||
let payload_len = match read_frame_payload_len(&mut self.inner).await {
|
||||
Ok(Some(value)) => value,
|
||||
Ok(None) => {
|
||||
tracing::debug!(target: TRACING_TARGET, "remote WebTransport primary stream closed cleanly");
|
||||
return Ok(game_realtime_transport_lib::TransportReceive::Closed);
|
||||
},
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
let payload = match read_frame_payload(&mut self.inner, payload_len).await {
|
||||
Ok(value) => value,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
tracing::trace!(target: TRACING_TARGET, payload_len = payload.len(), "framed WebTransport payload received");
|
||||
return Ok(game_realtime_transport_lib::TransportReceive::Message(game_realtime_transport_lib::TransportMessage::new(payload)));
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/// Send half of the primary reliable WebTransport stream.
|
||||
pub struct WebTransportSender {
|
||||
inner: web_transport_quinn::SendStream,
|
||||
_session: web_transport_quinn::Session,
|
||||
}
|
||||
|
||||
impl game_realtime_transport_lib::RealtimeSender for WebTransportSender {
|
||||
type CloseFuture<'a>
|
||||
= std::pin::Pin<Box<dyn core::future::Future<Output = Result<(), game_realtime_transport_lib::TransportError>> + 'a>>
|
||||
where
|
||||
Self: 'a;
|
||||
type SendFuture<'a>
|
||||
= std::pin::Pin<Box<dyn core::future::Future<Output = Result<(), game_realtime_transport_lib::TransportError>> + 'a>>
|
||||
where
|
||||
Self: 'a;
|
||||
|
||||
fn close(&mut self) -> Self::CloseFuture<'_> {
|
||||
return Box::pin(async move {
|
||||
return match self.inner.finish() {
|
||||
Ok(()) => {
|
||||
tracing::debug!(target: TRACING_TARGET, "local WebTransport primary stream close initiated");
|
||||
Ok(())
|
||||
},
|
||||
Err(error) => {
|
||||
let mapped = transport_error(game_realtime_transport_lib::TransportErrorKind::Closed, error.to_string());
|
||||
tracing::warn!(target: TRACING_TARGET, detail = mapped.detail(), "WebTransport primary stream close failed");
|
||||
Err(mapped)
|
||||
},
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
fn send(&mut self, message: game_realtime_transport_lib::TransportMessage) -> Self::SendFuture<'_> {
|
||||
return Box::pin(async move {
|
||||
let payload_len = message.len();
|
||||
let frame_header = match frame_header(payload_len) {
|
||||
Ok(value) => value,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
if let Err(error) = self.inner.write_all(&frame_header).await {
|
||||
let mapped = transport_error(game_realtime_transport_lib::TransportErrorKind::Protocol, error.to_string());
|
||||
tracing::warn!(target: TRACING_TARGET, payload_len = payload_len, detail = mapped.detail(), "WebTransport frame header send failed");
|
||||
return Err(mapped);
|
||||
}
|
||||
if let Err(error) = self.inner.write_all(message.as_bytes()).await {
|
||||
let mapped = transport_error(game_realtime_transport_lib::TransportErrorKind::Protocol, error.to_string());
|
||||
tracing::warn!(target: TRACING_TARGET, payload_len = payload_len, detail = mapped.detail(), "WebTransport frame payload send failed");
|
||||
return Err(mapped);
|
||||
}
|
||||
tracing::trace!(target: TRACING_TARGET, payload_len = payload_len, "framed WebTransport payload sent");
|
||||
return Ok(());
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/// Bound native WebTransport server endpoint that accepts HTTP/3 WebTransport sessions.
|
||||
pub struct WebTransportListener {
|
||||
server: web_transport_quinn::Server,
|
||||
@@ -270,10 +418,75 @@ fn certificate_hash(certificate_der: &[u8]) -> Result<WebTransportCertificateHas
|
||||
return Ok(WebTransportCertificateHash::from_sha256(bytes));
|
||||
}
|
||||
|
||||
fn frame_header(payload_len: usize) -> Result<[u8; PRIMARY_FRAME_HEADER_SIZE], game_realtime_transport_lib::TransportError> {
|
||||
if payload_len > PRIMARY_FRAME_MAX_PAYLOAD_SIZE {
|
||||
return Err(message_too_large(payload_len));
|
||||
}
|
||||
let payload_len = match u32::try_from(payload_len) {
|
||||
Ok(value) => value,
|
||||
Err(_) => return Err(message_too_large(payload_len)),
|
||||
};
|
||||
return Ok(payload_len.to_be_bytes());
|
||||
}
|
||||
|
||||
fn invalid_configuration(detail: impl Into<String>) -> game_realtime_transport_lib::TransportError {
|
||||
return transport_error(game_realtime_transport_lib::TransportErrorKind::InvalidConfiguration, detail);
|
||||
}
|
||||
|
||||
fn message_too_large(payload_len: usize) -> game_realtime_transport_lib::TransportError {
|
||||
return transport_error(
|
||||
game_realtime_transport_lib::TransportErrorKind::MessageTooLarge,
|
||||
format!("WebTransport primary frame payload size {payload_len} exceeds {PRIMARY_FRAME_MAX_PAYLOAD_SIZE} bytes"),
|
||||
);
|
||||
}
|
||||
|
||||
async fn read_frame_payload(stream: &mut web_transport_quinn::RecvStream, payload_len: usize) -> Result<Vec<u8>, game_realtime_transport_lib::TransportError> {
|
||||
let mut payload = vec![0_u8; payload_len];
|
||||
let mut offset = 0_usize;
|
||||
while offset < payload_len {
|
||||
let read = match stream.read(&mut payload[offset..]).await {
|
||||
Ok(Some(value)) => value,
|
||||
Ok(None) => return Err(protocol_error("WebTransport primary stream closed in the middle of a frame payload")),
|
||||
Err(error) => return Err(protocol_error(error.to_string())),
|
||||
};
|
||||
if read == 0 {
|
||||
return Err(protocol_error("WebTransport primary stream returned an empty read in the middle of a frame payload"));
|
||||
}
|
||||
offset += read;
|
||||
}
|
||||
return Ok(payload);
|
||||
}
|
||||
|
||||
async fn read_frame_payload_len(stream: &mut web_transport_quinn::RecvStream) -> Result<Option<usize>, game_realtime_transport_lib::TransportError> {
|
||||
let mut header = [0_u8; PRIMARY_FRAME_HEADER_SIZE];
|
||||
let mut offset = 0_usize;
|
||||
while offset < PRIMARY_FRAME_HEADER_SIZE {
|
||||
let read = match stream.read(&mut header[offset..]).await {
|
||||
Ok(Some(value)) => value,
|
||||
Ok(None) => {
|
||||
if offset == 0 {
|
||||
return Ok(None);
|
||||
}
|
||||
return Err(protocol_error("WebTransport primary stream closed in the middle of a frame header"));
|
||||
},
|
||||
Err(error) => return Err(protocol_error(error.to_string())),
|
||||
};
|
||||
if read == 0 {
|
||||
return Err(protocol_error("WebTransport primary stream returned an empty read in the middle of a frame header"));
|
||||
}
|
||||
offset += read;
|
||||
}
|
||||
let payload_len = u32::from_be_bytes(header) as usize;
|
||||
if payload_len > PRIMARY_FRAME_MAX_PAYLOAD_SIZE {
|
||||
return Err(message_too_large(payload_len));
|
||||
}
|
||||
return Ok(Some(payload_len));
|
||||
}
|
||||
|
||||
fn protocol_error(detail: impl Into<String>) -> game_realtime_transport_lib::TransportError {
|
||||
return transport_error(game_realtime_transport_lib::TransportErrorKind::Protocol, detail);
|
||||
}
|
||||
|
||||
fn transport_error(kind: game_realtime_transport_lib::TransportErrorKind, detail: impl Into<String>) -> game_realtime_transport_lib::TransportError {
|
||||
return game_realtime_transport_lib::TransportError::new(kind, detail);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user