use futures_util::{SinkExt as _, StreamExt as _};
use tokio::net::TcpStream;
use tokio_tungstenite::{
MaybeTlsStream, WebSocketStream, connect_async_with_config,
tungstenite::{
Error as WsError, Message,
client::IntoClientRequest as _,
http::{HeaderValue, header::SEC_WEBSOCKET_PROTOCOL},
protocol::WebSocketConfig,
},
};
use super::framing::{CONTROL_LINK_SAME, LineBuffer, SCHEME_SEPARATOR, redact};
use super::{EncodedRequest, StreamOpen, Transport, TransportError, TransportProperties};
use crate::protocol::request::{TlcpRequest as _, WsOk, encode_ws_batch};
const STREAM_PATH: &str = "/lightstreamer";
const SUBPROTOCOL: &str = "TLCP-2.5.0.lightstreamer.com";
const MAX_MESSAGE_BYTES: usize = super::framing::MAX_LINE_BYTES;
#[derive(Debug, PartialEq, Eq)]
enum Reception {
Buffered,
Ignored,
PeerClosed {
description: Option<String>,
},
}
fn receive(buffer: &mut LineBuffer, message: Message) -> Result<Reception, TransportError> {
match message {
Message::Text(text) => {
buffer.push_message(text.as_str())?;
Ok(Reception::Buffered)
}
Message::Ping(_) | Message::Pong(_) => Ok(Reception::Ignored),
Message::Close(frame) => Ok(Reception::PeerClosed {
description: frame.map(|frame| frame.to_string()),
}),
Message::Binary(payload) => Err(TransportError::MalformedFrame {
reason: format!(
"binary websocket message of {} bytes; TLCP is text only",
payload.len()
),
}),
Message::Frame(_) => Err(TransportError::MalformedFrame {
reason: "raw websocket frame".to_owned(),
}),
}
}
#[cold]
#[inline(never)]
fn map_ws_error(error: &WsError) -> TransportError {
match error {
WsError::Utf8(reason) => TransportError::MalformedFrame {
reason: format!("invalid UTF-8 in text message: {reason}"),
},
other => TransportError::ConnectionLost {
reason: other.to_string(),
},
}
}
#[must_use]
fn is_fatal(error: &WsError) -> bool {
!matches!(error, WsError::Capacity(_) | WsError::WriteBufferFull(_))
}
#[cold]
#[inline(never)]
fn connect_error(target: &str, reason: impl Into<String>) -> TransportError {
TransportError::Connect {
target: target.to_owned(),
source: Box::new(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
reason.into(),
)),
}
}
fn stream_uri(address: &str, control_link: Option<&str>) -> Result<String, TransportError> {
let address = address.trim();
let redacted = redact(address);
let Some((scheme, rest)) = address.split_once(SCHEME_SEPARATOR) else {
return Err(connect_error(
&redacted,
"server address must start with ws://, wss://, http:// or https://",
));
};
let scheme = match scheme.to_ascii_lowercase().as_str() {
"ws" | "http" => "ws",
"wss" | "https" => "wss",
other => {
return Err(connect_error(
&redacted,
format!("unsupported scheme `{other}`; expected ws, wss, http or https"),
));
}
};
let authority = match control_link {
Some(link) if !link.trim().is_empty() && link.trim() != CONTROL_LINK_SAME => {
let link = link.trim();
link.split_once(SCHEME_SEPARATOR)
.map_or(link, |(_, rest)| rest)
}
_ => rest,
};
let authority = authority.trim_end_matches('/');
if authority.is_empty() {
return Err(connect_error(&redacted, "server address has no host"));
}
if authority.ends_with(STREAM_PATH) {
Ok(format!("{scheme}{SCHEME_SEPARATOR}{authority}"))
} else {
Ok(format!(
"{scheme}{SCHEME_SEPARATOR}{authority}{STREAM_PATH}"
))
}
}
type Socket = WebSocketStream<MaybeTlsStream<TcpStream>>;
#[derive(Debug)]
pub(crate) struct WsTransport {
address: String,
control_link: Option<String>,
socket: Option<Socket>,
buffer: LineBuffer,
truncated: Option<TransportError>,
}
impl WsTransport {
#[must_use = "a transport does nothing until a stream is opened"]
pub(crate) fn try_new(address: impl Into<String>) -> Result<Self, TransportError> {
let address = address.into();
stream_uri(&address, None)?;
Ok(Self {
address,
control_link: None,
socket: None,
buffer: LineBuffer::default(),
truncated: None,
})
}
async fn connect(&mut self) -> Result<(), TransportError> {
let uri = stream_uri(&self.address, self.control_link.as_deref())?;
let target = redact(&uri);
let mut request =
uri.as_str()
.into_client_request()
.map_err(|error| TransportError::Connect {
target: target.clone(),
source: Box::new(error),
})?;
request.headers_mut().insert(
SEC_WEBSOCKET_PROTOCOL,
HeaderValue::from_static(SUBPROTOCOL),
);
tracing::debug!(target = %target, subprotocol = SUBPROTOCOL, "opening TLCP websocket");
let config = WebSocketConfig::default()
.max_message_size(Some(MAX_MESSAGE_BYTES))
.max_frame_size(Some(MAX_MESSAGE_BYTES));
let (socket, response) = connect_async_with_config(request, Some(config), false)
.await
.map_err(|error| TransportError::Connect {
target: target.clone(),
source: Box::new(error),
})?;
match response.headers().get(SEC_WEBSOCKET_PROTOCOL) {
Some(value) if value.as_bytes() == SUBPROTOCOL.as_bytes() => {}
Some(other) => {
return Err(connect_error(
&target,
format!(
"server negotiated subprotocol `{}`, expected `{SUBPROTOCOL}`",
String::from_utf8_lossy(other.as_bytes())
),
));
}
None => {
return Err(connect_error(
&target,
format!("server negotiated no subprotocol, expected `{SUBPROTOCOL}`"),
));
}
}
self.socket = Some(socket);
self.send_message(WsOk::NAME, WsOk::ws_message().to_owned())
.await
}
async fn send_message(
&mut self,
name: &'static str,
message: String,
) -> Result<(), TransportError> {
let bytes = message.len();
let socket = self.socket.as_mut().ok_or_else(|| TransportError::Send {
name,
reason: "no websocket is open".to_owned(),
})?;
match socket.send(Message::text(message)).await {
Ok(()) => {
tracing::trace!(request = name, bytes, "sent TLCP request");
Ok(())
}
Err(error) => {
let reason = error.to_string();
if is_fatal(&error) {
self.socket = None;
}
Err(TransportError::Send { name, reason })
}
}
}
fn abandon_socket(&mut self) {
self.socket = None;
self.buffer.clear();
self.truncated = None;
}
async fn acknowledge_peer_close(&mut self) -> Option<TransportError> {
if let Some(mut socket) = self.socket.take()
&& let Err(error) = socket.close(None).await
{
tracing::debug!(%error, "closing handshake did not complete");
}
self.truncated_line()
}
fn truncated_line(&mut self) -> Option<TransportError> {
let held = self.buffer.discard_partial()?;
tracing::warn!(
bytes = held,
"the connection ended on a line with no CR-LF; discarding it"
);
Some(TransportError::MalformedFrame {
reason: format!("the connection ended mid-line, after {held} bytes with no CR-LF"),
})
}
}
impl Transport for WsTransport {
fn properties(&self) -> TransportProperties {
TransportProperties {
control_shares_stream: true,
ends_on_content_length: false,
is_polling: false,
}
}
fn set_control_link(&mut self, host: Option<&str>) {
self.control_link = match host {
Some(host) if !host.trim().is_empty() && host.trim() != CONTROL_LINK_SAME => {
Some(host.trim().to_owned())
}
_ => None,
};
tracing::debug!(
control_link = self
.control_link
.as_deref()
.unwrap_or("<configured address>"),
"control link recorded; applies to the next websocket established"
);
}
async fn open_stream(&mut self, request: StreamOpen) -> Result<(), TransportError> {
let (name, message) = match &request {
StreamOpen::Create(create) => {
let message = create.ws_message().map_err(|_| TransportError::Send {
name: crate::protocol::request::CreateSession::NAME,
reason: "session creation parameters cannot be encoded".to_owned(),
})?;
(crate::protocol::request::CreateSession::NAME, message)
}
StreamOpen::Bind(bind) => {
let message = bind.ws_message().map_err(|error| TransportError::Send {
name: crate::protocol::request::BindSession::NAME,
reason: error.to_string(),
})?;
(crate::protocol::request::BindSession::NAME, message)
}
};
let reuse = self.socket.is_some() && matches!(request, StreamOpen::Bind(_));
if !reuse {
if let Some(mut socket) = self.socket.take()
&& let Err(error) = socket.close(None).await
{
tracing::debug!(%error, "could not close the previous websocket cleanly");
}
self.buffer.clear();
self.truncated = None;
self.connect().await?;
} else {
tracing::debug!("rebinding on the live websocket");
}
self.send_message(name, message).await
}
async fn next_line(&mut self) -> Option<Result<String, TransportError>> {
loop {
if let Some(line) = self.buffer.pop_line() {
return Some(line);
}
if let Some(error) = self.truncated.take() {
return Some(Err(error));
}
let socket = self.socket.as_mut()?;
match socket.next().await {
Some(Ok(message)) => match receive(&mut self.buffer, message) {
Ok(Reception::Buffered | Reception::Ignored) => {}
Ok(Reception::PeerClosed { description }) => {
tracing::debug!(
close = description.as_deref().unwrap_or("<no close frame>"),
"peer closed the websocket"
);
if let Some(error) = self.acknowledge_peer_close().await {
self.truncated = Some(error);
}
}
Err(error) => {
self.abandon_socket();
return Some(Err(error));
}
},
Some(Err(error)) => {
let error = map_ws_error(&error);
self.abandon_socket();
return Some(Err(error));
}
None => {
self.socket = None;
if let Some(error) = self.truncated_line() {
self.truncated = Some(error);
}
}
}
}
}
async fn send_control(&mut self, request: EncodedRequest) -> Result<(), TransportError> {
let message = encode_ws_batch(request.name, std::slice::from_ref(&request.parameters))
.map_err(|error| TransportError::Send {
name: request.name,
reason: error.to_string(),
})?;
self.send_message(request.name, message).await
}
async fn close(&mut self) -> Result<(), TransportError> {
self.buffer.clear();
self.truncated = None;
let Some(mut socket) = self.socket.take() else {
return Ok(());
};
tracing::debug!("closing the TLCP websocket");
match socket.close(None).await {
Ok(()) => Ok(()),
Err(error) => Err(TransportError::ConnectionLost {
reason: error.to_string(),
}),
}
}
}
#[cfg(test)]
mod tests {
#![allow(clippy::unwrap_used, clippy::expect_used)]
#![allow(clippy::result_large_err)]
use super::*;
use crate::protocol::request::PROTOCOL_VERSION;
use std::net::SocketAddr;
use tokio::net::TcpListener;
use tokio_tungstenite::{
accept_hdr_async,
tungstenite::handshake::server::{Request as ServerRequest, Response as ServerResponse},
};
fn line(buffer: &mut LineBuffer) -> Option<String> {
buffer.pop_line().map(|result| result.expect("valid UTF-8"))
}
#[test]
fn test_receive_text_buffers_its_lines() {
let mut buffer = LineBuffer::default();
let outcome = receive(
&mut buffer,
Message::text("WSOK\r\nCONOK,S1,50000,5000,*\r\n"),
);
assert!(matches!(outcome, Ok(Reception::Buffered)));
assert_eq!(line(&mut buffer).as_deref(), Some("WSOK"));
assert_eq!(line(&mut buffer).as_deref(), Some("CONOK,S1,50000,5000,*"));
}
#[test]
fn test_receive_binary_message_is_malformed_frame() {
let mut buffer = LineBuffer::default();
let outcome = receive(&mut buffer, Message::binary(vec![0x00, 0x01, 0x02]));
assert!(matches!(
outcome,
Err(TransportError::MalformedFrame { .. })
));
}
#[test]
fn test_receive_ping_and_pong_are_ignored() {
let mut buffer = LineBuffer::default();
assert!(matches!(
receive(&mut buffer, Message::Ping(Vec::new().into())),
Ok(Reception::Ignored)
));
assert!(matches!(
receive(&mut buffer, Message::Pong(Vec::new().into())),
Ok(Reception::Ignored)
));
assert!(buffer.pop_line().is_none());
}
#[test]
fn test_receive_close_reports_peer_close() {
let mut buffer = LineBuffer::default();
assert!(matches!(
receive(&mut buffer, Message::Close(None)),
Ok(Reception::PeerClosed { description: None })
));
}
#[test]
fn test_receive_raw_frame_is_malformed_frame() {
use tokio_tungstenite::tungstenite::protocol::frame::Frame;
let mut buffer = LineBuffer::default();
let outcome = receive(&mut buffer, Message::Frame(Frame::close(None)));
assert!(matches!(
outcome,
Err(TransportError::MalformedFrame { .. })
));
}
#[test]
fn test_map_ws_error_reset_without_handshake_is_connection_lost() {
let error = WsError::Protocol(
tokio_tungstenite::tungstenite::error::ProtocolError::ResetWithoutClosingHandshake,
);
assert!(matches!(
map_ws_error(&error),
TransportError::ConnectionLost { .. }
));
assert!(is_fatal(&error));
}
#[test]
fn test_map_ws_error_invalid_utf8_is_malformed_frame() {
let error = WsError::Utf8("invalid utf-8 sequence".to_owned());
assert!(matches!(
map_ws_error(&error),
TransportError::MalformedFrame { .. }
));
}
#[test]
fn test_map_ws_error_io_failure_is_connection_lost() {
let error = WsError::Io(std::io::Error::new(
std::io::ErrorKind::ConnectionReset,
"reset",
));
assert!(matches!(
map_ws_error(&error),
TransportError::ConnectionLost { .. }
));
}
#[test]
fn test_is_fatal_capacity_error_does_not_kill_the_socket() {
let error = WsError::Capacity(
tokio_tungstenite::tungstenite::error::CapacityError::MessageTooLong {
size: 10,
max_size: 5,
},
);
assert!(!is_fatal(&error));
}
#[test]
fn test_subprotocol_carries_the_encoder_protocol_version() {
assert!(
SUBPROTOCOL.starts_with(PROTOCOL_VERSION),
"subprotocol and LS_protocol must name the same TLCP version"
);
assert_eq!(SUBPROTOCOL, "TLCP-2.5.0.lightstreamer.com");
}
#[test]
fn test_stream_uri_appends_the_lightstreamer_path() {
assert_eq!(
stream_uri("wss://push.example.com", None).unwrap(),
"wss://push.example.com/lightstreamer"
);
}
#[test]
fn test_stream_uri_maps_https_to_wss_and_http_to_ws() {
assert_eq!(
stream_uri("https://push.example.com:443", None).unwrap(),
"wss://push.example.com:443/lightstreamer"
);
assert_eq!(
stream_uri("HTTP://push.example.com:8080", None).unwrap(),
"ws://push.example.com:8080/lightstreamer"
);
}
#[test]
fn test_stream_uri_ignores_a_trailing_slash() {
assert_eq!(
stream_uri("wss://push.example.com/", None).unwrap(),
"wss://push.example.com/lightstreamer"
);
}
#[test]
fn test_stream_uri_does_not_duplicate_an_explicit_endpoint_path() {
assert_eq!(
stream_uri("wss://push.example.com/lightstreamer", None).unwrap(),
"wss://push.example.com/lightstreamer"
);
}
#[test]
fn test_stream_uri_without_scheme_is_rejected() {
let error = stream_uri("push.example.com", None).unwrap_err();
assert!(matches!(error, TransportError::Connect { .. }));
}
#[test]
fn test_stream_uri_with_unsupported_scheme_is_rejected() {
let error = stream_uri("ftp://push.example.com", None).unwrap_err();
assert!(matches!(error, TransportError::Connect { .. }));
}
#[test]
fn test_stream_uri_control_link_replaces_the_authority() {
assert_eq!(
stream_uri("wss://push.example.com", Some("node4.example.com:443")).unwrap(),
"wss://node4.example.com:443/lightstreamer"
);
}
#[test]
fn test_stream_uri_control_link_star_keeps_the_configured_address() {
assert_eq!(
stream_uri("wss://push.example.com", Some(CONTROL_LINK_SAME)).unwrap(),
"wss://push.example.com/lightstreamer"
);
}
#[test]
fn test_stream_uri_control_link_cannot_change_the_scheme() {
assert_eq!(
stream_uri("wss://push.example.com", Some("ws://node4.example.com")).unwrap(),
"wss://node4.example.com/lightstreamer"
);
}
#[test]
fn test_properties_declare_ws_behavior() {
let transport = WsTransport::try_new("wss://push.example.com").unwrap();
assert_eq!(
transport.properties(),
TransportProperties {
control_shares_stream: true,
ends_on_content_length: false,
is_polling: false,
}
);
}
#[test]
fn test_try_new_rejects_an_unusable_address() {
assert!(WsTransport::try_new("push.example.com").is_err());
}
#[test]
fn test_set_control_link_records_and_clears() {
let mut transport = WsTransport::try_new("wss://push.example.com").unwrap();
transport.set_control_link(Some("node4.example.com"));
assert_eq!(transport.control_link.as_deref(), Some("node4.example.com"));
transport.set_control_link(Some(CONTROL_LINK_SAME));
assert_eq!(transport.control_link, None);
transport.set_control_link(Some("node4.example.com"));
transport.set_control_link(None);
assert_eq!(transport.control_link, None);
}
#[tokio::test]
async fn test_send_control_without_a_socket_fails_without_panicking() {
let mut transport = WsTransport::try_new("wss://push.example.com").unwrap();
let error = transport
.send_control(EncodedRequest {
name: "control",
path: "/lightstreamer/control.txt",
parameters: "LS_reqId=1&LS_op=delete&LS_subId=1".to_owned(),
})
.await
.unwrap_err();
assert!(matches!(error, TransportError::Send { .. }));
}
#[tokio::test]
async fn test_close_without_a_socket_is_a_no_op() {
let mut transport = WsTransport::try_new("wss://push.example.com").unwrap();
assert!(transport.close().await.is_ok());
}
#[tokio::test]
async fn test_next_line_without_a_socket_is_end_of_stream() {
let mut transport = WsTransport::try_new("wss://push.example.com").unwrap();
assert!(transport.next_line().await.is_none());
}
fn echo_subprotocol(
_request: &ServerRequest,
mut response: ServerResponse,
) -> Result<ServerResponse, tokio_tungstenite::tungstenite::handshake::server::ErrorResponse>
{
response.headers_mut().insert(
SEC_WEBSOCKET_PROTOCOL,
HeaderValue::from_static(SUBPROTOCOL),
);
Ok(response)
}
async fn bind_loopback() -> (TcpListener, SocketAddr) {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
(listener, address)
}
#[tokio::test]
async fn test_open_stream_sends_wsok_then_creation_and_yields_each_line() {
let (listener, address) = bind_loopback().await;
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
let mut socket = accept_hdr_async(stream, echo_subprotocol).await.unwrap();
let first = socket.next().await.unwrap().unwrap();
let second = socket.next().await.unwrap().unwrap();
socket
.send(Message::text(
"WSOK\r\nCONOK,S1aa6c792585db57aT1726545,50000,5000,*\r\n",
))
.await
.unwrap();
socket
.send(Message::text("SERVNAME,Lightstreamer HTTP Server\r\n"))
.await
.unwrap();
socket.close(None).await.unwrap();
while socket.next().await.is_some() {}
(
first.into_text().unwrap().as_str().to_owned(),
second.into_text().unwrap().as_str().to_owned(),
)
});
let mut transport = WsTransport::try_new(format!("ws://{address}")).unwrap();
transport
.open_stream(StreamOpen::Create(Box::default()))
.await
.unwrap();
assert_eq!(transport.next_line().await.unwrap().unwrap(), "WSOK");
assert_eq!(
transport.next_line().await.unwrap().unwrap(),
"CONOK,S1aa6c792585db57aT1726545,50000,5000,*"
);
assert_eq!(
transport.next_line().await.unwrap().unwrap(),
"SERVNAME,Lightstreamer HTTP Server"
);
assert!(
transport.next_line().await.is_none(),
"a clean close ends the stream with None"
);
let (wsok, creation) = server.await.unwrap();
assert_eq!(wsok, "wsok", "the establishment check goes first (§15)");
assert!(
creation.starts_with("create_session\r\n"),
"request name on its own line (§6.2.2), got {creation:?}"
);
}
#[tokio::test]
async fn test_abrupt_disconnect_reports_connection_lost() {
let (listener, address) = bind_loopback().await;
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
let mut socket = accept_hdr_async(stream, echo_subprotocol).await.unwrap();
let _wsok = socket.next().await.unwrap().unwrap();
let _create = socket.next().await.unwrap().unwrap();
socket.send(Message::text("WSOK\r\n")).await.unwrap();
});
let mut transport = WsTransport::try_new(format!("ws://{address}")).unwrap();
transport
.open_stream(StreamOpen::Create(Box::default()))
.await
.unwrap();
assert_eq!(transport.next_line().await.unwrap().unwrap(), "WSOK");
let error = transport.next_line().await.unwrap().unwrap_err();
assert!(
matches!(error, TransportError::ConnectionLost { .. }),
"expected ConnectionLost, got {error:?}"
);
server.await.unwrap();
}
#[tokio::test]
async fn test_binary_message_from_server_is_malformed_frame() {
let (listener, address) = bind_loopback().await;
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
let mut socket = accept_hdr_async(stream, echo_subprotocol).await.unwrap();
let _wsok = socket.next().await.unwrap().unwrap();
let _create = socket.next().await.unwrap().unwrap();
socket
.send(Message::binary(vec![0xDE, 0xAD, 0xBE, 0xEF]))
.await
.unwrap();
while socket.next().await.is_some() {}
});
let mut transport = WsTransport::try_new(format!("ws://{address}")).unwrap();
transport
.open_stream(StreamOpen::Create(Box::default()))
.await
.unwrap();
let error = transport.next_line().await.unwrap().unwrap_err();
assert!(
matches!(error, TransportError::MalformedFrame { .. }),
"expected MalformedFrame, got {error:?}"
);
drop(server);
}
#[tokio::test]
async fn test_wrong_negotiated_subprotocol_is_refused() {
let (listener, address) = bind_loopback().await;
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
let _ = accept_hdr_async(stream, |_: &ServerRequest, mut response: ServerResponse| {
response.headers_mut().insert(
SEC_WEBSOCKET_PROTOCOL,
HeaderValue::from_static("something-else"),
);
Ok(response)
})
.await;
});
let mut transport = WsTransport::try_new(format!("ws://{address}")).unwrap();
let error = transport
.open_stream(StreamOpen::Create(Box::default()))
.await
.unwrap_err();
assert!(
matches!(error, TransportError::Connect { .. }),
"expected Connect, got {error:?}"
);
drop(server);
}
#[tokio::test]
async fn test_absent_negotiated_subprotocol_is_refused() {
let (listener, address) = bind_loopback().await;
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
let _ = accept_hdr_async(stream, |_: &ServerRequest, response: ServerResponse| {
Ok(response)
})
.await;
});
let mut transport = WsTransport::try_new(format!("ws://{address}")).unwrap();
let error = transport
.open_stream(StreamOpen::Create(Box::default()))
.await
.unwrap_err();
assert!(
matches!(error, TransportError::Connect { .. }),
"expected Connect, got {error:?}"
);
drop(server);
}
#[tokio::test]
async fn test_a_clean_close_after_an_unterminated_line_is_malformed() {
let (listener, address) = bind_loopback().await;
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
let mut socket = accept_hdr_async(stream, echo_subprotocol).await.unwrap();
let _wsok = socket.next().await.unwrap().unwrap();
let _create = socket.next().await.unwrap().unwrap();
socket
.send(Message::text("WSOK\r\nEND,31,session destro"))
.await
.unwrap();
socket.close(None).await.unwrap();
while socket.next().await.is_some() {}
});
let mut transport = WsTransport::try_new(format!("ws://{address}")).unwrap();
transport
.open_stream(StreamOpen::Create(Box::default()))
.await
.unwrap();
assert_eq!(transport.next_line().await.unwrap().unwrap(), "WSOK");
let error = transport.next_line().await.unwrap().unwrap_err();
assert!(
matches!(error, TransportError::MalformedFrame { .. }),
"expected a malformed line, got {error:?}"
);
assert!(transport.next_line().await.is_none());
server.await.unwrap();
}
#[tokio::test]
async fn test_send_control_frames_the_request_on_the_stream_connection() {
let (listener, address) = bind_loopback().await;
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
let mut socket = accept_hdr_async(stream, echo_subprotocol).await.unwrap();
let _wsok = socket.next().await.unwrap().unwrap();
let _create = socket.next().await.unwrap().unwrap();
let control = socket.next().await.unwrap().unwrap();
socket.send(Message::text("REQOK,1\r\n")).await.unwrap();
while socket.next().await.is_some() {}
control.into_text().unwrap().as_str().to_owned()
});
let mut transport = WsTransport::try_new(format!("ws://{address}")).unwrap();
transport
.open_stream(StreamOpen::Create(Box::default()))
.await
.unwrap();
transport
.send_control(EncodedRequest {
name: "control",
path: "/lightstreamer/control.txt",
parameters: "LS_reqId=1&LS_op=delete&LS_subId=1".to_owned(),
})
.await
.unwrap();
assert_eq!(transport.next_line().await.unwrap().unwrap(), "REQOK,1");
transport.close().await.unwrap();
assert_eq!(
server.await.unwrap(),
"control\r\nLS_reqId=1&LS_op=delete&LS_subId=1"
);
}
}