Skip to main content

mkit_server/connect/
health.rs

1//! `grpc.health.v1.Health` over [`Pipeline::health`]: the generated trait,
2//! hand-implemented. The `connectrpc-health` crate would force
3//! `connectrpc/server`, which is not wasm-clean (see `build.rs`).
4
5use 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
19/// `Check` answers for the whole server (`""`) and for
20/// `mkit.transport.v1.TransportService`: `SERVING` when both stores answer
21/// their probe, else `NOT_SERVING`. Any other name is `not_found`. `Watch`
22/// is `unimplemented`, which tells a Watch-capable client not to retry.
23pub 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    /// Health over `pipeline`'s stores.
35    #[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}