pub(crate) mod sql {
pub(crate) const SESSION_SETUP: Option<&str> = Some("PRAGMA busy_timeout = 5000");
pub(crate) const BEGIN: &str = "BEGIN IMMEDIATE";
pub(crate) const CLAIM_PICK: &str = r#"SELECT id, kind, version, payload, attempts, max_attempts
FROM arcature_jobs
WHERE status = 'pending'
AND available_at <= CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
ORDER BY available_at, id
LIMIT ?"#;
pub(crate) const CLAIM_MARK: &str = r#"UPDATE arcature_jobs
SET status = 'running', attempts = attempts + 1, locked_by = ?,
locked_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER),
claim_token = ?, lease_seconds = ?, last_error = NULL,
last_error_kind = NULL,
updated_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
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 = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
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 = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
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 = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER),
updated_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
WHERE id = ? AND status = 'running' AND claim_token = ?"#;
pub(crate) const HEARTBEAT: &str = r#"UPDATE arcature_jobs
SET locked_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER),
lease_seconds = ?,
updated_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
WHERE id = ? AND status = 'running' AND claim_token = ?"#;
pub(crate) const RELEASE_CLAIM: &str = r#"UPDATE arcature_jobs
SET status = 'pending',
available_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER),
locked_by = NULL, locked_at = NULL, claim_token = NULL,
updated_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
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 = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER),
updated_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
WHERE id IN (
SELECT id FROM arcature_jobs
WHERE status = 'running'
AND locked_at IS NOT NULL
AND locked_at + lease_seconds * 1000
< CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
AND attempts >= max_attempts
LIMIT ?
)"#;
pub(crate) const SWEEP_REQUEUE: &str = r#"UPDATE arcature_jobs
SET status = 'pending',
available_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER),
locked_by = NULL, locked_at = NULL, claim_token = NULL,
updated_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
WHERE id IN (
SELECT id FROM arcature_jobs
WHERE status = 'running'
AND locked_at IS NOT NULL
AND locked_at + lease_seconds * 1000
< CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
AND attempts < max_attempts
LIMIT ?
)"#;
pub(crate) const CANCEL: &str = r#"UPDATE arcature_jobs
SET status = 'cancelled', locked_by = NULL, locked_at = NULL,
claim_token = NULL,
updated_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
WHERE id = ? AND status IN ('pending', 'running')"#;
pub(crate) const REQUEUE_DEAD: &str = r#"UPDATE arcature_jobs
SET status = 'pending', attempts = 0,
available_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER),
locked_by = NULL, locked_at = NULL, claim_token = NULL,
last_error = NULL, last_error_kind = NULL, failed_at = NULL,
updated_at = CAST((julianday('now')-2440587.5)*86400000 AS INTEGER)
WHERE id = ? AND status = 'dead'"#;
pub(crate) const CREATE_HISTORY: &str = r#"CREATE TABLE IF NOT EXISTS arcature_jobs_schema_migrations (
version TEXT PRIMARY KEY,
applied_at INTEGER NOT NULL
DEFAULT (CAST((julianday('now')-2440587.5)*86400000 AS INTEGER))
)"#;
pub(crate) const COUNT_APPLIED: &str =
"SELECT COUNT(*) FROM arcature_jobs_schema_migrations WHERE version = ?";
pub(crate) const RECORD_APPLIED: &str =
"INSERT OR IGNORE INTO arcature_jobs_schema_migrations (version) VALUES (?)";
pub(crate) const LOCK: Option<&str> = None;
pub(crate) const UNLOCK: Option<&str> = None;
pub(crate) const SCHEMA: &str = include_str!("../migrations/sqlite/0001_jobs.sql");
}
#[cfg(all(test, feature = "test-kit"))]
mod tests {
use crate::jobs::test_support::{
JOBS, WORKERS, assert_claimed_exactly_once, drain_concurrently, enqueue, queue,
};
#[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 claimers_wait_for_the_write_lock_rather_than_failing() {
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;
}
}