0.3.4-alpha.3

This commit is contained in:
2026-09-21 17:40:56 +02:00
parent 4e2be3cc9a
commit 4fb521f54e
9 changed files with 636 additions and 9 deletions

View File

@@ -0,0 +1,24 @@
# file: crates/common/game-realtime-websocket-lib/Cargo.toml
# version: 1
[package]
name = "game-realtime-websocket-lib"
version.workspace = true
edition.workspace = true
license.workspace = true
repository.workspace = true
authors.workspace = true
publish.workspace = true
[dependencies]
futures-util = { workspace = true, default-features = false, features = ["sink", "std"] }
game-realtime-transport-lib = { path = "../game-realtime-transport-lib" }
tokio = { workspace = true, features = ["net"] }
tokio-tungstenite = { workspace = true, default-features = false, features = ["connect", "handshake"] }
tracing.workspace = true
[dev-dependencies]
tokio = { workspace = true, features = ["macros", "rt", "time"] }
[lints]
workspace = true

View File

@@ -0,0 +1,21 @@
// file: crates/common/game-realtime-websocket-lib/src/lib.rs
// version: 1
#![warn(missing_docs)]
#![deny(unreachable_pub)]
#![forbid(unsafe_code)]
//! Tokio/tokio-tungstenite WebSocket backend for the games.sasedev realtime transport contract.
mod websocket;
/// Re-export of the WebSocket client connection constructor.
pub use self::websocket::connect;
/// Re-export of an established WebSocket transport connection.
pub use self::websocket::WebSocketConnection;
/// Re-export of the receive half of an established WebSocket connection.
pub use self::websocket::WebSocketReceiver;
/// Re-export of a bound WebSocket server listener.
pub use self::websocket::WebSocketListener;
/// Re-export of the send half of an established WebSocket connection.
pub use self::websocket::WebSocketSender;

View File

