0.3.5-alpha.4
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
# file: crates/common/game-realtime-webtransport-lib/Cargo.toml
|
||||
# version: 1
|
||||
# version: 2
|
||||
|
||||
[package]
|
||||
name = "game-realtime-webtransport-lib"
|
||||
@@ -13,6 +13,7 @@ publish.workspace = true
|
||||
[dependencies]
|
||||
game-realtime-transport-lib = { path = "../game-realtime-transport-lib" }
|
||||
rcgen = { workspace = true, features = ["ring"] }
|
||||
tokio = { workspace = true, features = ["time"] }
|
||||
tracing.workspace = true
|
||||
url.workspace = true
|
||||
web-transport-quinn = { workspace = true, features = ["ring"] }
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
<!-- file: crates/common/game-realtime-webtransport-lib/README.md -->
|
||||
<!-- version: 2 -->
|
||||
<!-- version: 3 -->
|
||||
|
||||
# game-realtime-webtransport-lib
|
||||
|
||||
@@ -18,9 +18,14 @@ La frontière native disponible couvre désormais :
|
||||
- établissement HTTP/3 WebTransport client/server natif ;
|
||||
- sélection d'un unique stream bidirectionnel fiable comme chemin realtime principal ;
|
||||
- framing privé `u32` big-endian + payload binaire ;
|
||||
- borne POC de 1 MiB vérifiée avant allocation côté réception et avant écriture côté émission ;
|
||||
- limite de message configurable, 1 MiB par défaut, vérifiée avant allocation côté réception et avant écriture côté émission ;
|
||||
- deadlines configurables pour la connexion, l'ouverture/accept du stream primaire et un envoi complet ;
|
||||
- adaptation `RealtimeConnection` / `RealtimeSender` / `RealtimeReceiver` ;
|
||||
- fermeture propre du chemin logique par FIN du stream primaire ;
|
||||
- FIN propre via `RealtimeSender::close()` ;
|
||||
- reset/STOP_SENDING backend-spécifiques via `WebTransportSender::abort(...)` et `WebTransportReceiver::abort(...)` ;
|
||||
- cancellation/drop terminale : un sender abandonné est reset plutôt que transformé implicitement en FIN ;
|
||||
- parseur de framing réception incrémental conservant son état si une future `receive()` est annulée ;
|
||||
- mapping stable des erreurs reset/close/session/protocole vers `TransportErrorKind` ;
|
||||
- tracing sous `games::realtime::webtransport`.
|
||||
|
||||
## TLS de développement
|
||||
@@ -48,23 +53,56 @@ one WebTransport session
|
||||
|
||||
`RealtimeConnection::split()` conserve la session WebTransport dans les deux moitiés afin que la session ne soit pas fermée au moment où l'objet connexion est consommé.
|
||||
|
||||
## Limite de frame POC
|
||||
## Limites et deadlines
|
||||
|
||||
La borne actuelle du framing fiable est volontairement interne à cette tranche : 1 MiB par message. Elle empêche une longueur `u32` hostile de provoquer une allocation arbitraire et rejette aussi l'émission hors limite avec `TransportErrorKind::MessageTooLarge`.
|
||||
`WebTransportConfig::default()` conserve la baseline de 1 MiB par message. La limite peut être réduite ou augmentée tant qu'elle reste strictement positive et représentable dans le champ de longueur `u32` du framing.
|
||||
|
||||
Cette valeur n'est pas encore une configuration produit. La tranche de robustesse suivante doit décider si la limite devient configurable avec les deadlines, la backpressure, les resets, l'abort/cancellation et le mapping d'erreurs détaillé.
|
||||
Les deadlines configurables couvrent :
|
||||
|
||||
- connexion client et réponse finale à une requête WebTransport déjà surfacée côté serveur ;
|
||||
- ouverture ou accept du stream bidirectionnel principal ;
|
||||
- écriture complète header + payload d'une frame.
|
||||
|
||||
L'attente d'un nouveau pair sur le listener reste volontairement non bornée : un serveur inactif ne doit pas produire périodiquement une erreur uniquement parce qu'aucun client ne se présente.
|
||||
|
||||
QUIC applique sa propre flow-control. Le backend n'ajoute pas une seconde file applicative : si un envoi reste bloqué par flow-control/réseau au-delà de `send_timeout`, l'opération retourne `TransportErrorKind::Timeout` et le stream est reset afin qu'une frame partiellement transmise ne puisse pas être suivie d'une nouvelle frame invalide.
|
||||
|
||||
## Lifecycle, abort et cancellation
|
||||
|
||||
`RealtimeSender::close()` reste la fermeture propre de la direction d'émission et produit un FIN. À l'inverse :
|
||||
|
||||
- `WebTransportSender::abort(code)` envoie un `RESET_STREAM` WebTransport ;
|
||||
- `WebTransportReceiver::abort(code)` envoie un `STOP_SENDING` WebTransport ;
|
||||
- dropper un `WebTransportSender` encore actif provoque un reset explicite ;
|
||||
- dropper un `WebTransportReceiver` encore actif provoque un stop explicite ;
|
||||
- annuler une future `send()` en cours provoque également un reset via une garde de cancellation.
|
||||
|
||||
Une erreur terminale de lecture/écriture rend la moitié concernée indisponible pour une réutilisation silencieuse.
|
||||
|
||||
La réception n'utilise plus une lecture exacte monolithique. Le header et le payload sont lus progressivement avec l'API de lecture cancel-safe de Quinn ; `header_read`/`payload_read` restent dans le receiver. Une future `receive()` annulée peut donc être relancée sans perdre les octets déjà consommés ni décaler le framing.
|
||||
|
||||
## Mapping d'erreurs
|
||||
|
||||
Le backend distingue notamment :
|
||||
|
||||
- payload hors limite -> `MessageTooLarge` ;
|
||||
- reset/STOP_SENDING valide -> `Aborted` ;
|
||||
- stream déjà fermé -> `Closed` ;
|
||||
- fermeture de session WebTransport explicite -> `Closed` ;
|
||||
- erreur de session/connexion non classée comme fermeture propre -> `Io` ;
|
||||
- reset/stop invalide ou framing tronqué -> `Protocol` ;
|
||||
- deadline dépassée -> `Timeout`.
|
||||
|
||||
Une longueur entrante hors limite ou un framing tronqué provoque aussi l'arrêt de la direction de réception afin d'éviter de poursuivre sur un flux désynchronisé.
|
||||
|
||||
## Frontières actuelles
|
||||
|
||||
La crate ne possède toujours pas :
|
||||
|
||||
- d'API datagram transport-neutral ;
|
||||
- de deadlines applicatives WebTransport ;
|
||||
- de politique de backpressure explicite ;
|
||||
- de reset/abort/cancellation produit ;
|
||||
- de mapping fin de toutes les erreurs Quinn/WebTransport ;
|
||||
- de chemin navigateur/WASM ;
|
||||
- de smoke executable public WebTransport ;
|
||||
- de fallback WebSocket ;
|
||||
- de benchmark WebSocket/WebTransport.
|
||||
|
||||
Ces responsabilités restent réservées aux tranches suivantes du plan `0.3.5`.
|
||||
|
||||
@@ -1,9 +1,25 @@
|
||||
<!-- file: crates/common/game-realtime-webtransport-lib/USAGE.md -->
|
||||
<!-- version: 1 -->
|
||||
<!-- version: 2 -->
|
||||
|
||||
# Utilisation de game-realtime-webtransport-lib
|
||||
|
||||
Ce guide décrit le chemin natif fiable actuellement exposé par `game-realtime-webtransport-lib`. Il ne décrit ni gameplay, ni protocole wire métier, ni datagrams.
|
||||
Ce guide décrit le chemin natif fiable exposé par `game-realtime-webtransport-lib`. Il ne décrit ni gameplay, ni protocole wire métier, ni datagrams.
|
||||
|
||||
## Configuration transport
|
||||
|
||||
`WebTransportConfig` porte les limites et deadlines du chemin fiable. La configuration par défaut garde une limite de 1 MiB par message et des deadlines bornées pour connexion, stream primaire et send.
|
||||
|
||||
Exemple de configuration plus stricte :
|
||||
|
||||
```rust
|
||||
let transport = game_realtime_webtransport_lib::WebTransportConfig::default()
|
||||
.with_max_message_size(256 * 1024)
|
||||
.with_connect_timeout(std::time::Duration::from_secs(5))
|
||||
.with_primary_stream_timeout(std::time::Duration::from_secs(2))
|
||||
.with_send_timeout(std::time::Duration::from_secs(2));
|
||||
```
|
||||
|
||||
La validation effective se fait lors du bind serveur ou de la connexion client. Une limite nulle, une limite non représentable en `u32` ou une deadline nulle est rejetée comme `InvalidConfiguration`.
|
||||
|
||||
## Serveur natif
|
||||
|
||||
@@ -17,7 +33,8 @@ let identity = match game_realtime_webtransport_lib::WebTransportServerIdentity:
|
||||
let config = game_realtime_webtransport_lib::WebTransportServerConfig::new(
|
||||
std::net::SocketAddr::from(([127, 0, 0, 1], 4433)),
|
||||
identity,
|
||||
);
|
||||
)
|
||||
.with_transport_config(transport);
|
||||
let mut listener = match game_realtime_webtransport_lib::WebTransportListener::bind(config) {
|
||||
Ok(value) => value,
|
||||
Err(error) => return Err(error),
|
||||
@@ -32,6 +49,8 @@ let connection = match session.accept_primary_connection().await {
|
||||
};
|
||||
```
|
||||
|
||||
L'attente du prochain client dans `listener.accept()` n'a pas de timeout périodique. Une fois une requête WebTransport surfacée, sa réponse finale utilise la deadline de connexion configurée.
|
||||
|
||||
`open_primary_connection()` écrit l'en-tête WebTransport requis pour identifier le stream avant de retourner. Le serveur peut donc attendre `accept_primary_connection()` puis commencer les échanges applicatifs ; aucune frame artificielle n'est nécessaire pour rendre le stream visible.
|
||||
|
||||
## Client natif
|
||||
@@ -43,7 +62,7 @@ let config = match game_realtime_webtransport_lib::WebTransportClientConfig::new
|
||||
"https://127.0.0.1:4433/game",
|
||||
certificate_hash,
|
||||
) {
|
||||
Ok(value) => value,
|
||||
Ok(value) => value.with_transport_config(transport),
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
let session = match game_realtime_webtransport_lib::connect(&config).await {
|
||||
@@ -60,7 +79,7 @@ Le pinning est obligatoire dans cette API native ; il n'existe pas de variante q
|
||||
|
||||
## Contrat realtime
|
||||
|
||||
Une fois le stream primaire sélectionné, utiliser uniquement les traits de `game-realtime-transport-lib` :
|
||||
Une fois le stream primaire sélectionné, utiliser les traits de `game-realtime-transport-lib` pour le chemin fiable normal :
|
||||
|
||||
```rust
|
||||
let (mut sender, mut receiver) = game_realtime_transport_lib::RealtimeConnection::split(connection);
|
||||
@@ -82,8 +101,30 @@ if let Err(error) = game_realtime_transport_lib::RealtimeSender::close(&mut send
|
||||
|
||||
Le backend encode chaque `TransportMessage` sous la forme `u32` big-endian + payload. Le consommateur ne doit pas reproduire ce framing lui-même.
|
||||
|
||||
## Fermeture actuelle
|
||||
La flow-control QUIC est respectée naturellement par l'écriture asynchrone. Un send qui dépasse sa deadline est considéré terminal : le stream est reset et le même sender ne doit pas être réutilisé.
|
||||
|
||||
## Fermeture et abort
|
||||
|
||||
`RealtimeSender::close()` termine proprement la direction d'émission du stream primaire. Le pair observe ensuite `TransportReceive::Closed` lorsqu'il atteint le FIN après les messages déjà écrits.
|
||||
|
||||
La fermeture de session complète, les resets, les aborts, les timeouts et les cas d'annulation appartiennent à la tranche de robustesse suivante et ne doivent pas être simulés par le consommateur.
|
||||
Pour abandonner explicitement une direction WebTransport :
|
||||
|
||||
```rust
|
||||
if let Err(error) = sender.abort(42) {
|
||||
return Err(error);
|
||||
}
|
||||
```
|
||||
|
||||
ou côté réception :
|
||||
|
||||
```rust
|
||||
if let Err(error) = receiver.abort(43) {
|
||||
return Err(error);
|
||||
}
|
||||
```
|
||||
|
||||
Ces deux méthodes sont backend-spécifiques : elles ne sont pas ajoutées au contrat commun car WebSocket n'expose pas la même primitive QUIC de reset/stop.
|
||||
|
||||
Dropper un sender actif ou annuler une future `send()` en cours provoque un reset explicite. Dropper un receiver actif provoque un stop explicite. Cela évite qu'une cancellation d'écriture partielle soit interprétée comme une fermeture propre ou qu'une frame suivante reprenne au mauvais offset.
|
||||
|
||||
Une future `receive()` peut en revanche être annulée puis relancée : le backend conserve l'état partiel du header/payload et reprend le framing à l'octet correct.
|
||||
|
||||
109
crates/common/game-realtime-webtransport-lib/src/config.rs
Normal file
109
crates/common/game-realtime-webtransport-lib/src/config.rs
Normal file
@@ -0,0 +1,109 @@
|
||||
// file: crates/common/game-realtime-webtransport-lib/src/config.rs
|
||||
// version: 1
|
||||
|
||||
const DEFAULT_CONNECT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
|
||||
const DEFAULT_MAX_MESSAGE_SIZE: usize = 1024 * 1024;
|
||||
const DEFAULT_PRIMARY_STREAM_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
|
||||
const DEFAULT_SEND_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
|
||||
|
||||
/// Product-facing limits and operation deadlines for the native WebTransport reliable path.
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
pub struct WebTransportConfig {
|
||||
connect_timeout: std::time::Duration,
|
||||
max_message_size: usize,
|
||||
primary_stream_timeout: std::time::Duration,
|
||||
send_timeout: std::time::Duration,
|
||||
}
|
||||
|
||||
impl Default for WebTransportConfig {
|
||||
fn default() -> Self {
|
||||
return Self {
|
||||
connect_timeout: DEFAULT_CONNECT_TIMEOUT,
|
||||
max_message_size: DEFAULT_MAX_MESSAGE_SIZE,
|
||||
primary_stream_timeout: DEFAULT_PRIMARY_STREAM_TIMEOUT,
|
||||
send_timeout: DEFAULT_SEND_TIMEOUT,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
impl WebTransportConfig {
|
||||
/// Returns a copy with a different client/session handshake deadline.
|
||||
#[must_use]
|
||||
pub fn with_connect_timeout(mut self, value: std::time::Duration) -> Self {
|
||||
self.connect_timeout = value;
|
||||
return self;
|
||||
}
|
||||
|
||||
/// Returns a copy with a different maximum framed binary message size.
|
||||
#[must_use]
|
||||
pub fn with_max_message_size(mut self, value: usize) -> Self {
|
||||
self.max_message_size = value;
|
||||
return self;
|
||||
}
|
||||
|
||||
/// Returns a copy with a different deadline for opening or accepting the primary stream.
|
||||
#[must_use]
|
||||
pub fn with_primary_stream_timeout(mut self, value: std::time::Duration) -> Self {
|
||||
self.primary_stream_timeout = value;
|
||||
return self;
|
||||
}
|
||||
|
||||
/// Returns a copy with a different deadline for one complete framed send operation.
|
||||
#[must_use]
|
||||
pub fn with_send_timeout(mut self, value: std::time::Duration) -> Self {
|
||||
self.send_timeout = value;
|
||||
return self;
|
||||
}
|
||||
|
||||
/// Returns the configured client/session handshake deadline.
|
||||
#[must_use]
|
||||
pub fn connect_timeout(&self) -> std::time::Duration {
|
||||
return self.connect_timeout;
|
||||
}
|
||||
|
||||
/// Returns the maximum framed binary message size.
|
||||
#[must_use]
|
||||
pub fn max_message_size(&self) -> usize {
|
||||
return self.max_message_size;
|
||||
}
|
||||
|
||||
/// Returns the deadline for opening or accepting the primary stream.
|
||||
#[must_use]
|
||||
pub fn primary_stream_timeout(&self) -> std::time::Duration {
|
||||
return self.primary_stream_timeout;
|
||||
}
|
||||
|
||||
/// Returns the deadline for one complete framed send operation.
|
||||
#[must_use]
|
||||
pub fn send_timeout(&self) -> std::time::Duration {
|
||||
return self.send_timeout;
|
||||
}
|
||||
|
||||
/// Validates all limits and deadlines required by the reliable WebTransport path.
|
||||
pub fn validate(&self) -> Result<(), game_realtime_transport_lib::TransportError> {
|
||||
if self.max_message_size == 0 {
|
||||
return Err(invalid_configuration("max_message_size must be greater than zero"));
|
||||
}
|
||||
if u32::try_from(self.max_message_size).is_err() {
|
||||
return Err(invalid_configuration("max_message_size must fit the u32 framing length field"));
|
||||
}
|
||||
if self.connect_timeout.is_zero() {
|
||||
return Err(invalid_configuration("connect_timeout must be greater than zero"));
|
||||
}
|
||||
if self.primary_stream_timeout.is_zero() {
|
||||
return Err(invalid_configuration("primary_stream_timeout must be greater than zero"));
|
||||
}
|
||||
if self.send_timeout.is_zero() {
|
||||
return Err(invalid_configuration("send_timeout must be greater than zero"));
|
||||
}
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
|
||||
fn invalid_configuration(detail: &str) -> game_realtime_transport_lib::TransportError {
|
||||
return game_realtime_transport_lib::TransportError::new(game_realtime_transport_lib::TransportErrorKind::InvalidConfiguration, detail);
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "../unit_tests/config.rs"]
|
||||
mod tests;
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/common/game-realtime-webtransport-lib/src/lib.rs
|
||||
// version: 2
|
||||
// version: 3
|
||||
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
@@ -7,8 +7,11 @@
|
||||
|
||||
//! Native WebTransport/QUIC backend candidate for the transport-neutral realtime contract.
|
||||
|
||||
mod config;
|
||||
mod webtransport;
|
||||
|
||||
/// Re-export of product-facing limits and operation deadlines for the reliable WebTransport path.
|
||||
pub use self::config::WebTransportConfig;
|
||||
/// Re-export of the pinned SHA-256 certificate fingerprint used by the native client.
|
||||
pub use self::webtransport::WebTransportCertificateHash;
|
||||
/// Re-export of native WebTransport client configuration.
|
||||
|
||||
@@ -1,11 +1,10 @@
|
||||
// file: crates/common/game-realtime-webtransport-lib/src/webtransport.rs
|
||||
// version: 2
|
||||
// version: 3
|
||||
|
||||
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.
|
||||
@@ -91,11 +90,12 @@ impl WebTransportServerIdentity {
|
||||
}
|
||||
}
|
||||
|
||||
/// Native WebTransport client endpoint and pinned server-certificate fingerprint.
|
||||
/// Native WebTransport client 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 {
|
||||
@@ -111,7 +111,14 @@ impl WebTransportClientConfig {
|
||||
if parsed.host().is_none() {
|
||||
return Err(invalid_configuration("WebTransport endpoint must contain a host"));
|
||||
}
|
||||
return Ok(Self { endpoint: parsed, certificate_hash });
|
||||
return Ok(Self { endpoint: parsed, certificate_hash, transport: crate::WebTransportConfig::default() });
|
||||
}
|
||||
|
||||
/// Returns a copy with explicit reliable-path limits and deadlines.
|
||||
#[must_use]
|
||||
pub fn with_transport_config(mut self, transport: crate::WebTransportConfig) -> Self {
|
||||
self.transport = transport;
|
||||
return self;
|
||||
}
|
||||
|
||||
/// Returns the validated WebTransport endpoint URL.
|
||||
@@ -125,19 +132,33 @@ impl WebTransportClientConfig {
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
||||
/// Native WebTransport server bind address and TLS identity.
|
||||
/// Native WebTransport server bind address, TLS identity and reliable-path configuration.
|
||||
pub struct WebTransportServerConfig {
|
||||
bind_address: std::net::SocketAddr,
|
||||
identity: WebTransportServerIdentity,
|
||||
transport: crate::WebTransportConfig,
|
||||
}
|
||||
|
||||
impl WebTransportServerConfig {
|
||||
/// Creates native server configuration for the requested bind address and TLS identity.
|
||||
#[must_use]
|
||||
pub fn new(bind_address: std::net::SocketAddr, identity: WebTransportServerIdentity) -> Self {
|
||||
return Self { bind_address, identity };
|
||||
return Self { bind_address, identity, transport: crate::WebTransportConfig::default() };
|
||||
}
|
||||
|
||||
/// Returns a copy with explicit reliable-path limits and deadlines.
|
||||
#[must_use]
|
||||
pub fn with_transport_config(mut self, transport: crate::WebTransportConfig) -> Self {
|
||||
self.transport = transport;
|
||||
return self;
|
||||
}
|
||||
|
||||
/// Returns the requested UDP bind address.
|
||||
@@ -151,16 +172,23 @@ impl WebTransportServerConfig {
|
||||
pub fn certificate_hash(&self) -> &WebTransportCertificateHash {
|
||||
return self.identity.certificate_hash();
|
||||
}
|
||||
|
||||
/// Returns the reliable-path limits and deadlines.
|
||||
#[must_use]
|
||||
pub fn transport_config(&self) -> crate::WebTransportConfig {
|
||||
return self.transport;
|
||||
}
|
||||
}
|
||||
|
||||
/// Established native WebTransport session before or while the single primary application stream is selected.
|
||||
pub struct WebTransportSession {
|
||||
inner: web_transport_quinn::Session,
|
||||
transport: crate::WebTransportConfig,
|
||||
}
|
||||
|
||||
impl WebTransportSession {
|
||||
fn new(inner: web_transport_quinn::Session) -> Self {
|
||||
return Self { inner };
|
||||
fn new(inner: web_transport_quinn::Session, transport: crate::WebTransportConfig) -> Self {
|
||||
return Self { inner, transport };
|
||||
}
|
||||
|
||||
/// Accepts the peer-created primary bidirectional stream and adapts it to the transport-neutral realtime contract.
|
||||
@@ -168,30 +196,44 @@ impl WebTransportSession {
|
||||
/// 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 timeout = self.transport.primary_stream_timeout();
|
||||
let accepted = tokio::time::timeout(timeout, self.inner.accept_bi()).await;
|
||||
let (sender, receiver) = match accepted {
|
||||
Ok(Ok(value)) => value,
|
||||
Ok(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);
|
||||
},
|
||||
Err(_) => {
|
||||
let mapped = timeout_error("WebTransport primary bidirectional stream accept", timeout);
|
||||
tracing::warn!(target: TRACING_TARGET, timeout_ms = duration_millis(timeout), "WebTransport primary bidirectional stream accept timed out");
|
||||
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));
|
||||
return Ok(WebTransportConnection::new(self.inner, sender, receiver, self.transport));
|
||||
}
|
||||
|
||||
/// 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 timeout = self.transport.primary_stream_timeout();
|
||||
let opened = tokio::time::timeout(timeout, self.inner.open_bi()).await;
|
||||
let (sender, receiver) = match opened {
|
||||
Ok(Ok(value)) => value,
|
||||
Ok(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);
|
||||
},
|
||||
Err(_) => {
|
||||
let mapped = timeout_error("WebTransport primary bidirectional stream open", timeout);
|
||||
tracing::warn!(target: TRACING_TARGET, timeout_ms = duration_millis(timeout), "WebTransport primary bidirectional stream open timed out");
|
||||
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));
|
||||
return Ok(WebTransportConnection::new(self.inner, sender, receiver, self.transport));
|
||||
}
|
||||
|
||||
/// Returns the remote UDP socket address backing the established QUIC connection.
|
||||
@@ -215,11 +257,17 @@ pub struct WebTransportConnection {
|
||||
receiver: web_transport_quinn::RecvStream,
|
||||
sender: web_transport_quinn::SendStream,
|
||||
session: web_transport_quinn::Session,
|
||||
transport: crate::WebTransportConfig,
|
||||
}
|
||||
|
||||
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 };
|
||||
fn new(
|
||||
session: web_transport_quinn::Session,
|
||||
sender: web_transport_quinn::SendStream,
|
||||
receiver: web_transport_quinn::RecvStream,
|
||||
transport: crate::WebTransportConfig,
|
||||
) -> Self {
|
||||
return Self { receiver, sender, session, transport };
|
||||
}
|
||||
}
|
||||
|
||||
@@ -230,8 +278,24 @@ impl game_realtime_transport_lib::RealtimeConnection for WebTransportConnection
|
||||
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 },
|
||||
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,
|
||||
},
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -240,6 +304,137 @@ impl game_realtime_transport_lib::RealtimeConnection for WebTransportConnection
|
||||
pub struct WebTransportReceiver {
|
||||
inner: web_transport_quinn::RecvStream,
|
||||
_session: web_transport_quinn::Session,
|
||||
max_message_size: usize,
|
||||
header: [u8; PRIMARY_FRAME_HEADER_SIZE],
|
||||
header_read: usize,
|
||||
payload: Vec<u8>,
|
||||
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, "WebTransport receiver is already terminal"));
|
||||
}
|
||||
return match self.inner.stop(code) {
|
||||
Ok(()) => {
|
||||
self.terminal = true;
|
||||
tracing::debug!(target: TRACING_TARGET, code = code, "WebTransport primary receive stream aborted");
|
||||
Ok(())
|
||||
},
|
||||
Err(error) => {
|
||||
self.terminal = true;
|
||||
let mapped = transport_error(game_realtime_transport_lib::TransportErrorKind::Closed, error.to_string());
|
||||
tracing::warn!(target: TRACING_TARGET, code = code, detail = mapped.detail(), "WebTransport primary receive stream abort failed");
|
||||
Err(mapped)
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
async fn receive_frame(&mut self) -> Result<game_realtime_transport_lib::TransportReceive, game_realtime_transport_lib::TransportError> {
|
||||
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,
|
||||
"WebTransport receiver is unavailable after a terminal stream failure or abort",
|
||||
));
|
||||
}
|
||||
loop {
|
||||
if self.header_read < PRIMARY_FRAME_HEADER_SIZE {
|
||||
let read = self.inner.read(&mut self.header[self.header_read..]).await;
|
||||
match read {
|
||||
Ok(Some(0)) => return self.fail_protocol("WebTransport primary stream returned an empty read in the middle of a frame header"),
|
||||
Ok(Some(value)) => {
|
||||
self.header_read += value;
|
||||
continue;
|
||||
},
|
||||
Ok(None) => {
|
||||
if self.header_read == 0 {
|
||||
self.clean_closed = true;
|
||||
tracing::debug!(target: TRACING_TARGET, "remote WebTransport primary stream closed cleanly");
|
||||
return Ok(game_realtime_transport_lib::TransportReceive::Closed);
|
||||
}
|
||||
return self.fail_protocol("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 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 read = self.inner.read(&mut self.payload[self.payload_read..]).await;
|
||||
match read {
|
||||
Ok(Some(0)) => return self.fail_protocol("WebTransport primary stream returned an empty read in the middle of a frame payload"),
|
||||
Ok(Some(value)) => {
|
||||
self.payload_read += value;
|
||||
if self.payload_read < self.payload.len() {
|
||||
continue;
|
||||
}
|
||||
},
|
||||
Ok(None) => return self.fail_protocol("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 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<game_realtime_transport_lib::TransportReceive, game_realtime_transport_lib::TransportError> {
|
||||
let error = protocol_error(detail);
|
||||
self.stop_after_failure(FRAME_PROTOCOL_ERROR_CODE);
|
||||
return Err(error);
|
||||
}
|
||||
|
||||
fn fail_read(
|
||||
&mut self,
|
||||
error: web_transport_quinn::ReadError,
|
||||
) -> Result<game_realtime_transport_lib::TransportReceive, game_realtime_transport_lib::TransportError> {
|
||||
self.terminal = true;
|
||||
let mapped = map_read_error(error);
|
||||
tracing::warn!(target: TRACING_TARGET, kind = %mapped.kind(), detail = mapped.detail(), "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) {
|
||||
let _ = self.inner.stop(code);
|
||||
self.terminal = true;
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for WebTransportReceiver {
|
||||
fn drop(&mut self) {
|
||||
if !self.clean_closed && !self.terminal {
|
||||
let _ = self.inner.stop(STREAM_CANCELLED_ERROR_CODE);
|
||||
self.terminal = true;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl game_realtime_transport_lib::RealtimeReceiver for WebTransportReceiver {
|
||||
@@ -251,22 +446,7 @@ impl game_realtime_transport_lib::RealtimeReceiver for WebTransportReceiver {
|
||||
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)));
|
||||
});
|
||||
return Box::pin(async move { self.receive_frame().await });
|
||||
}
|
||||
}
|
||||
|
||||
@@ -274,6 +454,39 @@ impl game_realtime_transport_lib::RealtimeReceiver for WebTransportReceiver {
|
||||
pub struct WebTransportSender {
|
||||
inner: web_transport_quinn::SendStream,
|
||||
_session: web_transport_quinn::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, "WebTransport sender is already terminal"));
|
||||
}
|
||||
self.terminal = true;
|
||||
return match self.inner.reset(code) {
|
||||
Ok(()) => {
|
||||
tracing::debug!(target: TRACING_TARGET, code = code, "WebTransport primary send stream aborted");
|
||||
Ok(())
|
||||
},
|
||||
Err(error) => {
|
||||
let mapped = transport_error(game_realtime_transport_lib::TransportErrorKind::Closed, error.to_string());
|
||||
tracing::warn!(target: TRACING_TARGET, code = code, detail = mapped.detail(), "WebTransport primary send stream abort failed");
|
||||
Err(mapped)
|
||||
},
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for WebTransportSender {
|
||||
fn drop(&mut self) {
|
||||
if !self.terminal {
|
||||
let _ = self.inner.reset(STREAM_CANCELLED_ERROR_CODE);
|
||||
self.terminal = true;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl game_realtime_transport_lib::RealtimeSender for WebTransportSender {
|
||||
@@ -288,6 +501,10 @@ impl game_realtime_transport_lib::RealtimeSender for WebTransportSender {
|
||||
|
||||
fn close(&mut self) -> Self::CloseFuture<'_> {
|
||||
return Box::pin(async move {
|
||||
if self.terminal {
|
||||
return Err(transport_error(game_realtime_transport_lib::TransportErrorKind::Closed, "WebTransport sender is already terminal"));
|
||||
}
|
||||
self.terminal = true;
|
||||
return match self.inner.finish() {
|
||||
Ok(()) => {
|
||||
tracing::debug!(target: TRACING_TARGET, "local WebTransport primary stream close initiated");
|
||||
@@ -304,36 +521,102 @@ impl game_realtime_transport_lib::RealtimeSender for WebTransportSender {
|
||||
|
||||
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,
|
||||
"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) {
|
||||
let frame_header = match frame_header(payload_len, self.max_message_size) {
|
||||
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(());
|
||||
let send_timeout = self.send_timeout;
|
||||
let mut guard = SendOperationGuard::new(&mut self.inner, &mut self.terminal);
|
||||
let operation = guard.write_frame(&frame_header, message.as_bytes());
|
||||
let result = tokio::time::timeout(send_timeout, operation).await;
|
||||
return match result {
|
||||
Ok(Ok(())) => {
|
||||
guard.complete();
|
||||
tracing::trace!(target: TRACING_TARGET, payload_len = payload_len, "framed WebTransport payload sent");
|
||||
Ok(())
|
||||
},
|
||||
Ok(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(), "WebTransport framed send failed");
|
||||
Err(mapped)
|
||||
},
|
||||
Err(_) => {
|
||||
guard.abort(SEND_TIMEOUT_ERROR_CODE);
|
||||
let mapped = timeout_error("WebTransport framed send", send_timeout);
|
||||
tracing::warn!(target: TRACING_TARGET, payload_len = payload_len, timeout_ms = duration_millis(send_timeout), "WebTransport framed send timed out under flow control/backpressure");
|
||||
Err(mapped)
|
||||
},
|
||||
};
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
struct SendOperationGuard<'a> {
|
||||
stream: &'a mut web_transport_quinn::SendStream,
|
||||
terminal: &'a mut bool,
|
||||
armed: bool,
|
||||
}
|
||||
|
||||
impl<'a> SendOperationGuard<'a> {
|
||||
fn new(stream: &'a mut web_transport_quinn::SendStream, terminal: &'a mut bool) -> Self {
|
||||
return Self { stream, terminal, armed: true };
|
||||
}
|
||||
|
||||
async fn write_frame(&mut self, header: &[u8; PRIMARY_FRAME_HEADER_SIZE], payload: &[u8]) -> Result<(), web_transport_quinn::WriteError> {
|
||||
if let Err(error) = self.stream.write_all(header).await {
|
||||
return Err(error);
|
||||
}
|
||||
if let Err(error) = self.stream.write_all(payload).await {
|
||||
return Err(error);
|
||||
}
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
fn abort(&mut self, code: u32) {
|
||||
if self.armed {
|
||||
let _ = self.stream.reset(code);
|
||||
*self.terminal = true;
|
||||
self.armed = false;
|
||||
}
|
||||
}
|
||||
|
||||
fn complete(&mut self) {
|
||||
self.armed = false;
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for SendOperationGuard<'_> {
|
||||
fn drop(&mut self) {
|
||||
if self.armed {
|
||||
let _ = self.stream.reset(STREAM_CANCELLED_ERROR_CODE);
|
||||
*self.terminal = true;
|
||||
self.armed = false;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Bound native WebTransport server endpoint that accepts HTTP/3 WebTransport sessions.
|
||||
pub struct WebTransportListener {
|
||||
server: web_transport_quinn::Server,
|
||||
local_addr: std::net::SocketAddr,
|
||||
transport: crate::WebTransportConfig,
|
||||
}
|
||||
|
||||
impl WebTransportListener {
|
||||
/// Binds a native WebTransport server using TLS 1.3 and the configured certificate identity.
|
||||
pub fn bind(config: WebTransportServerConfig) -> Result<Self, game_realtime_transport_lib::TransportError> {
|
||||
if let Err(error) = config.transport.validate() {
|
||||
return Err(error);
|
||||
}
|
||||
let transport = config.transport;
|
||||
let certificate = web_transport_quinn::quinn::rustls::pki_types::CertificateDer::from(config.identity.certificate_der);
|
||||
let private_key = web_transport_quinn::quinn::rustls::pki_types::PrivatePkcs8KeyDer::from(config.identity.private_key_pkcs8_der);
|
||||
let private_key = web_transport_quinn::quinn::rustls::pki_types::PrivateKeyDer::Pkcs8(private_key);
|
||||
@@ -354,7 +637,7 @@ impl WebTransportListener {
|
||||
},
|
||||
};
|
||||
tracing::info!(target: TRACING_TARGET, address = %local_addr, "WebTransport listener bound");
|
||||
return Ok(Self { server, local_addr });
|
||||
return Ok(Self { server, local_addr, transport });
|
||||
}
|
||||
|
||||
/// Returns the concrete UDP socket address, including an ephemeral port selected by the OS.
|
||||
@@ -364,6 +647,9 @@ impl WebTransportListener {
|
||||
}
|
||||
|
||||
/// Accepts one native WebTransport CONNECT request and returns the established session.
|
||||
///
|
||||
/// Waiting for the next peer remains intentionally unbounded. Once a WebTransport request is surfaced, the final server
|
||||
/// response is bounded by the configured connection deadline.
|
||||
pub async fn accept(&mut self) -> Result<WebTransportSession, game_realtime_transport_lib::TransportError> {
|
||||
let request = match self.server.accept().await {
|
||||
Some(value) => value,
|
||||
@@ -374,37 +660,60 @@ impl WebTransportListener {
|
||||
},
|
||||
};
|
||||
let peer = request.conn().remote_address();
|
||||
let session = match request.ok().await {
|
||||
Ok(value) => value,
|
||||
Err(error) => {
|
||||
let timeout = self.transport.connect_timeout();
|
||||
let accepted = tokio::time::timeout(timeout, request.ok()).await;
|
||||
let session = match accepted {
|
||||
Ok(Ok(value)) => value,
|
||||
Ok(Err(error)) => {
|
||||
let mapped = transport_error(game_realtime_transport_lib::TransportErrorKind::Accept, error.to_string());
|
||||
tracing::warn!(target: TRACING_TARGET, peer = %peer, detail = mapped.detail(), "WebTransport server handshake failed");
|
||||
return Err(mapped);
|
||||
},
|
||||
Err(_) => {
|
||||
let mapped = timeout_error("WebTransport server handshake response", timeout);
|
||||
tracing::warn!(target: TRACING_TARGET, peer = %peer, timeout_ms = duration_millis(timeout), "WebTransport server handshake response timed out");
|
||||
return Err(mapped);
|
||||
},
|
||||
};
|
||||
tracing::info!(target: TRACING_TARGET, peer = %peer, "WebTransport peer accepted");
|
||||
return Ok(WebTransportSession::new(session));
|
||||
return Ok(WebTransportSession::new(session, self.transport));
|
||||
}
|
||||
}
|
||||
|
||||
/// Establishes one native WebTransport session using an exact SHA-256 certificate pin.
|
||||
pub async fn connect(config: &WebTransportClientConfig) -> Result<WebTransportSession, game_realtime_transport_lib::TransportError> {
|
||||
if let Err(error) = config.transport.validate() {
|
||||
return Err(error);
|
||||
}
|
||||
let client = match web_transport_quinn::ClientBuilder::new().with_server_certificate_hashes(vec![config.certificate_hash.as_bytes().to_vec()]) {
|
||||
Ok(value) => value,
|
||||
Err(error) => return Err(invalid_configuration(error.to_string())),
|
||||
};
|
||||
let session = match client.connect(config.endpoint.clone()).await {
|
||||
Ok(value) => value,
|
||||
Err(error) => {
|
||||
let timeout = config.transport.connect_timeout();
|
||||
let connected = tokio::time::timeout(timeout, client.connect(config.endpoint.clone())).await;
|
||||
let session = match connected {
|
||||
Ok(Ok(value)) => value,
|
||||
Ok(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(), "WebTransport client connection failed");
|
||||
return Err(mapped);
|
||||
},
|
||||
Err(_) => {
|
||||
let mapped = timeout_error("WebTransport client connect", timeout);
|
||||
tracing::warn!(target: TRACING_TARGET, endpoint = config.endpoint.as_str(), timeout_ms = duration_millis(timeout), "WebTransport client connection timed out");
|
||||
return Err(mapped);
|
||||
},
|
||||
};
|
||||
tracing::info!(target: TRACING_TARGET, endpoint = config.endpoint.as_str(), peer = %session.remote_address(), "WebTransport client connected");
|
||||
return Ok(WebTransportSession::new(session));
|
||||
return Ok(WebTransportSession::new(session, config.transport));
|
||||
}
|
||||
|
||||
const FRAME_PROTOCOL_ERROR_CODE: u32 = 0x10;
|
||||
const FRAME_TOO_LARGE_ERROR_CODE: u32 = 0x11;
|
||||
const SEND_FAILURE_ERROR_CODE: u32 = 0x12;
|
||||
const SEND_TIMEOUT_ERROR_CODE: u32 = 0x13;
|
||||
const STREAM_CANCELLED_ERROR_CODE: u32 = 0x14;
|
||||
|
||||
fn certificate_hash(certificate_der: &[u8]) -> Result<WebTransportCertificateHash, game_realtime_transport_lib::TransportError> {
|
||||
let certificate = web_transport_quinn::quinn::rustls::pki_types::CertificateDer::from(certificate_der.to_vec());
|
||||
let provider = web_transport_quinn::crypto::default_provider();
|
||||
@@ -418,13 +727,13 @@ 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));
|
||||
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)),
|
||||
Err(_) => return Err(message_too_large(payload_len, max_message_size)),
|
||||
};
|
||||
return Ok(payload_len.to_be_bytes());
|
||||
}
|
||||
@@ -433,60 +742,57 @@ fn invalid_configuration(detail: impl Into<String>) -> game_realtime_transport_l
|
||||
return transport_error(game_realtime_transport_lib::TransportErrorKind::InvalidConfiguration, detail);
|
||||
}
|
||||
|
||||
fn message_too_large(payload_len: usize) -> game_realtime_transport_lib::TransportError {
|
||||
fn map_read_error(error: web_transport_quinn::ReadError) -> game_realtime_transport_lib::TransportError {
|
||||
let kind = match &error {
|
||||
web_transport_quinn::ReadError::Reset(_) => game_realtime_transport_lib::TransportErrorKind::Aborted,
|
||||
web_transport_quinn::ReadError::ClosedStream => game_realtime_transport_lib::TransportErrorKind::Closed,
|
||||
web_transport_quinn::ReadError::SessionError(session) => map_session_error_kind(session),
|
||||
web_transport_quinn::ReadError::InvalidReset(_) | web_transport_quinn::ReadError::IllegalOrderedRead => {
|
||||
game_realtime_transport_lib::TransportErrorKind::Protocol
|
||||
},
|
||||
};
|
||||
return transport_error(kind, error.to_string());
|
||||
}
|
||||
|
||||
fn map_session_error_kind(error: &web_transport_quinn::SessionError) -> game_realtime_transport_lib::TransportErrorKind {
|
||||
if matches!(error, web_transport_quinn::SessionError::WebTransportError(web_transport_quinn::WebTransportError::Closed(_, _))) {
|
||||
return game_realtime_transport_lib::TransportErrorKind::Closed;
|
||||
}
|
||||
return game_realtime_transport_lib::TransportErrorKind::Io;
|
||||
}
|
||||
|
||||
fn map_write_error(error: web_transport_quinn::WriteError) -> game_realtime_transport_lib::TransportError {
|
||||
let kind = match &error {
|
||||
web_transport_quinn::WriteError::Stopped(_) => game_realtime_transport_lib::TransportErrorKind::Aborted,
|
||||
web_transport_quinn::WriteError::ClosedStream => game_realtime_transport_lib::TransportErrorKind::Closed,
|
||||
web_transport_quinn::WriteError::SessionError(session) => map_session_error_kind(session),
|
||||
web_transport_quinn::WriteError::InvalidStopped(_) => game_realtime_transport_lib::TransportErrorKind::Protocol,
|
||||
};
|
||||
return transport_error(kind, error.to_string());
|
||||
}
|
||||
|
||||
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 primary frame payload size {payload_len} exceeds {PRIMARY_FRAME_MAX_PAYLOAD_SIZE} bytes"),
|
||||
format!("WebTransport primary frame payload size {payload_len} exceeds configured maximum {max_message_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 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", duration_millis(timeout)),
|
||||
);
|
||||
}
|
||||
|
||||
fn duration_millis(duration: std::time::Duration) -> u128 {
|
||||
return duration.as_millis();
|
||||
}
|
||||
|
||||
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);
|
||||
}
|
||||
|
||||
146
crates/common/game-realtime-webtransport-lib/tests/robustness.rs
Normal file
146
crates/common/game-realtime-webtransport-lib/tests/robustness.rs
Normal file
@@ -0,0 +1,146 @@
|
||||
// file: crates/common/game-realtime-webtransport-lib/tests/robustness.rs
|
||||
// version: 2
|
||||
|
||||
//! Negative and bounded lifecycle tests for the native reliable WebTransport backend.
|
||||
|
||||
const SHORT_OPERATION_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(50);
|
||||
const SMALL_MESSAGE_LIMIT: usize = 32;
|
||||
const TEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(3);
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn outbound_payload_over_the_configured_limit_is_rejected_before_write() {
|
||||
let config = game_realtime_webtransport_lib::WebTransportConfig::default().with_max_message_size(SMALL_MESSAGE_LIMIT);
|
||||
let (_server_connection, client_connection) = establish_backend_pair(config, config).await;
|
||||
let (mut client_sender, _client_receiver) = game_realtime_transport_lib::RealtimeConnection::split(client_connection);
|
||||
let oversized = game_realtime_transport_lib::TransportMessage::new(vec![7; SMALL_MESSAGE_LIMIT + 1]);
|
||||
let result = game_realtime_transport_lib::RealtimeSender::send(&mut client_sender, oversized).await;
|
||||
match result {
|
||||
Ok(()) => panic!("oversized outbound WebTransport payload was accepted"),
|
||||
Err(error) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::MessageTooLarge),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn inbound_payload_over_the_configured_limit_is_rejected_before_allocation() {
|
||||
let server_config = game_realtime_webtransport_lib::WebTransportConfig::default().with_max_message_size(SMALL_MESSAGE_LIMIT);
|
||||
let client_config = game_realtime_webtransport_lib::WebTransportConfig::default().with_max_message_size(SMALL_MESSAGE_LIMIT + 1);
|
||||
let (server_connection, client_connection) = establish_backend_pair(server_config, client_config).await;
|
||||
let (_server_sender, mut server_receiver) = game_realtime_transport_lib::RealtimeConnection::split(server_connection);
|
||||
let (mut client_sender, _client_receiver) = game_realtime_transport_lib::RealtimeConnection::split(client_connection);
|
||||
let payload = game_realtime_transport_lib::TransportMessage::new(vec![3; SMALL_MESSAGE_LIMIT + 1]);
|
||||
if let Err(error) = game_realtime_transport_lib::RealtimeSender::send(&mut client_sender, payload).await {
|
||||
panic!("client failed to send payload allowed by its local bound: {error}");
|
||||
}
|
||||
let receive = tokio::time::timeout(TEST_TIMEOUT, game_realtime_transport_lib::RealtimeReceiver::receive(&mut server_receiver)).await;
|
||||
match receive {
|
||||
Ok(Ok(value)) => panic!("oversized inbound WebTransport payload produced a successful receive: {value:?}"),
|
||||
Ok(Err(error)) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::MessageTooLarge),
|
||||
Err(_) => panic!("oversized inbound WebTransport payload did not complete within the test timeout"),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn explicit_sender_abort_is_observed_as_aborted_receive() {
|
||||
let config = game_realtime_webtransport_lib::WebTransportConfig::default();
|
||||
let (server_connection, client_connection) = establish_backend_pair(config, config).await;
|
||||
let (_server_sender, mut server_receiver) = game_realtime_transport_lib::RealtimeConnection::split(server_connection);
|
||||
let (mut client_sender, _client_receiver) = game_realtime_transport_lib::RealtimeConnection::split(client_connection);
|
||||
if let Err(error) = client_sender.abort(0x41) {
|
||||
panic!("client sender abort failed: {error}");
|
||||
}
|
||||
let receive = tokio::time::timeout(TEST_TIMEOUT, game_realtime_transport_lib::RealtimeReceiver::receive(&mut server_receiver)).await;
|
||||
match receive {
|
||||
Ok(Ok(value)) => panic!("reset WebTransport stream produced a successful receive: {value:?}"),
|
||||
Ok(Err(error)) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::Aborted),
|
||||
Err(_) => panic!("peer reset was not observed within the test timeout"),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn explicit_receiver_abort_makes_the_local_receive_half_terminal() {
|
||||
let config = game_realtime_webtransport_lib::WebTransportConfig::default();
|
||||
let (server_connection, _client_connection) = establish_backend_pair(config, config).await;
|
||||
let (_server_sender, mut server_receiver) = game_realtime_transport_lib::RealtimeConnection::split(server_connection);
|
||||
if let Err(error) = server_receiver.abort(0x42) {
|
||||
panic!("server receiver abort failed: {error}");
|
||||
}
|
||||
let result = game_realtime_transport_lib::RealtimeReceiver::receive(&mut server_receiver).await;
|
||||
match result {
|
||||
Ok(value) => panic!("aborted WebTransport receiver produced a successful receive: {value:?}"),
|
||||
Err(error) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::Aborted),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn dropping_sender_without_close_resets_the_stream_instead_of_synthesizing_fin() {
|
||||
let config = game_realtime_webtransport_lib::WebTransportConfig::default();
|
||||
let (server_connection, client_connection) = establish_backend_pair(config, config).await;
|
||||
let (_server_sender, mut server_receiver) = game_realtime_transport_lib::RealtimeConnection::split(server_connection);
|
||||
let (client_sender, _client_receiver) = game_realtime_transport_lib::RealtimeConnection::split(client_connection);
|
||||
drop(client_sender);
|
||||
let receive = tokio::time::timeout(TEST_TIMEOUT, game_realtime_transport_lib::RealtimeReceiver::receive(&mut server_receiver)).await;
|
||||
match receive {
|
||||
Ok(Ok(value)) => panic!("dropped WebTransport sender produced a clean receive result: {value:?}"),
|
||||
Ok(Err(error)) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::Aborted),
|
||||
Err(_) => panic!("sender drop reset was not observed within the test timeout"),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn primary_stream_accept_honors_the_configured_deadline() {
|
||||
let config = game_realtime_webtransport_lib::WebTransportConfig::default().with_primary_stream_timeout(SHORT_OPERATION_TIMEOUT);
|
||||
let (server_session, _client_session) = establish_sessions(config, config).await;
|
||||
let result = server_session.accept_primary_connection().await;
|
||||
match result {
|
||||
Ok(_) => panic!("primary stream accept unexpectedly succeeded without a peer-created stream"),
|
||||
Err(error) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::Timeout),
|
||||
}
|
||||
}
|
||||
|
||||
async fn establish_backend_pair(
|
||||
server_transport: game_realtime_webtransport_lib::WebTransportConfig,
|
||||
client_transport: game_realtime_webtransport_lib::WebTransportConfig,
|
||||
) -> (game_realtime_webtransport_lib::WebTransportConnection, game_realtime_webtransport_lib::WebTransportConnection) {
|
||||
let (server_session, client_session) = establish_sessions(server_transport, client_transport).await;
|
||||
let client_connection = match client_session.open_primary_connection().await {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("client primary stream open failed: {error}"),
|
||||
};
|
||||
let server_connection = match server_session.accept_primary_connection().await {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("server primary stream accept failed: {error}"),
|
||||
};
|
||||
return (server_connection, client_connection);
|
||||
}
|
||||
|
||||
async fn establish_sessions(
|
||||
server_transport: game_realtime_webtransport_lib::WebTransportConfig,
|
||||
client_transport: game_realtime_webtransport_lib::WebTransportConfig,
|
||||
) -> (game_realtime_webtransport_lib::WebTransportSession, game_realtime_webtransport_lib::WebTransportSession) {
|
||||
let identity = match game_realtime_webtransport_lib::WebTransportServerIdentity::generate_loopback() {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("loopback identity generation failed: {error}"),
|
||||
};
|
||||
let certificate_hash = identity.certificate_hash().clone();
|
||||
let server_config = game_realtime_webtransport_lib::WebTransportServerConfig::new(std::net::SocketAddr::from(([127, 0, 0, 1], 0)), identity)
|
||||
.with_transport_config(server_transport);
|
||||
let mut listener = match game_realtime_webtransport_lib::WebTransportListener::bind(server_config) {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("WebTransport listener bind failed: {error}"),
|
||||
};
|
||||
let endpoint = format!("https://{}/robustness", listener.local_addr());
|
||||
let client_config = match game_realtime_webtransport_lib::WebTransportClientConfig::new(endpoint.as_str(), certificate_hash) {
|
||||
Ok(value) => value.with_transport_config(client_transport),
|
||||
Err(error) => panic!("WebTransport client configuration failed: {error}"),
|
||||
};
|
||||
let sessions = tokio::time::timeout(TEST_TIMEOUT, async {
|
||||
return tokio::join!(listener.accept(), game_realtime_webtransport_lib::connect(&client_config));
|
||||
})
|
||||
.await;
|
||||
return match sessions {
|
||||
Ok((Ok(server), Ok(client))) => (server, client),
|
||||
Ok((Err(error), _)) => panic!("WebTransport server establishment failed: {error}"),
|
||||
Ok((_, Err(error))) => panic!("WebTransport client establishment failed: {error}"),
|
||||
Err(_) => panic!("WebTransport loopback establishment timed out"),
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
// file: crates/common/game-realtime-webtransport-lib/unit_tests/config.rs
|
||||
// version: 1
|
||||
|
||||
#[test]
|
||||
fn defaults_are_valid_and_preserve_the_one_mib_baseline() {
|
||||
let config = super::WebTransportConfig::default();
|
||||
assert!(config.validate().is_ok());
|
||||
assert_eq!(config.max_message_size(), 1024 * 1024);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn zero_limits_and_deadlines_are_rejected() {
|
||||
let cases = [
|
||||
super::WebTransportConfig::default().with_max_message_size(0),
|
||||
super::WebTransportConfig::default().with_connect_timeout(std::time::Duration::ZERO),
|
||||
super::WebTransportConfig::default().with_primary_stream_timeout(std::time::Duration::ZERO),
|
||||
super::WebTransportConfig::default().with_send_timeout(std::time::Duration::ZERO),
|
||||
];
|
||||
for config in cases {
|
||||
let result = config.validate();
|
||||
match result {
|
||||
Ok(()) => panic!("invalid WebTransport configuration unexpectedly accepted"),
|
||||
Err(error) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::InvalidConfiguration),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,7 @@
|
||||
// file: crates/common/game-realtime-webtransport-lib/unit_tests/webtransport.rs
|
||||
// version: 2
|
||||
// version: 4
|
||||
|
||||
const TEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(3);
|
||||
|
||||
#[test]
|
||||
fn certificate_hash_preserves_exact_sha256_bytes() {
|
||||
@@ -26,19 +28,33 @@ fn client_config_accepts_https_and_rejects_non_secure_schemes() {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn frame_header_is_big_endian_and_payload_bound_is_enforced() {
|
||||
let header = match super::frame_header(0x00_01_02_03) {
|
||||
fn frame_header_is_big_endian_and_configured_payload_bound_is_enforced() {
|
||||
let header = match super::frame_header(0x00_01_02_03, 1024 * 1024) {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("valid frame header rejected: {error}"),
|
||||
};
|
||||
assert_eq!(header, [0x00, 0x01, 0x02, 0x03]);
|
||||
let oversized = super::frame_header(super::PRIMARY_FRAME_MAX_PAYLOAD_SIZE + 1);
|
||||
let oversized = super::frame_header(33, 32);
|
||||
match oversized {
|
||||
Ok(_) => panic!("oversized WebTransport frame unexpectedly accepted"),
|
||||
Err(error) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::MessageTooLarge),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stream_error_mapping_distinguishes_abort_close_and_protocol_failures() {
|
||||
let stopped = super::map_write_error(web_transport_quinn::WriteError::Stopped(7));
|
||||
assert_eq!(stopped.kind(), game_realtime_transport_lib::TransportErrorKind::Aborted);
|
||||
let reset = super::map_read_error(web_transport_quinn::ReadError::Reset(8));
|
||||
assert_eq!(reset.kind(), game_realtime_transport_lib::TransportErrorKind::Aborted);
|
||||
let closed = super::map_read_error(web_transport_quinn::ReadError::SessionError(web_transport_quinn::SessionError::WebTransportError(
|
||||
web_transport_quinn::WebTransportError::Closed(9, "done".to_owned()),
|
||||
)));
|
||||
assert_eq!(closed.kind(), game_realtime_transport_lib::TransportErrorKind::Closed);
|
||||
let protocol = super::map_read_error(web_transport_quinn::ReadError::IllegalOrderedRead);
|
||||
assert_eq!(protocol.kind(), game_realtime_transport_lib::TransportErrorKind::Protocol);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn generated_loopback_identity_has_sha256_fingerprint() {
|
||||
let identity = match super::WebTransportServerIdentity::generate_loopback() {
|
||||
@@ -61,3 +77,118 @@ fn injected_identity_rejects_empty_certificate_or_key() {
|
||||
Err(error) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::InvalidConfiguration),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn cancelled_receive_preserves_partial_frame_state() {
|
||||
let (server_session, client_session) = establish_private_sessions().await;
|
||||
let (mut raw_sender, _raw_receiver) = match client_session.inner.open_bi().await {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("raw client stream open failed: {error}"),
|
||||
};
|
||||
let (server_sender, server_receiver) = match server_session.inner.accept_bi().await {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("raw server stream accept failed: {error}"),
|
||||
};
|
||||
let server_connection = super::WebTransportConnection::new(server_session.inner, server_sender, server_receiver, server_session.transport);
|
||||
let (_server_sender, mut server_receiver) = game_realtime_transport_lib::RealtimeConnection::split(server_connection);
|
||||
if let Err(error) = raw_sender.write_all(&[0x00, 0x00]).await {
|
||||
panic!("partial frame header write failed: {error}");
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
|
||||
let first_receive =
|
||||
tokio::time::timeout(std::time::Duration::from_millis(100), game_realtime_transport_lib::RealtimeReceiver::receive(&mut server_receiver)).await;
|
||||
assert!(first_receive.is_err());
|
||||
assert_eq!(server_receiver.header_read, 2);
|
||||
if let Err(error) = raw_sender.write_all(&[0x00, 0x03, 0x10, 0x20, 0x30]).await {
|
||||
panic!("remaining frame write failed: {error}");
|
||||
}
|
||||
let second_receive = tokio::time::timeout(TEST_TIMEOUT, game_realtime_transport_lib::RealtimeReceiver::receive(&mut server_receiver)).await;
|
||||
match second_receive {
|
||||
Ok(Ok(game_realtime_transport_lib::TransportReceive::Message(message))) => assert_eq!(message.as_bytes(), &[0x10, 0x20, 0x30]),
|
||||
Ok(Ok(value)) => panic!("resumed receive returned unexpected result: {value:?}"),
|
||||
Ok(Err(error)) => panic!("resumed receive failed: {error}"),
|
||||
Err(_) => panic!("resumed receive timed out"),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn truncated_frame_payload_is_a_protocol_failure() {
|
||||
let (server_session, client_session) = establish_private_sessions().await;
|
||||
let (mut raw_sender, _raw_receiver) = match client_session.inner.open_bi().await {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("raw client stream open failed: {error}"),
|
||||
};
|
||||
let (server_sender, server_receiver) = match server_session.inner.accept_bi().await {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("raw server stream accept failed: {error}"),
|
||||
};
|
||||
let server_connection = super::WebTransportConnection::new(server_session.inner, server_sender, server_receiver, server_session.transport);
|
||||
let (_server_sender, mut server_receiver) = game_realtime_transport_lib::RealtimeConnection::split(server_connection);
|
||||
if let Err(error) = raw_sender.write_all(&[0x00, 0x00, 0x00, 0x03, 0x10]).await {
|
||||
panic!("truncated frame payload write failed: {error}");
|
||||
}
|
||||
if let Err(error) = raw_sender.finish() {
|
||||
panic!("raw client stream finish failed: {error}");
|
||||
}
|
||||
let receive = tokio::time::timeout(TEST_TIMEOUT, game_realtime_transport_lib::RealtimeReceiver::receive(&mut server_receiver)).await;
|
||||
match receive {
|
||||
Ok(Ok(value)) => panic!("truncated frame payload produced a successful receive: {value:?}"),
|
||||
Ok(Err(error)) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::Protocol),
|
||||
Err(_) => panic!("truncated frame payload was not rejected within the test timeout"),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn truncated_frame_header_is_a_protocol_failure() {
|
||||
let (server_session, client_session) = establish_private_sessions().await;
|
||||
let (mut raw_sender, _raw_receiver) = match client_session.inner.open_bi().await {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("raw client stream open failed: {error}"),
|
||||
};
|
||||
let (server_sender, server_receiver) = match server_session.inner.accept_bi().await {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("raw server stream accept failed: {error}"),
|
||||
};
|
||||
let server_connection = super::WebTransportConnection::new(server_session.inner, server_sender, server_receiver, server_session.transport);
|
||||
let (_server_sender, mut server_receiver) = game_realtime_transport_lib::RealtimeConnection::split(server_connection);
|
||||
if let Err(error) = raw_sender.write_all(&[0x00, 0x00]).await {
|
||||
panic!("partial frame header write failed: {error}");
|
||||
}
|
||||
if let Err(error) = raw_sender.finish() {
|
||||
panic!("raw client stream finish failed: {error}");
|
||||
}
|
||||
let receive = tokio::time::timeout(TEST_TIMEOUT, game_realtime_transport_lib::RealtimeReceiver::receive(&mut server_receiver)).await;
|
||||
match receive {
|
||||
Ok(Ok(value)) => panic!("truncated frame header produced a successful receive: {value:?}"),
|
||||
Ok(Err(error)) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::Protocol),
|
||||
Err(_) => panic!("truncated frame header was not rejected within the test timeout"),
|
||||
}
|
||||
}
|
||||
|
||||
async fn establish_private_sessions() -> (super::WebTransportSession, super::WebTransportSession) {
|
||||
let identity = match super::WebTransportServerIdentity::generate_loopback() {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("loopback identity generation failed: {error}"),
|
||||
};
|
||||
let certificate_hash = identity.certificate_hash().clone();
|
||||
let server_config = super::WebTransportServerConfig::new(std::net::SocketAddr::from(([127, 0, 0, 1], 0)), identity);
|
||||
let mut listener = match super::WebTransportListener::bind(server_config) {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("WebTransport listener bind failed: {error}"),
|
||||
};
|
||||
let endpoint = format!("https://{}/private-test", listener.local_addr());
|
||||
let client_config = match super::WebTransportClientConfig::new(endpoint.as_str(), certificate_hash) {
|
||||
Ok(value) => value,
|
||||
Err(error) => panic!("WebTransport client configuration failed: {error}"),
|
||||
};
|
||||
let pair = tokio::time::timeout(TEST_TIMEOUT, async {
|
||||
return tokio::join!(listener.accept(), super::connect(&client_config));
|
||||
})
|
||||
.await;
|
||||
return match pair {
|
||||
Ok((Ok(server), Ok(client))) => (server, client),
|
||||
Ok((Err(error), _)) => panic!("WebTransport server establishment failed: {error}"),
|
||||
Ok((_, Err(error))) => panic!("WebTransport client establishment failed: {error}"),
|
||||
Err(_) => panic!("WebTransport private loopback establishment timed out"),
|
||||
};
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user