* chore: Reorganize examples and interop crates * fix interop tests
83 lines
2.1 KiB
Rust
83 lines
2.1 KiB
Rust
pub mod pb {
|
|
tonic::include_proto!("grpc.examples.echo");
|
|
}
|
|
|
|
use futures::Stream;
|
|
use std::net::SocketAddr;
|
|
use std::pin::Pin;
|
|
use tokio::sync::mpsc;
|
|
use tonic::{transport::Server, Request, Response, Status, Streaming};
|
|
|
|
use pb::{EchoRequest, EchoResponse};
|
|
|
|
type EchoResult<T> = Result<Response<T>, Status>;
|
|
type ResponseStream = Pin<Box<dyn Stream<Item = Result<EchoResponse, Status>> + Send + Sync>>;
|
|
|
|
#[derive(Debug)]
|
|
pub struct EchoServer {
|
|
addr: SocketAddr,
|
|
}
|
|
|
|
#[tonic::async_trait]
|
|
impl pb::echo_server::Echo for EchoServer {
|
|
async fn unary_echo(&self, request: Request<EchoRequest>) -> EchoResult<EchoResponse> {
|
|
let message = format!("{} (from {})", request.into_inner().message, self.addr);
|
|
|
|
Ok(Response::new(EchoResponse { message }))
|
|
}
|
|
|
|
type ServerStreamingEchoStream = ResponseStream;
|
|
|
|
async fn server_streaming_echo(
|
|
&self,
|
|
_: Request<EchoRequest>,
|
|
) -> EchoResult<Self::ServerStreamingEchoStream> {
|
|
Err(Status::unimplemented("not implemented"))
|
|
}
|
|
|
|
async fn client_streaming_echo(
|
|
&self,
|
|
_: Request<Streaming<EchoRequest>>,
|
|
) -> EchoResult<EchoResponse> {
|
|
Err(Status::unimplemented("not implemented"))
|
|
}
|
|
|
|
type BidirectionalStreamingEchoStream = ResponseStream;
|
|
|
|
async fn bidirectional_streaming_echo(
|
|
&self,
|
|
_: Request<Streaming<EchoRequest>>,
|
|
) -> EchoResult<Self::BidirectionalStreamingEchoStream> {
|
|
Err(Status::unimplemented("not implemented"))
|
|
}
|
|
}
|
|
|
|
#[tokio::main]
|
|
async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
|
let addrs = ["[::1]:50051", "[::1]:50052"];
|
|
|
|
let (tx, mut rx) = mpsc::unbounded_channel();
|
|
|
|
for addr in &addrs {
|
|
let addr = addr.parse()?;
|
|
let tx = tx.clone();
|
|
|
|
let server = EchoServer { addr };
|
|
let serve = Server::builder()
|
|
.add_service(pb::echo_server::EchoServer::new(server))
|
|
.serve(addr);
|
|
|
|
tokio::spawn(async move {
|
|
if let Err(e) = serve.await {
|
|
eprintln!("Error = {:?}", e);
|
|
}
|
|
|
|
tx.send(()).unwrap();
|
|
});
|
|
}
|
|
|
|
rx.recv().await;
|
|
|
|
Ok(())
|
|
}
|