use awa::model::{
cron::{
atomic_enqueue, delete_cron_job, list_cron_jobs, pause_cron_job, resume_cron_job,
trigger_cron_job, upsert_cron_job,
},
migrations,
};
use awa::{
Client, CronMissedFirePolicy, JobArgs, JobContext, JobResult, JobState, PeriodicJob,
QueueConfig,
};
use chrono::Utc;
use serde::{Deserialize, Serialize};
use sqlx::postgres::PgPoolOptions;
use sqlx::PgPool;
fn database_url() -> String {
std::env::var("DATABASE_URL")
.unwrap_or_else(|_| "postgres://postgres:test@localhost:15432/awa_test".to_string())
}
async fn pool() -> PgPool {
PgPoolOptions::new()
.max_connections(10)
.connect(&database_url())
.await
.expect("Failed to connect to database")
}
async fn setup() -> PgPool {
let pool = pool().await;
migrations::run(&pool)
.await
.expect("Failed to run migrations");
sqlx::query(
r#"
UPDATE awa.storage_transition_state
SET current_engine = 'canonical',
prepared_engine = NULL,
state = 'canonical',
transition_epoch = transition_epoch + 1,
details = '{}'::jsonb,
updated_at = now(),
finalized_at = NULL
WHERE singleton
"#,
)
.execute(&pool)
.await
.expect("Failed to reset storage transition state");
sqlx::query("DELETE FROM awa.runtime_storage_backends WHERE backend = 'queue_storage'")
.execute(&pool)
.await
.expect("Failed to reset active runtime backend");
pool
}
async fn clean_cron_names(pool: &PgPool, names: &[&str]) {
for name in names {
let _ = sqlx::query("DELETE FROM awa.cron_jobs WHERE name = $1")
.bind(name)
.execute(pool)
.await;
}
}
async fn clean_queue(pool: &PgPool, queue: &str) {
sqlx::query("DELETE FROM awa.jobs WHERE queue = $1")
.bind(queue)
.execute(pool)
.await
.expect("Failed to clean queue");
}
async fn wait_for_leader(client: &Client, timeout: std::time::Duration) {
let start = std::time::Instant::now();
loop {
if client.health_check().await.leader {
return;
}
assert!(
start.elapsed() < timeout,
"Timed out waiting for periodic test client to become leader"
);
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
}
#[derive(Debug, Serialize, Deserialize, JobArgs)]
struct DailyReport {
format: String,
}
#[tokio::test]
async fn test_canonical_schema_creates_cron_jobs_table() {
let pool = setup().await;
let version = migrations::current_version(&pool).await.unwrap();
assert_eq!(version, migrations::CURRENT_VERSION);
let has_cron_table: bool = sqlx::query_scalar(
"SELECT EXISTS(SELECT 1 FROM information_schema.tables WHERE table_schema = 'awa' AND table_name = 'cron_jobs')",
)
.fetch_one(&pool)
.await
.unwrap();
assert!(
has_cron_table,
"awa.cron_jobs table must exist after canonical schema setup"
);
}
#[tokio::test]
async fn test_migration_idempotent() {
let pool = setup().await;
migrations::run(&pool).await.unwrap();
migrations::run(&pool).await.unwrap();
let version = migrations::current_version(&pool).await.unwrap();
assert_eq!(version, migrations::CURRENT_VERSION);
}
#[tokio::test]
async fn test_upsert_inserts_new_schedule() {
let pool = setup().await;
clean_cron_names(&pool, &["test_insert_new"]).await;
let job = PeriodicJob::builder("test_insert_new", "0 9 * * *")
.build_raw(
"daily_report".to_string(),
serde_json::json!({"format": "pdf"}),
)
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
let rows = list_cron_jobs(&pool).await.unwrap();
let found = rows.iter().find(|r| r.name == "test_insert_new").unwrap();
assert_eq!(found.cron_expr, "0 9 * * *");
assert_eq!(found.kind, "daily_report");
assert_eq!(found.timezone, "UTC");
assert_eq!(found.queue, "default");
assert_eq!(found.priority, 2);
assert_eq!(found.missed_fire_policy, "coalesce");
assert!(found.last_enqueued_at.is_none());
}
#[tokio::test]
async fn test_upsert_updates_existing_schedule() {
let pool = setup().await;
clean_cron_names(&pool, &["test_upsert_update"]).await;
let job_v1 = PeriodicJob::builder("test_upsert_update", "0 9 * * *")
.build_raw(
"daily_report".to_string(),
serde_json::json!({"format": "pdf"}),
)
.unwrap();
upsert_cron_job(&pool, &job_v1).await.unwrap();
let job_v2 = PeriodicJob::builder("test_upsert_update", "30 8 * * *")
.timezone("Pacific/Auckland")
.queue("reports")
.missed_fire_policy(CronMissedFirePolicy::CatchUp)
.build_raw(
"daily_report".to_string(),
serde_json::json!({"format": "csv"}),
)
.unwrap();
upsert_cron_job(&pool, &job_v2).await.unwrap();
let rows = list_cron_jobs(&pool).await.unwrap();
let found = rows
.iter()
.find(|r| r.name == "test_upsert_update")
.unwrap();
assert_eq!(found.cron_expr, "30 8 * * *");
assert_eq!(found.timezone, "Pacific/Auckland");
assert_eq!(found.queue, "reports");
assert_eq!(found.args, serde_json::json!({"format": "csv"}));
assert_eq!(found.missed_fire_policy, "catch_up");
}
#[tokio::test]
async fn test_multi_deployment_no_orphan_deletion() {
let pool = setup().await;
clean_cron_names(&pool, &["deploy_a_job", "deploy_b_job"]).await;
let job_a = PeriodicJob::builder("deploy_a_job", "0 * * * *")
.build_raw("sync_a".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job_a).await.unwrap();
let job_b = PeriodicJob::builder("deploy_b_job", "30 * * * *")
.build_raw("sync_b".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job_b).await.unwrap();
let rows = list_cron_jobs(&pool).await.unwrap();
let names: Vec<&str> = rows.iter().map(|r| r.name.as_str()).collect();
assert!(
names.contains(&"deploy_a_job"),
"Deployment A's schedule must be preserved"
);
assert!(
names.contains(&"deploy_b_job"),
"Deployment B's schedule must be preserved"
);
}
#[tokio::test]
async fn test_atomic_cte_mark_and_insert() {
let pool = setup().await;
clean_cron_names(&pool, &["test_atomic_enqueue"]).await;
let queue = "cron_atomic_cte";
clean_queue(&pool, queue).await;
let job = PeriodicJob::builder("test_atomic_enqueue", "* * * * *")
.queue(queue)
.build_raw(
"daily_report".to_string(),
serde_json::json!({"format": "pdf"}),
)
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
let fire_time = Utc::now();
let result = atomic_enqueue(&pool, "test_atomic_enqueue", fire_time, None)
.await
.unwrap();
assert!(result.is_some(), "Should have enqueued a job");
let job_row = result.unwrap();
assert_eq!(job_row.kind, "daily_report");
assert_eq!(job_row.queue, queue);
assert_eq!(job_row.state, JobState::Available);
assert_eq!(job_row.priority, 2);
assert_eq!(
job_row.metadata.get("cron_name").and_then(|v| v.as_str()),
Some("test_atomic_enqueue")
);
assert!(job_row.metadata.get("cron_fire_time").is_some());
let rows = list_cron_jobs(&pool).await.unwrap();
let cron_row = rows
.iter()
.find(|r| r.name == "test_atomic_enqueue")
.unwrap();
assert!(cron_row.last_enqueued_at.is_some());
}
#[tokio::test]
async fn test_atomic_cte_dedup_second_call() {
let pool = setup().await;
clean_cron_names(&pool, &["test_dedup_cte"]).await;
let queue = "cron_atomic_dedup";
clean_queue(&pool, queue).await;
let job = PeriodicJob::builder("test_dedup_cte", "* * * * *")
.queue(queue)
.build_raw("daily_report".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
let fire_time = Utc::now();
let result1 = atomic_enqueue(&pool, "test_dedup_cte", fire_time, None)
.await
.unwrap();
assert!(result1.is_some());
let result2 = atomic_enqueue(&pool, "test_dedup_cte", fire_time, None)
.await
.unwrap();
assert!(result2.is_none(), "Second call should return None (dedup)");
}
#[tokio::test]
async fn test_no_backfill_only_latest_fire() {
let pool = setup().await;
clean_cron_names(&pool, &["test_no_backfill"]).await;
let queue = "cron_no_backfill";
clean_queue(&pool, queue).await;
let job = PeriodicJob::builder("test_no_backfill", "* * * * *")
.queue(queue)
.build_raw("hourly_sync".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
let one_hour_ago = Utc::now() - chrono::Duration::hours(1);
sqlx::query("UPDATE awa.cron_jobs SET last_enqueued_at = $1 WHERE name = $2")
.bind(one_hour_ago)
.bind("test_no_backfill")
.execute(&pool)
.await
.unwrap();
let now = Utc::now();
let fire = job.latest_fire_time(now, Some(one_hour_ago));
assert!(fire.is_some(), "Should find a fire time");
let fire_time = fire.unwrap();
assert!(
(now - fire_time).num_seconds() < 120,
"Fire time should be within the last 2 minutes, not backfilled. Got {:?}",
fire_time
);
let result = atomic_enqueue(&pool, "test_no_backfill", fire_time, Some(one_hour_ago))
.await
.unwrap();
assert!(result.is_some());
let count: (i64,) = sqlx::query_as("SELECT count(*)::bigint FROM awa.jobs WHERE queue = $1")
.bind(queue)
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(
count.0, 1,
"Should create exactly 1 job, not backfill all missed fires"
);
}
#[tokio::test]
async fn test_delete_cron_job() {
let pool = setup().await;
clean_cron_names(&pool, &["test_delete"]).await;
let job = PeriodicJob::builder("test_delete", "0 9 * * *")
.build_raw("test_job".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
let deleted = delete_cron_job(&pool, "test_delete").await.unwrap();
assert!(deleted);
let deleted_again = delete_cron_job(&pool, "test_delete").await.unwrap();
assert!(!deleted_again, "Second delete should return false");
}
#[tokio::test]
async fn test_end_to_end_periodic_job_enqueued() {
let pool = setup().await;
clean_cron_names(&pool, &["e2e_test"]).await;
let queue = "cron_e2e";
clean_queue(&pool, queue).await;
let client = Client::builder(pool.clone())
.queue(queue, QueueConfig::default())
.leader_election_interval(std::time::Duration::from_secs(1))
.register::<DailyReport, _, _>(|_args: DailyReport, _ctx: &JobContext| async move {
Ok(JobResult::Completed)
})
.periodic(
PeriodicJob::builder("e2e_test", "* * * * *")
.queue(queue)
.build(&DailyReport {
format: "pdf".into(),
})
.unwrap(),
)
.build()
.unwrap();
client.start().await.unwrap();
wait_for_leader(&client, std::time::Duration::from_secs(15)).await;
let start = std::time::Instant::now();
let timeout = std::time::Duration::from_secs(120);
let jobs = loop {
let found: Vec<awa::JobRow> =
sqlx::query_as("SELECT * FROM awa.jobs WHERE queue = $1 AND kind = 'daily_report'")
.bind(queue)
.fetch_all(&pool)
.await
.unwrap();
if !found.is_empty() {
break found;
}
assert!(
start.elapsed() < timeout,
"Periodic job was not enqueued within {timeout:?}"
);
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
};
let job = &jobs[0];
assert_eq!(
job.metadata.get("cron_name").and_then(|v| v.as_str()),
Some("e2e_test")
);
assert!(job.metadata.get("cron_fire_time").is_some());
let cron_rows = list_cron_jobs(&pool).await.unwrap();
let found = cron_rows.iter().find(|r| r.name == "e2e_test");
assert!(found.is_some(), "Schedule should be synced to DB");
client.shutdown(std::time::Duration::from_secs(2)).await;
}
#[tokio::test]
async fn test_periodic_job_builder_with_job_args_trait() {
let job = PeriodicJob::builder("typed_test", "0 9 * * *")
.timezone("Pacific/Auckland")
.queue("reports")
.priority(1)
.build(&DailyReport {
format: "csv".into(),
})
.unwrap();
assert_eq!(job.name, "typed_test");
assert_eq!(job.kind, "daily_report");
assert_eq!(job.args, serde_json::json!({"format": "csv"}));
assert_eq!(job.timezone, "Pacific/Auckland");
assert_eq!(job.queue, "reports");
assert_eq!(job.priority, 1);
}
#[tokio::test]
async fn test_cron_job_with_tags_and_metadata() {
let pool = setup().await;
clean_cron_names(&pool, &["tagged_job"]).await;
let queue = "cron_tags_meta";
clean_queue(&pool, queue).await;
let job = PeriodicJob::builder("tagged_job", "0 9 * * *")
.queue(queue)
.tags(vec!["important".to_string(), "daily".to_string()])
.metadata(serde_json::json!({"team": "analytics"}))
.build_raw("report".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
let rows = list_cron_jobs(&pool).await.unwrap();
let found = rows.iter().find(|r| r.name == "tagged_job").unwrap();
assert_eq!(found.tags, vec!["important", "daily"]);
assert_eq!(
found.metadata.get("team").and_then(|v| v.as_str()),
Some("analytics")
);
let fire_time = Utc::now();
let result = atomic_enqueue(&pool, "tagged_job", fire_time, None)
.await
.unwrap();
let job_row = result.unwrap();
assert_eq!(job_row.tags, vec!["important", "daily"]);
assert_eq!(
job_row.metadata.get("team").and_then(|v| v.as_str()),
Some("analytics")
);
assert_eq!(
job_row.metadata.get("cron_name").and_then(|v| v.as_str()),
Some("tagged_job")
);
}
#[tokio::test]
async fn test_pause_sets_paused_at_and_paused_by() {
let pool = setup().await;
clean_cron_names(&pool, &["pause_sets_fields"]).await;
let job = PeriodicJob::builder("pause_sets_fields", "* * * * *")
.build_raw("pause_test".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
let updated = pause_cron_job(&pool, "pause_sets_fields", Some("alice"))
.await
.unwrap();
assert!(updated, "pause should report row updated");
let row = list_cron_jobs(&pool)
.await
.unwrap()
.into_iter()
.find(|r| r.name == "pause_sets_fields")
.unwrap();
assert!(row.is_paused(), "row should report paused");
assert_eq!(row.paused_by.as_deref(), Some("alice"));
assert!(row.paused_at.is_some());
}
#[tokio::test]
async fn test_pause_unknown_returns_false() {
let pool = setup().await;
clean_cron_names(&pool, &["pause_unknown"]).await;
let updated = pause_cron_job(&pool, "pause_unknown", None).await.unwrap();
assert!(
!updated,
"pause on a non-existent schedule should report no row updated"
);
}
#[tokio::test]
async fn test_atomic_enqueue_refuses_while_paused() {
let pool = setup().await;
clean_cron_names(&pool, &["pause_blocks_enqueue"]).await;
let queue = "cron_pause_blocks";
clean_queue(&pool, queue).await;
let job = PeriodicJob::builder("pause_blocks_enqueue", "* * * * *")
.queue(queue)
.build_raw("pause_test".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
pause_cron_job(&pool, "pause_blocks_enqueue", Some("bob"))
.await
.unwrap();
let fire_time = Utc::now();
let result = atomic_enqueue(&pool, "pause_blocks_enqueue", fire_time, None)
.await
.unwrap();
assert!(
result.is_none(),
"atomic_enqueue must return None while the schedule is paused"
);
let row = list_cron_jobs(&pool)
.await
.unwrap()
.into_iter()
.find(|r| r.name == "pause_blocks_enqueue")
.unwrap();
assert!(
row.last_enqueued_at.is_none(),
"last_enqueued_at must not advance while paused"
);
}
#[tokio::test]
async fn test_resume_re_enables_enqueue() {
let pool = setup().await;
clean_cron_names(&pool, &["resume_enables"]).await;
let queue = "cron_resume_enables";
clean_queue(&pool, queue).await;
let job = PeriodicJob::builder("resume_enables", "* * * * *")
.queue(queue)
.build_raw("pause_test".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
pause_cron_job(&pool, "resume_enables", None).await.unwrap();
let updated = resume_cron_job(&pool, "resume_enables").await.unwrap();
assert!(updated, "resume should report row updated");
let fire_time = Utc::now();
let result = atomic_enqueue(&pool, "resume_enables", fire_time, None)
.await
.unwrap();
assert!(
result.is_some(),
"atomic_enqueue should succeed after resume"
);
let row = list_cron_jobs(&pool)
.await
.unwrap()
.into_iter()
.find(|r| r.name == "resume_enables")
.unwrap();
assert!(!row.is_paused());
assert!(row.paused_at.is_none());
assert!(row.paused_by.is_none());
}
#[tokio::test]
async fn test_pause_preserves_last_enqueued_at() {
let pool = setup().await;
clean_cron_names(&pool, &["pause_preserves"]).await;
let queue = "cron_pause_preserves";
clean_queue(&pool, queue).await;
let job = PeriodicJob::builder("pause_preserves", "* * * * *")
.queue(queue)
.build_raw("pause_test".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
let fire_time = Utc::now();
atomic_enqueue(&pool, "pause_preserves", fire_time, None)
.await
.unwrap()
.expect("first enqueue should succeed");
let before = list_cron_jobs(&pool)
.await
.unwrap()
.into_iter()
.find(|r| r.name == "pause_preserves")
.unwrap();
let last_before = before.last_enqueued_at.expect("should be set after fire");
pause_cron_job(&pool, "pause_preserves", None)
.await
.unwrap();
resume_cron_job(&pool, "pause_preserves").await.unwrap();
let after = list_cron_jobs(&pool)
.await
.unwrap()
.into_iter()
.find(|r| r.name == "pause_preserves")
.unwrap();
assert_eq!(
after.last_enqueued_at,
Some(last_before),
"last_enqueued_at must survive a pause/resume cycle so missed_fire_policy decides catch-up"
);
}
#[tokio::test]
async fn test_upsert_preserves_pause_state() {
let pool = setup().await;
clean_cron_names(&pool, &["upsert_preserves_pause"]).await;
let job = PeriodicJob::builder("upsert_preserves_pause", "0 9 * * *")
.build_raw("pause_test".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
pause_cron_job(&pool, "upsert_preserves_pause", Some("ops"))
.await
.unwrap();
let job_v2 = PeriodicJob::builder("upsert_preserves_pause", "30 8 * * *")
.build_raw(
"pause_test".to_string(),
serde_json::json!({"version": "v2"}),
)
.unwrap();
upsert_cron_job(&pool, &job_v2).await.unwrap();
let row = list_cron_jobs(&pool)
.await
.unwrap()
.into_iter()
.find(|r| r.name == "upsert_preserves_pause")
.unwrap();
assert!(
row.is_paused(),
"upsert must not clear paused state — operators expect pause to survive deploys"
);
assert_eq!(row.paused_by.as_deref(), Some("ops"));
assert_eq!(row.cron_expr, "30 8 * * *");
assert_eq!(row.args, serde_json::json!({"version": "v2"}));
}
#[tokio::test]
async fn test_trigger_works_on_paused_schedule() {
let pool = setup().await;
clean_cron_names(&pool, &["trigger_on_paused"]).await;
let queue = "cron_trigger_paused";
clean_queue(&pool, queue).await;
let job = PeriodicJob::builder("trigger_on_paused", "0 9 * * *")
.queue(queue)
.build_raw("pause_test".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
pause_cron_job(&pool, "trigger_on_paused", None)
.await
.unwrap();
let job_row = trigger_cron_job(&pool, "trigger_on_paused").await.unwrap();
assert_eq!(job_row.queue, queue);
assert_eq!(
job_row
.metadata
.get("triggered_manually")
.and_then(|v| v.as_bool()),
Some(true)
);
let row = list_cron_jobs(&pool)
.await
.unwrap()
.into_iter()
.find(|r| r.name == "trigger_on_paused")
.unwrap();
assert!(row.is_paused());
}
#[tokio::test]
async fn test_resume_unknown_returns_false() {
let pool = setup().await;
clean_cron_names(&pool, &["resume_unknown"]).await;
let updated = resume_cron_job(&pool, "resume_unknown").await.unwrap();
assert!(
!updated,
"resume on a non-existent schedule should report no row updated"
);
}
#[tokio::test]
async fn test_delete_clears_paused_row() {
let pool = setup().await;
clean_cron_names(&pool, &["delete_paused"]).await;
let job = PeriodicJob::builder("delete_paused", "0 9 * * *")
.build_raw("pause_test".to_string(), serde_json::json!({}))
.unwrap();
upsert_cron_job(&pool, &job).await.unwrap();
pause_cron_job(&pool, "delete_paused", None).await.unwrap();
let deleted = delete_cron_job(&pool, "delete_paused").await.unwrap();
assert!(deleted);
let still_there = list_cron_jobs(&pool)
.await
.unwrap()
.into_iter()
.any(|r| r.name == "delete_paused");
assert!(!still_there);
}