use anyhow::Context;
use ar_pe_ce::{Result, Stream};
use futures::{FutureExt, TryStreamExt};
use serde::{Deserialize, Serialize};
#[derive(Debug, Deserialize, Serialize)]
struct Package {
payload: Vec<u8>,
}
#[ar_pe_ce::rpc]
trait Performance {
#[rpc(server_streaming)] async fn get_data(&self, arg: ()) -> Result<Stream<Package>>;
}
struct PerformanceServer;
#[ar_pe_ce::async_trait]
impl Performance for PerformanceServer {
#[tracing::instrument(skip(self))]
async fn get_data(&self, _: ()) -> Result<Stream<Package>> {
let stream = async_stream::stream! {
loop {
tracing::info!("Sending msg");
yield Ok(Package { payload: vec![0; 100_000_000] }); }
};
Ok(Box::pin(stream))
}
}
async fn client_task() -> anyhow::Result<()> {
let client = PerformanceClient::new("http://localhost:3000".parse()?);
let mut data_stream = client
.get_data(())
.await
.context("Could not get data stream")?;
loop {
tracing::info!("Waiting for message");
match data_stream
.try_next()
.await
.context("Could not retrieve message from data stream")?
{
Some(s) => {
tracing::info!(len = s.payload.len(), "Got message");
}
None => {
tracing::warn!("Got no message. Stream is closed");
break;
}
}
tokio::time::sleep(tokio::time::Duration::from_secs(5)).await;
}
Ok(())
}
#[tokio::main]
async fn main() -> anyhow::Result<()> {
tracing_subscriber::fmt().pretty().compact().init();
let mut client = tokio::spawn(async { client_task().await }).fuse();
use std::net::SocketAddr;
let addr = SocketAddr::from(([0, 0, 0, 0], 3000));
let server = PerformanceServer.serve(addr).fuse();
futures::pin_mut!(server);
futures::select! {
server = server => server?,
client = client => client??
};
Ok(())
}