Skip to main content

acme_proxy/notify/
mod.rs

1//! Operator notifications for ACME lifecycle events.
2//!
3//! The shape mirrors [`filter`](crate::filter)/[`signer`](crate::signer): a
4//! trait, an error type, and a [`from_config`] selector building the
5//! configured set at startup. Three backends exist today —
6//! [`email`], [`webhook`] (any HTTP endpoint, with the URL, method, headers
7//! and body all configured, which is what makes a chat provider configuration
8//! rather than code), and [`custom`] (an external script, for a channel that
9//! is not an HTTP request at all).
10//!
11//! ## Fire-and-forget, always — and now durable
12//!
13//! A notification can never affect the ACME response that triggered it.
14//! [`NotifyDispatcher::dispatch`] returns `()` and cannot be `?`'d: it writes
15//! the delivery to the [durable queue](crate::jobs) and returns. That is a
16//! deliberate API shape, not a convention callers must remember — there is no
17//! method on this type whose result a handler *could* propagate.
18//!
19//! What changed is what happens *after* it returns. Delivery used to be a bare
20//! `tokio::spawn`, so a refused SMTP connection or a 503 from a webhook was
21//! logged once and the notification was gone, and a restart lost everything in
22//! flight — the operator never heard that a certificate had been issued, and
23//! nothing recorded that they hadn't. Now each delivery is a `notify_deliver`
24//! job row: a transport failure is retried under `jobs.max_attempts` and the
25//! shared backoff, and a row outlives the process that queued it.
26//!
27//! Two consequences worth not rediscovering:
28//!
29//! - **One job per (occurrence × backend)**, never one per event. A retry must
30//!   not re-send to a backend that already succeeded, or one flaky webhook
31//!   produces a duplicate email on every attempt.
32//! - **[`NotifyError`] carries whether it is worth retrying.** A template that
33//!   does not render and a 400 from a webhook will fail identically for ever,
34//!   so they are permanent and refused on the first attempt; a connection
35//!   refused, a timeout and a 503 have decided nothing, so they are retried.
36//!   That is [`crate::jobs`]'s `Retry`/`Failed` split, applied at the source.
37//!
38//! ## Per-backend event filtering
39//!
40//! Unlike [`FilterConfig::rules`](crate::config::FilterConfig), where every
41//! filter must agree, notify backends are independent broadcast side-channels
42//! — an operator plausibly wants email only for issuance/revocation and a chat
43//! webhook for everything including failures. Each backend's own `events`
44//! list (defaulting to all six kinds) decides what reaches it, and the list
45//! lives on its [`BackendSlot`] so the check happens once, in
46//! [`NotifyDispatcher::dispatch`], rather than in each backend's own delivery
47//! code. Filtering *there* rather than in a wrapper around `send` is what stops
48//! a job being queued for a delivery that would immediately do nothing.
49//!
50//! ## Per-profile, like every other subsystem
51//!
52//! [`build_registry`] builds one [`NotifyDispatcher`] per resolved profile
53//! (not deduplicated by configuration identity like signer backends are —
54//! dispatchers are stateless side-channels, so two profiles with identical
55//! `[notify]` sections simply get two independent instances). The
56//! asynchronous `relay` signer backend, whose completion happens outside
57//! any HTTP handler, is handed the whole `profile name -> dispatcher` map so
58//! it can notify the right profile once an order settles — see
59//! `signer::relay::flow::settle`.
60
61use std::collections::{HashMap, HashSet};
62use std::sync::{Arc, LazyLock};
63
64use async_trait::async_trait;
65use tracing::info;
66
67use crate::config::{ALL_NOTIFY_EVENTS, NotifyConfig, ProfileConfig};
68use crate::jobs::{JobQueue, JobSpec};
69
70pub mod custom;
71pub mod email;
72pub mod job;
73pub mod webhook;
74
75pub use job::{NOTIFY_JOB_KIND, NotifyJob};
76
77/// A pluggable notification channel.
78#[async_trait]
79pub trait NotifyBackend: Send + Sync {
80    /// The configuration name this backend runs under, used in logs.
81    fn name(&self) -> &'static str;
82
83    /// Delivers `event`. Failure is always non-fatal to the *caller* — see the
84    /// module docs — but it is no longer thrown away: [`NotifyJob`] retries it
85    /// unless the [`NotifyError`] says the attempt could never have worked.
86    async fn send(&self, event: &NotifyEvent) -> Result<(), NotifyError>;
87}
88
89/// Why a notify backend failed to deliver, and whether asking again could help.
90///
91/// The `retryable` half is what [`NotifyJob`] turns into
92/// [`JobOutcome::Retry`](crate::jobs::JobOutcome::Retry) or
93/// [`JobOutcome::Failed`](crate::jobs::JobOutcome::Failed), so the distinction
94/// has to be drawn where the failure happens rather than guessed from a string
95/// afterwards. The default — [`NotifyError::new`] — is *retryable*, because
96/// most of these are transport; a backend that knows better says so with
97/// [`NotifyError::permanent`].
98#[derive(Debug, thiserror::Error)]
99#[error("{detail}")]
100pub struct NotifyError {
101    detail: String,
102    retryable: bool,
103}
104
105impl NotifyError {
106    /// A failure that may not recur: a refused connection, a timeout, a 503.
107    pub fn new(detail: impl Into<String>) -> Self {
108        Self {
109            detail: detail.into(),
110            retryable: true,
111        }
112    }
113
114    /// A failure that will repeat identically however many times it is tried:
115    /// a template that does not render, a URL that does not parse, a webhook
116    /// that answers 400.
117    pub fn permanent(detail: impl Into<String>) -> Self {
118        Self {
119            detail: detail.into(),
120            retryable: false,
121        }
122    }
123
124    /// Whether another attempt could plausibly succeed.
125    #[must_use]
126    pub fn retryable(&self) -> bool {
127        self.retryable
128    }
129}
130
131/// Context common to every event: the profile it happened on and the client
132/// address, when a request was in scope.
133///
134/// `client_ip` is `None` on the one firing site with no request in scope at
135/// all — the `relay` signer backend's asynchronous completion, which
136/// runs in a background task long after any handler returned. Templates must
137/// treat it as optional.
138#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
139pub struct ProfileMountedData {
140    pub profile: String,
141}
142
143/// Payload of [`NotifyEvent::AccountCreated`]: a client registered a new
144/// account at this endpoint.
145///
146/// `contact` is what the client supplied, which may legitimately be empty —
147/// RFC 8555 does not require one.
148#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
149pub struct AccountCreatedData {
150    pub profile: String,
151    pub account_id: String,
152    pub contact: Vec<String>,
153    pub client_ip: Option<String>,
154}
155
156/// Payload of [`NotifyEvent::AccountDeactivated`]: an account was deactivated,
157/// either by the client (§7.3.6) or by `acme-proxy account deactivate`.
158///
159/// Deactivation is permanent, so this event has no counterpart.
160#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
161pub struct AccountDeactivatedData {
162    pub profile: String,
163    pub account_id: String,
164    pub client_ip: Option<String>,
165}
166
167/// Payload of [`NotifyEvent::CertificateIssued`]: a certificate was signed.
168///
169/// `cert_serial` is the hex serial, the same value `POST /revokeCert` and the
170/// audit trail identify a certificate by. `identifiers` are the names the
171/// certificate covers.
172#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
173pub struct CertificateIssuedData {
174    pub profile: String,
175    pub order_id: String,
176    pub account_id: String,
177    pub cert_serial: String,
178    pub identifiers: Vec<String>,
179    pub client_ip: Option<String>,
180}
181
182/// Payload of [`NotifyEvent::CertificateRevoked`]: a certificate was withdrawn.
183///
184/// `reason` is the RFC 5280 §5.3.1 code the caller supplied, and is `None` when
185/// none was given — which is not the same as `Some(0)`. It reaches a `custom`
186/// script only in the JSON on stdin, since it has no environment variable.
187#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
188pub struct CertificateRevokedData {
189    pub profile: String,
190    pub order_id: String,
191    pub account_id: String,
192    pub cert_serial: String,
193    pub reason: Option<u32>,
194    pub client_ip: Option<String>,
195}
196
197/// Payload of [`NotifyEvent::ChallengeFailed`]: a validation attempt did not
198/// succeed.
199///
200/// This fires on the failure of one *challenge*, which is not necessarily the
201/// failure of the order: a client may have another enabled type left to try.
202/// `error` is the human-readable detail, the same text the challenge object's
203/// `error` member carries back to the client.
204#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
205pub struct ChallengeFailedData {
206    pub profile: String,
207    pub order_id: String,
208    pub account_id: String,
209    pub authz_id: String,
210    pub challenge_id: String,
211    pub challenge_type: String,
212    pub identifier: String,
213    pub error: String,
214    pub client_ip: Option<String>,
215}
216
217/// One lifecycle event, carrying everything a template or `custom` script
218/// needs to describe it.
219///
220/// **Internally tagged on `hook`**, which is load-bearing twice over. It is the
221/// `custom` backend's stdin contract — a script reads `.hook` to tell one event
222/// from another — and it is what lets a queued delivery survive a restart, since
223/// a `notify_deliver` job payload is this enum and nothing else. The tag and the
224/// variant renaming reproduce exactly what [`Self::payload`] used to assemble by
225/// hand, so neither the script contract nor a row already in the queue changes
226/// shape; `payload_is_tagged_with_its_own_hook` pins that.
227#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
228#[serde(tag = "hook", rename_all = "snake_case")]
229pub enum NotifyEvent {
230    ProfileMounted(ProfileMountedData),
231    AccountCreated(AccountCreatedData),
232    AccountDeactivated(AccountDeactivatedData),
233    CertificateIssued(CertificateIssuedData),
234    CertificateRevoked(CertificateRevokedData),
235    ChallengeFailed(ChallengeFailedData),
236}
237
238impl NotifyEvent {
239    /// The event kind name — matches [`ALL_NOTIFY_EVENTS`] and the template
240    /// file stem used to render it (e.g. `email/certificate_issued.body.j2`).
241    pub fn kind(&self) -> &'static str {
242        match self {
243            Self::ProfileMounted(_) => "profile_mounted",
244            Self::AccountCreated(_) => "account_created",
245            Self::AccountDeactivated(_) => "account_deactivated",
246            Self::CertificateIssued(_) => "certificate_issued",
247            Self::CertificateRevoked(_) => "certificate_revoked",
248            Self::ChallengeFailed(_) => "challenge_failed",
249        }
250    }
251
252    /// The profile this event happened on.
253    pub fn profile(&self) -> &str {
254        match self {
255            Self::ProfileMounted(data) => &data.profile,
256            Self::AccountCreated(data) => &data.profile,
257            Self::AccountDeactivated(data) => &data.profile,
258            Self::CertificateIssued(data) => &data.profile,
259            Self::CertificateRevoked(data) => &data.profile,
260            Self::ChallengeFailed(data) => &data.profile,
261        }
262    }
263
264    /// The template rendering context: this event's own data, serialized.
265    pub(crate) fn context(&self) -> minijinja::Value {
266        match self {
267            Self::ProfileMounted(data) => minijinja::Value::from_serialize(data),
268            Self::AccountCreated(data) => minijinja::Value::from_serialize(data),
269            Self::AccountDeactivated(data) => minijinja::Value::from_serialize(data),
270            Self::CertificateIssued(data) => minijinja::Value::from_serialize(data),
271            Self::CertificateRevoked(data) => minijinja::Value::from_serialize(data),
272            Self::ChallengeFailed(data) => minijinja::Value::from_serialize(data),
273        }
274    }
275
276    /// This event's own data as a JSON object, tagged with `"hook"` — the
277    /// `custom` backend's stdin payload, and the body of a queued
278    /// `notify_deliver` job. Not used by the templating backends, which render
279    /// from [`Self::context`] instead.
280    ///
281    /// This is the enum's own `Serialize`: the internal tag *is* the `"hook"`
282    /// member that used to be spliced in here after the fact.
283    pub(crate) fn payload(&self) -> serde_json::Value {
284        serde_json::to_value(self).expect("notify event data always serializes to a JSON object")
285    }
286
287    /// This event's client address, when a request was in scope — `None` on
288    /// the `relay` async-completion path. Used by the `custom` backend
289    /// to fill `ACME_NOTIFY_CLIENT_IP`.
290    fn client_ip(&self) -> Option<&str> {
291        match self {
292            Self::ProfileMounted(_) => None,
293            Self::AccountCreated(data) => data.client_ip.as_deref(),
294            Self::AccountDeactivated(data) => data.client_ip.as_deref(),
295            Self::CertificateIssued(data) => data.client_ip.as_deref(),
296            Self::CertificateRevoked(data) => data.client_ip.as_deref(),
297            Self::ChallengeFailed(data) => data.client_ip.as_deref(),
298        }
299    }
300
301    /// This event's account id, when it has one. Used by the `custom` backend
302    /// to fill `ACME_NOTIFY_ACCOUNT_ID`.
303    fn account_id(&self) -> Option<&str> {
304        match self {
305            Self::ProfileMounted(_) => None,
306            Self::AccountCreated(data) => Some(&data.account_id),
307            Self::AccountDeactivated(data) => Some(&data.account_id),
308            Self::CertificateIssued(data) => Some(&data.account_id),
309            Self::CertificateRevoked(data) => Some(&data.account_id),
310            Self::ChallengeFailed(data) => Some(&data.account_id),
311        }
312    }
313
314    /// This event's order id, when it has one. Used by the `custom` backend
315    /// to fill `ACME_NOTIFY_ORDER_ID`.
316    fn order_id(&self) -> Option<&str> {
317        // Enumerated rather than `_ => None`, like every other accessor here:
318        // a wildcard would let a seventh variant carrying an order id compile,
319        // ship, and write an empty `ACME_NOTIFY_ORDER_ID` — the exact silent
320        // failure the accessor tests exist to make impossible.
321        match self {
322            Self::ProfileMounted(_) => None,
323            Self::AccountCreated(_) => None,
324            Self::AccountDeactivated(_) => None,
325            Self::CertificateIssued(data) => Some(&data.order_id),
326            Self::CertificateRevoked(data) => Some(&data.order_id),
327            Self::ChallengeFailed(data) => Some(&data.order_id),
328        }
329    }
330
331    /// This event's certificate serial, when it has one. Used by the `custom`
332    /// backend to fill `ACME_NOTIFY_CERT_SERIAL`.
333    fn cert_serial(&self) -> Option<&str> {
334        match self {
335            Self::ProfileMounted(_) => None,
336            Self::AccountCreated(_) => None,
337            Self::AccountDeactivated(_) => None,
338            Self::CertificateIssued(data) => Some(&data.cert_serial),
339            Self::CertificateRevoked(data) => Some(&data.cert_serial),
340            Self::ChallengeFailed(_) => None,
341        }
342    }
343
344    /// This event's requested identifiers, comma-joined. Used by the `custom`
345    /// backend to fill `ACME_NOTIFY_IDENTIFIERS`.
346    fn identifiers_joined(&self) -> String {
347        match self {
348            Self::ProfileMounted(_) => String::new(),
349            Self::AccountCreated(_) => String::new(),
350            Self::AccountDeactivated(_) => String::new(),
351            Self::CertificateIssued(data) => data.identifiers.join(","),
352            Self::CertificateRevoked(_) => String::new(),
353            Self::ChallengeFailed(_) => String::new(),
354        }
355    }
356}
357
358/// One configured backend, under the id a queued job addresses it by.
359///
360/// The id is **not** [`NotifyBackend::name`]: that is a `&'static str` a backend
361/// type answers, so every `custom` entry reports `"custom"` and two of them
362/// would be indistinguishable to a job row. It is built from the configuration
363/// instead — `email`, `webhook:<entry>`, `custom:<entry>` — which makes it
364/// stable across a restart, the property a durable payload needs. The same
365/// reasoning gives `filter::custom` its `ACME_FILTER_CHECK_NAME`.
366pub struct BackendSlot {
367    id: String,
368    /// The event kinds this backend's own `events` list admits. Filtering here
369    /// rather than inside a wrapper around `send` is what stops a job being
370    /// queued for a delivery that would immediately no-op.
371    events: HashSet<String>,
372    backend: Arc<dyn NotifyBackend>,
373}
374
375impl BackendSlot {
376    /// The id a job payload names this backend by.
377    #[must_use]
378    pub fn id(&self) -> &str {
379        &self.id
380    }
381
382    /// Whether this backend's `events` list admits `event`.
383    #[must_use]
384    pub fn wants(&self, event: &NotifyEvent) -> bool {
385        self.events.contains(event.kind())
386    }
387
388    /// Builds a slot directly.
389    ///
390    /// [`from_config`] is how one is normally made — it derives the ids from the
391    /// configuration, which is what keeps them stable across a restart. This is
392    /// for a caller assembling a dispatcher over a backend of its own, which in
393    /// practice means a test.
394    #[must_use]
395    pub fn new(id: impl Into<String>, backend: Arc<dyn NotifyBackend>, events: &[String]) -> Self {
396        Self {
397            id: id.into(),
398            events: events.iter().cloned().collect(),
399            backend,
400        }
401    }
402}
403
404/// The configured notify backends for one profile.
405///
406/// Cheap to clone behind the `Arc` it is always held in (`Profile::notify`,
407/// and the `profile name -> dispatcher` map handed to the `relay` signer
408/// backend and to [`NotifyJob`]).
409pub struct NotifyDispatcher {
410    profile: String,
411    slots: Vec<BackendSlot>,
412    jobs: JobQueue,
413}
414
415impl std::fmt::Debug for NotifyDispatcher {
416    /// `dyn NotifyBackend` is not `Debug`, so show the slot ids — the only part
417    /// worth reading anyway. Mirrors `FilterPolicy`'s own `Debug` impl.
418    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
419        formatter
420            .debug_struct("NotifyDispatcher")
421            .field("profile", &self.profile)
422            .field(
423                "backends",
424                &self.slots.iter().map(BackendSlot::id).collect::<Vec<_>>(),
425            )
426            .finish()
427    }
428}
429
430impl NotifyDispatcher {
431    /// Builds a dispatcher over already-constructed slots.
432    #[must_use]
433    pub fn new(profile: impl Into<String>, slots: Vec<BackendSlot>, jobs: JobQueue) -> Self {
434        Self {
435            profile: profile.into(),
436            slots,
437            jobs,
438        }
439    }
440
441    /// A dispatcher with no backends — every `dispatch` is a no-op. The shape a
442    /// profile with an empty `notify.enabled` gets, and what tests that do not
443    /// care about notifications want.
444    #[must_use]
445    pub fn disabled(jobs: JobQueue) -> Self {
446        Self::new("default", Vec::new(), jobs)
447    }
448
449    /// The profile whose `[notify]` section this dispatcher was built from.
450    #[must_use]
451    pub fn profile(&self) -> &str {
452        &self.profile
453    }
454
455    /// The backend registered under `id`, if it is still configured. `None`
456    /// after a configuration change removed a backend a queued job still names.
457    #[must_use]
458    pub fn slot(&self, id: &str) -> Option<&BackendSlot> {
459        self.slots.iter().find(|slot| slot.id == id)
460    }
461
462    /// Queues one delivery per backend that wants this event, and returns.
463    ///
464    /// Not `?`-able and with nothing useful to return: a notification can never
465    /// affect the ACME response that triggered it, which is why this goes
466    /// through [`JobQueue::enqueue_or_log`] — a database failure here is logged
467    /// and swallowed exactly like the delivery failure it stands in for.
468    ///
469    /// The `delivery_id` is minted per call, so the queue's identity index never
470    /// refuses a genuine second occurrence of the same event: two dispatches
471    /// mean two notifications, as they always did. It is shared across the
472    /// backends of one call purely so an operator can correlate them in a log.
473    pub async fn dispatch(&self, event: NotifyEvent) {
474        if self.slots.is_empty() {
475            return;
476        }
477
478        let kind = event.kind();
479        let payload = event.payload();
480        let delivery_id = uuid::Uuid::new_v4().to_string();
481        for slot in &self.slots {
482            if !slot.wants(&event) {
483                continue;
484            }
485            let spec = JobSpec::now(NOTIFY_JOB_KIND, format!("{delivery_id}:{}", slot.id))
486                .with_payload(serde_json::json!({
487                    "profile": self.profile,
488                    "backend": slot.id,
489                    "event": payload,
490                }));
491            if self.jobs.enqueue_or_log(spec).await {
492                info!(
493                    event = "notify_delivery_queued",
494                    outcome = "progress",
495                    profile = %self.profile,
496                    backend = %slot.id,
497                    kind,
498                    delivery_id = %delivery_id,
499                );
500            }
501        }
502    }
503
504    /// Runs one delivery, now, against the backend registered under `id`.
505    ///
506    /// What [`NotifyJob`] calls for each claimed row, and what a test calls when
507    /// it wants delivery without a runner in the way. `Ok(None)` means no such
508    /// backend is configured any more — a decision for the caller, since a
509    /// handler must retire that job rather than retry it for ever.
510    pub(crate) async fn deliver(
511        &self,
512        id: &str,
513        event: &NotifyEvent,
514    ) -> Option<Result<(), NotifyError>> {
515        let slot = self.slot(id)?;
516        Some(slot.backend.send(event).await)
517    }
518}
519
520/// An `events` list naming something outside [`ALL_NOTIFY_EVENTS`] is a
521/// startup error, the same treatment `filter.enabled`/`challenge.enabled`
522/// give an unknown name — a typo here should stop the server, not silently
523/// mean "never fires".
524fn validate_events(field: &str, events: &[String]) -> anyhow::Result<()> {
525    for event in events {
526        anyhow::ensure!(
527            ALL_NOTIFY_EVENTS.contains(&event.as_str()),
528            "{field}: unknown event `{event}` (expected one of {ALL_NOTIFY_EVENTS:?})"
529        );
530    }
531    Ok(())
532}
533
534/// Builds the configured notify dispatcher. Called once per profile at
535/// startup (via [`build_registry`]), so it may fail fast.
536///
537/// `jobs` is the enqueue side of the durable queue: every delivery is a row on
538/// it, so a dispatcher cannot be built without one.
539pub fn from_config(
540    profile: &str,
541    cfg: &NotifyConfig,
542    outbound: crate::http_client::Outbound,
543    jobs: &JobQueue,
544) -> anyhow::Result<Arc<NotifyDispatcher>> {
545    // Built once and shared: both templating backends render from the same
546    // `template_dir` override, and `Environment` is cheap to clone (it holds
547    // an `Arc`-like handle to its loader internally) but not to construct.
548    let env = Arc::new(build_environment(&cfg.template_dir));
549
550    let mut slots: Vec<BackendSlot> = Vec::with_capacity(cfg.enabled.len());
551    for name in &cfg.enabled {
552        let built: Vec<BackendSlot> = match name.as_str() {
553            "email" => {
554                validate_events("notify.email.events", &cfg.email.events)?;
555                vec![BackendSlot::new(
556                    "email",
557                    Arc::new(email::EmailNotifier::from_config(&cfg.email, env.clone())?),
558                    &cfg.email.events,
559                )]
560            }
561            "webhook" => build_webhook_slots(cfg, &env, outbound.clone())?,
562            "custom" => build_custom_slots(cfg)?,
563            // Refused by name, the `signer.backend = "acme_proxy"` -> `relay`
564            // treatment: the `mattermost` backend was one provider's payload
565            // shape frozen into a copy of the webhook transport, and every
566            // other part of it now lives in `webhook`. An unmigrated
567            // configuration stops the server rather than coming up looking
568            // configured and notifying nobody.
569            "mattermost" => anyhow::bail!(
570                "notify.enabled: `mattermost` was replaced by `webhook`. Use \
571                 notify.enabled = [\"webhook\"] with a [notify.webhook.<name>] entry \
572                 whose `url` is the incoming webhook URL; the default `body` is \
573                 already the payload Mattermost accepts. `channel` and `username` \
574                 move into that `body`"
575            ),
576            other => anyhow::bail!("unknown notify backend: {other}"),
577        };
578        slots.extend(built);
579    }
580
581    if slots.is_empty() {
582        info!(
583            event = "notify_disabled",
584            outcome = "success",
585            "no notification backends configured"
586        );
587    } else {
588        info!(event = "notify_enabled", outcome = "success", backends = ?cfg.enabled);
589    }
590
591    Ok(Arc::new(NotifyDispatcher::new(
592        profile,
593        slots,
594        jobs.clone(),
595    )))
596}
597
598/// One slot per `notify.custom` entry, addressed `custom:<entry>`.
599///
600/// The entry name is what makes two custom backends tell-apart-able: every
601/// [`custom::CustomScriptNotifier`] answers `"custom"` to
602/// [`NotifyBackend::name`], so a job payload naming that alone could not say
603/// which script it meant. `resolve_named_entries` has already refused a name
604/// that is not a valid environment-variable segment, so the id is safe to build
605/// from it.
606fn build_custom_slots(cfg: &NotifyConfig) -> anyhow::Result<Vec<BackendSlot>> {
607    crate::config::resolve_named_entries(
608        "notify.custom",
609        "notify.custom_enabled",
610        "custom",
611        &cfg.custom,
612        &cfg.custom_enabled,
613    )?
614    .into_iter()
615    .map(|(name, script)| -> anyhow::Result<BackendSlot> {
616        validate_events(&format!("notify.custom.{name}.events"), &script.events)?;
617        let backend = custom::CustomScriptNotifier::from_config(script)?;
618        Ok(BackendSlot::new(
619            format!("custom:{name}"),
620            Arc::new(backend),
621            &script.events,
622        ))
623    })
624    .collect()
625}
626
627/// One slot per selected `notify.webhook` entry, addressed `webhook:<entry>`.
628///
629/// Same shape and same reasoning as [`build_custom_slots`] — every
630/// [`webhook::WebhookNotifier`] answers `"webhook"` to
631/// [`NotifyBackend::name`], so the entry name is what tells two of them apart
632/// in a durable job payload. The environment is passed by reference rather than
633/// cloned in: each notifier compiles its own `body` into a clone of it, so a
634/// template that does not parse is a startup error.
635fn build_webhook_slots(
636    cfg: &NotifyConfig,
637    env: &minijinja::Environment<'static>,
638    outbound: crate::http_client::Outbound,
639) -> anyhow::Result<Vec<BackendSlot>> {
640    crate::config::resolve_named_entries(
641        "notify.webhook",
642        "notify.webhook_enabled",
643        "webhook",
644        &cfg.webhook,
645        &cfg.webhook_enabled,
646    )?
647    .into_iter()
648    .map(|(name, entry)| -> anyhow::Result<BackendSlot> {
649        validate_events(&format!("notify.webhook.{name}.events"), &entry.events)?;
650        let backend = webhook::WebhookNotifier::from_config(name, entry, env, outbound.clone())?;
651        Ok(BackendSlot::new(
652            format!("webhook:{name}"),
653            Arc::new(backend),
654            &entry.events,
655        ))
656    })
657    .collect()
658}
659
660/// One generation's `profile name -> dispatcher` map.
661///
662/// Named because it appears in four signatures and clippy is right that the
663/// spelled-out form is unreadable in all of them.
664pub type DispatcherMap = HashMap<String, Arc<NotifyDispatcher>>;
665
666/// The writing half of a [`Notifiers`] handle.
667pub type NotifiersSender = tokio::sync::watch::Sender<Arc<DispatcherMap>>;
668
669/// The `profile name -> dispatcher` map, as a handle that survives a
670/// configuration reload.
671///
672/// The map itself is rebuilt whole on every generation, but two of its readers
673/// outlive a generation: a signer backend captures it at construction (and
674/// signer backends are carried across reloads rather than rebuilt), and
675/// [`NotifyJob`] captures it at registration. Handing those two a plain `Arc`
676/// pinned them to generation zero, which is worse than stale — a request served
677/// by a *new* router writes a `notify_deliver` row naming a slot id from the
678/// *new* configuration, and a `NotifyJob` still holding the old map would answer
679/// [`JobOutcome::Failed`](crate::jobs::JobOutcome::Failed) for it. Permanently:
680/// an unknown backend id is retired rather than retried, by design.
681///
682/// [`tokio::sync::watch::Receiver::borrow`] takes `&self` and does not mark the
683/// value seen, so this stays `Clone + Send + Sync` and [`get`](Self::get) is
684/// callable from any task. A generation is published with a single synchronous
685/// [`tokio::sync::watch::Sender::send_replace`], so no reader can observe a
686/// half-built map.
687#[derive(Clone)]
688pub struct Notifiers(tokio::sync::watch::Receiver<Arc<DispatcherMap>>);
689
690impl Notifiers {
691    /// The dispatcher for `profile` in the current generation, if it is mounted.
692    #[must_use]
693    pub fn get(&self, profile: &str) -> Option<Arc<NotifyDispatcher>> {
694        // The `Ref` guard dies at the semicolon, so nothing holds the lock
695        // across an await — which is the whole reason this returns an owned
696        // `Arc` rather than lending one out.
697        self.0.borrow().get(profile).cloned()
698    }
699}
700
701/// A fixed map, for every caller that has no reload to serve — the tests, and
702/// any construction that predates the first generation being published.
703///
704/// Dropping the sender leaves `borrow` working for ever (only `changed()` ever
705/// errors on a closed channel), so this is a genuine constant, not a channel
706/// that will go quiet. It is also what lets every existing call site pass an
707/// `Arc<HashMap<..>>` unchanged.
708impl From<Arc<DispatcherMap>> for Notifiers {
709    fn from(map: Arc<DispatcherMap>) -> Self {
710        Self(tokio::sync::watch::channel(map).1)
711    }
712}
713
714impl From<DispatcherMap> for Notifiers {
715    fn from(map: DispatcherMap) -> Self {
716        Arc::new(map).into()
717    }
718}
719
720/// One cell the reload path replaces, and every reader sees the new map at the
721/// next `get`.
722#[must_use]
723pub fn notifiers_channel(initial: DispatcherMap) -> (NotifiersSender, Notifiers) {
724    let (sender, receiver) = tokio::sync::watch::channel(Arc::new(initial));
725    (sender, Notifiers(receiver))
726}
727
728/// Builds one [`NotifyDispatcher`] per resolved profile, keyed by profile
729/// name — the map the `relay` signer backend needs to notify the right
730/// profile from its background completion task, where there is no
731/// `AppState`/`Profile` to reach through.
732pub fn build_registry(
733    profiles: &[ProfileConfig],
734    outbound: crate::http_client::Outbound,
735    jobs: &JobQueue,
736) -> anyhow::Result<HashMap<String, Arc<NotifyDispatcher>>> {
737    let mut registry = HashMap::with_capacity(profiles.len());
738    for profile in profiles {
739        let dispatcher = from_config(
740            &profile.name,
741            &profile.sections.notify,
742            outbound.clone(),
743            jobs,
744        )
745        .map_err(|error| anyhow::anyhow!("profile `{}`: {error}", profile.name))?;
746        registry.insert(profile.name.clone(), dispatcher);
747    }
748    Ok(registry)
749}
750
751/// One embedded template: the name it is known by *is* the path it lives at.
752///
753/// Each entry used to spell the filename twice — once as the key and once
754/// inside `include_str!` — across 18 entries here and 35 in
755/// `webadmin::pages::templates`. Two spellings of one name is a mismatch
756/// waiting to happen, and the failure would be a template that renders as
757/// missing at delivery time rather than at build time.
758macro_rules! embed {
759    ($name:literal) => {
760        ($name, include_str!(concat!("templates/", $name)))
761    };
762}
763
764/// Every default template, embedded so the server needs no external
765/// `templates/` directory to run. Keyed the same way [`build_environment`]'s
766/// loader looks them up: `"<backend>/<event>.<subject|body>.j2"` for email,
767/// `"<backend>/<event>.j2"` for the webhook message.
768static EMBEDDED_TEMPLATES: LazyLock<HashMap<&'static str, &'static str>> = LazyLock::new(|| {
769    HashMap::from([
770        embed!("email/profile_mounted.subject.j2"),
771        embed!("email/profile_mounted.body.j2"),
772        embed!("email/account_created.subject.j2"),
773        embed!("email/account_created.body.j2"),
774        embed!("email/account_deactivated.subject.j2"),
775        embed!("email/account_deactivated.body.j2"),
776        embed!("email/certificate_issued.subject.j2"),
777        embed!("email/certificate_issued.body.j2"),
778        embed!("email/certificate_revoked.subject.j2"),
779        embed!("email/certificate_revoked.body.j2"),
780        embed!("email/challenge_failed.subject.j2"),
781        embed!("email/challenge_failed.body.j2"),
782        embed!("webhook/profile_mounted.j2"),
783        embed!("webhook/account_created.j2"),
784        embed!("webhook/account_deactivated.j2"),
785        embed!("webhook/certificate_issued.j2"),
786        embed!("webhook/certificate_revoked.j2"),
787        embed!("webhook/challenge_failed.j2"),
788    ])
789});
790
791/// Builds a template environment: `template_dir` (if set) is checked for each
792/// named template before falling back to the compiled-in default, so an
793/// operator can override a single message and leave every other one at its
794/// default.
795pub(crate) fn build_environment(template_dir: &str) -> minijinja::Environment<'static> {
796    crate::templating::loader_env(template_dir, &EMBEDDED_TEMPLATES)
797}
798
799/// Renders one named template against `event`'s own data.
800///
801/// Both failures are **permanent**: a template that is absent or does not
802/// compile will be just as absent on the fifth attempt, so retrying only delays
803/// the log line that tells the operator to fix it.
804pub(crate) fn render(
805    env: &minijinja::Environment<'static>,
806    template_name: &str,
807    event: &NotifyEvent,
808) -> Result<String, NotifyError> {
809    let template = env.get_template(template_name).map_err(|error| {
810        NotifyError::permanent(format!("template `{template_name}` not found: {error}"))
811    })?;
812    template.render(event.context()).map_err(|error| {
813        NotifyError::permanent(format!("template `{template_name}` failed: {error}"))
814    })
815}
816
817#[cfg(test)]
818pub(crate) mod tests {
819    use super::*;
820    use crate::config::CustomNotifyConfig;
821    use crate::sqlite::db::Database;
822    use crate::sqlite::job::Job;
823    use std::collections::BTreeMap;
824    use std::sync::Mutex;
825
826    /// The shared resolver `Profile::build_all` supplies at startup.
827    fn test_resolver() -> Arc<dyn crate::dns::Resolver> {
828        Arc::new(crate::dns::HickoryResolver::from_system_uncached().unwrap())
829    }
830
831    /// A queue over an in-memory database, for the assertions that read back
832    /// the rows `dispatch` wrote.
833    pub(crate) async fn test_queue() -> JobQueue {
834        let database = Arc::new(Database::connect_in_memory().await.unwrap());
835        JobQueue::new(database, &crate::config::JobsConfig::default())
836    }
837
838    /// How a backend failed, when it is configured to.
839    #[derive(Default, Clone, Copy)]
840    enum Failure {
841        #[default]
842        None,
843        Retryable,
844        Permanent,
845    }
846
847    /// A backend recording every event it received, for asserting dispatch
848    /// behavior without a real SMTP/HTTP/script target.
849    #[derive(Default)]
850    pub(crate) struct RecordingNotifyBackend {
851        pub(crate) events: Mutex<Vec<NotifyEvent>>,
852        fail: Failure,
853    }
854
855    impl RecordingNotifyBackend {
856        /// Fails the way a refused connection does: worth another attempt.
857        pub(crate) fn failing() -> Self {
858            Self {
859                events: Mutex::new(Vec::new()),
860                fail: Failure::Retryable,
861            }
862        }
863
864        /// Fails the way a missing template does: another attempt is pointless.
865        pub(crate) fn failing_permanently() -> Self {
866            Self {
867                events: Mutex::new(Vec::new()),
868                fail: Failure::Permanent,
869            }
870        }
871    }
872
873    #[async_trait]
874    impl NotifyBackend for RecordingNotifyBackend {
875        fn name(&self) -> &'static str {
876            "recording"
877        }
878
879        async fn send(&self, event: &NotifyEvent) -> Result<(), NotifyError> {
880            self.events.lock().unwrap().push(event.clone());
881            match self.fail {
882                Failure::None => Ok(()),
883                Failure::Retryable => Err(NotifyError::new("recording backend configured to fail")),
884                Failure::Permanent => Err(NotifyError::permanent(
885                    "recording backend configured to fail permanently",
886                )),
887            }
888        }
889    }
890
891    fn profile_mounted(profile: &str) -> NotifyEvent {
892        NotifyEvent::ProfileMounted(ProfileMountedData {
893            profile: profile.to_string(),
894        })
895    }
896
897    /// One of every variant, all on the same profile — the input to the
898    /// accessor sweep below, which is what the `custom` backend's environment
899    /// and both templating backends' contexts are built from.
900    fn every_event() -> Vec<NotifyEvent> {
901        vec![
902            profile_mounted("p"),
903            NotifyEvent::AccountCreated(AccountCreatedData {
904                profile: "p".to_string(),
905                account_id: "acct-1".to_string(),
906                contact: vec!["mailto:a@example.com".to_string()],
907                client_ip: Some("203.0.113.5".to_string()),
908            }),
909            NotifyEvent::AccountDeactivated(AccountDeactivatedData {
910                profile: "p".to_string(),
911                account_id: "acct-1".to_string(),
912                client_ip: Some("203.0.113.5".to_string()),
913            }),
914            NotifyEvent::CertificateIssued(CertificateIssuedData {
915                profile: "p".to_string(),
916                order_id: "ord-1".to_string(),
917                account_id: "acct-1".to_string(),
918                cert_serial: "0a0b".to_string(),
919                identifiers: vec!["a.example.com".to_string(), "b.example.com".to_string()],
920                client_ip: Some("203.0.113.5".to_string()),
921            }),
922            NotifyEvent::CertificateRevoked(CertificateRevokedData {
923                profile: "p".to_string(),
924                order_id: "ord-1".to_string(),
925                account_id: "acct-1".to_string(),
926                cert_serial: "0a0b".to_string(),
927                reason: Some(1),
928                client_ip: None,
929            }),
930            NotifyEvent::ChallengeFailed(ChallengeFailedData {
931                profile: "p".to_string(),
932                order_id: "ord-1".to_string(),
933                account_id: "acct-1".to_string(),
934                authz_id: "authz-1".to_string(),
935                challenge_id: "chall-1".to_string(),
936                challenge_type: "http-01".to_string(),
937                identifier: "a.example.com".to_string(),
938                error: "connection refused".to_string(),
939                client_ip: Some("203.0.113.5".to_string()),
940            }),
941        ]
942    }
943
944    /// Every variant answers every accessor. These are wide `match`es over an
945    /// enum that grows, so a new variant added to only some of them would
946    /// otherwise surface as a silently empty template field.
947    #[test]
948    fn every_event_answers_every_accessor() {
949        let events = every_event();
950        assert_eq!(
951            events.len(),
952            ALL_NOTIFY_EVENTS.len(),
953            "every declared event kind needs a sample here"
954        );
955
956        for event in &events {
957            assert_eq!(event.profile(), "p");
958            assert!(
959                ALL_NOTIFY_EVENTS.contains(&event.kind()),
960                "{}",
961                event.kind()
962            );
963
964            // The rendering context and the `custom` backend's stdin payload
965            // are built from the same data by two different routes.
966            assert!(!event.context().is_undefined());
967            let payload = event.payload();
968            assert_eq!(
969                payload.get("hook").and_then(|v| v.as_str()),
970                Some(event.kind()),
971                "the payload must name its own hook"
972            );
973            assert_eq!(payload.get("profile").and_then(|v| v.as_str()), Some("p"));
974        }
975
976        // `profile_mounted` is the one event with no account, order or client
977        // behind it: it happens at startup, outside any request.
978        let mounted = &events[0];
979        assert_eq!(mounted.client_ip(), None);
980        assert_eq!(mounted.account_id(), None);
981        assert_eq!(mounted.order_id(), None);
982        assert_eq!(mounted.cert_serial(), None);
983        assert_eq!(mounted.identifiers_joined(), "");
984
985        for event in &events[1..] {
986            assert_eq!(event.account_id(), Some("acct-1"));
987        }
988        // A revocation reached through the admin CLI has no client address.
989        assert_eq!(events[4].client_ip(), None);
990        assert_eq!(events[1].client_ip(), Some("203.0.113.5"));
991
992        // Order and serial only exist once there is a certificate to name.
993        assert_eq!(events[1].order_id(), None);
994        assert_eq!(events[2].cert_serial(), None);
995        for event in &events[3..] {
996            assert_eq!(event.order_id(), Some("ord-1"));
997        }
998        assert_eq!(events[3].cert_serial(), Some("0a0b"));
999        assert_eq!(events[4].cert_serial(), Some("0a0b"));
1000        assert_eq!(events[5].cert_serial(), None);
1001
1002        // Only issuance carries the names the certificate is for.
1003        assert_eq!(
1004            events[3].identifiers_joined(),
1005            "a.example.com,b.example.com"
1006        );
1007        assert_eq!(events[5].identifiers_joined(), "");
1008    }
1009
1010    /// `dyn NotifyBackend` is not `Debug`, so the dispatcher renders the slot
1011    /// ids instead — the part a startup log is read for, and now also the part a
1012    /// queued job addresses.
1013    #[tokio::test]
1014    async fn the_dispatcher_debug_names_its_backends() {
1015        let queue = test_queue().await;
1016        let dispatcher = NotifyDispatcher::new(
1017            "le",
1018            vec![BackendSlot::new(
1019                "recording",
1020                Arc::new(RecordingNotifyBackend::default()),
1021                &every_kind(),
1022            )],
1023            queue.clone(),
1024        );
1025        let rendered = format!("{dispatcher:?}");
1026        assert!(rendered.contains("NotifyDispatcher"), "{rendered}");
1027        assert!(rendered.contains("recording"), "{rendered}");
1028        assert!(rendered.contains("le"), "{rendered}");
1029
1030        assert!(format!("{:?}", NotifyDispatcher::disabled(queue)).contains("[]"));
1031    }
1032
1033    /// The reload property: a reader that took its handle before the swap sees
1034    /// the map that came *after* it.
1035    ///
1036    /// This is what a signer backend and [`NotifyJob`] rely on — both are built
1037    /// once and outlive a configuration generation, so a captured `Arc` would
1038    /// pin them to whatever was configured when the process started.
1039    #[tokio::test]
1040    async fn a_handle_taken_before_a_swap_reads_the_map_after_it() {
1041        let queue = test_queue().await;
1042        let (sender, notifiers) = notifiers_channel(HashMap::new());
1043        assert!(notifiers.get("le").is_none());
1044
1045        let mut next = HashMap::new();
1046        next.insert(
1047            "le".to_string(),
1048            Arc::new(NotifyDispatcher::disabled(queue.clone())),
1049        );
1050        sender.send_replace(Arc::new(next));
1051
1052        assert!(notifiers.get("le").is_some());
1053        // A profile the new generation does not mount is absent, not stale.
1054        assert!(notifiers.get("staging").is_none());
1055
1056        // And a swap back is seen too: this is a cell, not a latch.
1057        sender.send_replace(Arc::new(HashMap::new()));
1058        assert!(notifiers.get("le").is_none());
1059    }
1060
1061    /// A fixed map keeps answering after its sender is gone.
1062    ///
1063    /// The `From` impl drops the sender on the spot, which is what lets every
1064    /// caller with no reload to serve — the tests, and anything built before the
1065    /// first generation is published — pass a plain map. If a closed channel
1066    /// made `borrow` fail, that conversion would be a trap rather than a
1067    /// convenience.
1068    #[tokio::test]
1069    async fn a_fixed_map_survives_its_sender_being_dropped() {
1070        let queue = test_queue().await;
1071        let mut map = HashMap::new();
1072        map.insert(
1073            "le".to_string(),
1074            Arc::new(NotifyDispatcher::disabled(queue)),
1075        );
1076
1077        let notifiers: Notifiers = map.into();
1078        assert!(notifiers.get("le").is_some());
1079        // Cloned handles are the same cell, and equally durable.
1080        assert!(notifiers.clone().get("le").is_some());
1081    }
1082
1083    /// Every event kind, as a backend's own `events` list would spell them.
1084    fn every_kind() -> Vec<String> {
1085        ALL_NOTIFY_EVENTS.iter().map(|k| (*k).to_string()).collect()
1086    }
1087
1088    /// A dispatcher over one recording backend, plus the queue its rows land in.
1089    async fn recording_dispatcher(
1090        events: &[String],
1091    ) -> (Arc<NotifyDispatcher>, Arc<RecordingNotifyBackend>, JobQueue) {
1092        let queue = test_queue().await;
1093        let recorder = Arc::new(RecordingNotifyBackend::default());
1094        let dispatcher = Arc::new(NotifyDispatcher::new(
1095            "le",
1096            vec![BackendSlot::new("recording", recorder.clone(), events)],
1097            queue.clone(),
1098        ));
1099        (dispatcher, recorder, queue)
1100    }
1101
1102    /// The property the whole change turns on: `dispatch` delivers nothing
1103    /// itself, it writes a row. Nothing has reached the backend when it returns.
1104    #[tokio::test]
1105    async fn dispatch_queues_a_row_rather_than_delivering() {
1106        let (dispatcher, recorder, queue) = recording_dispatcher(&every_kind()).await;
1107
1108        dispatcher.dispatch(profile_mounted("le")).await;
1109
1110        assert!(
1111            recorder.events.lock().unwrap().is_empty(),
1112            "dispatch must not deliver inline"
1113        );
1114        let queued = Job::count_live(NOTIFY_JOB_KIND, queue.database())
1115            .await
1116            .unwrap();
1117        assert_eq!(queued, 1, "one backend, one row");
1118    }
1119
1120    /// The `events` list is applied at *enqueue*, so a backend that does not
1121    /// want an event costs no row at all — not a row that runs and no-ops.
1122    #[tokio::test]
1123    async fn a_backend_that_does_not_want_the_event_gets_no_row() {
1124        let (dispatcher, _recorder, queue) =
1125            recording_dispatcher(&["certificate_issued".to_string()]).await;
1126
1127        dispatcher.dispatch(profile_mounted("le")).await;
1128
1129        assert_eq!(
1130            Job::count_live(NOTIFY_JOB_KIND, queue.database())
1131                .await
1132                .unwrap(),
1133            0
1134        );
1135    }
1136
1137    /// One row **per backend**, so a retry against a failing one never re-sends
1138    /// through a healthy one that already delivered.
1139    #[tokio::test]
1140    async fn one_dispatch_queues_one_row_per_wanting_backend() {
1141        let queue = test_queue().await;
1142        let dispatcher = NotifyDispatcher::new(
1143            "le",
1144            vec![
1145                BackendSlot::new(
1146                    "email",
1147                    Arc::new(RecordingNotifyBackend::default()),
1148                    &every_kind(),
1149                ),
1150                BackendSlot::new(
1151                    "custom:webhook",
1152                    Arc::new(RecordingNotifyBackend::default()),
1153                    &every_kind(),
1154                ),
1155                BackendSlot::new(
1156                    "custom:pager",
1157                    Arc::new(RecordingNotifyBackend::default()),
1158                    &["certificate_revoked".to_string()],
1159                ),
1160            ],
1161            queue.clone(),
1162        );
1163
1164        dispatcher.dispatch(profile_mounted("le")).await;
1165
1166        assert_eq!(
1167            Job::count_live(NOTIFY_JOB_KIND, queue.database())
1168                .await
1169                .unwrap(),
1170            2,
1171            "the third backend does not want this kind"
1172        );
1173    }
1174
1175    /// Two dispatches of the same event are two notifications, as they always
1176    /// were: the per-call `delivery_id` keeps the identity index from mistaking
1177    /// the second for a duplicate of the first.
1178    #[tokio::test]
1179    async fn the_same_event_dispatched_twice_queues_twice() {
1180        let (dispatcher, _recorder, queue) = recording_dispatcher(&every_kind()).await;
1181
1182        dispatcher.dispatch(profile_mounted("le")).await;
1183        dispatcher.dispatch(profile_mounted("le")).await;
1184
1185        assert_eq!(
1186            Job::count_live(NOTIFY_JOB_KIND, queue.database())
1187                .await
1188                .unwrap(),
1189            2
1190        );
1191    }
1192
1193    /// A dispatcher with nothing configured writes nothing at all — the queue is
1194    /// not a place to park work no backend will ever ask for.
1195    #[tokio::test]
1196    async fn a_disabled_dispatcher_queues_nothing() {
1197        let queue = test_queue().await;
1198        let dispatcher = NotifyDispatcher::disabled(queue.clone());
1199
1200        dispatcher.dispatch(profile_mounted("le")).await;
1201
1202        assert_eq!(
1203            Job::count_live(NOTIFY_JOB_KIND, queue.database())
1204                .await
1205                .unwrap(),
1206            0
1207        );
1208    }
1209
1210    /// A database that cannot take the row must not become a failed ACME
1211    /// request: `dispatch` returns `()` and there is nowhere to put the error.
1212    #[tokio::test]
1213    async fn a_database_failure_is_swallowed_by_dispatch() {
1214        let (dispatcher, _recorder, queue) = recording_dispatcher(&every_kind()).await;
1215        queue.database().pool.close().await;
1216
1217        dispatcher.dispatch(profile_mounted("le")).await;
1218    }
1219
1220    /// `deliver` is the seam the job handler runs through, and the answer for an
1221    /// id nobody has is `None` rather than an error — the handler has to tell
1222    /// "this backend refused" from "this backend is gone".
1223    #[tokio::test]
1224    async fn deliver_reaches_one_backend_and_reports_an_unknown_id() {
1225        let (dispatcher, recorder, _queue) = recording_dispatcher(&every_kind()).await;
1226
1227        let outcome = dispatcher
1228            .deliver("recording", &profile_mounted("le"))
1229            .await;
1230        assert!(matches!(outcome, Some(Ok(()))));
1231        assert_eq!(recorder.events.lock().unwrap().len(), 1);
1232
1233        assert!(
1234            dispatcher
1235                .deliver("carrier-pigeon", &profile_mounted("le"))
1236                .await
1237                .is_none()
1238        );
1239        assert_eq!(
1240            recorder.events.lock().unwrap().len(),
1241            1,
1242            "an unknown id must reach no backend at all"
1243        );
1244    }
1245
1246    /// The selector's own arms: each name reaches its constructor and lands in
1247    /// the dispatcher. Neither backend touches the network at build time —
1248    /// `lettre` only assembles a transport and `webhook` only parses a URL and
1249    /// compiles a template — so this is a pure configuration test.
1250    #[tokio::test]
1251    async fn each_backend_name_builds_its_own_backend() {
1252        let cfg = NotifyConfig {
1253            enabled: vec!["email".to_string(), "webhook".to_string()],
1254            email: crate::config::EmailNotifyConfig {
1255                smtp_host: "smtp.example.com".to_string(),
1256                from: "acme@example.com".to_string(),
1257                to: vec!["ops@example.com".to_string()],
1258                // Every `smtp_security` value builds a different transport.
1259                smtp_security: "none".to_string(),
1260                smtp_username: "user".to_string(),
1261                smtp_password: "pass".to_string(),
1262                ..crate::config::EmailNotifyConfig::default()
1263            },
1264            webhook_enabled: vec!["chat".to_string()],
1265            webhook: BTreeMap::from([("chat".to_string(), webhook_entry())]),
1266            ..NotifyConfig::default()
1267        };
1268
1269        let dispatcher = from_config(
1270            "le",
1271            &cfg,
1272            crate::testutil::outbound_with(test_resolver()),
1273            &test_queue().await,
1274        )
1275        .expect("both backends must build");
1276        let rendered = format!("{dispatcher:?}");
1277        assert!(rendered.contains("email"), "{rendered}");
1278        assert!(rendered.contains("webhook:chat"), "{rendered}");
1279    }
1280
1281    fn webhook_entry() -> crate::config::WebhookNotifyConfig {
1282        crate::config::WebhookNotifyConfig {
1283            url: "https://chat.example.com/hooks/abc".to_string(),
1284            ..crate::config::WebhookNotifyConfig::default()
1285        }
1286    }
1287
1288    /// Two webhook entries are two slots with distinct ids — the `custom`
1289    /// property this backend inherits and needs for the same reason: every
1290    /// entry answers `"webhook"` to `NotifyBackend::name`, so a job row naming
1291    /// that alone could not say which endpoint it meant, and a retry would
1292    /// re-send through the one that already succeeded.
1293    #[tokio::test]
1294    async fn two_webhook_entries_get_distinct_slot_ids() {
1295        let cfg = NotifyConfig {
1296            enabled: vec!["webhook".to_string()],
1297            webhook_enabled: vec!["slack".to_string(), "teams".to_string()],
1298            webhook: BTreeMap::from([
1299                ("slack".to_string(), webhook_entry()),
1300                ("teams".to_string(), webhook_entry()),
1301            ]),
1302            ..NotifyConfig::default()
1303        };
1304
1305        let dispatcher = from_config(
1306            "le",
1307            &cfg,
1308            crate::testutil::outbound_with(test_resolver()),
1309            &test_queue().await,
1310        )
1311        .expect("both entries must build");
1312
1313        assert!(dispatcher.slot("webhook:slack").is_some());
1314        assert!(dispatcher.slot("webhook:teams").is_some());
1315        assert!(dispatcher.slot("webhook").is_none());
1316    }
1317
1318    /// The `mattermost` backend is gone and is refused **by name**, so an
1319    /// unmigrated configuration stops the server rather than coming up looking
1320    /// configured and notifying nobody.
1321    #[tokio::test]
1322    async fn the_removed_mattermost_backend_is_refused_by_name() {
1323        let cfg = NotifyConfig {
1324            enabled: vec!["mattermost".to_string()],
1325            ..NotifyConfig::default()
1326        };
1327        let error = from_config(
1328            "le",
1329            &cfg,
1330            crate::testutil::outbound_with(test_resolver()),
1331            &test_queue().await,
1332        )
1333        .unwrap_err()
1334        .to_string();
1335        assert!(error.contains("mattermost"), "{error}");
1336        assert!(error.contains("webhook"), "{error}");
1337    }
1338
1339    /// The two other `smtp_security` values, which each pick a different
1340    /// `lettre` builder, plus the one that is not a value at all.
1341    #[tokio::test]
1342    async fn every_smtp_security_mode_is_recognised() {
1343        for mode in ["starttls", "tls", "none"] {
1344            let cfg = email_config(mode);
1345            from_config(
1346                "le",
1347                &cfg,
1348                crate::testutil::outbound_with(test_resolver()),
1349                &test_queue().await,
1350            )
1351            .unwrap_or_else(|error| panic!("`{mode}` must build: {error}"));
1352        }
1353
1354        let error = from_config(
1355            "le",
1356            &email_config("carrier-pigeon"),
1357            crate::testutil::outbound_with(test_resolver()),
1358            &test_queue().await,
1359        )
1360        .unwrap_err()
1361        .to_string();
1362        assert!(error.contains("smtp_security"), "{error}");
1363    }
1364
1365    fn email_config(smtp_security: &str) -> NotifyConfig {
1366        NotifyConfig {
1367            enabled: vec!["email".to_string()],
1368            email: crate::config::EmailNotifyConfig {
1369                smtp_host: "smtp.example.com".to_string(),
1370                from: "acme@example.com".to_string(),
1371                to: vec!["ops@example.com".to_string()],
1372                smtp_security: smtp_security.to_string(),
1373                ..crate::config::EmailNotifyConfig::default()
1374            },
1375            ..NotifyConfig::default()
1376        }
1377    }
1378
1379    /// An event name nobody recognises is caught per backend, before the
1380    /// backend itself is built — otherwise a typo would silently mean "never
1381    /// notify" for that channel.
1382    #[tokio::test]
1383    async fn an_unknown_event_name_is_caught_on_each_backend() {
1384        let mut email = email_config("none");
1385        email.email.events = vec!["certificate_exploded".to_string()];
1386        let error = from_config(
1387            "le",
1388            &email,
1389            crate::testutil::outbound_with(test_resolver()),
1390            &test_queue().await,
1391        )
1392        .unwrap_err()
1393        .to_string();
1394        assert!(error.contains("notify.email.events"), "{error}");
1395
1396        let webhook = NotifyConfig {
1397            enabled: vec!["webhook".to_string()],
1398            webhook_enabled: vec!["chat".to_string()],
1399            webhook: BTreeMap::from([(
1400                "chat".to_string(),
1401                crate::config::WebhookNotifyConfig {
1402                    events: vec!["certificate_exploded".to_string()],
1403                    ..webhook_entry()
1404                },
1405            )]),
1406            ..NotifyConfig::default()
1407        };
1408        let error = from_config(
1409            "le",
1410            &webhook,
1411            crate::testutil::outbound_with(test_resolver()),
1412            &test_queue().await,
1413        )
1414        .unwrap_err()
1415        .to_string();
1416        assert!(error.contains("notify.webhook.chat.events"), "{error}");
1417    }
1418
1419    /// `notify.custom_enabled` names entries in `notify.custom`; a name with no
1420    /// entry behind it is a startup error rather than a silently missing hook.
1421    #[tokio::test]
1422    async fn a_custom_name_with_no_entry_is_a_startup_error() {
1423        let cfg = NotifyConfig {
1424            enabled: vec!["custom".to_string()],
1425            custom_enabled: vec!["webhook".to_string()],
1426            ..NotifyConfig::default()
1427        };
1428        let error = from_config(
1429            "le",
1430            &cfg,
1431            crate::testutil::outbound_with(test_resolver()),
1432            &test_queue().await,
1433        )
1434        .unwrap_err()
1435        .to_string();
1436        assert!(
1437            error.contains("notify.custom_enabled names `webhook`"),
1438            "{error}"
1439        );
1440    }
1441
1442    /// A `notify.custom` key that is not a valid environment-variable segment
1443    /// could name a different entry through `ACME_PROXY_…` than in the file.
1444    #[tokio::test]
1445    async fn an_invalid_custom_key_name_is_a_startup_error() {
1446        let mut custom = std::collections::BTreeMap::new();
1447        custom.insert("Web Hook".to_string(), CustomNotifyConfig::default());
1448        let cfg = NotifyConfig {
1449            enabled: vec!["custom".to_string()],
1450            custom_enabled: vec!["Web Hook".to_string()],
1451            custom,
1452            ..NotifyConfig::default()
1453        };
1454        let error = from_config(
1455            "le",
1456            &cfg,
1457            crate::testutil::outbound_with(test_resolver()),
1458            &test_queue().await,
1459        )
1460        .unwrap_err()
1461        .to_string();
1462        assert!(error.contains("invalid name"), "{error}");
1463    }
1464
1465    #[tokio::test]
1466    async fn unknown_backend_name_is_a_startup_error() {
1467        let cfg = NotifyConfig {
1468            enabled: vec!["carrier-pigeon".to_string()],
1469            ..NotifyConfig::default()
1470        };
1471        let error = from_config(
1472            "le",
1473            &cfg,
1474            crate::testutil::outbound_with(test_resolver()),
1475            &test_queue().await,
1476        )
1477        .unwrap_err()
1478        .to_string();
1479        assert!(error.contains("unknown notify backend"), "{error}");
1480    }
1481
1482    #[tokio::test]
1483    async fn custom_enabled_empty_is_a_startup_error() {
1484        let cfg = NotifyConfig {
1485            enabled: vec!["custom".to_string()],
1486            ..NotifyConfig::default()
1487        };
1488        let error = from_config(
1489            "le",
1490            &cfg,
1491            crate::testutil::outbound_with(test_resolver()),
1492            &test_queue().await,
1493        )
1494        .unwrap_err()
1495        .to_string();
1496        assert!(error.contains("notify.custom_enabled is empty"), "{error}");
1497    }
1498
1499    #[tokio::test]
1500    async fn an_unknown_event_name_is_a_startup_error() {
1501        let cfg = NotifyConfig {
1502            enabled: vec!["email".to_string()],
1503            email: crate::config::EmailNotifyConfig {
1504                smtp_host: "localhost".to_string(),
1505                events: vec!["orders_shipped".to_string()],
1506                ..crate::config::EmailNotifyConfig::default()
1507            },
1508            ..NotifyConfig::default()
1509        };
1510        let error = from_config(
1511            "le",
1512            &cfg,
1513            crate::testutil::outbound_with(test_resolver()),
1514            &test_queue().await,
1515        )
1516        .unwrap_err()
1517        .to_string();
1518        assert!(error.contains("unknown event"), "{error}");
1519    }
1520
1521    /// The `events` list decides membership of the *queue*, not of the delivery:
1522    /// a backend outside the list never gets a row, so `wants` is asserted on
1523    /// both sides rather than on what arrived at a backend afterwards.
1524    #[tokio::test]
1525    async fn a_backend_only_accepts_events_it_is_configured_for() {
1526        let wide = BackendSlot::new(
1527            "wide",
1528            Arc::new(RecordingNotifyBackend::default()),
1529            &every_kind(),
1530        );
1531        let narrow = BackendSlot::new(
1532            "narrow",
1533            Arc::new(RecordingNotifyBackend::default()),
1534            &["certificate_issued".to_string()],
1535        );
1536
1537        assert!(wide.wants(&profile_mounted("default")));
1538        assert!(!narrow.wants(&profile_mounted("default")));
1539
1540        let queue = test_queue().await;
1541        let dispatcher = NotifyDispatcher::new("le", vec![wide, narrow], queue.clone());
1542        dispatcher.dispatch(profile_mounted("default")).await;
1543
1544        assert_eq!(
1545            Job::count_live(NOTIFY_JOB_KIND, queue.database())
1546                .await
1547                .unwrap(),
1548            1,
1549            "only the wide backend is queued for"
1550        );
1551    }
1552
1553    /// One backend's failure is another's business, and the queue is what keeps
1554    /// them apart now: two rows, settled independently, so the healthy one is
1555    /// `done` while the failing one is still being retried.
1556    #[tokio::test]
1557    async fn a_failing_backend_does_not_stop_another_from_receiving_the_event() {
1558        let failing = Arc::new(RecordingNotifyBackend::failing());
1559        let healthy = Arc::new(RecordingNotifyBackend::default());
1560        let queue = test_queue().await;
1561        let dispatcher = NotifyDispatcher::new(
1562            "le",
1563            vec![
1564                BackendSlot::new("failing", failing.clone(), &every_kind()),
1565                BackendSlot::new("healthy", healthy.clone(), &every_kind()),
1566            ],
1567            queue,
1568        );
1569
1570        assert!(matches!(
1571            dispatcher
1572                .deliver("failing", &profile_mounted("default"))
1573                .await,
1574            Some(Err(_))
1575        ));
1576        assert!(matches!(
1577            dispatcher
1578                .deliver("healthy", &profile_mounted("default"))
1579                .await,
1580            Some(Ok(()))
1581        ));
1582
1583        assert_eq!(failing.events.lock().unwrap().len(), 1);
1584        assert_eq!(healthy.events.lock().unwrap().len(), 1);
1585    }
1586
1587    #[tokio::test]
1588    async fn build_registry_builds_one_dispatcher_per_profile() {
1589        let profiles = vec![
1590            ProfileConfig {
1591                name: "a".to_string(),
1592                sections: crate::config::ProfileSections::default(),
1593            },
1594            ProfileConfig {
1595                name: "b".to_string(),
1596                sections: crate::config::ProfileSections::default(),
1597            },
1598        ];
1599        let registry = build_registry(
1600            &profiles,
1601            crate::testutil::outbound_with(test_resolver()),
1602            &test_queue().await,
1603        )
1604        .unwrap();
1605        assert_eq!(registry.len(), 2);
1606        assert_eq!(registry["a"].profile(), "a");
1607        assert_eq!(registry["b"].profile(), "b");
1608    }
1609
1610    /// Two `custom` entries are two backends, and a job row has to be able to
1611    /// say which one it means. `NotifyBackend::name` answers `"custom"` for
1612    /// both, so the slot id is built from the configuration key instead.
1613    #[tokio::test]
1614    async fn two_custom_entries_get_distinct_slot_ids() {
1615        let dir = crate::testutil::TempDir::new("notify-slot");
1616        let script = crate::testutil::write_script(&dir, "notify.sh", "#!/bin/sh\nexit 0\n");
1617        let entry = || CustomNotifyConfig {
1618            script_path: script.display().to_string(),
1619            ..CustomNotifyConfig::default()
1620        };
1621        let mut custom = std::collections::BTreeMap::new();
1622        custom.insert("webhook".to_string(), entry());
1623        custom.insert("pager".to_string(), entry());
1624
1625        let cfg = NotifyConfig {
1626            enabled: vec!["custom".to_string()],
1627            custom_enabled: vec!["webhook".to_string(), "pager".to_string()],
1628            custom,
1629            ..NotifyConfig::default()
1630        };
1631
1632        let dispatcher = from_config(
1633            "le",
1634            &cfg,
1635            crate::testutil::outbound_with(test_resolver()),
1636            &test_queue().await,
1637        )
1638        .expect("both custom entries must build");
1639
1640        assert!(dispatcher.slot("custom:webhook").is_some());
1641        assert!(dispatcher.slot("custom:pager").is_some());
1642        assert!(dispatcher.slot("custom").is_none());
1643    }
1644
1645    /// The durable payload has to survive a restart, so every variant must come
1646    /// back out of JSON as the variant that went in. A wide `match` over an enum
1647    /// that grows is exactly where a new variant gets forgotten.
1648    #[test]
1649    fn every_event_round_trips_through_its_payload() {
1650        for event in every_event() {
1651            let encoded = event.payload();
1652            let decoded: NotifyEvent = serde_json::from_value(encoded.clone())
1653                .unwrap_or_else(|error| panic!("{} must decode: {error}", event.kind()));
1654            assert_eq!(decoded.kind(), event.kind());
1655            assert_eq!(decoded.profile(), event.profile());
1656            assert_eq!(decoded.payload(), encoded, "re-encoding must be stable");
1657        }
1658    }
1659
1660    /// The `custom` backend's stdin contract: the tag is a `"hook"` member
1661    /// carrying the event kind, sitting flat beside the event's own fields. That
1662    /// used to be spliced in by hand and is now serde's internal tag — a script
1663    /// in the field must not be able to tell the difference.
1664    #[test]
1665    fn payload_is_tagged_with_its_own_hook() {
1666        let event = NotifyEvent::CertificateIssued(CertificateIssuedData {
1667            profile: "le".to_string(),
1668            order_id: "ord-1".to_string(),
1669            account_id: "acct-1".to_string(),
1670            cert_serial: "0a0b".to_string(),
1671            identifiers: vec!["a.example.com".to_string()],
1672            client_ip: Some("203.0.113.5".to_string()),
1673        });
1674
1675        assert_eq!(
1676            event.payload(),
1677            serde_json::json!({
1678                "hook": "certificate_issued",
1679                "profile": "le",
1680                "order_id": "ord-1",
1681                "account_id": "acct-1",
1682                "cert_serial": "0a0b",
1683                "identifiers": ["a.example.com"],
1684                "client_ip": "203.0.113.5",
1685            })
1686        );
1687    }
1688
1689    #[test]
1690    fn template_dir_override_wins_over_the_embedded_default() {
1691        let dir = crate::testutil::TempDir::new("notify");
1692        std::fs::create_dir_all(dir.join("email")).unwrap();
1693        std::fs::write(
1694            dir.join("email/profile_mounted.subject.j2"),
1695            "override: {{ profile }}",
1696        )
1697        .unwrap();
1698
1699        let env = build_environment(dir.path().to_str().unwrap());
1700        let rendered = render(
1701            &env,
1702            "email/profile_mounted.subject.j2",
1703            &profile_mounted("default"),
1704        )
1705        .unwrap();
1706        assert_eq!(rendered, "override: default");
1707
1708        // A template not present in `template_dir` still falls back to the
1709        // compiled-in default rather than failing outright.
1710        let rendered = render(
1711            &env,
1712            "email/profile_mounted.body.j2",
1713            &profile_mounted("default"),
1714        )
1715        .unwrap();
1716        assert!(rendered.contains("default"));
1717    }
1718
1719    /// A template that is not there will not be there next time either, so the
1720    /// failure must not spend a retry budget getting to the same answer.
1721    #[test]
1722    fn a_template_failure_is_permanent() {
1723        let env = build_environment("");
1724
1725        let missing = render(&env, "email/no_such_event.body.j2", &profile_mounted("le"))
1726            .expect_err("there is no such template");
1727        assert!(!missing.retryable(), "{missing}");
1728
1729        let dir = crate::testutil::TempDir::new("notify-broken");
1730        std::fs::create_dir_all(dir.join("email")).unwrap();
1731        std::fs::write(dir.join("email/profile_mounted.body.j2"), "{{ unclosed").unwrap();
1732        let env = build_environment(dir.path().to_str().unwrap());
1733        let broken = render(
1734            &env,
1735            "email/profile_mounted.body.j2",
1736            &profile_mounted("le"),
1737        )
1738        .expect_err("the template does not compile");
1739        assert!(!broken.retryable(), "{broken}");
1740    }
1741
1742    #[test]
1743    fn embedded_defaults_render_with_no_template_dir() {
1744        let env = build_environment("");
1745        let event = NotifyEvent::CertificateIssued(CertificateIssuedData {
1746            profile: "le".to_string(),
1747            order_id: "ord-1".to_string(),
1748            account_id: "acc-1".to_string(),
1749            cert_serial: "AA:BB".to_string(),
1750            identifiers: vec!["example.com".to_string()],
1751            client_ip: Some("203.0.113.1".to_string()),
1752        });
1753
1754        let subject = render(&env, "email/certificate_issued.subject.j2", &event).unwrap();
1755        assert!(subject.contains("le"), "{subject}");
1756
1757        let body = render(&env, "email/certificate_issued.body.j2", &event).unwrap();
1758        assert!(body.contains("example.com"), "{body}");
1759        assert!(body.contains("203.0.113.1"), "{body}");
1760    }
1761}