use std::collections::{HashMap, HashSet};
use std::sync::{Arc, LazyLock};
use async_trait::async_trait;
use tracing::info;
use crate::config::{ALL_NOTIFY_EVENTS, NotifyConfig, ProfileConfig};
use crate::jobs::{JobQueue, JobSpec};
pub mod custom;
pub mod email;
pub mod expiry;
pub mod job;
pub mod webhook;
pub use job::{NOTIFY_JOB_KIND, NotifyJob};
#[async_trait]
pub trait NotifyBackend: Send + Sync {
fn name(&self) -> &'static str;
async fn send(&self, event: &NotifyEvent) -> Result<(), NotifyError>;
}
#[derive(Debug, thiserror::Error)]
#[error("{detail}")]
pub struct NotifyError {
detail: String,
retryable: bool,
}
impl NotifyError {
pub fn new(detail: impl Into<String>) -> Self {
Self {
detail: detail.into(),
retryable: true,
}
}
pub fn permanent(detail: impl Into<String>) -> Self {
Self {
detail: detail.into(),
retryable: false,
}
}
#[must_use]
pub fn retryable(&self) -> bool {
self.retryable
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct ProfileMountedData {
pub profile: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AccountCreatedData {
pub profile: String,
pub account_id: String,
pub contact: Vec<String>,
pub client_ip: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AccountDeactivatedData {
pub profile: String,
pub account_id: String,
pub client_ip: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct CertificateIssuedData {
pub profile: String,
pub order_id: String,
pub account_id: String,
pub cert_serial: String,
pub identifiers: Vec<String>,
pub client_ip: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct CertificateRevokedData {
pub profile: String,
pub order_id: String,
pub account_id: String,
pub cert_serial: String,
pub reason: Option<u32>,
pub client_ip: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct ChallengeFailedData {
pub profile: String,
pub order_id: String,
pub account_id: String,
pub authz_id: String,
pub challenge_id: String,
pub challenge_type: String,
pub identifier: String,
pub error: String,
pub client_ip: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct ExpiringCertificate {
pub order_id: String,
pub account_id: String,
pub cert_serial: String,
pub identifiers: Vec<String>,
pub not_after: i64,
pub days_remaining: i64,
pub superseded_by: Option<SupersededBy>,
}
pub use crate::admin::SupersededBy;
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct CertificatesExpiringData {
pub profile: String,
pub generated_at: i64,
pub lead_days: u64,
pub total: i64,
pub certificates: Vec<ExpiringCertificate>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "hook", rename_all = "snake_case")]
pub enum NotifyEvent {
ProfileMounted(ProfileMountedData),
AccountCreated(AccountCreatedData),
AccountDeactivated(AccountDeactivatedData),
CertificateIssued(CertificateIssuedData),
CertificateRevoked(CertificateRevokedData),
ChallengeFailed(ChallengeFailedData),
CertificatesExpiring(CertificatesExpiringData),
}
impl NotifyEvent {
pub fn kind(&self) -> &'static str {
match self {
Self::ProfileMounted(_) => "profile_mounted",
Self::AccountCreated(_) => "account_created",
Self::AccountDeactivated(_) => "account_deactivated",
Self::CertificateIssued(_) => "certificate_issued",
Self::CertificateRevoked(_) => "certificate_revoked",
Self::ChallengeFailed(_) => "challenge_failed",
Self::CertificatesExpiring(_) => "certificates_expiring",
}
}
pub fn profile(&self) -> &str {
match self {
Self::ProfileMounted(data) => &data.profile,
Self::AccountCreated(data) => &data.profile,
Self::AccountDeactivated(data) => &data.profile,
Self::CertificateIssued(data) => &data.profile,
Self::CertificateRevoked(data) => &data.profile,
Self::ChallengeFailed(data) => &data.profile,
Self::CertificatesExpiring(data) => &data.profile,
}
}
pub(crate) fn context(&self) -> minijinja::Value {
match self {
Self::ProfileMounted(data) => minijinja::Value::from_serialize(data),
Self::AccountCreated(data) => minijinja::Value::from_serialize(data),
Self::AccountDeactivated(data) => minijinja::Value::from_serialize(data),
Self::CertificateIssued(data) => minijinja::Value::from_serialize(data),
Self::CertificateRevoked(data) => minijinja::Value::from_serialize(data),
Self::ChallengeFailed(data) => minijinja::Value::from_serialize(data),
Self::CertificatesExpiring(data) => minijinja::Value::from_serialize(data),
}
}
pub(crate) fn payload(&self) -> serde_json::Value {
serde_json::to_value(self).expect("notify event data always serializes to a JSON object")
}
fn client_ip(&self) -> Option<&str> {
match self {
Self::ProfileMounted(_) => None,
Self::AccountCreated(data) => data.client_ip.as_deref(),
Self::AccountDeactivated(data) => data.client_ip.as_deref(),
Self::CertificateIssued(data) => data.client_ip.as_deref(),
Self::CertificateRevoked(data) => data.client_ip.as_deref(),
Self::ChallengeFailed(data) => data.client_ip.as_deref(),
Self::CertificatesExpiring(_) => None,
}
}
fn account_id(&self) -> Option<&str> {
match self {
Self::ProfileMounted(_) => None,
Self::AccountCreated(data) => Some(&data.account_id),
Self::AccountDeactivated(data) => Some(&data.account_id),
Self::CertificateIssued(data) => Some(&data.account_id),
Self::CertificateRevoked(data) => Some(&data.account_id),
Self::ChallengeFailed(data) => Some(&data.account_id),
Self::CertificatesExpiring(_) => None,
}
}
fn order_id(&self) -> Option<&str> {
match self {
Self::ProfileMounted(_) => None,
Self::AccountCreated(_) => None,
Self::AccountDeactivated(_) => None,
Self::CertificateIssued(data) => Some(&data.order_id),
Self::CertificateRevoked(data) => Some(&data.order_id),
Self::ChallengeFailed(data) => Some(&data.order_id),
Self::CertificatesExpiring(_) => None,
}
}
fn cert_serial(&self) -> Option<&str> {
match self {
Self::ProfileMounted(_) => None,
Self::AccountCreated(_) => None,
Self::AccountDeactivated(_) => None,
Self::CertificateIssued(data) => Some(&data.cert_serial),
Self::CertificateRevoked(data) => Some(&data.cert_serial),
Self::ChallengeFailed(_) => None,
Self::CertificatesExpiring(_) => None,
}
}
fn identifiers_joined(&self) -> String {
match self {
Self::ProfileMounted(_) => String::new(),
Self::AccountCreated(_) => String::new(),
Self::AccountDeactivated(_) => String::new(),
Self::CertificateIssued(data) => data.identifiers.join(","),
Self::CertificateRevoked(_) => String::new(),
Self::ChallengeFailed(_) => String::new(),
Self::CertificatesExpiring(_) => String::new(),
}
}
}
pub struct BackendSlot {
id: String,
events: HashSet<String>,
backend: Arc<dyn NotifyBackend>,
}
impl BackendSlot {
#[must_use]
pub fn id(&self) -> &str {
&self.id
}
#[must_use]
pub fn wants(&self, event: &NotifyEvent) -> bool {
self.events.contains(event.kind())
}
#[must_use]
pub fn new(id: impl Into<String>, backend: Arc<dyn NotifyBackend>, events: &[String]) -> Self {
Self {
id: id.into(),
events: events.iter().cloned().collect(),
backend,
}
}
}
pub struct NotifyDispatcher {
profile: String,
slots: Vec<BackendSlot>,
jobs: JobQueue,
}
impl std::fmt::Debug for NotifyDispatcher {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("NotifyDispatcher")
.field("profile", &self.profile)
.field(
"backends",
&self.slots.iter().map(BackendSlot::id).collect::<Vec<_>>(),
)
.finish()
}
}
impl NotifyDispatcher {
#[must_use]
pub fn new(profile: impl Into<String>, slots: Vec<BackendSlot>, jobs: JobQueue) -> Self {
Self {
profile: profile.into(),
slots,
jobs,
}
}
#[must_use]
pub fn disabled(jobs: JobQueue) -> Self {
Self::new("default", Vec::new(), jobs)
}
#[must_use]
pub fn profile(&self) -> &str {
&self.profile
}
#[must_use]
pub fn slot(&self, id: &str) -> Option<&BackendSlot> {
self.slots.iter().find(|slot| slot.id == id)
}
pub async fn dispatch(&self, event: NotifyEvent) {
if self.slots.is_empty() {
return;
}
let kind = event.kind();
let payload = event.payload();
let delivery_id = uuid::Uuid::now_v7().to_string();
for slot in &self.slots {
if !slot.wants(&event) {
continue;
}
let spec = JobSpec::now(NOTIFY_JOB_KIND, format!("{delivery_id}:{}", slot.id))
.with_payload(serde_json::json!({
"profile": self.profile,
"backend": slot.id,
"event": payload,
}));
if self.jobs.enqueue_or_log(spec).await {
info!(
event = "notify_delivery_queued",
outcome = "progress",
profile = %self.profile,
backend = %slot.id,
kind,
delivery_id = %delivery_id,
);
}
}
}
pub(crate) async fn deliver(
&self,
id: &str,
event: &NotifyEvent,
) -> Option<Result<(), NotifyError>> {
let slot = self.slot(id)?;
Some(slot.backend.send(event).await)
}
}
fn validate_events(field: &str, events: &[String]) -> anyhow::Result<()> {
for event in events {
anyhow::ensure!(
ALL_NOTIFY_EVENTS.contains(&event.as_str()),
"{field}: unknown event `{event}` (expected one of {ALL_NOTIFY_EVENTS:?})"
);
}
Ok(())
}
pub fn from_config(
profile: &str,
cfg: &NotifyConfig,
outbound: crate::http_client::Outbound,
jobs: &JobQueue,
) -> anyhow::Result<Arc<NotifyDispatcher>> {
let env = Arc::new(build_environment(&cfg.template_dir));
let mut slots: Vec<BackendSlot> = Vec::with_capacity(cfg.enabled.len());
for name in &cfg.enabled {
let built: Vec<BackendSlot> = match name.as_str() {
"email" => {
validate_events("notify.email.events", &cfg.email.events)?;
vec![BackendSlot::new(
"email",
Arc::new(email::EmailNotifier::from_config(&cfg.email, env.clone())?),
&cfg.email.events,
)]
}
"webhook" => build_webhook_slots(cfg, &env, outbound.clone())?,
"custom" => build_custom_slots(cfg)?,
"mattermost" => anyhow::bail!(
"notify.enabled: `mattermost` was replaced by `webhook`. Use \
notify.enabled = [\"webhook\"] with a [notify.webhook.<name>] entry \
whose `url` is the incoming webhook URL; the default `body` is \
already the payload Mattermost accepts. `channel` and `username` \
move into that `body`"
),
other => anyhow::bail!("unknown notify backend: {other}"),
};
slots.extend(built);
}
if slots.is_empty() {
info!(
event = "notify_disabled",
outcome = "success",
"no notification backends configured"
);
} else {
info!(event = "notify_enabled", outcome = "success", backends = ?cfg.enabled);
}
Ok(Arc::new(NotifyDispatcher::new(
profile,
slots,
jobs.clone(),
)))
}
fn build_custom_slots(cfg: &NotifyConfig) -> anyhow::Result<Vec<BackendSlot>> {
crate::config::resolve_named_entries(
"notify.custom",
"notify.custom_enabled",
"custom",
&cfg.custom,
&cfg.custom_enabled,
)?
.into_iter()
.map(|(name, script)| -> anyhow::Result<BackendSlot> {
validate_events(&format!("notify.custom.{name}.events"), &script.events)?;
let backend = custom::CustomScriptNotifier::from_config(script)?;
Ok(BackendSlot::new(
format!("custom:{name}"),
Arc::new(backend),
&script.events,
))
})
.collect()
}
fn build_webhook_slots(
cfg: &NotifyConfig,
env: &minijinja::Environment<'static>,
outbound: crate::http_client::Outbound,
) -> anyhow::Result<Vec<BackendSlot>> {
crate::config::resolve_named_entries(
"notify.webhook",
"notify.webhook_enabled",
"webhook",
&cfg.webhook,
&cfg.webhook_enabled,
)?
.into_iter()
.map(|(name, entry)| -> anyhow::Result<BackendSlot> {
validate_events(&format!("notify.webhook.{name}.events"), &entry.events)?;
let backend = webhook::WebhookNotifier::from_config(name, entry, env, outbound.clone())?;
Ok(BackendSlot::new(
format!("webhook:{name}"),
Arc::new(backend),
&entry.events,
))
})
.collect()
}
pub type DispatcherMap = HashMap<String, Arc<NotifyDispatcher>>;
pub type NotifiersSender = tokio::sync::watch::Sender<Arc<DispatcherMap>>;
#[derive(Clone)]
pub struct Notifiers(tokio::sync::watch::Receiver<Arc<DispatcherMap>>);
impl Notifiers {
#[must_use]
pub fn get(&self, profile: &str) -> Option<Arc<NotifyDispatcher>> {
self.0.borrow().get(profile).cloned()
}
}
impl From<Arc<DispatcherMap>> for Notifiers {
fn from(map: Arc<DispatcherMap>) -> Self {
Self(tokio::sync::watch::channel(map).1)
}
}
impl From<DispatcherMap> for Notifiers {
fn from(map: DispatcherMap) -> Self {
Arc::new(map).into()
}
}
#[must_use]
pub fn notifiers_channel(initial: DispatcherMap) -> (NotifiersSender, Notifiers) {
let (sender, receiver) = tokio::sync::watch::channel(Arc::new(initial));
(sender, Notifiers(receiver))
}
pub fn build_registry(
profiles: &[ProfileConfig],
outbound: crate::http_client::Outbound,
jobs: &JobQueue,
) -> anyhow::Result<HashMap<String, Arc<NotifyDispatcher>>> {
let mut registry = HashMap::with_capacity(profiles.len());
for profile in profiles {
let dispatcher = from_config(
&profile.name,
&profile.sections.notify,
outbound.clone(),
jobs,
)
.map_err(|error| anyhow::anyhow!("profile `{}`: {error}", profile.name))?;
registry.insert(profile.name.clone(), dispatcher);
}
Ok(registry)
}
macro_rules! embed {
($name:literal) => {
($name, include_str!(concat!("templates/", $name)))
};
}
static EMBEDDED_TEMPLATES: LazyLock<HashMap<&'static str, &'static str>> = LazyLock::new(|| {
HashMap::from([
embed!("email/profile_mounted.subject.j2"),
embed!("email/profile_mounted.body.j2"),
embed!("email/account_created.subject.j2"),
embed!("email/account_created.body.j2"),
embed!("email/account_deactivated.subject.j2"),
embed!("email/account_deactivated.body.j2"),
embed!("email/certificate_issued.subject.j2"),
embed!("email/certificate_issued.body.j2"),
embed!("email/certificate_revoked.subject.j2"),
embed!("email/certificate_revoked.body.j2"),
embed!("email/challenge_failed.subject.j2"),
embed!("email/challenge_failed.body.j2"),
embed!("email/certificates_expiring.subject.j2"),
embed!("email/certificates_expiring.body.j2"),
embed!("webhook/profile_mounted.j2"),
embed!("webhook/account_created.j2"),
embed!("webhook/account_deactivated.j2"),
embed!("webhook/certificate_issued.j2"),
embed!("webhook/certificate_revoked.j2"),
embed!("webhook/challenge_failed.j2"),
embed!("webhook/certificates_expiring.j2"),
])
});
pub(crate) fn build_environment(template_dir: &str) -> minijinja::Environment<'static> {
crate::templating::loader_env(template_dir, &EMBEDDED_TEMPLATES)
}
#[must_use]
pub fn template_names() -> Vec<&'static str> {
let mut names: Vec<&'static str> = EMBEDDED_TEMPLATES.keys().copied().collect();
names.sort_unstable();
names
}
pub(crate) fn render(
env: &minijinja::Environment<'static>,
template_name: &str,
event: &NotifyEvent,
) -> Result<String, NotifyError> {
let template = env.get_template(template_name).map_err(|error| {
NotifyError::permanent(format!("template `{template_name}` not found: {error}"))
})?;
template.render(event.context()).map_err(|error| {
NotifyError::permanent(format!("template `{template_name}` failed: {error}"))
})
}
#[cfg(test)]
pub(crate) mod tests {
use super::*;
use crate::config::CustomNotifyConfig;
use crate::sqlite::db::Database;
use crate::sqlite::job::Job;
use std::collections::BTreeMap;
use std::sync::Mutex;
fn test_resolver() -> Arc<dyn crate::dns::Resolver> {
Arc::new(crate::dns::HickoryResolver::from_system_uncached().unwrap())
}
pub(crate) async fn test_queue() -> JobQueue {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
JobQueue::new(database, &crate::config::JobsConfig::default())
}
#[derive(Default, Clone, Copy)]
enum Failure {
#[default]
None,
Retryable,
Permanent,
}
#[derive(Default)]
pub(crate) struct RecordingNotifyBackend {
pub(crate) events: Mutex<Vec<NotifyEvent>>,
fail: Failure,
}
impl RecordingNotifyBackend {
pub(crate) fn failing() -> Self {
Self {
events: Mutex::new(Vec::new()),
fail: Failure::Retryable,
}
}
pub(crate) fn failing_permanently() -> Self {
Self {
events: Mutex::new(Vec::new()),
fail: Failure::Permanent,
}
}
}
#[async_trait]
impl NotifyBackend for RecordingNotifyBackend {
fn name(&self) -> &'static str {
"recording"
}
async fn send(&self, event: &NotifyEvent) -> Result<(), NotifyError> {
self.events.lock().unwrap().push(event.clone());
match self.fail {
Failure::None => Ok(()),
Failure::Retryable => Err(NotifyError::new("recording backend configured to fail")),
Failure::Permanent => Err(NotifyError::permanent(
"recording backend configured to fail permanently",
)),
}
}
}
fn profile_mounted(profile: &str) -> NotifyEvent {
NotifyEvent::ProfileMounted(ProfileMountedData {
profile: profile.to_string(),
})
}
fn every_event() -> Vec<NotifyEvent> {
vec![
profile_mounted("p"),
NotifyEvent::AccountCreated(AccountCreatedData {
profile: "p".to_string(),
account_id: "acct-1".to_string(),
contact: vec!["mailto:a@example.com".to_string()],
client_ip: Some("203.0.113.5".to_string()),
}),
NotifyEvent::AccountDeactivated(AccountDeactivatedData {
profile: "p".to_string(),
account_id: "acct-1".to_string(),
client_ip: Some("203.0.113.5".to_string()),
}),
NotifyEvent::CertificateIssued(CertificateIssuedData {
profile: "p".to_string(),
order_id: "ord-1".to_string(),
account_id: "acct-1".to_string(),
cert_serial: "0a0b".to_string(),
identifiers: vec!["a.example.com".to_string(), "b.example.com".to_string()],
client_ip: Some("203.0.113.5".to_string()),
}),
NotifyEvent::CertificateRevoked(CertificateRevokedData {
profile: "p".to_string(),
order_id: "ord-1".to_string(),
account_id: "acct-1".to_string(),
cert_serial: "0a0b".to_string(),
reason: Some(1),
client_ip: None,
}),
NotifyEvent::ChallengeFailed(ChallengeFailedData {
profile: "p".to_string(),
order_id: "ord-1".to_string(),
account_id: "acct-1".to_string(),
authz_id: "authz-1".to_string(),
challenge_id: "chall-1".to_string(),
challenge_type: "http-01".to_string(),
identifier: "a.example.com".to_string(),
error: "connection refused".to_string(),
client_ip: Some("203.0.113.5".to_string()),
}),
NotifyEvent::CertificatesExpiring(CertificatesExpiringData {
profile: "p".to_string(),
generated_at: 1_700_000_000,
lead_days: 14,
total: 3,
certificates: vec![ExpiringCertificate {
order_id: "ord-1".to_string(),
account_id: "acct-1".to_string(),
cert_serial: "0a0b".to_string(),
identifiers: vec!["a.example.com".to_string()],
not_after: 1_700_600_000,
days_remaining: 6,
superseded_by: None,
}],
}),
]
}
#[test]
fn every_event_answers_every_accessor() {
let events = every_event();
assert_eq!(
events.len(),
ALL_NOTIFY_EVENTS.len(),
"every declared event kind needs a sample here"
);
for event in &events {
assert_eq!(event.profile(), "p");
assert!(
ALL_NOTIFY_EVENTS.contains(&event.kind()),
"{}",
event.kind()
);
assert!(!event.context().is_undefined());
let payload = event.payload();
assert_eq!(
payload.get("hook").and_then(|v| v.as_str()),
Some(event.kind()),
"the payload must name its own hook"
);
assert_eq!(payload.get("profile").and_then(|v| v.as_str()), Some("p"));
}
let subjectless = [&events[0], events.last().unwrap()];
for event in subjectless {
assert_eq!(event.client_ip(), None, "{}", event.kind());
assert_eq!(event.account_id(), None, "{}", event.kind());
assert_eq!(event.order_id(), None, "{}", event.kind());
assert_eq!(event.cert_serial(), None, "{}", event.kind());
assert_eq!(event.identifiers_joined(), "", "{}", event.kind());
}
let per_subject = &events[1..events.len() - 1];
for event in per_subject {
assert_eq!(event.account_id(), Some("acct-1"));
}
assert_eq!(events[4].client_ip(), None);
assert_eq!(events[1].client_ip(), Some("203.0.113.5"));
assert_eq!(events[1].order_id(), None);
assert_eq!(events[2].cert_serial(), None);
for event in &per_subject[2..] {
assert_eq!(event.order_id(), Some("ord-1"));
}
assert_eq!(events[3].cert_serial(), Some("0a0b"));
assert_eq!(events[4].cert_serial(), Some("0a0b"));
assert_eq!(events[5].cert_serial(), None);
assert_eq!(
events[3].identifiers_joined(),
"a.example.com,b.example.com"
);
assert_eq!(events[5].identifiers_joined(), "");
let digest = events.last().unwrap().payload();
assert_eq!(digest["total"], 3);
assert_eq!(digest["certificates"][0]["order_id"], "ord-1");
assert_eq!(digest["certificates"][0]["days_remaining"], 6);
}
#[tokio::test]
async fn the_dispatcher_debug_names_its_backends() {
let queue = test_queue().await;
let dispatcher = NotifyDispatcher::new(
"le",
vec![BackendSlot::new(
"recording",
Arc::new(RecordingNotifyBackend::default()),
&every_kind(),
)],
queue.clone(),
);
let rendered = format!("{dispatcher:?}");
assert!(rendered.contains("NotifyDispatcher"), "{rendered}");
assert!(rendered.contains("recording"), "{rendered}");
assert!(rendered.contains("le"), "{rendered}");
assert!(format!("{:?}", NotifyDispatcher::disabled(queue)).contains("[]"));
}
#[tokio::test]
async fn a_handle_taken_before_a_swap_reads_the_map_after_it() {
let queue = test_queue().await;
let (sender, notifiers) = notifiers_channel(HashMap::new());
assert!(notifiers.get("le").is_none());
let mut next = HashMap::new();
next.insert(
"le".to_string(),
Arc::new(NotifyDispatcher::disabled(queue.clone())),
);
sender.send_replace(Arc::new(next));
assert!(notifiers.get("le").is_some());
assert!(notifiers.get("staging").is_none());
sender.send_replace(Arc::new(HashMap::new()));
assert!(notifiers.get("le").is_none());
}
#[tokio::test]
async fn a_fixed_map_survives_its_sender_being_dropped() {
let queue = test_queue().await;
let mut map = HashMap::new();
map.insert(
"le".to_string(),
Arc::new(NotifyDispatcher::disabled(queue)),
);
let notifiers: Notifiers = map.into();
assert!(notifiers.get("le").is_some());
assert!(notifiers.clone().get("le").is_some());
}
fn every_kind() -> Vec<String> {
ALL_NOTIFY_EVENTS.iter().map(|k| (*k).to_string()).collect()
}
async fn recording_dispatcher(
events: &[String],
) -> (Arc<NotifyDispatcher>, Arc<RecordingNotifyBackend>, JobQueue) {
let queue = test_queue().await;
let recorder = Arc::new(RecordingNotifyBackend::default());
let dispatcher = Arc::new(NotifyDispatcher::new(
"le",
vec![BackendSlot::new("recording", recorder.clone(), events)],
queue.clone(),
));
(dispatcher, recorder, queue)
}
#[tokio::test]
async fn dispatch_queues_a_row_rather_than_delivering() {
let (dispatcher, recorder, queue) = recording_dispatcher(&every_kind()).await;
dispatcher.dispatch(profile_mounted("le")).await;
assert!(
recorder.events.lock().unwrap().is_empty(),
"dispatch must not deliver inline"
);
let queued = Job::count_live(NOTIFY_JOB_KIND, queue.database())
.await
.unwrap();
assert_eq!(queued, 1, "one backend, one row");
}
#[tokio::test]
async fn a_backend_that_does_not_want_the_event_gets_no_row() {
let (dispatcher, _recorder, queue) =
recording_dispatcher(&["certificate_issued".to_string()]).await;
dispatcher.dispatch(profile_mounted("le")).await;
assert_eq!(
Job::count_live(NOTIFY_JOB_KIND, queue.database())
.await
.unwrap(),
0
);
}
#[tokio::test]
async fn one_dispatch_queues_one_row_per_wanting_backend() {
let queue = test_queue().await;
let dispatcher = NotifyDispatcher::new(
"le",
vec![
BackendSlot::new(
"email",
Arc::new(RecordingNotifyBackend::default()),
&every_kind(),
),
BackendSlot::new(
"custom:webhook",
Arc::new(RecordingNotifyBackend::default()),
&every_kind(),
),
BackendSlot::new(
"custom:pager",
Arc::new(RecordingNotifyBackend::default()),
&["certificate_revoked".to_string()],
),
],
queue.clone(),
);
dispatcher.dispatch(profile_mounted("le")).await;
assert_eq!(
Job::count_live(NOTIFY_JOB_KIND, queue.database())
.await
.unwrap(),
2,
"the third backend does not want this kind"
);
}
#[tokio::test]
async fn the_same_event_dispatched_twice_queues_twice() {
let (dispatcher, _recorder, queue) = recording_dispatcher(&every_kind()).await;
dispatcher.dispatch(profile_mounted("le")).await;
dispatcher.dispatch(profile_mounted("le")).await;
assert_eq!(
Job::count_live(NOTIFY_JOB_KIND, queue.database())
.await
.unwrap(),
2
);
}
#[tokio::test]
async fn a_disabled_dispatcher_queues_nothing() {
let queue = test_queue().await;
let dispatcher = NotifyDispatcher::disabled(queue.clone());
dispatcher.dispatch(profile_mounted("le")).await;
assert_eq!(
Job::count_live(NOTIFY_JOB_KIND, queue.database())
.await
.unwrap(),
0
);
}
#[tokio::test]
async fn a_database_failure_is_swallowed_by_dispatch() {
let (dispatcher, _recorder, queue) = recording_dispatcher(&every_kind()).await;
queue.database().pool.close().await;
dispatcher.dispatch(profile_mounted("le")).await;
}
#[tokio::test]
async fn deliver_reaches_one_backend_and_reports_an_unknown_id() {
let (dispatcher, recorder, _queue) = recording_dispatcher(&every_kind()).await;
let outcome = dispatcher
.deliver("recording", &profile_mounted("le"))
.await;
assert!(matches!(outcome, Some(Ok(()))));
assert_eq!(recorder.events.lock().unwrap().len(), 1);
assert!(
dispatcher
.deliver("carrier-pigeon", &profile_mounted("le"))
.await
.is_none()
);
assert_eq!(
recorder.events.lock().unwrap().len(),
1,
"an unknown id must reach no backend at all"
);
}
#[tokio::test]
async fn each_backend_name_builds_its_own_backend() {
let cfg = NotifyConfig {
enabled: vec!["email".to_string(), "webhook".to_string()],
email: crate::config::EmailNotifyConfig {
smtp_host: "smtp.example.com".to_string(),
from: "acme@example.com".to_string(),
to: vec!["ops@example.com".to_string()],
smtp_security: "none".to_string(),
smtp_username: "user".to_string(),
smtp_password: "pass".to_string(),
..crate::config::EmailNotifyConfig::default()
},
webhook_enabled: vec!["chat".to_string()],
webhook: BTreeMap::from([("chat".to_string(), webhook_entry())]),
..NotifyConfig::default()
};
let dispatcher = from_config(
"le",
&cfg,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.expect("both backends must build");
let rendered = format!("{dispatcher:?}");
assert!(rendered.contains("email"), "{rendered}");
assert!(rendered.contains("webhook:chat"), "{rendered}");
}
fn webhook_entry() -> crate::config::WebhookNotifyConfig {
crate::config::WebhookNotifyConfig {
url: "https://chat.example.com/hooks/abc".to_string(),
..crate::config::WebhookNotifyConfig::default()
}
}
#[tokio::test]
async fn two_webhook_entries_get_distinct_slot_ids() {
let cfg = NotifyConfig {
enabled: vec!["webhook".to_string()],
webhook_enabled: vec!["slack".to_string(), "teams".to_string()],
webhook: BTreeMap::from([
("slack".to_string(), webhook_entry()),
("teams".to_string(), webhook_entry()),
]),
..NotifyConfig::default()
};
let dispatcher = from_config(
"le",
&cfg,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.expect("both entries must build");
assert!(dispatcher.slot("webhook:slack").is_some());
assert!(dispatcher.slot("webhook:teams").is_some());
assert!(dispatcher.slot("webhook").is_none());
}
#[tokio::test]
async fn the_removed_mattermost_backend_is_refused_by_name() {
let cfg = NotifyConfig {
enabled: vec!["mattermost".to_string()],
..NotifyConfig::default()
};
let error = from_config(
"le",
&cfg,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.unwrap_err()
.to_string();
assert!(error.contains("mattermost"), "{error}");
assert!(error.contains("webhook"), "{error}");
}
#[tokio::test]
async fn every_smtp_security_mode_is_recognised() {
for mode in ["starttls", "tls", "none"] {
let cfg = email_config(mode);
from_config(
"le",
&cfg,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.unwrap_or_else(|error| panic!("`{mode}` must build: {error}"));
}
let error = from_config(
"le",
&email_config("carrier-pigeon"),
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.unwrap_err()
.to_string();
assert!(error.contains("smtp_security"), "{error}");
}
fn email_config(smtp_security: &str) -> NotifyConfig {
NotifyConfig {
enabled: vec!["email".to_string()],
email: crate::config::EmailNotifyConfig {
smtp_host: "smtp.example.com".to_string(),
from: "acme@example.com".to_string(),
to: vec!["ops@example.com".to_string()],
smtp_security: smtp_security.to_string(),
..crate::config::EmailNotifyConfig::default()
},
..NotifyConfig::default()
}
}
#[tokio::test]
async fn an_unknown_event_name_is_caught_on_each_backend() {
let mut email = email_config("none");
email.email.events = vec!["certificate_exploded".to_string()];
let error = from_config(
"le",
&email,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.unwrap_err()
.to_string();
assert!(error.contains("notify.email.events"), "{error}");
let webhook = NotifyConfig {
enabled: vec!["webhook".to_string()],
webhook_enabled: vec!["chat".to_string()],
webhook: BTreeMap::from([(
"chat".to_string(),
crate::config::WebhookNotifyConfig {
events: vec!["certificate_exploded".to_string()],
..webhook_entry()
},
)]),
..NotifyConfig::default()
};
let error = from_config(
"le",
&webhook,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.unwrap_err()
.to_string();
assert!(error.contains("notify.webhook.chat.events"), "{error}");
}
#[tokio::test]
async fn a_custom_name_with_no_entry_is_a_startup_error() {
let cfg = NotifyConfig {
enabled: vec!["custom".to_string()],
custom_enabled: vec!["webhook".to_string()],
..NotifyConfig::default()
};
let error = from_config(
"le",
&cfg,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.unwrap_err()
.to_string();
assert!(
error.contains("notify.custom_enabled names `webhook`"),
"{error}"
);
}
#[tokio::test]
async fn an_invalid_custom_key_name_is_a_startup_error() {
let mut custom = std::collections::BTreeMap::new();
custom.insert("Web Hook".to_string(), CustomNotifyConfig::default());
let cfg = NotifyConfig {
enabled: vec!["custom".to_string()],
custom_enabled: vec!["Web Hook".to_string()],
custom,
..NotifyConfig::default()
};
let error = from_config(
"le",
&cfg,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.unwrap_err()
.to_string();
assert!(error.contains("invalid name"), "{error}");
}
#[tokio::test]
async fn unknown_backend_name_is_a_startup_error() {
let cfg = NotifyConfig {
enabled: vec!["carrier-pigeon".to_string()],
..NotifyConfig::default()
};
let error = from_config(
"le",
&cfg,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.unwrap_err()
.to_string();
assert!(error.contains("unknown notify backend"), "{error}");
}
#[tokio::test]
async fn custom_enabled_empty_is_a_startup_error() {
let cfg = NotifyConfig {
enabled: vec!["custom".to_string()],
..NotifyConfig::default()
};
let error = from_config(
"le",
&cfg,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.unwrap_err()
.to_string();
assert!(error.contains("notify.custom_enabled is empty"), "{error}");
}
#[tokio::test]
async fn an_unknown_event_name_is_a_startup_error() {
let cfg = NotifyConfig {
enabled: vec!["email".to_string()],
email: crate::config::EmailNotifyConfig {
smtp_host: "localhost".to_string(),
events: vec!["orders_shipped".to_string()],
..crate::config::EmailNotifyConfig::default()
},
..NotifyConfig::default()
};
let error = from_config(
"le",
&cfg,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.unwrap_err()
.to_string();
assert!(error.contains("unknown event"), "{error}");
}
#[tokio::test]
async fn a_backend_only_accepts_events_it_is_configured_for() {
let wide = BackendSlot::new(
"wide",
Arc::new(RecordingNotifyBackend::default()),
&every_kind(),
);
let narrow = BackendSlot::new(
"narrow",
Arc::new(RecordingNotifyBackend::default()),
&["certificate_issued".to_string()],
);
assert!(wide.wants(&profile_mounted("default")));
assert!(!narrow.wants(&profile_mounted("default")));
let queue = test_queue().await;
let dispatcher = NotifyDispatcher::new("le", vec![wide, narrow], queue.clone());
dispatcher.dispatch(profile_mounted("default")).await;
assert_eq!(
Job::count_live(NOTIFY_JOB_KIND, queue.database())
.await
.unwrap(),
1,
"only the wide backend is queued for"
);
}
#[tokio::test]
async fn a_failing_backend_does_not_stop_another_from_receiving_the_event() {
let failing = Arc::new(RecordingNotifyBackend::failing());
let healthy = Arc::new(RecordingNotifyBackend::default());
let queue = test_queue().await;
let dispatcher = NotifyDispatcher::new(
"le",
vec![
BackendSlot::new("failing", failing.clone(), &every_kind()),
BackendSlot::new("healthy", healthy.clone(), &every_kind()),
],
queue,
);
assert!(matches!(
dispatcher
.deliver("failing", &profile_mounted("default"))
.await,
Some(Err(_))
));
assert!(matches!(
dispatcher
.deliver("healthy", &profile_mounted("default"))
.await,
Some(Ok(()))
));
assert_eq!(failing.events.lock().unwrap().len(), 1);
assert_eq!(healthy.events.lock().unwrap().len(), 1);
}
#[tokio::test]
async fn build_registry_builds_one_dispatcher_per_profile() {
let profiles = vec![
ProfileConfig {
name: "a".to_string(),
sections: crate::config::ProfileSections::default(),
},
ProfileConfig {
name: "b".to_string(),
sections: crate::config::ProfileSections::default(),
},
];
let registry = build_registry(
&profiles,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.unwrap();
assert_eq!(registry.len(), 2);
assert_eq!(registry["a"].profile(), "a");
assert_eq!(registry["b"].profile(), "b");
}
#[tokio::test]
async fn two_custom_entries_get_distinct_slot_ids() {
let dir = crate::testutil::TempDir::new("notify-slot");
let script = crate::testutil::write_script(&dir, "notify.sh", "#!/bin/sh\nexit 0\n");
let entry = || CustomNotifyConfig {
script_path: script.display().to_string(),
..CustomNotifyConfig::default()
};
let mut custom = std::collections::BTreeMap::new();
custom.insert("webhook".to_string(), entry());
custom.insert("pager".to_string(), entry());
let cfg = NotifyConfig {
enabled: vec!["custom".to_string()],
custom_enabled: vec!["webhook".to_string(), "pager".to_string()],
custom,
..NotifyConfig::default()
};
let dispatcher = from_config(
"le",
&cfg,
crate::testutil::outbound_with(test_resolver()),
&test_queue().await,
)
.expect("both custom entries must build");
assert!(dispatcher.slot("custom:webhook").is_some());
assert!(dispatcher.slot("custom:pager").is_some());
assert!(dispatcher.slot("custom").is_none());
}
#[test]
fn every_event_round_trips_through_its_payload() {
for event in every_event() {
let encoded = event.payload();
let decoded: NotifyEvent = serde_json::from_value(encoded.clone())
.unwrap_or_else(|error| panic!("{} must decode: {error}", event.kind()));
assert_eq!(decoded.kind(), event.kind());
assert_eq!(decoded.profile(), event.profile());
assert_eq!(decoded.payload(), encoded, "re-encoding must be stable");
}
}
#[test]
fn payload_is_tagged_with_its_own_hook() {
let event = NotifyEvent::CertificateIssued(CertificateIssuedData {
profile: "le".to_string(),
order_id: "ord-1".to_string(),
account_id: "acct-1".to_string(),
cert_serial: "0a0b".to_string(),
identifiers: vec!["a.example.com".to_string()],
client_ip: Some("203.0.113.5".to_string()),
});
assert_eq!(
event.payload(),
serde_json::json!({
"hook": "certificate_issued",
"profile": "le",
"order_id": "ord-1",
"account_id": "acct-1",
"cert_serial": "0a0b",
"identifiers": ["a.example.com"],
"client_ip": "203.0.113.5",
})
);
}
#[test]
fn template_dir_override_wins_over_the_embedded_default() {
let dir = crate::testutil::TempDir::new("notify");
std::fs::create_dir_all(dir.join("email")).unwrap();
std::fs::write(
dir.join("email/profile_mounted.subject.j2"),
"override: {{ profile }}",
)
.unwrap();
let env = build_environment(dir.path().to_str().unwrap());
let rendered = render(
&env,
"email/profile_mounted.subject.j2",
&profile_mounted("default"),
)
.unwrap();
assert_eq!(rendered, "override: default");
let rendered = render(
&env,
"email/profile_mounted.body.j2",
&profile_mounted("default"),
)
.unwrap();
assert!(rendered.contains("default"));
}
#[test]
fn a_template_failure_is_permanent() {
let env = build_environment("");
let missing = render(&env, "email/no_such_event.body.j2", &profile_mounted("le"))
.expect_err("there is no such template");
assert!(!missing.retryable(), "{missing}");
let dir = crate::testutil::TempDir::new("notify-broken");
std::fs::create_dir_all(dir.join("email")).unwrap();
std::fs::write(dir.join("email/profile_mounted.body.j2"), "{{ unclosed").unwrap();
let env = build_environment(dir.path().to_str().unwrap());
let broken = render(
&env,
"email/profile_mounted.body.j2",
&profile_mounted("le"),
)
.expect_err("the template does not compile");
assert!(!broken.retryable(), "{broken}");
}
#[test]
fn embedded_defaults_render_with_no_template_dir() {
let env = build_environment("");
let event = NotifyEvent::CertificateIssued(CertificateIssuedData {
profile: "le".to_string(),
order_id: "ord-1".to_string(),
account_id: "acc-1".to_string(),
cert_serial: "AA:BB".to_string(),
identifiers: vec!["example.com".to_string()],
client_ip: Some("203.0.113.1".to_string()),
});
let subject = render(&env, "email/certificate_issued.subject.j2", &event).unwrap();
assert!(subject.contains("le"), "{subject}");
let body = render(&env, "email/certificate_issued.body.j2", &event).unwrap();
assert!(body.contains("example.com"), "{body}");
assert!(body.contains("203.0.113.1"), "{body}");
}
#[test]
fn the_digest_templates_render_both_kinds_of_entry() {
let env = build_environment("");
let event = NotifyEvent::CertificatesExpiring(CertificatesExpiringData {
profile: "le".to_string(),
generated_at: 1_700_000_000,
lead_days: 14,
total: 7,
certificates: vec![
ExpiringCertificate {
order_id: "ord-1".to_string(),
account_id: "acc-1".to_string(),
cert_serial: "0a0b".to_string(),
identifiers: vec!["renew-me.example.com".to_string()],
not_after: 1_700_600_000,
days_remaining: 6,
superseded_by: None,
},
ExpiringCertificate {
order_id: "ord-2".to_string(),
account_id: "acc-1".to_string(),
cert_serial: "0c0d".to_string(),
identifiers: vec!["already-done.example.com".to_string()],
not_after: 1_700_600_000,
days_remaining: 6,
superseded_by: Some(SupersededBy {
order_id: "ord-3".to_string(),
cert_serial: "0e0f".to_string(),
not_after: 1_800_000_000,
via: "replaces".to_string(),
}),
},
],
});
let subject = render(&env, "email/certificates_expiring.subject.j2", &event).unwrap();
assert!(subject.contains('7'), "the count, not the page: {subject}");
assert!(subject.contains("le"), "{subject}");
let body = render(&env, "email/certificates_expiring.body.j2", &event).unwrap();
assert!(body.contains("renew-me.example.com"), "{body}");
assert!(body.contains("already-done.example.com"), "{body}");
assert!(body.contains("ord-3"), "the successor is named: {body}");
assert!(
body.contains("and 5 more"),
"a truncated digest says how many it did not name: {body}"
);
let hook = render(&env, "webhook/certificates_expiring.j2", &event).unwrap();
assert!(
!hook.contains('\n'),
"a webhook message is one line, since an entry's `body` wraps it: {hook}"
);
assert!(hook.contains("already replaced"), "{hook}");
}
}