// file: crates/common/game-realtime-webtransport-lib/src/webtransport_wasm.rs // version: 3 const CERTIFICATE_HASH_SIZE: usize = 32; const FRAME_PROTOCOL_ERROR_CODE: u32 = 0x10; const FRAME_TOO_LARGE_ERROR_CODE: u32 = 0x11; const PRIMARY_FRAME_HEADER_SIZE: usize = 4; const SEND_FAILURE_ERROR_CODE: u32 = 0x12; const SEND_TIMEOUT_ERROR_CODE: u32 = 0x13; const STREAM_CANCELLED_ERROR_CODE: u32 = 0x14; const TRACING_TARGET: &str = "games::realtime::webtransport"; /// SHA-256 fingerprint of one certificate accepted by the browser WebTransport client. #[derive(Clone, Debug, Eq, PartialEq)] pub struct WebTransportCertificateHash { bytes: [u8; CERTIFICATE_HASH_SIZE], } impl WebTransportCertificateHash { /// Creates a fingerprint from an already-computed SHA-256 digest. #[must_use] pub fn from_sha256(bytes: [u8; CERTIFICATE_HASH_SIZE]) -> Self { return Self { bytes }; } /// Returns the exact 32-byte SHA-256 digest. #[must_use] pub fn as_bytes(&self) -> &[u8; CERTIFICATE_HASH_SIZE] { return &self.bytes; } } /// Browser WebTransport endpoint, certificate pin and reliable-path configuration. #[derive(Clone, Debug, Eq, PartialEq)] pub struct WebTransportClientConfig { endpoint: url::Url, certificate_hash: WebTransportCertificateHash, transport: crate::WebTransportConfig, } impl WebTransportClientConfig { /// Parses and validates a secure WebTransport endpoint with one pinned SHA-256 certificate fingerprint. pub fn new(endpoint: &str, certificate_hash: WebTransportCertificateHash) -> Result { let parsed = match url::Url::parse(endpoint) { Ok(value) => value, Err(error) => return Err(invalid_configuration(error.to_string())), }; if parsed.scheme() != "https" { return Err(invalid_configuration("WebTransport endpoint scheme must be https")); } if parsed.host().is_none() { return Err(invalid_configuration("WebTransport endpoint must contain a host")); } return Ok(Self { endpoint: parsed, certificate_hash, transport: crate::WebTransportConfig::default() }); } /// Returns a copy with explicit reliable-path limits and deadlines. /// /// The browser path applies the message-size limit and operation deadlines to connect, primary-stream open and send. #[must_use] pub fn with_transport_config(mut self, transport: crate::WebTransportConfig) -> Self { self.transport = transport; return self; } /// Returns the validated WebTransport endpoint URL. #[must_use] pub fn endpoint(&self) -> &str { return self.endpoint.as_str(); } /// Returns the pinned SHA-256 server-certificate fingerprint. #[must_use] pub fn certificate_hash(&self) -> &WebTransportCertificateHash { return &self.certificate_hash; } /// Returns the reliable-path limits and deadlines. #[must_use] pub fn transport_config(&self) -> crate::WebTransportConfig { return self.transport; } } /// Established browser WebTransport session before the primary application stream is selected. pub struct WebTransportSession { inner: web_transport_wasm::Session, transport: crate::WebTransportConfig, } impl WebTransportSession { fn new(inner: web_transport_wasm::Session, transport: crate::WebTransportConfig) -> Self { return Self { inner, transport }; } /// Opens the single primary bidirectional stream and adapts it to the transport-neutral realtime contract. pub async fn open_primary_connection(self) -> Result { let timeout = self.transport.primary_stream_timeout(); let timeout_millis = match browser_timeout_millis(timeout, "primary_stream_timeout") { Ok(value) => value, Err(error) => return Err(error), }; let opened = await_with_timeout(self.inner.open_bi(), timeout_millis).await; let (sender, receiver) = match opened { Some(Ok(value)) => value, Some(Err(error)) => { let mapped = transport_error(game_realtime_transport_lib::TransportErrorKind::Protocol, error.to_string()); tracing::warn!(target: TRACING_TARGET, detail = mapped.detail(), "browser WebTransport primary bidirectional stream open failed"); return Err(mapped); }, None => { let mapped = timeout_error("browser WebTransport primary bidirectional stream open", timeout); tracing::warn!(target: TRACING_TARGET, timeout_ms = timeout.as_millis(), "browser WebTransport primary bidirectional stream open timed out"); return Err(mapped); }, }; tracing::debug!(target: TRACING_TARGET, endpoint = self.inner.url().as_str(), "browser WebTransport primary bidirectional stream opened"); return Ok(WebTransportConnection::new(self.inner, sender, receiver, self.transport)); } /// Returns the endpoint URL backing this browser WebTransport session. #[must_use] pub fn endpoint(&self) -> &str { return self.inner.url().as_str(); } /// Returns the maximum payload size accepted by the browser WebTransport datagram path. #[must_use] pub fn max_datagram_size(&self) -> usize { return self.inner.max_datagram_size(); } /// Sends one browser WebTransport datagram without reliability or ordering guarantees. /// /// This capability is intentionally not part of `RealtimeConnection`. pub async fn send_datagram(&self, payload: &[u8]) -> Result<(), game_realtime_transport_lib::TransportError> { let max_datagram_size = self.max_datagram_size(); if payload.len() > max_datagram_size { return Err(datagram_too_large(payload.len(), max_datagram_size)); } return match self.inner.send_datagram(payload.to_vec().into()).await { Ok(()) => Ok(()), Err(error) => Err(map_write_error(error)), }; } /// Receives one browser WebTransport datagram without reliability or ordering guarantees. /// /// The caller is responsible for applying any operation deadline required by its use case. pub async fn receive_datagram(&self) -> Result, game_realtime_transport_lib::TransportError> { return match self.inner.recv_datagram().await { Ok(payload) => Ok(payload.to_vec()), Err(error) => Err(map_read_error(error)), }; } } /// Established browser WebTransport connection carrying the single reliable primary stream. pub struct WebTransportConnection { receiver: web_transport_wasm::RecvStream, sender: web_transport_wasm::SendStream, session: web_transport_wasm::Session, transport: crate::WebTransportConfig, } impl WebTransportConnection { fn new( session: web_transport_wasm::Session, sender: web_transport_wasm::SendStream, receiver: web_transport_wasm::RecvStream, transport: crate::WebTransportConfig, ) -> Self { return Self { receiver, sender, session, transport }; } } 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, max_message_size: self.transport.max_message_size(), send_timeout: self.transport.send_timeout(), terminal: false, }, crate::WebTransportReceiver { inner: self.receiver, _session: receiver_session, max_message_size: self.transport.max_message_size(), header: [0_u8; PRIMARY_FRAME_HEADER_SIZE], header_read: 0, payload: Vec::new(), payload_read: 0, clean_closed: false, terminal: false, }, ); } } /// Receive half of the browser primary reliable WebTransport stream. pub struct WebTransportReceiver { inner: web_transport_wasm::RecvStream, _session: web_transport_wasm::Session, max_message_size: usize, header: [u8; PRIMARY_FRAME_HEADER_SIZE], header_read: usize, payload: Vec, payload_read: usize, clean_closed: bool, terminal: bool, } impl WebTransportReceiver { /// Abruptly stops the reliable receive direction with one WebTransport application error code. pub fn abort(&mut self, code: u32) -> Result<(), game_realtime_transport_lib::TransportError> { if self.clean_closed || self.terminal { return Err(transport_error(game_realtime_transport_lib::TransportErrorKind::Closed, "browser WebTransport receiver is already terminal")); } self.inner.stop(code); self.terminal = true; tracing::debug!(target: TRACING_TARGET, code = code, "browser WebTransport primary receive stream aborted"); return Ok(()); } async fn receive_frame(&mut self) -> Result { if self.clean_closed { return Ok(game_realtime_transport_lib::TransportReceive::Closed); } if self.terminal { return Err(transport_error( game_realtime_transport_lib::TransportErrorKind::Aborted, "browser WebTransport receiver is unavailable after a terminal stream failure or abort", )); } loop { if self.header_read < PRIMARY_FRAME_HEADER_SIZE { let remaining = PRIMARY_FRAME_HEADER_SIZE - self.header_read; let read = self.inner.read(remaining).await; match read { Ok(Some(chunk)) => { if chunk.is_empty() { return self.fail_protocol("browser WebTransport primary stream returned an empty read in the middle of a frame header"); } let end = self.header_read + chunk.len(); self.header[self.header_read..end].copy_from_slice(chunk.as_ref()); self.header_read = end; continue; }, Ok(None) => { if self.header_read == 0 { self.clean_closed = true; tracing::debug!(target: TRACING_TARGET, "remote browser WebTransport primary stream closed cleanly"); return Ok(game_realtime_transport_lib::TransportReceive::Closed); } return self.fail_protocol("browser WebTransport primary stream closed in the middle of a frame header"); }, Err(error) => return self.fail_read(error), } } if self.payload.is_empty() && self.payload_read == 0 { let payload_len = u32::from_be_bytes(self.header) as usize; if payload_len > self.max_message_size { let error = message_too_large(payload_len, self.max_message_size); self.stop_after_failure(FRAME_TOO_LARGE_ERROR_CODE); return Err(error); } if payload_len == 0 { self.reset_frame_state(); tracing::trace!(target: TRACING_TARGET, payload_len = 0, "framed browser WebTransport payload received"); return Ok(game_realtime_transport_lib::TransportReceive::Message(game_realtime_transport_lib::TransportMessage::new(Vec::new()))); } self.payload = vec![0_u8; payload_len]; } if self.payload_read < self.payload.len() { let remaining = self.payload.len() - self.payload_read; let read = self.inner.read(remaining).await; match read { Ok(Some(chunk)) => { if chunk.is_empty() { return self.fail_protocol("browser WebTransport primary stream returned an empty read in the middle of a frame payload"); } let end = self.payload_read + chunk.len(); self.payload[self.payload_read..end].copy_from_slice(chunk.as_ref()); self.payload_read = end; if self.payload_read < self.payload.len() { continue; } }, Ok(None) => return self.fail_protocol("browser WebTransport primary stream closed in the middle of a frame payload"), Err(error) => return self.fail_read(error), } } let payload = core::mem::take(&mut self.payload); self.reset_frame_state(); tracing::trace!(target: TRACING_TARGET, payload_len = payload.len(), "framed browser WebTransport payload received"); return Ok(game_realtime_transport_lib::TransportReceive::Message(game_realtime_transport_lib::TransportMessage::new(payload))); } } fn fail_protocol(&mut self, detail: &str) -> Result { let error = protocol_error(detail); self.stop_after_failure(FRAME_PROTOCOL_ERROR_CODE); return Err(error); } fn fail_read( &mut self, error: web_transport_wasm::Error, ) -> Result { self.terminal = true; let mapped = map_read_error(error); tracing::warn!(target: TRACING_TARGET, kind = %mapped.kind(), detail = mapped.detail(), "browser WebTransport primary stream receive failed"); return Err(mapped); } fn reset_frame_state(&mut self) { self.header = [0_u8; PRIMARY_FRAME_HEADER_SIZE]; self.header_read = 0; self.payload.clear(); self.payload_read = 0; } fn stop_after_failure(&mut self, code: u32) { self.inner.stop(code); self.terminal = true; } } impl Drop for WebTransportReceiver { fn drop(&mut self) { if !self.clean_closed && !self.terminal { self.inner.stop(STREAM_CANCELLED_ERROR_CODE); self.terminal = true; } } } impl game_realtime_transport_lib::RealtimeReceiver for WebTransportReceiver { type ReceiveFuture<'a> = std::pin::Pin< Box> + 'a>, > where Self: 'a; fn receive(&mut self) -> Self::ReceiveFuture<'_> { return Box::pin(async move { return self.receive_frame().await }); } } /// Send half of the browser primary reliable WebTransport stream. pub struct WebTransportSender { inner: web_transport_wasm::SendStream, _session: web_transport_wasm::Session, max_message_size: usize, send_timeout: std::time::Duration, terminal: bool, } impl WebTransportSender { /// Abruptly resets the reliable send direction with one WebTransport application error code. pub fn abort(&mut self, code: u32) -> Result<(), game_realtime_transport_lib::TransportError> { if self.terminal { return Err(transport_error(game_realtime_transport_lib::TransportErrorKind::Closed, "browser WebTransport sender is already terminal")); } self.inner.reset(code); self.terminal = true; tracing::debug!(target: TRACING_TARGET, code = code, "browser WebTransport primary send stream aborted"); return Ok(()); } } impl Drop for WebTransportSender { fn drop(&mut self) { if !self.terminal { self.inner.reset(STREAM_CANCELLED_ERROR_CODE); self.terminal = true; } } } impl game_realtime_transport_lib::RealtimeSender for WebTransportSender { type CloseFuture<'a> = std::pin::Pin> + 'a>> where Self: 'a; type SendFuture<'a> = std::pin::Pin> + 'a>> where Self: 'a; fn close(&mut self) -> Self::CloseFuture<'_> { return Box::pin(async move { if self.terminal { return Err(transport_error(game_realtime_transport_lib::TransportErrorKind::Closed, "browser WebTransport sender is already terminal")); } self.terminal = true; return match self.inner.finish() { Ok(()) => { tracing::debug!(target: TRACING_TARGET, "local browser WebTransport primary stream close initiated"); Ok(()) }, Err(error) => { let mapped = map_write_error(error); tracing::warn!(target: TRACING_TARGET, detail = mapped.detail(), "browser WebTransport primary stream close failed"); Err(mapped) }, }; }); } fn send(&mut self, message: game_realtime_transport_lib::TransportMessage) -> Self::SendFuture<'_> { return Box::pin(async move { if self.terminal { return Err(transport_error( game_realtime_transport_lib::TransportErrorKind::Aborted, "browser WebTransport sender is unavailable after close, abort, cancellation or terminal send failure", )); } let payload_len = message.len(); let frame_header = match frame_header(payload_len, self.max_message_size) { Ok(value) => value, Err(error) => return Err(error), }; let send_timeout = self.send_timeout; let timeout_millis = match browser_timeout_millis(send_timeout, "send_timeout") { Ok(value) => value, Err(error) => return Err(error), }; let mut guard = SendOperationGuard::new(&mut self.inner, &mut self.terminal); let operation = guard.write_frame(&frame_header, message.as_bytes()); let result = await_with_timeout(operation, timeout_millis).await; return match result { Some(Ok(())) => { guard.complete(); tracing::trace!(target: TRACING_TARGET, payload_len = payload_len, "framed browser WebTransport payload sent"); Ok(()) }, Some(Err(error)) => { guard.abort(SEND_FAILURE_ERROR_CODE); let mapped = map_write_error(error); tracing::warn!(target: TRACING_TARGET, payload_len = payload_len, kind = %mapped.kind(), detail = mapped.detail(), "browser WebTransport framed send failed"); Err(mapped) }, None => { guard.abort(SEND_TIMEOUT_ERROR_CODE); let mapped = timeout_error("browser WebTransport framed send", send_timeout); tracing::warn!(target: TRACING_TARGET, payload_len = payload_len, timeout_ms = send_timeout.as_millis(), "browser WebTransport framed send timed out under flow control/backpressure"); Err(mapped) }, }; }); } } struct SendOperationGuard<'a> { inner: &'a mut web_transport_wasm::SendStream, terminal: &'a mut bool, armed: bool, } impl<'a> SendOperationGuard<'a> { fn new(inner: &'a mut web_transport_wasm::SendStream, terminal: &'a mut bool) -> Self { return Self { inner, terminal, armed: true }; } fn abort(&mut self, code: u32) { if self.armed { self.inner.reset(code); *self.terminal = true; self.armed = false; } } fn complete(&mut self) { self.armed = false; } async fn write_frame(&mut self, frame_header: &[u8; PRIMARY_FRAME_HEADER_SIZE], payload: &[u8]) -> Result<(), web_transport_wasm::Error> { match self.inner.write(frame_header).await { Ok(()) => {}, Err(error) => return Err(error), } if !payload.is_empty() { match self.inner.write(payload).await { Ok(()) => {}, Err(error) => return Err(error), } } return Ok(()); } } impl Drop for SendOperationGuard<'_> { fn drop(&mut self) { self.abort(STREAM_CANCELLED_ERROR_CODE); } } /// Establishes one browser WebTransport session using an exact SHA-256 certificate pin. pub async fn connect(config: &WebTransportClientConfig) -> Result { if let Err(error) = config.transport.validate() { return Err(error); } if let Err(error) = validate_browser_deadlines(config.transport) { return Err(error); } let client = web_transport_wasm::ClientBuilder::new() .with_pooling(false) .with_unreliable(true) .with_server_certificate_hashes(vec![config.certificate_hash.as_bytes().to_vec()]); let timeout = config.transport.connect_timeout(); let timeout_millis = match browser_timeout_millis(timeout, "connect_timeout") { Ok(value) => value, Err(error) => return Err(error), }; let connected = await_with_timeout(client.connect(config.endpoint.clone()), timeout_millis).await; let session = match connected { Some(Ok(value)) => value, Some(Err(error)) => { let mapped = transport_error(game_realtime_transport_lib::TransportErrorKind::Connect, error.to_string()); tracing::warn!( target: TRACING_TARGET, endpoint = config.endpoint.as_str(), detail = mapped.detail(), "browser WebTransport client connection failed" ); return Err(mapped); }, None => { let mapped = timeout_error("browser WebTransport client connect", timeout); tracing::warn!(target: TRACING_TARGET, endpoint = config.endpoint.as_str(), timeout_ms = timeout.as_millis(), "browser WebTransport client connection timed out"); return Err(mapped); }, }; tracing::info!(target: TRACING_TARGET, endpoint = config.endpoint.as_str(), "browser WebTransport client connected"); return Ok(WebTransportSession::new(session, config.transport)); } async fn await_with_timeout(operation: F, timeout_millis: u32) -> Option where F: core::future::Future, { let operation = Box::pin(operation); let timeout = Box::pin(gloo_timers::future::TimeoutFuture::new(timeout_millis)); return match futures_util::future::select(operation, timeout).await { futures_util::future::Either::Left((output, _)) => Some(output), futures_util::future::Either::Right(((), _)) => None, }; } fn browser_timeout_millis(duration: std::time::Duration, name: &str) -> Result { return match u32::try_from(duration.as_millis()) { Ok(value) if value > 0 => Ok(value), Ok(_) => Err(invalid_configuration(format!("{name} must resolve to at least one browser timer millisecond"))), Err(_) => Err(invalid_configuration(format!("{name} exceeds the browser timer range of u32 milliseconds"))), }; } fn timeout_error(operation: &str, timeout: std::time::Duration) -> game_realtime_transport_lib::TransportError { return transport_error( game_realtime_transport_lib::TransportErrorKind::Timeout, format!("{operation} exceeded configured deadline of {} ms", timeout.as_millis()), ); } fn validate_browser_deadlines(config: crate::WebTransportConfig) -> Result<(), game_realtime_transport_lib::TransportError> { if let Err(error) = browser_timeout_millis(config.connect_timeout(), "connect_timeout") { return Err(error); } if let Err(error) = browser_timeout_millis(config.primary_stream_timeout(), "primary_stream_timeout") { return Err(error); } if let Err(error) = browser_timeout_millis(config.send_timeout(), "send_timeout") { return Err(error); } return Ok(()); } fn frame_header(payload_len: usize, max_message_size: usize) -> Result<[u8; PRIMARY_FRAME_HEADER_SIZE], game_realtime_transport_lib::TransportError> { if payload_len > max_message_size { return Err(message_too_large(payload_len, max_message_size)); } let payload_len = match u32::try_from(payload_len) { Ok(value) => value, Err(_) => return Err(message_too_large(payload_len, max_message_size)), }; return Ok(payload_len.to_be_bytes()); } fn invalid_configuration(detail: impl Into) -> game_realtime_transport_lib::TransportError { return transport_error(game_realtime_transport_lib::TransportErrorKind::InvalidConfiguration, detail); } fn map_read_error(error: web_transport_wasm::Error) -> game_realtime_transport_lib::TransportError { let kind = match &error { web_transport_wasm::Error::Closed => game_realtime_transport_lib::TransportErrorKind::Closed, web_transport_wasm::Error::Session { .. } => game_realtime_transport_lib::TransportErrorKind::Closed, web_transport_wasm::Error::Stream(_) => game_realtime_transport_lib::TransportErrorKind::Aborted, web_transport_wasm::Error::Unknown(_) => game_realtime_transport_lib::TransportErrorKind::Io, }; return transport_error(kind, error.to_string()); } fn map_write_error(error: web_transport_wasm::Error) -> game_realtime_transport_lib::TransportError { let kind = match &error { web_transport_wasm::Error::Closed => game_realtime_transport_lib::TransportErrorKind::Closed, web_transport_wasm::Error::Session { .. } => game_realtime_transport_lib::TransportErrorKind::Closed, web_transport_wasm::Error::Stream(_) => game_realtime_transport_lib::TransportErrorKind::Aborted, web_transport_wasm::Error::Unknown(_) => game_realtime_transport_lib::TransportErrorKind::Io, }; return transport_error(kind, error.to_string()); } fn datagram_too_large(payload_len: usize, max_datagram_size: usize) -> game_realtime_transport_lib::TransportError { return transport_error( game_realtime_transport_lib::TransportErrorKind::MessageTooLarge, format!("WebTransport datagram payload size {payload_len} exceeds current session maximum {max_datagram_size} bytes"), ); } fn message_too_large(payload_len: usize, max_message_size: usize) -> game_realtime_transport_lib::TransportError { return transport_error( game_realtime_transport_lib::TransportErrorKind::MessageTooLarge, format!("WebTransport framed payload length {payload_len} exceeds configured maximum {max_message_size}"), ); } fn protocol_error(detail: impl Into) -> 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) -> game_realtime_transport_lib::TransportError { return game_realtime_transport_lib::TransportError::new(kind, detail); }