mod error;
mod health;
mod interceptor;
mod service;
use std::sync::Arc;
use connectrpc::{ConnectRpcService, Router};
pub use error::from_upload_error;
pub use health::ConnectHealth;
pub use interceptor::AuthInterceptor;
pub use service::ConnectTransport;
use crate::pipeline::{Authenticated, HookSet, Pipeline};
use crate::store::{MultipartBlobStore, NamespaceStore};
pub use mkit_rpc::transport as proto;
pub fn router<B, N, H>(pipeline: Arc<Pipeline<B, N, H>>) -> Router
where
B: MultipartBlobStore + 'static,
N: NamespaceStore + 'static,
H: HookSet + 'static,
{
use proto::grpc::health::v1::HealthExt;
use proto::mkit::transport::v1::TransportServiceExt;
let router = Arc::new(ConnectTransport::new(pipeline.clone())).register(Router::new());
Arc::new(ConnectHealth::new(pipeline)).register(router)
}
pub fn service<B, N, H>(pipeline: Arc<Pipeline<B, N, H>>) -> ConnectRpcService
where
B: MultipartBlobStore + 'static,
N: NamespaceStore + 'static,
H: HookSet + 'static,
{
ConnectRpcService::new(router(pipeline.clone()))
.with_interceptor(AuthInterceptor::new(pipeline))
}
struct Shared<P> {
#[cfg(not(target_arch = "wasm32"))]
pipe: Arc<P>,
#[cfg(target_arch = "wasm32")]
pipe: send_wrapper::SendWrapper<Arc<P>>,
}
impl<P> Shared<P> {
fn new(pipe: Arc<P>) -> Self {
Self {
#[cfg(not(target_arch = "wasm32"))]
pipe,
#[cfg(target_arch = "wasm32")]
pipe: send_wrapper::SendWrapper::new(pipe),
}
}
fn get(&self) -> &P {
&self.pipe
}
fn arc(&self) -> Arc<P> {
let pipe: &Arc<P> = &self.pipe;
Arc::clone(pipe)
}
}
fn authenticated(ctx: &connectrpc::RequestContext) -> Result<Authenticated, crate::ServerError> {
ctx.extensions()
.get::<Authenticated>()
.cloned()
.ok_or_else(|| crate::ServerError::unauthenticated("missing authorization"))
}