@@ -0,0 +1,300 @@
// file: crates/common/game-realtime-websocket-lib/src/websocket.rs
// version: 1
use futures_util::SinkExt; // rust-rules: trait-import
use futures_util::StreamExt; // rust-rules: trait-import
const TRACING_TARGET: &str = "games::realtime::websocket";
type ClientStream = tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>;
type ServerStream = tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>;
type ClientSink = futures_util::stream::SplitSink<ClientStream, tokio_tungstenite::tungstenite::Message>;
type ServerSink = futures_util::stream::SplitSink<ServerStream, tokio_tungstenite::tungstenite::Message>;
type ClientReceiver = futures_util::stream::SplitStream<ClientStream>;
type ServerReceiver = futures_util::stream::SplitStream<ServerStream>;
enum WebSocketStreamKind {
Client(ClientStream),
Server(ServerStream),
}
enum WebSocketSinkKind {
Client(ClientSink),
Server(ServerSink),
}
enum WebSocketReceiverKind {
Client(ClientReceiver),
Server(ServerReceiver),
}
/// Established WebSocket connection implementing the transport-neutral realtime contract.
pub struct WebSocketConnection {
inner: WebSocketStreamKind,
}
impl WebSocketConnection {
fn from_client(stream: ClientStream) -> Self {
return Self { inner: WebSocketStreamKind::Client(stream) };
}
fn from_server(stream: ServerStream) -> Self {
return Self { inner: WebSocketStreamKind::Server(stream) };
}
}
impl game_realtime_transport_lib::RealtimeConnection for WebSocketConnection {
type Sender = crate::WebSocketSender;
type Receiver = crate::WebSocketReceiver;
fn split(self) -> (Self::Sender, Self::Receiver) {
return match self.inner {
WebSocketStreamKind::Client(stream) => {
let (sender, receiver) = stream.split();
(
crate::WebSocketSender { inner: WebSocketSinkKind::Client(sender) },
crate::WebSocketReceiver { inner: WebSocketReceiverKind::Client(receiver) },
)
},
WebSocketStreamKind::Server(stream) => {
let (sender, receiver) = stream.split();
(
crate::WebSocketSender { inner: WebSocketSinkKind::Server(sender) },
crate::WebSocketReceiver { inner: WebSocketReceiverKind::Server(receiver) },
)
},
};
}
}
/// Send half of an established WebSocket transport connection.
pub struct WebSocketSender {
inner: WebSocketSinkKind,
}
impl game_realtime_transport_lib::RealtimeSender for WebSocketSender {
type SendFuture<'a> = futures_util::future::LocalBoxFuture<'a, Result<(), game_realtime_transport_lib::TransportError>> where Self: 'a;
type CloseFuture<'a> = futures_util::future::LocalBoxFuture<'a, Result<(), game_realtime_transport_lib::TransportError>> where Self: 'a;
fn send(&mut self, message: game_realtime_transport_lib::TransportMessage) -> Self::SendFuture<'_> {
return Box::pin(async move {
let payload_len = message.len();
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,
};
return match result {
Ok(()) => {
tracing::trace!(target: TRACING_TARGET, payload_len = payload_len, "binary WebSocket payload sent");
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)
},
};
});
}
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,
};
return match result {
Ok(()) => {
tracing::debug!(target: TRACING_TARGET, "local WebSocket close initiated");
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)
},
};
});
}
}
/// Receive half of an established WebSocket transport connection.
pub struct WebSocketReceiver {
inner: WebSocketReceiverKind,
}
impl game_realtime_transport_lib::RealtimeReceiver for WebSocketReceiver {
type ReceiveFuture<'a> =
futures_util::future::LocalBoxFuture<'a, Result<game_realtime_transport_lib::TransportReceive, game_realtime_transport_lib::TransportError>>
where
Self: 'a;
fn receive(&mut self) -> Self::ReceiveFuture<'_> {
return Box::pin(async move {
loop {
let next_message = match &mut self.inner {
WebSocketReceiverKind::Client(receiver) => receiver.next().await,
WebSocketReceiverKind::Server(receiver) => receiver.next().await,
};
match next_message {
Some(Ok(tokio_tungstenite::tungstenite::Message::Binary(bytes))) => {
tracing::trace!(target: TRACING_TARGET, payload_len = bytes.len(), "binary WebSocket payload received");
return Ok(game_realtime_transport_lib::TransportReceive::Message(
game_realtime_transport_lib::TransportMessage::new(bytes.to_vec()),
));
},
Some(Ok(tokio_tungstenite::tungstenite::Message::Close(_))) => {
tracing::debug!(target: TRACING_TARGET, "remote WebSocket close observed");
return Ok(game_realtime_transport_lib::TransportReceive::Closed);
},
Some(Ok(tokio_tungstenite::tungstenite::Message::Ping(_)))
| Some(Ok(tokio_tungstenite::tungstenite::Message::Pong(_))) => {},
Some(Ok(tokio_tungstenite::tungstenite::Message::Text(_))) => {
let error = game_realtime_transport_lib::TransportError::new(
game_realtime_transport_lib::TransportErrorKind::Protocol,
"text WebSocket messages are not part of the binary transport contract",
);
tracing::warn!(target: TRACING_TARGET, kind = %error.kind(), "unsupported WebSocket text message received");
return Err(error);
},
Some(Ok(tokio_tungstenite::tungstenite::Message::Frame(_))) => {
let error = game_realtime_transport_lib::TransportError::new(
game_realtime_transport_lib::TransportErrorKind::Protocol,
"unexpected raw WebSocket frame surfaced by the backend",
);
tracing::warn!(target: TRACING_TARGET, kind = %error.kind(), "unexpected raw WebSocket frame received");
return Err(error);
},
Some(Err(error)) => {
let mapped = map_stream_error(error);
tracing::warn!(target: TRACING_TARGET, kind = %mapped.kind(), detail = mapped.detail(), "WebSocket receive failed");
return Err(mapped);
},
None => {
tracing::debug!(target: TRACING_TARGET, "WebSocket stream ended");
return Ok(game_realtime_transport_lib::TransportReceive::Closed);
},
}
}
});
}
}
/// Bound TCP listener that upgrades accepted peers to WebSocket connections.
pub struct WebSocketListener {
listener: tokio::net::TcpListener,
local_addr: std::net::SocketAddr,
}
impl WebSocketListener {
/// Binds a WebSocket listener to one concrete socket address.
pub async fn bind(address: std::net::SocketAddr) -> Result<Self, game_realtime_transport_lib::TransportError> {
let listener = match tokio::net::TcpListener::bind(address).await {
Ok(value) => value,
Err(error) => {
let mapped = game_realtime_transport_lib::TransportError::new(
game_realtime_transport_lib::TransportErrorKind::Bind,
error.to_string(),
);
tracing::warn!(target: TRACING_TARGET, address = %address, detail = mapped.detail(), "WebSocket listener bind failed");
return Err(mapped);
},
};
let local_addr = match listener.local_addr() {
Ok(value) => value,
Err(error) => {
let mapped = game_realtime_transport_lib::TransportError::new(
game_realtime_transport_lib::TransportErrorKind::Bind,
error.to_string(),
);
tracing::warn!(target: TRACING_TARGET, detail = mapped.detail(), "bound WebSocket listener address lookup failed");
return Err(mapped);
},
};
tracing::info!(target: TRACING_TARGET, address = %local_addr, "WebSocket listener bound");
return Ok(Self { listener, local_addr });
}
/// Returns the concrete local socket address, including an ephemeral port selected by the OS.
#[must_use]
pub fn local_addr(&self) -> std::net::SocketAddr {
return self.local_addr;
}
/// Accepts one TCP peer and completes the 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,
Err(error) => {
let mapped = game_realtime_transport_lib::TransportError::new(
game_realtime_transport_lib::TransportErrorKind::Accept,
error.to_string(),
);
tracing::warn!(target: TRACING_TARGET, detail = mapped.detail(), "WebSocket TCP accept failed");
return Err(mapped);
},
};
let websocket = match tokio_tungstenite::accept_async(stream).await {
Ok(value) => value,
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);
},
};
tracing::info!(target: TRACING_TARGET, peer = %peer_addr, "WebSocket peer accepted");
return Ok(crate::WebSocketConnection::from_server(websocket));
}
}
/// Connects a client to one plain `ws://` endpoint and completes the WebSocket handshake.
pub async fn connect(endpoint: &str) -> 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",
));
}
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 mapped = map_connect_error(error);
tracing::warn!(
target: TRACING_TARGET,
endpoint = endpoint,
kind = %mapped.kind(),
detail = mapped.detail(),
"WebSocket client connect failed"
);
return Err(mapped);
},
};
tracing::info!(target: TRACING_TARGET, endpoint = endpoint, "WebSocket client connected");
return Ok(crate::WebSocketConnection::from_client(stream));
}
fn map_connect_error(error: tokio_tungstenite::tungstenite::Error) -> game_realtime_transport_lib::TransportError {
let kind = match &error {
tokio_tungstenite::tungstenite::Error::Url(_) => game_realtime_transport_lib::TransportErrorKind::InvalidConfiguration,
_ => game_realtime_transport_lib::TransportErrorKind::Connect,
};
return game_realtime_transport_lib::TransportError::new(kind, error.to_string());
}
fn map_stream_error(error: tokio_tungstenite::tungstenite::Error) -> game_realtime_transport_lib::TransportError {
let kind = match &error {
tokio_tungstenite::tungstenite::Error::ConnectionClosed | tokio_tungstenite::tungstenite::Error::AlreadyClosed => {
game_realtime_transport_lib::TransportErrorKind::Closed
},
tokio_tungstenite::tungstenite::Error::Io(_) => game_realtime_transport_lib::TransportErrorKind::Io,
tokio_tungstenite::tungstenite::Error::Capacity(_) => game_realtime_transport_lib::TransportErrorKind::MessageTooLarge,
tokio_tungstenite::tungstenite::Error::WriteBufferFull(_) => game_realtime_transport_lib::TransportErrorKind::Backpressure,
_ => game_realtime_transport_lib::TransportErrorKind::Protocol,
};
return game_realtime_transport_lib::TransportError::new(kind, error.to_string());
}

