echo_server/
echo_server.rs1use 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}