forked from restatedev/restate
-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Make grpc server reusable by other metadata store implementations
- Loading branch information
1 parent
49dec91
commit 0dd0374
Showing
10 changed files
with
279 additions
and
152 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,76 @@ | ||
// Copyright (c) 2024 - Restate Software, Inc., Restate GmbH. | ||
// All rights reserved. | ||
// | ||
// Use of this software is governed by the Business Source License | ||
// included in the LICENSE file. | ||
// | ||
// As of the Change Date specified in that file, in accordance with | ||
// the Business Source License, use of this software will be governed | ||
// by the Apache License, Version 2.0. | ||
|
||
use http::Request; | ||
use hyper::body::Incoming; | ||
use hyper_util::service::TowerToHyperService; | ||
use restate_core::network::net_util; | ||
use restate_core::ShutdownError; | ||
use restate_types::health::HealthStatus; | ||
use restate_types::net::BindAddress; | ||
use restate_types::protobuf::common::MetadataServerStatus; | ||
use tonic::body::boxed; | ||
use tonic::service::Routes; | ||
use tower::ServiceExt; | ||
use tower_http::classify::{GrpcCode, GrpcErrorsAsFailures, SharedClassifier}; | ||
|
||
pub struct GrpcServer { | ||
bind_address: BindAddress, | ||
routes: Routes, | ||
} | ||
|
||
#[derive(Debug, thiserror::Error)] | ||
pub enum Error { | ||
#[error("failed running grpc server: {0}")] | ||
GrpcServer(#[from] net_util::Error), | ||
#[error("system is shutting down")] | ||
Shutdown(#[from] ShutdownError), | ||
} | ||
|
||
impl GrpcServer { | ||
pub fn new(bind_address: BindAddress, routes: Routes) -> Self { | ||
Self { | ||
bind_address, | ||
routes, | ||
} | ||
} | ||
|
||
pub async fn run(self, health_status: HealthStatus<MetadataServerStatus>) -> Result<(), Error> { | ||
let span_factory = tower_http::trace::DefaultMakeSpan::new() | ||
.include_headers(true) | ||
.level(tracing::Level::ERROR); | ||
|
||
let trace_layer = tower_http::trace::TraceLayer::new(SharedClassifier::new( | ||
GrpcErrorsAsFailures::new().with_success(GrpcCode::FailedPrecondition), | ||
)) | ||
.make_span_with(span_factory); | ||
|
||
let server_builder = tonic::transport::Server::builder() | ||
.layer(trace_layer) | ||
.add_routes(self.routes); | ||
|
||
let service = TowerToHyperService::new( | ||
server_builder | ||
.into_service() | ||
.map_request(|req: Request<Incoming>| req.map(boxed)), | ||
); | ||
|
||
net_util::run_hyper_server( | ||
&self.bind_address, | ||
service, | ||
"metadata-store-grpc", | ||
|| health_status.update(MetadataServerStatus::Ready), | ||
|| health_status.update(MetadataServerStatus::Unknown), | ||
) | ||
.await?; | ||
|
||
Ok(()) | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,81 @@ | ||
// Copyright (c) 2024 - Restate Software, Inc., Restate GmbH. | ||
// All rights reserved. | ||
// | ||
// Use of this software is governed by the Business Source License | ||
// included in the LICENSE file. | ||
// | ||
// As of the Change Date specified in that file, in accordance with | ||
// the Business Source License, use of this software will be governed | ||
// by the Apache License, Version 2.0. | ||
|
||
use http::{Request, Response}; | ||
use std::convert::Infallible; | ||
use tonic::body::BoxBody; | ||
use tonic::server::NamedService; | ||
use tonic::service::{Routes, RoutesBuilder}; | ||
use tonic_health::ServingStatus; | ||
use tower::Service; | ||
|
||
#[derive(Debug)] | ||
pub struct GrpcServiceBuilder<'a> { | ||
reflection_service_builder: Option<tonic_reflection::server::Builder<'a>>, | ||
routes_builder: RoutesBuilder, | ||
svc_names: Vec<&'static str>, | ||
} | ||
|
||
impl<'a> Default for GrpcServiceBuilder<'a> { | ||
fn default() -> Self { | ||
let routes_builder = RoutesBuilder::default(); | ||
|
||
Self { | ||
reflection_service_builder: Some(tonic_reflection::server::Builder::configure()), | ||
routes_builder, | ||
svc_names: Vec::default(), | ||
} | ||
} | ||
} | ||
|
||
impl<'a> GrpcServiceBuilder<'a> { | ||
pub fn add_service<S>(&mut self, svc: S) | ||
where | ||
S: Service<Request<BoxBody>, Response = Response<BoxBody>, Error = Infallible> | ||
+ NamedService | ||
+ Clone | ||
+ Send | ||
+ 'static, | ||
S::Future: Send + 'static, | ||
{ | ||
self.svc_names.push(S::NAME); | ||
self.routes_builder.add_service(svc); | ||
} | ||
|
||
pub fn register_file_descriptor_set_for_reflection<'b: 'a>( | ||
&mut self, | ||
encoded_file_descriptor_set: &'b [u8], | ||
) { | ||
self.reflection_service_builder = Some( | ||
self.reflection_service_builder | ||
.take() | ||
.expect("be present") | ||
.register_encoded_file_descriptor_set(encoded_file_descriptor_set), | ||
); | ||
} | ||
|
||
pub async fn build(mut self) -> Result<Routes, tonic_reflection::server::Error> { | ||
let (mut health_reporter, health_service) = tonic_health::server::health_reporter(); | ||
|
||
for svc_name in self.svc_names { | ||
health_reporter | ||
.set_service_status(svc_name, ServingStatus::Serving) | ||
.await; | ||
} | ||
|
||
self.routes_builder.add_service(health_service); | ||
self.routes_builder.add_service( | ||
self.reflection_service_builder | ||
.expect("be present") | ||
.build_v1()?, | ||
); | ||
Ok(self.routes_builder.routes()) | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.