View File

@@ -0,0 +1,63 @@
// file: crates/common/game-realtime-websocket-lib/tests/loopback.rs
// version: 1
//! Deterministic localhost proof for the Tokio/tokio-tungstenite backend.
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 TEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(3);
#[tokio::test(flavor = "current_thread")]
async fn binary_round_trip_and_clean_close_work_on_loopback() {
let bind_address = std::net::SocketAddr::from(([127, 0, 0, 1], 0));
let listener = match game_realtime_websocket_lib::WebSocketListener::bind(bind_address).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(endpoint.as_str()));
})
.await;
let (server_connection, client_connection) = match pair {
Ok((Ok(server), Ok(client))) => (server, client),
Ok((_server, _client)) => panic!("loopback connection establishment failed"),
Err(_) => panic!("loopback connection establishment timed out"),
};
let (mut server_sender, mut server_receiver) = server_connection.split();
let (mut client_sender, mut client_receiver) = client_connection.split();
let client_payload = game_realtime_transport_lib::TransportMessage::new(vec![0, 1, 2, 3, 255]);
let client_send = tokio::time::timeout(TEST_TIMEOUT, client_sender.send(client_payload)).await;
assert!(matches!(client_send, Ok(Ok(()))));
let server_receive = tokio::time::timeout(TEST_TIMEOUT, server_receiver.receive()).await;
match server_receive {
Ok(Ok(received)) => assert_eq!(
received,
game_realtime_transport_lib::TransportReceive::Message(game_realtime_transport_lib::TransportMessage::new(vec![0, 1, 2, 3, 255]))
),
Ok(Err(error)) => panic!("server receive failed: {error}"),
Err(_) => panic!("server receive timed out"),
}
let server_payload = game_realtime_transport_lib::TransportMessage::new(vec![9, 8, 7, 6]);
let server_send = tokio::time::timeout(TEST_TIMEOUT, server_sender.send(server_payload)).await;
assert!(matches!(server_send, Ok(Ok(()))));
let client_receive = tokio::time::timeout(TEST_TIMEOUT, client_receiver.receive()).await;
match client_receive {
Ok(Ok(received)) => assert_eq!(
received,
game_realtime_transport_lib::TransportReceive::Message(game_realtime_transport_lib::TransportMessage::new(vec![9, 8, 7, 6]))
),
Ok(Err(error)) => panic!("client receive failed: {error}"),
Err(_) => panic!("client receive timed out"),
}
let client_close = tokio::time::timeout(TEST_TIMEOUT, client_sender.close()).await;
assert!(matches!(client_close, Ok(Ok(()))));
let server_close_observation = tokio::time::timeout(TEST_TIMEOUT, server_receiver.receive()).await;
match server_close_observation {
Ok(Ok(received)) => assert_eq!(received, game_realtime_transport_lib::TransportReceive::Closed),
Ok(Err(error)) => panic!("server close observation failed: {error}"),
Err(_) => panic!("server close observation timed out"),
};
}