use crate::sql::Row;
use serde_json::Value;
use tracing::{debug, info, warn};
use uuid::Uuid;
use crate::db::Database;
use crate::nonce::now_secs;
use crate::status::JobStatus;
#[derive(Debug, Clone)]
pub struct Job {
pub id: Uuid,
pub kind: String,
pub dedup_key: String,
pub payload: Value,
pub status: String,
pub run_at: i64,
pub attempts: i64,
pub max_attempts: i64,
pub deadline: Option<i64>,
pub lease_until: Option<i64>,
pub lease_owner: Option<String>,
pub last_error: Option<String>,
pub created_at: i64,
pub updated_at: i64,
}
const COLUMNS: &str = "id, kind, dedup_key, payload, status, run_at, attempts, max_attempts, \
deadline, lease_until, lease_owner, last_error, created_at, updated_at";
#[derive(Debug)]
pub struct NewJob<'a> {
pub id: Uuid,
pub kind: &'a str,
pub dedup_key: &'a str,
pub payload: &'a Value,
pub run_at: i64,
pub deadline: Option<i64>,
pub max_attempts: i64,
}
#[derive(Debug, Clone, Default)]
pub struct JobQuery {
pub kind: Option<String>,
pub status: Option<crate::status::JobStatus>,
pub limit: i64,
pub offset: i64,
}
impl JobQuery {
fn push_predicates(&self, builder: &mut crate::sql::Builder) {
crate::query::push_equalities(
builder,
crate::query::WHERE,
&[
("kind = ", self.kind.as_deref()),
("status = ", self.status.map(JobStatus::as_str)),
],
);
}
}
impl Job {
fn from_row(row: Row) -> Result<Self, sqlx::Error> {
let payload_json: String = row.try_get("payload")?;
let payload: Value = serde_json::from_str(&payload_json)
.map_err(|error| sqlx::Error::Decode(Box::new(error)))?;
Ok(Self {
id: row.try_get("id")?,
kind: row.try_get("kind")?,
dedup_key: row.try_get("dedup_key")?,
payload,
status: row.try_get("status")?,
run_at: row.try_get("run_at")?,
attempts: row.try_get("attempts")?,
max_attempts: row.try_get("max_attempts")?,
deadline: row.try_get("deadline")?,
lease_until: row.try_get("lease_until")?,
lease_owner: row.try_get("lease_owner")?,
last_error: row.try_get("last_error")?,
created_at: row.try_get("created_at")?,
updated_at: row.try_get("updated_at")?,
})
}
pub async fn enqueue(row: NewJob<'_>, database: &Database) -> Result<bool, sqlx::Error> {
Self::enqueue_on(row, database).await
}
pub async fn enqueue_on<'e>(
row: NewJob<'_>,
executor: impl Into<crate::sql::Exec<'e>>,
) -> Result<bool, sqlx::Error> {
let now = now_secs();
let queued = crate::sql::query(
"INSERT INTO jobs \
(id, kind, dedup_key, payload, status, run_at, attempts, max_attempts, \
deadline, created_at, updated_at) \
VALUES (?, ?, ?, ?, 'ready', ?, 0, ?, ?, ?, ?) \
ON CONFLICT DO NOTHING;",
)
.bind(row.id)
.bind(row.kind)
.bind(row.dedup_key)
.bind(row.payload.to_string())
.bind(row.run_at)
.bind(row.max_attempts)
.bind(row.deadline)
.bind(now)
.bind(now)
.execute(executor)
.await?
.rows_affected()
== 1;
if queued {
debug!(
event = "db_job_enqueued",
outcome = "success",
job_id = %row.id,
job_kind = %row.kind,
dedup_key = %row.dedup_key,
run_at = row.run_at,
);
}
Ok(queued)
}
pub async fn claim_next(
runner_id: &str,
kinds: &[&str],
lease_until: i64,
now: i64,
database: &Database,
) -> Result<Option<Self>, sqlx::Error> {
if kinds.is_empty() {
return Ok(None);
}
let placeholders = std::iter::repeat_n("?", kinds.len())
.collect::<Vec<_>>()
.join(", ");
let sql = format!(
"UPDATE jobs \
SET status = 'running', attempts = attempts + 1, lease_owner = ?, \
lease_until = ?, updated_at = ? \
WHERE id = (SELECT id FROM jobs \
WHERE status = 'ready' AND run_at <= ? AND kind IN ({placeholders}) \
ORDER BY run_at ASC, created_at ASC LIMIT 1) \
AND status = 'ready' \
RETURNING {COLUMNS};"
);
let mut query = crate::sql::query(sqlx::AssertSqlSafe(sql))
.bind(runner_id)
.bind(lease_until)
.bind(now)
.bind(now);
for kind in kinds {
query = query.bind(*kind);
}
let row = query.fetch_optional(database).await?;
let job = row.map(Self::from_row).transpose()?;
if let Some(job) = &job {
debug!(
event = "db_job_claimed",
outcome = "success",
job_id = %job.id,
job_kind = %job.kind,
attempts = job.attempts,
lease_until = lease_until,
);
}
Ok(job)
}
pub async fn complete(
id: Uuid,
runner_id: &str,
database: &Database,
) -> Result<bool, sqlx::Error> {
Self::settle(id, runner_id, "done", None, None, false, database).await
}
pub async fn retry(
id: Uuid,
runner_id: &str,
run_at: i64,
error: &str,
database: &Database,
) -> Result<bool, sqlx::Error> {
Self::settle(
id,
runner_id,
"ready",
Some(run_at),
Some(error),
false,
database,
)
.await
}
pub async fn reschedule(
id: Uuid,
runner_id: &str,
run_at: i64,
database: &Database,
) -> Result<bool, sqlx::Error> {
Self::settle(id, runner_id, "ready", Some(run_at), None, true, database).await
}
pub async fn abandon(
id: Uuid,
runner_id: &str,
error: &str,
database: &Database,
) -> Result<bool, sqlx::Error> {
Self::settle(id, runner_id, "failed", None, Some(error), false, database).await
}
async fn settle(
id: Uuid,
runner_id: &str,
status: &str,
run_at: Option<i64>,
error: Option<&str>,
reset_attempts: bool,
database: &Database,
) -> Result<bool, sqlx::Error> {
let now = now_secs();
let sql = format!(
"UPDATE jobs \
SET status = ?, last_error = ?, lease_owner = NULL, lease_until = NULL, \
run_at = COALESCE(?, run_at), updated_at = ?{} \
WHERE id = ? AND status = 'running' AND lease_owner = ?;",
if reset_attempts { ", attempts = 0" } else { "" }
);
let settled = crate::sql::query(sqlx::AssertSqlSafe(sql))
.bind(status)
.bind(error)
.bind(run_at)
.bind(now)
.bind(id)
.bind(runner_id)
.execute(database)
.await?
.rows_affected()
== 1;
Ok(settled)
}
pub async fn reclaim_expired(now: i64, database: &Database) -> Result<u64, sqlx::Error> {
let reclaimed = crate::sql::query(
"UPDATE jobs SET status = 'ready', lease_owner = NULL, lease_until = NULL, \
updated_at = ? WHERE status = 'running' AND lease_until <= ?;",
)
.bind(now)
.bind(now)
.execute(database)
.await?
.rows_affected();
if reclaimed > 0 {
warn!(
event = "db_job_leases_reclaimed",
outcome = "advisory",
rows_reclaimed = reclaimed,
);
}
Ok(reclaimed)
}
pub async fn release_owned(runner_id: &str, database: &Database) -> Result<u64, sqlx::Error> {
let released = crate::sql::query(
"UPDATE jobs SET status = 'ready', lease_owner = NULL, lease_until = NULL, \
updated_at = ? WHERE status = 'running' AND lease_owner = ?;",
)
.bind(now_secs())
.bind(runner_id)
.execute(database)
.await?
.rows_affected();
if released > 0 {
debug!(
event = "db_job_leases_released",
outcome = "success",
rows_released = released,
);
}
Ok(released)
}
pub async fn find_by_id(id: Uuid, database: &Database) -> Result<Option<Self>, sqlx::Error> {
let mut query = crate::sql::Builder::new(
database.dialect(),
format!("SELECT {COLUMNS} FROM jobs WHERE id = "),
);
query.push_bind(id);
let row = query.build().fetch_optional(database).await?;
row.map(Self::from_row).transpose()
}
pub async fn find_live(
kind: &str,
dedup_key: &str,
database: &Database,
) -> Result<Option<Self>, sqlx::Error> {
let mut query = crate::sql::Builder::new(
database.dialect(),
format!(
"SELECT {COLUMNS} FROM jobs \
WHERE status IN ('ready', 'running') AND kind = "
),
);
query.push_bind(kind);
query.push(" AND dedup_key = ");
query.push_bind(dedup_key);
let row = query.build().fetch_optional(database).await?;
row.map(Self::from_row).transpose()
}
pub async fn count_live(kind: &str, database: &Database) -> Result<i64, sqlx::Error> {
let count: i64 = crate::sql::query(
"SELECT COUNT(*) FROM jobs WHERE status IN ('ready', 'running') AND kind = ?;",
)
.bind(kind)
.fetch_one(database)
.await?
.try_get(0usize)?;
Ok(count)
}
pub async fn search(
query: &JobQuery,
database: &Database,
) -> Result<(Vec<Self>, i64), sqlx::Error> {
debug!(
event = "db_job_search_started",
outcome = "progress",
job_kind = ?query.kind,
status = ?query.status,
limit = query.limit,
offset = query.offset,
);
let mut page =
crate::sql::Builder::new(database.dialect(), format!("SELECT {COLUMNS} FROM jobs"));
query.push_predicates(&mut page);
page.push(" ORDER BY created_at DESC, id DESC LIMIT ");
page.push_bind(query.limit);
page.push(" OFFSET ");
page.push_bind(query.offset);
let rows = page.build().fetch_all(database).await?;
let jobs: Vec<Self> = rows
.into_iter()
.map(Self::from_row)
.collect::<Result<_, _>>()?;
let mut count = crate::sql::Builder::new(database.dialect(), "SELECT COUNT(*) FROM jobs");
query.push_predicates(&mut count);
let total: i64 = count.build().fetch_one(database).await?.try_get::<i64>(0)?;
Ok((jobs, total))
}
pub async fn find_latest_by_dedup(
kind: &str,
dedup_key: &str,
database: &Database,
) -> Result<Option<Self>, sqlx::Error> {
let mut query = crate::sql::Builder::new(
database.dialect(),
format!("SELECT {COLUMNS} FROM jobs WHERE kind = "),
);
query.push_bind(kind);
query.push(" AND dedup_key = ");
query.push_bind(dedup_key);
query.push(" ORDER BY created_at DESC, id DESC LIMIT 1");
let row = query.build().fetch_optional(database).await?;
row.map(Self::from_row).transpose()
}
pub async fn cancel_row(
id: Uuid,
from: JobStatus,
database: &Database,
) -> Result<Option<Self>, sqlx::Error> {
debug_assert_ne!(
from,
JobStatus::Running,
"a running job is the runner's; cancelling it would strand a lease"
);
let sql = format!(
"UPDATE jobs \
SET status = 'cancelled', updated_at = ?, lease_owner = NULL, lease_until = NULL \
WHERE id = ? AND status = ? \
RETURNING {COLUMNS};"
);
let row = crate::sql::query(sqlx::AssertSqlSafe(sql))
.bind(now_secs())
.bind(id)
.bind(from.as_str())
.fetch_optional(database)
.await?;
let job = row.map(Self::from_row).transpose()?;
if job.is_some() {
info!(event = "db_job_cancelled", outcome = "success", job_id = %id, from = from.as_str());
}
Ok(job)
}
pub async fn advance_row(id: Uuid, database: &Database) -> Result<Option<Self>, sqlx::Error> {
let now = now_secs();
let sql = format!(
"UPDATE jobs SET run_at = ?, updated_at = ? \
WHERE id = ? AND status = 'ready' RETURNING {COLUMNS};"
);
let row = crate::sql::query(sqlx::AssertSqlSafe(sql))
.bind(now)
.bind(now)
.bind(id)
.fetch_optional(database)
.await?;
let job = row.map(Self::from_row).transpose()?;
if job.is_some() {
info!(event = "db_job_advanced", outcome = "success", job_id = %id);
}
Ok(job)
}
pub async fn revive_row(id: Uuid, database: &Database) -> Result<Option<Self>, sqlx::Error> {
let now = now_secs();
let sql = format!(
"UPDATE jobs \
SET status = 'ready', run_at = ?, updated_at = ?, \
attempts = CASE WHEN max_attempts - 1 > 0 \
THEN max_attempts - 1 ELSE 0 END, \
lease_owner = NULL, lease_until = NULL \
WHERE id = ? AND status = 'failed' RETURNING {COLUMNS};"
);
let row = crate::sql::query(sqlx::AssertSqlSafe(sql))
.bind(now)
.bind(now)
.bind(id)
.fetch_optional(database)
.await?;
let job = row.map(Self::from_row).transpose()?;
if let Some(job) = &job {
info!(
event = "db_job_revived",
outcome = "success",
job_id = %id,
attempts = job.attempts,
);
}
Ok(job)
}
pub async fn cleanup(cutoff: i64, database: &Database) -> Result<u64, sqlx::Error> {
let deleted = crate::sql::query(
"DELETE FROM jobs \
WHERE status IN ('done', 'failed', 'cancelled') AND updated_at < ?;",
)
.bind(cutoff)
.execute(database)
.await?
.rows_affected();
info!(
event = "db_job_cleanup_completed",
outcome = "success",
rows_removed = deleted,
cutoff = cutoff,
);
Ok(deleted)
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use std::sync::Arc;
async fn db() -> Arc<Database> {
Arc::new(Database::connect_for_test().await.unwrap())
}
fn job_id(name: &str) -> Uuid {
let mut bytes = [0u8; 16];
let name = name.as_bytes();
let take = name.len().min(16);
bytes[..take].copy_from_slice(&name[..take]);
Uuid::from_bytes(bytes)
}
async fn enqueue(id: Uuid, key: &str, run_at: i64, database: &Database) -> bool {
Job::enqueue(
NewJob {
id,
kind: "test",
dedup_key: key,
payload: &json!({"n": 1}),
run_at,
deadline: None,
max_attempts: 3,
},
database,
)
.await
.unwrap()
}
async fn claim(runner: &str, database: &Database) -> Option<Job> {
Job::claim_next(runner, &["test"], now_secs() + 60, now_secs(), database)
.await
.unwrap()
}
#[tokio::test]
async fn a_queued_job_round_trips_through_every_column() {
let database = db().await;
let deadline = now_secs() + 900;
assert!(
Job::enqueue(
NewJob {
id: job_id("job-1"),
kind: "relay",
dedup_key: "ord-1",
payload: &json!({"order_id": "ord-1"}),
run_at: 1_234,
deadline: Some(deadline),
max_attempts: 7,
},
&database,
)
.await
.unwrap()
);
let job = Job::find_by_id(job_id("job-1"), &database)
.await
.unwrap()
.unwrap();
assert_eq!(job.kind, "relay");
assert_eq!(job.dedup_key, "ord-1");
assert_eq!(job.payload, json!({"order_id": "ord-1"}));
assert_eq!(job.status, "ready");
assert_eq!(job.run_at, 1_234);
assert_eq!(job.attempts, 0);
assert_eq!(job.max_attempts, 7);
assert_eq!(job.deadline, Some(deadline));
assert!(job.lease_until.is_none());
assert!(job.lease_owner.is_none());
assert!(job.last_error.is_none());
}
#[tokio::test]
async fn a_live_job_holds_its_identity_and_a_settled_one_releases_it() {
let database = db().await;
assert!(enqueue(job_id("job-1"), "ord-1", now_secs(), &database).await);
assert!(
!enqueue(job_id("job-2"), "ord-1", now_secs(), &database).await,
"a second live job must not take an identity already held"
);
let job = claim("runner-a", &database).await.unwrap();
assert!(
!enqueue(job_id("job-3"), "ord-1", now_secs(), &database).await,
"'running' holds the identity exactly as 'ready' does"
);
Job::complete(job.id, "runner-a", &database).await.unwrap();
assert!(
enqueue(job_id("job-4"), "ord-1", now_secs(), &database).await,
"a settled job releases its identity"
);
}
#[tokio::test]
async fn claiming_takes_the_oldest_eligible_row_and_only_once() {
let database = db().await;
assert!(enqueue(job_id("job-new"), "b", now_secs() - 10, &database).await);
assert!(enqueue(job_id("job-old"), "a", now_secs() - 100, &database).await);
let first = claim("runner-a", &database).await.unwrap();
assert_eq!(first.id, job_id("job-old"), "oldest `run_at` first");
assert_eq!(first.attempts, 1, "the counter moves at claim");
assert_eq!(first.lease_owner.as_deref(), Some("runner-a"));
let second = claim("runner-b", &database).await.unwrap();
assert_eq!(second.id, job_id("job-new"));
assert!(
claim("runner-c", &database).await.is_none(),
"a claimed row is not claimable again"
);
}
#[tokio::test]
async fn claiming_skips_a_job_whose_run_at_has_not_arrived() {
let database = db().await;
assert!(enqueue(job_id("job-1"), "a", now_secs() + 3_600, &database).await);
assert!(claim("runner-a", &database).await.is_none());
}
#[tokio::test]
async fn claiming_skips_a_kind_the_runner_does_not_hold() {
let database = db().await;
assert!(enqueue(job_id("job-1"), "a", now_secs(), &database).await);
let claimed = Job::claim_next(
"runner-a",
&["something-else"],
now_secs() + 60,
now_secs(),
&database,
)
.await
.unwrap();
assert!(
claimed.is_none(),
"an unregistered kind is left alone, not mis-run"
);
}
#[tokio::test]
async fn claiming_with_no_registered_kinds_asks_the_database_nothing() {
let database = db().await;
assert!(enqueue(job_id("job-1"), "a", now_secs(), &database).await);
assert!(
Job::claim_next("runner-a", &[], now_secs() + 60, now_secs(), &database)
.await
.unwrap()
.is_none()
);
}
#[tokio::test]
async fn a_settlement_from_a_runner_that_lost_the_lease_writes_nothing() {
let database = db().await;
assert!(enqueue(job_id("job-1"), "a", now_secs(), &database).await);
let job = claim("runner-a", &database).await.unwrap();
for settled in [
Job::complete(job.id, "runner-b", &database).await.unwrap(),
Job::retry(job.id, "runner-b", now_secs(), "x", &database)
.await
.unwrap(),
Job::reschedule(job.id, "runner-b", now_secs(), &database)
.await
.unwrap(),
Job::abandon(job.id, "runner-b", "x", &database)
.await
.unwrap(),
] {
assert!(!settled, "the lease guard must refuse a foreign settlement");
}
let after = Job::find_by_id(job_id("job-1"), &database)
.await
.unwrap()
.unwrap();
assert_eq!(after.status, "running");
assert_eq!(after.lease_owner.as_deref(), Some("runner-a"));
}
#[tokio::test]
async fn retry_returns_the_job_at_its_new_time_and_keeps_the_attempt_count() {
let database = db().await;
assert!(enqueue(job_id("job-1"), "a", now_secs(), &database).await);
let job = claim("runner-a", &database).await.unwrap();
assert!(
Job::retry(job.id, "runner-a", 9_999, "upstream timed out", &database)
.await
.unwrap()
);
let after = Job::find_by_id(job_id("job-1"), &database)
.await
.unwrap()
.unwrap();
assert_eq!(after.status, "ready");
assert_eq!(after.run_at, 9_999);
assert_eq!(after.attempts, 1, "a retry does not forgive the attempt");
assert_eq!(after.last_error.as_deref(), Some("upstream timed out"));
assert!(after.lease_owner.is_none());
}
#[tokio::test]
async fn reschedule_resets_the_attempt_count_and_clears_the_last_error() {
let database = db().await;
assert!(enqueue(job_id("job-1"), "a", now_secs(), &database).await);
let job = claim("runner-a", &database).await.unwrap();
Job::retry(job.id, "runner-a", now_secs(), "a failure", &database)
.await
.unwrap();
let job = claim("runner-a", &database).await.unwrap();
assert_eq!(job.attempts, 2);
assert!(
Job::reschedule(job.id, "runner-a", 5_000, &database)
.await
.unwrap()
);
let after = Job::find_by_id(job_id("job-1"), &database)
.await
.unwrap()
.unwrap();
assert_eq!(after.status, "ready");
assert_eq!(after.run_at, 5_000);
assert_eq!(after.attempts, 0, "a fresh occurrence starts fresh");
assert!(after.last_error.is_none());
}
#[tokio::test]
async fn abandon_is_terminal_and_records_why() {
let database = db().await;
assert!(enqueue(job_id("job-1"), "a", now_secs(), &database).await);
let job = claim("runner-a", &database).await.unwrap();
assert!(
Job::abandon(job.id, "runner-a", "the order no longer exists", &database)
.await
.unwrap()
);
let after = Job::find_by_id(job_id("job-1"), &database)
.await
.unwrap()
.unwrap();
assert_eq!(after.status, "failed");
assert_eq!(
after.last_error.as_deref(),
Some("the order no longer exists")
);
}
#[tokio::test]
async fn reclaim_takes_only_leases_that_have_expired() {
let database = db().await;
assert!(enqueue(job_id("job-live"), "a", now_secs(), &database).await);
assert!(enqueue(job_id("job-dead"), "b", now_secs(), &database).await);
Job::claim_next(
"runner-a",
&["test"],
now_secs() + 600,
now_secs(),
&database,
)
.await
.unwrap();
Job::claim_next("runner-a", &["test"], now_secs() - 1, now_secs(), &database)
.await
.unwrap();
assert_eq!(
Job::reclaim_expired(now_secs(), &database).await.unwrap(),
1
);
let reclaimed = Job::find_by_id(job_id("job-dead"), &database)
.await
.unwrap()
.unwrap();
assert_eq!(reclaimed.status, "ready");
assert!(reclaimed.lease_owner.is_none());
assert_eq!(
reclaimed.attempts, 1,
"the attempt was spent, and the row must keep saying so"
);
let held = Job::find_by_id(job_id("job-live"), &database)
.await
.unwrap()
.unwrap();
assert_eq!(held.status, "running");
}
#[tokio::test]
async fn release_owned_is_scoped_to_one_runner() {
let database = db().await;
assert!(enqueue(job_id("job-1"), "a", now_secs(), &database).await);
assert!(enqueue(job_id("job-2"), "b", now_secs(), &database).await);
claim("runner-a", &database).await.unwrap();
claim("runner-b", &database).await.unwrap();
assert_eq!(Job::release_owned("runner-a", &database).await.unwrap(), 1);
assert_eq!(
Job::find_by_id(job_id("job-1"), &database)
.await
.unwrap()
.unwrap()
.status,
"ready"
);
assert_eq!(
Job::find_by_id(job_id("job-2"), &database)
.await
.unwrap()
.unwrap()
.status,
"running"
);
}
#[tokio::test]
async fn find_live_sees_a_queued_job_and_not_a_settled_one() {
let database = db().await;
assert!(enqueue(job_id("job-1"), "ord-1", now_secs(), &database).await);
assert!(
Job::find_live("test", "ord-1", &database)
.await
.unwrap()
.is_some()
);
let job = claim("runner-a", &database).await.unwrap();
Job::complete(job.id, "runner-a", &database).await.unwrap();
assert!(
Job::find_live("test", "ord-1", &database)
.await
.unwrap()
.is_none()
);
}
#[tokio::test]
async fn cleanup_removes_settled_rows_strictly_older_than_the_cutoff() {
let database = db().await;
assert!(enqueue(job_id("job-done"), "a", now_secs(), &database).await);
assert!(enqueue(job_id("job-live"), "b", now_secs() + 3_600, &database).await);
let job = claim("runner-a", &database).await.unwrap();
Job::complete(job.id, "runner-a", &database).await.unwrap();
let settled = Job::find_by_id(job_id("job-done"), &database)
.await
.unwrap()
.unwrap();
assert_eq!(
Job::cleanup(settled.updated_at, &database).await.unwrap(),
0
);
assert_eq!(
Job::cleanup(settled.updated_at + 1, &database)
.await
.unwrap(),
1
);
assert!(
Job::find_by_id(job_id("job-done"), &database)
.await
.unwrap()
.is_none()
);
assert!(
Job::find_by_id(job_id("job-live"), &database)
.await
.unwrap()
.is_some(),
"a job still queued is not old, however long ago it was written"
);
}
#[tokio::test]
async fn an_unknown_status_is_refused_by_the_check_constraint() {
let database = db().await;
let error = crate::sql::query(
"INSERT INTO jobs (id, kind, dedup_key, payload, status, run_at, attempts, \
max_attempts, created_at, updated_at) \
VALUES (?, 'test', 'k', '{}', 'halfway', 0, 0, 1, 0, 0);",
)
.bind(crate::id::mint())
.execute(&database)
.await
.unwrap_err();
assert!(crate::sql::is_check_violation(&error), "{error}");
}
use crate::status::JobStatus;
async fn enqueue_kind(
id: Uuid,
kind: &str,
key: &str,
run_at: i64,
max_attempts: i64,
database: &Database,
) -> bool {
Job::enqueue(
NewJob {
id,
kind,
dedup_key: key,
payload: &json!({}),
run_at,
deadline: None,
max_attempts,
},
database,
)
.await
.unwrap()
}
async fn failed_job_max(
id: Uuid,
kind: &str,
key: &str,
max_attempts: i64,
database: &Database,
) -> Job {
assert!(enqueue_kind(id, kind, key, now_secs(), max_attempts, database).await);
crate::sql::query(
"UPDATE jobs SET status = 'failed', attempts = 1, \
last_error = 'upstream said no', updated_at = ? WHERE id = ?;",
)
.bind(now_secs())
.bind(id)
.execute(database)
.await
.unwrap();
Job::find_by_id(id, database).await.unwrap().unwrap()
}
async fn failed_job(id: Uuid, kind: &str, key: &str, database: &Database) -> Job {
failed_job_max(id, kind, key, 3, database).await
}
#[tokio::test]
async fn search_filters_kind_and_status_and_pages_without_overlap() {
let database = db().await;
for i in 0..3 {
assert!(
enqueue_kind(
job_id(&format!("relay-{i}")),
"relay",
&format!("ord-{i}"),
now_secs() - i64::from(10 - i),
3,
&database,
)
.await
);
}
for i in 0..2 {
assert!(
enqueue_kind(
job_id(&format!("sweep-{i}")),
"sweep",
&format!("s-{i}"),
now_secs(),
3,
&database,
)
.await
);
}
failed_job(job_id("relay-failed"), "relay", "ord-f", &database).await;
let (rows, total) = Job::search(
&JobQuery {
kind: Some("relay".to_string()),
limit: 50,
..JobQuery::default()
},
&database,
)
.await
.unwrap();
assert_eq!(total, 4);
assert_eq!(rows.len(), 4);
let (rows, total) = Job::search(
&JobQuery {
kind: Some("relay".to_string()),
status: Some(JobStatus::Failed),
limit: 50,
offset: 0,
},
&database,
)
.await
.unwrap();
assert_eq!(total, 1);
assert_eq!(rows[0].id, job_id("relay-failed"));
let mut seen = std::collections::BTreeSet::new();
for offset in 0..4 {
let (rows, total) = Job::search(
&JobQuery {
kind: Some("relay".to_string()),
limit: 1,
offset,
..JobQuery::default()
},
&database,
)
.await
.unwrap();
assert_eq!(total, 4);
assert_eq!(rows.len(), 1);
assert!(seen.insert(rows[0].id), "a row appeared on two pages");
}
assert_eq!(seen.len(), 4);
}
#[tokio::test]
async fn a_kind_filter_value_is_bound_not_interpolated() {
let database = db().await;
assert!(enqueue_kind(job_id("job-1"), "relay", "k", now_secs(), 3, &database).await);
let (rows, total) = Job::search(
&JobQuery {
kind: Some("' OR 1=1 --".to_string()),
limit: 50,
..JobQuery::default()
},
&database,
)
.await
.unwrap();
assert_eq!(total, 0);
assert!(rows.is_empty());
assert!(
Job::find_by_id(job_id("job-1"), &database)
.await
.unwrap()
.is_some()
);
}
#[tokio::test]
async fn cancel_row_moves_ready_and_failed_and_refuses_the_rest() {
let database = db().await;
assert!(enqueue_kind(job_id("ready"), "test", "a", now_secs(), 3, &database).await);
let cancelled = Job::cancel_row(job_id("ready"), JobStatus::Ready, &database)
.await
.unwrap()
.expect("a ready job cancels");
assert_eq!(cancelled.status, "cancelled");
assert!(
Job::cancel_row(job_id("ready"), JobStatus::Ready, &database)
.await
.unwrap()
.is_none()
);
failed_job(job_id("failed"), "test", "b", &database).await;
assert!(
Job::cancel_row(job_id("failed"), JobStatus::Failed, &database)
.await
.unwrap()
.is_some()
);
assert!(enqueue_kind(job_id("running"), "test", "c", now_secs(), 3, &database).await);
Job::claim_next("r", &["test"], now_secs() + 60, now_secs(), &database)
.await
.unwrap();
assert!(
Job::cancel_row(job_id("running"), JobStatus::Ready, &database)
.await
.unwrap()
.is_none()
);
assert_eq!(
Job::find_by_id(job_id("running"), &database)
.await
.unwrap()
.unwrap()
.status,
"running"
);
assert!(enqueue_kind(job_id("done"), "test", "d", now_secs(), 3, &database).await);
let job = Job::claim_next("r", &["test"], now_secs() + 60, now_secs(), &database)
.await
.unwrap()
.unwrap();
Job::complete(job.id, "r", &database).await.unwrap();
assert!(
Job::cancel_row(job_id("done"), JobStatus::Ready, &database)
.await
.unwrap()
.is_none()
);
}
#[tokio::test]
async fn advance_row_nudges_a_ready_job_and_refuses_a_running_one() {
let database = db().await;
assert!(
enqueue_kind(
job_id("ready"),
"test",
"a",
now_secs() + 3_600,
3,
&database
)
.await
);
let advanced = Job::advance_row(job_id("ready"), &database)
.await
.unwrap()
.expect("a ready job advances");
assert!(advanced.run_at <= now_secs() + 1);
assert_eq!(advanced.attempts, 0, "advancing does not spend an attempt");
Job::claim_next("r", &["test"], now_secs() + 60, now_secs(), &database)
.await
.unwrap();
assert!(
Job::advance_row(job_id("ready"), &database)
.await
.unwrap()
.is_none()
);
failed_job(job_id("failed"), "test", "b", &database).await;
assert!(
Job::advance_row(job_id("failed"), &database)
.await
.unwrap()
.is_none()
);
}
#[tokio::test]
async fn revive_row_sets_attempts_to_max_minus_one_and_refuses_a_live_job() {
let database = db().await;
let failed = failed_job(job_id("failed"), "test", "a", &database).await;
assert_eq!(failed.max_attempts, 3);
let revived = Job::revive_row(job_id("failed"), &database)
.await
.unwrap()
.expect("a failed job revives");
assert_eq!(revived.status, "ready");
assert_eq!(revived.attempts, 2, "exactly one attempt left");
assert!(revived.run_at <= now_secs() + 1);
assert_eq!(
revived.last_error.as_deref(),
Some("upstream said no"),
"the last error stays as the record of why it stopped"
);
assert!(enqueue_kind(job_id("ready"), "test", "b", now_secs(), 3, &database).await);
assert!(
Job::revive_row(job_id("ready"), &database)
.await
.unwrap()
.is_none()
);
let one = failed_job_max(job_id("one"), "test", "c", 1, &database).await;
assert_eq!(one.max_attempts, 1);
let revived = Job::revive_row(job_id("one"), &database)
.await
.unwrap()
.unwrap();
assert_eq!(revived.attempts, 0);
}
#[tokio::test]
async fn find_latest_by_dedup_sees_a_settled_job_where_find_live_does_not() {
let database = db().await;
assert!(enqueue_kind(job_id("job-1"), "relay", "ord-1", now_secs(), 3, &database).await);
let job = Job::claim_next("r", &["relay"], now_secs() + 60, now_secs(), &database)
.await
.unwrap()
.unwrap();
Job::complete(job.id, "r", &database).await.unwrap();
assert!(
Job::find_live("relay", "ord-1", &database)
.await
.unwrap()
.is_none()
);
let latest = Job::find_latest_by_dedup("relay", "ord-1", &database)
.await
.unwrap()
.expect("the settled job is still resolvable");
assert_eq!(latest.id, job_id("job-1"));
assert_eq!(latest.status, "done");
}
#[tokio::test]
async fn revive_then_a_second_failure_stays_failed_without_another_nudge() {
let database = db().await;
failed_job(job_id("j"), "test", "a", &database).await;
Job::revive_row(job_id("j"), &database)
.await
.unwrap()
.unwrap();
let job = Job::claim_next("r", &["test"], now_secs() + 60, now_secs(), &database)
.await
.unwrap()
.unwrap();
assert_eq!(job.attempts, 3);
Job::abandon(job.id, "r", "again", &database).await.unwrap();
let after = Job::find_by_id(job_id("j"), &database)
.await
.unwrap()
.unwrap();
assert_eq!(after.status, "failed");
assert!(
Job::advance_row(job_id("j"), &database)
.await
.unwrap()
.is_none()
);
assert!(
Job::revive_row(job_id("j"), &database)
.await
.unwrap()
.is_some()
);
}
}