use std::sync::Arc;
use thubo::Bytes;
use tokio::{
net::{TcpListener, TcpStream},
sync::Barrier,
};
const ADDR: &str = "localhost:9999";
const N: u64 = 1_000;
async fn pong(barrier: Arc<Barrier>) {
let listener = TcpListener::bind(ADDR).await.unwrap();
barrier.wait().await;
let (stream, _addr) = listener.accept().await.unwrap();
let (stream_reader, stream_writer) = stream.into_split();
let (mut sender, _sender_task) = thubo::sender(stream_writer).build();
let (mut receiver, _receiver_task) = thubo::receiver(stream_reader).build();
loop {
let (data, _) = receiver.recv().await.unwrap();
sender.send(&data).await.unwrap();
}
}
#[cfg(not(miri))]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn base() {
let barrier = Arc::new(Barrier::new(2));
tokio::task::spawn(pong(barrier.clone()));
barrier.wait().await;
let stream = TcpStream::connect(ADDR).await.unwrap();
let (stream_reader, stream_writer) = stream.into_split();
let (mut sender, _sender_task) = thubo::sender(stream_writer).build();
let (mut receiver, _receiver_task) = thubo::receiver(stream_reader).build();
for size in [8, 32_000, 1_000_000] {
let payload = Bytes::from(vec![0u8; size]);
for _ in 0..N {
sender.send(&payload).await.unwrap();
let _ = receiver.recv().await.unwrap();
}
}
}