pub(crate) mod sql {
pub(crate) const SESSION_SETUP: Option<&str> = None;
pub(crate) const BEGIN: &str = "START TRANSACTION";
pub(crate) const CLAIM_PICK: &str = r#"SELECT id, kind, version, payload, attempts, max_attempts
FROM arcature_jobs
WHERE status = 'pending' AND available_at <= UTC_TIMESTAMP(6)
ORDER BY available_at, id
LIMIT ?
FOR UPDATE SKIP LOCKED"#;
pub(crate) const CLAIM_MARK: &str = r#"UPDATE arcature_jobs
SET status = 'running', attempts = attempts + 1, locked_by = ?,
locked_at = UTC_TIMESTAMP(6), claim_token = ?, lease_seconds = ?,
last_error = NULL, last_error_kind = NULL,
updated_at = UTC_TIMESTAMP(6)
WHERE id = ? AND status = 'pending'"#;
pub(crate) const INSERT: &str = r#"INSERT INTO arcature_jobs
(id, kind, version, payload, max_attempts, run_at, available_at)
VALUES (?, ?, ?, ?, ?, ?, ?)"#;
pub(crate) const MARK_SUCCEEDED: &str = r#"UPDATE arcature_jobs
SET status = 'succeeded', locked_by = NULL, locked_at = NULL,
claim_token = NULL, last_error = NULL, last_error_kind = NULL,
failed_at = NULL, updated_at = UTC_TIMESTAMP(6)
WHERE id = ? AND status = 'running' AND claim_token = ?"#;
pub(crate) const MARK_RETRY: &str = r#"UPDATE arcature_jobs
SET status = 'pending', available_at = ?, locked_by = NULL,
locked_at = NULL, claim_token = NULL, last_error = ?,
last_error_kind = 'retryable', failed_at = NULL,
updated_at = UTC_TIMESTAMP(6)
WHERE id = ? AND status = 'running' AND claim_token = ?"#;
pub(crate) const MARK_DEAD: &str = r#"UPDATE arcature_jobs
SET status = 'dead', locked_by = NULL, locked_at = NULL,
claim_token = NULL, last_error = ?, last_error_kind = ?,
failed_at = UTC_TIMESTAMP(6), updated_at = UTC_TIMESTAMP(6)
WHERE id = ? AND status = 'running' AND claim_token = ?"#;
pub(crate) const HEARTBEAT: &str = r#"UPDATE arcature_jobs
SET locked_at = UTC_TIMESTAMP(6), lease_seconds = ?,
updated_at = UTC_TIMESTAMP(6)
WHERE id = ? AND status = 'running' AND claim_token = ?"#;
pub(crate) const RELEASE_CLAIM: &str = r#"UPDATE arcature_jobs
SET status = 'pending', available_at = UTC_TIMESTAMP(6),
locked_by = NULL, locked_at = NULL, claim_token = NULL,
updated_at = UTC_TIMESTAMP(6)
WHERE id = ? AND status = 'running' AND claim_token = ?"#;
pub(crate) const SWEEP_DEAD: &str = r#"UPDATE arcature_jobs
SET status = 'dead', locked_by = NULL, locked_at = NULL,
claim_token = NULL,
last_error = 'job exhausted its max_attempts after a crash',
last_error_kind = 'exhausted', failed_at = UTC_TIMESTAMP(6),
updated_at = UTC_TIMESTAMP(6)
WHERE status = 'running'
AND locked_at IS NOT NULL
AND locked_at + INTERVAL lease_seconds SECOND < UTC_TIMESTAMP(6)
AND attempts >= max_attempts
ORDER BY locked_at
LIMIT ?"#;
pub(crate) const SWEEP_REQUEUE: &str = r#"UPDATE arcature_jobs
SET status = 'pending', available_at = UTC_TIMESTAMP(6),
locked_by = NULL, locked_at = NULL, claim_token = NULL,
updated_at = UTC_TIMESTAMP(6)
WHERE status = 'running'
AND locked_at IS NOT NULL
AND locked_at + INTERVAL lease_seconds SECOND < UTC_TIMESTAMP(6)
AND attempts < max_attempts
ORDER BY locked_at
LIMIT ?"#;
pub(crate) const CANCEL: &str = r#"UPDATE arcature_jobs
SET status = 'cancelled', locked_by = NULL, locked_at = NULL,
claim_token = NULL, updated_at = UTC_TIMESTAMP(6)
WHERE id = ? AND status IN ('pending', 'running')"#;
pub(crate) const REQUEUE_DEAD: &str = r#"UPDATE arcature_jobs
SET status = 'pending', attempts = 0,
available_at = UTC_TIMESTAMP(6), locked_by = NULL, locked_at = NULL,
claim_token = NULL, last_error = NULL, last_error_kind = NULL,
failed_at = NULL, updated_at = UTC_TIMESTAMP(6)
WHERE id = ? AND status = 'dead'"#;
pub(crate) const CREATE_HISTORY: &str = r#"CREATE TABLE IF NOT EXISTS arcature_jobs_schema_migrations (
version VARCHAR(191) NOT NULL PRIMARY KEY,
applied_at DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6)
) ENGINE=InnoDB"#;
pub(crate) const COUNT_APPLIED: &str =
"SELECT COUNT(*) FROM arcature_jobs_schema_migrations WHERE version = ?";
pub(crate) const RECORD_APPLIED: &str =
"INSERT IGNORE INTO arcature_jobs_schema_migrations (version) VALUES (?)";
pub(crate) const LOCK: Option<&str> = Some("SELECT GET_LOCK('arcature_jobs_migrate', 10)");
pub(crate) const UNLOCK: Option<&str> = Some("SELECT RELEASE_LOCK('arcature_jobs_migrate')");
pub(crate) const SCHEMA: &str = include_str!("../migrations/mysql/0001_jobs.sql");
}
#[cfg(all(test, feature = "test-kit"))]
mod tests {
use std::time::Duration;
use crate::jobs::test_support::{
JOBS, WORKERS, assert_claimed_exactly_once, drain_concurrently, enqueue, queue, rows,
};
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn every_job_goes_to_exactly_one_worker() {
let Some(fixture) = queue().await else {
return;
};
let pool = fixture.pool();
let enqueued = enqueue(pool, JOBS).await;
let claimed = drain_concurrently(pool, (JOBS / WORKERS) as i64).await;
assert_claimed_exactly_once(pool, &enqueued, &claimed).await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_single_row_batch_is_still_exclusive() {
let Some(fixture) = queue().await else {
return;
};
let pool = fixture.pool();
let enqueued = enqueue(pool, JOBS).await;
let claimed = drain_concurrently(pool, 1).await;
assert_claimed_exactly_once(pool, &enqueued, &claimed).await;
}
#[tokio::test]
async fn marking_a_row_that_is_no_longer_pending_changes_nothing() {
let Some(fixture) = queue().await else {
return;
};
let pool = fixture.pool();
let enqueued = enqueue(pool, 1).await;
let claimed = crate::jobs::admin::claim_jobs(pool, "first", Duration::from_secs(60), 1)
.await
.expect("claim");
assert_eq!(claimed.len(), 1, "the only pending job was not claimed");
let marked = sqlx::query(super::sql::CLAIM_MARK)
.bind("second")
.bind(uuid::Uuid::new_v4())
.bind(60_i32)
.bind(enqueued[0])
.execute(pool)
.await
.expect("run the mark")
.rows_affected();
assert_eq!(marked, 0, "a claimed row was marked a second time");
let observed = rows(pool).await;
assert_eq!(
observed,
vec![(enqueued[0], "running".to_owned(), 1)],
"the row moved when the mark should have been refused"
);
}
}