1. Protocol Buffers
// proto/product.proto
syntax = "proto3";
package product;
service ProductService {
rpc GetProduct (GetProductRequest) returns (Product);
rpc ListProducts (ListProductsRequest) returns (stream Product);
rpc CreateProduct (CreateProductRequest) returns (Product);
}
message Product {
string id = 1;
string name = 2;
double price = 3;
optional string description = 4;
}
message GetProductRequest { string id = 1; }
message ListProductsRequest { int32 limit = 1; int32 offset = 2; }
message CreateProductRequest { string name = 1; double price = 2; }
// build.rs
fn main() -> Result<(), Box<dyn std::error::Error>> {
tonic_build::compile_protos("proto/product.proto")?;
Ok(())
}
2. gRPC Server
use tonic::{transport::Server, Request, Response, Status};
pub mod product_proto {
tonic::include_proto!("product");
}
use product_proto::product_service_server::{ProductService, ProductServiceServer};
use product_proto::{Product, GetProductRequest, CreateProductRequest};
#[derive(Default)]
pub struct MyProductService;
#[tonic::async_trait]
impl ProductService for MyProductService {
async fn get_product(
&self,
request: Request<GetProductRequest>,
) -> Result<Response<Product>, Status> {
let id = request.into_inner().id;
// Lookup from DB
let product = Product {
id: id.clone(),
name: "Sample".into(),
price: 29.99,
description: Some("A sample product".into()),
};
Ok(Response::new(product))
}
type ListProductsStream = tokio_stream::wrappers::ReceiverStream<Result<Product, Status>>;
async fn list_products(
&self,
request: Request<ListProductsRequest>,
) -> Result<Response<Self::ListProductsStream>, Status> {
let (tx, rx) = tokio::sync::mpsc::channel(128);
tokio::spawn(async move {
for i in 0..10 {
tx.send(Ok(Product {
id: format!("prod-{}", i),
name: format!("Product {}", i),
price: i as f64 * 10.0,
description: None,
})).await.unwrap();
}
});
Ok(Response::new(tokio_stream::wrappers::ReceiverStream::new(rx)))
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
Server::builder()
.add_service(ProductServiceServer::new(MyProductService::default()))
.serve("0.0.0.0:50051".parse()?)
.await?;
Ok(())
}
3. gRPC Client
use product_proto::product_service_client::ProductServiceClient;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let mut client = ProductServiceClient::connect("http://localhost:50051").await?;
let response = client.get_product(GetProductRequest { id: "prod-1".into() }).await?;
println!("Product: {:?}", response.into_inner());
// Streaming
let mut stream = client.list_products(ListProductsRequest { limit: 10, offset: 0 }).await?.into_inner();
while let Some(product) = stream.message().await? {
println!("Received: {:?}", product);
}
Ok(())
}
Next article: Message Queues & Background Jobs.