1pub use isb_core::net;
20pub mod provider;
21mod rate;
22pub mod smtp;
23
24use std::collections::{BTreeMap, VecDeque};
25use std::path::{Path, PathBuf};
26use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
27use std::sync::{Arc, Condvar, Mutex};
28use std::time::{Duration, Instant};
29
30use serde::{Deserialize, Serialize};
31
32use crate::error::{Error, Result};
33use crate::org::OrgId;
34use crate::secrets::Secrets;
35use crate::stack::controller::{Controller, Event, now_ms};
36use net::{Net, SendError};
37pub use provider::{Message, Provider};
38pub use rate::{BACKOFF_MAX, RateLimit, backoff};
39
40pub const QUEUE_MAX: usize = 100;
42pub const LOG_KEPT: usize = 50;
44pub const ATTEMPTS: u32 = 6;
46pub const PER_MINUTE: usize = 20;
48const BACKOFF: Duration = Duration::from_secs(5);
50const IDLE: Duration = Duration::from_secs(60);
52
53#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
57#[serde(deny_unknown_fields)]
58pub struct Rule {
59 #[serde(default = "all_events")]
60 pub events: Vec<String>,
61 #[serde(default, skip_serializing_if = "Vec::is_empty")]
63 pub projects: Vec<String>,
64 #[serde(default, skip_serializing_if = "Vec::is_empty")]
66 pub apps: Vec<String>,
67 #[serde(default, skip_serializing_if = "Vec::is_empty")]
69 pub stacks: Vec<String>,
70}
71
72fn all_events() -> Vec<String> {
73 vec!["*".into()]
74}
75
76impl Default for Rule {
77 fn default() -> Rule {
78 Rule {
79 events: all_events(),
80 projects: Vec::new(),
81 apps: Vec::new(),
82 stacks: Vec::new(),
83 }
84 }
85}
86
87pub fn glob(pattern: &str, s: &str) -> bool {
89 let (p, s) = (pattern.as_bytes(), s.as_bytes());
90 let (mut pi, mut si) = (0, 0);
91 let (mut star, mut mark) = (None, 0);
92 while si < s.len() {
93 if pi < p.len() && p[pi] == b'*' {
94 star = Some(pi);
95 mark = si;
96 pi += 1;
97 } else if pi < p.len() && p[pi] == s[si] {
98 pi += 1;
99 si += 1;
100 } else if let Some(st) = star {
101 pi = st + 1;
102 mark += 1;
103 si = mark;
104 } else {
105 return false;
106 }
107 }
108 p[pi..].iter().all(|c| *c == b'*')
109}
110
111#[derive(Debug, Clone, Default)]
113pub struct Subject<'a> {
114 pub kind: &'a str,
115 pub stack: &'a str,
117 pub service: &'a str,
118 pub project: Option<&'a str>,
120}
121
122impl Rule {
123 pub fn matches(&self, s: &Subject) -> bool {
124 self.events.iter().any(|p| glob(p, s.kind))
125 && (self.stacks.is_empty() || self.stacks.iter().any(|x| x == s.stack))
126 && (self.apps.is_empty() || self.apps.iter().any(|x| x == s.service))
127 && (self.projects.is_empty()
128 || s.project
129 .is_some_and(|p| self.projects.iter().any(|x| x == p)))
130 }
131
132 fn validate(&self) -> Result<()> {
133 if self.events.is_empty() {
134 return Err(Error::invalid("a rule needs at least one event pattern"));
135 }
136 for e in &self.events {
137 if e.is_empty()
138 || e.len() > 64
139 || !e
140 .chars()
141 .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || "._-*".contains(c))
142 {
143 return Err(Error::invalid(format!(
144 "event pattern {e:?}: [a-z0-9._-] and * globs, like deploy.* or *.failed"
145 )));
146 }
147 }
148 Ok(())
149 }
150}
151
152#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
154#[serde(deny_unknown_fields)]
155pub struct Channel {
156 pub name: String,
157 pub provider: Provider,
158 #[serde(default = "yes")]
159 pub enabled: bool,
160 #[serde(default = "default_rules")]
162 pub rules: Vec<Rule>,
163 #[serde(default)]
164 pub created_at: u64,
165 #[serde(default)]
166 pub updated_at: u64,
167}
168
169fn yes() -> bool {
170 true
171}
172
173fn default_rules() -> Vec<Rule> {
174 vec![Rule::default()]
175}
176
177impl Channel {
178 pub fn matches(&self, s: &Subject) -> bool {
179 self.enabled && self.rules.iter().any(|r| r.matches(s))
180 }
181
182 pub fn validate(&self) -> Result<()> {
183 validate_name(&self.name)?;
184 self.provider.validate().map_err(Error::Invalid)?;
185 if self.rules.len() > 20 {
186 return Err(Error::invalid("at most 20 rules per channel"));
187 }
188 for r in &self.rules {
189 r.validate()?;
190 }
191 Ok(())
192 }
193}
194
195pub fn validate_name(n: &str) -> Result<()> {
197 let ok = !n.is_empty()
198 && n.len() <= 40
199 && n.starts_with(|c: char| c.is_ascii_lowercase())
200 && n.chars()
201 .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '-');
202 if ok {
203 Ok(())
204 } else {
205 Err(Error::invalid(format!(
206 "channel name {n:?}: [a-z0-9-], starting with a letter, at most 40 characters"
207 )))
208 }
209}
210
211#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
213pub struct Delivery {
214 pub id: String,
215 pub at: u64,
217 pub kind: String,
218 pub seq: u64,
219 pub summary: String,
220 pub status: String,
223 pub attempts: u32,
224 #[serde(default, skip_serializing_if = "Option::is_none")]
225 pub http_status: Option<u16>,
226 #[serde(default, skip_serializing_if = "Option::is_none")]
227 pub error: Option<String>,
228 #[serde(default, skip_serializing_if = "Option::is_none")]
229 pub finished_at: Option<u64>,
230 #[serde(default)]
231 pub test: bool,
232}
233
234#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
236pub struct Settings {
237 #[serde(default)]
239 pub allow_private_targets: bool,
240}
241
242pub type Resolve = Arc<dyn Fn(&OrgId, &str, &str) -> Option<String> + Send + Sync>;
244
245pub type Details = Arc<dyn Fn(&OrgId, &Event) -> Option<serde_json::Value> + Send + Sync>;
248
249pub fn event_org(stack: &str) -> (OrgId, &str) {
252 match stack.split_once('/') {
253 Some((o, s)) => (OrgId::new(o).unwrap_or_else(|_| OrgId::default_org()), s),
254 None => (OrgId::default_org(), stack),
255 }
256}
257
258struct Job {
259 org: OrgId,
260 channel: String,
261 msg: Message,
262}
263
264struct Queue {
265 jobs: Mutex<(VecDeque<Job>, bool)>,
266 wake: Condvar,
267}
268
269struct Inner {
270 state: PathBuf,
271 secrets: Arc<Secrets>,
272 settings: Mutex<Settings>,
273 edit: Mutex<()>,
275 queues: Mutex<BTreeMap<(OrgId, String), Arc<Queue>>>,
276 logs: Mutex<BTreeMap<(OrgId, String), VecDeque<Delivery>>>,
277 resolve: Resolve,
278 details: Mutex<Option<Details>>,
279 next_id: AtomicU64,
280 stop: AtomicBool,
281 backoff: Duration,
282 per_minute: usize,
283 tls: Option<Arc<rustls::ClientConfig>>,
285}
286
287#[derive(Clone)]
289pub struct Notifier {
290 inner: Arc<Inner>,
291}
292
293fn settings_path(state: &Path) -> PathBuf {
294 state.join("notify.json")
295}
296
297fn write_atomic(path: &Path, data: &[u8]) -> Result<()> {
298 if let Some(d) = path.parent() {
299 std::fs::create_dir_all(d)?;
300 }
301 let tmp = path.with_extension("tmp");
302 std::fs::write(&tmp, data)?;
303 std::fs::rename(&tmp, path)?;
304 Ok(())
305}
306
307impl Notifier {
308 pub fn new(state: &Path, secrets: Arc<Secrets>, resolve: Resolve) -> Result<Notifier> {
311 let settings = match std::fs::read(settings_path(state)) {
312 Ok(b) => serde_json::from_slice(&b)
313 .map_err(|e| Error::invalid(format!("{}: {e}", settings_path(state).display())))?,
314 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Settings::default(),
315 Err(e) => return Err(e.into()),
316 };
317 Ok(Notifier {
318 inner: Arc::new(Inner {
319 state: state.to_path_buf(),
320 secrets,
321 settings: Mutex::new(settings),
322 edit: Mutex::new(()),
323 queues: Mutex::new(BTreeMap::new()),
324 logs: Mutex::new(BTreeMap::new()),
325 resolve,
326 details: Mutex::new(None),
327 next_id: AtomicU64::new(1),
328 stop: AtomicBool::new(false),
329 backoff: BACKOFF,
330 per_minute: PER_MINUTE,
331 tls: None,
332 }),
333 })
334 }
335
336 pub fn start(&self, ctl: Controller) {
338 let me = self.clone();
339 let _ = std::thread::Builder::new()
340 .name("isb-notify".into())
341 .spawn(move || {
342 let mut since = 0;
343 while !me.inner.stop.load(Ordering::SeqCst) {
344 let (seq, evs) = ctl.wait_events(since, 1000, Duration::from_secs(5));
345 for e in &evs {
346 me.route(e);
347 }
348 since = since.max(seq);
351 }
352 });
353 }
354
355 pub fn set_details(&self, d: Details) {
357 *self.inner.details.lock().unwrap() = Some(d);
358 }
359
360 pub fn shutdown(&self) {
361 self.inner.stop.store(true, Ordering::SeqCst);
362 for q in self.inner.queues.lock().unwrap().values() {
363 q.wake.notify_all();
364 }
365 }
366
367 pub fn settings(&self) -> Settings {
368 self.inner.settings.lock().unwrap().clone()
369 }
370
371 pub fn set_settings(&self, s: Settings) -> Result<()> {
372 write_atomic(
373 &settings_path(&self.inner.state),
374 &serde_json::to_vec_pretty(&s)?,
375 )?;
376 *self.inner.settings.lock().unwrap() = s;
377 Ok(())
378 }
379
380 fn net(&self) -> Net {
381 let allow = self.inner.settings.lock().unwrap().allow_private_targets;
382 Net {
383 allow_private: allow,
384 tls: self.inner.tls.clone().unwrap_or_else(net::default_tls),
385 }
386 }
387
388 fn dir(&self, org: &OrgId) -> PathBuf {
389 org.dir(&self.inner.state).join("notify")
390 }
391
392 pub fn list(&self, org: &OrgId) -> Result<Vec<Channel>> {
394 match std::fs::read(self.dir(org).join("channels.json")) {
395 Ok(b) => Ok(serde_json::from_slice(&b)?),
396 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(Vec::new()),
397 Err(e) => Err(e.into()),
398 }
399 }
400
401 fn save(&self, org: &OrgId, chans: &[Channel]) -> Result<()> {
402 write_atomic(
403 &self.dir(org).join("channels.json"),
404 &serde_json::to_vec_pretty(chans)?,
405 )
406 }
407
408 pub fn get(&self, org: &OrgId, name: &str) -> Result<Channel> {
409 self.list(org)?
410 .into_iter()
411 .find(|c| c.name == name)
412 .ok_or_else(|| Error::NotFound(format!("notification channel {name}")))
413 }
414
415 fn check_secrets(&self, org: &OrgId, c: &Channel) -> Result<()> {
418 for s in c.provider.secrets() {
419 self.inner.secrets.inspect(org, s).map_err(|e| {
420 if e.is_not_found() {
421 Error::invalid(format!(
422 "no secret {s} in org {org}: create it first (isb secret create {s})"
423 ))
424 } else {
425 e
426 }
427 })?;
428 }
429 if let Provider::Webhook { url_secret, .. }
430 | Provider::Slack { url_secret }
431 | Provider::Discord { url_secret } = &c.provider
432 {
433 let v = self
434 .secret_value(org, url_secret)
435 .map_err(|e| Error::Invalid(e.message))?;
436 provider::check_url_value(&c.provider, &v)
437 .map_err(|e| Error::invalid(format!("secret {url_secret}: {e}")))?;
438 }
439 Ok(())
440 }
441
442 pub fn create(&self, org: &OrgId, mut c: Channel) -> Result<Channel> {
443 c.validate()?;
444 self.check_secrets(org, &c)?;
445 let _g = self.inner.edit.lock().unwrap();
446 let mut all = self.list(org)?;
447 if all.iter().any(|x| x.name == c.name) {
448 return Err(Error::invalid(format!(
449 "notification channel {} exists",
450 c.name
451 )));
452 }
453 if all.len() >= 50 {
454 return Err(Error::invalid("at most 50 notification channels per org"));
455 }
456 let now = now_ms() / 1000;
457 c.created_at = now;
458 c.updated_at = now;
459 all.push(c.clone());
460 self.save(org, &all)?;
461 Ok(c)
462 }
463
464 pub fn update(
466 &self,
467 org: &OrgId,
468 name: &str,
469 provider: Option<Provider>,
470 rules: Option<Vec<Rule>>,
471 enabled: Option<bool>,
472 ) -> Result<Channel> {
473 let _g = self.inner.edit.lock().unwrap();
474 let mut all = self.list(org)?;
475 let c = all
476 .iter_mut()
477 .find(|c| c.name == name)
478 .ok_or_else(|| Error::NotFound(format!("notification channel {name}")))?;
479 let mut n = c.clone();
480 if let Some(p) = provider {
481 n.provider = p;
482 }
483 if let Some(r) = rules {
484 n.rules = r;
485 }
486 if let Some(e) = enabled {
487 n.enabled = e;
488 }
489 n.validate()?;
490 self.check_secrets(org, &n)?;
491 n.updated_at = now_ms() / 1000;
492 *c = n.clone();
493 self.save(org, &all)?;
494 Ok(n)
495 }
496
497 pub fn delete(&self, org: &OrgId, name: &str) -> Result<()> {
498 let _g = self.inner.edit.lock().unwrap();
499 let mut all = self.list(org)?;
500 let before = all.len();
501 all.retain(|c| c.name != name);
502 if all.len() == before {
503 return Err(Error::NotFound(format!("notification channel {name}")));
504 }
505 self.save(org, &all)?;
506 let key = (org.clone(), name.to_string());
507 self.inner.logs.lock().unwrap().remove(&key);
508 let _ = std::fs::remove_file(self.log_path(org, name));
509 if let Some(q) = self.inner.queues.lock().unwrap().get(&key) {
510 q.jobs.lock().unwrap().0.clear();
511 }
512 Ok(())
513 }
514
515 fn log_path(&self, org: &OrgId, name: &str) -> PathBuf {
516 self.dir(org)
517 .join("deliveries")
518 .join(format!("{name}.json"))
519 }
520
521 pub fn deliveries(&self, org: &OrgId, name: &str) -> Result<Vec<Delivery>> {
523 self.get(org, name)?;
524 let mut logs = self.inner.logs.lock().unwrap();
525 let log = self.log_of(&mut logs, org, name);
526 Ok(log.iter().rev().cloned().collect())
527 }
528
529 fn log_of<'a>(
530 &self,
531 logs: &'a mut BTreeMap<(OrgId, String), VecDeque<Delivery>>,
532 org: &OrgId,
533 name: &str,
534 ) -> &'a mut VecDeque<Delivery> {
535 logs.entry((org.clone(), name.to_string()))
536 .or_insert_with(|| {
537 std::fs::read(self.log_path(org, name))
538 .ok()
539 .and_then(|b| serde_json::from_slice(&b).ok())
540 .unwrap_or_default()
541 })
542 }
543
544 fn record(&self, org: &OrgId, name: &str, d: &Delivery) {
546 let mut logs = self.inner.logs.lock().unwrap();
547 let log = self.log_of(&mut logs, org, name);
548 match log.iter_mut().find(|x| x.id == d.id) {
549 Some(x) => *x = d.clone(),
550 None => {
551 log.push_back(d.clone());
552 while log.len() > LOG_KEPT {
553 log.pop_front();
554 }
555 }
556 }
557 let data = serde_json::to_vec(&*log).unwrap_or_default();
558 if let Err(e) = write_atomic(&self.log_path(org, name), &data) {
559 eprintln!("isb serve: notify: cannot save the delivery log of {org}/{name}: {e}");
560 }
561 }
562
563 fn new_id(&self) -> String {
564 format!(
565 "{}-{}",
566 now_ms(),
567 self.inner.next_id.fetch_add(1, Ordering::SeqCst)
568 )
569 }
570
571 pub fn route(&self, e: &Event) {
573 let Some(kind) = e.kind.as_deref() else {
574 return;
575 };
576 let (org, stack) = event_org(&e.stack);
577 let Ok(chans) = self.list(&org) else { return };
578 if chans.is_empty() {
579 return;
580 }
581 let project = if e.service.is_empty() {
582 None
583 } else {
584 (self.inner.resolve)(&org, stack, &e.service)
585 };
586 let subject = Subject {
587 kind,
588 stack,
589 service: &e.service,
590 project: project.as_deref(),
591 };
592 let details = self
593 .inner
594 .details
595 .lock()
596 .unwrap()
597 .clone()
598 .and_then(|d| d(&org, e));
599 for c in chans.iter().filter(|c| c.matches(&subject)) {
600 let msg = Message {
601 id: self.new_id(),
602 org: org.to_string(),
603 kind: kind.to_string(),
604 level: e.level.clone(),
605 stack: stack.to_string(),
606 service: e.service.clone(),
607 project: project.clone(),
608 instance: e.instance.clone(),
609 message: e.message.clone(),
610 details: details.clone(),
611 at: e.at,
612 seq: e.seq,
613 test: false,
614 };
615 self.enqueue(Job {
616 org: org.clone(),
617 channel: c.name.clone(),
618 msg,
619 });
620 }
621 }
622
623 fn delivery(msg: &Message, status: &str) -> Delivery {
624 Delivery {
625 id: msg.id.clone(),
626 at: now_ms(),
627 kind: msg.kind.clone(),
628 seq: msg.seq,
629 summary: msg.message.chars().take(200).collect(),
630 status: status.into(),
631 attempts: 0,
632 http_status: None,
633 error: None,
634 finished_at: None,
635 test: msg.test,
636 }
637 }
638
639 fn enqueue(&self, job: Job) {
640 let key = (job.org.clone(), job.channel.clone());
641 let q = self
642 .inner
643 .queues
644 .lock()
645 .unwrap()
646 .entry(key.clone())
647 .or_insert_with(|| {
648 Arc::new(Queue {
649 jobs: Mutex::new((VecDeque::new(), false)),
650 wake: Condvar::new(),
651 })
652 })
653 .clone();
654 self.record(&job.org, &job.channel, &Self::delivery(&job.msg, "queued"));
655 let mut dropped = None;
656 let spawn = {
657 let mut g = q.jobs.lock().unwrap();
658 if g.0.len() >= QUEUE_MAX {
659 dropped = g.0.pop_front();
660 }
661 g.0.push_back(job);
662 q.wake.notify_all();
663 !std::mem::replace(&mut g.1, true)
664 };
665 if let Some(d) = dropped {
666 let mut r = Self::delivery(&d.msg, "dropped");
667 r.error = Some(format!("the channel's queue was full ({QUEUE_MAX})"));
668 r.finished_at = Some(now_ms());
669 self.record(&d.org, &d.channel, &r);
670 }
671 if spawn {
672 let me = self.clone();
673 let r = std::thread::Builder::new()
674 .name(format!("isb-notify-{}", key.1))
675 .spawn(move || me.sender(q));
676 if let Err(e) = r {
677 eprintln!("isb serve: notify: cannot start a sender: {e}");
678 }
679 }
680 }
681
682 fn sender(&self, q: Arc<Queue>) {
684 let mut rate = RateLimit::default();
685 loop {
686 let job = {
687 let mut g = q.jobs.lock().unwrap();
688 let started = Instant::now();
689 loop {
690 if self.inner.stop.load(Ordering::SeqCst) {
691 g.1 = false;
692 return;
693 }
694 if let Some(j) = g.0.pop_front() {
695 break j;
696 }
697 let left = IDLE.saturating_sub(started.elapsed());
698 if left.is_zero() {
699 g.1 = false;
700 return;
701 }
702 g = q.wake.wait_timeout(g, left).unwrap().0;
703 }
704 };
705 let wait = rate.wait(Instant::now(), self.inner.per_minute);
706 if !wait.is_zero() {
707 self.sleep(wait);
708 }
709 rate.record(Instant::now());
710 self.deliver(&job);
711 }
712 }
713
714 fn sleep(&self, d: Duration) {
716 let end = Instant::now() + d;
717 while !self.inner.stop.load(Ordering::SeqCst) {
718 let left = end.saturating_duration_since(Instant::now());
719 if left.is_zero() {
720 return;
721 }
722 std::thread::sleep(left.min(Duration::from_millis(500)));
723 }
724 }
725
726 fn deliver(&self, job: &Job) {
728 let mut d = Self::delivery(&job.msg, "queued");
729 if let Some(prev) = self
731 .inner
732 .logs
733 .lock()
734 .unwrap()
735 .get(&(job.org.clone(), job.channel.clone()))
736 .and_then(|l| l.iter().find(|x| x.id == job.msg.id).cloned())
737 {
738 d.at = prev.at;
739 }
740 for attempt in 1..=ATTEMPTS {
741 let ch = match self.get(&job.org, &job.channel) {
743 Ok(c) if c.enabled => c,
744 Ok(_) | Err(_) => {
745 d.status = "skipped".into();
746 d.error = Some("the channel was disabled or removed".into());
747 d.finished_at = Some(now_ms());
748 self.record(&job.org, &job.channel, &d);
749 return;
750 }
751 };
752 d.attempts = attempt;
753 let r = self.send(&job.org, &ch, &job.msg);
754 match r {
755 Ok(status) => {
756 d.status = "sent".into();
757 d.http_status = status;
758 d.error = None;
759 d.finished_at = Some(now_ms());
760 self.record(&job.org, &job.channel, &d);
761 return;
762 }
763 Err((status, e, retry_after)) => {
764 d.http_status = status;
765 d.error = Some(e.message.clone());
766 if !e.retryable || attempt == ATTEMPTS || self.inner.stop.load(Ordering::SeqCst)
767 {
768 d.status = "failed".into();
769 d.finished_at = Some(now_ms());
770 self.record(&job.org, &job.channel, &d);
771 eprintln!(
772 "isb serve: notify {}/{}: {} not delivered: {}",
773 job.org, job.channel, job.msg.kind, e.message
774 );
775 return;
776 }
777 d.status = "retrying".into();
778 self.record(&job.org, &job.channel, &d);
779 self.sleep(backoff(self.inner.backoff, attempt, retry_after));
780 }
781 }
782 }
783 }
784
785 fn secret_value(&self, org: &OrgId, name: &str) -> std::result::Result<String, SendError> {
786 let (v, _) = self.inner.secrets.get(org, name).map_err(|e| {
787 if e.is_not_found() {
788 SendError::permanent(format!("secret {name} is gone"))
789 } else {
790 SendError::transient(format!("secret {name}: {e}"))
791 }
792 })?;
793 String::from_utf8(v)
794 .map(|s| s.trim().to_string())
795 .map_err(|_| SendError::permanent(format!("secret {name} is not UTF-8 text")))
796 }
797
798 fn send(
801 &self,
802 org: &OrgId,
803 ch: &Channel,
804 msg: &Message,
805 ) -> std::result::Result<Option<u16>, (Option<u16>, SendError, Option<Duration>)> {
806 let net = self.net();
807 let secret = |n: &str| self.secret_value(org, n);
808 if let Provider::Email {
809 host,
810 port,
811 tls,
812 username,
813 password_secret,
814 from,
815 to,
816 } = &ch.provider
817 {
818 let password = match password_secret {
819 Some(p) => Some(secret(p).map_err(|e| (None, e, None))?),
820 None => None,
821 };
822 let body = format!(
823 "{}{}\n\nlevel: {}\norg: {}\nstack: {}\n{}{}at: {}\n",
824 msg.message,
825 provider::detail_lines(msg),
826 msg.level,
827 msg.org,
828 msg.stack,
829 if msg.service.is_empty() {
830 String::new()
831 } else {
832 format!("service: {}\n", msg.service)
833 },
834 msg.project
835 .as_ref()
836 .map(|p| format!("project: {p}\n"))
837 .unwrap_or_default(),
838 smtp::rfc2822(msg.at / 1000),
839 );
840 let mid = format!("{}@isb", msg.id);
841 return smtp::send(
842 &net,
843 &smtp::Mail {
844 host,
845 port: port.unwrap_or(tls.default_port()),
846 tls: *tls,
847 username: username.as_deref(),
848 password: password.as_deref(),
849 from,
850 to,
851 subject: &msg.title(),
852 body: &body,
853 date: now_ms() / 1000,
854 message_id: &mid,
855 },
856 )
857 .map(|_| None)
858 .map_err(|e| (None, e, None));
859 }
860 let req = provider::request(&ch.provider, msg, &secret).map_err(|e| (None, e, None))?;
861 let r = net::post(&net, &req).map_err(|e| (None, e, None))?;
862 if (200..300).contains(&r.status) {
863 return Ok(Some(r.status));
864 }
865 let retryable = r.status == 429 || r.status == 408 || r.status >= 500;
866 let what = if (300..400).contains(&r.status) {
867 "a redirect (not followed)".to_string()
868 } else {
869 let b: String = r.body.chars().take(200).collect();
870 format!(
871 "HTTP {}{}",
872 r.status,
873 if b.is_empty() {
874 String::new()
875 } else {
876 format!(": {b}")
877 }
878 )
879 };
880 Err((
881 Some(r.status),
882 SendError {
883 message: format!("{} answered {what}", ch.provider.kind()),
884 retryable,
885 },
886 r.retry_after,
887 ))
888 }
889
890 pub fn test(&self, org: &OrgId, name: &str, by: &str) -> Result<Delivery> {
892 let ch = self.get(org, name)?;
893 let msg = Message {
894 id: self.new_id(),
895 org: org.to_string(),
896 kind: "test".into(),
897 level: "info".into(),
898 stack: "-".into(),
899 service: String::new(),
900 project: None,
901 instance: None,
902 message: format!("A test notification from isb for channel {name}, sent by {by}."),
903 details: None,
904 at: now_ms(),
905 seq: 0,
906 test: true,
907 };
908 let mut d = Self::delivery(&msg, "sent");
909 d.attempts = 1;
910 match self.send(org, &ch, &msg) {
911 Ok(s) => d.http_status = s,
912 Err((s, e, _)) => {
913 d.status = "failed".into();
914 d.http_status = s;
915 d.error = Some(e.message);
916 }
917 }
918 d.finished_at = Some(now_ms());
919 self.record(org, name, &d);
920 Ok(d)
921 }
922}
923
924#[cfg(test)]
925mod tests {
926 use super::*;
927 use std::io::{Read, Write};
928 use std::net::TcpListener;
929
930 #[test]
931 fn globs() {
932 for (p, s) in [
933 ("*", "deploy.failed"),
934 ("deploy.*", "deploy.failed"),
935 ("*.failed", "backup.failed"),
936 ("deploy.failed", "deploy.failed"),
937 ("*.*", "a.b"),
938 ("d*y.*d", "deploy.failed"),
939 ] {
940 assert!(glob(p, s), "{p} {s}");
941 }
942 for (p, s) in [
943 ("deploy.*", "health.unhealthy"),
944 ("*.failed", "deploy.succeeded"),
945 ("deploy", "deploy.failed"),
946 ("deploy.failed", "deploy.failed2"),
947 ] {
948 assert!(!glob(p, s), "{p} {s}");
949 }
950 }
951
952 #[test]
953 fn rules() {
954 let s = Subject {
955 kind: "deploy.failed",
956 stack: "shop-production",
957 service: "web",
958 project: Some("shop"),
959 };
960 assert!(Rule::default().matches(&s));
961 let r = |j: serde_json::Value| -> Rule { serde_json::from_value(j).unwrap() };
962 assert!(r(serde_json::json!({"events": ["deploy.*"]})).matches(&s));
963 assert!(!r(serde_json::json!({"events": ["health.*"]})).matches(&s));
964 assert!(r(serde_json::json!({"events": ["*"], "projects": ["shop"]})).matches(&s));
965 assert!(!r(serde_json::json!({"events": ["*"], "projects": ["blog"]})).matches(&s));
966 assert!(r(serde_json::json!({"apps": ["web", "api"]})).matches(&s));
967 assert!(!r(serde_json::json!({"apps": ["api"]})).matches(&s));
968 assert!(r(serde_json::json!({"stacks": ["shop-production"]})).matches(&s));
969 assert!(!r(serde_json::json!({"stacks": ["shop-staging"]})).matches(&s));
970 let plain = Subject {
972 project: None,
973 ..s.clone()
974 };
975 assert!(!r(serde_json::json!({"projects": ["shop"]})).matches(&plain));
976 assert!(!r(serde_json::json!({"events": ["deploy.*"], "apps": ["api"]})).matches(&s));
978 assert!(serde_json::from_value::<Rule>(serde_json::json!({"evnts": ["x"]})).is_err());
979 assert!(
980 r(serde_json::json!({"events": ["Deploy"]}))
981 .validate()
982 .is_err()
983 );
984 assert!(r(serde_json::json!({"events": []})).validate().is_err());
985 let ch = Channel {
987 name: "ops".into(),
988 provider: Provider::Slack {
989 url_secret: "S".into(),
990 },
991 enabled: false,
992 rules: default_rules(),
993 created_at: 0,
994 updated_at: 0,
995 };
996 assert!(!ch.matches(&s));
997 }
998
999 #[test]
1000 fn event_orgs() {
1001 let (o, s) = event_org("acme/shop-production");
1002 assert_eq!((o.as_str(), s), ("acme", "shop-production"));
1003 let (o, s) = event_org("web");
1004 assert_eq!((o.as_str(), s), ("default", "web"));
1005 }
1006
1007 #[test]
1008 fn backoff_and_rate() {
1009 let b = Duration::from_secs(5);
1010 assert_eq!(backoff(b, 1, None), Duration::from_secs(5));
1011 assert_eq!(backoff(b, 2, None), Duration::from_secs(10));
1012 assert_eq!(backoff(b, 5, None), Duration::from_secs(80));
1013 assert_eq!(backoff(b, 30, None), BACKOFF_MAX);
1014 assert_eq!(
1015 backoff(b, 1, Some(Duration::from_secs(42))),
1016 Duration::from_secs(42)
1017 );
1018 assert_eq!(backoff(b, 1, Some(Duration::from_secs(9999))), BACKOFF_MAX);
1019 let mut r = RateLimit::default();
1020 let t0 = Instant::now();
1021 for i in 0..3 {
1022 assert!(r.wait(t0 + Duration::from_secs(i), 3).is_zero());
1023 r.record(t0 + Duration::from_secs(i));
1024 }
1025 assert_eq!(
1026 r.wait(t0 + Duration::from_secs(10), 3),
1027 Duration::from_secs(50)
1028 );
1029 assert!(r.wait(t0 + Duration::from_secs(60), 3).is_zero());
1030 }
1031
1032 #[expect(
1034 clippy::excessive_nesting,
1035 reason = "predates the lint ratchet; split it when next changed"
1036 )]
1037 fn receiver(statuses: Vec<u16>) -> (u16, std::thread::JoinHandle<Vec<String>>) {
1038 let l = TcpListener::bind("127.0.0.1:0").unwrap();
1039 let port = l.local_addr().unwrap().port();
1040 let h = std::thread::spawn(move || {
1041 let mut got = Vec::new();
1042 for st in statuses {
1043 let (mut s, _) = l.accept().unwrap();
1044 let mut buf = Vec::new();
1045 let mut b = [0u8; 4096];
1046 loop {
1047 let n = s.read(&mut b).unwrap();
1048 buf.extend_from_slice(&b[..n]);
1049 let t = String::from_utf8_lossy(&buf).to_string();
1050 if let Some(i) = t.find("\r\n\r\n") {
1051 let len: usize = t
1052 .lines()
1053 .find_map(|l| l.strip_prefix("Content-Length: "))
1054 .and_then(|v| v.trim().parse().ok())
1055 .unwrap_or(0);
1056 if buf.len() >= i + 4 + len {
1057 break;
1058 }
1059 }
1060 }
1061 got.push(String::from_utf8_lossy(&buf).to_string());
1062 write!(s, "HTTP/1.1 {st} X\r\nContent-Length: 0\r\n\r\n").unwrap();
1063 }
1064 got
1065 });
1066 (port, h)
1067 }
1068
1069 fn notifier(dir: &Path) -> (Notifier, Arc<Secrets>, OrgId) {
1070 let keyring = Arc::new(crate::secrets::keys::Keyring::new(
1071 age::x25519::Identity::generate(),
1072 vec![],
1073 ));
1074 let secrets = Arc::new(Secrets::new(crate::secrets::local::LocalDriver::new(
1075 dir, keyring,
1076 )));
1077 let mut n = Notifier::new(
1078 dir,
1079 secrets.clone(),
1080 Arc::new(|_, _, _| Some("shop".into())),
1081 )
1082 .unwrap();
1083 let inner = Arc::get_mut(&mut n.inner).unwrap();
1084 inner.backoff = Duration::from_millis(20);
1085 (n, secrets, OrgId::new("acme").unwrap())
1086 }
1087
1088 #[test]
1089 fn retries_then_delivers_with_a_signature() {
1090 let dir = tempfile::tempdir().unwrap();
1091 let (n, secrets, org) = notifier(dir.path());
1092 let (port, h) = receiver(vec![503, 500, 200]);
1093 secrets
1094 .create(
1095 &org,
1096 "HOOK",
1097 None,
1098 format!("http://127.0.0.1:{port}/in").as_bytes(),
1099 &Default::default(),
1100 )
1101 .unwrap();
1102 secrets
1103 .create(&org, "SIGN", None, b"k3y", &Default::default())
1104 .unwrap();
1105 let ch = Channel {
1106 name: "ops".into(),
1107 provider: Provider::Webhook {
1108 url_secret: "HOOK".into(),
1109 signing_secret: Some("SIGN".into()),
1110 },
1111 enabled: true,
1112 rules: vec![
1113 serde_json::from_value(
1114 serde_json::json!({"events": ["deploy.*"], "projects": ["shop"]}),
1115 )
1116 .unwrap(),
1117 ],
1118 created_at: 0,
1119 updated_at: 0,
1120 };
1121 n.create(&org, ch.clone()).unwrap();
1123 let d = n.test(&org, "ops", "tester").unwrap();
1124 assert_eq!(d.status, "failed");
1125 assert!(d.error.unwrap().contains("private targets are off"));
1126 n.set_settings(Settings {
1127 allow_private_targets: true,
1128 })
1129 .unwrap();
1130 let ev = |kind: &str, seq| Event {
1131 seq,
1132 at: 1,
1133 level: "error".into(),
1134 stack: "acme/shop-production".into(),
1135 service: "web".into(),
1136 instance: None,
1137 message: "deployment 3 failed".into(),
1138 kind: Some(kind.into()),
1139 };
1140 n.route(&ev("health.unhealthy", 1));
1142 n.route(&ev("deploy.failed", 2));
1143 let got = h.join().unwrap();
1144 assert_eq!(got.len(), 3);
1145 let req = &got[2];
1146 let body = &req[req.find("\r\n\r\n").unwrap() + 4..];
1147 let sig = req
1148 .lines()
1149 .find_map(|l| l.strip_prefix("X-Isb-Signature: "))
1150 .unwrap();
1151 assert_eq!(
1152 sig,
1153 format!("sha256={}", net::hmac_sha256_hex(b"k3y", body.as_bytes()))
1154 );
1155 let v: serde_json::Value = serde_json::from_str(body).unwrap();
1156 assert_eq!(v["kind"], "deploy.failed");
1157 assert_eq!(v["project"], "shop");
1158 let deadline = Instant::now() + Duration::from_secs(5);
1160 loop {
1161 let log = n.deliveries(&org, "ops").unwrap();
1162 if let Some(d) = log
1163 .iter()
1164 .find(|d| d.kind == "deploy.failed" && d.status == "sent")
1165 {
1166 assert_eq!(d.attempts, 3);
1167 assert_eq!(d.http_status, Some(200));
1168 break;
1169 }
1170 assert!(Instant::now() < deadline, "{log:?}");
1171 std::thread::sleep(Duration::from_millis(20));
1172 }
1173 let (n2, _, _) = notifier(dir.path());
1175 assert!(n2.deliveries(&org, "ops").unwrap().len() >= 2);
1176 n.shutdown();
1177 }
1178
1179 #[test]
1180 fn permanent_failures_are_not_retried() {
1181 let dir = tempfile::tempdir().unwrap();
1182 let (n, secrets, org) = notifier(dir.path());
1183 n.set_settings(Settings {
1184 allow_private_targets: true,
1185 })
1186 .unwrap();
1187 let (port, h) = receiver(vec![404]);
1188 secrets
1189 .create(
1190 &org,
1191 "HOOK",
1192 None,
1193 format!("http://127.0.0.1:{port}/").as_bytes(),
1194 &Default::default(),
1195 )
1196 .unwrap();
1197 n.create(
1198 &org,
1199 Channel {
1200 name: "x".into(),
1201 provider: Provider::Webhook {
1202 url_secret: "HOOK".into(),
1203 signing_secret: None,
1204 },
1205 enabled: true,
1206 rules: default_rules(),
1207 created_at: 0,
1208 updated_at: 0,
1209 },
1210 )
1211 .unwrap();
1212 let d = n.test(&org, "x", "t").unwrap();
1213 assert_eq!((d.status.as_str(), d.http_status), ("failed", Some(404)));
1214 h.join().unwrap();
1215 }
1216
1217 #[test]
1218 fn channels_check_their_secrets() {
1219 let dir = tempfile::tempdir().unwrap();
1220 let (n, secrets, org) = notifier(dir.path());
1221 let slack = |s: &str| Channel {
1222 name: "s".into(),
1223 provider: Provider::Slack {
1224 url_secret: s.into(),
1225 },
1226 enabled: true,
1227 rules: default_rules(),
1228 created_at: 0,
1229 updated_at: 0,
1230 };
1231 let e = n.create(&org, slack("NOPE")).unwrap_err();
1232 assert!(e.to_string().contains("no secret NOPE"), "{e}");
1233 secrets
1234 .create(
1235 &org,
1236 "BAD",
1237 None,
1238 b"https://example.com/services/TOKEN",
1239 &Default::default(),
1240 )
1241 .unwrap();
1242 let e = n.create(&org, slack("BAD")).unwrap_err().to_string();
1243 assert!(e.contains("hooks.slack.com") && !e.contains("TOKEN"), "{e}");
1244 secrets
1245 .create(
1246 &org,
1247 "GOOD",
1248 None,
1249 b"https://hooks.slack.com/services/T/B/x\n",
1250 &Default::default(),
1251 )
1252 .unwrap();
1253 n.create(&org, slack("GOOD")).unwrap();
1254 assert!(n.create(&org, slack("GOOD")).is_err(), "duplicate");
1255 let c = n.update(&org, "s", None, None, Some(false)).unwrap();
1256 assert!(!c.enabled);
1257 assert_eq!(n.list(&org).unwrap().len(), 1);
1258 assert!(n.list(&OrgId::new("beta").unwrap()).unwrap().is_empty());
1260 n.delete(&org, "s").unwrap();
1261 assert!(n.get(&org, "s").is_err());
1262 }
1263}