Skip to main content

acme_proxy/sqlite/
job.rs

1//! The `jobs` model: the durable queue behind [`crate::jobs`].
2//!
3//! Every statement here is written so the *database* decides a race, never a
4//! read-then-write in Rust. Two shapes carry all of it:
5//!
6//! - **`INSERT OR IGNORE` + `rows_affected() == 0`** for enqueue, against the
7//!   partial unique index on `(kind, dedup_key)`. A `0` means a live job already
8//!   holds that identity — the same guard [`crate::sqlite::upstream_order::UpstreamOrder::create`]
9//!   takes on its primary key.
10//! - **A guarded `UPDATE … RETURNING`** for the claim, and a guarded `UPDATE`
11//!   for every settlement, each carrying `AND status = 'running' AND lease_owner
12//!   = ?`. A runner whose lease expired and was reclaimed therefore cannot
13//!   overwrite the row a second runner now owns: its write affects zero rows and
14//!   says so.
15//!
16//! See `migrations/20260815120000_add_jobs.sql` for why the table has no foreign
17//! key, why `kind` carries no `CHECK`, and why the identity index is partial.
18
19use serde_json::Value;
20use sqlx::Row;
21use sqlx::sqlite::SqliteRow;
22use tracing::{debug, info, warn};
23
24use crate::sqlite::db::Database;
25use crate::sqlite::nonce::now_secs;
26
27/// One stored job row.
28///
29/// `status` and `kind` come back as the strings they were stored as, for the
30/// reason [`crate::sqlite::audit::AuditEntry`] keeps `event` a `String`: an
31/// older binary meeting a row a newer one wrote should render it, not refuse to
32/// load. The runner never claims a `kind` its registry does not hold, so an
33/// unrecognised one is simply left alone.
34#[derive(Debug, Clone)]
35pub struct Job {
36    pub id: String,
37    pub kind: String,
38    pub dedup_key: String,
39    pub payload: Value,
40    pub status: String,
41    pub run_at: i64,
42    pub attempts: i64,
43    pub max_attempts: i64,
44    pub deadline: Option<i64>,
45    pub lease_until: Option<i64>,
46    pub lease_owner: Option<String>,
47    pub last_error: Option<String>,
48    pub created_at: i64,
49    pub updated_at: i64,
50}
51
52/// Every column, in one place: the claim's `RETURNING` and the single-row read
53/// must select the same set or `from_row` fails on whichever forgot one.
54const COLUMNS: &str = "id, kind, dedup_key, payload, status, run_at, attempts, max_attempts, \
55                       deadline, lease_until, lease_owner, last_error, created_at, updated_at";
56
57/// The columns [`Job::enqueue`] writes.
58///
59/// A struct rather than eight positional parameters — which needed
60/// `#[allow(clippy::too_many_arguments)]`, and put four `&str`/`i64` values in
61/// a row where transposing two would still compile. `JobSpec` (the
62/// `src/jobs/` half) is the caller-facing shape and this is the row it becomes;
63/// they are deliberately separate, so the storage layer names no `src/jobs/`
64/// type.
65#[derive(Debug)]
66pub struct NewJob<'a> {
67    pub id: &'a str,
68    pub kind: &'a str,
69    pub dedup_key: &'a str,
70    pub payload: &'a Value,
71    pub run_at: i64,
72    pub deadline: Option<i64>,
73    /// Frozen onto the row here rather than read per attempt, so raising
74    /// `jobs.max_attempts` applies to new work and not to a waiting backlog.
75    pub max_attempts: i64,
76}
77
78impl Job {
79    fn from_row(row: SqliteRow) -> Result<Self, sqlx::Error> {
80        let payload_json: String = row.try_get("payload")?;
81        let payload: Value = serde_json::from_str(&payload_json)
82            .map_err(|error| sqlx::Error::Decode(Box::new(error)))?;
83
84        Ok(Self {
85            id: row.try_get("id")?,
86            kind: row.try_get("kind")?,
87            dedup_key: row.try_get("dedup_key")?,
88            payload,
89            status: row.try_get("status")?,
90            run_at: row.try_get("run_at")?,
91            attempts: row.try_get("attempts")?,
92            max_attempts: row.try_get("max_attempts")?,
93            deadline: row.try_get("deadline")?,
94            lease_until: row.try_get("lease_until")?,
95            lease_owner: row.try_get("lease_owner")?,
96            last_error: row.try_get("last_error")?,
97            created_at: row.try_get("created_at")?,
98            updated_at: row.try_get("updated_at")?,
99        })
100    }
101
102    /// Queues one job, unless a live one already holds `(kind, dedup_key)`.
103    ///
104    /// `Ok(false)` is "already queued", not an error: it is what two racing
105    /// callers must both be able to survive, and the partial unique index is
106    /// what decides which of them won. A job that already reached `done` or
107    /// `failed` releases the identity, so the same key can be queued again —
108    /// which is what makes a retried order and a periodic sweep both expressible
109    /// without a second table.
110    pub async fn enqueue(row: NewJob<'_>, database: &Database) -> Result<bool, sqlx::Error> {
111        let now = now_secs();
112        let queued = sqlx::query(
113            "INSERT OR IGNORE INTO jobs \
114             (id, kind, dedup_key, payload, status, run_at, attempts, max_attempts, \
115              deadline, created_at, updated_at) \
116             VALUES (?, ?, ?, ?, 'ready', ?, 0, ?, ?, ?, ?);",
117        )
118        .bind(row.id)
119        .bind(row.kind)
120        .bind(row.dedup_key)
121        .bind(row.payload.to_string())
122        .bind(row.run_at)
123        .bind(row.max_attempts)
124        .bind(row.deadline)
125        .bind(now)
126        .bind(now)
127        .execute(&database.pool)
128        .await?
129        .rows_affected()
130            == 1;
131
132        if queued {
133            debug!(
134                event = "db_job_enqueued",
135                outcome = "success",
136                job_id = %row.id,
137                job_kind = %row.kind,
138                dedup_key = %row.dedup_key,
139                run_at = row.run_at,
140            );
141        }
142        Ok(queued)
143    }
144
145    /// Claims the oldest eligible job of one of `kinds`, or `None`.
146    ///
147    /// One statement, and that is the point: the subselect picks a candidate and
148    /// the outer `AND status = 'ready'` is what makes the pick binding, so two
149    /// runners choosing the same row end with exactly one write. `None` collapses
150    /// "the queue was empty" and "somebody else won" — which is correct, because
151    /// the caller does the same thing either way.
152    ///
153    /// `attempts` increments here rather than at completion: a job that reliably
154    /// kills the process must still exhaust its budget, and nothing reports back
155    /// from a process that died.
156    ///
157    /// **`deadline` is deliberately not filtered here.** Skipping a dead row
158    /// would leave it `ready` and re-read on every tick for ever; the runner
159    /// claims it and retires it on the spot.
160    pub async fn claim_next(
161        runner_id: &str,
162        kinds: &[&str],
163        lease_until: i64,
164        now: i64,
165        database: &Database,
166    ) -> Result<Option<Self>, sqlx::Error> {
167        if kinds.is_empty() {
168            return Ok(None);
169        }
170
171        // sqlx has no array binding for SQLite, so the IN list is built from one
172        // placeholder per kind — never from the values themselves. Same shape as
173        // `UpstreamOrder::list_processing`, and `AssertSqlSafe` for the same
174        // reason: sqlx refuses a non-`'static` query string, and the only
175        // runtime part of this one is the count of `?`.
176        let placeholders = std::iter::repeat_n("?", kinds.len())
177            .collect::<Vec<_>>()
178            .join(", ");
179        let sql = format!(
180            "UPDATE jobs \
181             SET status = 'running', attempts = attempts + 1, lease_owner = ?, \
182                 lease_until = ?, updated_at = ? \
183             WHERE id = (SELECT id FROM jobs \
184                          WHERE status = 'ready' AND run_at <= ? AND kind IN ({placeholders}) \
185                          ORDER BY run_at ASC, created_at ASC LIMIT 1) \
186               AND status = 'ready' \
187             RETURNING {COLUMNS};"
188        );
189
190        let mut query = sqlx::query(sqlx::AssertSqlSafe(sql))
191            .bind(runner_id)
192            .bind(lease_until)
193            .bind(now)
194            .bind(now);
195        for kind in kinds {
196            query = query.bind(*kind);
197        }
198        let row = query.fetch_optional(&database.pool).await?;
199        let job = row.map(Self::from_row).transpose()?;
200
201        if let Some(job) = &job {
202            debug!(
203                event = "db_job_claimed",
204                outcome = "success",
205                job_id = %job.id,
206                job_kind = %job.kind,
207                attempts = job.attempts,
208                lease_until = lease_until,
209            );
210        }
211        Ok(job)
212    }
213
214    /// Marks a claimed job finished. `Ok(false)` means the lease was lost.
215    pub async fn complete(
216        id: &str,
217        runner_id: &str,
218        database: &Database,
219    ) -> Result<bool, sqlx::Error> {
220        Self::settle(id, runner_id, "done", None, None, false, database).await
221    }
222
223    /// Returns a claimed job to the queue, to run again at `run_at`.
224    pub async fn retry(
225        id: &str,
226        runner_id: &str,
227        run_at: i64,
228        error: &str,
229        database: &Database,
230    ) -> Result<bool, sqlx::Error> {
231        Self::settle(
232            id,
233            runner_id,
234            "ready",
235            Some(run_at),
236            Some(error),
237            false,
238            database,
239        )
240        .await
241    }
242
243    /// Returns a claimed job to the queue as a *fresh* occurrence: the attempt
244    /// counter goes back to zero and the last error is cleared.
245    ///
246    /// That reset is what separates a periodic job from a retried one. A sweep
247    /// that runs every day for a year must not accumulate 365 attempts and
248    /// retire itself, and a successful occurrence must not leave the previous
249    /// failure's text sitting on the row as though it were current.
250    pub async fn reschedule(
251        id: &str,
252        runner_id: &str,
253        run_at: i64,
254        database: &Database,
255    ) -> Result<bool, sqlx::Error> {
256        Self::settle(id, runner_id, "ready", Some(run_at), None, true, database).await
257    }
258
259    /// Retires a claimed job permanently, recording why.
260    pub async fn abandon(
261        id: &str,
262        runner_id: &str,
263        error: &str,
264        database: &Database,
265    ) -> Result<bool, sqlx::Error> {
266        Self::settle(id, runner_id, "failed", None, Some(error), false, database).await
267    }
268
269    /// The one guarded write every settlement goes through.
270    ///
271    /// `AND status = 'running' AND lease_owner = ?` is the guard, and it is why
272    /// these four are one function: a settlement that forgot it would let a
273    /// runner whose lease had already been reclaimed overwrite the row a second
274    /// runner was working, and the four call sites would each have had to
275    /// remember. `rows_affected() == 1` is the caller's answer.
276    async fn settle(
277        id: &str,
278        runner_id: &str,
279        status: &str,
280        run_at: Option<i64>,
281        error: Option<&str>,
282        reset_attempts: bool,
283        database: &Database,
284    ) -> Result<bool, sqlx::Error> {
285        let now = now_secs();
286        let sql = format!(
287            "UPDATE jobs \
288             SET status = ?, last_error = ?, lease_owner = NULL, lease_until = NULL, \
289                 run_at = COALESCE(?, run_at), updated_at = ?{} \
290             WHERE id = ? AND status = 'running' AND lease_owner = ?;",
291            if reset_attempts { ", attempts = 0" } else { "" }
292        );
293        let settled = sqlx::query(sqlx::AssertSqlSafe(sql))
294            .bind(status)
295            .bind(error)
296            .bind(run_at)
297            .bind(now)
298            .bind(id)
299            .bind(runner_id)
300            .execute(&database.pool)
301            .await?
302            .rows_affected()
303            == 1;
304        Ok(settled)
305    }
306
307    /// Returns to the queue every job whose runner died holding its lease.
308    ///
309    /// `attempts` is deliberately left alone: the attempt really was spent, and
310    /// a counter rewritten to look better would let a job that crashes the
311    /// process loop for ever.
312    pub async fn reclaim_expired(now: i64, database: &Database) -> Result<u64, sqlx::Error> {
313        let reclaimed = sqlx::query(
314            "UPDATE jobs SET status = 'ready', lease_owner = NULL, lease_until = NULL, \
315             updated_at = ? WHERE status = 'running' AND lease_until <= ?;",
316        )
317        .bind(now)
318        .bind(now)
319        .execute(&database.pool)
320        .await?
321        .rows_affected();
322
323        if reclaimed > 0 {
324            warn!(
325                event = "db_job_leases_reclaimed",
326                outcome = "advisory",
327                rows_reclaimed = reclaimed,
328            );
329        }
330        Ok(reclaimed)
331    }
332
333    /// Releases every lease this runner holds, without settling the jobs.
334    ///
335    /// What a graceful shutdown runs, so a restart re-claims its own work
336    /// immediately instead of waiting out a full lease.
337    pub async fn release_owned(runner_id: &str, database: &Database) -> Result<u64, sqlx::Error> {
338        let released = sqlx::query(
339            "UPDATE jobs SET status = 'ready', lease_owner = NULL, lease_until = NULL, \
340             updated_at = ? WHERE status = 'running' AND lease_owner = ?;",
341        )
342        .bind(now_secs())
343        .bind(runner_id)
344        .execute(&database.pool)
345        .await?
346        .rows_affected();
347
348        if released > 0 {
349            debug!(
350                event = "db_job_leases_released",
351                outcome = "success",
352                rows_released = released,
353            );
354        }
355        Ok(released)
356    }
357
358    /// One row by id.
359    pub async fn find_by_id(id: &str, database: &Database) -> Result<Option<Self>, sqlx::Error> {
360        // A `QueryBuilder` rather than `sqlx::query`, which takes only
361        // `&'static str` and so cannot be handed the shared `COLUMNS`. `id`
362        // still goes through `push_bind`, so nothing is interpolated.
363        let mut query = sqlx::QueryBuilder::new(format!("SELECT {COLUMNS} FROM jobs WHERE id = "));
364        query.push_bind(id);
365        let row = query.build().fetch_optional(&database.pool).await?;
366        row.map(Self::from_row).transpose()
367    }
368
369    /// The live job holding `(kind, dedup_key)`, if there is one.
370    pub async fn find_live(
371        kind: &str,
372        dedup_key: &str,
373        database: &Database,
374    ) -> Result<Option<Self>, sqlx::Error> {
375        let mut query = sqlx::QueryBuilder::new(format!(
376            "SELECT {COLUMNS} FROM jobs \
377             WHERE status IN ('ready', 'running') AND kind = "
378        ));
379        query.push_bind(kind);
380        query.push(" AND dedup_key = ");
381        query.push_bind(dedup_key);
382        let row = query.build().fetch_optional(&database.pool).await?;
383        row.map(Self::from_row).transpose()
384    }
385
386    /// How many live jobs of `kind` are queued or running.
387    ///
388    /// The kind-wide counterpart to [`Self::find_live`], for the callers that
389    /// want "is there work of this sort outstanding?" without naming a
390    /// `dedup_key` — a `notify_deliver` key is a per-occurrence uuid, so there
391    /// is no single key to ask about.
392    pub async fn count_live(kind: &str, database: &Database) -> Result<i64, sqlx::Error> {
393        let (count,): (i64,) = sqlx::query_as(
394            "SELECT COUNT(*) FROM jobs WHERE status IN ('ready', 'running') AND kind = ?;",
395        )
396        .bind(kind)
397        .fetch_one(&database.pool)
398        .await?;
399        Ok(count)
400    }
401
402    /// Deletes terminal rows settled before `cutoff`, returning how many went.
403    ///
404    /// Only terminal ones: a `ready` job scheduled far in the future is not old,
405    /// however long ago it was written, and a `running` one is somebody's work.
406    pub async fn cleanup(cutoff: i64, database: &Database) -> Result<u64, sqlx::Error> {
407        let deleted = sqlx::query(
408            "DELETE FROM jobs \
409             WHERE status IN ('done', 'failed', 'cancelled') AND updated_at < ?;",
410        )
411        .bind(cutoff)
412        .execute(&database.pool)
413        .await?
414        .rows_affected();
415        info!(
416            event = "db_job_cleanup_completed",
417            outcome = "success",
418            rows_removed = deleted,
419            cutoff = cutoff,
420        );
421        Ok(deleted)
422    }
423}
424
425#[cfg(test)]
426mod tests {
427    use super::*;
428    use serde_json::json;
429    use std::sync::Arc;
430
431    async fn db() -> Arc<Database> {
432        Arc::new(Database::connect_in_memory().await.unwrap())
433    }
434
435    /// Queues one job with sensible defaults, returning whether it was queued.
436    async fn enqueue(id: &str, key: &str, run_at: i64, database: &Database) -> bool {
437        Job::enqueue(
438            NewJob {
439                id,
440                kind: "test",
441                dedup_key: key,
442                payload: &json!({"n": 1}),
443                run_at,
444                deadline: None,
445                max_attempts: 3,
446            },
447            database,
448        )
449        .await
450        .unwrap()
451    }
452
453    async fn claim(runner: &str, database: &Database) -> Option<Job> {
454        Job::claim_next(runner, &["test"], now_secs() + 60, now_secs(), database)
455            .await
456            .unwrap()
457    }
458
459    #[tokio::test]
460    async fn a_queued_job_round_trips_through_every_column() {
461        let database = db().await;
462        let deadline = now_secs() + 900;
463        assert!(
464            Job::enqueue(
465                NewJob {
466                    id: "job-1",
467                    kind: "relay",
468                    dedup_key: "ord-1",
469                    payload: &json!({"order_id": "ord-1"}),
470                    run_at: 1_234,
471                    deadline: Some(deadline),
472                    max_attempts: 7,
473                },
474                &database,
475            )
476            .await
477            .unwrap()
478        );
479
480        let job = Job::find_by_id("job-1", &database).await.unwrap().unwrap();
481        assert_eq!(job.kind, "relay");
482        assert_eq!(job.dedup_key, "ord-1");
483        assert_eq!(job.payload, json!({"order_id": "ord-1"}));
484        assert_eq!(job.status, "ready");
485        assert_eq!(job.run_at, 1_234);
486        assert_eq!(job.attempts, 0);
487        assert_eq!(job.max_attempts, 7);
488        assert_eq!(job.deadline, Some(deadline));
489        assert!(job.lease_until.is_none());
490        assert!(job.lease_owner.is_none());
491        assert!(job.last_error.is_none());
492    }
493
494    /// The partial index's whole point: an identity is held by a *live* job and
495    /// released the moment one settles.
496    #[tokio::test]
497    async fn a_live_job_holds_its_identity_and_a_settled_one_releases_it() {
498        let database = db().await;
499        assert!(enqueue("job-1", "ord-1", now_secs(), &database).await);
500        assert!(
501            !enqueue("job-2", "ord-1", now_secs(), &database).await,
502            "a second live job must not take an identity already held"
503        );
504
505        // Claimed is still live.
506        let job = claim("runner-a", &database).await.unwrap();
507        assert!(
508            !enqueue("job-3", "ord-1", now_secs(), &database).await,
509            "'running' holds the identity exactly as 'ready' does"
510        );
511
512        Job::complete(&job.id, "runner-a", &database).await.unwrap();
513        assert!(
514            enqueue("job-4", "ord-1", now_secs(), &database).await,
515            "a settled job releases its identity"
516        );
517    }
518
519    #[tokio::test]
520    async fn claiming_takes_the_oldest_eligible_row_and_only_once() {
521        let database = db().await;
522        assert!(enqueue("job-new", "b", now_secs() - 10, &database).await);
523        assert!(enqueue("job-old", "a", now_secs() - 100, &database).await);
524
525        let first = claim("runner-a", &database).await.unwrap();
526        assert_eq!(first.id, "job-old", "oldest `run_at` first");
527        assert_eq!(first.attempts, 1, "the counter moves at claim");
528        assert_eq!(first.lease_owner.as_deref(), Some("runner-a"));
529
530        let second = claim("runner-b", &database).await.unwrap();
531        assert_eq!(second.id, "job-new");
532        assert!(
533            claim("runner-c", &database).await.is_none(),
534            "a claimed row is not claimable again"
535        );
536    }
537
538    #[tokio::test]
539    async fn claiming_skips_a_job_whose_run_at_has_not_arrived() {
540        let database = db().await;
541        assert!(enqueue("job-1", "a", now_secs() + 3_600, &database).await);
542        assert!(claim("runner-a", &database).await.is_none());
543    }
544
545    #[tokio::test]
546    async fn claiming_skips_a_kind_the_runner_does_not_hold() {
547        let database = db().await;
548        assert!(enqueue("job-1", "a", now_secs(), &database).await);
549
550        let claimed = Job::claim_next(
551            "runner-a",
552            &["something-else"],
553            now_secs() + 60,
554            now_secs(),
555            &database,
556        )
557        .await
558        .unwrap();
559        assert!(
560            claimed.is_none(),
561            "an unregistered kind is left alone, not mis-run"
562        );
563    }
564
565    #[tokio::test]
566    async fn claiming_with_no_registered_kinds_asks_the_database_nothing() {
567        let database = db().await;
568        assert!(enqueue("job-1", "a", now_secs(), &database).await);
569        assert!(
570            Job::claim_next("runner-a", &[], now_secs() + 60, now_secs(), &database)
571                .await
572                .unwrap()
573                .is_none()
574        );
575    }
576
577    /// The guard that makes a reclaimed lease safe: the runner that lost it
578    /// writes nothing, rather than overwriting the row somebody else now owns.
579    #[tokio::test]
580    async fn a_settlement_from_a_runner_that_lost_the_lease_writes_nothing() {
581        let database = db().await;
582        assert!(enqueue("job-1", "a", now_secs(), &database).await);
583        let job = claim("runner-a", &database).await.unwrap();
584
585        for settled in [
586            Job::complete(&job.id, "runner-b", &database).await.unwrap(),
587            Job::retry(&job.id, "runner-b", now_secs(), "x", &database)
588                .await
589                .unwrap(),
590            Job::reschedule(&job.id, "runner-b", now_secs(), &database)
591                .await
592                .unwrap(),
593            Job::abandon(&job.id, "runner-b", "x", &database)
594                .await
595                .unwrap(),
596        ] {
597            assert!(!settled, "the lease guard must refuse a foreign settlement");
598        }
599
600        let after = Job::find_by_id("job-1", &database).await.unwrap().unwrap();
601        assert_eq!(after.status, "running");
602        assert_eq!(after.lease_owner.as_deref(), Some("runner-a"));
603    }
604
605    #[tokio::test]
606    async fn retry_returns_the_job_at_its_new_time_and_keeps_the_attempt_count() {
607        let database = db().await;
608        assert!(enqueue("job-1", "a", now_secs(), &database).await);
609        let job = claim("runner-a", &database).await.unwrap();
610
611        assert!(
612            Job::retry(&job.id, "runner-a", 9_999, "upstream timed out", &database)
613                .await
614                .unwrap()
615        );
616
617        let after = Job::find_by_id("job-1", &database).await.unwrap().unwrap();
618        assert_eq!(after.status, "ready");
619        assert_eq!(after.run_at, 9_999);
620        assert_eq!(after.attempts, 1, "a retry does not forgive the attempt");
621        assert_eq!(after.last_error.as_deref(), Some("upstream timed out"));
622        assert!(after.lease_owner.is_none());
623    }
624
625    #[tokio::test]
626    async fn reschedule_resets_the_attempt_count_and_clears_the_last_error() {
627        let database = db().await;
628        assert!(enqueue("job-1", "a", now_secs(), &database).await);
629        let job = claim("runner-a", &database).await.unwrap();
630        Job::retry(&job.id, "runner-a", now_secs(), "a failure", &database)
631            .await
632            .unwrap();
633        let job = claim("runner-a", &database).await.unwrap();
634        assert_eq!(job.attempts, 2);
635
636        assert!(
637            Job::reschedule(&job.id, "runner-a", 5_000, &database)
638                .await
639                .unwrap()
640        );
641
642        let after = Job::find_by_id("job-1", &database).await.unwrap().unwrap();
643        assert_eq!(after.status, "ready");
644        assert_eq!(after.run_at, 5_000);
645        assert_eq!(after.attempts, 0, "a fresh occurrence starts fresh");
646        assert!(after.last_error.is_none());
647    }
648
649    #[tokio::test]
650    async fn abandon_is_terminal_and_records_why() {
651        let database = db().await;
652        assert!(enqueue("job-1", "a", now_secs(), &database).await);
653        let job = claim("runner-a", &database).await.unwrap();
654
655        assert!(
656            Job::abandon(&job.id, "runner-a", "the order no longer exists", &database)
657                .await
658                .unwrap()
659        );
660
661        let after = Job::find_by_id("job-1", &database).await.unwrap().unwrap();
662        assert_eq!(after.status, "failed");
663        assert_eq!(
664            after.last_error.as_deref(),
665            Some("the order no longer exists")
666        );
667    }
668
669    #[tokio::test]
670    async fn reclaim_takes_only_leases_that_have_expired() {
671        let database = db().await;
672        assert!(enqueue("job-live", "a", now_secs(), &database).await);
673        assert!(enqueue("job-dead", "b", now_secs(), &database).await);
674
675        // One long lease, one already expired.
676        Job::claim_next(
677            "runner-a",
678            &["test"],
679            now_secs() + 600,
680            now_secs(),
681            &database,
682        )
683        .await
684        .unwrap();
685        Job::claim_next("runner-a", &["test"], now_secs() - 1, now_secs(), &database)
686            .await
687            .unwrap();
688
689        assert_eq!(
690            Job::reclaim_expired(now_secs(), &database).await.unwrap(),
691            1
692        );
693
694        let reclaimed = Job::find_by_id("job-dead", &database)
695            .await
696            .unwrap()
697            .unwrap();
698        assert_eq!(reclaimed.status, "ready");
699        assert!(reclaimed.lease_owner.is_none());
700        assert_eq!(
701            reclaimed.attempts, 1,
702            "the attempt was spent, and the row must keep saying so"
703        );
704
705        let held = Job::find_by_id("job-live", &database)
706            .await
707            .unwrap()
708            .unwrap();
709        assert_eq!(held.status, "running");
710    }
711
712    #[tokio::test]
713    async fn release_owned_is_scoped_to_one_runner() {
714        let database = db().await;
715        assert!(enqueue("job-1", "a", now_secs(), &database).await);
716        assert!(enqueue("job-2", "b", now_secs(), &database).await);
717        claim("runner-a", &database).await.unwrap();
718        claim("runner-b", &database).await.unwrap();
719
720        assert_eq!(Job::release_owned("runner-a", &database).await.unwrap(), 1);
721        assert_eq!(
722            Job::find_by_id("job-1", &database)
723                .await
724                .unwrap()
725                .unwrap()
726                .status,
727            "ready"
728        );
729        assert_eq!(
730            Job::find_by_id("job-2", &database)
731                .await
732                .unwrap()
733                .unwrap()
734                .status,
735            "running"
736        );
737    }
738
739    #[tokio::test]
740    async fn find_live_sees_a_queued_job_and_not_a_settled_one() {
741        let database = db().await;
742        assert!(enqueue("job-1", "ord-1", now_secs(), &database).await);
743        assert!(
744            Job::find_live("test", "ord-1", &database)
745                .await
746                .unwrap()
747                .is_some()
748        );
749
750        let job = claim("runner-a", &database).await.unwrap();
751        Job::complete(&job.id, "runner-a", &database).await.unwrap();
752        assert!(
753            Job::find_live("test", "ord-1", &database)
754                .await
755                .unwrap()
756                .is_none()
757        );
758    }
759
760    #[tokio::test]
761    async fn cleanup_removes_settled_rows_strictly_older_than_the_cutoff() {
762        let database = db().await;
763        assert!(enqueue("job-done", "a", now_secs(), &database).await);
764        assert!(enqueue("job-live", "b", now_secs() + 3_600, &database).await);
765        let job = claim("runner-a", &database).await.unwrap();
766        Job::complete(&job.id, "runner-a", &database).await.unwrap();
767
768        // The cutoff is `updated_at`, which `complete` has just stamped to now.
769        assert_eq!(Job::cleanup(now_secs(), &database).await.unwrap(), 0);
770        assert_eq!(Job::cleanup(now_secs() + 1, &database).await.unwrap(), 1);
771        assert!(
772            Job::find_by_id("job-done", &database)
773                .await
774                .unwrap()
775                .is_none()
776        );
777        assert!(
778            Job::find_by_id("job-live", &database)
779                .await
780                .unwrap()
781                .is_some(),
782            "a job still queued is not old, however long ago it was written"
783        );
784    }
785
786    #[tokio::test]
787    async fn an_unknown_status_is_refused_by_the_check_constraint() {
788        let database = db().await;
789        let error = sqlx::query(
790            "INSERT INTO jobs (id, kind, dedup_key, payload, status, run_at, attempts, \
791             max_attempts, created_at, updated_at) \
792             VALUES ('x', 'test', 'k', '{}', 'halfway', 0, 0, 1, 0, 0);",
793        )
794        .execute(&database.pool)
795        .await
796        .unwrap_err();
797        assert!(error.to_string().contains("CHECK constraint failed"));
798    }
799}