1 pub mod pb { 2 tonic::include_proto!("grpc.examples.echo"); 3 } 4 5 use futures::Stream; 6 use pb::{EchoRequest, EchoResponse}; 7 use std::pin::Pin; 8 use tonic::{body::BoxBody, transport::Server, Request, Response, Status, Streaming}; 9 use tower::Service; 10 11 type EchoResult<T> = Result<Response<T>, Status>; 12 type ResponseStream = Pin<Box<dyn Stream<Item = Result<EchoResponse, Status>> + Send + Sync>>; 13 14 #[derive(Default)] 15 pub struct EchoServer; 16 17 #[tonic::async_trait] 18 impl pb::echo_server::Echo for EchoServer { 19 async fn unary_echo(&self, request: Request<EchoRequest>) -> EchoResult<EchoResponse> { 20 let message = request.into_inner().message; 21 Ok(Response::new(EchoResponse { message })) 22 } 23 24 type ServerStreamingEchoStream = ResponseStream; 25 26 async fn server_streaming_echo( 27 &self, 28 _: Request<EchoRequest>, 29 ) -> EchoResult<Self::ServerStreamingEchoStream> { 30 Err(Status::unimplemented("not implemented")) 31 } 32 33 async fn client_streaming_echo( 34 &self, 35 _: Request<Streaming<EchoRequest>>, 36 ) -> EchoResult<EchoResponse> { 37 Err(Status::unimplemented("not implemented")) 38 } 39 40 type BidirectionalStreamingEchoStream = ResponseStream; 41 42 async fn bidirectional_streaming_echo( 43 &self, 44 _: Request<Streaming<EchoRequest>>, 45 ) -> EchoResult<Self::BidirectionalStreamingEchoStream> { 46 Err(Status::unimplemented("not implemented")) 47 } 48 } 49 50 #[tokio::main] 51 async fn main() -> Result<(), Box<dyn std::error::Error>> { 52 let addr = "[::1]:50051".parse().unwrap(); 53 let server = EchoServer::default(); 54 55 Server::builder() 56 .interceptor_fn(move |svc, req| { 57 let auth_header = req.headers().get("authorization").clone(); 58 59 let authed = if let Some(auth_header) = auth_header { 60 auth_header == "Bearer some-secret-token" 61 } else { 62 false 63 }; 64 65 let fut = svc.call(req); 66 67 async move { 68 if authed { 69 fut.await 70 } else { 71 // Cancel the inner future since we never await it 72 // the IO never gets registered. 73 drop(fut); 74 let res = http::Response::builder() 75 .header("grpc-status", "16") 76 .body(BoxBody::empty()) 77 .unwrap(); 78 Ok(res) 79 } 80 } 81 }) 82 .add_service(pb::echo_server::EchoServer::new(server)) 83 .serve(addr) 84 .await?; 85 86 Ok(()) 87 } 88