use futures::{stream::FuturesUnordered, StreamExt}; use rspc_dev_utilities::test_data::{ dur_to_str, make_test_data, TestData, TestDataClient, TestDataServer, CALLS_PER_THREAD, DATASIZE, THREADS, }; use rspc::transport::serde::{TcpClient, TcpServer}; #[tokio::main] async fn main() -> Result<(), Box> { let data: TestData = make_test_data(DATASIZE); let mut server = TestDataServer::from(data); let t = TcpServer::new(&"127.0.0.1:6543").await.unwrap(); let srv_thread = tokio::spawn(async move { server.listen(t).await }); tokio::time::sleep(tokio::time::Duration::from_millis(10)).await; let t = TcpClient::connect("127.0.0.1:6543").await.unwrap(); let client = t.spawn().await; let client = TestDataClient::new(client); let now = std::time::Instant::now(); { let set = FuturesUnordered::new(); for _ in 0..THREADS { set.push(async { for _ in 0..CALLS_PER_THREAD { client.heavy_calc().await.unwrap(); } }); } let _: Vec<_> = set.collect().await; } println!("time: {}", dur_to_str(now.elapsed())); client.stop().await.unwrap(); srv_thread.await.unwrap().unwrap(); Ok(()) }