use async_trait::async_trait;
use serde_json::Value;
use tracing::{info, warn};
use super::{Notifiers, NotifyEvent};
use crate::jobs::{JobHandler, JobOutcome};
use crate::sqlite::job::Job;
pub const NOTIFY_JOB_KIND: &str = "notify_deliver";
pub struct NotifyJob(Notifiers);
impl NotifyJob {
#[must_use]
pub fn new(dispatchers: impl Into<Notifiers>) -> Self {
Self(dispatchers.into())
}
}
struct Delivery {
profile: String,
backend: String,
event: NotifyEvent,
}
impl Delivery {
fn parse(payload: &Value) -> Result<Self, String> {
let profile = payload
.get("profile")
.and_then(Value::as_str)
.ok_or("the job payload names no profile")?
.to_string();
let backend = payload
.get("backend")
.and_then(Value::as_str)
.ok_or("the job payload names no backend")?
.to_string();
let event = payload
.get("event")
.ok_or("the job payload carries no event")?;
let event: NotifyEvent = serde_json::from_value(event.clone())
.map_err(|error| format!("the job payload's event does not parse: {error}"))?;
Ok(Self {
profile,
backend,
event,
})
}
}
#[async_trait]
impl JobHandler for NotifyJob {
fn kind(&self) -> &'static str {
NOTIFY_JOB_KIND
}
async fn run(&self, job: &Job) -> JobOutcome {
let delivery = match Delivery::parse(&job.payload) {
Ok(delivery) => delivery,
Err(reason) => return JobOutcome::Failed(reason),
};
let Some(dispatcher) = self.0.get(&delivery.profile) else {
return JobOutcome::Failed(format!(
"no profile `{}` is mounted, so its notification cannot be delivered",
delivery.profile
));
};
let kind = delivery.event.kind();
match dispatcher.deliver(&delivery.backend, &delivery.event).await {
Some(Ok(())) => {
info!(
event = "notify_delivered",
outcome = "success",
profile = %delivery.profile,
backend = %delivery.backend,
kind,
attempt = job.attempts,
);
JobOutcome::Done
}
Some(Err(error)) => {
warn!(
event = "notify_delivery_failed",
outcome = "failure",
profile = %delivery.profile,
backend = %delivery.backend,
kind,
attempt = job.attempts,
retryable = error.retryable(),
error = %error,
);
if error.retryable() {
JobOutcome::Retry(error.to_string())
} else {
JobOutcome::Failed(error.to_string())
}
}
None => JobOutcome::Failed(format!(
"profile `{}` has no notify backend `{}` configured any more",
delivery.profile, delivery.backend
)),
}
}
async fn abandon(&self, job: &Job, reason: &str) {
let (profile, backend, kind) = match Delivery::parse(&job.payload) {
Ok(delivery) => (
delivery.profile,
delivery.backend,
delivery.event.kind().to_string(),
),
Err(_) => (String::new(), String::new(), String::new()),
};
warn!(
event = "notify_delivery_abandoned",
outcome = "failure",
profile = %profile,
backend = %backend,
kind = %kind,
attempts = job.attempts,
reason = %reason,
);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::ALL_NOTIFY_EVENTS;
use crate::jobs::{JobQueue, JobSpec};
use crate::notify::tests::RecordingNotifyBackend;
use crate::notify::{BackendSlot, NotifyDispatcher, NotifyError, ProfileMountedData};
use crate::sqlite::db::Database;
use serde_json::json;
use std::collections::HashMap;
use std::sync::Arc;
async fn queue() -> JobQueue {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
JobQueue::new(database, &crate::config::JobsConfig::default())
}
fn every_kind() -> Vec<String> {
ALL_NOTIFY_EVENTS.iter().map(|k| (*k).to_string()).collect()
}
fn mounted(profile: &str) -> NotifyEvent {
NotifyEvent::ProfileMounted(ProfileMountedData {
profile: profile.to_string(),
})
}
fn handler(
profile: &str,
recorder: Arc<RecordingNotifyBackend>,
queue: JobQueue,
) -> (NotifyJob, Arc<NotifyDispatcher>) {
let dispatcher = Arc::new(NotifyDispatcher::new(
profile,
vec![BackendSlot::new("recording", recorder, &every_kind())],
queue,
));
let mut map = HashMap::new();
map.insert(profile.to_string(), dispatcher.clone());
(NotifyJob::new(Arc::new(map)), dispatcher)
}
fn row(payload: Value) -> Job {
Job {
id: crate::sqlite::id::mint(),
kind: NOTIFY_JOB_KIND.to_string(),
dedup_key: "k".to_string(),
payload,
status: "running".to_string(),
run_at: 0,
attempts: 1,
max_attempts: 5,
deadline: None,
lease_until: None,
lease_owner: None,
last_error: None,
created_at: 0,
updated_at: 0,
}
}
fn delivery(profile: &str, backend: &str, event: &NotifyEvent) -> Value {
json!({ "profile": profile, "backend": backend, "event": event.payload() })
}
#[tokio::test]
async fn a_delivered_notification_is_done() {
let recorder = Arc::new(RecordingNotifyBackend::default());
let (job, _dispatcher) = handler("le", recorder.clone(), queue().await);
let outcome = job
.run(&row(delivery("le", "recording", &mounted("le"))))
.await;
assert!(matches!(outcome, JobOutcome::Done), "{outcome:?}");
assert_eq!(recorder.events.lock().unwrap().len(), 1);
}
#[tokio::test]
async fn a_retryable_failure_retries_and_a_permanent_one_does_not() {
let queue = queue().await;
let flaky = Arc::new(RecordingNotifyBackend::failing());
let (job, _d) = handler("le", flaky, queue.clone());
let outcome = job
.run(&row(delivery("le", "recording", &mounted("le"))))
.await;
assert!(matches!(outcome, JobOutcome::Retry(_)), "{outcome:?}");
let broken = Arc::new(RecordingNotifyBackend::failing_permanently());
let (job, _d) = handler("le", broken, queue);
let outcome = job
.run(&row(delivery("le", "recording", &mounted("le"))))
.await;
assert!(matches!(outcome, JobOutcome::Failed(_)), "{outcome:?}");
}
#[tokio::test]
async fn an_unknown_profile_or_backend_is_retired_rather_than_retried() {
let recorder = Arc::new(RecordingNotifyBackend::default());
let (job, _d) = handler("le", recorder.clone(), queue().await);
let outcome = job
.run(&row(delivery("staging", "recording", &mounted("le"))))
.await;
match outcome {
JobOutcome::Failed(reason) => assert!(reason.contains("staging"), "{reason}"),
other => panic!("{other:?}"),
}
let outcome = job
.run(&row(delivery("le", "carrier-pigeon", &mounted("le"))))
.await;
match outcome {
JobOutcome::Failed(reason) => assert!(reason.contains("carrier-pigeon"), "{reason}"),
other => panic!("{other:?}"),
}
assert!(
recorder.events.lock().unwrap().is_empty(),
"neither case may reach a backend"
);
}
#[tokio::test]
async fn a_payload_that_cannot_be_read_is_never_retried() {
let recorder = Arc::new(RecordingNotifyBackend::default());
let (job, _d) = handler("le", recorder, queue().await);
for payload in [
json!({ "backend": "recording", "event": mounted("le").payload() }),
json!({ "profile": "le", "event": mounted("le").payload() }),
json!({ "profile": "le", "backend": "recording" }),
json!({ "profile": "le", "backend": "recording", "event": {"hook": "not_an_event"} }),
] {
let outcome = job.run(&row(payload.clone())).await;
assert!(
matches!(outcome, JobOutcome::Failed(_)),
"{payload}: {outcome:?}"
);
}
}
#[tokio::test]
async fn abandon_survives_the_payload_that_retired_the_job() {
let recorder = Arc::new(RecordingNotifyBackend::default());
let (job, _d) = handler("le", recorder, queue().await);
job.abandon(&row(json!({})), "the payload names no profile")
.await;
job.abandon(
&row(delivery("le", "recording", &mounted("le"))),
"the attempts ran out",
)
.await;
}
#[tokio::test]
async fn recover_queues_nothing() {
let queue = queue().await;
let recorder = Arc::new(RecordingNotifyBackend::default());
let (job, _d) = handler("le", recorder, queue.clone());
job.recover(&queue).await;
assert!(
Job::find_live(NOTIFY_JOB_KIND, "k", queue.database())
.await
.unwrap()
.is_none()
);
}
#[test]
fn the_error_split_survives_the_round_trip() {
assert!(NotifyError::new("connection refused").retryable());
assert!(!NotifyError::permanent("template missing").retryable());
}
#[tokio::test]
async fn the_queued_spec_is_addressed_at_one_backend() {
let queue = queue().await;
let spec = JobSpec::now(NOTIFY_JOB_KIND, "delivery-1:email");
assert!(queue.enqueue(spec).await.unwrap());
assert!(
Job::find_live(NOTIFY_JOB_KIND, "delivery-1:email", queue.database())
.await
.unwrap()
.is_some()
);
}
}