Skip to main content

isb_apps/notify/
mod.rs

1//! Notifications: per-org channels (webhook, Slack, Discord, Telegram,
2//! email) told about events of chosen kinds (`deploy.*`, `health.*`,
3//! `backup.*`, `job.*`, `cert.*`, `monitor.*`; see
4//! [`crate::stack::controller::Event`]).
5//!
6//! A dispatcher thread follows the controller's event feed and queues a
7//! delivery per matching channel. Each channel has its own bounded queue and
8//! sender thread, so a slow or failing channel never holds up another: a
9//! delivery is retried with exponential backoff, a channel is rate limited,
10//! and the last deliveries per channel are kept with their outcome.
11//!
12//! The feed lives in memory and starts at 1 with each daemon run, so the
13//! dispatcher starts from the beginning of the run: nothing from before a
14//! restart is sent again (and deliveries still queued at a restart are lost).
15//!
16//! Channels and delivery logs are kept under `<state>/orgs/<org>/notify/`.
17//! Destinations are held to [`net`]'s address policy.
18
19pub 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
40/// Deliveries waiting per channel; the oldest is dropped past this.
41pub const QUEUE_MAX: usize = 100;
42/// Deliveries kept in a channel's log.
43pub const LOG_KEPT: usize = 50;
44/// Attempts per delivery, the first included.
45pub const ATTEMPTS: u32 = 6;
46/// Messages a channel sends per minute at most.
47pub const PER_MINUTE: usize = 20;
48/// The first retry's wait; it doubles per attempt.
49const BACKOFF: Duration = Duration::from_secs(5);
50/// A channel's sender thread exits after this long with nothing to send.
51const IDLE: Duration = Duration::from_secs(60);
52
53/// Which events a channel hears about. An event matches a rule when its kind
54/// matches one of `events` (globs: `deploy.*`, `*.failed`, `*`) and it
55/// passes every non-empty filter.
56#[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    /// App projects (an app's events only).
62    #[serde(default, skip_serializing_if = "Vec::is_empty")]
63    pub projects: Vec<String>,
64    /// Apps (a service of the same name in the app's stack).
65    #[serde(default, skip_serializing_if = "Vec::is_empty")]
66    pub apps: Vec<String>,
67    /// Stacks, by their name in the org.
68    #[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
87/// `*` matches any run of characters, everything else itself.
88pub 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/// What a rule is matched against.
112#[derive(Debug, Clone, Default)]
113pub struct Subject<'a> {
114    pub kind: &'a str,
115    /// The stack's name in its org.
116    pub stack: &'a str,
117    pub service: &'a str,
118    /// The app's project, when the service is an app.
119    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/// A notification channel of an org.
153#[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    /// Any rule matching is enough. Default: every event with a kind.
161    #[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
195/// A channel name: [a-z0-9-], starts with a letter, at most 40.
196pub 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/// One delivery to one channel, as its log shows it.
212#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
213pub struct Delivery {
214    pub id: String,
215    /// Unix milliseconds it was queued.
216    pub at: u64,
217    pub kind: String,
218    pub seq: u64,
219    pub summary: String,
220    /// `queued`, `retrying`, `sent`, `failed`, `dropped` (queue full) or
221    /// `skipped` (channel disabled or gone).
222    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/// Server-wide settings (platform admins).
235#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
236pub struct Settings {
237    /// Let channels reach loopback, private and other non-public addresses.
238    #[serde(default)]
239    pub allow_private_targets: bool,
240}
241
242/// The app project a service belongs to: `(org, stack, service)`.
243pub type Resolve = Arc<dyn Fn(&OrgId, &str, &str) -> Option<String> + Send + Sync>;
244
245/// Structured details of an event, from its producer (a monitor's URL,
246/// status and latency), for [`Message::details`].
247pub type Details = Arc<dyn Fn(&OrgId, &Event) -> Option<serde_json::Value> + Send + Sync>;
248
249/// The org of an event and its stack's own name. Events name their stack
250/// `org/stack`, or just `stack` in the default org.
251pub 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    /// Held across read-modify-write of a channels file.
274    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    /// Trust these TLS roots instead of the public set (tests).
284    tls: Option<Arc<rustls::ClientConfig>>,
285}
286
287/// The notification service of a daemon.
288#[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    /// Open the notifier over a state directory. Nothing is sent until
309    /// [`Notifier::start`].
310    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    /// Follow the controller's events from the start of this run.
337    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                    // The feed returns everything after `since` (it keeps
349                    // fewer than 1000), so `seq` is where to go on from.
350                    since = since.max(seq);
351                }
352            });
353    }
354
355    /// Where events' details come from.
356    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    /// The org's channels.
393    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    /// Check a channel against the org's secrets: they must exist, and URL
416    /// secrets must hold a URL fit for the provider. Values are never shown.
417    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    /// Replace a channel's provider, rules or enabled flag (each optional).
465    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    /// A channel's recent deliveries, newest first.
522    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    /// Insert or replace a delivery in its channel's log, and persist it.
545    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    /// Queue a delivery to every channel of the event's org that wants it.
572    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    /// A channel's sender: one delivery at a time, in order, until idle.
683    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    /// Sleep in short steps so a shutdown is not held up.
715    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    /// Send one queued delivery, retrying what may succeed later.
727    fn deliver(&self, job: &Job) {
728        let mut d = Self::delivery(&job.msg, "queued");
729        // The delivery was recorded at enqueue time: keep its id and time.
730        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            // Read the channel each time: it may have been edited or removed.
742            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    /// One attempt. Ok: the HTTP status, if HTTP. Err: the HTTP status (if
799    /// any), why, and how long the server asked us to wait.
800    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    /// Send a test message to a channel now, once, and log it.
891    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        // A project filter never matches a plain stack's events.
971        let plain = Subject {
972            project: None,
973            ..s.clone()
974        };
975        assert!(!r(serde_json::json!({"projects": ["shop"]})).matches(&plain));
976        // All filters must hold.
977        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        // A disabled channel hears nothing.
986        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    /// A fake HTTP receiver answering each request with the next status.
1033    #[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        // Loopback is refused until a platform admin allows private targets.
1122        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        // Not matching: no delivery.
1141        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        // The log shows the delivery, sent on the third attempt.
1159        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        // The log survives a restart.
1174        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        // Another org sees none of it.
1259        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}