mod receiver;
mod sender;
pub use receiver::AsyncReceiver;
pub use sender::AsyncSender;
#[cfg(test)]
mod tests {
use super::*;
use crate::{Profile, ReceiverOptions};
use ::tokio::io::AsyncReadExt;
use ::tokio::time::timeout;
use std::time::Duration;
#[tokio::test]
async fn test_receiver_bind() {
let _guard = crate::TEST_MUTEX.lock().unwrap();
let url = format!("rist://@:{}", crate::next_test_port());
let result = AsyncReceiver::bind(Profile::Main, &url);
assert!(result.is_ok());
}
#[tokio::test]
async fn test_receiver_bind_with_options() {
let _guard = crate::TEST_MUTEX.lock().unwrap();
let url = format!("rist://@:{}", crate::next_test_port());
let options = ReceiverOptions::new().fifo_size(2048);
let result = AsyncReceiver::bind_with_options(Profile::Main, &url, options);
assert!(result.is_ok());
}
#[tokio::test]
async fn test_receiver_bind_rejects_invalid_fifo_size() {
let _guard = crate::TEST_MUTEX.lock().unwrap();
let url = format!("rist://@:{}", crate::next_test_port());
let options = ReceiverOptions::new().fifo_size(3);
let result = AsyncReceiver::bind_with_options(Profile::Main, &url, options);
assert!(result.is_err());
}
#[tokio::test]
async fn test_receiver_recv_timeout() {
let _guard = crate::TEST_MUTEX.lock().unwrap();
let url = format!("rist://@:{}", crate::next_test_port());
let receiver = AsyncReceiver::bind(Profile::Main, &url).unwrap();
let result = timeout(
Duration::from_millis(200),
receiver.recv_timeout(Duration::from_millis(100)),
)
.await;
assert!(result.is_ok());
let inner = result.unwrap();
assert!(inner.is_ok());
assert!(inner.unwrap().is_none()); }
#[tokio::test]
async fn test_sender_connect_empty_url() {
let _guard = crate::TEST_MUTEX.lock().unwrap();
let result = AsyncSender::connect(Profile::Main, "").await;
assert!(result.is_err());
}
#[tokio::test]
async fn test_async_client_server() {
let _guard = crate::TEST_MUTEX.lock().unwrap();
let port = crate::next_test_port();
let receiver_url = format!("rist://@:{port}");
let sender_url = format!("rist://127.0.0.1:{port}");
let receiver = AsyncReceiver::bind(Profile::Main, &receiver_url).unwrap();
let sender = AsyncSender::connect(Profile::Main, &sender_url)
.await
.unwrap();
let total_packets = 50;
for _ in 0..total_packets {
let payload = [0x47u8; 1316]; let sent = sender.send(&payload).await.unwrap();
assert_eq!(sent, 1316);
}
::tokio::time::sleep(Duration::from_millis(100)).await;
let mut received_count = 0;
for _ in 0..total_packets {
let result = receiver.recv_timeout(Duration::from_millis(100)).await;
if let Ok(Some(data)) = result {
assert_eq!(data.payload().len(), 1316);
assert_eq!(data.payload()[0], 0x47);
received_count += 1;
}
}
assert!(received_count > 0, "expected to receive some packets");
let _ = receiver.raw_stats();
}
#[tokio::test]
async fn test_stats_available() {
let _guard = crate::TEST_MUTEX.lock().unwrap();
let url = format!("rist://@:{}", crate::next_test_port());
let receiver = AsyncReceiver::bind(Profile::Main, &url).unwrap();
let stats = receiver.raw_stats();
drop(stats);
}
#[tokio::test]
async fn test_stream_api() {
let _guard = crate::TEST_MUTEX.lock().unwrap();
let port = crate::next_test_port();
let receiver_url = format!("rist://@:{port}");
let sender_url = format!("rist://127.0.0.1:{port}");
let mut receiver = AsyncReceiver::bind(Profile::Main, &receiver_url).unwrap();
let sender = AsyncSender::connect(Profile::Main, &sender_url)
.await
.unwrap();
let test_data = b"Hello, RIST stream!";
sender.send(test_data).await.unwrap();
::tokio::time::sleep(Duration::from_millis(100)).await;
let mut buf = vec![0u8; 1024];
let read_result = timeout(Duration::from_millis(500), receiver.read(&mut buf)).await;
assert!(read_result.is_ok() || read_result.is_err());
}
}