use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use renox::Kernel;
use renox::db::Dialect;
use renox::prelude::*;
use serde::{Deserialize, Serialize};
async fn log(state: &AppState, message: &str) -> Result {
renox::db::sql("INSERT INTO log (message) VALUES (?)")
.bind(message.to_owned())
.execute(&state.db)
.await?;
Ok(())
}
async fn logged(kernel: &Kernel) -> Vec<String> {
renox::db::sql("SELECT message FROM log ORDER BY id")
.scalars(kernel.db())
.await
.unwrap()
}
#[derive(Serialize, Deserialize)]
struct Greet {
name: String,
}
impl Job for Greet {
const NAME: &'static str = "greet";
async fn handle(self, ctx: JobContext) -> Result {
log(&ctx.state, &format!("hello {}", self.name)).await
}
}
#[derive(Serialize, Deserialize)]
struct Flaky;
impl Job for Flaky {
const NAME: &'static str = "flaky";
fn backoff(_: u32) -> Duration {
Duration::ZERO
}
async fn handle(self, ctx: JobContext) -> Result {
log(&ctx.state, &format!("attempt {}", ctx.attempt)).await?;
if ctx.attempt < 3 {
return Err(anyhow_error("not yet"));
}
Ok(())
}
}
#[derive(Serialize, Deserialize)]
struct Broken;
impl Job for Broken {
const NAME: &'static str = "broken";
const MAX_ATTEMPTS: u32 = 2;
fn backoff(_: u32) -> Duration {
Duration::ZERO
}
async fn handle(self, _: JobContext) -> Result {
Err(anyhow_error("the printer is on fire"))
}
}
#[derive(Serialize, Deserialize)]
struct Slow;
impl Job for Slow {
const NAME: &'static str = "slow";
const MAX_ATTEMPTS: u32 = 1;
const TIMEOUT: Duration = Duration::from_millis(50);
async fn handle(self, _: JobContext) -> Result {
tokio::time::sleep(Duration::from_secs(5)).await;
Ok(())
}
}
#[derive(Serialize, Deserialize)]
struct Email;
impl Job for Email {
const NAME: &'static str = "email";
const QUEUE: &'static str = "emails";
async fn handle(self, ctx: JobContext) -> Result {
log(&ctx.state, "email sent").await
}
}
#[derive(Serialize, Deserialize)]
struct Orphan;
impl Job for Orphan {
const NAME: &'static str = "orphan";
async fn handle(self, _: JobContext) -> Result {
Ok(())
}
}
fn anyhow_error(message: &str) -> Error {
Error::Internal(renox_anyhow(message))
}
fn renox_anyhow(message: &str) -> anyhow::Error {
anyhow::anyhow!(message.to_owned())
}
#[derive(Clone)]
struct OrderPlaced {
id: i64,
}
impl Event for OrderPlaced {}
struct Shop;
impl Module for Shop {
fn name(&self) -> &'static str {
"shop"
}
fn register(&self, app: &mut Registry) {
app.job::<Email>()
.listen(|e: OrderPlaced, state| async move {
log(&state, &format!("shop saw order {}", e.id)).await
});
}
}
fn config() -> Config {
{
let mut c = Config::default();
c.env = Environment::Testing;
c.key = Some(renox::generate_key());
c
}
}
async fn kernel_with(app: impl FnOnce(App) -> App) -> Kernel {
let kernel = app(App::with_config(config())
.module(Shop)
.job::<Greet>()
.job::<Flaky>()
.job::<Broken>()
.job::<Slow>())
.boot()
.await
.unwrap();
kernel.migrate().await.unwrap();
renox::db::sql(match kernel.db().dialect() {
Dialect::Sqlite => "CREATE TABLE log (id INTEGER PRIMARY KEY, message TEXT NOT NULL)",
Dialect::Postgres => {
"CREATE TABLE log (id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY, \
message TEXT NOT NULL)"
}
})
.execute(kernel.db())
.await
.unwrap();
kernel
}
async fn kernel() -> Kernel {
kernel_with(|app| app).await
}
#[tokio::test]
async fn dispatched_jobs_run() {
let kernel = kernel().await;
let queue = &kernel.state().queue;
kernel
.state()
.dispatch(Greet {
name: "Alex".into(),
})
.await
.unwrap();
queue.dispatch(Greet { name: "Ben".into() }).await.unwrap();
assert_eq!(queue.pending().await.unwrap(), 2);
assert_eq!(kernel.run_jobs().await.unwrap(), 2);
assert_eq!(logged(&kernel).await, ["hello Alex", "hello Ben"]);
assert_eq!(queue.pending().await.unwrap(), 0);
}
#[tokio::test]
async fn failing_jobs_are_retried() {
let kernel = kernel().await;
kernel.state().dispatch(Flaky).await.unwrap();
kernel.run_jobs().await.unwrap();
assert_eq!(
logged(&kernel).await,
["attempt 1", "attempt 2", "attempt 3"]
);
assert!(kernel.state().queue.failed().await.unwrap().is_empty());
}
#[tokio::test]
async fn jobs_that_keep_failing_can_be_retried_later() {
let kernel = kernel().await;
let queue = &kernel.state().queue;
queue.dispatch(Broken).await.unwrap();
kernel.run_jobs().await.unwrap();
let failed = queue.failed().await.unwrap();
assert_eq!(failed.len(), 1);
assert_eq!(failed[0].job, "broken");
assert!(
failed[0].error.contains("the printer is on fire"),
"{}",
failed[0].error
);
assert_eq!(queue.pending().await.unwrap(), 0);
assert!(queue.retry(failed[0].id).await.unwrap());
assert_eq!(queue.pending().await.unwrap(), 1);
let attempts: i64 = renox::db::sql("SELECT max_attempts FROM jobs")
.scalar(kernel.db())
.await
.unwrap();
assert_eq!(attempts, 2, "a retried job gets its own attempts back");
kernel.run_jobs().await.unwrap();
assert_eq!(queue.flush_failed().await.unwrap(), 1);
assert!(queue.failed().await.unwrap().is_empty());
}
#[tokio::test]
async fn slow_and_unknown_jobs_fail() {
let kernel = kernel().await;
let queue = &kernel.state().queue;
queue.dispatch(Slow).await.unwrap();
queue.dispatch(Orphan).await.unwrap();
kernel.run_jobs().await.unwrap();
let errors: Vec<String> = queue
.failed()
.await
.unwrap()
.into_iter()
.map(|f| f.error)
.collect();
assert!(errors[0].contains("timed out"), "{errors:?}");
assert_eq!(errors[1], "no handler registered for job `orphan`");
}
#[tokio::test]
async fn delayed_jobs_wait() {
let kernel = kernel().await;
let queue = &kernel.state().queue;
queue
.dispatch_after(
Greet {
name: "later".into(),
},
Duration::from_secs(3600),
)
.await
.unwrap();
assert_eq!(kernel.run_jobs().await.unwrap(), 0);
assert_eq!(queue.pending().await.unwrap(), 1);
}
#[tokio::test]
async fn workers_can_be_limited_to_queues() {
let kernel = kernel().await;
kernel.state().dispatch(Email).await.unwrap();
kernel
.state()
.dispatch(Greet {
name: "Alex".into(),
})
.await
.unwrap();
assert_eq!(
kernel.worker(vec!["emails".into()]).drain().await.unwrap(),
1
);
assert_eq!(logged(&kernel).await, ["email sent"]);
assert_eq!(kernel.state().queue.pending().await.unwrap(), 1);
}
#[tokio::test]
async fn events_reach_every_listener() {
let calls = Arc::new(AtomicUsize::new(0));
let counter = calls.clone();
let kernel = kernel_with(move |app| {
app.listen(move |e: OrderPlaced, state| {
let counter = counter.clone();
async move {
counter.fetch_add(1, Ordering::SeqCst);
log(&state, &format!("app saw order {}", e.id)).await?;
Err(anyhow_error("listener broke"))
}
})
})
.await;
let result = kernel.state().emit(OrderPlaced { id: 7 }).await;
assert!(format!("{:?}", result.unwrap_err()).contains("listener broke"));
assert_eq!(
logged(&kernel).await,
["app saw order 7", "shop saw order 7"]
);
assert_eq!(calls.load(Ordering::SeqCst), 1);
#[derive(Clone)]
struct Nobody;
impl Event for Nobody {}
assert!(kernel.state().emit(Nobody).await.is_ok());
}
#[tokio::test]
async fn background_workers_pick_up_jobs_and_stop_cleanly() {
let kernel = kernel().await;
let (stop, stopped) = tokio::sync::watch::channel(false);
let worker = tokio::spawn(kernel.worker(Vec::new()).run(2, stopped));
kernel
.state()
.dispatch(Greet {
name: "background".into(),
})
.await
.unwrap();
let tries = match kernel.db().dialect() {
Dialect::Sqlite => 50,
Dialect::Postgres => 150,
};
for _ in 0..tries {
if !logged(&kernel).await.is_empty() {
break;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
assert_eq!(
logged(&kernel).await,
["hello background"],
"woken without waiting for the poll"
);
stop.send(true).unwrap();
tokio::time::timeout(Duration::from_secs(3), worker)
.await
.unwrap()
.unwrap();
}
#[tokio::test]
async fn misconfiguration_fails_at_boot() {
let twice = App::with_config(config())
.job::<Greet>()
.job::<Greet>()
.boot()
.await;
assert!(format!("{:?}", twice.err().unwrap()).contains("job `greet` is registered twice"));
let bad_time = App::with_config(config())
.schedule(|s| {
s.daily_at("25:00", "nope", |_| async { Ok(()) });
})
.boot()
.await;
assert!(format!("{:?}", bad_time.err().unwrap()).contains("task `nope`"));
}
#[tokio::test]
async fn scheduled_tasks_run_on_demand_with_their_hooks() {
static RUNS: AtomicUsize = AtomicUsize::new(0);
let kernel = kernel_with(|app| {
app.schedule(|s| {
s.cron("0 3 * * *", "report", |state| async move {
RUNS.fetch_add(1, Ordering::SeqCst);
assert!(renox::context::app().is_some(), "runs in the app's context");
state.cache.put("report.ran", &true, None).await
})
.timezone("Asia/Jakarta")
.on_success(|state| async move {
let _ = state.cache.put("report.ok", &true, None).await;
});
s.daily_at("04:00", "broken", |_| async {
Err(abort(StatusCode::INTERNAL_SERVER_ERROR, "disk full"))
})
.on_failure(|err, state| async move {
let _ = state
.cache
.put("broken.error", &format!("{err:?}"), None)
.await;
});
s.hourly("panics", |_| async { panic!("boom") })
.on_failure(|_, state| async move {
let _ = state.cache.put("panics.caught", &true, None).await;
});
})
})
.await;
let cache = &kernel.state().cache;
kernel.run_scheduled("report").await.unwrap();
assert_eq!(RUNS.load(Ordering::SeqCst), 1);
assert_eq!(cache.get::<bool>("report.ok").await.unwrap(), Some(true));
assert!(kernel.run_scheduled("broken").await.is_err());
let error = cache.get::<String>("broken.error").await.unwrap().unwrap();
assert!(error.contains("disk full"), "{error}");
assert!(kernel.run_scheduled("panics").await.is_err());
assert_eq!(
cache.get::<bool>("panics.caught").await.unwrap(),
Some(true)
);
assert!(kernel.run_scheduled("nope").await.is_err());
}
#[tokio::test]
async fn each_scheduled_run_is_claimed_once() {
let kernel = kernel().await;
let state = kernel.state();
assert!(renox::schedule::claim(state, "report", 1_800_000_000, 120).await);
assert!(!renox::schedule::claim(state, "report", 1_800_000_000, 120).await);
assert!(renox::schedule::claim(state, "report", 1_800_000_060, 120).await);
assert!(renox::schedule::claim(state, "backup", 1_800_000_000, 120).await);
}