1use 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#[async_trait]
79pub trait NotifyBackend: Send + Sync {
80 fn name(&self) -> &'static str;
82
83 async fn send(&self, event: &NotifyEvent) -> Result<(), NotifyError>;
87}
88
89#[derive(Debug, thiserror::Error)]
99#[error("{detail}")]
100pub struct NotifyError {
101 detail: String,
102 retryable: bool,
103}
104
105impl NotifyError {
106 pub fn new(detail: impl Into<String>) -> Self {
108 Self {
109 detail: detail.into(),
110 retryable: true,
111 }
112 }
113
114 pub fn permanent(detail: impl Into<String>) -> Self {
118 Self {
119 detail: detail.into(),
120 retryable: false,
121 }
122 }
123
124 #[must_use]
126 pub fn retryable(&self) -> bool {
127 self.retryable
128 }
129}
130
131#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
139pub struct ProfileMountedData {
140 pub profile: String,
141}
142
143#[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#[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#[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#[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#[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#[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 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 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 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 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 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 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 fn order_id(&self) -> Option<&str> {
317 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 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 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
358pub struct BackendSlot {
367 id: String,
368 events: HashSet<String>,
372 backend: Arc<dyn NotifyBackend>,
373}
374
375impl BackendSlot {
376 #[must_use]
378 pub fn id(&self) -> &str {
379 &self.id
380 }
381
382 #[must_use]
384 pub fn wants(&self, event: &NotifyEvent) -> bool {
385 self.events.contains(event.kind())
386 }
387
388 #[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
404pub struct NotifyDispatcher {
410 profile: String,
411 slots: Vec<BackendSlot>,
412 jobs: JobQueue,
413}
414
415impl std::fmt::Debug for NotifyDispatcher {
416 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 #[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 #[must_use]
445 pub fn disabled(jobs: JobQueue) -> Self {
446 Self::new("default", Vec::new(), jobs)
447 }
448
449 #[must_use]
451 pub fn profile(&self) -> &str {
452 &self.profile
453 }
454
455 #[must_use]
458 pub fn slot(&self, id: &str) -> Option<&BackendSlot> {
459 self.slots.iter().find(|slot| slot.id == id)
460 }
461
462 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 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
520fn 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
534pub fn from_config(
540 profile: &str,
541 cfg: &NotifyConfig,
542 outbound: crate::http_client::Outbound,
543 jobs: &JobQueue,
544) -> anyhow::Result<Arc<NotifyDispatcher>> {
545 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 "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
598fn 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
627fn 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
660pub type DispatcherMap = HashMap<String, Arc<NotifyDispatcher>>;
665
666pub type NotifiersSender = tokio::sync::watch::Sender<Arc<DispatcherMap>>;
668
669#[derive(Clone)]
688pub struct Notifiers(tokio::sync::watch::Receiver<Arc<DispatcherMap>>);
689
690impl Notifiers {
691 #[must_use]
693 pub fn get(&self, profile: &str) -> Option<Arc<NotifyDispatcher>> {
694 self.0.borrow().get(profile).cloned()
698 }
699}
700
701impl 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#[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
728pub 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
751macro_rules! embed {
759 ($name:literal) => {
760 ($name, include_str!(concat!("templates/", $name)))
761 };
762}
763
764static 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
791pub(crate) fn build_environment(template_dir: &str) -> minijinja::Environment<'static> {
796 crate::templating::loader_env(template_dir, &EMBEDDED_TEMPLATES)
797}
798
799pub(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 fn test_resolver() -> Arc<dyn crate::dns::Resolver> {
828 Arc::new(crate::dns::HickoryResolver::from_system_uncached().unwrap())
829 }
830
831 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 #[derive(Default, Clone, Copy)]
840 enum Failure {
841 #[default]
842 None,
843 Retryable,
844 Permanent,
845 }
846
847 #[derive(Default)]
850 pub(crate) struct RecordingNotifyBackend {
851 pub(crate) events: Mutex<Vec<NotifyEvent>>,
852 fail: Failure,
853 }
854
855 impl RecordingNotifyBackend {
856 pub(crate) fn failing() -> Self {
858 Self {
859 events: Mutex::new(Vec::new()),
860 fail: Failure::Retryable,
861 }
862 }
863
864 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 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 #[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 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 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 assert_eq!(events[4].client_ip(), None);
990 assert_eq!(events[1].client_ip(), Some("203.0.113.5"));
991
992 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 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 #[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 #[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 assert!(notifiers.get("staging").is_none());
1055
1056 sender.send_replace(Arc::new(HashMap::new()));
1058 assert!(notifiers.get("le").is_none());
1059 }
1060
1061 #[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 assert!(notifiers.clone().get("le").is_some());
1081 }
1082
1083 fn every_kind() -> Vec<String> {
1085 ALL_NOTIFY_EVENTS.iter().map(|k| (*k).to_string()).collect()
1086 }
1087
1088 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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 #[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 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 #[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}