use fastwebsockets::{Frame, OpCode, Payload, Role, WebSocket};
use h2ts_server::bridge;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::time::{sleep, Duration};
const BUF: usize = 16 * 1024;
const BIG: usize = 1024 * 1024;
#[tokio::test]
async fn peer_to_ws_producer_blocks_until_consumer_drains() {
let (client_io, server_io) = tokio::io::duplex(BUF);
let (peer_for_bridge, mut peer_test) = tokio::io::duplex(BUF);
tokio::spawn(async move {
let _ = bridge(server_io, peer_for_bridge).await;
});
let mut client_ws = WebSocket::after_handshake(client_io, Role::Client);
let payload: Vec<u8> = (0..BIG).map(|i| (i & 0xff) as u8).collect();
let expected = payload.clone();
let writer = tokio::spawn(async move {
peer_test.write_all(&payload).await.unwrap();
peer_test
});
sleep(Duration::from_millis(50)).await;
assert!(
!writer.is_finished(),
"producer finished with the consumer paused — the bridge buffered unboundedly instead of applying backpressure"
);
let mut got = Vec::with_capacity(expected.len());
while got.len() < expected.len() {
let frame = client_ws.read_frame().await.unwrap();
match frame.opcode {
OpCode::Binary | OpCode::Continuation => got.extend_from_slice(&frame.payload),
OpCode::Close => break,
_ => {}
}
}
let _peer_test = writer.await.unwrap();
assert_eq!(got.len(), expected.len());
assert_eq!(got, expected, "a backpressured stream must stay byte-exact");
}
#[tokio::test]
async fn ws_to_peer_producer_blocks_until_consumer_drains() {
let (client_io, server_io) = tokio::io::duplex(BUF);
let (peer_for_bridge, mut peer_test) = tokio::io::duplex(BUF);
tokio::spawn(async move {
let _ = bridge(server_io, peer_for_bridge).await;
});
let payload: Vec<u8> = (0..BIG).map(|i| (i & 0xff) as u8).collect();
let expected = payload.clone();
let writer = tokio::spawn(async move {
let mut client_ws = WebSocket::after_handshake(client_io, Role::Client);
client_ws
.write_frame(Frame::binary(Payload::Owned(payload)))
.await
.unwrap();
client_ws });
sleep(Duration::from_millis(50)).await;
assert!(
!writer.is_finished(),
"WS producer finished with the peer paused — no backpressure"
);
let mut got = vec![0u8; expected.len()];
peer_test.read_exact(&mut got).await.unwrap();
let _client_ws = writer.await.unwrap();
assert_eq!(got, expected, "a backpressured stream must stay byte-exact");
}