Skip to main content

renox_core/queue/
mod.rs

1//! Background jobs stored in SQLite, with retries and a failed-jobs table.
2//!
3//! ```
4//! # use renox::prelude::*;
5//! # #[derive(Model, serde::Serialize, Default)] struct Order { id: i64 }
6//! use serde::{Deserialize, Serialize};
7//!
8//! #[derive(Serialize, Deserialize)]
9//! struct SendReceipt { order_id: i64 }
10//!
11//! impl Job for SendReceipt {
12//!     const NAME: &'static str = "send-receipt";
13//!     const MAX_ATTEMPTS: u32 = 5;
14//!
15//!     async fn handle(self, ctx: JobContext) -> Result {
16//!         let order = Order::find_or_404(&ctx.state.db, self.order_id).await?;
17//!         // … send it
18//! #       let _ = order;
19//!         Ok(())
20//!     }
21//! }
22//!
23//! # let _ =
24//! App::new().job::<SendReceipt>();            // register the handler
25//! # async fn demo(state: AppState, order_id: i64) -> Result {
26//! state.dispatch(SendReceipt { order_id }).await?;   // in a handler
27//! # Ok(()) }
28//! ```
29//!
30//! `serve` runs workers in the same process (`QUEUE_WORKERS`, default 2; 0 to
31//! turn off); `my-app queue:work` runs them on their own, and
32//! `queue:work --queue high,default` drains `high` before `default`.
33//!
34//! More than one job at a time:
35//!
36//! ```
37//! # use renox::prelude::*;
38//! # #[derive(serde::Serialize, serde::Deserialize)] struct Import { file: String }
39//! # impl Job for Import { const NAME: &'static str = "import"; async fn handle(self, _: JobContext) -> Result { Ok(()) } }
40//! # #[derive(serde::Serialize, serde::Deserialize)] struct Notify { user_id: i64 }
41//! # impl Job for Notify { const NAME: &'static str = "notify"; async fn handle(self, _: JobContext) -> Result { Ok(()) } }
42//! # async fn demo(state: AppState) -> Result {
43//! // One after another; a failure stops the rest.
44//! state.queue.chain()
45//!     .then(Import { file: "a.csv".into() })
46//!     .then(Notify { user_id: 7 })
47//!     .dispatch()
48//!     .await?;
49//!
50//! // Side by side, with progress and follow-up jobs.
51//! let batch = state.queue.batch("import-october")
52//!     .push(Import { file: "a.csv".into() })
53//!     .push(Import { file: "b.csv".into() })
54//!     .then(Notify { user_id: 7 })   // all succeeded
55//!     .catch(Notify { user_id: 1 })  // the first failure (which cancels the rest)
56//!     .dispatch()
57//!     .await?;
58//! let status = state.queue.batch_status(batch).await?; // total, pending, failed, progress()
59//! # let _ = status; Ok(()) }
60//! ```
61
62mod dashboard;
63mod worker;
64
65use std::collections::HashMap;
66use std::future::Future;
67use std::pin::Pin;
68use std::sync::Arc;
69use std::time::Duration;
70
71use serde::de::DeserializeOwned;
72use serde::{Deserialize, Serialize};
73use tokio::sync::Notify;
74
75pub use dashboard::{Dashboard, GATE as DASHBOARD_GATE, QueueCounts, QueueStats};
76pub use worker::Worker;
77
78use crate::db::{Db, Migration, Transaction};
79use crate::{AppState, Error, Result};
80
81pub(crate) const MIGRATIONS: &[Migration] = &[
82    crate::db::framework_migration!("queue", "00010101000100_create_jobs_table"),
83    crate::db::framework_migration!("queue", "00010101000110_add_chains_and_batches_to_jobs"),
84    crate::db::framework_migration!("queue", "00010101000120_add_callback_of_to_jobs"),
85];
86
87/// Payloads of `ENCRYPTED` jobs start with this; JSON never does.
88const SEALED: &str = "enc:";
89/// Cache rows claiming a unique job, holding its id.
90const UNIQUE_PREFIX: &str = "renox:unique:";
91
92/// A unit of background work. It is stored as JSON, so keep it to ids and
93/// small values rather than whole models.
94pub trait Job: Serialize + DeserializeOwned + Send + Sync + 'static {
95    /// A stable name; stored with each job, so don't rename it while jobs are queued.
96    const NAME: &'static str;
97    /// Its queue; `queue:work --queue high,default` drains queues in order.
98    const QUEUE: &'static str = "default";
99    /// Runs before the job is moved to `failed_jobs`.
100    const MAX_ATTEMPTS: u32 = 3;
101    /// How long one attempt may run.
102    const TIMEOUT: Duration = Duration::from_secs(60);
103    /// Store the payload encrypted with `APP_KEY` (it holds personal data or
104    /// secrets). The worker decrypts it; `failed_jobs` keeps it encrypted.
105    const ENCRYPTED: bool = false;
106    /// While a job with the same [`Job::unique_id`] is queued or running,
107    /// dispatching another returns the queued one's id instead, for up to
108    /// this long (the claim is dropped when the job finishes).
109    const UNIQUE_FOR: Option<Duration> = None;
110
111    /// Wait before retrying after the given failed attempt (1-based).
112    fn backoff(attempt: u32) -> Duration {
113        Duration::from_secs(10 * u64::from(attempt))
114    }
115
116    /// What makes two jobs "the same" for [`Job::UNIQUE_FOR`], e.g. the
117    /// order id. Empty (the default) makes every job of the type the same.
118    fn unique_id(&self) -> String {
119        String::new()
120    }
121
122    /// Conditions checked before each attempt; a job that can't run yet is
123    /// put back without using up an attempt.
124    fn middleware(&self) -> Vec<Middleware> {
125        Vec::new()
126    }
127
128    /// Does the work; an error fails the attempt (retried unless `Error::permanent`).
129    fn handle(self, ctx: JobContext) -> impl Future<Output = Result> + Send;
130
131    /// Runs once the job has failed for good (attempts used up, or a
132    /// permanent error), with the last error, e.g. to tell the user. `ctx`
133    /// is the last attempt's: the state, the job's id, the attempt, its batch.
134    fn failed(self, ctx: JobContext, error: String) -> impl Future<Output = ()> + Send {
135        let _ = (ctx, error);
136        async {}
137    }
138}
139
140/// A condition checked before a job's attempt; see [`Job::middleware`].
141///
142/// ```
143/// # use renox::prelude::*;
144/// use renox::queue::Middleware;
145/// # use std::time::Duration;
146/// # #[derive(serde::Serialize, serde::Deserialize)] struct SyncStock { shop_id: i64 }
147/// impl Job for SyncStock {
148///     const NAME: &'static str = "sync-stock";
149///
150///     fn middleware(&self) -> Vec<Middleware> {
151///         vec![
152///             // One sync per shop at a time; others wait 10 s and try again.
153///             Middleware::without_overlapping(format!("shop:{}", self.shop_id))
154///                 .release_after(Duration::from_secs(10)),
155///             // The supplier's API allows 60 calls a minute.
156///             Middleware::rate_limited("supplier-api", 60, Duration::from_secs(60)),
157///         ]
158///     }
159///
160///     async fn handle(self, _ctx: JobContext) -> Result { Ok(()) }
161/// }
162/// ```
163///
164/// Both use the cache (`CACHE_STORE`): with `memory`, they hold within one
165/// process; use `database` when several processes run workers.
166#[derive(Debug, Clone)]
167pub struct Middleware(MiddlewareKind);
168
169#[derive(Debug, Clone)]
170enum MiddlewareKind {
171    WithoutOverlapping {
172        key: String,
173        release_after: Duration,
174    },
175    RateLimited {
176        key: String,
177        max: u32,
178        per: Duration,
179    },
180}
181
182impl Middleware {
183    /// One job with this `key` at a time; others are put back for 5 seconds
184    /// (see [`Middleware::release_after`]).
185    pub fn without_overlapping(key: impl Into<String>) -> Self {
186        Self(MiddlewareKind::WithoutOverlapping {
187            key: key.into(),
188            release_after: Duration::from_secs(5),
189        })
190    }
191
192    /// At most `max` attempts with this `key` per `per` (a fixed window);
193    /// the rest are put back until the window ends.
194    pub fn rate_limited(key: impl Into<String>, max: u32, per: Duration) -> Self {
195        Self(MiddlewareKind::RateLimited {
196            key: key.into(),
197            max,
198            per: per.max(Duration::from_secs(1)),
199        })
200    }
201
202    /// How long a job waits before trying again when it would overlap.
203    pub fn release_after(mut self, wait: Duration) -> Self {
204        if let MiddlewareKind::WithoutOverlapping { release_after, .. } = &mut self.0 {
205            *release_after = wait;
206        }
207        self
208    }
209}
210
211/// What a job gets when it runs.
212#[non_exhaustive]
213pub struct JobContext {
214    /// The app's state (database, mailer, queue…).
215    pub state: AppState,
216    /// 1 on the first try.
217    pub attempt: u32,
218    /// The job's id (0 for [`AppState::dispatch_sync`]).
219    pub id: i64,
220    /// The batch it belongs to, or, for a batch's `then`/`catch`/`finally`
221    /// job, the batch it follows (read it with `queue.batch_status`).
222    pub batch_id: Option<i64>,
223}
224
225type RunFn =
226    Arc<dyn Fn(String, JobContext) -> Pin<Box<dyn Future<Output = Result> + Send>> + Send + Sync>;
227type FailedFn = Arc<
228    dyn Fn(String, JobContext, String) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + Sync,
229>;
230
231/// A registered job type: how to run it and when to retry it.
232#[derive(Clone)]
233pub(crate) struct JobHandler {
234    pub(crate) run: RunFn,
235    pub(crate) failed: FailedFn,
236    pub(crate) backoff: fn(u32) -> Duration,
237    pub(crate) timeout: Duration,
238    /// The unique claim's cache key, from the (decrypted) payload.
239    pub(crate) unique_key: fn(&str) -> Option<String>,
240    pub(crate) middleware: fn(&str) -> Vec<Middleware>,
241}
242
243/// The handler for `J`; `Registry::job` registers it.
244pub(crate) fn handler<J: Job>() -> JobHandler {
245    JobHandler {
246        run: Arc::new(|payload, ctx| {
247            Box::pin(async move {
248                // A payload that doesn't decode now never will.
249                let job: J = serde_json::from_str(&payload).map_err(crate::Error::permanent)?;
250                job.handle(ctx).await
251            })
252        }),
253        failed: Arc::new(|payload, ctx, error| {
254            Box::pin(async move {
255                if let Ok(job) = serde_json::from_str::<J>(&payload) {
256                    job.failed(ctx, error).await;
257                }
258            })
259        }),
260        backoff: J::backoff,
261        timeout: J::TIMEOUT,
262        unique_key: |payload| {
263            J::UNIQUE_FOR?;
264            let job: J = serde_json::from_str(payload).ok()?;
265            Some(unique_key::<J>(&job))
266        },
267        middleware: |payload| {
268            serde_json::from_str::<J>(payload)
269                .map(|job| job.middleware())
270                .unwrap_or_default()
271        },
272    }
273}
274
275fn unique_key<J: Job>(job: &J) -> String {
276    format!("{UNIQUE_PREFIX}{}:{}", J::NAME, job.unique_id())
277}
278
279pub(crate) type Handlers = Arc<HashMap<&'static str, JobHandler>>;
280
281/// A job ready to be stored: for chains, batches and their follow-ups.
282#[derive(Debug, Clone, Serialize, Deserialize)]
283pub(crate) struct Encoded {
284    queue: String,
285    job: String,
286    payload: String,
287    max_attempts: u32,
288}
289
290/// Dispatches jobs and wakes the workers of this process.
291#[derive(Clone)]
292pub struct Queue {
293    db: Db,
294    key: cookie::Key,
295    wake: Arc<Notify>,
296}
297
298pub(crate) fn unix_now() -> i64 {
299    crate::clock::unix_secs()
300}
301
302impl Queue {
303    pub(crate) fn new(db: Db, key: cookie::Key) -> Self {
304        Self {
305            db,
306            key,
307            wake: Arc::new(Notify::new()),
308        }
309    }
310
311    pub(crate) fn wake(&self) -> &Notify {
312        &self.wake
313    }
314
315    pub(crate) fn wake_workers(&self) {
316        self.wake.notify_waiters();
317    }
318
319    fn encode<J: Job>(&self, job: &J, queue: &str) -> Result<Encoded> {
320        let json = serde_json::to_string(job)?;
321        let payload = if J::ENCRYPTED {
322            format!("{SEALED}{}", crate::crypto::seal(&self.key, &json))
323        } else {
324            json
325        };
326        Ok(Encoded {
327            queue: queue.to_owned(),
328            job: J::NAME.to_owned(),
329            payload,
330            max_attempts: J::MAX_ATTEMPTS.max(1),
331        })
332    }
333
334    /// The JSON of a stored payload, decrypting an encrypted one.
335    pub(crate) fn open(&self, payload: &str) -> Result<String> {
336        match payload.strip_prefix(SEALED) {
337            Some(sealed) => Ok(crate::crypto::open(&self.key, sealed).map_err(Error::permanent)?),
338            None => Ok(payload.to_owned()),
339        }
340    }
341
342    /// Queues a job to run as soon as a worker is free; returns its id.
343    pub async fn dispatch<J: Job>(&self, job: J) -> Result<i64> {
344        self.dispatch_after(job, Duration::ZERO).await
345    }
346
347    /// Queues a job to run after `delay`; returns its id.
348    pub async fn dispatch_after<J: Job>(&self, job: J, delay: Duration) -> Result<i64> {
349        self.push(job, J::QUEUE, delay).await
350    }
351
352    /// Queues a job on `queue` instead of its [`Job::QUEUE`].
353    pub async fn dispatch_on<J: Job>(&self, queue: &str, job: J) -> Result<i64> {
354        self.push(job, queue, Duration::ZERO).await
355    }
356
357    async fn push<J: Job>(&self, job: J, queue: &str, delay: Duration) -> Result<i64> {
358        let encoded = self.encode(&job, queue)?;
359        let id = match J::UNIQUE_FOR {
360            None => insert(&self.db, &encoded, delay, None, None, None).await?,
361            Some(ttl) => {
362                let mut tx = self.db.begin().await?;
363                let id =
364                    insert_unique(&mut tx, &unique_key::<J>(&job), ttl, &encoded, delay).await?;
365                tx.commit().await?;
366                id
367            }
368        };
369        self.wake.notify_waiters();
370        Ok(id)
371    }
372
373    /// Queues a job inside `tx`, so it only exists if the transaction
374    /// commits; workers pick it up within a second of the commit. Use it
375    /// instead of `dispatch` while a transaction is open: on SQLite, which
376    /// writes one transaction at a time, `dispatch` would wait for `tx`.
377    ///
378    /// ```
379    /// # use renox::prelude::*;
380    /// # #[derive(serde::Serialize, serde::Deserialize)] struct SendReceipt { order_id: i64 }
381    /// # impl Job for SendReceipt { const NAME: &'static str = "r"; async fn handle(self, _: JobContext) -> Result { Ok(()) } }
382    /// # async fn demo(state: AppState) -> Result {
383    /// let mut tx = state.db.begin().await?;
384    /// let order_id: i64 = renox::db::sql("INSERT INTO orders (total) VALUES (?) RETURNING id")
385    ///     .bind(75_000)
386    ///     .scalar(&mut tx)
387    ///     .await?;
388    /// state.queue.dispatch_in(&mut tx, SendReceipt { order_id }).await?;
389    /// tx.commit().await?; // no order, no receipt
390    /// # Ok(()) }
391    /// ```
392    pub async fn dispatch_in<J: Job>(&self, tx: &mut Transaction, job: J) -> Result<i64> {
393        let encoded = self.encode(&job, J::QUEUE)?;
394        match J::UNIQUE_FOR {
395            None => insert(tx, &encoded, Duration::ZERO, None, None, None).await,
396            Some(ttl) => {
397                insert_unique(tx, &unique_key::<J>(&job), ttl, &encoded, Duration::ZERO).await
398            }
399        }
400    }
401
402    /// Jobs that run one after another: the next is queued when the one
403    /// before succeeds, and a job that fails for good stops the chain
404    /// (`queue:retry` resumes it).
405    pub fn chain(&self) -> Chain {
406        Chain {
407            queue: self.clone(),
408            jobs: Vec::new(),
409            error: None,
410        }
411    }
412
413    /// Jobs that run side by side, tracked together; see [`Batch`].
414    pub fn batch(&self, name: &str) -> Batch {
415        Batch {
416            queue: self.clone(),
417            name: name.to_owned(),
418            jobs: Vec::new(),
419            then: None,
420            catch: None,
421            finally: None,
422            allow_failures: false,
423            error: None,
424        }
425    }
426
427    /// How a batch is doing; `None` if there's no such batch.
428    pub async fn batch_status(&self, id: i64) -> Result<Option<BatchStatus>> {
429        let row = crate::db::sql(
430            "SELECT id, name, total, pending, failed, cancelled_at, finished_at, created_at \
431             FROM job_batches WHERE id = ?",
432        )
433        .bind(id)
434        .fetch_optional(&self.db)
435        .await?;
436        let Some(row) = row else {
437            return Ok(None);
438        };
439        Ok(Some(BatchStatus {
440            id: row.try_get("id")?,
441            name: row.try_get("name")?,
442            total: row.try_get("total")?,
443            pending: row.try_get("pending")?,
444            failed: row.try_get("failed")?,
445            cancelled: row.try_get::<Option<i64>>("cancelled_at")?.is_some(),
446            finished: row.try_get::<Option<i64>>("finished_at")?.is_some(),
447            created_at: crate::db::from_unix(row.try_get("created_at")?),
448        }))
449    }
450
451    /// Cancels a batch: its jobs that haven't run are skipped.
452    pub async fn cancel_batch(&self, id: i64) -> Result<bool> {
453        Ok(crate::db::sql(
454            "UPDATE job_batches SET cancelled_at = ? WHERE id = ? AND cancelled_at IS NULL",
455        )
456        .bind(unix_now())
457        .bind(id)
458        .execute(&self.db)
459        .await?
460            == 1)
461    }
462
463    /// Jobs waiting or running.
464    pub async fn pending(&self) -> Result<i64> {
465        Ok(crate::db::sql("SELECT COUNT(*) FROM jobs")
466            .scalar(&self.db)
467            .await?)
468    }
469
470    /// Jobs that failed for good, oldest first.
471    pub async fn failed(&self) -> Result<Vec<FailedJob>> {
472        let rows = crate::db::sql(
473            "SELECT id, queue, job, payload, error, failed_at FROM failed_jobs ORDER BY id",
474        )
475        .fetch_all(&self.db)
476        .await?;
477        rows.iter()
478            .map(|r| {
479                Ok(FailedJob {
480                    id: r.try_get("id")?,
481                    queue: r.try_get("queue")?,
482                    job: r.try_get("job")?,
483                    payload: r.try_get("payload")?,
484                    error: r.try_get("error")?,
485                    failed_at: crate::db::from_unix(r.try_get("failed_at")?),
486                })
487            })
488            .collect()
489    }
490
491    /// Puts the failed job `id` back on its queue with fresh attempts (and
492    /// its chain and batch); returns whether there was such a job.
493    pub async fn retry(&self, id: i64) -> Result<bool> {
494        Ok(self.move_failed(Some(id)).await? == 1)
495    }
496
497    /// Puts every failed job back on its queue; returns how many.
498    pub async fn retry_all(&self) -> Result<u64> {
499        self.move_failed(None).await
500    }
501
502    async fn move_failed(&self, id: Option<i64>) -> Result<u64> {
503        let mut tx = self.db.begin().await?;
504        let filter = if id.is_some() { " WHERE id = ?" } else { "" };
505        let mut batches = crate::db::sql(format!(
506            "UPDATE job_batches SET failed = failed - 1, pending = pending + 1, finished_at = NULL \
507             WHERE id IN (SELECT batch_id FROM failed_jobs{filter})"
508        ));
509        if let Some(id) = id {
510            batches = batches.bind(id);
511        }
512        batches.execute(&mut tx).await?;
513        let insert = format!(
514            "INSERT INTO jobs (queue, job, payload, max_attempts, available_at, created_at, chain, \
515             batch_id, callback_of) SELECT queue, job, payload, max_attempts, ?, ?, chain, batch_id, \
516             callback_of FROM failed_jobs{filter}"
517        );
518        let mut query = crate::db::sql(insert).bind(unix_now()).bind(unix_now());
519        if let Some(id) = id {
520            query = query.bind(id);
521        }
522        let moved = query.execute(&mut tx).await?;
523        let mut delete = crate::db::sql(format!("DELETE FROM failed_jobs{filter}"));
524        if let Some(id) = id {
525            delete = delete.bind(id);
526        }
527        delete.execute(&mut tx).await?;
528        tx.commit().await?;
529        self.wake.notify_waiters();
530        Ok(moved)
531    }
532
533    /// Deletes one failed job; returns whether it existed.
534    pub async fn forget_failed(&self, id: i64) -> Result<bool> {
535        Ok(crate::db::sql("DELETE FROM failed_jobs WHERE id = ?")
536            .bind(id)
537            .execute(&self.db)
538            .await?
539            == 1)
540    }
541
542    /// Deletes failed jobs older than `age`; returns how many.
543    pub async fn prune_failed(&self, age: Duration) -> Result<u64> {
544        Ok(
545            crate::db::sql("DELETE FROM failed_jobs WHERE failed_at < ?")
546                .bind(unix_now() - age.as_secs() as i64)
547                .execute(&self.db)
548                .await?,
549        )
550    }
551
552    /// Deletes batches that finished (or were cancelled) more than `age` ago.
553    pub async fn prune_batches(&self, age: Duration) -> Result<u64> {
554        let before = unix_now() - age.as_secs() as i64;
555        Ok(crate::db::sql(
556            "DELETE FROM job_batches WHERE (finished_at IS NOT NULL AND finished_at < ?) \
557             OR (cancelled_at IS NOT NULL AND cancelled_at < ? AND pending = 0)",
558        )
559        .bind(before)
560        .bind(before)
561        .execute(&self.db)
562        .await?)
563    }
564
565    /// Deletes failed jobs.
566    pub async fn flush_failed(&self) -> Result<u64> {
567        Ok(crate::db::sql("DELETE FROM failed_jobs")
568            .execute(&self.db)
569            .await?)
570    }
571}
572
573/// Stores one job; returns its id.
574async fn insert<'c>(
575    db: impl crate::db::Executor<'c>,
576    job: &Encoded,
577    delay: Duration,
578    chain: Option<String>,
579    batch_id: Option<i64>,
580    callback_of: Option<i64>,
581) -> Result<i64> {
582    let now = unix_now();
583    let id: i64 = crate::db::sql(
584        "INSERT INTO jobs (queue, job, payload, max_attempts, available_at, created_at, chain, \
585         batch_id, callback_of) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) RETURNING id",
586    )
587    .bind(&job.queue)
588    .bind(&job.job)
589    .bind(&job.payload)
590    .bind(i64::from(job.max_attempts))
591    .bind(now + delay.as_secs() as i64)
592    .bind(now)
593    .bind(chain)
594    .bind(batch_id)
595    .bind(callback_of)
596    .scalar(db)
597    .await?;
598    Ok(id)
599}
600
601/// Stores a unique job unless its claim is held; returns its id or the
602/// holder's. The claim is a `cache` row (whatever `CACHE_STORE` is), so it
603/// holds across processes.
604async fn insert_unique(
605    tx: &mut Transaction,
606    key: &str,
607    ttl: Duration,
608    job: &Encoded,
609    delay: Duration,
610) -> Result<i64> {
611    let now = unix_now();
612    let claimed = crate::db::sql(
613        "INSERT INTO cache (key, value, expires_at) VALUES (?, '0', ?) \
614         ON CONFLICT (key) DO UPDATE SET value = '0', expires_at = excluded.expires_at \
615         WHERE cache.expires_at IS NOT NULL AND cache.expires_at <= ?",
616    )
617    .bind(key)
618    .bind(now + ttl.as_secs().max(1) as i64)
619    .bind(now)
620    .execute(&mut *tx)
621    .await?;
622    if claimed == 0 {
623        let holder: String = crate::db::sql("SELECT value FROM cache WHERE key = ?")
624            .bind(key)
625            .scalar(&mut *tx)
626            .await?;
627        return Ok(holder.parse().unwrap_or_default());
628    }
629    let id = insert(&mut *tx, job, delay, None, None, None).await?;
630    crate::db::sql("UPDATE cache SET value = ? WHERE key = ?")
631        .bind(id.to_string())
632        .bind(key)
633        .execute(&mut *tx)
634        .await?;
635    Ok(id)
636}
637
638/// Jobs that run one after another; from [`Queue::chain`].
639#[must_use = "a chain does nothing until dispatched"]
640pub struct Chain {
641    queue: Queue,
642    jobs: Vec<Encoded>,
643    error: Option<Error>,
644}
645
646impl Chain {
647    /// Adds a job to the end of the chain.
648    pub fn then<J: Job>(mut self, job: J) -> Self {
649        match self.queue.encode(&job, J::QUEUE) {
650            Ok(encoded) => self.jobs.push(encoded),
651            Err(err) => {
652                self.error.get_or_insert(err);
653            }
654        }
655        self
656    }
657
658    /// Queues the first job; returns its id (0 for an empty chain).
659    pub async fn dispatch(self) -> Result<i64> {
660        if let Some(err) = self.error {
661            return Err(err);
662        }
663        let mut jobs = self.jobs.into_iter();
664        let Some(first) = jobs.next() else {
665            return Ok(0);
666        };
667        let rest: Vec<Encoded> = jobs.collect();
668        let chain = (!rest.is_empty())
669            .then(|| serde_json::to_string(&rest))
670            .transpose()?;
671        let id = insert(&self.queue.db, &first, Duration::ZERO, chain, None, None).await?;
672        self.queue.wake_workers();
673        Ok(id)
674    }
675}
676
677/// Jobs that run side by side and are tracked together; from
678/// [`Queue::batch`]. The first job that fails for good cancels the batch
679/// (jobs that haven't run are skipped) unless [`Batch::allow_failures`].
680#[must_use = "a batch does nothing until dispatched"]
681pub struct Batch {
682    queue: Queue,
683    name: String,
684    jobs: Vec<Encoded>,
685    then: Option<Encoded>,
686    catch: Option<Encoded>,
687    finally: Option<Encoded>,
688    allow_failures: bool,
689    error: Option<Error>,
690}
691
692impl Batch {
693    fn encode<J: Job>(&mut self, job: J) -> Option<Encoded> {
694        match self.queue.encode(&job, J::QUEUE) {
695            Ok(encoded) => Some(encoded),
696            Err(err) => {
697                self.error.get_or_insert(err);
698                None
699            }
700        }
701    }
702
703    /// Adds a job to the batch.
704    pub fn push<J: Job>(mut self, job: J) -> Self {
705        if let Some(job) = self.encode(job) {
706            self.jobs.push(job);
707        }
708        self
709    }
710
711    /// Queued when every job has succeeded.
712    pub fn then<J: Job>(mut self, job: J) -> Self {
713        self.then = self.encode(job);
714        self
715    }
716
717    /// Queued when the first job fails for good.
718    pub fn catch<J: Job>(mut self, job: J) -> Self {
719        self.catch = self.encode(job);
720        self
721    }
722
723    /// Queued when every job has run or been skipped, however it went.
724    pub fn finally<J: Job>(mut self, job: J) -> Self {
725        self.finally = self.encode(job);
726        self
727    }
728
729    /// A job failing for good doesn't cancel the others.
730    pub fn allow_failures(mut self) -> Self {
731        self.allow_failures = true;
732        self
733    }
734
735    /// Stores the batch and its jobs; returns the batch id.
736    pub async fn dispatch(self) -> Result<i64> {
737        if let Some(err) = self.error {
738            return Err(err);
739        }
740        let json = |job: &Option<Encoded>| job.as_ref().map(serde_json::to_string).transpose();
741        let (then, catch, finally) = (json(&self.then)?, json(&self.catch)?, json(&self.finally)?);
742        let total = self.jobs.len() as i64;
743        let mut tx = self.queue.db.begin().await?;
744        let id: i64 = crate::db::sql(
745            "INSERT INTO job_batches (name, total, pending, allow_failures, then_job, catch_job, \
746             finally_job, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?) RETURNING id",
747        )
748        .bind(&self.name)
749        .bind(total)
750        .bind(total)
751        .bind(self.allow_failures)
752        .bind(then)
753        .bind(catch)
754        .bind(finally)
755        .bind(unix_now())
756        .scalar(&mut tx)
757        .await?;
758        for job in &self.jobs {
759            insert(&mut tx, job, Duration::ZERO, None, Some(id), None).await?;
760        }
761        if self.jobs.is_empty() {
762            finish_batch(&mut tx, id, 0).await?;
763        }
764        tx.commit().await?;
765        self.queue.wake_workers();
766        Ok(id)
767    }
768}
769
770/// Counts one of a batch's jobs as done (`failed`: for good) and queues the
771/// follow-up jobs that are due.
772pub(crate) async fn batch_job_done(tx: &mut Transaction, batch_id: i64, failed: bool) -> Result {
773    let row: Option<(i64, i64, Option<String>)> = crate::db::sql(
774        "UPDATE job_batches SET pending = pending - 1, failed = failed + ?, \
775         cancelled_at = CASE WHEN ? AND NOT allow_failures THEN COALESCE(cancelled_at, ?) \
776         ELSE cancelled_at END \
777         WHERE id = ? RETURNING pending, failed, catch_job",
778    )
779    .bind(i64::from(failed))
780    .bind(failed)
781    .bind(unix_now())
782    .bind(batch_id)
783    .fetch_as(&mut *tx)
784    .await?
785    .into_iter()
786    .next();
787    let Some((pending, failures, catch)) = row else {
788        return Ok(()); // pruned meanwhile
789    };
790    if failed && failures == 1 {
791        queue_follow_up(tx, catch, batch_id).await?;
792    }
793    if pending <= 0 {
794        finish_batch(tx, batch_id, failures).await?;
795    }
796    Ok(())
797}
798
799async fn finish_batch(tx: &mut Transaction, batch_id: i64, failures: i64) -> Result {
800    let (then, finally, cancelled): (Option<String>, Option<String>, Option<i64>) = crate::db::sql(
801        "UPDATE job_batches SET finished_at = ? WHERE id = ? \
802         RETURNING then_job, finally_job, cancelled_at",
803    )
804    .bind(unix_now())
805    .bind(batch_id)
806    .fetch_as(&mut *tx)
807    .await?
808    .into_iter()
809    .next()
810    .unwrap_or_default();
811    if failures == 0 && cancelled.is_none() {
812        queue_follow_up(tx, then, batch_id).await?;
813    }
814    queue_follow_up(tx, finally, batch_id).await
815}
816
817/// Queues a batch's `then`/`catch`/`finally` job, which sees the batch in
818/// `JobContext::batch_id` but isn't counted in it.
819async fn queue_follow_up(tx: &mut Transaction, job: Option<String>, batch_id: i64) -> Result {
820    if let Some(job) = job {
821        let job: Encoded = serde_json::from_str(&job)?;
822        insert(&mut *tx, &job, Duration::ZERO, None, None, Some(batch_id)).await?;
823    }
824    Ok(())
825}
826
827/// Queues the next job of a chain, handing it the rest.
828pub(crate) async fn continue_chain(tx: &mut Transaction, chain: &str) -> Result {
829    let rest: Vec<Encoded> = serde_json::from_str(chain)?;
830    let mut rest = rest.into_iter();
831    if let Some(next) = rest.next() {
832        let rest: Vec<Encoded> = rest.collect();
833        let chain = (!rest.is_empty())
834            .then(|| serde_json::to_string(&rest))
835            .transpose()?;
836        insert(&mut *tx, &next, Duration::ZERO, chain, None, None).await?;
837    }
838    Ok(())
839}
840
841/// Drops a unique job's claim.
842pub(crate) async fn release_unique(tx: &mut Transaction, key: &str) -> Result {
843    crate::db::sql("DELETE FROM cache WHERE key = ?")
844        .bind(key)
845        .execute(&mut *tx)
846        .await?;
847    Ok(())
848}
849
850/// How a batch is doing; from [`Queue::batch_status`].
851#[derive(Debug, Clone, Serialize)]
852#[non_exhaustive]
853pub struct BatchStatus {
854    /// The `job_batches` row id.
855    pub id: i64,
856    /// The name given to [`Queue::batch`].
857    pub name: String,
858    /// Jobs in the batch.
859    pub total: i64,
860    /// Jobs not run yet (or being retried).
861    pub pending: i64,
862    /// Jobs that failed for good.
863    pub failed: i64,
864    /// Cancelled by a failure (without `allow_failures`) or `cancel_batch`.
865    pub cancelled: bool,
866    /// Every job has run or been skipped.
867    pub finished: bool,
868    /// When it was made.
869    pub created_at: crate::db::DateTime,
870}
871
872impl BatchStatus {
873    /// Percent of the jobs that have run, 0–100.
874    pub fn progress(&self) -> u8 {
875        if self.total == 0 {
876            return 100;
877        }
878        ((self.total - self.pending) * 100 / self.total).clamp(0, 100) as u8
879    }
880}
881
882/// A job that used up its attempts.
883#[derive(Debug, Clone, Serialize)]
884#[non_exhaustive]
885pub struct FailedJob {
886    /// The `failed_jobs` row id, for `queue:retry` and `queue:forget`.
887    pub id: i64,
888    /// The queue it ran on.
889    pub queue: String,
890    /// The job type's [`Job::NAME`].
891    pub job: String,
892    /// The serialized job (JSON, or `enc:…` when encrypted).
893    pub payload: String,
894    /// The last attempt's error.
895    pub error: String,
896    /// When it failed for good.
897    pub failed_at: crate::db::DateTime,
898}
899
900impl AppState {
901    /// Shorthand for `state.queue.dispatch(job)`.
902    pub async fn dispatch<J: Job>(&self, job: J) -> Result<i64> {
903        self.queue.dispatch(job).await
904    }
905
906    /// Runs `job` now, in this task, instead of queueing it (no retries,
907    /// middleware or `failed` hook): the error, if any, is returned.
908    pub async fn dispatch_sync<J: Job>(&self, job: J) -> Result {
909        let ctx = JobContext {
910            state: self.clone(),
911            attempt: 1,
912            id: 0,
913            batch_id: None,
914        };
915        crate::context::scope_app(self.clone(), job.handle(ctx)).await
916    }
917}