mirror of
https://github.com/chatmail/relay.git
synced 2026-08-11 02:50:53 +00:00
feat: https transport channel (#122)
Transport mode: Implements an additional mail delivery channel over HTTPS. Incoming mode: Adds a http server listening for incoming messages delivered over HTTPS. Signed-off-by: Jagoda Ślązak <jslazak@jslazak.com>
This commit is contained in:
committed by
GitHub
parent
1b1585eb7d
commit
0a652e61ee
@@ -14,6 +14,8 @@ pub struct Config {
|
||||
pub filtermail_smtp_port: u16,
|
||||
#[serde(default = "Config::default_filtermail_smtp_port_incoming")]
|
||||
pub filtermail_smtp_port_incoming: u16,
|
||||
#[serde(default = "Config::default_filtermail_http_port_incoming")]
|
||||
pub filtermail_http_port_incoming: u16,
|
||||
#[serde(default = "Config::default_filtermail_lmtp_port_transport")]
|
||||
pub filtermail_lmtp_port_transport: u16,
|
||||
#[serde(default = "Config::default_postfix_host")]
|
||||
@@ -110,6 +112,9 @@ impl Config {
|
||||
const fn default_filtermail_smtp_port_incoming() -> u16 {
|
||||
10081
|
||||
}
|
||||
const fn default_filtermail_http_port_incoming() -> u16 {
|
||||
10082
|
||||
}
|
||||
const fn default_filtermail_lmtp_port_transport() -> u16 {
|
||||
10083
|
||||
}
|
||||
@@ -143,6 +148,7 @@ impl Default for Config {
|
||||
filtermail_host: Self::default_filtermail_host(),
|
||||
filtermail_smtp_port: Self::default_filtermail_smtp_port(),
|
||||
filtermail_smtp_port_incoming: Self::default_filtermail_smtp_port_incoming(),
|
||||
filtermail_http_port_incoming: Self::default_filtermail_http_port_incoming(),
|
||||
filtermail_lmtp_port_transport: Self::default_filtermail_lmtp_port_transport(),
|
||||
postfix_host: Self::default_postfix_host(),
|
||||
postfix_reinject_port: Self::default_postfix_reinject_port(),
|
||||
|
||||
@@ -27,6 +27,12 @@ pub enum Error {
|
||||
Tls(#[from] rustls::Error),
|
||||
#[error(transparent)]
|
||||
InvalidDnsName(#[from] rustls::pki_types::InvalidDnsNameError),
|
||||
#[error(transparent)]
|
||||
Hyper(#[from] hyper::Error),
|
||||
#[error(transparent)]
|
||||
HyperHttp(#[from] hyper::http::Error),
|
||||
#[error(transparent)]
|
||||
HyperClient(#[from] hyper_util::client::legacy::Error),
|
||||
}
|
||||
|
||||
impl Error {
|
||||
|
||||
@@ -0,0 +1,150 @@
|
||||
use crate::smtp_server::{Envelope, SmtpHandler};
|
||||
use http_body_util::combinators::BoxBody;
|
||||
use http_body_util::{BodyExt, Full};
|
||||
use hyper::body::{Bytes, Incoming};
|
||||
use hyper::service::Service;
|
||||
use hyper::{Request, Response};
|
||||
use hyper_util::rt::TokioIo;
|
||||
use std::convert::Infallible;
|
||||
use std::pin::Pin;
|
||||
use std::sync::Arc;
|
||||
use tokio::net::{TcpListener, TcpStream};
|
||||
|
||||
/// Runs the HTTP server on the specified address with the given handler and maximum message size.
|
||||
pub async fn run_http_server<H>(
|
||||
addr: &impl tokio::net::ToSocketAddrs,
|
||||
handler: Arc<H>,
|
||||
max_size: usize,
|
||||
) -> Result<(), crate::error::Error>
|
||||
where
|
||||
H: SmtpHandler + 'static,
|
||||
{
|
||||
let listener = TcpListener::bind(addr).await?;
|
||||
loop {
|
||||
let (socket, _) = listener.accept().await?;
|
||||
|
||||
// Disable Nagle's algorithm.
|
||||
socket.set_nodelay(true)?;
|
||||
|
||||
let handler = handler.clone();
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = handle_connection(socket, handler, max_size).await {
|
||||
log::error!("Error handling connection: {e}");
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/// Handles a single HTTP connection.
|
||||
async fn handle_connection<H>(
|
||||
socket: TcpStream,
|
||||
handler: Arc<H>,
|
||||
max_size: usize,
|
||||
) -> Result<(), String>
|
||||
where
|
||||
H: SmtpHandler + 'static,
|
||||
{
|
||||
let service = MxDelivService::new(handler, max_size);
|
||||
|
||||
hyper_util::server::conn::auto::Builder::new(hyper_util::rt::TokioExecutor::new())
|
||||
.serve_connection(TokioIo::new(socket), service)
|
||||
.await
|
||||
.map_err(|e| e.to_string())?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
struct MxDelivService<H: SmtpHandler> {
|
||||
handler: Arc<H>,
|
||||
max_size: usize,
|
||||
}
|
||||
|
||||
impl<H: SmtpHandler> MxDelivService<H> {
|
||||
/// Creates a new [`MxDelivService`].
|
||||
fn new(handler: Arc<H>, max_size: usize) -> Self {
|
||||
Self { handler, max_size }
|
||||
}
|
||||
}
|
||||
|
||||
impl<H: SmtpHandler + 'static> Service<Request<Incoming>> for MxDelivService<H> {
|
||||
type Response = Response<BoxBody<Bytes, Infallible>>;
|
||||
type Error = crate::error::Error;
|
||||
type Future = Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send>>;
|
||||
|
||||
fn call(&self, req: Request<Incoming>) -> Self::Future {
|
||||
let handler = self.handler.clone();
|
||||
let max_size = self.max_size;
|
||||
|
||||
let fut = async move {
|
||||
if req.method() != hyper::Method::POST {
|
||||
return Ok(Response::builder().status(405).body(
|
||||
// This is client's implementation error if it happens,
|
||||
// so we don't care about sending a proper smtp response.
|
||||
Full::new(Bytes::from("Method Not Allowed")).boxed(),
|
||||
)?);
|
||||
}
|
||||
|
||||
let mail_from = req
|
||||
.headers()
|
||||
.get(crate::transport::HEADER_MAIL_FROM)
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.unwrap_or("")
|
||||
.to_string();
|
||||
|
||||
match handler.handle_mail(&mail_from) {
|
||||
Ok(_) => {}
|
||||
Err(e) => {
|
||||
return Ok(Response::builder()
|
||||
.status(400)
|
||||
.body(Full::new(Bytes::from(e)).boxed())?);
|
||||
}
|
||||
};
|
||||
|
||||
let rcpt_to = req
|
||||
.headers()
|
||||
.get_all(crate::transport::HEADER_RCPT_TO)
|
||||
.iter()
|
||||
.filter_map(|v| v.to_str().ok())
|
||||
.map(ToString::to_string)
|
||||
.collect();
|
||||
|
||||
let body_limited = http_body_util::Limited::new(req.into_body(), max_size);
|
||||
let body_bytes = match body_limited.collect().await {
|
||||
Ok(body) => body.to_bytes(),
|
||||
Err(_) => {
|
||||
return Ok(Response::builder().status(413).body(
|
||||
Full::new(Bytes::from("552 Message exceeds maximum size")).boxed(),
|
||||
)?);
|
||||
}
|
||||
};
|
||||
|
||||
let mut envelope = Envelope {
|
||||
origin_ip: "".to_string(),
|
||||
mail_from,
|
||||
rcpt_to,
|
||||
data: body_bytes.to_vec(),
|
||||
};
|
||||
|
||||
log::debug!("(HTTP) MAIL FROM:<{}>", envelope.mail_from);
|
||||
for rcpt in &envelope.rcpt_to {
|
||||
log::debug!("(HTTP) RCPT TO:<{}>", rcpt);
|
||||
}
|
||||
|
||||
log::trace!(
|
||||
"(HTTP) DATA:\n{:?}",
|
||||
String::from_utf8_lossy(&envelope.data)
|
||||
);
|
||||
|
||||
match handler.handle_data(&mut envelope).await {
|
||||
Ok(response) => Ok(Response::builder()
|
||||
.status(201)
|
||||
.body(Full::new(Bytes::from(response)).boxed())?),
|
||||
Err(e) => Ok(Response::builder()
|
||||
.status(400)
|
||||
.body(Full::new(Bytes::from(e)).boxed())?),
|
||||
}
|
||||
};
|
||||
|
||||
Box::pin(fut)
|
||||
}
|
||||
}
|
||||
@@ -27,6 +27,7 @@
|
||||
mod config;
|
||||
mod dkim_verifier;
|
||||
pub(crate) mod error;
|
||||
mod http_server;
|
||||
pub(crate) mod inbound;
|
||||
pub(crate) mod message;
|
||||
pub(crate) mod openpgp;
|
||||
@@ -37,6 +38,7 @@ mod tls;
|
||||
mod transport;
|
||||
pub(crate) mod utils;
|
||||
|
||||
use crate::http_server::run_http_server;
|
||||
use crate::transport::TransportHandler;
|
||||
use config::Config;
|
||||
use env_logger::Env;
|
||||
@@ -153,6 +155,17 @@ async fn main() -> Result<(), error::Error> {
|
||||
addr_smtp.1
|
||||
);
|
||||
|
||||
let addr_http = (config.filtermail_host, config.filtermail_http_port_incoming);
|
||||
let handler_http = handler.clone();
|
||||
|
||||
server_set
|
||||
.spawn(async move { run_http_server(&addr_http, handler_http, max_size).await });
|
||||
log::debug!(
|
||||
"Incoming HTTP server listening on {}:{}",
|
||||
addr_http.0,
|
||||
addr_http.1
|
||||
);
|
||||
|
||||
while let Some(result) = server_set.join_next().await {
|
||||
if let Err(e) = result {
|
||||
eprintln!("Server error: {}", e);
|
||||
|
||||
+14
-6
@@ -16,8 +16,19 @@ pub async fn wrap_rustls<IO>(
|
||||
where
|
||||
IO: AsyncRead + AsyncWrite + Unpin,
|
||||
{
|
||||
let mut root_cert_store = rustls::RootCertStore::empty();
|
||||
root_cert_store.extend(webpki_roots::TLS_SERVER_ROOTS.iter().cloned());
|
||||
let config = configure_rustls(resumption_store, dangerous_no_cert_verification);
|
||||
let tls = tokio_rustls::TlsConnector::from(Arc::new(config));
|
||||
let name = rustls::pki_types::ServerName::try_from(hostname)?.to_owned();
|
||||
let tls_stream = tls.connect(name, stream).await?;
|
||||
Ok(tls_stream.into())
|
||||
}
|
||||
|
||||
pub fn configure_rustls(
|
||||
resumption_store: Arc<ClientSessionMemoryCache>,
|
||||
dangerous_no_cert_verification: bool,
|
||||
) -> rustls::ClientConfig {
|
||||
let root_cert_store =
|
||||
rustls::RootCertStore::from_iter(webpki_roots::TLS_SERVER_ROOTS.iter().cloned());
|
||||
|
||||
let mut config = rustls::ClientConfig::builder()
|
||||
.with_root_certificates(root_cert_store)
|
||||
@@ -39,8 +50,5 @@ where
|
||||
.set_certificate_verifier(Arc::new(NoCertificateVerification::default()));
|
||||
}
|
||||
|
||||
let tls = tokio_rustls::TlsConnector::from(Arc::new(config));
|
||||
let name = rustls::pki_types::ServerName::try_from(hostname)?.to_owned();
|
||||
let tls_stream = tls.connect(name, stream).await?;
|
||||
Ok(tls_stream.into())
|
||||
config
|
||||
}
|
||||
|
||||
+172
-3
@@ -1,45 +1,123 @@
|
||||
use crate::config::Config;
|
||||
use crate::smtp_client::{SmtpConnectionPool, TlsConfig};
|
||||
use crate::smtp_server::{Envelope, SmtpHandler};
|
||||
use crate::tls;
|
||||
use crate::utils::{AddressDomain, build_resolver};
|
||||
use async_trait::async_trait;
|
||||
use hickory_resolver::TokioResolver;
|
||||
use http_body_util::BodyExt;
|
||||
use hyper::body::Bytes;
|
||||
use hyper_rustls::HttpsConnector;
|
||||
use hyper_util::client::legacy::connect::HttpConnector;
|
||||
use std::collections::BTreeMap;
|
||||
use std::str::FromStr;
|
||||
use std::sync::Arc;
|
||||
use tokio::task::JoinSet;
|
||||
use std::time::Duration;
|
||||
use tokio::task::{JoinHandle, JoinSet};
|
||||
use tokio_rustls::rustls;
|
||||
|
||||
pub const HEADER_MAIL_FROM: &str = "X-MAIL-FROM";
|
||||
pub const HEADER_RCPT_TO: &str = "X-MAIL-TO";
|
||||
|
||||
/// Cheaply clonable HTTPS client.
|
||||
///
|
||||
/// Holds regular secure variant and relaxed - without certificate verification.
|
||||
///
|
||||
/// Connection pool handled internally by [`hyper_util::client::legacy::Client`].
|
||||
#[derive(Clone)]
|
||||
struct HttpsClient {
|
||||
pub secure: hyper_util::client::legacy::Client<
|
||||
HttpsConnector<HttpConnector>,
|
||||
http_body_util::Full<Bytes>,
|
||||
>,
|
||||
pub relaxed: hyper_util::client::legacy::Client<
|
||||
HttpsConnector<HttpConnector>,
|
||||
http_body_util::Full<Bytes>,
|
||||
>,
|
||||
}
|
||||
|
||||
impl HttpsClient {
|
||||
/// Creates a new `[HttpsClient]`.
|
||||
pub fn new(tls_resumption_store: Arc<rustls::client::ClientSessionMemoryCache>) -> Self {
|
||||
let tls_client_config = tls::configure_rustls(tls_resumption_store.clone(), false);
|
||||
let https_connector = hyper_rustls::HttpsConnectorBuilder::new()
|
||||
.with_tls_config(tls_client_config)
|
||||
.https_only()
|
||||
.enable_http1()
|
||||
.enable_http2()
|
||||
.build();
|
||||
let https_client =
|
||||
hyper_util::client::legacy::Client::builder(hyper_util::rt::TokioExecutor::new())
|
||||
.build(https_connector);
|
||||
|
||||
let tls_client_config_relaxed = tls::configure_rustls(tls_resumption_store, true);
|
||||
let https_connector_relaxed = hyper_rustls::HttpsConnectorBuilder::new()
|
||||
.with_tls_config(tls_client_config_relaxed)
|
||||
.https_only()
|
||||
.enable_http1()
|
||||
.enable_http2()
|
||||
.build();
|
||||
let https_client_relaxed =
|
||||
hyper_util::client::legacy::Client::builder(hyper_util::rt::TokioExecutor::new())
|
||||
.build(https_connector_relaxed);
|
||||
|
||||
Self {
|
||||
secure: https_client,
|
||||
relaxed: https_client_relaxed,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub struct TransportHandler {
|
||||
config: Config,
|
||||
dns_resolver: Arc<TokioResolver>,
|
||||
tls_resumption_store: Arc<rustls::client::ClientSessionMemoryCache>,
|
||||
smtp_connection_pool: Arc<SmtpConnectionPool>,
|
||||
https_client: HttpsClient,
|
||||
mxdeliv_unsupported_hosts: Arc<retainer::Cache<String, bool>>,
|
||||
monitor_handle: JoinHandle<()>,
|
||||
}
|
||||
|
||||
impl TransportHandler {
|
||||
pub fn new(config: Config) -> Result<Self, crate::error::Error> {
|
||||
let dns_resolver = Arc::new(build_resolver()?);
|
||||
let tls_resumption_store = Arc::new(rustls::client::ClientSessionMemoryCache::new(256));
|
||||
let https_client = HttpsClient::new(tls_resumption_store.clone());
|
||||
|
||||
let mxdeliv_cache = Arc::new(retainer::Cache::new());
|
||||
let mxdeliv_cache_clone = mxdeliv_cache.clone();
|
||||
|
||||
let monitor_handle = tokio::spawn(async move {
|
||||
mxdeliv_cache_clone
|
||||
.monitor(4, 0.25, Duration::from_secs(10))
|
||||
.await
|
||||
});
|
||||
|
||||
Ok(Self {
|
||||
config,
|
||||
dns_resolver,
|
||||
tls_resumption_store,
|
||||
smtp_connection_pool: SmtpConnectionPool::new(),
|
||||
https_client,
|
||||
mxdeliv_unsupported_hosts: mxdeliv_cache,
|
||||
monitor_handle,
|
||||
})
|
||||
}
|
||||
|
||||
/// Handles a single email transaction for a single recipient domain.
|
||||
#[expect(clippy::too_many_arguments)]
|
||||
async fn handle_single_domain(
|
||||
tls_resumption_store: Arc<rustls::client::ClientSessionMemoryCache>,
|
||||
smtp_connection_pool: Arc<SmtpConnectionPool>,
|
||||
mxdeliv_unsupported_hosts: Arc<retainer::Cache<String, bool>>,
|
||||
https_client: HttpsClient,
|
||||
dns_resolver: Arc<TokioResolver>,
|
||||
domain: AddressDomain,
|
||||
envelope: Envelope,
|
||||
client_hostname: String,
|
||||
) -> Result<String, String> {
|
||||
let mut allow_invalid_cert = false;
|
||||
let mut skip_tls = false;
|
||||
let mut skip_tls = false; // only respected by smtp channel
|
||||
|
||||
let mx_hosts = match domain {
|
||||
// no-DNS setup; assume the ip from email address is the destination.
|
||||
@@ -110,6 +188,34 @@ impl TransportHandler {
|
||||
// but the IPv4 and IPv6 connections (after `smtp_client::send` resolves mx hostname)
|
||||
// happens in parallel.
|
||||
'try_relay: for (_, mx_host) in mx_hosts {
|
||||
let skip_mxdeliv = mxdeliv_unsupported_hosts
|
||||
.get(&mx_host)
|
||||
.await
|
||||
.map(|guard| *guard.value())
|
||||
.unwrap_or(false);
|
||||
|
||||
// HTTPS channel
|
||||
if skip_mxdeliv {
|
||||
log::debug!("Skipping HTTP delivery to host that failed recently: {mx_host}");
|
||||
} else {
|
||||
match Self::https_delivery(
|
||||
https_client.clone(),
|
||||
mx_host.clone(),
|
||||
&envelope,
|
||||
allow_invalid_cert,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => {
|
||||
return Ok("250 Ok (HTTPS)".to_string());
|
||||
}
|
||||
Err(e) => {
|
||||
log::debug!("HTTPS delivery to {mx_host} failed: {e}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// SMTP channel (fallback)
|
||||
match crate::smtp_client::send(
|
||||
&mx_host,
|
||||
25,
|
||||
@@ -122,7 +228,14 @@ impl TransportHandler {
|
||||
.await
|
||||
{
|
||||
Ok(_) => {
|
||||
return Ok("250 Ok".to_string());
|
||||
// Switches this host to SMTP for 30 minutes.
|
||||
// Note: this MUST happen only after a successful SMTP delivery,
|
||||
// or otherwise any http error will lock us out of any way to
|
||||
// deliver to a relay with a blocked port 25 for 30 minutes.
|
||||
mxdeliv_unsupported_hosts
|
||||
.insert(mx_host.clone(), true, Duration::from_mins(30))
|
||||
.await;
|
||||
return Ok("250 Ok (SMTP)".to_string());
|
||||
}
|
||||
Err(error) => {
|
||||
match error {
|
||||
@@ -152,6 +265,60 @@ impl TransportHandler {
|
||||
|
||||
Err("421 Failed to connect to any mail server".to_string())
|
||||
}
|
||||
|
||||
/// Performs mail delivery to `mx_host` over HTTPS.
|
||||
///
|
||||
/// Times out after 60s.
|
||||
async fn https_delivery(
|
||||
https_client: HttpsClient,
|
||||
mx_host: String,
|
||||
envelope: &Envelope,
|
||||
allow_invalid_cert: bool,
|
||||
) -> Result<(), crate::error::Error> {
|
||||
let request: hyper::Request<http_body_util::Full<Bytes>> = {
|
||||
let mut builder = hyper::Request::builder()
|
||||
.method(hyper::Method::POST)
|
||||
.uri(format!("https://{mx_host}/mxdeliv"));
|
||||
|
||||
if !envelope.mail_from.is_empty() {
|
||||
builder = builder.header(HEADER_MAIL_FROM, &envelope.mail_from);
|
||||
}
|
||||
|
||||
for rcpt_to in &envelope.rcpt_to {
|
||||
builder = builder.header(HEADER_RCPT_TO, rcpt_to);
|
||||
}
|
||||
|
||||
builder.body(http_body_util::Full::from(envelope.data.clone()))?
|
||||
};
|
||||
|
||||
let client = if allow_invalid_cert {
|
||||
https_client.relaxed
|
||||
} else {
|
||||
https_client.secure
|
||||
};
|
||||
|
||||
let response = tokio::time::timeout(Duration::from_secs(60), client.request(request))
|
||||
.await
|
||||
.map_err(|_| crate::error::Error::MailSend {
|
||||
context: "HTTPS delivery".to_string(),
|
||||
raw_smtp_answer: "[timeout]".to_string(),
|
||||
})??;
|
||||
if response.status().is_success() {
|
||||
Ok(())
|
||||
} else {
|
||||
let response_body = response.collect().await?.to_bytes();
|
||||
Err(crate::error::Error::MailSend {
|
||||
context: "HTTPS delivery".to_string(),
|
||||
raw_smtp_answer: String::from_utf8_lossy(&response_body).into(),
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for TransportHandler {
|
||||
fn drop(&mut self) {
|
||||
self.monitor_handle.abort();
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
@@ -201,6 +368,8 @@ impl SmtpHandler for TransportHandler {
|
||||
.spawn(Self::handle_single_domain(
|
||||
self.tls_resumption_store.clone(),
|
||||
self.smtp_connection_pool.clone(),
|
||||
self.mxdeliv_unsupported_hosts.clone(),
|
||||
self.https_client.clone(),
|
||||
self.dns_resolver.clone(),
|
||||
rcpt_domain.clone(),
|
||||
domain_envelope,
|
||||
|
||||
Reference in New Issue
Block a user