mod config;
#[cfg(not(any(feature = "transport-lettre", feature = "transport-resend")))]
compile_error!(
"at least one transport feature (`transport-lettre` or `transport-resend`) must be enabled"
);
use std::collections::BTreeMap;
use std::path::PathBuf;
use std::sync::Arc;
use anyhow::{Context as _, Result, bail};
use clap::Parser;
#[cfg(feature = "attachment-opendal")]
use email_kit::attachment::opendal::OpendalResolver;
use email_kit::attachment::{
AttachmentLimits, AttachmentResolveError, AttachmentResolver, AttachmentResolvingTransport,
ResolvedAttachment, SchemeRouter,
};
use email_kit::message::AttachmentReference;
#[cfg(feature = "transport-lettre")]
use email_kit::transport::lettre::LettreTransport;
#[cfg(feature = "transport-resend")]
use email_kit::transport::resend::ResendTransport;
use email_kit::transport::{Transport, transport_option_registry};
use figment::Figment;
use figment::providers::{Env, Format, Json, Toml, Yaml};
use restate_email::{Service, StaticTransportRegistry};
use restate_sdk::{endpoint::Endpoint, http_server::HttpServer, service::IntoServiceDefinition};
use tracing_subscriber::EnvFilter;
#[cfg(feature = "attachment-opendal")]
use crate::config::ResolverConfig;
use crate::config::{AttachmentConfig, Config, TransportConfig};
#[derive(Parser, Debug)]
#[command(version)]
struct Cli {
#[arg(long, value_name = "FILE", env = "CONFIG_FILE")]
config: Option<PathBuf>,
#[arg(long, default_value = "9080", env = "PORT")]
port: u16,
}
impl Cli {
fn load_config(&self) -> Result<Config> {
let mut figment = Figment::new();
if let Some(path) = self.config.as_deref() {
if !path.exists() {
bail!("config file not found: {}", path.display());
}
figment = match path.extension().and_then(|extension| extension.to_str()) {
Some("toml") => figment.merge(Toml::file(path)),
Some("json") => figment.merge(Json::file(path)),
Some("yaml" | "yml") => figment.merge(Yaml::file(path)),
_ => bail!("unsupported config file format; use .toml, .json, .yaml, or .yml"),
};
}
figment = figment.merge(Env::prefixed("RESTATE_EMAIL_").split("__"));
figment.extract().context("failed to parse configuration")
}
}
#[derive(Clone)]
struct SharedAttachmentResolver(Arc<SchemeRouter>);
impl AttachmentResolver for SharedAttachmentResolver {
async fn resolve(
&self,
reference: &AttachmentReference,
) -> Result<ResolvedAttachment, AttachmentResolveError> {
self.0.resolve(reference).await
}
}
#[tokio::main]
async fn main() -> Result<()> {
let cli = Cli::parse();
let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info"));
tracing_subscriber::fmt().with_env_filter(filter).init();
let Config {
transports,
attachments,
identity_keys,
} = cli.load_config()?;
let registry = create_registry(transports, attachments)?;
let option_registry = transport_option_registry();
for provider in option_registry.provider_keys() {
tracing::info!(provider, "registered transport option provider");
}
let service = Service::new(registry)
.with_transport_options(option_registry)
.into_service_definition();
let mut endpoint = Endpoint::builder().bind(service);
for identity_key in &identity_keys {
endpoint = endpoint
.identity_key(identity_key)
.with_context(|| format!("invalid Restate identity key `{identity_key}`"))?;
}
if !identity_keys.is_empty() {
tracing::info!(
keys = identity_keys.len(),
"request identity verification enabled"
);
}
let bind_addr = format!("0.0.0.0:{}", cli.port);
tracing::info!(%bind_addr, "starting Restate email endpoint");
HttpServer::new(endpoint.build())
.listen_and_serve(bind_addr.parse()?)
.await;
Ok(())
}
fn create_registry(
transports: BTreeMap<String, TransportConfig>,
attachments: Option<AttachmentConfig>,
) -> Result<StaticTransportRegistry> {
if transports.is_empty() {
bail!("at least one transport must be configured");
}
let attachment_preparation = attachments.map(create_attachment_preparation).transpose()?;
let mut registry = StaticTransportRegistry::new();
for (key, transport) in transports {
let provider = transport.provider_name();
match transport {
#[cfg(feature = "transport-resend")]
TransportConfig::Resend { api_key, base_url } => {
let mut builder = ResendTransport::builder(api_key);
if let Some(base_url) = base_url {
builder = builder.base_url(base_url);
}
insert_transport(
&mut registry,
&key,
builder.build(),
attachment_preparation.as_ref(),
);
}
#[cfg(feature = "transport-lettre")]
TransportConfig::Smtp { url } => {
let transport = LettreTransport::from_url(url.as_str()).with_context(|| {
format!("failed to configure SMTP transport `{key}` from its connection URL")
})?;
insert_transport(
&mut registry,
&key,
transport,
attachment_preparation.as_ref(),
);
}
}
tracing::info!(transport = %key, provider, "registered email transport");
}
Ok(registry)
}
fn insert_transport<T>(
registry: &mut StaticTransportRegistry,
key: &str,
transport: T,
attachment_preparation: Option<&(Arc<SchemeRouter>, AttachmentLimits)>,
) where
T: Transport + 'static,
{
if let Some((resolver, limits)) = attachment_preparation {
registry.insert(
key,
AttachmentResolvingTransport::new(
transport,
SharedAttachmentResolver(Arc::clone(resolver)),
)
.with_limits(*limits),
);
} else {
registry.insert(key, transport);
}
}
#[cfg_attr(
not(feature = "attachment-opendal"),
allow(clippy::needless_pass_by_value, clippy::unnecessary_wraps)
)]
fn create_attachment_preparation(
config: AttachmentConfig,
) -> Result<(Arc<SchemeRouter>, AttachmentLimits)> {
let limits = AttachmentLimits::new()
.with_max_attachment_bytes(config.max_attachment_bytes)
.with_max_total_bytes(config.max_total_bytes);
let router = SchemeRouter::new();
#[cfg(feature = "attachment-opendal")]
let router = {
let mut router = router;
let resolver_max_bytes = config.max_attachment_bytes.unwrap_or(usize::MAX);
for (scheme, resolver) in config.resolvers {
match resolver {
ResolverConfig::Opendal { service, options } => {
let operator = opendal::Operator::via_iter(&service, options)
.with_context(|| {
format!(
"failed to configure attachment resolver `{scheme}` with OpenDAL service `{service}`"
)
})?;
router.register(&scheme, OpendalResolver::new(operator, resolver_max_bytes));
tracing::info!(resolver = %scheme, service, "registered attachment resolver");
}
}
}
router
};
Ok((Arc::new(router), limits))
}
#[cfg(all(test, feature = "attachment-opendal", feature = "transport-resend"))]
mod tests {
use std::collections::BTreeMap;
use email_kit::attachment::ResolveErrorKind;
use email_kit::message::{
Address, Attachment, AttachmentReference, Body, ContentType, Message,
};
use restate_email::{SendOptions, SendRequest, TransportKey};
use serde_json::json;
use tempfile::tempdir;
use wiremock::matchers::{body_partial_json, method, path};
use wiremock::{Mock, MockServer, ResponseTemplate};
use super::*;
use crate::config::{AttachmentConfig, ResolverConfig};
fn reference_request(reference: &str) -> SendRequest {
let message = Message::builder(Body::text("body"))
.from_mailbox("sender@example.com".parse().expect("sender should parse"))
.to(vec![Address::Mailbox(
"recipient@example.com"
.parse()
.expect("recipient should parse"),
)])
.subject("Reference attachment")
.add_attachment(
Attachment::reference(
ContentType::try_from("application/octet-stream")
.expect("content type should parse"),
AttachmentReference::new(reference),
)
.with_filename("report.bin"),
)
.build_outbound()
.expect("message should validate");
SendRequest {
transport: TransportKey::new_unchecked("transactional"),
message,
options: SendOptions::default(),
}
}
fn byte_request(bytes: &[u8]) -> SendRequest {
let message = Message::builder(Body::text("body"))
.from_mailbox("sender@example.com".parse().expect("sender should parse"))
.to(vec![Address::Mailbox(
"recipient@example.com"
.parse()
.expect("recipient should parse"),
)])
.subject("Byte attachment")
.add_attachment(
Attachment::bytes(
ContentType::try_from("application/octet-stream")
.expect("content type should parse"),
bytes.to_vec(),
)
.with_filename("report.bin"),
)
.build_outbound()
.expect("message should validate");
SendRequest {
transport: TransportKey::new_unchecked("transactional"),
message,
options: SendOptions::default(),
}
}
fn resend_transport_config(server: &MockServer) -> BTreeMap<String, TransportConfig> {
BTreeMap::from([(
String::from("transactional"),
TransportConfig::Resend {
api_key: String::from("test-key"),
base_url: Some(
format!("{}/", server.uri())
.parse()
.expect("mock URL should parse"),
),
},
)])
}
struct TransientResolver;
impl AttachmentResolver for TransientResolver {
async fn resolve(
&self,
_reference: &AttachmentReference,
) -> Result<ResolvedAttachment, AttachmentResolveError> {
Err(AttachmentResolveError::new(
ResolveErrorKind::Transient,
"attachment store is temporarily unavailable",
))
}
}
#[tokio::test]
async fn configured_resolver_delivers_reference_as_bytes() {
let server = MockServer::start().await;
let payload: &[u8] = b"resolved attachment";
Mock::given(method("POST"))
.and(path("/emails"))
.and(body_partial_json(json!({
"attachments": [{
"filename": "report.bin",
"content": payload,
"contentType": "application/octet-stream",
}]
})))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({"id": "resolved-1"})))
.mount(&server)
.await;
let directory = tempdir().expect("temporary directory should be created");
std::fs::write(directory.path().join("report.bin"), payload)
.expect("attachment fixture should be written");
let mut resolvers = BTreeMap::new();
resolvers.insert(
String::from("docs"),
ResolverConfig::Opendal {
service: String::from("fs"),
options: BTreeMap::from([(
String::from("root"),
directory.path().display().to_string(),
)]),
},
);
let service = Service::new(
create_registry(
resend_transport_config(&server),
Some(AttachmentConfig {
max_attachment_bytes: Some(1024),
max_total_bytes: Some(2048),
resolvers,
}),
)
.expect("registry should build"),
);
let response = service
.send_request(&reference_request("docs:report.bin"))
.await
.expect("reference-backed request should send");
assert_eq!(
response.report.provider_message_id.as_deref(),
Some("resolved-1")
);
}
#[tokio::test]
async fn absent_preparation_preserves_bytes_and_rejects_references_terminally() {
let server = MockServer::start().await;
let payload: &[u8] = b"inline attachment";
Mock::given(method("POST"))
.and(path("/emails"))
.and(body_partial_json(json!({
"attachments": [{
"filename": "report.bin",
"content": payload,
"contentType": "application/octet-stream",
}]
})))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({"id": "bytes-1"})))
.expect(1)
.mount(&server)
.await;
let service = Service::new(
create_registry(resend_transport_config(&server), None).expect("registry should build"),
);
let response = service
.send_request(&byte_request(payload))
.await
.expect("byte-backed request should send");
assert_eq!(
response.report.provider_message_id.as_deref(),
Some("bytes-1")
);
let error = service
.send_request(&reference_request("docs:report.bin"))
.await
.expect_err("unresolved reference should fail");
let source: &(dyn std::error::Error + Send + Sync + 'static) = error.as_ref();
assert!(
source.to_string().starts_with("Terminal error [400]"),
"unsupported references should be terminal: {source}"
);
}
#[tokio::test]
async fn resolver_failures_keep_the_existing_handler_dispositions() {
let server = MockServer::start().await;
let directory = tempdir().expect("temporary directory should be created");
std::fs::write(directory.path().join("large.bin"), b"five!")
.expect("attachment fixture should be written");
let resolvers = BTreeMap::from([(
String::from("docs"),
ResolverConfig::Opendal {
service: String::from("fs"),
options: BTreeMap::from([(
String::from("root"),
directory.path().display().to_string(),
)]),
},
)]);
let service = Service::new(
create_registry(
resend_transport_config(&server),
Some(AttachmentConfig {
max_attachment_bytes: Some(4),
max_total_bytes: Some(8),
resolvers,
}),
)
.expect("registry should build"),
);
for reference in ["docs:missing.bin", "unknown:large.bin", "docs:large.bin"] {
let error = service
.send_request(&reference_request(reference))
.await
.expect_err("resolver failure should fail the send");
let source: &(dyn std::error::Error + Send + Sync + 'static) = error.as_ref();
assert!(
source.to_string().starts_with("Terminal error"),
"{reference} should fail terminally: {source}"
);
}
let mut registry = StaticTransportRegistry::new();
registry.insert(
"transactional",
AttachmentResolvingTransport::new(
ResendTransport::builder("test-key").build(),
TransientResolver,
),
);
let transient_service = Service::new(registry);
let error = transient_service
.send_request(&reference_request("docs:report.bin"))
.await
.expect_err("transient resolver failure should fail the send");
let source: &(dyn std::error::Error + Send + Sync + 'static) = error.as_ref();
assert!(
source.to_string().starts_with("Retryable error"),
"transient resolver failures should remain retryable: {source}"
);
}
}