use std::sync::Arc;
use async_trait::async_trait;
use dataflow_rs::engine::error::DataflowError;
use dataflow_rs::engine::functions::AsyncFunctionHandler;
use dataflow_rs::engine::task_context::TaskContext;
use dataflow_rs::engine::task_outcome::TaskOutcome;
use mail_builder::MessageBuilder;
use mail_builder::headers::address::Address;
use mail_builder::headers::raw::Raw;
use mail_send::smtp::AssertReply;
use serde_json::{Value, json};
use super::connector_helpers::{
ConnectorCall, apply_output, require_smtp_connector, resolve_optional_str, resolve_value,
to_connect_error,
};
use super::schema::{FieldKind, FieldSchema};
use crate::connector::smtp_pool::{PooledClient, SmtpPoolCache};
use crate::connector::{ConnectorRegistry, SmtpConnectorConfig};
const NAME: &str = "send_email";
pub(crate) const PROTECTED_HEADERS: &[&str] = &[
"from",
"to",
"cc",
"bcc",
"subject",
"date",
"message-id",
"content-type",
"content-transfer-encoding",
"mime-version",
"reply-to",
];
pub struct SendEmailHandler {
pub registry: Arc<ConnectorRegistry>,
pub smtp_pool: Arc<SmtpPoolCache>,
}
#[async_trait]
impl AsyncFunctionHandler for SendEmailHandler {
type Input = Value;
async fn execute(
&self,
ctx: &mut TaskContext<'_>,
input: &Value,
) -> dataflow_rs::Result<TaskOutcome> {
let call = ConnectorCall::begin(NAME, input, ctx)?;
check_headers_field(input)?;
let to = address_list(input, "to", ctx)?
.ok_or_else(|| validation("requires 'to' (an address or an array of addresses)"))?;
let cc = address_list(input, "cc", ctx)?.unwrap_or_default();
let bcc = address_list(input, "bcc", ctx)?.unwrap_or_default();
let subject = resolve_optional_str(input, "subject", NAME, ctx)?
.ok_or_else(|| validation("requires 'subject'"))?;
let text = resolve_optional_str(input, "text", NAME, ctx)?;
let html = resolve_optional_str(input, "html", NAME, ctx)?;
if text.is_none() && html.is_none() {
return Err(validation("requires at least one of 'text' or 'html'"));
}
let from_override = resolve_optional_str(input, "from", NAME, ctx)?;
let reply_to = resolve_optional_str(input, "reply_to", NAME, ctx)?
.map(|s| parse_mailbox("reply_to", &s))
.transpose()?;
call.run(&self.registry, async {
let connector_config = call.resolve(&self.registry, None).await?;
let smtp_config = require_smtp_connector(&connector_config, call.connector)?;
let from = sender(smtp_config, from_override.as_deref(), call.connector)?;
let bare_id = generated_message_id(&from);
let message = build_message(MessageParts {
from: &from,
message_id: &bare_id,
to: &to,
cc: &cc,
subject: &subject,
reply_to,
headers: input.get("headers"),
text,
html,
})?;
let recipients: Vec<&str> = to
.iter()
.chain(cc.iter())
.chain(bcc.iter())
.map(|mbox| mbox.email.as_str())
.collect();
let pool = self
.smtp_pool
.get_pool(call.connector, smtp_config)
.await
.map_err(to_connect_error)?;
let mut client = pool.checkout().await.map_err(to_connect_error)?;
let response = match deliver(&mut client, &from, &recipients, &message).await {
Ok(reply) => {
pool.checkin(client).await;
reply
}
Err(e) => {
return Err(DataflowError::function_execution(
format!("SMTP send via '{}' failed: {e}", call.connector),
None,
));
}
};
let result = json!({
"message_id": format!("<{bare_id}>"),
"response": response,
});
apply_output(ctx, call.output, result);
Ok(TaskOutcome::Success)
})
.await
}
}
fn validation(msg: &str) -> DataflowError {
DataflowError::Validation(format!("{NAME}: {msg}"))
}
fn address_list(
input: &Value,
field: &str,
ctx: &TaskContext<'_>,
) -> Result<Option<Vec<Mailbox>>, DataflowError> {
let raw = match input.get(field) {
None | Some(Value::Null) => return Ok(None),
Some(raw) => resolve_value(raw, ctx),
};
let parsed = match raw {
Value::String(s) => vec![parse_mailbox(field, &s)?],
Value::Array(items) => {
if items.is_empty() {
return Ok(None);
}
items
.iter()
.enumerate()
.map(|(i, item)| match item {
Value::String(s) => parse_mailbox(&format!("{field}[{i}]"), s),
_ => Err(validation(&format!("'{field}[{i}]' must be a string"))),
})
.collect::<Result<Vec<_>, _>>()?
}
Value::Null => return Ok(None),
_ => {
return Err(validation(&format!(
"'{field}' must resolve to an address or an array of addresses"
)));
}
};
Ok(Some(parsed))
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct Mailbox {
pub name: Option<String>,
pub email: String,
}
impl Mailbox {
fn domain(&self) -> &str {
self.email.rsplit_once('@').map_or("", |(_, domain)| domain)
}
fn to_address(&self) -> Address<'static> {
Address::new_address(self.name.clone(), self.email.clone())
}
}
pub(crate) fn parse_mailbox(field: &str, s: &str) -> Result<Mailbox, DataflowError> {
let s = s.trim();
let invalid = |why: &str| validation(&format!("'{field}' is not a valid email address: {why}"));
if s.is_empty() {
return Err(invalid("empty"));
}
if s.contains(['\r', '\n']) {
return Err(invalid("contains a line break"));
}
let raw = format!("To:{s}\r\n\r\n");
let parsed = mail_parser::MessageParser::new()
.with_address_headers()
.parse_headers(raw.as_bytes())
.ok_or_else(|| invalid("unparseable"))?;
let Some(mail_parser::Address::List(list)) = parsed.to() else {
return Err(invalid("expected a single address"));
};
let [addr] = list.as_slice() else {
return Err(invalid("expected exactly one address"));
};
let email = addr
.address
.as_deref()
.map(str::trim)
.filter(|email| !email.is_empty())
.ok_or_else(|| invalid("no address part"))?;
let Some((local, domain)) = email.split_once('@') else {
return Err(invalid("missing '@'"));
};
if local.is_empty() || domain.is_empty() || domain.contains('@') {
return Err(invalid("malformed address"));
}
if email.contains(|c: char| c.is_whitespace() || c.is_control() || "<>,;:\"\\".contains(c)) {
return Err(invalid("malformed address"));
}
Ok(Mailbox {
name: addr
.name
.as_deref()
.map(str::trim)
.filter(|name| !name.is_empty())
.map(str::to_string),
email: email.to_string(),
})
}
fn sender(
config: &SmtpConnectorConfig,
from_override: Option<&str>,
connector: &str,
) -> Result<Mailbox, DataflowError> {
match from_override {
Some(from) => {
if !config.allow_from_override {
return Err(validation(&format!(
"connector '{connector}' does not allow a per-send 'from' \
(set allow_from_override on the connector)"
)));
}
parse_mailbox("from", from)
}
None => parse_mailbox("connector 'from'", &config.from),
}
}
fn generated_message_id(from: &Mailbox) -> String {
format!("{}@{}", uuid::Uuid::new_v4(), from.domain())
}
fn check_headers_field(input: &Value) -> Result<(), DataflowError> {
match input.get("headers") {
None | Some(Value::Null) => Ok(()),
Some(Value::Object(map)) => {
for (name, value) in map {
if PROTECTED_HEADERS
.iter()
.any(|p| p.eq_ignore_ascii_case(name))
{
return Err(validation(&format!(
"'headers' may not set '{name}' — use the structured field instead"
)));
}
if !value.is_string() {
return Err(validation(&format!("'headers.{name}' must be a string")));
}
if name.is_empty() || !name.bytes().all(|b| b.is_ascii_graphic() && b != b':') {
return Err(validation(&format!(
"'headers.{name}' is not a valid header name"
)));
}
if value.as_str().is_some_and(|v| v.contains(['\r', '\n'])) {
return Err(validation(&format!(
"'headers.{name}' may not contain a line break"
)));
}
}
Ok(())
}
Some(_) => Err(validation("'headers' must be an object of string values")),
}
}
struct MessageParts<'a> {
from: &'a Mailbox,
message_id: &'a str,
to: &'a [Mailbox],
cc: &'a [Mailbox],
subject: &'a str,
reply_to: Option<Mailbox>,
headers: Option<&'a Value>,
text: Option<String>,
html: Option<String>,
}
fn build_message(parts: MessageParts<'_>) -> Result<Vec<u8>, DataflowError> {
let MessageParts {
from,
message_id,
to,
cc,
subject,
reply_to,
headers,
text,
html,
} = parts;
let addresses =
|boxes: &[Mailbox]| Address::List(boxes.iter().map(Mailbox::to_address).collect());
let mut builder = MessageBuilder::new()
.from(from.to_address())
.subject(subject)
.message_id(message_id.to_string());
if !to.is_empty() {
builder = builder.to(addresses(to));
}
if !cc.is_empty() {
builder = builder.cc(addresses(cc));
}
if let Some(mbox) = reply_to {
builder = builder.reply_to(mbox.to_address());
}
builder = match (text, html) {
(Some(text), Some(html)) => builder.text_body(text).html_body(html),
(Some(text), None) => builder.text_body(text),
(None, Some(html)) => builder.html_body(html),
(None, None) => unreachable!("checked in the prologue"),
};
if let Some(Value::Object(map)) = headers {
for (name, value) in map {
builder = builder.header(
name.clone(),
Raw::new(value.as_str().unwrap_or_default().to_string()),
);
}
}
builder.write_to_vec().map_err(|e| {
DataflowError::Validation(format!("{NAME}: message could not be assembled: {e}"))
})
}
async fn deliver(
client: &mut PooledClient,
from: &Mailbox,
recipients: &[&str],
message: &[u8],
) -> Result<String, mail_send::Error> {
client
.cmd(format!("MAIL FROM:<{}>\r\n", from.email).as_bytes())
.await?
.assert_positive_completion()?;
for rcpt in recipients {
client
.cmd(format!("RCPT TO:<{rcpt}>\r\n").as_bytes())
.await?
.assert_positive_completion()?;
}
client.cmd(b"DATA\r\n").await?.assert_code(354)?;
let timeout = client.timeout;
let reply = tokio::time::timeout(timeout, async {
client.write_message(message).await?;
client.read().await
})
.await
.map_err(|_| mail_send::Error::Timeout)??;
if !reply.is_positive_completion() {
return Err(mail_send::Error::UnexpectedReply(reply));
}
Ok(reply.message)
}
pub(super) fn validate_static_input(
obj: &serde_json::Map<String, Value>,
) -> Vec<(&'static str, &'static str, String)> {
let mut errors: Vec<(&'static str, &'static str, String)> = Vec::new();
let input = Value::Object(obj.clone());
if obj.get("text").is_none_or(Value::is_null) && obj.get("html").is_none_or(Value::is_null) {
errors.push((
"",
"REQUIRED",
"send_email requires at least one of 'text' or 'html'".to_string(),
));
}
if let Err(e) = check_headers_field(&input) {
errors.push(("headers", "INVALID", strip_name(&e)));
}
for field in ["to", "cc", "bcc", "from", "reply_to"] {
let Some(raw) = obj.get(field) else { continue };
let addresses: Vec<&str> = match raw {
Value::String(s) => vec![s.as_str()],
Value::Array(items) => items.iter().filter_map(Value::as_str).collect(),
_ => continue,
};
for s in addresses {
if let Err(e) = parse_mailbox(field, s) {
errors.push((field_name(field), "INVALID", strip_name(&e)));
}
}
}
errors
}
fn field_name(key: &str) -> &'static str {
super::schema::static_field_name(SEND_EMAIL_FIELDS, key, "to")
}
fn strip_name(e: &DataflowError) -> String {
super::schema::strip_handler_prefix(NAME, e)
}
pub(super) const SEND_EMAIL_FIELDS: &[FieldSchema] = &[
FieldSchema {
name: "connector",
description: "Name of the SMTP connector to send through.",
kind: FieldKind::String,
required: true,
resolvable: false,
alias: None,
},
FieldSchema {
name: "to",
description: "Recipient address or array of addresses; each is 'addr@example.com' \
or 'Name <addr@example.com>'.",
kind: FieldKind::Any,
required: true,
resolvable: true,
alias: None,
},
FieldSchema {
name: "cc",
description: "Carbon-copy recipients; same forms as 'to'.",
kind: FieldKind::Any,
required: false,
resolvable: true,
alias: None,
},
FieldSchema {
name: "bcc",
description: "Blind-carbon-copy recipients; same forms as 'to'.",
kind: FieldKind::Any,
required: false,
resolvable: true,
alias: None,
},
FieldSchema {
name: "subject",
description: "Message subject (UTF-8).",
kind: FieldKind::String,
required: true,
resolvable: true,
alias: None,
},
FieldSchema {
name: "text",
description: "Plain-text body. At least one of 'text'/'html' is required; both \
together send multipart/alternative.",
kind: FieldKind::String,
required: false,
resolvable: true,
alias: None,
},
FieldSchema {
name: "html",
description: "HTML body. At least one of 'text'/'html' is required.",
kind: FieldKind::String,
required: false,
resolvable: true,
alias: None,
},
FieldSchema {
name: "from",
description: "Per-send sender override; honored only when the connector sets \
allow_from_override. Default: the connector's 'from'.",
kind: FieldKind::String,
required: false,
resolvable: true,
alias: None,
},
FieldSchema {
name: "reply_to",
description: "Reply-To address.",
kind: FieldKind::String,
required: false,
resolvable: true,
alias: None,
},
FieldSchema {
name: "headers",
description: "Extra headers (string values). Structured names (From, To, Subject, \
Content-Type, ...) are rejected; intended for List-Unsubscribe, \
Auto-Submitted, correlation IDs.",
kind: FieldKind::Object,
required: false,
resolvable: false,
alias: None,
},
FieldSchema {
name: "output",
description: "Dotted path where { message_id, response } is stored. Defaults to \
\"data\".",
kind: FieldKind::String,
required: false,
resolvable: false,
alias: None,
},
];
#[cfg(test)]
mod tests {
use super::*;
use crate::connector::{ConnectorConfig, SmtpAuth, SmtpTls};
use serde_json::json;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
async fn spawn_smtp_mock() -> (std::net::SocketAddr, tokio::sync::oneshot::Receiver<String>) {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("test");
let addr = listener.local_addr().expect("test");
let (tx, rx) = tokio::sync::oneshot::channel();
tokio::spawn(async move {
let (stream, _) = listener.accept().await.expect("test");
let (read, mut write) = stream.into_split();
let mut lines = BufReader::new(read).lines();
write.write_all(b"220 mock ESMTP\r\n").await.expect("test");
let mut data = String::new();
let mut in_data = false;
while let Ok(Some(line)) = lines.next_line().await {
if in_data {
if line == "." {
in_data = false;
write
.write_all(b"250 2.0.0 OK queued as mock-42\r\n")
.await
.expect("test");
} else {
data.push_str(&line);
data.push('\n');
}
continue;
}
let upper = line.to_ascii_uppercase();
let reply: &[u8] = if upper.starts_with("EHLO") || upper.starts_with("HELO") {
b"250-mock\r\n250 8BITMIME\r\n"
} else if upper.starts_with("MAIL FROM") || upper.starts_with("RCPT TO") {
b"250 OK\r\n"
} else if upper.starts_with("DATA") {
in_data = true;
b"354 go ahead\r\n"
} else if upper.starts_with("QUIT") {
write.write_all(b"221 bye\r\n").await.expect("test");
break;
} else {
b"250 OK\r\n"
};
write.write_all(reply).await.expect("test");
}
let _ = tx.send(data);
});
(addr, rx)
}
#[derive(Default)]
struct MockLog {
connections: usize,
transcript: String,
}
async fn spawn_logging_smtp_mock() -> (std::net::SocketAddr, Arc<std::sync::Mutex<MockLog>>) {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("test");
let addr = listener.local_addr().expect("test");
let log = Arc::new(std::sync::Mutex::new(MockLog::default()));
let accept_log = log.clone();
tokio::spawn(async move {
while let Ok((stream, _)) = listener.accept().await {
accept_log.lock().expect("test").connections += 1;
let session_log = accept_log.clone();
tokio::spawn(async move {
let (read, mut write) = stream.into_split();
let mut lines = BufReader::new(read).lines();
let _ = write.write_all(b"220 mock ESMTP\r\n").await;
let mut in_data = false;
while let Ok(Some(line)) = lines.next_line().await {
{
let mut log = session_log.lock().expect("test");
log.transcript.push_str(&line);
log.transcript.push('\n');
}
if in_data {
if line == "." {
in_data = false;
let _ =
write.write_all(b"250 2.0.0 OK queued as mock-42\r\n").await;
}
continue;
}
let upper = line.to_ascii_uppercase();
let reply: &[u8] = if upper.starts_with("EHLO") || upper.starts_with("HELO")
{
b"250-mock\r\n250 8BITMIME\r\n"
} else if upper.starts_with("DATA") {
in_data = true;
b"354 go ahead\r\n"
} else if upper.starts_with("QUIT") {
let _ = write.write_all(b"221 bye\r\n").await;
break;
} else {
b"250 OK\r\n"
};
let _ = write.write_all(reply).await;
}
});
}
});
(addr, log)
}
fn smtp_config(addr: std::net::SocketAddr) -> SmtpConnectorConfig {
SmtpConnectorConfig {
host: addr.ip().to_string(),
port: addr.port(),
tls: SmtpTls::None,
auth: SmtpAuth::None,
from: "Orion Test <noreply@example.test>".to_string(),
allow_from_override: false,
allow_private_urls: true, timeout_ms: 5_000,
}
}
async fn run(input: Value, config: SmtpConnectorConfig, data: Value) -> Result<Value, String> {
let registry =
std::sync::Arc::new(crate::connector::ConnectorRegistry::new(Default::default()));
registry
.insert_for_test("mailer", ConnectorConfig::Smtp(config))
.await;
crate::engine::functions::run_test_task(
NAME,
Box::new(SendEmailHandler {
registry,
smtp_pool: std::sync::Arc::new(SmtpPoolCache::new(4)),
}),
input,
data,
)
.await
}
#[tokio::test]
async fn sends_a_multipart_message_through_a_real_smtp_exchange() {
let (addr, rx) = spawn_smtp_mock().await;
let out = run(
json!({
"connector": "mailer",
"to": {"var": "data.email"},
"cc": ["Audit <audit@example.test>"],
"subject": "Your verification code",
"text": "Your OTP is 123456",
"html": "<p>Your OTP is <b>123456</b></p>",
"headers": {"Auto-Submitted": "auto-generated"},
"output": "data.mail"
}),
smtp_config(addr),
json!({"email": "User <user@example.test>"}),
)
.await
.expect("test");
let message_id = out["mail"]["message_id"].as_str().expect("test");
assert!(
message_id.starts_with('<') && message_id.ends_with("@example.test>"),
"{message_id}"
);
assert!(
out["mail"]["response"]
.as_str()
.expect("test")
.contains("mock-42"),
"{}",
out["mail"]
);
let wire = rx.await.expect("test");
assert!(wire.contains("multipart/alternative"), "{wire}");
assert!(wire.contains("Auto-Submitted: auto-generated"), "{wire}");
assert!(wire.contains("Your OTP is 123456"), "{wire}");
assert!(wire.contains("Subject: Your verification code"), "{wire}");
assert!(wire.contains("Cc: "), "{wire}");
assert!(wire.contains(message_id), "{wire}");
}
#[tokio::test]
async fn from_override_needs_the_connector_gate() {
let (addr, _rx) = spawn_smtp_mock().await;
let input = json!({
"connector": "mailer",
"to": "user@example.test",
"subject": "s",
"text": "b",
"from": "spoof@example.test"
});
let err = run(input.clone(), smtp_config(addr), json!({}))
.await
.expect_err("test");
assert!(err.contains("allow_from_override"), "{err}");
let (addr, _rx) = spawn_smtp_mock().await;
let mut config = smtp_config(addr);
config.allow_from_override = true;
run(input, config, json!({})).await.expect("test");
}
#[tokio::test]
async fn message_shape_errors_are_named() {
let (addr, _rx) = spawn_smtp_mock().await;
let config = smtp_config(addr);
let err = run(
json!({"connector": "mailer", "to": "a@b.test", "subject": "s"}),
config.clone(),
json!({}),
)
.await
.expect_err("test");
assert!(err.contains("'text' or 'html'"), "{err}");
let err = run(
json!({"connector": "mailer", "to": ["ok@b.test", "not an address"],
"subject": "s", "text": "b"}),
config.clone(),
json!({}),
)
.await
.expect_err("test");
assert!(err.contains("to[1]"), "{err}");
let err = run(
json!({"connector": "mailer", "to": "a@b.test", "subject": "s",
"text": "b", "headers": {"Subject": "override"}}),
config,
json!({}),
)
.await
.expect_err("test");
assert!(err.contains("structured field"), "{err}");
}
#[test]
fn static_validation_reads_the_same_rules() {
let obj = json!({"connector": "m", "to": "user@example.test", "subject": "s"});
let errs = validate_static_input(obj.as_object().expect("test"));
assert!(
errs.iter()
.any(|(_, c, m)| *c == "REQUIRED" && m.contains("'text' or 'html'")),
"{errs:?}"
);
let obj = json!({"connector": "m", "to": "not an address",
"subject": "s", "text": "b"});
let errs = validate_static_input(obj.as_object().expect("test"));
assert!(
errs.iter().any(|(f, c, _)| *f == "to" && *c == "INVALID"),
"{errs:?}"
);
let obj = json!({"connector": "m", "to": {"var": "data.email"},
"subject": "s", "text": {"var": "data.body"}});
let errs = validate_static_input(obj.as_object().expect("test"));
assert!(errs.is_empty(), "{errs:?}");
let obj = json!({"connector": "m", "to": "a@b.test", "subject": "s",
"text": "b", "headers": {"Message-ID": "<x@y>"}});
let errs = validate_static_input(obj.as_object().expect("test"));
assert!(
errs.iter()
.any(|(f, c, _)| *f == "headers" && *c == "INVALID"),
"{errs:?}"
);
}
async fn send_n(addr: std::net::SocketAddr, input: Value, count: usize) -> Arc<SmtpPoolCache> {
let registry =
std::sync::Arc::new(crate::connector::ConnectorRegistry::new(Default::default()));
registry
.insert_for_test("mailer", ConnectorConfig::Smtp(smtp_config(addr)))
.await;
let smtp_pool = std::sync::Arc::new(SmtpPoolCache::new(4));
for _ in 0..count {
crate::engine::functions::run_test_task(
NAME,
Box::new(SendEmailHandler {
registry: registry.clone(),
smtp_pool: smtp_pool.clone(),
}),
input.clone(),
json!({}),
)
.await
.expect("test");
}
smtp_pool
}
#[tokio::test]
async fn bcc_reaches_the_envelope_but_never_the_headers() {
let (addr, log) = spawn_logging_smtp_mock().await;
let pool = send_n(
addr,
json!({"connector": "mailer", "to": "user@example.test",
"bcc": ["blind@example.test"], "subject": "s", "text": "b"}),
1,
)
.await;
drop(pool);
let transcript = log.lock().expect("test").transcript.clone();
assert!(
transcript.contains("RCPT TO:<blind@example.test>"),
"bcc must be an envelope recipient: {transcript}"
);
assert!(
!transcript.to_ascii_lowercase().contains("bcc:"),
"bcc must never be a header: {transcript}"
);
}
#[tokio::test]
async fn a_second_send_reuses_the_pooled_connection() {
let (addr, log) = spawn_logging_smtp_mock().await;
let pool = send_n(
addr,
json!({"connector": "mailer", "to": "user@example.test",
"subject": "s", "text": "b"}),
2,
)
.await;
drop(pool);
let log = log.lock().expect("test");
assert_eq!(
log.connections, 1,
"two sends should share one connection: {}",
log.transcript
);
assert!(
log.transcript.contains("RSET"),
"reuse must probe with RSET: {}",
log.transcript
);
}
#[test]
fn header_injection_attempts_are_refused() {
for (name, value) in [
("X-Evil", "ok\r\nBcc: attacker@example.test"),
("X-Evil", "ok\nSubject: replaced"),
("X:Evil", "ok"),
("X Evil", "ok"),
] {
let obj = json!({"connector": "m", "to": "a@b.test", "subject": "s",
"text": "b", "headers": {name: value}});
let errs = validate_static_input(obj.as_object().expect("test"));
assert!(
errs.iter()
.any(|(f, c, _)| *f == "headers" && *c == "INVALID"),
"{name}: {value:?} should be refused, got {errs:?}"
);
}
}
#[test]
fn addresses_with_line_breaks_are_refused() {
let err = parse_mailbox("to", "a@b.test\r\nBcc: attacker@example.test")
.expect_err("test")
.to_string();
assert!(err.contains("line break"), "{err}");
}
#[test]
fn mailbox_forms_parse_and_name_the_field_on_failure() {
assert!(parse_mailbox("to", "user@example.test").is_ok());
assert!(parse_mailbox("to", "Ada Lovelace <ada@example.test>").is_ok());
let err = parse_mailbox("reply_to", "nope")
.expect_err("test")
.to_string();
assert!(err.contains("reply_to"), "{err}");
}
}