0.3.4-alpha.4

This commit is contained in:
2026-09-21 18:00:32 +02:00
parent bbecae5e0d
commit fa17ea0ca5
12 changed files with 792 additions and 47 deletions

View File

@@ -1,5 +1,5 @@
# file: crates/common/game-realtime-websocket-lib/Cargo.toml
# version: 2
# version: 3
[package]
name = "game-realtime-websocket-lib"
@@ -13,7 +13,7 @@ publish.workspace = true
[dependencies]
futures-util = { workspace = true, features = ["sink", "std"] }
game-realtime-transport-lib = { path = "../game-realtime-transport-lib" }
tokio = { workspace = true, features = ["net"] }
tokio = { workspace = true, features = ["net", "time"] }
tokio-tungstenite = { workspace = true, features = ["connect", "handshake"] }
tracing.workspace = true

View File

@@ -0,0 +1,167 @@
// file: crates/common/game-realtime-websocket-lib/src/config.rs
// version: 1
const DEFAULT_CLOSE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
const DEFAULT_CONNECT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
const DEFAULT_MAX_FRAME_SIZE: usize = 1024 * 1024;
const DEFAULT_MAX_MESSAGE_SIZE: usize = 1024 * 1024;
const DEFAULT_MAX_WRITE_BUFFER_SIZE: usize = 2 * 1024 * 1024;
const DEFAULT_SEND_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
const DEFAULT_WRITE_BUFFER_SIZE: usize = 64 * 1024;
/// Product-facing limits and operation deadlines for one WebSocket connection.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct WebSocketConfig {
max_message_size: usize,
max_frame_size: usize,
write_buffer_size: usize,
max_write_buffer_size: usize,
connect_timeout: std::time::Duration,
send_timeout: std::time::Duration,
close_timeout: std::time::Duration,
}
impl Default for WebSocketConfig {
fn default() -> Self {
return Self {
max_message_size: DEFAULT_MAX_MESSAGE_SIZE,
max_frame_size: DEFAULT_MAX_FRAME_SIZE,
write_buffer_size: DEFAULT_WRITE_BUFFER_SIZE,
max_write_buffer_size: DEFAULT_MAX_WRITE_BUFFER_SIZE,
connect_timeout: DEFAULT_CONNECT_TIMEOUT,
send_timeout: DEFAULT_SEND_TIMEOUT,
close_timeout: DEFAULT_CLOSE_TIMEOUT,
};
}
}
impl WebSocketConfig {
/// Returns a copy with a different maximum 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 maximum WebSocket frame payload size.
#[must_use]
pub fn with_max_frame_size(mut self, value: usize) -> Self {
self.max_frame_size = value;
return self;
}
/// Returns a copy with a different Tungstenite write-buffer target.
#[must_use]
pub fn with_write_buffer_size(mut self, value: usize) -> Self {
self.write_buffer_size = value;
return self;
}
/// Returns a copy with a different hard maximum for the Tungstenite write buffer.
#[must_use]
pub fn with_max_write_buffer_size(mut self, value: usize) -> Self {
self.max_write_buffer_size = value;
return self;
}
/// Returns a copy with a different connection or 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 deadline for one send operation.
#[must_use]
pub fn with_send_timeout(mut self, value: std::time::Duration) -> Self {
self.send_timeout = value;
return self;
}
/// Returns a copy with a different deadline for a clean local close operation.
#[must_use]
pub fn with_close_timeout(mut self, value: std::time::Duration) -> Self {
self.close_timeout = value;
return self;
}
/// Returns the configured maximum binary message size.
#[must_use]
pub fn max_message_size(&self) -> usize {
return self.max_message_size;
}
/// Returns the configured maximum WebSocket frame payload size.
#[must_use]
pub fn max_frame_size(&self) -> usize {
return self.max_frame_size;
}
/// Returns the configured Tungstenite write-buffer target.
#[must_use]
pub fn write_buffer_size(&self) -> usize {
return self.write_buffer_size;
}
/// Returns the configured hard maximum for the Tungstenite write buffer.
#[must_use]
pub fn max_write_buffer_size(&self) -> usize {
return self.max_write_buffer_size;
}
/// Returns the configured connection or handshake deadline.
#[must_use]
pub fn connect_timeout(&self) -> std::time::Duration {
return self.connect_timeout;
}
/// Returns the configured deadline for one send operation.
#[must_use]
pub fn send_timeout(&self) -> std::time::Duration {
return self.send_timeout;
}
/// Returns the configured deadline for a clean local close operation.
#[must_use]
pub fn close_timeout(&self) -> std::time::Duration {
return self.close_timeout;
}
/// Validates all invariants required before creating a WebSocket endpoint or connection.
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 self.max_frame_size == 0 {
return Err(invalid_configuration("max_frame_size must be greater than zero"));
}
if self.max_frame_size > self.max_message_size {
return Err(invalid_configuration("max_frame_size must not exceed max_message_size"));
}
let minimum_max_write_buffer_size = match self.write_buffer_size.checked_add(self.max_message_size) {
Some(value) => value,
None => return Err(invalid_configuration("write_buffer_size + max_message_size overflows usize")),
};
if self.max_write_buffer_size < minimum_max_write_buffer_size {
return Err(invalid_configuration("max_write_buffer_size must fit write_buffer_size plus one maximum-sized message"));
}
if self.connect_timeout.is_zero() {
return Err(invalid_configuration("connect_timeout must be greater than zero"));
}
if self.send_timeout.is_zero() {
return Err(invalid_configuration("send_timeout must be greater than zero"));
}
if self.close_timeout.is_zero() {
return Err(invalid_configuration("close_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;

View File

@@ -1,5 +1,5 @@
// file: crates/common/game-realtime-websocket-lib/src/lib.rs
// version: 1
// version: 2
#![warn(missing_docs)]
#![deny(unreachable_pub)]
@@ -7,8 +7,11 @@
//! Tokio/tokio-tungstenite WebSocket backend for the games.sasedev realtime transport contract.
mod config;
mod websocket;
/// Re-export of product-facing WebSocket limits and operation deadlines.
pub use self::config::WebSocketConfig;
/// Re-export of an established WebSocket transport connection.
pub use self::websocket::WebSocketConnection;
/// Re-export of a bound WebSocket server listener.
@@ -17,5 +20,7 @@ pub use self::websocket::WebSocketListener;
pub use self::websocket::WebSocketReceiver;
/// Re-export of the send half of an established WebSocket connection.
pub use self::websocket::WebSocketSender;
/// Re-export of the WebSocket client connection constructor.
/// Re-export of the WebSocket client connection constructor using baseline defaults.
pub use self::websocket::connect;
/// Re-export of the configurable WebSocket client connection constructor.
pub use self::websocket::connect_with_config;

View File

@@ -1,5 +1,5 @@
// file: crates/common/game-realtime-websocket-lib/src/websocket.rs
// version: 1
// version: 2
use futures_util::SinkExt; // rust-rules: trait-import
use futures_util::StreamExt; // rust-rules: trait-import
@@ -31,15 +31,16 @@ enum WebSocketReceiverKind {
/// Established WebSocket connection implementing the transport-neutral realtime contract.
pub struct WebSocketConnection {
inner: WebSocketStreamKind,
config: crate::WebSocketConfig,
}
impl WebSocketConnection {
fn from_client(stream: ClientStream) -> Self {
return Self { inner: WebSocketStreamKind::Client(stream) };
fn from_client(stream: ClientStream, config: crate::WebSocketConfig) -> Self {
return Self { inner: WebSocketStreamKind::Client(stream), config };
}
fn from_server(stream: ServerStream) -> Self {
return Self { inner: WebSocketStreamKind::Server(stream) };
fn from_server(stream: ServerStream, config: crate::WebSocketConfig) -> Self {
return Self { inner: WebSocketStreamKind::Server(stream), config };
}
}
@@ -48,18 +49,36 @@ impl game_realtime_transport_lib::RealtimeConnection for WebSocketConnection {
type Receiver = crate::WebSocketReceiver;
fn split(self) -> (Self::Sender, Self::Receiver) {
let max_frame_size = self.config.max_frame_size();
let max_message_size = self.config.max_message_size();
let send_timeout = self.config.send_timeout();
let close_timeout = self.config.close_timeout();
return match self.inner {
WebSocketStreamKind::Client(stream) => {
let (sender, receiver) = stream.split();
(
crate::WebSocketSender { inner: WebSocketSinkKind::Client(sender) },
crate::WebSocketSender {
inner: WebSocketSinkKind::Client(sender),
max_frame_size,
max_message_size,
send_timeout,
close_timeout,
send_timed_out: false,
},
crate::WebSocketReceiver { inner: WebSocketReceiverKind::Client(receiver) },
)
},
WebSocketStreamKind::Server(stream) => {
let (sender, receiver) = stream.split();
(
crate::WebSocketSender { inner: WebSocketSinkKind::Server(sender) },
crate::WebSocketSender {
inner: WebSocketSinkKind::Server(sender),
max_frame_size,
max_message_size,
send_timeout,
close_timeout,
send_timed_out: false,
},
crate::WebSocketReceiver { inner: WebSocketReceiverKind::Server(receiver) },
)
},
@@ -70,6 +89,11 @@ impl game_realtime_transport_lib::RealtimeConnection for WebSocketConnection {
/// Send half of an established WebSocket transport connection.
pub struct WebSocketSender {
inner: WebSocketSinkKind,
max_frame_size: usize,
max_message_size: usize,
send_timeout: std::time::Duration,
close_timeout: std::time::Duration,
send_timed_out: bool,
}
impl game_realtime_transport_lib::RealtimeSender for WebSocketSender {
@@ -84,42 +108,81 @@ impl game_realtime_transport_lib::RealtimeSender for WebSocketSender {
fn send(&mut self, message: game_realtime_transport_lib::TransportMessage) -> Self::SendFuture<'_> {
return Box::pin(async move {
if self.send_timed_out {
return Err(game_realtime_transport_lib::TransportError::new(
game_realtime_transport_lib::TransportErrorKind::Aborted,
"sender is unavailable after a previous send timeout",
));
}
let payload_len = message.len();
if payload_len > self.max_message_size || payload_len > self.max_frame_size {
let error = game_realtime_transport_lib::TransportError::new(
game_realtime_transport_lib::TransportErrorKind::MessageTooLarge,
format!("binary payload size {payload_len} exceeds configured message/frame maxima {}/{}", self.max_message_size, self.max_frame_size),
);
tracing::warn!(
target: TRACING_TARGET,
payload_len = payload_len,
max_message_size = self.max_message_size,
max_frame_size = self.max_frame_size,
"outbound WebSocket payload rejected"
);
return Err(error);
}
let websocket_message = tokio_tungstenite::tungstenite::Message::Binary(message.into_bytes().into());
let result = match &mut self.inner {
WebSocketSinkKind::Client(sender) => sender.send(websocket_message).await,
WebSocketSinkKind::Server(sender) => sender.send(websocket_message).await,
let send_timeout = self.send_timeout;
let send = async {
return match &mut self.inner {
WebSocketSinkKind::Client(sender) => sender.send(websocket_message).await,
WebSocketSinkKind::Server(sender) => sender.send(websocket_message).await,
};
};
let result = tokio::time::timeout(send_timeout, send).await;
return match result {
Ok(()) => {
Ok(Ok(())) => {
tracing::trace!(target: TRACING_TARGET, payload_len = payload_len, "binary WebSocket payload sent");
Ok(())
},
Err(error) => {
Ok(Err(error)) => {
let mapped = map_stream_error(error);
tracing::warn!(target: TRACING_TARGET, kind = %mapped.kind(), detail = mapped.detail(), "WebSocket send failed");
Err(mapped)
},
Err(_) => {
self.send_timed_out = true;
let error = timeout_error("WebSocket send", send_timeout);
tracing::warn!(target: TRACING_TARGET, timeout_ms = duration_millis(send_timeout), "WebSocket send timed out");
Err(error)
},
};
});
}
fn close(&mut self) -> Self::CloseFuture<'_> {
return Box::pin(async move {
let result = match &mut self.inner {
WebSocketSinkKind::Client(sender) => sender.close().await,
WebSocketSinkKind::Server(sender) => sender.close().await,
let close_timeout = self.close_timeout;
let close = async {
return match &mut self.inner {
WebSocketSinkKind::Client(sender) => sender.close().await,
WebSocketSinkKind::Server(sender) => sender.close().await,
};
};
let result = tokio::time::timeout(close_timeout, close).await;
return match result {
Ok(()) => {
Ok(Ok(())) => {
tracing::debug!(target: TRACING_TARGET, "local WebSocket close initiated");
Ok(())
},
Err(error) => {
Ok(Err(error)) => {
let mapped = map_stream_error(error);
tracing::warn!(target: TRACING_TARGET, kind = %mapped.kind(), detail = mapped.detail(), "WebSocket close failed");
Err(mapped)
},
Err(_) => {
let error = timeout_error("WebSocket close", close_timeout);
tracing::warn!(target: TRACING_TARGET, timeout_ms = duration_millis(close_timeout), "WebSocket close timed out");
Err(error)
},
};
});
}
@@ -188,11 +251,20 @@ impl game_realtime_transport_lib::RealtimeReceiver for WebSocketReceiver {
pub struct WebSocketListener {
listener: tokio::net::TcpListener,
local_addr: std::net::SocketAddr,
config: crate::WebSocketConfig,
}
impl WebSocketListener {
/// Binds a WebSocket listener to one concrete socket address.
/// Binds a WebSocket listener with the baseline configuration.
pub async fn bind(address: std::net::SocketAddr) -> Result<Self, game_realtime_transport_lib::TransportError> {
return Self::bind_with_config(address, crate::WebSocketConfig::default()).await;
}
/// Binds a WebSocket listener with explicit product-facing limits and deadlines.
pub async fn bind_with_config(address: std::net::SocketAddr, config: crate::WebSocketConfig) -> Result<Self, game_realtime_transport_lib::TransportError> {
if let Err(error) = config.validate() {
return Err(error);
}
let listener = match tokio::net::TcpListener::bind(address).await {
Ok(value) => value,
Err(error) => {
@@ -210,7 +282,7 @@ impl WebSocketListener {
},
};
tracing::info!(target: TRACING_TARGET, address = %local_addr, "WebSocket listener bound");
return Ok(Self { listener, local_addr });
return Ok(Self { listener, local_addr, config });
}
/// Returns the concrete local socket address, including an ephemeral port selected by the OS.
@@ -219,7 +291,7 @@ impl WebSocketListener {
return self.local_addr;
}
/// Accepts one TCP peer and completes the server-side WebSocket handshake.
/// Accepts one TCP peer and completes a bounded server-side WebSocket handshake.
pub async fn accept(&self) -> Result<crate::WebSocketConnection, game_realtime_transport_lib::TransportError> {
let (stream, peer_addr) = match self.listener.accept().await {
Ok(value) => value,
@@ -229,31 +301,51 @@ impl WebSocketListener {
return Err(mapped);
},
};
let websocket = match tokio_tungstenite::accept_async(stream).await {
Ok(value) => value,
Err(error) => {
let handshake = tokio_tungstenite::accept_async_with_config(stream, Some(tungstenite_config(&self.config)));
let result = tokio::time::timeout(self.config.connect_timeout(), handshake).await;
let websocket = match result {
Ok(Ok(value)) => value,
Ok(Err(error)) => {
let mapped = game_realtime_transport_lib::TransportError::new(game_realtime_transport_lib::TransportErrorKind::Accept, error.to_string());
tracing::warn!(target: TRACING_TARGET, peer = %peer_addr, detail = mapped.detail(), "WebSocket server handshake failed");
return Err(mapped);
},
Err(_) => {
let error = timeout_error("WebSocket server handshake", self.config.connect_timeout());
tracing::warn!(target: TRACING_TARGET, peer = %peer_addr, timeout_ms = duration_millis(self.config.connect_timeout()), "WebSocket server handshake timed out");
return Err(error);
},
};
tracing::info!(target: TRACING_TARGET, peer = %peer_addr, "WebSocket peer accepted");
return Ok(crate::WebSocketConnection::from_server(websocket));
return Ok(crate::WebSocketConnection::from_server(websocket, self.config));
}
}
/// Connects a client to one plain `ws://` endpoint and completes the WebSocket handshake.
/// Connects a client to one plain `ws://` endpoint with the baseline configuration.
pub async fn connect(endpoint: &str) -> Result<crate::WebSocketConnection, game_realtime_transport_lib::TransportError> {
return crate::connect_with_config(endpoint, crate::WebSocketConfig::default()).await;
}
/// Connects a client to one plain `ws://` endpoint with explicit limits and deadlines.
pub async fn connect_with_config(
endpoint: &str,
config: crate::WebSocketConfig,
) -> Result<crate::WebSocketConnection, game_realtime_transport_lib::TransportError> {
if !endpoint.starts_with("ws://") {
return Err(game_realtime_transport_lib::TransportError::new(
game_realtime_transport_lib::TransportErrorKind::InvalidConfiguration,
"the baseline WebSocket backend accepts only ws:// endpoints",
));
}
if let Err(error) = config.validate() {
return Err(error);
}
tracing::debug!(target: TRACING_TARGET, endpoint = endpoint, "connecting WebSocket client");
let (stream, _) = match tokio_tungstenite::connect_async(endpoint).await {
Ok(value) => value,
Err(error) => {
let handshake = tokio_tungstenite::connect_async_with_config(endpoint, Some(tungstenite_config(&config)), false);
let result = tokio::time::timeout(config.connect_timeout(), handshake).await;
let (stream, _) = match result {
Ok(Ok(value)) => value,
Ok(Err(error)) => {
let mapped = map_connect_error(error);
tracing::warn!(
target: TRACING_TARGET,
@@ -264,9 +356,33 @@ pub async fn connect(endpoint: &str) -> Result<crate::WebSocketConnection, game_
);
return Err(mapped);
},
Err(_) => {
let error = timeout_error("WebSocket client connect", config.connect_timeout());
tracing::warn!(target: TRACING_TARGET, endpoint = endpoint, timeout_ms = duration_millis(config.connect_timeout()), "WebSocket client connect timed out");
return Err(error);
},
};
tracing::info!(target: TRACING_TARGET, endpoint = endpoint, "WebSocket client connected");
return Ok(crate::WebSocketConnection::from_client(stream));
return Ok(crate::WebSocketConnection::from_client(stream, config));
}
fn tungstenite_config(config: &crate::WebSocketConfig) -> tokio_tungstenite::tungstenite::protocol::WebSocketConfig {
return tokio_tungstenite::tungstenite::protocol::WebSocketConfig::default()
.write_buffer_size(config.write_buffer_size())
.max_write_buffer_size(config.max_write_buffer_size())
.max_message_size(Some(config.max_message_size()))
.max_frame_size(Some(config.max_frame_size()));
}
fn timeout_error(operation: &str, timeout: std::time::Duration) -> game_realtime_transport_lib::TransportError {
return game_realtime_transport_lib::TransportError::new(
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 map_connect_error(error: tokio_tungstenite::tungstenite::Error) -> game_realtime_transport_lib::TransportError {
@@ -289,3 +405,7 @@ fn map_stream_error(error: tokio_tungstenite::tungstenite::Error) -> game_realti
};
return game_realtime_transport_lib::TransportError::new(kind, error.to_string());
}
#[cfg(test)]
#[path = "../unit_tests/websocket.rs"]
mod tests;

View File

@@ -0,0 +1,134 @@
// file: crates/common/game-realtime-websocket-lib/tests/robustness.rs
// version: 1
//! Negative and bounded lifecycle tests for the WebSocket transport backend.
use futures_util::SinkExt; // rust-rules: trait-import
use game_realtime_transport_lib::RealtimeConnection; // rust-rules: trait-import
use game_realtime_transport_lib::RealtimeReceiver; // rust-rules: trait-import
use game_realtime_transport_lib::RealtimeSender; // rust-rules: trait-import
const SMALL_MESSAGE_LIMIT: usize = 32;
const TEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(3);
type RawClient = tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>;
#[tokio::test(flavor = "current_thread")]
async fn outbound_payload_over_the_configured_limit_is_rejected_before_write() {
let config = small_message_config();
let (server_connection, client_connection) = establish_backend_pair(config).await;
let (_server_sender, _server_receiver) = server_connection.split();
let (mut client_sender, _client_receiver) = client_connection.split();
let oversized = game_realtime_transport_lib::TransportMessage::new(vec![7; SMALL_MESSAGE_LIMIT + 1]);
let result = client_sender.send(oversized).await;
match result {
Ok(()) => panic!("oversized outbound 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_by_tungstenite() {
let config = small_message_config();
let (server_connection, mut raw_client) = establish_backend_server_with_raw_client(config).await;
let (_server_sender, mut server_receiver) = server_connection.split();
let send = raw_client.send(tokio_tungstenite::tungstenite::Message::binary(vec![3; SMALL_MESSAGE_LIMIT + 1])).await;
assert!(send.is_ok());
let receive = tokio::time::timeout(TEST_TIMEOUT, server_receiver.receive()).await;
match receive {
Ok(Ok(value)) => panic!("oversized inbound payload produced a successful receive: {value:?}"),
Ok(Err(error)) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::MessageTooLarge),
Err(_) => panic!("oversized inbound payload did not complete within the test timeout"),
}
}
#[tokio::test(flavor = "current_thread")]
async fn text_message_is_rejected_by_the_binary_transport_contract() {
let (server_connection, mut raw_client) = establish_backend_server_with_raw_client(game_realtime_websocket_lib::WebSocketConfig::default()).await;
let (_server_sender, mut server_receiver) = server_connection.split();
let send = raw_client.send(tokio_tungstenite::tungstenite::Message::text("text is outside the transport contract")).await;
assert!(send.is_ok());
let receive = tokio::time::timeout(TEST_TIMEOUT, server_receiver.receive()).await;
match receive {
Ok(Ok(value)) => panic!("text message produced a successful receive: {value:?}"),
Ok(Err(error)) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::Protocol),
Err(_) => panic!("text-message rejection did not complete within the test timeout"),
}
}
#[tokio::test(flavor = "current_thread")]
async fn peer_drop_without_close_handshake_is_reported_as_protocol_failure() {
let (server_connection, raw_client) = establish_backend_server_with_raw_client(game_realtime_websocket_lib::WebSocketConfig::default()).await;
let (_server_sender, mut server_receiver) = server_connection.split();
drop(raw_client);
let receive = tokio::time::timeout(TEST_TIMEOUT, server_receiver.receive()).await;
match receive {
Ok(Ok(value)) => panic!("abrupt peer drop was reported as a successful receive: {value:?}"),
Ok(Err(error)) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::Protocol),
Err(_) => panic!("abrupt peer drop was not observed within the test timeout"),
}
}
#[tokio::test(flavor = "current_thread")]
async fn silent_tcp_peer_hits_the_server_handshake_deadline() {
let config = game_realtime_websocket_lib::WebSocketConfig::default().with_connect_timeout(std::time::Duration::from_millis(50));
let bind_address = std::net::SocketAddr::from(([127, 0, 0, 1], 0));
let listener = match game_realtime_websocket_lib::WebSocketListener::bind_with_config(bind_address, config).await {
Ok(value) => value,
Err(error) => panic!("bounded listener bind failed: {error}"),
};
let _silent_peer = match tokio::net::TcpStream::connect(listener.local_addr()).await {
Ok(value) => value,
Err(error) => panic!("silent TCP peer connection failed: {error}"),
};
let accept = tokio::time::timeout(TEST_TIMEOUT, listener.accept()).await;
match accept {
Ok(Ok(_connection)) => panic!("silent TCP peer unexpectedly completed a WebSocket handshake"),
Ok(Err(error)) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::Timeout),
Err(_) => panic!("server handshake timeout did not fire within the outer test timeout"),
}
}
fn small_message_config() -> game_realtime_websocket_lib::WebSocketConfig {
return game_realtime_websocket_lib::WebSocketConfig::default().with_max_message_size(SMALL_MESSAGE_LIMIT).with_max_frame_size(SMALL_MESSAGE_LIMIT);
}
async fn establish_backend_pair(
config: game_realtime_websocket_lib::WebSocketConfig,
) -> (game_realtime_websocket_lib::WebSocketConnection, game_realtime_websocket_lib::WebSocketConnection) {
let bind_address = std::net::SocketAddr::from(([127, 0, 0, 1], 0));
let listener = match game_realtime_websocket_lib::WebSocketListener::bind_with_config(bind_address, config).await {
Ok(value) => value,
Err(error) => panic!("loopback listener bind failed: {error}"),
};
let endpoint = format!("ws://{}/", listener.local_addr());
let pair = tokio::time::timeout(TEST_TIMEOUT, async {
return tokio::join!(listener.accept(), game_realtime_websocket_lib::connect_with_config(endpoint.as_str(), config));
})
.await;
return match pair {
Ok((Ok(server), Ok(client))) => (server, client),
Ok((_server, _client)) => panic!("loopback connection establishment failed"),
Err(_) => panic!("loopback connection establishment timed out"),
};
}
async fn establish_backend_server_with_raw_client(
config: game_realtime_websocket_lib::WebSocketConfig,
) -> (game_realtime_websocket_lib::WebSocketConnection, RawClient) {
let bind_address = std::net::SocketAddr::from(([127, 0, 0, 1], 0));
let listener = match game_realtime_websocket_lib::WebSocketListener::bind_with_config(bind_address, config).await {
Ok(value) => value,
Err(error) => panic!("loopback listener bind failed: {error}"),
};
let endpoint = format!("ws://{}/", listener.local_addr());
let pair = tokio::time::timeout(TEST_TIMEOUT, async {
return tokio::join!(listener.accept(), tokio_tungstenite::connect_async(endpoint.as_str()));
})
.await;
return match pair {
Ok((Ok(server), Ok((client, _response)))) => (server, client),
Ok((_server, _client)) => panic!("raw-client loopback connection establishment failed"),
Err(_) => panic!("raw-client loopback connection establishment timed out"),
};
}

View File

@@ -0,0 +1,40 @@
// file: crates/common/game-realtime-websocket-lib/unit_tests/config.rs
// version: 1
#[test]
fn default_configuration_matches_the_product_baseline() {
let config = crate::WebSocketConfig::default();
assert_eq!(config.max_message_size(), 1024 * 1024);
assert_eq!(config.max_frame_size(), 1024 * 1024);
assert_eq!(config.write_buffer_size(), 64 * 1024);
assert_eq!(config.max_write_buffer_size(), 2 * 1024 * 1024);
assert_eq!(config.connect_timeout(), std::time::Duration::from_secs(10));
assert_eq!(config.send_timeout(), std::time::Duration::from_secs(5));
assert_eq!(config.close_timeout(), std::time::Duration::from_secs(2));
assert!(config.validate().is_ok());
}
#[test]
fn invalid_message_and_write_buffer_bounds_are_rejected() {
let zero_message = crate::WebSocketConfig::default().with_max_message_size(0);
assert_invalid_configuration(zero_message);
let oversized_frame = crate::WebSocketConfig::default().with_max_frame_size(2 * 1024 * 1024);
assert_invalid_configuration(oversized_frame);
let insufficient_write_buffer = crate::WebSocketConfig::default().with_max_write_buffer_size(1024 * 1024);
assert_invalid_configuration(insufficient_write_buffer);
}
#[test]
fn zero_operation_deadlines_are_rejected() {
assert_invalid_configuration(crate::WebSocketConfig::default().with_connect_timeout(std::time::Duration::ZERO));
assert_invalid_configuration(crate::WebSocketConfig::default().with_send_timeout(std::time::Duration::ZERO));
assert_invalid_configuration(crate::WebSocketConfig::default().with_close_timeout(std::time::Duration::ZERO));
}
fn assert_invalid_configuration(config: crate::WebSocketConfig) {
let result = config.validate();
match result {
Ok(()) => panic!("invalid WebSocket configuration was accepted"),
Err(error) => assert_eq!(error.kind(), game_realtime_transport_lib::TransportErrorKind::InvalidConfiguration),
}
}

View File

@@ -0,0 +1,18 @@
// file: crates/common/game-realtime-websocket-lib/unit_tests/websocket.rs
// version: 1
#[test]
fn tungstenite_backpressure_maps_to_transport_backpressure() {
let message = tokio_tungstenite::tungstenite::Message::Binary(vec![1, 2, 3].into());
let backend_error = tokio_tungstenite::tungstenite::Error::WriteBufferFull(Box::new(message));
let mapped = super::map_stream_error(backend_error);
assert_eq!(mapped.kind(), game_realtime_transport_lib::TransportErrorKind::Backpressure);
}
#[test]
fn tungstenite_capacity_maps_to_message_too_large() {
let capacity = tokio_tungstenite::tungstenite::error::CapacityError::MessageTooLong { size: 65, max_size: 64 };
let backend_error = tokio_tungstenite::tungstenite::Error::Capacity(capacity);
let mapped = super::map_stream_error(backend_error);
assert_eq!(mapped.kind(), game_realtime_transport_lib::TransportErrorKind::MessageTooLarge);
}