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