Skip to main content

echo_server/
echo_server.rs

1use std::{net::SocketAddr, pin::Pin, time::Duration};
2
3use futures::Stream;
4use rpc::{
5    echo::{
6        echo_server::{Echo, EchoServer},
7        EchoRequest, EchoResponse,
8    },
9    health_reporter, RpcServer, RpcServerConfig,
10};
11use tonic::{Request, Response, Status};
12
13#[derive(Default)]
14struct EchoService;
15
16#[tonic::async_trait]
17impl Echo for EchoService {
18    type ServerStreamStream =
19        Pin<Box<dyn Stream<Item = Result<EchoResponse, Status>> + Send + 'static>>;
20    type BidirectionalStreamStream = Self::ServerStreamStream;
21
22    async fn echo(&self, request: Request<EchoRequest>) -> Result<Response<EchoResponse>, Status> {
23        Ok(Response::new(EchoResponse {
24            message: request.into_inner().message,
25        }))
26    }
27
28    async fn server_stream(
29        &self,
30        request: Request<EchoRequest>,
31    ) -> Result<Response<Self::ServerStreamStream>, Status> {
32        Ok(Response::new(Box::pin(tokio_stream::once(Ok(
33            EchoResponse {
34                message: request.into_inner().message,
35            },
36        )))))
37    }
38
39    async fn client_stream(
40        &self,
41        request: Request<tonic::Streaming<EchoRequest>>,
42    ) -> Result<Response<EchoResponse>, Status> {
43        let mut input = request.into_inner();
44        let mut messages = Vec::new();
45        while let Some(request) = input.message().await? {
46            messages.push(request.message);
47        }
48        Ok(Response::new(EchoResponse {
49            message: messages.join(","),
50        }))
51    }
52
53    async fn bidirectional_stream(
54        &self,
55        request: Request<tonic::Streaming<EchoRequest>>,
56    ) -> Result<Response<Self::BidirectionalStreamStream>, Status> {
57        let replies = futures::stream::unfold(request.into_inner(), |mut input| async move {
58            match input.message().await {
59                Ok(Some(request)) => Some((
60                    Ok(EchoResponse {
61                        message: request.message,
62                    }),
63                    input,
64                )),
65                Err(status) => Some((Err(status), input)),
66                Ok(None) => None,
67            }
68        });
69        Ok(Response::new(Box::pin(replies)))
70    }
71}
72
73#[tokio::main]
74async fn main() -> Result<(), Box<dyn std::error::Error>> {
75    let address: SocketAddr = "[::1]:50051".parse()?;
76    let server = RpcServer::new(
77        RpcServerConfig::new(address)
78            .with_request_timeout(Duration::from_secs(10))
79            .with_concurrency_limit(1_024),
80    );
81    let (mut reporter, health_service) = health_reporter();
82    reporter.set_serving::<EchoServer<EchoService>>().await;
83
84    server
85        .router()
86        .add_service(health_service)
87        .add_service(EchoServer::new(EchoService))
88        .serve(server.config().address())
89        .await?;
90    Ok(())
91}