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}