mkit_server/connect/
health.rs1use core::fmt;
6use std::sync::Arc;
7
8use connectrpc::{RequestContext, Response, ServiceRequest, ServiceResult, ServiceStream};
9
10use super::Shared;
11use super::proto::grpc::health::v1::health_check_response::ServingStatus;
12use super::proto::grpc::health::v1::{Health, HealthCheckRequest, HealthCheckResponse};
13use super::proto::mkit::transport::v1::TRANSPORT_SERVICE_SERVICE_NAME;
14use crate::error::ServerError;
15use crate::pipeline::{HookSet, Pipeline};
16use crate::rt::send_wrap;
17use crate::store::{MultipartBlobStore, NamespaceStore};
18
19pub struct ConnectHealth<B, N, H> {
24 pipe: Shared<Pipeline<B, N, H>>,
25}
26
27impl<B, N, H> fmt::Debug for ConnectHealth<B, N, H> {
28 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
29 f.debug_struct("ConnectHealth").finish_non_exhaustive()
30 }
31}
32
33impl<B, N, H> ConnectHealth<B, N, H> {
34 #[must_use]
36 pub fn new(pipeline: Arc<Pipeline<B, N, H>>) -> Self {
37 Self {
38 pipe: Shared::new(pipeline),
39 }
40 }
41}
42
43#[allow(refining_impl_trait)]
44impl<B, N, H> Health for ConnectHealth<B, N, H>
45where
46 B: MultipartBlobStore + 'static,
47 N: NamespaceStore + 'static,
48 H: HookSet + 'static,
49{
50 async fn check(
51 &self,
52 _ctx: RequestContext,
53 request: ServiceRequest<'_, HealthCheckRequest>,
54 ) -> ServiceResult<HealthCheckResponse> {
55 let service = request.service;
56 if !(service.is_empty() || service == TRANSPORT_SERVICE_SERVICE_NAME) {
57 return Err(ServerError::not_found("unknown service").into());
58 }
59 let pipe = self.pipe.arc();
60 let healthy = send_wrap(async move { pipe.health().await.is_healthy() }).await;
61 let status = if healthy {
62 ServingStatus::SERVING
63 } else {
64 ServingStatus::NOT_SERVING
65 };
66 Response::ok(HealthCheckResponse {
67 status: status.into(),
68 ..Default::default()
69 })
70 }
71
72 async fn watch(
73 &self,
74 _ctx: RequestContext,
75 _request: ServiceRequest<'_, HealthCheckRequest>,
76 ) -> ServiceResult<ServiceStream<HealthCheckResponse>> {
77 Err(ServerError::unimplemented("Watch is not supported by this server; use Check").into())
78 }
79}