Skip to main content

mkit_server/connect/
mod.rs

1//! The `mkit.transport.v1` Connect binding (feature `connect`,
2//! SPEC-TRANSPORT-CONNECT): [`service`] mounts `TransportService` and
3//! `grpc.health.v1.Health` over a [`Pipeline`], behind the
4//! [`AuthInterceptor`]. It is runtime-agnostic and wasm-clean (connectrpc
5//! without its `server` and `zstd` features), so the native adapter serves
6//! it through axum and the Workers adapter through its fetch bridge, both
7//! unchanged.
8//!
9//! Every RPC runs pipeline stage 0 in the interceptor, which sees the exact
10//! unary request bytes (the auth v2 body commitment), and hands the
11//! handler an [`Authenticated`] bound to the procedure called. A handler
12//! decodes its message with the shared wire helpers ([`crate::refs`],
13//! [`crate::upload`]), calls one pipeline entry point inside
14//! [`crate::send_wrap`] and encodes the answer. Errors cross the wire only
15//! through `From<ServerError> for ConnectError`: public message, code,
16//! HTTP status, response headers and typed details.
17//!
18//! `grpc.health.v1.Health` is not authenticated: `Check` reports whether
19//! both stores answer their probe (`SERVING` or `NOT_SERVING`), for load
20//! balancers and kubelet probes.
21//!
22//! A `connect-timeout-ms` header makes connectrpc compute a deadline with
23//! `Instant::now()`, which panics on wasm32. The Workers adapter strips the
24//! header before dispatch (`mkit_worker_common::adapter::is_deadline_header`);
25//! this binding does not.
26//!
27//! An upload is read message by message and stops at the `last` chunk.
28//! connectrpc 0.9.1 then drains at most 1 MiB (and, natively, 5 s) more of
29//! the request body before it resets the stream (RUSTSEC-2026-0304), so a
30//! client cannot hold the handler open with trailing bytes.
31
32mod error;
33mod health;
34mod interceptor;
35mod service;
36
37use std::sync::Arc;
38
39use connectrpc::{ConnectRpcService, Router};
40
41pub use error::from_upload_error;
42pub use health::ConnectHealth;
43pub use interceptor::AuthInterceptor;
44pub use service::ConnectTransport;
45
46use crate::pipeline::{Authenticated, HookSet, Pipeline};
47use crate::store::{MultipartBlobStore, NamespaceStore};
48
49/// Shared generated transport and health messages and service traits.
50pub use mkit_rpc::transport as proto;
51
52/// `TransportService` and `Health` over `pipeline`, without the
53/// interceptor: every authenticated transport RPC then fails
54/// `unauthenticated`; the M1 stub RPCs answer `unimplemented`. Mount
55/// [`service`] unless another layer installs [`AuthInterceptor`].
56pub fn router<B, N, H>(pipeline: Arc<Pipeline<B, N, H>>) -> Router
57where
58    B: MultipartBlobStore + 'static,
59    N: NamespaceStore + 'static,
60    H: HookSet + 'static,
61{
62    use proto::grpc::health::v1::HealthExt;
63    use proto::mkit::transport::v1::TransportServiceExt;
64
65    let router = Arc::new(ConnectTransport::new(pipeline.clone())).register(Router::new());
66    Arc::new(ConnectHealth::new(pipeline)).register(router)
67}
68
69/// [`router`] behind [`AuthInterceptor`]: what an adapter mounts. Apply
70/// deployment limits with `ConnectRpcService::with_limits`.
71pub fn service<B, N, H>(pipeline: Arc<Pipeline<B, N, H>>) -> ConnectRpcService
72where
73    B: MultipartBlobStore + 'static,
74    N: NamespaceStore + 'static,
75    H: HookSet + 'static,
76{
77    ConnectRpcService::new(router(pipeline.clone()))
78        .with_interceptor(AuthInterceptor::new(pipeline))
79}
80
81/// The pipeline as connectrpc's `Send + Sync` service objects hold it: an
82/// `Arc` on native targets. On wasm32 the pipeline is `!Send` (Workers
83/// handles are), so the `Arc` sits in a `SendWrapper`: Workers run
84/// single-threaded, and a use on another thread panics, never undefined
85/// behavior.
86struct Shared<P> {
87    #[cfg(not(target_arch = "wasm32"))]
88    pipe: Arc<P>,
89    #[cfg(target_arch = "wasm32")]
90    pipe: send_wrapper::SendWrapper<Arc<P>>,
91}
92
93impl<P> Shared<P> {
94    fn new(pipe: Arc<P>) -> Self {
95        Self {
96            #[cfg(not(target_arch = "wasm32"))]
97            pipe,
98            #[cfg(target_arch = "wasm32")]
99            pipe: send_wrapper::SendWrapper::new(pipe),
100        }
101    }
102
103    fn get(&self) -> &P {
104        &self.pipe
105    }
106
107    /// An owned handle, to move into a [`crate::send_wrap`]ped future.
108    fn arc(&self) -> Arc<P> {
109        // On wasm32 this derefs the `SendWrapper` (its thread check).
110        let pipe: &Arc<P> = &self.pipe;
111        Arc::clone(pipe)
112    }
113}
114
115/// The [`Authenticated`] the interceptor stored for this request.
116fn authenticated(ctx: &connectrpc::RequestContext) -> Result<Authenticated, crate::ServerError> {
117    ctx.extensions()
118        .get::<Authenticated>()
119        .cloned()
120        .ok_or_else(|| crate::ServerError::unauthenticated("missing authorization"))
121}