1use std::collections::BTreeMap;
33use std::sync::Arc;
34use std::time::Duration;
35
36use async_trait::async_trait;
37use http_body_util::{BodyExt, Full};
38use hyper::body::Bytes;
39use hyper::header::{HeaderMap, HeaderName, HeaderValue};
40use hyper::{Method, Request};
41use tracing::info;
42use url::Url;
43
44use super::{NotifyBackend, NotifyError, NotifyEvent, render};
45use crate::config::WebhookNotifyConfig;
46
47use crate::http_client::MAX_ERROR_BODY_CHARS;
48
49const ALLOWED_METHODS: [&str; 3] = ["POST", "PUT", "PATCH"];
56
57pub struct WebhookNotifier {
58 entry: String,
61 url: Url,
62 method: Method,
63 headers: HeaderMap,
64 body_template: String,
66 timeout: Duration,
67 tls: Arc<rustls::ClientConfig>,
68 outbound: crate::http_client::Outbound,
71 env: Arc<minijinja::Environment<'static>>,
73}
74
75impl std::fmt::Debug for WebhookNotifier {
76 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
82 formatter
83 .debug_struct("WebhookNotifier")
84 .field("entry", &self.entry)
85 .field("webhook_host", &self.url.host_str())
86 .field("method", &self.method.as_str())
87 .field(
88 "headers",
89 &self
90 .headers
91 .keys()
92 .map(HeaderName::as_str)
93 .collect::<Vec<_>>(),
94 )
95 .field("timeout", &self.timeout)
96 .finish()
97 }
98}
99
100impl WebhookNotifier {
101 pub fn from_config(
107 entry: &str,
108 cfg: &WebhookNotifyConfig,
109 env: &minijinja::Environment<'static>,
110 outbound: crate::http_client::Outbound,
111 ) -> anyhow::Result<Self> {
112 let key = format!("notify.webhook.{entry}");
113
114 anyhow::ensure!(
115 !cfg.url.trim().is_empty(),
116 "{key} is enabled but {key}.url is empty"
117 );
118 let url: Url = cfg
119 .url
120 .parse()
121 .map_err(|error| anyhow::anyhow!("{key}.url is not a valid URL: {error}"))?;
122 anyhow::ensure!(
123 matches!(url.scheme(), "http" | "https"),
124 "{key}.url must be http:// or https://"
125 );
126
127 let spelled = cfg.method.trim().to_ascii_uppercase();
128 anyhow::ensure!(
129 ALLOWED_METHODS.contains(&spelled.as_str()),
130 "{key}.method: unknown method `{}` (expected one of {ALLOWED_METHODS:?})",
131 cfg.method
132 );
133 let method = Method::from_bytes(spelled.as_bytes())
134 .map_err(|error| anyhow::anyhow!("{key}.method is not a valid HTTP method: {error}"))?;
135
136 let headers = build_headers(&key, &cfg.headers)?;
137
138 anyhow::ensure!(
139 !cfg.body.trim().is_empty(),
140 "{key}.body is empty; a webhook with no payload delivers nothing"
141 );
142 let body_template = format!("webhook:{entry}.body");
148 let mut env = env.clone();
149 env.add_template_owned(body_template.clone(), cfg.body.clone())
150 .map_err(|error| anyhow::anyhow!("{key}.body is not a valid template: {error}"))?;
151
152 info!(
153 event = "notify_webhook_loaded",
154 outcome = "success",
155 entry = %entry,
156 method = %method,
157 webhook_host = ?url.host_str(),
158 );
159
160 Ok(Self {
161 entry: entry.to_string(),
162 url,
163 method,
164 headers,
165 body_template,
166 timeout: Duration::from_millis(cfg.timeout_ms),
167 tls: Arc::new(crate::http_client::webpki_tls_config()),
168 outbound,
169 env: Arc::new(env),
170 })
171 }
172
173 fn body_for(&self, event: &NotifyEvent) -> Result<String, NotifyError> {
176 let message = render(&self.env, &format!("webhook/{}.j2", event.kind()), event)?;
177
178 let template = self
179 .env
180 .get_template(&self.body_template)
181 .map_err(|error| {
182 NotifyError::permanent(format!("notify.webhook.{}.body: {error}", self.entry))
183 })?;
184 template
185 .render(minijinja::context! {
186 message,
187 hook => event.kind(),
188 ..event.context()
189 })
190 .map_err(|error| {
191 NotifyError::permanent(format!("notify.webhook.{}.body: {error}", self.entry))
192 })
193 }
194}
195
196fn build_headers(key: &str, configured: &BTreeMap<String, String>) -> anyhow::Result<HeaderMap> {
203 let mut headers = HeaderMap::new();
204 headers.insert(
205 hyper::header::CONTENT_TYPE,
206 HeaderValue::from_static("application/json"),
207 );
208 headers.insert(
209 hyper::header::USER_AGENT,
210 HeaderValue::from_static("acme-proxy"),
211 );
212
213 for (name, value) in configured {
214 let name = HeaderName::from_bytes(name.as_bytes()).map_err(|error| {
215 anyhow::anyhow!("{key}.headers: `{name}` is not a header name: {error}")
216 })?;
217 let value = HeaderValue::from_str(value).map_err(|_| {
219 anyhow::anyhow!("{key}.headers.{name}: the value is not a valid header value")
220 })?;
221 headers.insert(name, value);
222 }
223
224 Ok(headers)
225}
226
227#[async_trait]
228impl NotifyBackend for WebhookNotifier {
229 fn name(&self) -> &'static str {
230 "webhook"
233 }
234
235 async fn send(&self, event: &NotifyEvent) -> Result<(), NotifyError> {
236 let body = Bytes::from(self.body_for(event)?.into_bytes());
237
238 let (status, excerpt) = tokio::time::timeout(
239 self.timeout,
240 send_request(
241 &self.tls,
242 &self.outbound,
243 &self.method,
244 &self.url,
245 &self.headers,
246 body,
247 ),
248 )
249 .await
250 .map_err(|_| NotifyError::new(format!("timed out after {:?}", self.timeout)))??;
251
252 if status.is_success() {
253 Ok(())
254 } else {
255 let detail = if excerpt.is_empty() {
256 format!("webhook returned {status}")
257 } else {
258 format!("webhook returned {status}: {excerpt}")
259 };
260 if retryable_status(status) {
261 Err(NotifyError::new(detail))
262 } else {
263 Err(NotifyError::permanent(detail))
264 }
265 }
266 }
267}
268
269fn retryable_status(status: hyper::StatusCode) -> bool {
278 status.is_server_error()
279 || status == hyper::StatusCode::TOO_MANY_REQUESTS
280 || status == hyper::StatusCode::REQUEST_TIMEOUT
281}
282
283async fn send_request(
286 tls: &Arc<rustls::ClientConfig>,
287 outbound: &crate::http_client::Outbound,
288 method: &Method,
289 url: &Url,
290 headers: &HeaderMap,
291 body: Bytes,
292) -> Result<(hyper::StatusCode, String), NotifyError> {
293 let endpoint = crate::http_client::Endpoint::from_url(url).map_err(NotifyError::permanent)?;
296
297 let mut connection = outbound
299 .connect(&endpoint, tls)
300 .await
301 .map_err(NotifyError::new)?;
302
303 let mut builder = Request::builder()
306 .method(method.clone())
307 .uri(connection.request_target(url))
308 .header(hyper::header::HOST, endpoint.authority());
309 for (name, value) in headers {
310 builder = builder.header(name, value);
311 }
312 let request = builder
313 .body(Full::new(body))
314 .map_err(|error| NotifyError::permanent(format!("failed to build request: {error}")))?;
315
316 let response = connection
317 .send_request(request)
318 .await
319 .map_err(|error| NotifyError::new(error.to_string()))?;
320 let status = response.status();
321 let body = response
324 .into_body()
325 .collect()
326 .await
327 .map(http_body_util::Collected::to_bytes)
328 .unwrap_or_default();
329 let excerpt = String::from_utf8_lossy(&body)
330 .chars()
331 .take(MAX_ERROR_BODY_CHARS)
332 .collect::<String>()
333 .trim()
334 .to_string();
335 Ok((status, excerpt))
336}
337
338#[cfg(test)]
339mod tests {
340 use super::*;
341 use crate::notify::{ChallengeFailedData, ProfileMountedData, build_environment};
342 use hyper_util::rt::TokioIo;
343
344 fn test_resolver() -> Arc<dyn crate::dns::Resolver> {
348 Arc::new(crate::dns::HickoryResolver::from_system_uncached().unwrap())
349 }
350
351 fn cfg() -> WebhookNotifyConfig {
352 WebhookNotifyConfig {
353 url: "https://chat.example.com/hooks/xyz".to_string(),
354 ..WebhookNotifyConfig::default()
355 }
356 }
357
358 fn build(cfg: &WebhookNotifyConfig) -> anyhow::Result<WebhookNotifier> {
359 WebhookNotifier::from_config(
360 "chat",
361 cfg,
362 &build_environment(""),
363 crate::testutil::outbound_with(test_resolver()),
364 )
365 }
366
367 fn mounted() -> NotifyEvent {
368 NotifyEvent::ProfileMounted(ProfileMountedData {
369 profile: "default".to_string(),
370 })
371 }
372
373 async fn serve_once(
376 status: hyper::StatusCode,
377 body: &'static str,
378 ) -> (
379 std::net::SocketAddr,
380 tokio::sync::oneshot::Receiver<(String, HeaderMap, Bytes)>,
381 ) {
382 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
383 let addr = listener.local_addr().unwrap();
384 let (tx, rx) = tokio::sync::oneshot::channel();
385
386 tokio::spawn(async move {
387 let (stream, _) = listener.accept().await.unwrap();
388 let io = TokioIo::new(stream);
389 let tx = std::sync::Mutex::new(Some(tx));
390 let service = hyper::service::service_fn(move |req: Request<hyper::body::Incoming>| {
391 let tx = tx.lock().unwrap().take();
392 async move {
393 let method = req.method().to_string();
394 let headers = req.headers().clone();
395 let received = BodyExt::collect(req.into_body()).await.unwrap().to_bytes();
396 if let Some(tx) = tx {
397 let _ = tx.send((method, headers, received));
398 }
399 Ok::<_, std::convert::Infallible>(
400 hyper::Response::builder()
401 .status(status)
402 .body(Full::new(Bytes::from_static(body.as_bytes())))
403 .unwrap(),
404 )
405 }
406 });
407 let _ = hyper::server::conn::http1::Builder::new()
408 .serve_connection(io, service)
409 .await;
410 });
411
412 (addr, rx)
413 }
414
415 #[test]
418 fn an_unusable_entry_is_a_startup_error() {
419 let cases: Vec<(WebhookNotifyConfig, &str)> = vec![
420 (
421 WebhookNotifyConfig {
422 url: String::new(),
423 ..cfg()
424 },
425 "url is empty",
426 ),
427 (
428 WebhookNotifyConfig {
429 url: "not a url".to_string(),
430 ..cfg()
431 },
432 "not a valid URL",
433 ),
434 (
435 WebhookNotifyConfig {
436 url: "ftp://chat.example.com/hooks/xyz".to_string(),
437 ..cfg()
438 },
439 "must be http",
440 ),
441 (
442 WebhookNotifyConfig {
443 method: "GET".to_string(),
444 ..cfg()
445 },
446 "unknown method `GET`",
447 ),
448 (
449 WebhookNotifyConfig {
450 headers: BTreeMap::from([("not a header".to_string(), "x".to_string())]),
451 ..cfg()
452 },
453 "is not a header name",
454 ),
455 (
456 WebhookNotifyConfig {
457 headers: BTreeMap::from([("x-token".to_string(), "bad\nvalue".to_string())]),
458 ..cfg()
459 },
460 "not a valid header value",
461 ),
462 (
463 WebhookNotifyConfig {
464 body: " ".to_string(),
465 ..cfg()
466 },
467 "body is empty",
468 ),
469 (
470 WebhookNotifyConfig {
471 body: "{{ message".to_string(),
472 ..cfg()
473 },
474 "not a valid template",
475 ),
476 ];
477
478 for (config, expected) in cases {
479 let error = build(&config).unwrap_err().to_string();
480 assert!(
481 error.contains(expected) && error.contains("notify.webhook.chat"),
482 "expected `{expected}` naming the entry, got: {error}"
483 );
484 }
485 }
486
487 #[test]
491 fn a_method_is_case_insensitive() {
492 let notifier = build(&WebhookNotifyConfig {
493 method: "put".to_string(),
494 ..cfg()
495 })
496 .unwrap();
497 assert_eq!(notifier.method, Method::PUT);
498 }
499
500 #[test]
505 fn a_message_holding_quotes_and_newlines_still_renders_valid_json() {
506 let notifier = build(&cfg()).unwrap();
507 let event = NotifyEvent::ChallengeFailed(ChallengeFailedData {
508 profile: "default".to_string(),
509 order_id: "o1".to_string(),
510 account_id: "a1".to_string(),
511 authz_id: "z1".to_string(),
512 challenge_id: "c1".to_string(),
513 challenge_type: "http-01".to_string(),
514 identifier: "www.example.com".to_string(),
515 error: "fetched \"nonsense\"\nand a backslash \\".to_string(),
516 client_ip: None,
517 });
518
519 let body = notifier.body_for(&event).unwrap();
520 let parsed: serde_json::Value =
521 serde_json::from_str(&body).unwrap_or_else(|error| panic!("{error}: {body}"));
522 let text = parsed["text"].as_str().unwrap();
523 assert!(text.contains("nonsense"), "{text}");
524 assert!(text.contains("www.example.com"), "{text}");
525 }
526
527 #[test]
532 fn a_body_template_sees_the_events_own_fields() {
533 let notifier = build(&WebhookNotifyConfig {
534 body: r#"{"chat_id": "-100", "hook": {{ hook | tojson }}, "profile": {{ profile | tojson }}, "text": {{ message | tojson }}}"#
535 .to_string(),
536 ..cfg()
537 })
538 .unwrap();
539
540 let body = notifier.body_for(&mounted()).unwrap();
541 let parsed: serde_json::Value = serde_json::from_str(&body).unwrap();
542 assert_eq!(parsed["hook"], "profile_mounted");
543 assert_eq!(parsed["profile"], "default");
544 assert_eq!(parsed["chat_id"], "-100");
545 }
546
547 #[tokio::test]
551 async fn send_uses_the_configured_method_headers_and_body() {
552 let (addr, rx) = serve_once(hyper::StatusCode::OK, "").await;
553
554 let notifier = build(&WebhookNotifyConfig {
555 url: format!("http://{addr}/hooks/xyz"),
556 method: "PUT".to_string(),
557 headers: BTreeMap::from([
558 ("Authorization".to_string(), "Bearer s3cret".to_string()),
559 (
560 "content-type".to_string(),
561 "application/vnd.chat".to_string(),
562 ),
563 ]),
564 ..cfg()
565 })
566 .unwrap();
567
568 assert_eq!(notifier.name(), "webhook");
569 notifier
570 .send(&mounted())
571 .await
572 .expect("the webhook accepted the request");
573
574 let (method, headers, body) = rx.await.unwrap();
575 assert_eq!(method, "PUT");
576 assert_eq!(headers["authorization"], "Bearer s3cret");
577 assert_eq!(headers["content-type"], "application/vnd.chat");
578 assert_eq!(headers["user-agent"], "acme-proxy");
579 let parsed: serde_json::Value = serde_json::from_slice(&body).unwrap();
580 assert!(
581 parsed["text"].as_str().unwrap().contains("default"),
582 "the rendered text must name the profile: {parsed}"
583 );
584 }
585
586 #[tokio::test]
589 async fn a_rejecting_webhook_reports_its_status_and_body() {
590 let (addr, _rx) = serve_once(hyper::StatusCode::BAD_REQUEST, "invalid_payload").await;
591
592 let notifier = build(&WebhookNotifyConfig {
593 url: format!("http://{addr}/hooks/xyz"),
594 ..cfg()
595 })
596 .unwrap();
597
598 let error = notifier
599 .send(&mounted())
600 .await
601 .expect_err("400 is not a delivery");
602 assert!(error.to_string().contains("400"), "{error}");
603 assert!(error.to_string().contains("invalid_payload"), "{error}");
604 assert!(
605 !error.retryable(),
606 "a 400 is the provider stating a reason, not a bad minute"
607 );
608 }
609
610 #[tokio::test]
614 async fn an_unusable_url_is_permanent_and_an_unreachable_host_is_not() {
615 let tls = Arc::new(crate::http_client::webpki_tls_config());
616
617 let url: Url = "ftp://chat.example.com/hooks/xyz".parse().unwrap();
618 let error = send_request(
619 &tls,
620 &crate::testutil::outbound_with(test_resolver()),
621 &Method::POST,
622 &url,
623 &HeaderMap::new(),
624 Bytes::from_static(b"{}"),
625 )
626 .await
627 .expect_err("ftp is not a webhook transport");
628 assert!(error.to_string().contains("unsupported scheme"), "{error}");
629 assert!(!error.retryable(), "{error}");
630
631 let url: Url = "http://127.0.0.1:1/hooks/xyz".parse().unwrap();
632 let error = send_request(
633 &tls,
634 &crate::testutil::outbound_with(test_resolver()),
635 &Method::POST,
636 &url,
637 &HeaderMap::new(),
638 Bytes::from_static(b"{}"),
639 )
640 .await
641 .expect_err("nothing is listening on port 1");
642 assert!(error.to_string().contains("connecting to"), "{error}");
643 assert!(
644 error.retryable(),
645 "a host that is down may come back: {error}"
646 );
647 }
648
649 #[test]
654 fn a_body_that_fails_at_render_time_is_permanent() {
655 let notifier = build(&WebhookNotifyConfig {
656 body: "{{ message.no_such_method() }}".to_string(),
657 ..cfg()
658 })
659 .unwrap();
660
661 let error = notifier.body_for(&mounted()).unwrap_err();
662 assert!(
663 error.to_string().contains("notify.webhook.chat.body"),
664 "{error}"
665 );
666 assert!(!error.retryable(), "{error}");
667 }
668
669 #[tokio::test]
673 async fn a_silent_webhook_times_out_and_is_retryable() {
674 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
675 let addr = listener.local_addr().unwrap();
676 tokio::spawn(async move {
677 let (stream, _) = listener.accept().await.unwrap();
678 std::future::pending::<()>().await;
680 drop(stream);
681 });
682
683 let notifier = build(&WebhookNotifyConfig {
684 url: format!("http://{addr}/hooks/xyz"),
685 timeout_ms: 100,
686 ..cfg()
687 })
688 .unwrap();
689
690 let error = notifier
691 .send(&mounted())
692 .await
693 .expect_err("a silent server is not a delivery");
694 assert!(error.to_string().contains("timed out"), "{error}");
695 assert!(error.retryable(), "{error}");
696 }
697
698 #[test]
702 fn only_a_transient_status_is_retried() {
703 use hyper::StatusCode;
704 for status in [
705 StatusCode::INTERNAL_SERVER_ERROR,
706 StatusCode::BAD_GATEWAY,
707 StatusCode::SERVICE_UNAVAILABLE,
708 StatusCode::TOO_MANY_REQUESTS,
709 StatusCode::REQUEST_TIMEOUT,
710 ] {
711 assert!(retryable_status(status), "{status} must be retried");
712 }
713 for status in [
714 StatusCode::BAD_REQUEST,
715 StatusCode::UNAUTHORIZED,
716 StatusCode::FORBIDDEN,
717 StatusCode::NOT_FOUND,
718 StatusCode::GONE,
719 ] {
720 assert!(!retryable_status(status), "{status} must not be retried");
721 }
722 }
723
724 #[test]
728 fn debug_renders_neither_the_url_path_nor_a_header_value() {
729 let notifier = build(&WebhookNotifyConfig {
730 url: "https://chat.example.com/hooks/T00/B00/s3cret-hook-id".to_string(),
731 headers: BTreeMap::from([("Authorization".to_string(), "Bearer s3cret".to_string())]),
732 ..cfg()
733 })
734 .unwrap();
735
736 let rendered = format!("{notifier:?}");
737 assert!(rendered.contains("chat.example.com"), "{rendered}");
738 assert!(rendered.contains("authorization"), "{rendered}");
739 assert!(!rendered.contains("s3cret"), "{rendered}");
740 }
741}