Skip to main content

acme_proxy/notify/
job.rs

1//! One queued notification delivery, as a durable job.
2//!
3//! The counterpart to [`NotifyDispatcher::dispatch`](super::NotifyDispatcher::dispatch):
4//! that side writes a row per (occurrence × backend), this side claims one and
5//! runs it. Lives beside its subsystem rather than in [`crate::jobs`], the way
6//! `signer::relay::flow::RelayJob` does — `jobs` owns the mechanism, a
7//! subsystem owns what its work means.
8
9use async_trait::async_trait;
10use serde_json::Value;
11use tracing::{info, warn};
12
13use super::{Notifiers, NotifyEvent};
14use crate::jobs::{JobHandler, JobOutcome};
15use crate::sqlite::job::Job;
16
17/// The `jobs.kind` one notification delivery is queued under.
18pub const NOTIFY_JOB_KIND: &str = "notify_deliver";
19
20/// Delivers queued notifications.
21///
22/// Holds the whole `profile name -> dispatcher` map rather than one dispatcher,
23/// because there is exactly one handler for a kind and a job row names its own
24/// profile. That is also why nothing deduplicates it at registration the way the
25/// signer handlers are deduplicated: there is only ever one of these.
26///
27/// The map arrives as a [`Notifiers`] handle rather than a plain `Arc` because
28/// this handler is registered once per generation but must see the *current*
29/// configuration: a row queued by a reloaded router names a slot id only the new
30/// map has, and an unknown one is retired rather than retried.
31pub struct NotifyJob(Notifiers);
32
33impl NotifyJob {
34    #[must_use]
35    pub fn new(dispatchers: impl Into<Notifiers>) -> Self {
36        Self(dispatchers.into())
37    }
38}
39
40/// What a claimed row says it wants delivered.
41struct Delivery {
42    profile: String,
43    backend: String,
44    event: NotifyEvent,
45}
46
47impl Delivery {
48    /// Reads a payload back. Every failure here is permanent by construction —
49    /// the bytes are already written and re-reading them cannot change them.
50    fn parse(payload: &Value) -> Result<Self, String> {
51        let profile = payload
52            .get("profile")
53            .and_then(Value::as_str)
54            .ok_or("the job payload names no profile")?
55            .to_string();
56        let backend = payload
57            .get("backend")
58            .and_then(Value::as_str)
59            .ok_or("the job payload names no backend")?
60            .to_string();
61        let event = payload
62            .get("event")
63            .ok_or("the job payload carries no event")?;
64        let event: NotifyEvent = serde_json::from_value(event.clone())
65            .map_err(|error| format!("the job payload's event does not parse: {error}"))?;
66        Ok(Self {
67            profile,
68            backend,
69            event,
70        })
71    }
72}
73
74#[async_trait]
75impl JobHandler for NotifyJob {
76    fn kind(&self) -> &'static str {
77        NOTIFY_JOB_KIND
78    }
79
80    /// One delivery attempt.
81    ///
82    /// Three things retire the job immediately rather than retrying, and they
83    /// share a shape: nothing about them can change between now and the fifth
84    /// attempt. A payload that does not parse never will; a profile or a backend
85    /// that is no longer configured is a configuration the operator changed
86    /// under a queued row, and re-reading it every thirty seconds until the
87    /// budget runs out would say nothing the first log line did not.
88    async fn run(&self, job: &Job) -> JobOutcome {
89        let delivery = match Delivery::parse(&job.payload) {
90            Ok(delivery) => delivery,
91            Err(reason) => return JobOutcome::Failed(reason),
92        };
93
94        let Some(dispatcher) = self.0.get(&delivery.profile) else {
95            return JobOutcome::Failed(format!(
96                "no profile `{}` is mounted, so its notification cannot be delivered",
97                delivery.profile
98            ));
99        };
100
101        let kind = delivery.event.kind();
102        match dispatcher.deliver(&delivery.backend, &delivery.event).await {
103            Some(Ok(())) => {
104                info!(
105                    event = "notify_delivered",
106                    outcome = "success",
107                    profile = %delivery.profile,
108                    backend = %delivery.backend,
109                    kind,
110                    attempt = job.attempts,
111                );
112                JobOutcome::Done
113            }
114            Some(Err(error)) => {
115                warn!(
116                    event = "notify_delivery_failed",
117                    outcome = "failure",
118                    profile = %delivery.profile,
119                    backend = %delivery.backend,
120                    kind,
121                    attempt = job.attempts,
122                    retryable = error.retryable(),
123                    error = %error,
124                );
125                if error.retryable() {
126                    JobOutcome::Retry(error.to_string())
127                } else {
128                    JobOutcome::Failed(error.to_string())
129                }
130            }
131            None => JobOutcome::Failed(format!(
132                "profile `{}` has no notify backend `{}` configured any more",
133                delivery.profile, delivery.backend
134            )),
135        }
136    }
137
138    /// The end of the line for one notification.
139    ///
140    /// There is nobody to tell, which is the whole difference from
141    /// `RelayJob::abandon`: a relay's subject is a client polling an order, and
142    /// this one's subject is the operator whose only channel is what just
143    /// failed. So this is a log line and nothing else — but it is the log line
144    /// to alert on, because it is the moment a notification is genuinely lost.
145    async fn abandon(&self, job: &Job, reason: &str) {
146        let (profile, backend, kind) = match Delivery::parse(&job.payload) {
147            Ok(delivery) => (
148                delivery.profile,
149                delivery.backend,
150                delivery.event.kind().to_string(),
151            ),
152            // Unparsable is itself one of the ways a job is retired here, so
153            // this arm is reachable and must still name what it can.
154            Err(_) => (String::new(), String::new(), String::new()),
155        };
156        warn!(
157            event = "notify_delivery_abandoned",
158            outcome = "failure",
159            profile = %profile,
160            backend = %backend,
161            kind = %kind,
162            attempts = job.attempts,
163            reason = %reason,
164        );
165    }
166}
167
168#[cfg(test)]
169mod tests {
170    use super::*;
171    use crate::config::ALL_NOTIFY_EVENTS;
172    use crate::jobs::{JobQueue, JobSpec};
173    use crate::notify::tests::RecordingNotifyBackend;
174    use crate::notify::{BackendSlot, NotifyDispatcher, NotifyError, ProfileMountedData};
175    use crate::sqlite::db::Database;
176    use serde_json::json;
177    use std::collections::HashMap;
178    use std::sync::Arc;
179
180    async fn queue() -> JobQueue {
181        let database = Arc::new(Database::connect_in_memory().await.unwrap());
182        JobQueue::new(database, &crate::config::JobsConfig::default())
183    }
184
185    fn every_kind() -> Vec<String> {
186        ALL_NOTIFY_EVENTS.iter().map(|k| (*k).to_string()).collect()
187    }
188
189    fn mounted(profile: &str) -> NotifyEvent {
190        NotifyEvent::ProfileMounted(ProfileMountedData {
191            profile: profile.to_string(),
192        })
193    }
194
195    /// A handler over one profile whose single backend is `recorder`.
196    fn handler(
197        profile: &str,
198        recorder: Arc<RecordingNotifyBackend>,
199        queue: JobQueue,
200    ) -> (NotifyJob, Arc<NotifyDispatcher>) {
201        let dispatcher = Arc::new(NotifyDispatcher::new(
202            profile,
203            vec![BackendSlot::new("recording", recorder, &every_kind())],
204            queue,
205        ));
206        let mut map = HashMap::new();
207        map.insert(profile.to_string(), dispatcher.clone());
208        (NotifyJob::new(Arc::new(map)), dispatcher)
209    }
210
211    /// A job row carrying `payload`, without going through the queue — the
212    /// handler only ever reads `payload` and `attempts`.
213    fn row(payload: Value) -> Job {
214        Job {
215            id: "job-1".to_string(),
216            kind: NOTIFY_JOB_KIND.to_string(),
217            dedup_key: "k".to_string(),
218            payload,
219            status: "running".to_string(),
220            run_at: 0,
221            attempts: 1,
222            max_attempts: 5,
223            deadline: None,
224            lease_until: None,
225            lease_owner: None,
226            last_error: None,
227            created_at: 0,
228            updated_at: 0,
229        }
230    }
231
232    fn delivery(profile: &str, backend: &str, event: &NotifyEvent) -> Value {
233        json!({ "profile": profile, "backend": backend, "event": event.payload() })
234    }
235
236    #[tokio::test]
237    async fn a_delivered_notification_is_done() {
238        let recorder = Arc::new(RecordingNotifyBackend::default());
239        let (job, _dispatcher) = handler("le", recorder.clone(), queue().await);
240
241        let outcome = job
242            .run(&row(delivery("le", "recording", &mounted("le"))))
243            .await;
244
245        assert!(matches!(outcome, JobOutcome::Done), "{outcome:?}");
246        assert_eq!(recorder.events.lock().unwrap().len(), 1);
247    }
248
249    /// The distinction the whole queue rests on, at this subsystem's boundary:
250    /// a refused connection is asked again, a template that will never render
251    /// is not.
252    #[tokio::test]
253    async fn a_retryable_failure_retries_and_a_permanent_one_does_not() {
254        let queue = queue().await;
255
256        let flaky = Arc::new(RecordingNotifyBackend::failing());
257        let (job, _d) = handler("le", flaky, queue.clone());
258        let outcome = job
259            .run(&row(delivery("le", "recording", &mounted("le"))))
260            .await;
261        assert!(matches!(outcome, JobOutcome::Retry(_)), "{outcome:?}");
262
263        let broken = Arc::new(RecordingNotifyBackend::failing_permanently());
264        let (job, _d) = handler("le", broken, queue);
265        let outcome = job
266            .run(&row(delivery("le", "recording", &mounted("le"))))
267            .await;
268        assert!(matches!(outcome, JobOutcome::Failed(_)), "{outcome:?}");
269    }
270
271    /// A configuration change under a queued row. Both of these would otherwise
272    /// burn the whole retry budget re-reading a map that cannot change while the
273    /// process lives.
274    #[tokio::test]
275    async fn an_unknown_profile_or_backend_is_retired_rather_than_retried() {
276        let recorder = Arc::new(RecordingNotifyBackend::default());
277        let (job, _d) = handler("le", recorder.clone(), queue().await);
278
279        let outcome = job
280            .run(&row(delivery("staging", "recording", &mounted("le"))))
281            .await;
282        match outcome {
283            JobOutcome::Failed(reason) => assert!(reason.contains("staging"), "{reason}"),
284            other => panic!("{other:?}"),
285        }
286
287        let outcome = job
288            .run(&row(delivery("le", "carrier-pigeon", &mounted("le"))))
289            .await;
290        match outcome {
291            JobOutcome::Failed(reason) => assert!(reason.contains("carrier-pigeon"), "{reason}"),
292            other => panic!("{other:?}"),
293        }
294
295        assert!(
296            recorder.events.lock().unwrap().is_empty(),
297            "neither case may reach a backend"
298        );
299    }
300
301    /// Three shapes of payload that can never be delivered. Written out rather
302    /// than folded into one case because each is a different missing member, and
303    /// a `parse` that stopped checking one of them would still pass the others.
304    #[tokio::test]
305    async fn a_payload_that_cannot_be_read_is_never_retried() {
306        let recorder = Arc::new(RecordingNotifyBackend::default());
307        let (job, _d) = handler("le", recorder, queue().await);
308
309        for payload in [
310            json!({ "backend": "recording", "event": mounted("le").payload() }),
311            json!({ "profile": "le", "event": mounted("le").payload() }),
312            json!({ "profile": "le", "backend": "recording" }),
313            json!({ "profile": "le", "backend": "recording", "event": {"hook": "not_an_event"} }),
314        ] {
315            let outcome = job.run(&row(payload.clone())).await;
316            assert!(
317                matches!(outcome, JobOutcome::Failed(_)),
318                "{payload}: {outcome:?}"
319            );
320        }
321    }
322
323    /// `abandon` must not panic on the payload that got the job retired in the
324    /// first place — the unparsable one is reachable from `run`'s own `Failed`.
325    #[tokio::test]
326    async fn abandon_survives_the_payload_that_retired_the_job() {
327        let recorder = Arc::new(RecordingNotifyBackend::default());
328        let (job, _d) = handler("le", recorder, queue().await);
329
330        job.abandon(&row(json!({})), "the payload names no profile")
331            .await;
332        job.abandon(
333            &row(delivery("le", "recording", &mounted("le"))),
334            "the attempts ran out",
335        )
336        .await;
337    }
338
339    /// The handler never enqueues anything of its own: a queued row is already
340    /// durable, so there is no external state for a restart to re-derive.
341    #[tokio::test]
342    async fn recover_queues_nothing() {
343        let queue = queue().await;
344        let recorder = Arc::new(RecordingNotifyBackend::default());
345        let (job, _d) = handler("le", recorder, queue.clone());
346
347        job.recover(&queue).await;
348
349        assert!(
350            Job::find_live(NOTIFY_JOB_KIND, "k", queue.database())
351                .await
352                .unwrap()
353                .is_none()
354        );
355    }
356
357    /// A permanent error reported by a backend must not be laundered into a
358    /// retry on the way through the handler.
359    #[test]
360    fn the_error_split_survives_the_round_trip() {
361        assert!(NotifyError::new("connection refused").retryable());
362        assert!(!NotifyError::permanent("template missing").retryable());
363    }
364
365    #[tokio::test]
366    async fn the_queued_spec_is_addressed_at_one_backend() {
367        let queue = queue().await;
368        let spec = JobSpec::now(NOTIFY_JOB_KIND, "delivery-1:email");
369        assert!(queue.enqueue(spec).await.unwrap());
370        assert!(
371            Job::find_live(NOTIFY_JOB_KIND, "delivery-1:email", queue.database())
372                .await
373                .unwrap()
374                .is_some()
375        );
376    }
377}