Files
Stalwart/crates/http/src/request.rs

842 lines
35 KiB
Rust

/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <hello@stalw.art>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*/
use crate::{
HttpSessionManager,
api::{ManagementApi, ToManageHttpResponse},
auth::{
authenticate::{Authenticator, HttpHeaders},
oauth::{
FormData, auth::OAuthApiHandler, openid::OpenIdHandler,
registration::ClientRegistrationHandler, token::TokenHandler,
},
},
form::FormHandler,
};
use common::{
BuildServer, Inner, KV_ACME, Server,
ipc::PushEvent,
manager::application::Resource,
network::{SessionData, SessionManager, SessionStream},
};
use dav::{DavMethod, request::DavRequestHandler};
use groupware::{DavResourceName, calendar::itip::ItipIngest};
use http_proto::{
DownloadResponse, HtmlResponse, HttpContext, HttpRequest, HttpResponse, HttpResponseBody,
HttpSessionData, JsonProblemResponse, ToHttpResponse, form_urlencoded, request::fetch_body,
};
use hyper::{
Method, StatusCode, body,
header::{self, CONTENT_TYPE},
server::conn::http1,
service::service_fn,
};
use hyper_util::rt::TokioIo;
use jmap::{
api::{
ToJmapHttpResponse, event_source::EventSourceHandler, request::RequestHandler,
session::SessionHandler,
},
blob::{download::BlobDownload, upload::BlobUpload},
websocket::upgrade::WebSocketUpgrade,
};
use jmap_proto::request::{Request, capability::Session};
use registry::schema::enums::Permission;
use std::{net::IpAddr, str::FromStr, sync::Arc};
use store::dispatch::lookup::KeyValue;
use trc::SecurityEvent;
use types::{blob::BlobId, id::Id};
pub trait ParseHttp: Sync + Send {
fn parse_http_request(
&self,
req: HttpRequest,
session: HttpSessionData,
) -> impl Future<Output = trc::Result<HttpResponse>> + Send;
}
impl ParseHttp for Server {
async fn parse_http_request(
&self,
mut req: HttpRequest,
session: HttpSessionData,
) -> trc::Result<HttpResponse> {
let mut path = req.uri().path().split('/');
path.next();
// Validate endpoint access
let ctx = HttpContext::new(&session, &req);
match ctx.has_endpoint_access(self).await {
StatusCode::OK => (),
status => {
// Allow loopback address to avoid lockouts
if !session.remote_ip.is_loopback() {
return Ok(JsonProblemResponse(status).into_http_response());
}
}
}
match path.next().unwrap_or_default() {
"jmap" => {
match (path.next().unwrap_or_default(), req.method()) {
("", &Method::POST) => {
// Authenticate request
let (_in_flight, access_token) =
self.authenticate_headers(&req, &session).await?;
let bytes = fetch_body(
&mut req,
if !access_token.has_permission(Permission::UnlimitedUploads) {
self.core.jmap.upload_max_size
} else {
0
},
session.session_id,
)
.await
.ok_or_else(|| trc::LimitEvent::SizeRequest.into_err())?;
return Ok(self
.handle_jmap_request(
Request::parse(
&bytes,
self.core.jmap.request_max_calls,
self.core.jmap.request_max_size,
)?,
&access_token,
&session,
)
.await
.into_http_response());
}
("download", &Method::GET) => {
// Authenticate request
let (_in_flight, access_token) =
self.authenticate_headers(&req, &session).await?;
if let (Some(_), Some(blob_id), Some(name)) = (
path.next().and_then(|p| Id::from_str(p).ok()),
path.next().and_then(BlobId::from_base32),
path.next(),
) {
return match self.blob_download(&blob_id, &access_token).await? {
Some(blob) => Ok(DownloadResponse {
filename: name.to_string(),
content_type: req
.uri()
.query()
.and_then(|q| {
form_urlencoded::parse(q.as_bytes())
.find(|(k, _)| k == "accept")
.map(|(_, v)| v.into_owned())
})
.unwrap_or("application/octet-stream".to_string()),
blob,
}
.into_http_response()),
None => Err(trc::ResourceEvent::NotFound.into_err()),
};
}
}
("upload", &Method::POST) => {
// Authenticate request
let (_in_flight, access_token) =
self.authenticate_headers(&req, &session).await?;
if let Some(account_id) = path.next().and_then(|p| Id::from_str(p).ok()) {
return match fetch_body(
&mut req,
if !access_token.has_permission(Permission::UnlimitedUploads) {
self.core.jmap.upload_max_size
} else {
0
},
session.session_id,
)
.await
{
Some(bytes) => Ok(self
.blob_upload(
account_id,
req.headers()
.get(CONTENT_TYPE)
.and_then(|h| h.to_str().ok())
.unwrap_or("application/octet-stream"),
&bytes,
&access_token,
)
.await?
.into_http_response()),
None => Err(trc::LimitEvent::SizeUpload.into_err()),
};
}
}
("eventsource", &Method::GET) => {
// Authenticate request
let (_in_flight, access_token) =
self.authenticate_headers(&req, &session).await?;
return self.handle_event_source(req, access_token).await;
}
("ws", &Method::GET) => {
// Authenticate request
let (_in_flight, access_token) =
self.authenticate_headers(&req, &session).await?;
return self
.upgrade_websocket_connection(req, access_token, session)
.await;
}
("session", &Method::GET) => {
return if req.headers().contains_key(header::AUTHORIZATION) {
// Authenticate request
let (_in_flight, access_token) =
self.authenticate_headers(&req, &session).await?;
self.handle_session_resource(
self.core.network.http.url_https.to_string(),
&access_token,
)
.await
.map(|s| s.into_http_response())
} else {
Ok(Session::new(
&self.core.network.http.url_https,
&self.core.jmap.capabilities,
)
.into_http_response())
};
}
(_, &Method::OPTIONS) => {
return Ok(HttpResponse::new(StatusCode::NO_CONTENT));
}
_ => (),
}
}
"dav" => {
let response = match (
path.next().and_then(DavResourceName::parse),
DavMethod::parse(req.method()),
) {
(Some(_), Some(DavMethod::OPTIONS)) => HttpResponse::new(StatusCode::OK)
.with_header(
"DAV",
concat!(
"1, 2, 3, access-control, extended-mkcol, calendar-access, ",
"calendar-auto-schedule, calendar-no-timezone, addressbook"
),
)
.with_header(
"Allow",
concat!(
"OPTIONS, GET, HEAD, POST, PUT, DELETE, COPY, MOVE, MKCALENDAR, ",
"MKCOL, PROPFIND, PROPPATCH, LOCK, UNLOCK, REPORT, ACL"
),
),
(Some(resource), Some(method)) => {
// Authenticate request
let (_in_flight, access_token) =
self.authenticate_headers(&req, &session).await?;
self.handle_dav_request(req, access_token, &session, resource, method)
.await
}
(_, None) => HttpResponse::new(StatusCode::METHOD_NOT_ALLOWED),
(None, _) => HttpResponse::new(StatusCode::NOT_FOUND),
};
return Ok(response);
}
".well-known" => match (path.next().unwrap_or_default(), req.method()) {
("jmap", &Method::GET) => {
return Ok(HttpResponse::new(StatusCode::TEMPORARY_REDIRECT)
.with_no_cache()
.with_location("/jmap/session"));
}
("caldav", _) => {
return Ok(HttpResponse::new(StatusCode::TEMPORARY_REDIRECT)
.with_no_cache()
.with_location(DavResourceName::Cal.base_path()));
}
("carddav", _) => {
return Ok(HttpResponse::new(StatusCode::TEMPORARY_REDIRECT)
.with_no_cache()
.with_location(DavResourceName::Card.base_path()));
}
("oauth-authorization-server", &Method::GET) => {
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return self.handle_oauth_metadata().await;
}
("openid-configuration", &Method::GET) => {
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return self.handle_oidc_metadata(false).await;
}
("acme-challenge", &Method::GET) if self.has_acme_http_providers() => {
if let Some(token) = path.next() {
return match self
.in_memory_store()
.key_get::<String>(KeyValue::<()>::build_key(KV_ACME, token))
.await?
{
Some(proof) => Ok(Resource::new("text/plain", proof.into_bytes())
.into_http_response()),
None => Err(trc::ResourceEvent::NotFound.into_err()),
};
}
}
("mta-sts.txt", &Method::GET) => {
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return if let Some(policy) = &self.core.smtp.session.mta_sts_policy {
Ok(Resource::new("text/plain", policy.to_string().into_bytes())
.into_http_response())
} else {
Err(trc::ResourceEvent::NotFound.into_err())
};
}
("user-agent-configuration.json", &Method::GET) => {
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return Ok(Resource::new(
"application/json",
self.get_pacc_for_fomain(
req.headers()
.get(header::HOST)
.and_then(|h| h.to_str().ok())
.map(|h| h.rsplit_once(':').map_or(h, |(h, _)| h))
.unwrap_or_default(),
)
.await?
.into_bytes(),
)
.into_http_response()
.with_cors_unrestricted());
}
("mail-v1.xml", &Method::GET) => {
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return self
.handle_autoconfig_request(req.uri().query())
.await
.map(|resource| resource.into_http_response());
}
("autoconfig", &Method::GET)
if path.next().unwrap_or_default() == "mail"
&& path.next().unwrap_or_default() == "config-v1.1.xml" =>
{
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return self
.handle_autoconfig_request(req.uri().query())
.await
.map(|resource| resource.into_http_response().with_cors_unrestricted());
}
(_, &Method::OPTIONS) => {
return Ok(HttpResponse::new(StatusCode::NO_CONTENT).with_cors_unrestricted());
}
_ => (),
},
"auth" => match (path.next().unwrap_or_default(), req.method()) {
("device", &Method::POST) => {
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return self.handle_device_auth(&mut req, &session).await;
}
("token", &Method::POST) => {
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return self.handle_token_request(&mut req, session).await;
}
("introspect", &Method::POST) => {
// Authenticate request
let (_in_flight, access_token) =
self.authenticate_headers(&req, &session).await?;
return self
.handle_token_introspect(&mut req, &access_token, session.session_id)
.await;
}
("userinfo", &Method::GET) => {
// Authenticate request
let (_in_flight, access_token) =
self.authenticate_headers(&req, &session).await?;
return self
.handle_userinfo_request(access_token.account_id())
.await;
}
("register", &Method::POST) => {
return self
.handle_oauth_registration_request(&mut req, session)
.await;
}
("jwks.json", &Method::GET) => {
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return Ok(self.core.oauth.oidc_jwks.clone().into_http_response());
}
(_, &Method::OPTIONS) => {
return Ok(HttpResponse::new(StatusCode::NO_CONTENT));
}
_ => (),
},
"api" => {
// Allow CORS preflight requests
if req.method() == Method::OPTIONS {
return Ok(HttpResponse::new(StatusCode::NO_CONTENT));
}
return self.handle_api_request(&mut req, &session).await;
}
"mail" => {
if req.method() == Method::GET
&& path.next().unwrap_or_default() == "config-v1.1.xml"
{
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return self
.handle_autoconfig_request(req.uri().query())
.await
.map(|resource| resource.into_http_response());
}
}
"calendar" => {
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
if self.core.groupware.itip_http_rsvp_url.is_some()
&& req.method() == Method::GET
&& path.next().unwrap_or_default() == "rsvp"
{
return self
.http_rsvp_handle(
req.uri().query().unwrap_or_default(),
req.headers()
.get(header::ACCEPT_LANGUAGE)
.and_then(|v| v.to_str().ok())
.map(|lang| {
let lang = lang.split_once(',').map_or(lang, |(l, _)| l);
lang.split_once(';').map_or(lang, |(l, _)| l)
})
.unwrap_or("en"),
)
.await
.map(|response| {
HtmlResponse::new(response)
.into_http_response()
.with_no_store()
});
}
}
"autodiscover" | "Autodiscover" | "AutoDiscover" => {
if req.method() == Method::POST
&& path
.next()
.unwrap_or_default()
.eq_ignore_ascii_case("autodiscover.xml")
{
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return self
.handle_autodiscover_request(
fetch_body(&mut req, 8192, session.session_id).await,
)
.await
.map(|resource| resource.into_http_response());
} else if req.method() == Method::POST
&& path
.next()
.unwrap_or_default()
.eq_ignore_ascii_case("autodiscover.json")
{
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return self
.handle_autodiscover_v2_request(req.uri().query())
.await
.map(|result| match result {
Ok(resource) => resource.into_http_response(),
Err(err) => HttpResponse::new(StatusCode::BAD_REQUEST)
.with_content_type("application/json; charset=utf-8")
.with_text_body(err),
});
}
}
"robots.txt" => {
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
return Ok(
Resource::new("text/plain", b"User-agent: *\nDisallow: /\n".to_vec())
.into_http_response(),
);
}
"healthz" => {
// Limit anonymous requests
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
match path.next().unwrap_or_default() {
"live" => {
return Ok(JsonProblemResponse(StatusCode::OK).into_http_response());
}
"ready" => {
return Ok(JsonProblemResponse({
if !self.core.storage.data.is_none() {
StatusCode::OK
} else {
StatusCode::SERVICE_UNAVAILABLE
}
})
.into_http_response());
}
_ => (),
}
}
"metrics" => match path.next().unwrap_or_default() {
"prometheus" => {
if let Some(prometheus) = &self.core.metrics.prometheus {
if let Some(auth) = &prometheus.auth
&& req
.authorization_basic()
.is_none_or(|secret| secret != auth)
{
return Err(trc::AuthEvent::Failed
.into_err()
.details("Invalid or missing credentials.")
.caused_by(trc::location!()));
}
return Ok(Resource::new(
"text/plain; version=0.0.4",
self.export_prometheus_metrics().await?.into_bytes(),
)
.into_http_response());
}
}
"otel" => {
// Reserved for future use
}
_ => (),
},
// SPDX-SnippetBegin
// SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <hello@stalw.art>
// SPDX-License-Identifier: LicenseRef-SEL
#[cfg(feature = "enterprise")]
"logo" if self.is_enterprise_edition() => {
match self
.logo_resource(
req.headers()
.get(header::HOST)
.and_then(|h| h.to_str().ok())
.map(|h| h.rsplit_once(':').map_or(h, |(h, _)| h))
.unwrap_or_default(),
)
.await
{
Ok(Some(resource)) => {
return Ok(resource.into_http_response());
}
Ok(None) => {
return Err(trc::ResourceEvent::NotFound.into_err());
}
Err(err) => {
trc::error!(err.span_id(session.session_id));
}
}
}
// SPDX-SnippetEnd
"form" => {
if let Some(form) = &self.core.network.contact_form {
match *req.method() {
Method::POST => {
self.is_http_anonymous_request_allowed(session.remote_ip)
.await?;
let form_data =
FormData::from_request(&mut req, form.max_size, session.session_id)
.await?;
return self.handle_contact_form(&session, form, form_data).await;
}
Method::OPTIONS => {
return Ok(HttpResponse::new(StatusCode::NO_CONTENT));
}
_ => {}
}
}
}
"login" | "device" => {
let page = include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/../../resources/html-templates/login.html.min"
));
return Ok(HtmlResponse::new(page.to_string()).into_http_response());
}
external => {
if path.next().is_none() {
if !external.is_empty() {
return Ok(HttpResponse::redirect(format!("/{external}/")));
}
} else if let Some(resource) = self
.inner
.data
.applications
.serve(
external,
req.uri()
.path()
.get(external.len() + 2..)
.unwrap_or_default(),
)
.await?
{
let response = resource.resource.into_http_response();
return Ok(if !resource.no_cache {
response.with_immutable_cache()
} else {
response.with_no_cache()
});
}
}
}
// Block dangerous URLs
let path = req.uri().path();
if self.is_http_banned_path(path, session.remote_ip).await? {
trc::event!(
Security(SecurityEvent::ScanBan),
SpanId = session.session_id,
RemoteIp = session.remote_ip,
Path = path.to_string(),
);
}
Err(trc::ResourceEvent::NotFound.into_err())
}
}
async fn handle_session<T: SessionStream>(inner: Arc<Inner>, session: SessionData<T>) {
let _in_flight = session.in_flight;
let is_tls = session.stream.is_tls();
if let Err(http_err) = http1::Builder::new()
.keep_alive(true)
.serve_connection(
TokioIo::new(session.stream),
service_fn(|req: hyper::Request<body::Incoming>| {
let instance = session.instance.clone();
let inner = inner.clone();
async move {
let server = inner.build_server();
// Obtain remote IP
let remote_ip = if !server.core.network.http.use_forwarded {
trc::event!(
Http(trc::HttpEvent::RequestUrl),
SpanId = session.session_id,
Url = req.uri().to_string(),
);
session.remote_ip
} else if let Some(forwarded_for) = req
.headers()
.get(header::FORWARDED)
.and_then(|h| h.to_str().ok())
.and_then(|h| {
let h = h.to_ascii_lowercase();
h.split_once("for=").and_then(|(_, rest)| {
let mut start_ip = usize::MAX;
let mut end_ip = usize::MAX;
for (pos, ch) in rest.char_indices() {
match ch {
'0'..='9' | 'a'..='f' | ':' | '.' => {
if start_ip == usize::MAX {
start_ip = pos;
}
end_ip = pos;
}
'"' | '[' | ' ' if start_ip == usize::MAX => {}
_ => {
break;
}
}
}
rest.get(start_ip..=end_ip)
.and_then(|h| h.parse::<IpAddr>().ok())
})
})
.or_else(|| {
req.headers()
.get("X-Forwarded-For")
.and_then(|h| h.to_str().ok())
.map(|h| h.split_once(',').map_or(h, |(ip, _)| ip).trim())
.and_then(|h| h.parse::<IpAddr>().ok())
})
{
// Check if the forwarded IP has been blocked
if server.is_ip_blocked(forwarded_for) {
trc::event!(
Security(trc::SecurityEvent::IpBlocked),
ListenerId = instance.id.clone(),
RemoteIp = forwarded_for,
SpanId = session.session_id,
);
return Ok::<_, hyper::Error>(
JsonProblemResponse(StatusCode::FORBIDDEN)
.into_http_response()
.build(),
);
}
trc::event!(
Http(trc::HttpEvent::RequestUrl),
SpanId = session.session_id,
RemoteIp = forwarded_for,
Url = req.uri().to_string(),
);
forwarded_for
} else {
trc::event!(
Http(trc::HttpEvent::XForwardedMissing),
SpanId = session.session_id,
);
session.remote_ip
};
// Parse HTTP request
let response = match Box::pin(server.parse_http_request(
req,
HttpSessionData {
instance,
local_ip: session.local_ip,
local_port: session.local_port,
remote_ip,
remote_port: session.remote_port,
is_tls,
session_id: session.session_id,
},
))
.await
{
Ok(response) => response,
Err(err) => {
let response = err.into_http_response();
trc::error!(err.span_id(session.session_id));
response
}
};
trc::event!(
Http(trc::HttpEvent::ResponseBody),
SpanId = session.session_id,
Contents = match response.body() {
HttpResponseBody::Text(value) =>
trc::Value::String(value.as_str().into()),
HttpResponseBody::Binary(_) =>
trc::Value::String("[binary data]".into()),
HttpResponseBody::Stream(_) => trc::Value::String("[stream]".into()),
_ => trc::Value::None,
},
Code = response.status().as_u16(),
Size = response.size(),
);
// Build response
let mut response = response.build();
// Add custom headers
if !server.core.network.http.response_headers.is_empty() {
let headers = response.headers_mut();
for (header, value) in &server.core.network.http.response_headers {
headers.insert(header.clone(), value.clone());
}
}
Ok::<_, hyper::Error>(response)
}
}),
)
.with_upgrades()
.await
{
if http_err.is_parse() {
let server = inner.build_server();
if !server.core.network.http.use_forwarded {
match server.is_scanner_fail2banned(session.remote_ip).await {
Ok(true) => {
trc::event!(
Security(SecurityEvent::ScanBan),
SpanId = session.session_id,
RemoteIp = session.remote_ip,
Reason = http_err.to_string(),
);
return;
}
Ok(false) => {}
Err(err) => {
trc::error!(
err.span_id(session.session_id)
.details("Failed to check for fail2ban")
);
}
}
}
}
trc::event!(
Http(trc::HttpEvent::Error),
SpanId = session.session_id,
Reason = http_err.to_string(),
);
}
}
impl SessionManager for HttpSessionManager {
fn handle<T: SessionStream>(self, session: SessionData<T>) -> impl Future<Output = ()> + Send {
handle_session(self.inner, session)
}
#[allow(clippy::manual_async_fn)]
fn shutdown(&self) -> impl std::future::Future<Output = ()> + Send {
async {
let _ = self.inner.ipc.push_tx.send(PushEvent::Stop).await;
}
}
}