diff --git a/Cargo.toml b/Cargo.toml index 5860b19d45..3f3e47d8cc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -69,13 +69,12 @@ rustls-pki-types = "1" rustls-platform-verifier = "0.5" serde = { version = "1", default-features = false, features = ["derive"] } serde_json = { version = "1.0.142", default-features = false, features = ["alloc", "raw_value"] } -soketto = "0.8.1" +yawc = { version = "0.3.3", default-features = false, features = ["rustls-ring"] } syn = { version = "2", default-features = false } thiserror = "2" tokio = "1.42" tokio-rustls = { version = "0.26", default-features = false } tokio-stream = "0.1.7" -tokio-util = "0.7" tower = "0.5" tower-http = "0.6" tracing = "0.1.34" diff --git a/client/transport/Cargo.toml b/client/transport/Cargo.toml index c641442d39..77b66eb505 100644 --- a/client/transport/Cargo.toml +++ b/client/transport/Cargo.toml @@ -24,7 +24,6 @@ thiserror = { workspace = true, optional = true } futures-util = { workspace = true, features = ["alloc"], optional = true } http = { workspace = true, optional = true } tracing = { workspace = true, optional = true } -tokio-util = { workspace = true, features = ["compat"], optional = true } tokio = { workspace = true, features = ["net", "time", "macros"], optional = true } pin-project = { workspace = true, optional = true } url = { workspace = true, optional = true } @@ -37,7 +36,7 @@ rustls-platform-verifier = { workspace = true, optional = true } rustls = { workspace = true, default-features = false, optional = true } # ws -soketto = { workspace = true, optional = true } +yawc = { workspace = true, optional = true } # web-sys [target.'cfg(target_arch = "wasm32")'.dependencies] @@ -53,8 +52,7 @@ ws = [ "futures-util", "http", "tokio", - "tokio-util", - "soketto", + "yawc", "pin-project", "thiserror", "tracing", diff --git a/client/transport/src/ws/mod.rs b/client/transport/src/ws/mod.rs index 8e00741c48..a94672f0e3 100644 --- a/client/transport/src/ws/mod.rs +++ b/client/transport/src/ws/mod.rs @@ -31,24 +31,21 @@ use std::net::SocketAddr; use std::time::Duration; use base64::Engine; -use futures_util::io::{BufReader, BufWriter}; +use futures_util::stream::{SplitSink, SplitStream}; +use futures_util::{SinkExt, StreamExt}; use jsonrpsee_core::Cow; use jsonrpsee_core::TEN_MB_SIZE_BYTES; use jsonrpsee_core::client::{ReceivedMessage, TransportReceiverT, TransportSenderT}; -use soketto::connection::CloseReason; -use soketto::connection::Error::Utf8; -use soketto::data::ByteSlice125; -use soketto::handshake::client::{Client as WsHandshakeClient, ServerResponse}; -use soketto::{Data, Incoming, connection}; use thiserror::Error; use tokio::net::TcpStream; -use tokio_util::compat::{Compat, TokioAsyncReadCompatExt}; +use yawc::frame::{Frame, OpCode}; +use yawc::{Options, WebSocket, WebSocketError}; pub use http::{HeaderMap, HeaderValue, Uri, uri::InvalidUri}; -pub use soketto::handshake::client::Header; pub use stream::EitherStream; pub use tokio::io::{AsyncRead, AsyncWrite}; pub use url::Url; +pub use yawc::{DeflateOptions, Options as WsOptions}; const LOG_TARGET: &str = "jsonrpsee-client"; @@ -69,20 +66,29 @@ pub enum CertificateStore { } /// Sending end of WebSocket transport. -#[derive(Debug)] pub struct Sender { - inner: connection::Sender>>, + inner: SplitSink, Frame>, max_request_size: u32, } +impl std::fmt::Debug for Sender { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("Sender").field("max_request_size", &self.max_request_size).finish() + } +} + /// Receiving end of WebSocket transport. -#[derive(Debug)] pub struct Receiver { - inner: connection::Receiver>>, + inner: SplitStream>, +} + +impl std::fmt::Debug for Receiver { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("Receiver").finish() + } } /// Builder for a WebSocket transport [`Sender`] and [`Receiver`] pair. -#[derive(Debug)] pub struct WsTransportClientBuilder { #[cfg(feature = "tls")] /// What certificate store to use @@ -101,6 +107,25 @@ pub struct WsTransportClientBuilder { pub max_redirections: usize, /// TCP no delay. pub tcp_no_delay: bool, + /// Custom WebSocket options (compression, etc). + pub ws_options: Option, +} + +impl std::fmt::Debug for WsTransportClientBuilder { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let mut s = f.debug_struct("WsTransportClientBuilder"); + #[cfg(feature = "tls")] + s.field("certificate_store", &self.certificate_store); + s.field("connection_timeout", &self.connection_timeout) + .field("headers", &self.headers) + .field("max_request_size", &self.max_request_size) + .field("max_response_size", &self.max_response_size) + .field("max_frame_size", &self.max_frame_size) + .field("max_redirections", &self.max_redirections) + .field("tcp_no_delay", &self.tcp_no_delay) + .field("ws_options", &self.ws_options.as_ref().map(|_| "..")) + .finish() + } } impl Default for WsTransportClientBuilder { @@ -115,6 +140,7 @@ impl Default for WsTransportClientBuilder { headers: http::HeaderMap::new(), max_redirections: 5, tcp_no_delay: true, + ws_options: None, } } } @@ -169,6 +195,15 @@ impl WsTransportClientBuilder { self.max_redirections = redirect; self } + + /// Set custom WebSocket options such as compression settings. + /// + /// By default, balanced compression is enabled. The `max_payload_read` setting + /// from these options will be overridden by [`WsTransportClientBuilder::max_response_size`]. + pub fn set_ws_options(mut self, options: Options) -> Self { + self.ws_options = Some(options); + self + } } /// Stream mode, either plain TCP or TLS. @@ -200,7 +235,7 @@ pub enum WsHandshakeError { /// Error in the transport layer. #[error("{0}")] - Transport(#[source] soketto::handshake::Error), + Transport(#[source] WebSocketError), /// Server rejected the handshake. #[error("Connection rejected with status code: {status_code}")] @@ -236,18 +271,18 @@ pub enum WsHandshakeError { pub enum WsError { /// Error in the WebSocket connection. #[error("{0}")] - Connection(#[source] soketto::connection::Error), + Connection(#[source] WebSocketError), /// Message was too large. #[error("The message was too large")] MessageTooLarge, /// Connection was closed. - #[error("Connection was closed: {0:?}")] - Closed(CloseReason), + #[error("Connection was closed")] + Closed, } impl TransportSenderT for Sender where - T: futures_util::io::AsyncRead + futures_util::io::AsyncWrite + Send + Unpin + 'static, + T: tokio::io::AsyncRead + tokio::io::AsyncWrite + Send + Unpin + 'static, { type Error = WsError; @@ -259,9 +294,7 @@ where return Err(WsError::MessageTooLarge); } - self.inner.send_text(body).await?; - self.inner.flush().await?; - Ok(()) + self.inner.send(Frame::text(body)).await.map_err(WsError::Connection) } } @@ -270,42 +303,39 @@ where fn send_ping(&mut self) -> impl Future> + Send { async { tracing::debug!(target: LOG_TARGET, "Send ping"); - // Submit empty slice as "optional" parameter. - let slice: &[u8] = &[]; - // Byte slice fails if the provided slice is larger than 125 bytes. - let byte_slice = ByteSlice125::try_from(slice).expect("Empty slice should fit into ByteSlice125"); - - self.inner.send_ping(byte_slice).await?; - self.inner.flush().await?; - Ok(()) + self.inner.send(Frame::ping(b"" as &[u8])).await.map_err(WsError::Connection) } } /// Send a close message and close the connection. fn close(&mut self) -> impl Future> + Send { - async { self.inner.close().await.map_err(Into::into) } + async { self.inner.close().await.map_err(WsError::Connection) } } } impl TransportReceiverT for Receiver where - T: futures_util::io::AsyncRead + futures_util::io::AsyncWrite + Unpin + Send + 'static, + T: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin + Send + 'static, { type Error = WsError; /// Returns a `Future` resolving when the server sent us something back. fn receive(&mut self) -> impl Future> + Send { async { - let mut message = Vec::new(); - - match self.inner.receive(&mut message).await? { - Incoming::Data(Data::Text(_)) => { - let s = String::from_utf8(message).map_err(|err| WsError::Connection(Utf8(err.utf8_error())))?; - Ok(ReceivedMessage::Text(s)) - } - Incoming::Data(Data::Binary(_)) => Ok(ReceivedMessage::Bytes(message)), - Incoming::Pong(_) => Ok(ReceivedMessage::Pong), - Incoming::Closed(c) => Err(WsError::Closed(c)), + match self.inner.next().await { + Some(frame) => match frame.opcode() { + OpCode::Text => { + let payload = frame.into_payload(); + let s = String::from_utf8(payload.to_vec()) + .map_err(|_| WsError::Connection(WebSocketError::InvalidUTF8))?; + Ok(ReceivedMessage::Text(s)) + } + OpCode::Binary => Ok(ReceivedMessage::Bytes(frame.into_payload().to_vec())), + OpCode::Pong => Ok(ReceivedMessage::Pong), + OpCode::Close => Err(WsError::Closed), + _ => Ok(ReceivedMessage::Pong), // treat other frames as activity + }, + None => Err(WsError::Closed), } } } @@ -315,10 +345,7 @@ impl WsTransportClientBuilder { /// Try to establish the connection. /// /// Uses the default connection over TCP. - pub async fn build( - self, - uri: Url, - ) -> Result<(Sender>, Receiver>), WsHandshakeError> { + pub async fn build(self, uri: Url) -> Result<(Sender, Receiver), WsHandshakeError> { self.try_connect_over_tcp(uri).await } @@ -327,12 +354,12 @@ impl WsTransportClientBuilder { self, uri: Url, data_stream: T, - ) -> Result<(Sender>, Receiver>), WsHandshakeError> + ) -> Result<(Sender, Receiver), WsHandshakeError> where - T: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin, + T: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin + Send + 'static, { let target: Target = uri.try_into()?; - self.try_connect(&target, data_stream.compat()).await + self.try_connect(&target, data_stream).await } #[cfg(feature = "tls")] @@ -354,7 +381,7 @@ impl WsTransportClientBuilder { async fn try_connect_over_tcp( &self, uri: Url, - ) -> Result<(Sender>, Receiver>), WsHandshakeError> { + ) -> Result<(Sender, Receiver), WsHandshakeError> { let mut target: Target = uri.clone().try_into()?; let mut err = None; @@ -399,7 +426,7 @@ impl WsTransportClientBuilder { } }; - match self.try_connect(&target, tcp_stream.compat()).await { + match self.try_connect(&target, tcp_stream).await { Ok(result) => return Ok(result), Err(WsHandshakeError::Redirected { status_code, location }) => { @@ -475,57 +502,49 @@ impl WsTransportClientBuilder { data_stream: T, ) -> Result<(Sender, Receiver), WsHandshakeError> where - T: futures_util::AsyncRead + futures_util::AsyncWrite + Unpin, + T: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin + Send + 'static, { - let mut client = WsHandshakeClient::new( - BufReader::new(BufWriter::new(data_stream)), - &target.host_header, - &target.path_and_query, - ); - - let headers: Vec<_> = match &target.basic_auth { - Some(basic_auth) if !self.headers.contains_key(http::header::AUTHORIZATION) => { - let it1 = - self.headers.iter().map(|(key, value)| Header { name: key.as_str(), value: value.as_bytes() }); - let it2 = std::iter::once(Header { - name: http::header::AUTHORIZATION.as_str(), - value: basic_auth.as_bytes(), - }); - - it1.chain(it2).collect() - } - _ => { - self.headers.iter().map(|(key, value)| Header { name: key.as_str(), value: value.as_bytes() }).collect() + let options = self + .ws_options + .clone() + .unwrap_or_else(|| Options::default().with_balanced_compression()) + .with_max_payload_read(self.max_response_size as usize); + + let mut builder = http::Request::builder(); + for (key, value) in &self.headers { + builder = builder.header(key, value); + } + + if let Some(basic_auth) = &target.basic_auth { + if !self.headers.contains_key(http::header::AUTHORIZATION) { + builder = builder.header(http::header::AUTHORIZATION, basic_auth); } - }; + } - client.set_headers(&headers); + // Build the URL from target components + let scheme = match target._mode { + Mode::Plain => "ws", + Mode::Tls => "wss", + }; + let ws_url: url::Url = format!("{scheme}://{}{}", target.host_header, target.path_and_query) + .parse() + .map_err(|e: url::ParseError| WsHandshakeError::Url(e.to_string().into()))?; - // Perform the initial handshake. - match client.handshake().await { - Ok(ServerResponse::Accepted { .. }) => { + match WebSocket::handshake_with_request(ws_url, data_stream, options, builder).await { + Ok(ws) => { tracing::debug!(target: LOG_TARGET, "Connection established to target: {:?}", target); - let mut builder = client.into_builder(); - builder.set_max_message_size(self.max_response_size as usize); - // Use the max frame size if any, otherwise let the underlying code use appropriate defaults. - if let Some(max_frame_size) = self.max_frame_size { - builder.set_max_frame_size(max_frame_size as usize); - } - let (sender, receiver) = builder.finish(); - Ok((Sender { inner: sender, max_request_size: self.max_request_size }, Receiver { inner: receiver })) + let (sink, stream) = ws.split(); + Ok((Sender { inner: sink, max_request_size: self.max_request_size }, Receiver { inner: stream })) } - - Ok(ServerResponse::Rejected { status_code }) => { - tracing::debug!(target: LOG_TARGET, "Connection rejected: {:?}", status_code); - Err(WsHandshakeError::Rejected { status_code }) - } - - Ok(ServerResponse::Redirect { status_code, location }) => { + Err(WebSocketError::Redirected { status_code, location }) => { tracing::debug!(target: LOG_TARGET, "Redirection: status_code: {}, location: {}", status_code, location); Err(WsHandshakeError::Redirected { status_code, location }) } - - Err(e) => Err(e.into()), + Err(WebSocketError::InvalidStatusCode(status_code)) => { + tracing::debug!(target: LOG_TARGET, "Connection rejected: {:?}", status_code); + Err(WsHandshakeError::Rejected { status_code }) + } + Err(e) => Err(WsHandshakeError::Transport(e)), } } } @@ -581,14 +600,14 @@ impl From for WsHandshakeError { } } -impl From for WsHandshakeError { - fn from(err: soketto::handshake::Error) -> WsHandshakeError { +impl From for WsHandshakeError { + fn from(err: WebSocketError) -> WsHandshakeError { WsHandshakeError::Transport(err) } } -impl From for WsError { - fn from(err: soketto::connection::Error) -> Self { +impl From for WsError { + fn from(err: WebSocketError) -> Self { WsError::Connection(err) } } diff --git a/client/ws-client/src/lib.rs b/client/ws-client/src/lib.rs index 7899a28deb..74315c86b1 100644 --- a/client/ws-client/src/lib.rs +++ b/client/ws-client/src/lib.rs @@ -47,6 +47,7 @@ use jsonrpsee_core::middleware::layer::RpcLoggerLayer; pub use jsonrpsee_types as types; use jsonrpsee_client_transport::ws::{AsyncRead, AsyncWrite, WsTransportClientBuilder}; +pub use jsonrpsee_client_transport::ws::{DeflateOptions, WsOptions}; use jsonrpsee_core::TEN_MB_SIZE_BYTES; use jsonrpsee_core::client::{ClientBuilder, Error, IdKind, MaybeSend, TransportReceiverT, TransportSenderT}; use std::time::Duration; @@ -85,7 +86,7 @@ use jsonrpsee_client_transport::ws::CertificateStore; /// } /// /// ``` -#[derive(Clone, Debug)] +#[derive(Clone)] pub struct WsClientBuilder { #[cfg(feature = "tls")] certificate_store: CertificateStore, @@ -101,9 +102,33 @@ pub struct WsClientBuilder { max_redirections: usize, id_kind: IdKind, tcp_no_delay: bool, + ws_options: Option, service_builder: RpcServiceBuilder, } +impl std::fmt::Debug for WsClientBuilder { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let mut s = f.debug_struct("WsClientBuilder"); + #[cfg(feature = "tls")] + s.field("certificate_store", &self.certificate_store); + s.field("max_request_size", &self.max_request_size) + .field("max_response_size", &self.max_response_size) + .field("max_frame_size", &self.max_frame_size) + .field("request_timeout", &self.request_timeout) + .field("connection_timeout", &self.connection_timeout) + .field("ping_config", &self.ping_config) + .field("headers", &self.headers) + .field("max_concurrent_requests", &self.max_concurrent_requests) + .field("max_buffer_capacity_per_subscription", &self.max_buffer_capacity_per_subscription) + .field("max_redirections", &self.max_redirections) + .field("id_kind", &self.id_kind) + .field("tcp_no_delay", &self.tcp_no_delay) + .field("ws_options", &self.ws_options.as_ref().map(|_| "..")) + .field("service_builder", &self.service_builder) + .finish() + } +} + impl Default for WsClientBuilder { fn default() -> Self { Self { @@ -121,6 +146,7 @@ impl Default for WsClientBuilder { max_redirections: 5, id_kind: IdKind::Number, tcp_no_delay: true, + ws_options: None, service_builder: RpcServiceBuilder::default().rpc_logger(1024), } } @@ -280,6 +306,12 @@ impl WsClientBuilder { self } + /// See documentation [`WsTransportClientBuilder::set_ws_options`] (default is balanced compression). + pub fn set_ws_options(mut self, options: WsOptions) -> Self { + self.ws_options = Some(options); + self + } + /// Set the RPC service builder. pub fn set_rpc_middleware(self, service_builder: RpcServiceBuilder) -> WsClientBuilder { WsClientBuilder { @@ -297,6 +329,7 @@ impl WsClientBuilder { max_redirections: self.max_redirections, id_kind: self.id_kind, tcp_no_delay: self.tcp_no_delay, + ws_options: self.ws_options, service_builder, } } @@ -358,6 +391,7 @@ impl WsClientBuilder { max_frame_size: self.max_frame_size, max_redirections: self.max_redirections, tcp_no_delay: self.tcp_no_delay, + ws_options: self.ws_options.clone(), }; let uri = Url::parse(url.as_ref()).map_err(|e| Error::Transport(e.into()))?; @@ -388,6 +422,7 @@ impl WsClientBuilder { max_frame_size: self.max_frame_size, max_redirections: self.max_redirections, tcp_no_delay: self.tcp_no_delay, + ws_options: self.ws_options.clone(), }; let uri = Url::parse(url.as_ref()).map_err(|e| Error::Transport(e.into()))?; diff --git a/server/Cargo.toml b/server/Cargo.toml index 97df82d42f..19fc78e914 100644 --- a/server/Cargo.toml +++ b/server/Cargo.toml @@ -17,7 +17,7 @@ publish = true workspace = true [dependencies] -futures-util = { workspace = true, features = ["io", "async-await-macro"] } +futures-util = { workspace = true, features = ["async-await-macro"] } http = { workspace = true } http-body = { workspace = true } http-body-util = { workspace = true } @@ -29,10 +29,9 @@ pin-project = "1.1.3" route-recognizer = "0.3.1" serde = "1" serde_json = { version = "1", features = ["raw_value"] } -soketto = { version = "0.8.1", features = ["http"] } +yawc = { workspace = true } thiserror = "2" tokio = { version = "1.23.1", features = ["net", "rt-multi-thread", "macros", "time"] } -tokio-util = { version = "0.7", features = ["compat"] } tokio-stream = { version = "0.1.7", features = ["sync"] } tower = { workspace = true, features = ["util"] } tracing = { workspace = true } diff --git a/server/src/lib.rs b/server/src/lib.rs index bfc1478469..03bfe8c037 100644 --- a/server/src/lib.rs +++ b/server/src/lib.rs @@ -51,6 +51,7 @@ pub use server::{ ServerConfigBuilder, TowerService, TowerServiceBuilder, TowerServiceNoHttp, }; pub use tracing; +pub use yawc::{DeflateOptions, Options as WsOptions}; pub use jsonrpsee_core::http_helpers::{Body as HttpBody, Request as HttpRequest, Response as HttpResponse}; pub use transport::http; diff --git a/server/src/server.rs b/server/src/server.rs index 87d2316149..fb375eae2d 100644 --- a/server/src/server.rs +++ b/server/src/server.rs @@ -39,8 +39,8 @@ use crate::transport::{http, ws}; use crate::utils::deserialize_with_ext; use crate::{Extensions, HttpBody, HttpRequest, HttpResponse, LOG_TARGET}; +use futures_util::StreamExt; use futures_util::future::{self, Either, FutureExt}; -use futures_util::io::{BufReader, BufWriter}; use hyper::body::Bytes; use hyper_util::rt::{TokioExecutor, TokioIo}; use jsonrpsee_core::id_providers::RandomIntegerIdProvider; @@ -53,10 +53,8 @@ use jsonrpsee_types::error::{ BATCHES_NOT_SUPPORTED_CODE, BATCHES_NOT_SUPPORTED_MSG, ErrorCode, reject_too_big_batch_request, }; use jsonrpsee_types::{ErrorObject, Id}; -use soketto::handshake::http::is_upgrade_request; use tokio::net::{TcpListener, TcpStream, ToSocketAddrs}; use tokio::sync::{OwnedSemaphorePermit, mpsc, watch}; -use tokio_util::compat::TokioAsyncReadCompatExt; use tower::layer::util::Identity; use tower::{Layer, Service}; use tracing::{Instrument, instrument}; @@ -170,7 +168,7 @@ where } /// Static server configuration which is shared per connection. -#[derive(Debug, Clone)] +#[derive(Clone)] pub struct ServerConfig { /// Maximum size in bytes of a request. pub(crate) max_request_body_size: u32, @@ -200,10 +198,34 @@ pub struct ServerConfig { pub(crate) keep_alive: Option, /// `KEEP_ALIVE_TIMEOUT` duration. pub(crate) keep_alive_timeout: Duration, + /// Custom WebSocket options (compression, etc). + pub(crate) ws_options: Option, +} + +impl std::fmt::Debug for ServerConfig { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("ServerConfig") + .field("max_request_body_size", &self.max_request_body_size) + .field("max_response_body_size", &self.max_response_body_size) + .field("max_connections", &self.max_connections) + .field("max_subscriptions_per_connection", &self.max_subscriptions_per_connection) + .field("batch_requests_config", &self.batch_requests_config) + .field("tokio_runtime", &self.tokio_runtime) + .field("enable_http", &self.enable_http) + .field("enable_ws", &self.enable_ws) + .field("message_buffer_capacity", &self.message_buffer_capacity) + .field("ping_config", &self.ping_config) + .field("id_provider", &self.id_provider) + .field("tcp_no_delay", &self.tcp_no_delay) + .field("keep_alive", &self.keep_alive) + .field("keep_alive_timeout", &self.keep_alive_timeout) + .field("ws_options", &self.ws_options.as_ref().map(|_| "..")) + .finish() + } } /// The builder to configure and create a JSON-RPC server configuration. -#[derive(Debug, Clone)] +#[derive(Clone)] pub struct ServerConfigBuilder { /// Maximum size in bytes of a request. max_request_body_size: u32, @@ -233,6 +255,30 @@ pub struct ServerConfigBuilder { keep_alive: Option, /// `KEEP_ALIVE_TIMEOUT` duration. keep_alive_timeout: std::time::Duration, + /// Custom WebSocket options (compression, etc). + ws_options: Option, +} + +impl std::fmt::Debug for ServerConfigBuilder { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("ServerConfigBuilder") + .field("max_request_body_size", &self.max_request_body_size) + .field("max_response_body_size", &self.max_response_body_size) + .field("max_connections", &self.max_connections) + .field("max_subscriptions_per_connection", &self.max_subscriptions_per_connection) + .field("batch_requests_config", &self.batch_requests_config) + .field("tokio_runtime", &self.tokio_runtime) + .field("enable_http", &self.enable_http) + .field("enable_ws", &self.enable_ws) + .field("message_buffer_capacity", &self.message_buffer_capacity) + .field("ping_config", &self.ping_config) + .field("id_provider", &self.id_provider) + .field("tcp_no_delay", &self.tcp_no_delay) + .field("keep_alive", &self.keep_alive) + .field("keep_alive_timeout", &self.keep_alive_timeout) + .field("ws_options", &self.ws_options.as_ref().map(|_| "..")) + .finish() + } } /// Builder for [`TowerService`]. @@ -374,6 +420,7 @@ impl Default for ServerConfigBuilder { keep_alive: None, //same as `hyper` default keep_alive_timeout: Duration::from_secs(20), + ws_options: None, } } } @@ -541,6 +588,29 @@ impl ServerConfigBuilder { self } + /// Set custom WebSocket options such as compression settings. + /// + /// By default, balanced compression is enabled. The `max_payload_read` setting + /// from these options will be overridden by [`ServerConfigBuilder::max_request_body_size`]. + /// + /// # Examples + /// + /// ```rust + /// use jsonrpsee_server::{ServerConfigBuilder, WsOptions}; + /// + /// // Disable compression + /// let ws_opts = WsOptions::default().without_compression(); + /// let builder = ServerConfigBuilder::default().set_ws_options(ws_opts); + /// + /// // High compression + /// let ws_opts = WsOptions::default().with_high_compression(); + /// let builder = ServerConfigBuilder::default().set_ws_options(ws_opts); + /// ``` + pub fn set_ws_options(mut self, options: yawc::Options) -> Self { + self.ws_options = Some(options); + self + } + /// Build the [`ServerConfig`]. pub fn build(self) -> ServerConfig { ServerConfig { @@ -558,6 +628,7 @@ impl ServerConfigBuilder { tcp_no_delay: self.tcp_no_delay, keep_alive: self.keep_alive, keep_alive_timeout: self.keep_alive_timeout, + ws_options: self.ws_options, } } } @@ -1039,15 +1110,20 @@ where req_ext.insert::(conn_guard.clone()); req_ext.insert::(conn.conn_id.into()); - let is_upgrade_request = is_upgrade_request(&request); + let is_upgrade_request = ws::is_upgrade_request(&request); if self.inner.server_cfg.enable_ws && is_upgrade_request { let this = self.inner.clone(); - let mut server = soketto::handshake::http::Server::new(); + let options = this + .server_cfg + .ws_options + .clone() + .unwrap_or_else(|| yawc::Options::default().with_balanced_compression()) + .with_max_payload_read(this.server_cfg.max_request_body_size as usize); - let response = match server.receive_request(&request) { - Ok(response) => { + let response = match yawc::WebSocket::upgrade_with_options(&mut request, options) { + Ok((response, upgrade_fut)) => { let (tx, rx) = mpsc::channel(this.server_cfg.message_buffer_capacity as usize); let sink = MethodSink::new(tx); @@ -1078,20 +1154,15 @@ where async move { let extensions = request.extensions().clone(); - let upgraded = match hyper::upgrade::on(request).await { - Ok(u) => u, + let ws = match upgrade_fut.await { + Ok(ws) => ws, Err(e) => { tracing::debug!(target: LOG_TARGET, "Could not upgrade connection: {}", e); return; } }; - let io = TokioIo::new(upgraded); - - let stream = BufReader::new(BufWriter::new(io.compat())); - let mut ws_builder = server.into_builder(stream); - ws_builder.set_max_message_size(this.server_cfg.max_request_body_size as usize); - let (sender, receiver) = ws_builder.finish(); + let (sender, receiver) = ws.split(); let params = BackgroundTaskParams { server_cfg: this.server_cfg, @@ -1111,7 +1182,7 @@ where .in_current_span(), ); - response.map(|()| HttpBody::empty()) + response.map(|_| HttpBody::empty()) } Err(e) => { tracing::debug!(target: LOG_TARGET, "Could not upgrade connection: {}", e); diff --git a/server/src/tests/ws.rs b/server/src/tests/ws.rs index d30d23f161..3847650bf5 100644 --- a/server/src/tests/ws.rs +++ b/server/src/tests/ws.rs @@ -56,14 +56,15 @@ async fn can_set_the_max_request_body_size() { let mut client = WebSocketTestClient::new(addr).await.unwrap(); - // Invalid: too long + // Invalid: too long — yawc terminates the connection when the payload exceeds the limit let req = format!(r#"{{"jsonrpc":"2.0","method":"{}","id":1}}"#, "a".repeat(100)); - let response = client.send_request_text(req).await.unwrap(); - assert_eq!(response, oversized_request(100)); + assert!(client.send_request_text(req).await.is_err()); - // Max request body size should not override the max response body size + // After an oversized request, the connection is closed; verify with a new connection + // that normal-sized requests still work fine. + let mut client2 = WebSocketTestClient::new(addr).await.unwrap(); let req = r#"{"jsonrpc":"2.0","method":"anything","id":1}"#; - let response = client.send_request_text(req).await.unwrap(); + let response = client2.send_request_text(req).await.unwrap(); assert_eq!(response, ok_response(JsonValue::String("a".repeat(100)), Id::Num(1))); handle.stop().unwrap(); diff --git a/server/src/transport/ws.rs b/server/src/transport/ws.rs index 87241cd6a3..6ab8176b73 100644 --- a/server/src/transport/ws.rs +++ b/server/src/transport/ws.rs @@ -7,45 +7,50 @@ use crate::server::{ConnectionState, ServerConfig, handle_rpc_call}; use crate::{HttpBody, HttpRequest, HttpResponse, LOG_TARGET, PingConfig}; use futures_util::future::{self, Either}; -use futures_util::io::{BufReader, BufWriter}; -use futures_util::{Future, StreamExt, TryStreamExt}; -use hyper::upgrade::Upgraded; -use hyper_util::rt::TokioIo; +use futures_util::{Future, SinkExt, StreamExt}; use jsonrpsee_core::middleware::{RpcServiceBuilder, RpcServiceT}; use jsonrpsee_core::server::{BoundedSubscriptions, MethodResponse, MethodSink, Methods}; use jsonrpsee_types::Id; -use jsonrpsee_types::error::{ErrorCode, reject_too_big_request}; +use jsonrpsee_types::error::ErrorCode; use serde_json::value::RawValue; -use soketto::connection::Error as SokettoError; -use soketto::data::ByteSlice125; use tokio::sync::{mpsc, oneshot}; use tokio::time::{interval, interval_at}; use tokio_stream::wrappers::ReceiverStream; -use tokio_util::compat::{Compat, TokioAsyncReadCompatExt}; - -pub(crate) type Sender = soketto::Sender>>>>; -pub(crate) type Receiver = soketto::Receiver>>>>; - -pub use soketto::handshake::http::is_upgrade_request; +use yawc::frame::{Frame, OpCode}; +use yawc::{HttpStream, WebSocket}; + +pub(crate) type Sender = futures_util::stream::SplitSink, Frame>; +pub(crate) type Receiver = futures_util::stream::SplitStream>; + +/// Checks whether the incoming request is a WebSocket upgrade request. +pub fn is_upgrade_request(req: &http::Request) -> bool { + let dominated_upgrade = req + .headers() + .get(http::header::UPGRADE) + .and_then(|v| v.to_str().ok()) + .map(|v| v.eq_ignore_ascii_case("websocket")) + .unwrap_or(false); + let has_connection_upgrade = req + .headers() + .get(http::header::CONNECTION) + .and_then(|v| v.to_str().ok()) + .map(|v| v.to_lowercase().contains("upgrade")) + .unwrap_or(false); + dominated_upgrade && has_connection_upgrade +} enum Incoming { Data(Vec), Pong, } -pub(crate) async fn send_message(sender: &mut Sender, response: Box) -> Result<(), SokettoError> { - sender.send_text_owned(String::from(Box::::from(response))).await?; - sender.flush().await +pub(crate) async fn send_message(sender: &mut Sender, response: Box) -> Result<(), yawc::WebSocketError> { + sender.send(Frame::text(String::from(Box::::from(response)))).await } -pub(crate) async fn send_ping(sender: &mut Sender) -> Result<(), SokettoError> { +pub(crate) async fn send_ping(sender: &mut Sender) -> Result<(), yawc::WebSocketError> { tracing::debug!(target: LOG_TARGET, "Send ping"); - // Submit empty slice as "optional" parameter. - let slice: &[u8] = &[]; - // Byte slice fails if the provided slice is larger than 125 bytes. - let byte_slice = ByteSlice125::try_from(slice).expect("Empty slice should fit into ByteSlice125"); - sender.send_ping(byte_slice).await?; - sender.flush().await + sender.send(Frame::ping(b"" as &[u8])).await } pub(crate) struct BackgroundTaskParams { @@ -83,7 +88,7 @@ where mut on_session_close, extensions, } = params; - let ServerConfig { ping_config, batch_requests_config, max_request_body_size, .. } = server_cfg; + let ServerConfig { ping_config, batch_requests_config, .. } = server_cfg; let (conn_tx, conn_rx) = oneshot::channel(); @@ -96,18 +101,18 @@ where tokio::pin!(stopped); - let ws_stream = futures_util::stream::unfold(ws_receiver, |mut receiver| async { - let mut data = Vec::new(); - match receiver.receive(&mut data).await { - Ok(soketto::Incoming::Data(_)) => Some((Ok(Incoming::Data(data)), receiver)), - Ok(soketto::Incoming::Pong(_)) => Some((Ok(Incoming::Pong), receiver)), - Ok(soketto::Incoming::Closed(_)) | Err(SokettoError::Closed) => None, - // The closing reason is already logged by `soketto` trace log level. - // Return the `Closed` error to avoid logging unnecessary warnings on clean shutdown. - Err(e) => Some((Err(e), receiver)), - } - }) - .fuse(); + // Convert the yawc Stream into Stream> + let ws_stream = ws_receiver + .filter_map(|frame| async move { + match frame.opcode() { + OpCode::Text | OpCode::Binary => Some(Ok(Incoming::Data(frame.into_payload().to_vec()))), + OpCode::Pong => Some(Ok(Incoming::Pong)), + OpCode::Close => None, + OpCode::Ping => None, // auto-ponged by yawc + _ => None, + } + }) + .fuse(); tokio::pin!(ws_stream); @@ -119,31 +124,8 @@ where stopped = stop; data } - Receive::Err(err, stop) => { - stopped = stop; - - match err { - SokettoError::Closed => { - break Ok(Shutdown::ConnectionClosed); - } - SokettoError::MessageTooLarge { current, maximum } => { - tracing::debug!( - target: LOG_TARGET, - "WS recv error: message too large current={}/max={}", - current, - maximum - ); - if sink.send_error(Id::Null, reject_too_big_request(max_request_body_size)).await.is_err() { - break Ok(Shutdown::ConnectionClosed); - } - - continue; - } - err => { - tracing::debug!(target: LOG_TARGET, "WS error: {}; terminate connection: {}", err, conn.conn_id); - break Err(err); - } - }; + Receive::Err(_stop) => { + break Ok(Shutdown::ConnectionClosed); } }; @@ -268,7 +250,7 @@ async fn send_task( enum Receive { ConnectionClosed, Stopped, - Err(SokettoError, S), + Err(S), Ok(Vec, S), } @@ -281,7 +263,7 @@ async fn try_recv( ) -> Receive where S: Future + Unpin, - T: StreamExt> + Unpin, + T: StreamExt> + Unpin, { let mut last_active = Instant::now(); let inactivity_check = match ping_config { @@ -306,7 +288,7 @@ where futs = futures_util::future::select(ws_stream.next(), inactive); } // Received an error, terminate the connection. - Either::Left((Either::Left((Some(Err(e)), _)), s)) => break Receive::Err(e, s), + Either::Left((Either::Left((Some(Err(_)), _)), s)) => break Receive::Err(s), // Max inactivity timeout fired, check if the connection has been idle too long. Either::Left((Either::Right((_instant, rcv)), s)) => { if let Some(p) = ping_config { @@ -343,27 +325,23 @@ pub(crate) enum Shutdown { /// /// This will return once the connection has been terminated or all pending calls have been executed. async fn graceful_shutdown( - result: Result, + result: Result, pending_calls: mpsc::Receiver<()>, ws_stream: S, mut conn_tx: oneshot::Sender<()>, send_task_handle: tokio::task::JoinHandle<()>, ) where - S: StreamExt> + Unpin, + S: StreamExt + Unpin, { let pending_calls = ReceiverStream::new(pending_calls); if let Ok(Shutdown::Stopped) = result { let graceful_shutdown = pending_calls.for_each(|_| async {}); - let disconnect = ws_stream.try_for_each(|_| async { Ok(()) }); + let disconnect = ws_stream.for_each(|_| async {}); tokio::select! { _ = graceful_shutdown => {} - res = disconnect => { - if let Err(err) = res { - tracing::warn!(target: LOG_TARGET, "Graceful shutdown terminated because of error: `{err}`"); - } - } + _ = disconnect => {} _ = conn_tx.closed() => {} } } @@ -435,10 +413,16 @@ where + Sync + 'static, { - let mut server = soketto::handshake::http::Server::new(); + let mut request = req.map(|_| http_body_util::Empty::::new()); - match server.receive_request(&req) { - Ok(response) => { + let options = server_cfg + .ws_options + .clone() + .unwrap_or_else(|| yawc::Options::default().with_balanced_compression()) + .with_max_payload_read(server_cfg.max_request_body_size as usize); + + match WebSocket::upgrade_with_options(&mut request, options) { + Ok((response, upgrade_fut)) => { let (tx, rx) = mpsc::channel(server_cfg.message_buffer_capacity as usize); let sink = MethodSink::new(tx); @@ -466,22 +450,17 @@ where // Note: This can't possibly be fulfilled until the HTTP response // is returned below, so that's why it's a separate async block let fut = async move { - let extensions = req.extensions().clone(); + let extensions = request.extensions().clone(); - let upgraded = match hyper::upgrade::on(req).await { - Ok(upgraded) => upgraded, + let ws = match upgrade_fut.await { + Ok(ws) => ws, Err(e) => { tracing::debug!(target: LOG_TARGET, "WS upgrade handshake failed: {}", e); return; } }; - let io = TokioIo::new(upgraded); - - let stream = BufReader::new(BufWriter::new(io.compat())); - let mut ws_builder = server.into_builder(stream); - ws_builder.set_max_message_size(server_cfg.max_response_body_size as usize); - let (sender, receiver) = ws_builder.finish(); + let (sender, receiver) = futures_util::StreamExt::split(ws); let params = BackgroundTaskParams { server_cfg, @@ -499,7 +478,7 @@ where background_task(params).await; }; - Ok((response.map(|()| HttpBody::default()), fut)) + Ok((response.map(|_| HttpBody::default()), fut)) } Err(e) => { tracing::debug!(target: LOG_TARGET, "WS upgrade handshake failed: {}", e); diff --git a/test-utils/Cargo.toml b/test-utils/Cargo.toml index 0bcfb81093..b0fe00d655 100644 --- a/test-utils/Cargo.toml +++ b/test-utils/Cargo.toml @@ -15,6 +15,6 @@ http-body-util = { workspace = true } tracing = { workspace = true } serde = { workspace = true, features = ["derive", "alloc"] } serde_json = { workspace = true } -soketto = { workspace = true, features = ["http"] } +url = { workspace = true } +yawc = { workspace = true } tokio = { workspace = true, features = ["net", "rt-multi-thread", "macros", "time"] } -tokio-util = { workspace = true, features = ["compat"] } diff --git a/test-utils/src/mocks.rs b/test-utils/src/mocks.rs index 70da9dc014..533b793f81 100644 --- a/test-utils/src/mocks.rs +++ b/test-utils/src/mocks.rs @@ -32,15 +32,14 @@ use std::time::Duration; use futures_channel::mpsc; use futures_channel::oneshot; use futures_util::future::FutureExt; -use futures_util::io::{BufReader, BufWriter}; use futures_util::sink::SinkExt; -use futures_util::stream::{self, StreamExt}; +use futures_util::stream::StreamExt; use futures_util::{pin_mut, select}; use hyper_util::rt::{TokioExecutor, TokioIo}; use serde::{Deserialize, Serialize}; -use soketto::handshake::{self, Error as SokettoError, Server, http::is_upgrade_request, server::Response}; use tokio::net::TcpStream; -use tokio_util::compat::{Compat, TokioAsyncReadCompatExt}; +use yawc::frame::{Frame, OpCode}; +use yawc::{Options, WebSocket, WebSocketError}; pub use hyper::{HeaderMap, StatusCode, Uri}; @@ -69,8 +68,8 @@ pub struct HttpResponse { /// WebSocket client to construct with arbitrary payload to construct bad payloads. pub struct WebSocketTestClient { - tx: soketto::Sender>>>, - rx: soketto::Receiver>>>, + tx: futures_util::stream::SplitSink, Frame>, + rx: futures_util::stream::SplitStream>, } impl std::fmt::Debug for WebSocketTestClient { @@ -83,54 +82,60 @@ impl std::fmt::Debug for WebSocketTestClient { pub enum WebSocketTestError { Redirect, RejectedWithStatusCode(u16), - Soketto(SokettoError), + Yawc(WebSocketError), } impl From for WebSocketTestError { fn from(err: io::Error) -> Self { - WebSocketTestError::Soketto(SokettoError::Io(err)) + WebSocketTestError::Yawc(WebSocketError::IoError(err)) } } impl WebSocketTestClient { pub async fn new(url: SocketAddr) -> Result { let socket = TcpStream::connect(url).await?; - let mut client = handshake::Client::new(BufReader::new(BufWriter::new(socket.compat())), "test-client", "/"); - match client.handshake().await { - Ok(handshake::ServerResponse::Accepted { .. }) => { - let (tx, rx) = client.into_builder().finish(); + let ws_url: url::Url = format!("ws://{url}/").parse().expect("valid url"); + match WebSocket::handshake(ws_url, socket, Options::default()).await { + Ok(ws) => { + let (tx, rx) = ws.split(); Ok(Self { tx, rx }) } - Ok(handshake::ServerResponse::Redirect { .. }) => Err(WebSocketTestError::Redirect), - Ok(handshake::ServerResponse::Rejected { status_code }) => { + Err(WebSocketError::Redirected { .. }) => Err(WebSocketTestError::Redirect), + Err(WebSocketError::InvalidStatusCode(status_code)) => { Err(WebSocketTestError::RejectedWithStatusCode(status_code)) } - Err(err) => Err(WebSocketTestError::Soketto(err)), + Err(err) => Err(WebSocketTestError::Yawc(err)), } } pub async fn send_request_text(&mut self, msg: impl AsRef) -> Result { self.send(msg).await?; - let mut data = Vec::new(); - self.rx.receive_data(&mut data).await?; - String::from_utf8(data).map_err(Into::into) + self.receive().await } pub async fn send(&mut self, msg: impl AsRef) -> Result<(), Error> { - self.tx.send_text(msg).await?; - self.tx.flush().await.map_err(Into::into) + self.tx.send(Frame::text(msg.as_ref().to_string())).await.map_err(Into::into) } pub async fn send_request_binary(&mut self, msg: &[u8]) -> Result { - self.tx.send_binary(msg).await?; - self.tx.flush().await?; + self.tx.send(Frame::binary(msg.to_vec())).await?; self.receive().await } pub async fn receive(&mut self) -> Result { - let mut data = Vec::new(); - self.rx.receive_data(&mut data).await?; - String::from_utf8(data).map_err(Into::into) + loop { + match self.rx.next().await { + Some(frame) => match frame.opcode() { + OpCode::Text | OpCode::Binary => { + return String::from_utf8(frame.into_payload().to_vec()).map_err(Into::into); + } + OpCode::Close => return Err("Connection closed".into()), + // Skip ping/pong frames + _ => continue, + }, + None => return Err("Connection closed".into()), + } + } } pub async fn close(&mut self) -> Result<(), Error> { @@ -241,28 +246,50 @@ async fn server_backend(listener: tokio::net::TcpListener, mut exit: mpsc::Unbou } async fn connection_task(socket: tokio::net::TcpStream, mode: ServerMode, mut exit: mpsc::UnboundedReceiver<()>) { - let mut server = Server::new(socket.compat()); + let io = TokioIo::new(socket); - let key = match server.receive_request().await { - Ok(req) => req.key(), - Err(_) => return, - }; + // Use a oneshot channel to send the established WebSocket to the main task. + let (ws_tx, ws_rx) = tokio::sync::oneshot::channel::>(); + let ws_tx = std::sync::Arc::new(std::sync::Mutex::new(Some(ws_tx))); - let accept = server.send_response(&Response::Accept { key, protocol: None }).await; + let mode_clone = mode.clone(); - if accept.is_err() { - return; - } + let service = hyper::service::service_fn(move |mut req: hyper::Request| { + let ws_tx = ws_tx.lock().unwrap().take().expect("service_fn called only once for HTTP upgrade"); + async move { + let (response, upgrade_fut) = yawc::WebSocket::upgrade(&mut req)?; - let (mut sender, receiver) = server.into_builder().finish(); + tokio::spawn(async move { + let _ = ws_tx.send(upgrade_fut.await); + }); - let ws_stream = stream::unfold(receiver, move |mut receiver| async { - let mut buf = Vec::new(); - let ret = match receiver.receive_data(&mut buf).await { - Ok(_) => Ok(buf), - Err(err) => Err(err), - }; - Some((ret, receiver)) + Ok::<_, yawc::WebSocketError>(response) + } + }); + + // We need to use hyper to handle the HTTP upgrade, then get the WebSocket. + // But since service_fn requires FnMut and we need to move ws_tx out, + // let's use a simpler approach: accept the TCP connection with hyper for the upgrade. + let builder = hyper_util::server::conn::auto::Builder::new(TokioExecutor::new()); + let conn = builder.serve_connection_with_upgrades(io, service); + + // Drive the HTTP connection to complete the upgrade + let _ = conn.await; + + // Now get the WebSocket + let ws = match ws_rx.await { + Ok(Ok(ws)) => ws, + _ => return, + }; + + let (mut sender, receiver) = ws.split(); + + let ws_stream = receiver.filter_map(|frame| async move { + match frame.opcode() { + OpCode::Text | OpCode::Binary => Some(frame.into_payload().to_vec()), + OpCode::Close => None, + _ => None, // skip pings/pongs + } }); pin_mut!(ws_stream); @@ -275,15 +302,15 @@ async fn connection_task(socket: tokio::net::TcpStream, mode: ServerMode, mut ex select! { _ = time_out => { - match &mode { + match &mode_clone { ServerMode::Subscription { subscription_response, .. } => { - if let Err(e) = sender.send_text(&subscription_response).await { + if let Err(e) = sender.send(Frame::text(subscription_response.clone())).await { tracing::warn!("send response to subscription: {:?}", e); break; } }, ServerMode::Notification(n) => { - if let Err(e) = sender.send_text(&n).await { + if let Err(e) = sender.send(Frame::text(n.clone())).await { tracing::warn!("send notification: {:?}", e); break; } @@ -294,16 +321,16 @@ async fn connection_task(socket: tokio::net::TcpStream, mode: ServerMode, mut ex ws = next_ws => { // Got a request on the connection but don't care about the contents. // Just send out the pre-configured hardcoded responses. - if let Some(Ok(_)) = ws { - match &mode { + if let Some(_) = ws { + match &mode_clone { ServerMode::Response(r) => { - if let Err(e) = sender.send_text(&r).await { + if let Err(e) = sender.send(Frame::text(r.clone())).await { tracing::warn!("send response to request error: {:?}", e); break; } }, ServerMode::Subscription { subscription_id, .. } => { - if let Err(e) = sender.send_text(&subscription_id).await { + if let Err(e) = sender.send(Frame::text(subscription_id.clone())).await { tracing::warn!("send subscription id error: {:?}", e); break; } @@ -319,6 +346,23 @@ async fn connection_task(socket: tokio::net::TcpStream, mode: ServerMode, mut ex } } +/// Checks whether the incoming request is a WebSocket upgrade request. +fn is_upgrade_request(req: &hyper::Request) -> bool { + let dominated_upgrade = req + .headers() + .get(hyper::header::UPGRADE) + .and_then(|v| v.to_str().ok()) + .map(|v| v.eq_ignore_ascii_case("websocket")) + .unwrap_or(false); + let has_connection_upgrade = req + .headers() + .get(hyper::header::CONNECTION) + .and_then(|v| v.to_str().ok()) + .map(|v| v.to_lowercase().contains("upgrade")) + .unwrap_or(false); + dominated_upgrade && has_connection_upgrade +} + // Run a WebSocket server running on localhost that redirects requests for testing. // Requests to any url except for `/myblock/two` will redirect one or two times (HTTP 301) and eventually end up in `/myblock/two`. pub async fn ws_server_with_redirect(other_server: String) -> String { diff --git a/tests/Cargo.toml b/tests/Cargo.toml index 02180beee8..2e8bde9676 100644 --- a/tests/Cargo.toml +++ b/tests/Cargo.toml @@ -24,7 +24,6 @@ serde = { workspace = true, features = ["alloc"] } serde_json = { workspace = true } tokio = { workspace = true, features = ["rt-multi-thread", "time"] } tokio-stream = { workspace = true } -tokio-util = { workspace = true, features = ["compat"]} tower = { workspace = true } tower-http = { workspace = true, features = ["cors"] } tracing = { workspace = true }