Skip to main content

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}