#[cfg(feature = "events")]
pub use events::Events;
#[cfg(feature = "jobs")]
pub use jobs::{JobOutcome, JobRecord, Jobs};
#[cfg(feature = "mail")]
pub use mail::{Mail, SentMail};
#[cfg(feature = "jobs")]
mod jobs {
use std::sync::{Arc, Mutex};
use crate::jobs::{Event, FailReason, Observer};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct JobRecord {
pub job_id: uuid::Uuid,
pub kind: String,
pub version: i16,
pub attempt: i32,
pub outcome: JobOutcome,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum JobOutcome {
Started,
Succeeded,
Retried,
Failed(FailReason),
}
impl std::fmt::Display for JobOutcome {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Started => f.write_str("started"),
Self::Succeeded => f.write_str("succeeded"),
Self::Retried => f.write_str("retried"),
Self::Failed(reason) => write!(f, "failed ({reason})"),
}
}
}
#[derive(Debug, Clone, Default)]
pub struct Jobs {
records: Arc<Mutex<Vec<JobRecord>>>,
}
impl Jobs {
#[must_use]
pub fn recorder() -> Self {
Self::default()
}
#[must_use]
pub fn records(&self) -> Vec<JobRecord> {
self.records
.lock()
.expect("job record lock was poisoned by an earlier panic")
.clone()
}
pub fn assert_job_ran(&self, kind: &str) -> &Self {
let records = self.records();
assert!(
records.iter().any(|record| record.kind == kind),
"no job of kind `{kind}` ran; {}",
summarise(&records)
);
self
}
pub fn assert_job_succeeded(&self, kind: &str) -> &Self {
let records = self.records();
assert!(
records
.iter()
.any(|r| r.kind == kind && r.outcome == JobOutcome::Succeeded),
"no job of kind `{kind}` succeeded; {}",
summarise(&records)
);
self
}
pub fn assert_job_failed(&self, kind: &str) -> &Self {
let records = self.records();
assert!(
records
.iter()
.any(|r| r.kind == kind && matches!(r.outcome, JobOutcome::Failed(_))),
"no job of kind `{kind}` failed; {}",
summarise(&records)
);
self
}
}
impl Observer for Jobs {
fn observe(&self, event: &Event) {
let Ok(mut records) = self.records.lock() else {
return;
};
let record = match event {
Event::Started {
job_id,
kind,
version,
attempt,
} => JobRecord {
job_id: *job_id,
kind: kind.clone(),
version: *version,
attempt: *attempt,
outcome: JobOutcome::Started,
},
Event::Succeeded {
job_id, attempt, ..
} => {
let (kind, version) = resolve_kind(&records, *job_id);
JobRecord {
job_id: *job_id,
kind,
version,
attempt: *attempt,
outcome: JobOutcome::Succeeded,
}
}
Event::Retried {
job_id, attempt, ..
} => {
let (kind, version) = resolve_kind(&records, *job_id);
JobRecord {
job_id: *job_id,
kind,
version,
attempt: *attempt,
outcome: JobOutcome::Retried,
}
}
Event::Failed {
job_id,
attempt,
reason,
..
} => {
let (kind, version) = resolve_kind(&records, *job_id);
JobRecord {
job_id: *job_id,
kind,
version,
attempt: *attempt,
outcome: JobOutcome::Failed(*reason),
}
}
};
records.push(record);
}
}
fn resolve_kind(records: &[JobRecord], job_id: uuid::Uuid) -> (String, i16) {
records
.iter()
.find(|record| record.job_id == job_id)
.map_or_else(
|| ("(unknown)".to_owned(), 0),
|record| (record.kind.clone(), record.version),
)
}
fn summarise(records: &[JobRecord]) -> String {
if records.is_empty() {
return "the worker ran no jobs at all".to_owned();
}
let lines: Vec<String> = records
.iter()
.map(|r| {
format!(
" {} v{} attempt {} {}",
r.kind, r.version, r.attempt, r.outcome
)
})
.collect();
format!("what ran:\n{}", lines.join("\n"))
}
}
#[cfg(feature = "mail")]
mod mail {
use crate::mail::{Mailer, parse_mailbox};
#[derive(Debug, Clone)]
pub struct SentMail {
pub to: Vec<String>,
pub raw: String,
}
impl SentMail {
#[must_use]
pub fn subject(&self) -> Option<&str> {
self.raw
.lines()
.take_while(|line| !line.is_empty())
.find_map(|line| line.strip_prefix("Subject: "))
}
#[must_use]
pub fn contains(&self, needle: &str) -> bool {
self.raw.contains(needle)
}
}
#[derive(Clone)]
pub struct Mail {
mailer: Mailer,
}
impl std::fmt::Debug for Mail {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Mail").field("capture", &true).finish()
}
}
impl Mail {
#[must_use]
pub fn fake() -> Self {
Self {
mailer: Mailer::capture_ok(),
}
}
#[must_use]
pub fn mailer(&self) -> Mailer {
self.mailer.clone()
}
pub fn sender(
&self,
from: &str,
) -> Result<crate::mail::Mail, lettre::address::AddressError> {
Ok(crate::mail::Mail::new(
self.mailer.clone(),
parse_mailbox(from)?,
))
}
}
}
#[cfg(feature = "mail")]
impl mail::Mail {
pub async fn sent(&self) -> Vec<SentMail> {
let captured = self
.mailer()
.captured()
.await
.expect("a fake mailer is always a capture mailer");
captured
.into_iter()
.map(|(envelope, raw)| SentMail {
to: envelope.to().iter().map(ToString::to_string).collect(),
raw,
})
.collect()
}
pub async fn assert_mail_sent(&self, address: &str) -> &Self {
let sent = self.sent().await;
let matched = sent
.iter()
.any(|message| message.to.iter().any(|to| to == address));
assert!(
matched,
"no mail was sent to `{address}`; {}",
describe(&sent)
);
self
}
pub async fn assert_mail_contains(&self, address: &str, needle: &str) -> &Self {
let sent = self.sent().await;
let matched = sent
.iter()
.any(|message| message.to.iter().any(|to| to == address) && message.contains(needle));
assert!(
matched,
"no mail to `{address}` contains `{needle}`; {}",
describe(&sent)
);
self
}
pub async fn assert_no_mail_sent(&self) -> &Self {
let sent = self.sent().await;
assert!(sent.is_empty(), "expected no mail; {}", describe(&sent));
self
}
}
#[cfg(feature = "mail")]
fn describe(sent: &[SentMail]) -> String {
if sent.is_empty() {
return "nothing was sent at all".to_owned();
}
let lines: Vec<String> = sent
.iter()
.map(|message| {
format!(
" to [{}] subject {:?}",
message.to.join(", "),
message.subject().unwrap_or("(none)")
)
})
.collect();
format!("what was sent:\n{}", lines.join("\n"))
}
#[cfg(feature = "jobs")]
pub async fn assert_job_enqueued(connection: &mut crate::database::Connection, kind: &str) {
let queued: Vec<String> =
sqlx::query_scalar::<crate::database::Driver, String>("SELECT kind FROM arcature_jobs")
.fetch_all(&mut *connection)
.await
.unwrap_or_else(|error| panic!("could not read `arcature_jobs`: {error}"));
assert!(
queued.iter().any(|queued| queued == kind),
"no job of kind `{kind}` is queued; the queue holds {}",
if queued.is_empty() {
"nothing".to_owned()
} else {
format!("[{}]", queued.join(", "))
}
);
}
#[cfg(feature = "events")]
mod events {
use crate::events::Dispatcher;
#[derive(Clone)]
pub struct Events {
dispatcher: Dispatcher,
}
impl std::fmt::Debug for Events {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Events")
.field("dispatched", &self.dispatcher.dispatched_events())
.finish()
}
}
impl Events {
#[must_use]
pub fn fake() -> Self {
Self {
dispatcher: Dispatcher::recording(),
}
}
#[must_use]
pub fn dispatcher(&self) -> Dispatcher {
self.dispatcher.clone()
}
#[must_use]
pub fn register<E, F, Fut>(self, listener: F) -> Self
where
E: crate::events::Event + serde::Serialize + serde::de::DeserializeOwned,
F: Fn(E) -> Fut + Send + Sync + 'static,
Fut: std::future::Future<Output = Result<(), crate::events::DispatchError>>
+ Send
+ 'static,
{
Self {
dispatcher: self.dispatcher.register(listener),
}
}
#[must_use]
pub fn dispatched(&self) -> Vec<String> {
self.dispatcher.dispatched_events()
}
pub fn assert_dispatched(&self, name: &str) -> &Self {
let dispatched = self.dispatched();
assert!(
dispatched.iter().any(|event| event == name),
"event `{name}` was not dispatched; {}",
if dispatched.is_empty() {
"nothing was dispatched at all".to_owned()
} else {
format!("what was: [{}]", dispatched.join(", "))
}
);
self
}
pub fn assert_not_dispatched(&self, name: &str) -> &Self {
let dispatched = self.dispatched();
assert!(
!dispatched.iter().any(|event| event == name),
"event `{name}` was dispatched; the full sequence was [{}]",
dispatched.join(", ")
);
self
}
}
}
#[cfg(all(test, feature = "jobs"))]
mod job_recorder_tests {
use super::{JobOutcome, Jobs};
use crate::jobs::{Event, FailReason, Observer};
use std::time::Duration;
fn started(id: uuid::Uuid, kind: &str) -> Event {
Event::Started {
job_id: id,
kind: kind.to_owned(),
version: 1,
attempt: 1,
}
}
#[test]
fn a_recorder_starts_empty_and_every_assertion_fails_on_it() {
let recorder = Jobs::recorder();
assert!(recorder.records().is_empty());
let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
recorder.assert_job_ran("SendWelcome")
}));
assert!(
outcome.is_err(),
"an empty recorder must fail every assertion, not pass vacuously"
);
}
#[test]
fn a_started_event_is_recorded_under_its_kind() {
let recorder = Jobs::recorder();
recorder.observe(&started(uuid::Uuid::nil(), "SendWelcome"));
recorder.assert_job_ran("SendWelcome");
}
#[test]
fn a_later_event_recovers_the_kind_from_the_start_of_the_same_job() {
let id = uuid::Uuid::from_u128(7);
let recorder = Jobs::recorder();
recorder.observe(&started(id, "SendWelcome"));
recorder.observe(&Event::Succeeded {
job_id: id,
attempt: 1,
duration: Duration::from_millis(3),
});
recorder.assert_job_succeeded("SendWelcome");
assert_eq!(recorder.records()[1].outcome, JobOutcome::Succeeded);
}
#[test]
fn a_job_that_only_started_has_not_succeeded() {
let recorder = Jobs::recorder();
recorder.observe(&started(uuid::Uuid::nil(), "SendWelcome"));
let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
recorder.assert_job_succeeded("SendWelcome")
}));
assert!(outcome.is_err(), "starting is not succeeding");
}
#[test]
fn a_permanent_failure_is_recorded_with_its_reason() {
let id = uuid::Uuid::from_u128(9);
let recorder = Jobs::recorder();
recorder.observe(&started(id, "ChargeCard"));
recorder.observe(&Event::Failed {
job_id: id,
attempt: 3,
duration: Duration::from_millis(1),
message: "declined".to_owned(),
reason: FailReason::Exhausted,
});
recorder.assert_job_failed("ChargeCard");
assert_eq!(
recorder.records()[1].outcome,
JobOutcome::Failed(FailReason::Exhausted)
);
}
#[test]
fn a_failure_message_lists_the_jobs_that_did_run() {
let recorder = Jobs::recorder();
recorder.observe(&started(uuid::Uuid::nil(), "SendWelcome"));
let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
recorder.assert_job_ran("ChargeCard")
}));
let payload = outcome.expect_err("the assertion must fail");
let message = payload
.downcast_ref::<String>()
.expect("panic payload is a String");
assert!(message.contains("SendWelcome"), "message: {message}");
}
}
#[cfg(all(test, feature = "events"))]
mod event_recorder_tests {
use super::Events;
#[test]
fn a_fresh_recorder_reports_nothing_dispatched() {
let events = Events::fake();
assert!(events.dispatched().is_empty());
events.assert_not_dispatched("UserRegistered");
}
#[test]
fn asserting_a_dispatch_on_a_fresh_recorder_fails() {
let events = Events::fake();
let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
events.assert_dispatched("UserRegistered")
}));
assert!(
outcome.is_err(),
"an empty recorder must not satisfy assert_dispatched"
);
}
}