/* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL * * Modified by Coffey Labs in 2026 for INBUXA. */ pub mod diagnose; use crate::{ api::diagnose::{DeliveryStage, spawn_delivery_diagnose}, auth::{ authenticate::Authenticator, oauth::auth::OAuthApiHandler, permissions::AccountApiHandler, }, }; use common::{ Server, auth::{AccessToken, oauth::GrantType}, manager::application::Resource, }; use groupware::calendar::itip::{ItipIngest, RsvpRequest}; use http_body_util::{StreamBody, combinators::BoxBody}; use http_proto::{ HttpRequest, HttpResponse, HttpSessionData, JsonResponse, ToHttpResponse, request::{decode_path_element, fetch_body}, }; use hyper::{ Method, StatusCode, header::{self, CONTENT_ENCODING}, }; use jmap::api::{ToJmapHttpResponse, ToRequestError}; use jmap_proto::error::request::RequestError; use registry::schema::enums::Permission; use std::time::Duration; use utils::url_params::UrlParams; pub trait ManagementApi: Sync + Send { fn handle_api_request( &self, req: &mut HttpRequest, session: &HttpSessionData, ) -> impl Future> + Send; fn management_access_token( &self, req: &HttpRequest, session: &HttpSessionData, ) -> impl Future> + Send; } impl ManagementApi for Server { #[allow(unused_variables)] async fn handle_api_request( &self, req: &mut HttpRequest, session: &HttpSessionData, ) -> trc::Result { let is_post = req.method() == Method::POST; let body = if is_post { fetch_body(req, 1024 * 1024, session.session_id).await } else { None }; let path = req.uri().path().split('/').skip(2).collect::>(); match path.first().copied().unwrap_or_default() { "auth" if is_post => { self.is_http_anonymous_request_allowed(session.remote_ip) .await?; Box::pin(self.handle_login_request( session, body.ok_or_else(|| trc::LimitEvent::SizeRequest.into_err())?, )) .await } "calendar" if is_post && path.get(1).copied() == Some("rsvp") && self.core.groupware.itip_http_rsvp_url.is_some() => { self.is_http_anonymous_request_allowed(session.remote_ip) .await?; let request = serde_json::from_slice::( &body.ok_or_else(|| trc::LimitEvent::SizeRequest.into_err())?, ) .map_err(|err| { trc::EventType::Resource(trc::ResourceEvent::BadParameters).from_json_error(err) })?; self.http_rsvp_handle(request, accept_language(req), session.remote_ip) .await .map(|response| JsonResponse::new(response).no_cache().into_http_response()) } "discover" => { if let Some(email) = path.get(1).copied() { self.is_http_anonymous_request_allowed(session.remote_ip) .await?; self.handle_discover_request(session, decode_path_element(email).as_ref()) .await } else { Err(trc::ResourceEvent::NotFound.into_err()) } } // inbuxa: EX-23, "Explain this", streamed as the model writes "explain" if is_post => { let (in_flight, access_token) = self.authenticate_headers(req, session).await?; jmap::inbuxa::explanation::assert_allowed(&access_token)?; let subject = body .as_deref() .and_then(|body| serde_json::from_slice::(body).ok()) .and_then(|mut body| body.get_mut("subject").map(serde_json::Value::take)) .ok_or_else(|| { trc::ResourceEvent::BadParameters .into_err() .details("Expected {\"subject\": …}") })?; let question = jmap::inbuxa::explanation::question(self, &access_token, &subject).await?; Ok(explain_stream(self.clone(), access_token, question, in_flight)) } // inbuxa: try a saved directory before anything signs in through it "directory" if is_post && path.get(1).copied() == Some("test") => { let (_in_flight, access_token) = self.authenticate_headers(req, session).await?; jmap::inbuxa::directory_test::assert_allowed(&access_token)?; let request = body .as_deref() .and_then(|body| serde_json::from_slice::(body).ok()) .unwrap_or_default(); let answer = jmap::inbuxa::directory_test::test(self, &request).await?; Ok(JsonResponse::new(answer).no_cache().into_http_response()) } // inbuxa: send one sample event to a saved webhook "webhook" if is_post && path.get(1).copied() == Some("test") => { let (_in_flight, access_token) = self.authenticate_headers(req, session).await?; jmap::inbuxa::webhook_test::assert_allowed(&access_token)?; let request = body .as_deref() .and_then(|body| serde_json::from_slice::(body).ok()) .unwrap_or_default(); let answer = jmap::inbuxa::webhook_test::test(self, &request).await?; Ok(JsonResponse::new(answer).no_cache().into_http_response()) } // inbuxa: whether the outside world reaches each node's ports "ports" if path.get(1).copied() == Some("check") => { let (_in_flight, access_token) = self.authenticate_headers(req, session).await?; if access_token.tenant_id().is_some() { return Err(trc::JmapEvent::Forbidden .into_err() .details("Port checks are for server-level administrators.")); } access_token.enforce_permission(Permission::SysNetworkListenerGet)?; let answer = common::reachability::report(self).await?; Ok(JsonResponse::new(answer).no_cache().into_http_response()) } "account" => { // Authenticate request let (_in_flight, access_token) = self.authenticate_headers(req, session).await?; self.handle_account_request(&access_token).await } "schema" => { // Authenticate request let (_in_flight, access_token) = self.authenticate_headers(req, session).await?; static SCHEMA_JSON: &[u8] = include_bytes!("../../../../resources/schema/schema.json.gz"); const SCHEMA_HASH: &str = include_str!("../../../../resources/schema/schema.json.sha256"); if path.get(1).is_some_and(|hash| hash == &SCHEMA_HASH) { // inbuxa: private, not public. This is behind // authenticate_headers and its CORS headers vary by // Origin, so a shared or origin-agnostic cache entry is // wrong -- and, being immutable, wrong for a year. Ok(Resource::new("application/json", SCHEMA_JSON.to_vec()) .into_http_response() .with_private_immutable_cache() .with_header(CONTENT_ENCODING, "gzip")) } else { Ok(HttpResponse::redirect(format!("/api/schema/{SCHEMA_HASH}"))) } } "token" => { let access_token = self.management_access_token(req, session).await?; let account_id = access_token.account_id(); match path.get(1).copied() { Some("delivery") => { // Validate the access token access_token.enforce_permission(Permission::LiveDeliveryTest)?; // Issue a live telemetry token valid for 60 seconds Ok(HttpResponse::new(StatusCode::OK) .with_no_cache() .with_text_body( self.encode_access_token( GrantType::LiveDelivery, account_id, self.account(account_id).await?.name(), 60, None, None, ) .await?, )) } // inbuxa: MON-23: a live telemetry token, valid 60 seconds Some(kind @ ("tracing" | "metrics")) => { let (grant, permission) = if kind == "tracing" { (GrantType::LiveTracing, Permission::LiveTracing) } else { (GrantType::LiveMetrics, Permission::LiveMetrics) }; access_token.enforce_permission(permission)?; crate::live::assert_server_level(&access_token)?; Ok(HttpResponse::new(StatusCode::OK) .with_no_cache() .with_text_body( self.encode_access_token( grant, account_id, self.account(account_id).await?.name(), 60, None, None, ) .await?, )) } _ => Err(trc::ResourceEvent::NotFound.into_err()), } } // inbuxa: the paths upstream's docs name, as aliases (MON-20, MON-22) "telemetry" if req.method() == Method::GET && path.get(2).copied() == Some("live") => { let access_token = self.management_access_token(req, session).await?; crate::live::assert_server_level(&access_token)?; match path.get(1).copied() { Some("traces") => { access_token.enforce_permission(Permission::LiveTracing)?; crate::live::live_tracing(req.uri().query()) } Some("metrics") => { access_token.enforce_permission(Permission::LiveMetrics)?; crate::live::live_metrics(&UrlParams::new(req.uri().query())) } _ => Err(trc::ResourceEvent::NotFound.into_err()), } } "live" => { let access_token = self.management_access_token(req, session).await?; let params = UrlParams::new(req.uri().query()); let account_id = access_token.account_id(); match ( path.get(1).copied().unwrap_or_default(), path.get(2).copied(), req.method(), ) { ("delivery", Some(target), &Method::GET) => { // Validate the access token access_token.enforce_permission(Permission::LiveDeliveryTest)?; let timeout = Duration::from_secs( params .parse::("timeout") .filter(|interval| *interval >= 1) .unwrap_or(30), ); let mut rx = spawn_delivery_diagnose( self.clone(), decode_path_element(target).to_lowercase(), timeout, ); Ok(HttpResponse::new(StatusCode::OK) .with_content_type("text/event-stream") .with_cache_control("no-store") .with_stream_body(BoxBody::new(StreamBody::new( async_stream::stream! { while let Some(stage) = rx.recv().await { yield Ok(stage.to_frame()); } yield Ok(DeliveryStage::Completed.to_frame()); }, )))) } // inbuxa: MON-20 to MON-22: live telemetry ("tracing", _, &Method::GET) => { access_token.enforce_permission(Permission::LiveTracing)?; crate::live::assert_server_level(&access_token)?; crate::live::live_tracing(req.uri().query()) } ("metrics", _, &Method::GET) => { access_token.enforce_permission(Permission::LiveMetrics)?; crate::live::assert_server_level(&access_token)?; crate::live::live_metrics(¶ms) } _ => Err(trc::ResourceEvent::NotFound.into_err()), } } _ => Err(trc::ResourceEvent::NotFound.into_err()), } } async fn management_access_token( &self, req: &HttpRequest, session: &HttpSessionData, ) -> trc::Result { let params = UrlParams::new(req.uri().query()); if let Some(token) = params.get("token") { let path = req.uri().path(); let grant = if path.starts_with("/api/live/delivery") { Some((GrantType::LiveDelivery, Permission::LiveDeliveryTest)) } else if path.starts_with("/api/live/tracing") || path.starts_with("/api/telemetry/traces/live") { // inbuxa: MON-23 Some((GrantType::LiveTracing, Permission::LiveTracing)) } else if path.starts_with("/api/live/metrics") || path.starts_with("/api/telemetry/metrics/live") { Some((GrantType::LiveMetrics, Permission::LiveMetrics)) } else { None }; if let Some((grant_type, permission)) = grant { self.validate_access_token(grant_type.into(), token) .await .map(|token_info| { AccessToken::from_permissions(token_info.account_id, [permission]) }) } else { self.authenticate_headers(req, session) .await .map(|(_, token)| token) } } else { self.authenticate_headers(req, session) .await .map(|(_, token)| token) } } } pub trait ToManageHttpResponse { fn into_http_response(self, challenge: AuthChallenge) -> HttpResponse; } impl ToManageHttpResponse for &trc::Error { fn into_http_response(self, challenge: AuthChallenge) -> HttpResponse { match self.as_ref() { trc::EventType::Auth( trc::AuthEvent::Failed | trc::AuthEvent::Error | trc::AuthEvent::TokenExpired, ) => HttpResponse::unauthorized(challenge), _ => self.to_request_error().into_http_response(), } } } pub fn accept_language(req: &HttpRequest) -> &str { req.headers() .get(header::ACCEPT_LANGUAGE) .and_then(|value| value.to_str().ok()) .map(|language| { let language = language.split_once(',').map_or(language, |(l, _)| l); language.split_once(';').map_or(language, |(l, _)| l).trim() }) .filter(|language| !language.is_empty()) .unwrap_or("en") } const BEARER_CHALLENGE: &str = concat!( concat!("Bearer realm=\"", types::brand_server!(), "\", "), "resource_metadata=\"/.well-known/oauth-protected-resource\"" ); const BASIC_CHALLENGE: &str = concat!("Basic realm=\"", types::brand_server!(), "\""); #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum AuthChallenge { Bearer, BearerAndBasic, } pub trait UnauthorizedResponse { fn unauthorized(challenge: AuthChallenge) -> Self; } impl UnauthorizedResponse for HttpResponse { fn unauthorized(challenge: AuthChallenge) -> Self { let response = HttpResponse::new(StatusCode::UNAUTHORIZED) .with_header(header::WWW_AUTHENTICATE, BEARER_CHALLENGE); if challenge == AuthChallenge::BearerAndBasic { response.with_header(header::WWW_AUTHENTICATE, BASIC_CHALLENGE) } else { response } .with_content_type("application/problem+json") .with_text_body(serde_json::to_string(&RequestError::unauthorized()).unwrap_or_default()) } } /// inbuxa: EX-23, the explanation as server-sent events: `delta` pieces as /// the model writes, then `done` with the whole explanation, or one `error`. /// The answer runs in its own task, so a client that goes away doesn't stop /// it: it finishes and is remembered (EX-24). fn explain_stream( server: Server, access_token: common::auth::AccessToken, question: Result< jmap::inbuxa::explanation::Question, jmap_proto::error::set::SetError< jmap_proto::object::inbuxa_explanation::ExplanationProperty, >, >, in_flight: Option, ) -> HttpResponse { use hyper::body::{Bytes, Frame}; use jmap::inbuxa::explanation::{answer, to_value}; fn event(name: &str, data: &serde_json::Value) -> Frame { Frame::data(Bytes::from(format!("event: {name}\ndata: {data}\n\n"))) } let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::(); let (done_tx, done_rx) = tokio::sync::oneshot::channel(); match question { Ok(question) => { tokio::spawn(async move { let result = answer(&server, &access_token, question, Some(tx)).await; let _ = done_tx.send(result); }); } Err(error) => { drop(tx); let _ = done_tx.send(Err(error)); } } HttpResponse::new(StatusCode::OK) .with_content_type("text/event-stream") .with_cache_control("no-store") .with_stream_body(BoxBody::new(StreamBody::new(async_stream::stream! { let _in_flight = in_flight; while let Some(text) = rx.recv().await { yield Ok(event("delta", &serde_json::json!({ "text": text }))); } match done_rx.await { Ok(Ok(answer)) => { let value = serde_json::to_value(to_value(answer)).unwrap_or_default(); yield Ok(event("done", &value)); } Ok(Err(error)) => { let value = serde_json::to_value(&error).unwrap_or_default(); yield Ok(event("error", &value)); } Err(_) => { yield Ok(event("error", &serde_json::json!({ "type": "serverFail", "description": "unavailable", }))); } } }))) }