ocre/runtime/mail.rs
1use std::time::Duration;
2
3use wasm_bindgen::{JsCast, JsValue};
4use worker::{
5 EmailAddress, EmailAttachment, Env, Fetch, ForwardableEmailMessage, Headers, Method, Request, RequestInit,
6 SendEmailBuilder, send::SendFuture,
7};
8
9use super::Ctx;
10use crate::{
11 Error, Result,
12 jobs::{DEFAULT_QUEUE, Payload},
13 mail::{
14 APP_URL, Adapter, Attachment, EMAIL_BINDING, Email, LOG_PREFIX, MAIL_ADAPTER, MAIL_FROM, Message, Outgoing,
15 RESEND_API_KEY, RESEND_URL, absolute_url, adapter_for, cloudflare_error, resend_error, resend_key,
16 },
17};
18
19/// Sends `email` now, from `MAIL_FROM`, with the adapter named by `MAIL_ADAPTER`.
20///
21/// `log` prints the whole email to the Worker console between
22/// [`LOG_PREFIX`](crate::mail::LOG_PREFIX) lines and sends nothing; `resend`
23/// makes one `POST https://api.resend.com/emails` subrequest; `cloudflare`
24/// calls the [`EMAIL_BINDING`](crate::mail::EMAIL_BINDING) send_email
25/// binding. The request waits for the provider; use
26/// [`deliver_later`](crate::mail::deliver_later) to answer first and retry failures.
27///
28/// The returned future is `Send`, so axum handlers can await it.
29///
30/// Free-plan limits (September 2026): Resend sends 100 emails a day and 3,000
31/// a month from one domain, to any recipient; Cloudflare Email Service on
32/// Workers Free only delivers to verified destination addresses of the account.
33///
34/// # Errors
35///
36/// - [`Error::BadRequest`](crate::Error::BadRequest) (400): a recipient (`to`, `cc`, `bcc`) or `reply_to` is not an email address.
37/// - [`Error::Internal`](crate::Error::Internal) (500), with a log message naming the fix:
38/// `MAIL_ADAPTER` unset or unknown; `MAIL_FROM` unset or not an address; a
39/// subject that is empty or spans several lines; the `RESEND_API_KEY`
40/// secret or the `EMAIL: bindings.sendEmail()` binding missing; the
41/// provider refused the email (e.g. Resend's quota reached, an unverified
42/// domain or recipient) or could not be reached.
43///
44/// # Examples
45///
46/// ```no_run
47/// use axum::extract::State;
48/// use ocre::mail::Email;
49/// use ocre::{Ctx, Result};
50///
51/// async fn invite(State(ctx): State<Ctx>) -> Result<&'static str> {
52/// let email = Email::new("ada@example.com", "You're invited", "Join us: https://example.com/join");
53/// ocre::mail::send(&ctx, email).await?;
54/// Ok("invited")
55/// }
56/// ```
57pub fn send(ctx: &Ctx, email: Email) -> impl Future<Output = Result<()>> + Send + use<> {
58 let env = ctx.env().clone();
59 SendFuture::new(async move { deliver(&env, email).await })
60}
61
62/// Checks `email` now and sends it later from the background jobs queue, like Rails' `deliver_later`.
63///
64/// The address, `MAIL_FROM` and `MAIL_ADAPTER` are checked right away, like
65/// [`send`](crate::mail::send), so a bad address is still a 400 for the
66/// request. The email is then put on the `JOBS` queue and the handler answers
67/// without waiting for the provider; [`consume`](crate::jobs::consume) sends
68/// it with [`send`](crate::mail::send), logs `[ocre jobs] mail done`, and on
69/// failure retries it with the jobs backoff (30 s, 1 min, 3 min, 9 min, 27
70/// min) before the dead-letter queue. The Resend key and the send_email
71/// binding are only looked up when the consumer sends.
72///
73/// Needs the `JOBS` queue: run `ocre g job <Name>` once to wire it.
74///
75/// Free-plan cost: one queue message, i.e. 3 of the 10,000 daily Queues
76/// operations (write, read, delete), one more read per retry and one more
77/// write if it is dead-lettered; plus the provider's limits of
78/// [`send`](crate::mail::send).
79///
80/// # Errors
81///
82/// - [`Error::BadRequest`](crate::Error::BadRequest) (400): a recipient or `reply_to` is not an email address.
83/// - [`Error::Internal`](crate::Error::Internal) (500): `MAIL_ADAPTER` unset or
84/// unknown, `MAIL_FROM` unset or not an address, a subject that is empty or
85/// spans several lines, an email over the 128 KB queue message limit, the
86/// `JOBS` queue binding missing from cloudflare.config.ts, or
87/// Queues refusing the message.
88///
89/// # Examples
90///
91/// ```no_run
92/// use axum::extract::State;
93/// use ocre::mail::Email;
94/// use ocre::{Ctx, Result};
95///
96/// async fn sign_up(State(ctx): State<Ctx>) -> Result<&'static str> {
97/// let email = Email::new("ada@example.com", "Welcome", "Hello Ada");
98/// ocre::mail::deliver_later(&ctx, email).await?;
99/// Ok("check your inbox")
100/// }
101/// ```
102pub fn deliver_later(ctx: &Ctx, email: Email) -> impl Future<Output = Result<()>> + Send + use<> {
103 deliver_in(ctx, email, Duration::ZERO)
104}
105
106/// Checks `email` now and sends it from the jobs queue after `delay`, like Rails' `deliver_later(wait:)`.
107///
108/// Works like [`deliver_later`](crate::mail::deliver_later), with the
109/// message due after `delay` (whole seconds, 24 hours at most, Queues'
110/// limit): a reminder an hour after sign-up, a digest at the end of the
111/// day. Same free-plan cost as [`deliver_later`](crate::mail::deliver_later).
112///
113/// # Errors
114///
115/// Those of [`deliver_later`](crate::mail::deliver_later), plus an
116/// [`Error::Internal`](crate::Error::Internal) when `delay` is over 24 hours.
117///
118/// # Examples
119///
120/// ```no_run
121/// use std::time::Duration;
122///
123/// use ocre::mail::Email;
124///
125/// async fn remind(ctx: &ocre::Ctx) -> ocre::Result<()> {
126/// let email = Email::new("ada@example.com", "Finish setting up your shop", "https://example.com/setup");
127/// ocre::mail::deliver_in(ctx, email, Duration::from_secs(3600)).await
128/// }
129/// ```
130pub fn deliver_in(ctx: &Ctx, email: Email, delay: Duration) -> impl Future<Output = Result<()>> + Send + use<> {
131 let env = ctx.env();
132 let checked = adapter_for(var(env, MAIL_ADAPTER).as_deref(), &email)
133 .and_then(|_| Outgoing::new(var(env, MAIL_FROM), email))
134 .map(|outgoing| Payload::Mail(Box::new(outgoing.email)));
135 super::jobs::send(env.clone(), DEFAULT_QUEUE, checked, delay)
136}
137
138fn var(env: &Env, name: &str) -> Option<String> {
139 env.var(name).ok().map(|value| value.to_string())
140}
141
142/// An absolute URL for `path` on the app's public address, the [`APP_URL`](crate::mail::APP_URL) variable (Rails' `_url` helpers in mailers).
143///
144/// For links and images in emails, which mail clients open outside the app:
145/// `url(ctx, "/posts/1")` is `https://shop.example.com/posts/1`. A `path`
146/// that is already absolute is returned as is. It reads a variable, no
147/// binding call, so mailers and jobs can use it outside any request.
148///
149/// # Errors
150///
151/// [`Error::Internal`](crate::Error::Internal) (500) when `APP_URL` is unset
152/// or not an `http(s)` address, naming the fix.
153///
154/// # Examples
155///
156/// ```no_run
157/// use ocre::{Ctx, Result, mail::{self, Email}};
158///
159/// fn reset_email(ctx: &Ctx, to: &str, token: &str) -> Result<Email> {
160/// let link = mail::url(ctx, &format!("/password/reset/{token}"))?;
161/// let logo = mail::url(ctx, "/images/logo.png")?;
162/// Ok(Email::new(to, "Reset your password", format!("Open {link}"))
163/// .html(format!("<img src=\"{logo}\" alt=\"\"><p><a href=\"{link}\">Reset your password</a></p>")))
164/// }
165/// # let _ = reset_email;
166/// ```
167pub fn url(ctx: &Ctx, path: &str) -> Result<String> {
168 absolute_url(var(ctx.env(), APP_URL), path)
169}
170
171pub(crate) async fn deliver(env: &Env, email: Email) -> Result<()> {
172 let adapter = adapter_for(var(env, MAIL_ADAPTER).as_deref(), &email)?;
173 let outgoing = Outgoing::new(var(env, MAIL_FROM), email)?;
174 match adapter {
175 Adapter::Log => {
176 worker::console_log!("{}", outgoing.log_text());
177 crate::mail::capture(outgoing);
178 Ok(())
179 }
180 Adapter::Resend => resend(env, &outgoing).await,
181 Adapter::Cloudflare => cloudflare(env, &outgoing).await,
182 }
183}
184
185async fn resend(env: &Env, outgoing: &Outgoing) -> Result<()> {
186 let key = resend_key(super::secrets::lookup(env, RESEND_API_KEY).await?)?;
187 let headers = Headers::new();
188 headers.set("Authorization", &format!("Bearer {key}"))?;
189 headers.set("Content-Type", "application/json")?;
190 let mut init = RequestInit::new();
191 init.with_method(Method::Post)
192 .with_headers(headers)
193 .with_body(Some(JsValue::from_str(&outgoing.resend_json().to_string())));
194 let mut response = Fetch::Request(Request::new_with_init(RESEND_URL, &init)?).send().await?;
195 let status = response.status_code();
196 if (200..300).contains(&status) {
197 Ok(())
198 } else {
199 Err(resend_error(status, &response.text().await.unwrap_or_default()))
200 }
201}
202
203async fn cloudflare(env: &Env, outgoing: &Outgoing) -> Result<()> {
204 let binding = env.send_email(EMAIL_BINDING).map_err(|err| {
205 Error::internal(format!(
206 "cannot send email: the send_email binding `{EMAIL_BINDING}` is missing ({err}). Fix: add \
207 `{EMAIL_BINDING}: bindings.sendEmail(),` to worker.env in cloudflare.config.ts"
208 ))
209 })?;
210 let Outgoing { from, email } = outgoing;
211 let to = Outgoing::addresses(&email.to);
212 let builder = match &from.name {
213 Some(name) => SendEmailBuilder::builder_with_email_address_and_slice(
214 &EmailAddress::new(name, &from.address),
215 &to,
216 &email.subject,
217 ),
218 None => SendEmailBuilder::builder_with_str_and_slice(&from.address, &to, &email.subject),
219 };
220 let mut builder = builder.text(&email.text);
221 if let Some(html) = &email.html {
222 builder = builder.html(html);
223 }
224 if let Some(reply_to) = &email.reply_to {
225 builder = builder.reply_to(reply_to);
226 }
227 if !email.cc.is_empty() {
228 builder = builder.cc_with_slice(&Outgoing::addresses(&email.cc));
229 }
230 if !email.bcc.is_empty() {
231 builder = builder.bcc_with_slice(&Outgoing::addresses(&email.bcc));
232 }
233 if !email.headers.is_empty() {
234 let headers = worker::js_sys::Object::new();
235 for (name, value) in &email.headers {
236 // Setting a property on a plain object cannot fail.
237 let _ = worker::js_sys::Reflect::set(&headers, &JsValue::from_str(name), &JsValue::from_str(value));
238 }
239 builder = builder.headers(headers.unchecked_ref());
240 }
241 if !email.attachments.is_empty() {
242 let files: Vec<EmailAttachment> = email
243 .attachments
244 .iter()
245 .map(|file| {
246 let bytes = worker::js_sys::Uint8Array::from(file.content.as_slice());
247 match &file.content_id {
248 Some(id) => {
249 EmailAttachment::new_inline_with_typed_array(id, &file.filename, &file.content_type, &bytes)
250 }
251 None => {
252 EmailAttachment::new_attachment_with_typed_array(&file.filename, &file.content_type, &bytes)
253 }
254 }
255 })
256 .collect();
257 builder = builder.attachments(&files);
258 }
259 binding.send_with_builder(&builder.build()).await.map_err(|err| cloudflare_error(&js_error(&err)))?;
260 Ok(())
261}
262
263/// `message` of a JS error, plus its `code` (`E_SENDER_NOT_VERIFIED`, ...) when set.
264fn js_error(err: &worker::js_sys::Error) -> String {
265 let message = String::from(err.message());
266 match worker::js_sys::Reflect::get(err, &JsValue::from_str("code")).ok().and_then(|code| code.as_string()) {
267 Some(code) => format!("{code}: {message}"),
268 None => message,
269 }
270}
271
272/// An email that Cloudflare Email Routing delivered to the Worker, handed to the app's mailbox.
273///
274/// [`receive`](crate::mail::receive) builds it. The envelope ([`from`](Self::from), [`to`](Self::to)) comes from
275/// Cloudflare; headers and bodies are parsed from the raw message
276/// (multipart, quoted-printable, base64, RFC 2047 encoded words; files in
277/// [`attachments`](Self::attachments), the bytes in [`raw`](Self::raw)). Bounce it with
278/// [`reject`](Self::reject) or pass it on with [`forward`](Self::forward).
279/// Receiving is free and unlimited on every plan.
280///
281/// # Examples
282///
283/// ```no_run
284/// use ocre::mail::InboundEmail;
285/// use ocre::{Ctx, Result};
286///
287/// // src/mailbox.rs
288/// pub async fn receive(_ctx: Ctx, email: InboundEmail) -> Result<()> {
289/// match email.to() {
290/// "support@example.com" => email.forward("team@example.com").await?,
291/// _ => email.reject("Unknown address"),
292/// }
293/// Ok(())
294/// }
295/// ```
296pub struct InboundEmail {
297 inner: ForwardableEmailMessage,
298 from: String,
299 to: String,
300 raw: Vec<u8>,
301 message: Message,
302}
303
304impl InboundEmail {
305 /// Returns the envelope sender (SMTP `MAIL FROM`), checked by Cloudflare.
306 ///
307 /// The `From` header may differ: read it with `email.header("From")`.
308 ///
309 /// # Examples
310 ///
311 /// ```no_run
312 /// # async fn mailbox(email: ocre::mail::InboundEmail) {
313 /// if email.from().ends_with("@example.com") {
314 /// // a colleague
315 /// }
316 /// # }
317 /// ```
318 pub fn from(&self) -> &str {
319 &self.from
320 }
321
322 /// Returns the envelope recipient: the address of this app that received the email.
323 ///
324 /// Match on it to route several addresses to one Worker.
325 ///
326 /// # Examples
327 ///
328 /// ```no_run
329 /// # async fn mailbox(email: ocre::mail::InboundEmail) {
330 /// let team = email.to().starts_with("support@");
331 /// # let _ = team;
332 /// # }
333 /// ```
334 pub fn to(&self) -> &str {
335 &self.to
336 }
337
338 /// Returns the decoded `Subject` header, or `""` when it is missing.
339 ///
340 /// # Examples
341 ///
342 /// ```no_run
343 /// # async fn mailbox(email: ocre::mail::InboundEmail) {
344 /// let urgent = email.subject().contains("URGENT");
345 /// # let _ = urgent;
346 /// # }
347 /// ```
348 pub fn subject(&self) -> &str {
349 self.message.header("Subject").unwrap_or("")
350 }
351
352 /// Returns the first header with this name (case-insensitive), decoded.
353 ///
354 /// Folded lines are joined and RFC 2047 encoded words decoded. `None`
355 /// when the message has no such header.
356 ///
357 /// # Examples
358 ///
359 /// ```no_run
360 /// # async fn mailbox(email: ocre::mail::InboundEmail) {
361 /// let id = email.header("message-id").unwrap_or("");
362 /// # let _ = id;
363 /// # }
364 /// ```
365 pub fn header(&self, name: &str) -> Option<&str> {
366 self.message.header(name)
367 }
368
369 /// Returns every header, in message order, as decoded `(name, value)` pairs.
370 ///
371 /// Repeated headers (`Received`, ...) appear once per occurrence.
372 ///
373 /// # Examples
374 ///
375 /// ```no_run
376 /// # async fn mailbox(email: ocre::mail::InboundEmail) {
377 /// let hops = email.headers().iter().filter(|(name, _)| name.eq_ignore_ascii_case("Received")).count();
378 /// # let _ = hops;
379 /// # }
380 /// ```
381 pub fn headers(&self) -> &[(String, String)] {
382 &self.message.headers
383 }
384
385 /// Returns the first `text/plain` part, decoded to UTF-8, or `None` when there is none.
386 ///
387 /// # Examples
388 ///
389 /// ```no_run
390 /// # async fn mailbox(email: ocre::mail::InboundEmail) {
391 /// let body = email.text().unwrap_or_default();
392 /// # let _ = body;
393 /// # }
394 /// ```
395 pub fn text(&self) -> Option<&str> {
396 self.message.text.as_deref()
397 }
398
399 /// Returns the files of the email, decoded: every part that is not the first text or HTML body.
400 ///
401 /// Inline images carry their `content_id`. Store them in R2
402 /// (`ocre::storage`) rather than D1; they are the sender's files, so
403 /// check the type and size before keeping them.
404 ///
405 /// # Examples
406 ///
407 /// ```no_run
408 /// # async fn mailbox(email: ocre::mail::InboundEmail) {
409 /// for file in email.attachments() {
410 /// worker::console_log!("{} ({}, {} bytes)", file.filename, file.content_type, file.content.len());
411 /// }
412 /// # }
413 /// ```
414 pub fn attachments(&self) -> &[Attachment] {
415 &self.message.attachments
416 }
417
418 /// Returns the first `text/html` part, decoded to UTF-8, or `None` when there is none.
419 ///
420 /// It is the sender's HTML: never render it unescaped.
421 ///
422 /// # Examples
423 ///
424 /// ```no_run
425 /// # async fn mailbox(email: ocre::mail::InboundEmail) {
426 /// let has_html = email.html().is_some();
427 /// # let _ = has_html;
428 /// # }
429 /// ```
430 pub fn html(&self) -> Option<&str> {
431 self.message.html.as_deref()
432 }
433
434 /// Returns the whole message as received (RFC 5322 bytes), e.g. to store it in R2 or read attachments.
435 ///
436 /// # Examples
437 ///
438 /// ```no_run
439 /// # async fn mailbox(email: ocre::mail::InboundEmail) {
440 /// let size = email.raw().len();
441 /// # let _ = size;
442 /// # }
443 /// ```
444 pub fn raw(&self) -> &[u8] {
445 &self.raw
446 }
447
448 /// Bounces the email: the sending server gets a permanent SMTP error with `reason`.
449 ///
450 /// The mailbox handler still returns normally; [`receive`](crate::mail::receive)
451 /// also bounces the email when the handler returns an `Err`.
452 ///
453 /// # Examples
454 ///
455 /// ```no_run
456 /// # async fn mailbox(email: ocre::mail::InboundEmail) {
457 /// if email.to() != "support@example.com" {
458 /// email.reject("Unknown address");
459 /// }
460 /// # }
461 /// ```
462 pub fn reject(&self, reason: &str) {
463 self.inner.set_reject(reason);
464 }
465
466 /// Forwards the email unchanged to `to`, which must be a verified destination address of the Cloudflare account.
467 ///
468 /// # Errors
469 ///
470 /// [`Error::Internal`](crate::Error::Internal) when Cloudflare refuses, typically because
471 /// `to` is not a verified destination address in Email Routing (the
472 /// message says so).
473 ///
474 /// # Examples
475 ///
476 /// ```no_run
477 /// # async fn mailbox(email: ocre::mail::InboundEmail) -> ocre::Result<()> {
478 /// email.forward("team@example.com").await?;
479 /// # Ok(())
480 /// # }
481 /// ```
482 pub async fn forward(&self, to: &str) -> Result<()> {
483 self.inner.forward(to).await.map_err(|err| {
484 Error::internal(format!(
485 "forwarding the email to {to} failed ({}). Fix: add {to} as a verified destination address \
486 in Cloudflare Email Routing",
487 js_error(&err)
488 ))
489 })?;
490 Ok(())
491 }
492}
493
494/// Runs the app's mailbox `handler` for an email from Cloudflare Email Routing; the Worker's `email` entry point.
495///
496/// Reads the raw message, parses it into an [`InboundEmail`], logs one
497/// `[ocre mail] received from <from> to <to>: <subject>` line and calls
498/// `handler(ctx, email)` with a fresh [`Ctx`](crate::Ctx). When the handler
499/// returns an `Err`, the error is logged (`[ocre mail] the mailbox failed:
500/// ...`) and the email bounced ("The message could not be processed"), so the
501/// sender knows it was not handled. Routing rules in the dashboard (Email
502/// Routing > Routing rules) send an address to the Worker; receiving is free
503/// and unlimited on every plan.
504///
505/// # Errors
506///
507/// A `worker::Error` only when the raw message cannot be read; handler
508/// errors are logged and bounce the email instead.
509///
510/// # Examples
511///
512/// `ocre g mailbox` writes the entry point in `src/lib.rs` and the handler in `src/mailbox.rs`:
513///
514/// ```no_run
515/// mod mailbox {
516/// use ocre::{Ctx, Result, mail::InboundEmail};
517///
518/// pub async fn receive(_ctx: Ctx, email: InboundEmail) -> Result<()> {
519/// email.forward("team@example.com").await
520/// }
521/// }
522///
523/// #[worker::event(email)]
524/// async fn email(
525/// message: worker::ForwardableEmailMessage,
526/// env: worker::Env,
527/// _ctx: worker::Context,
528/// ) -> worker::Result<()> {
529/// ocre::mail::receive(message, env, mailbox::receive).await
530/// }
531/// # fn main() {}
532/// ```
533pub async fn receive<F, Fut>(message: ForwardableEmailMessage, env: Env, handler: F) -> worker::Result<()>
534where
535 F: FnOnce(Ctx, InboundEmail) -> Fut,
536 Fut: Future<Output = Result<()>>,
537{
538 let raw = message.raw_bytes().await?;
539 let parsed = Message::parse(&raw);
540 let email = InboundEmail { from: message.from(), to: message.to(), raw, message: parsed, inner: message.clone() };
541 worker::console_log!("{LOG_PREFIX} received from {} to {}: {}", email.from(), email.to(), email.subject());
542 let ctx = Ctx::new(env);
543 if let Err(err) = handler(ctx.clone(), email).await {
544 worker::console_error!("{LOG_PREFIX} the mailbox failed: {err}");
545 super::jobs::report(&ctx, "ocre.mailbox", &err, []);
546 message.set_reject("The message could not be processed");
547 }
548 super::errors::flush(&ctx).await;
549 Ok(())
550}