1use 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
17pub const NOTIFY_JOB_KIND: &str = "notify_deliver";
19
20pub 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
40struct Delivery {
42 profile: String,
43 backend: String,
44 event: NotifyEvent,
45}
46
47impl Delivery {
48 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 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 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 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 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 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 #[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 #[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 #[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 #[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 #[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 #[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}