Skip to main content

apalis_sqlite/queries/
reenqueue_orphaned.rs

1use apalis_core::{backend::WorkerFilter, task::task_id::coalesce_ids};
2use sqlx::Executor;
3
4/// Re-enqueue tasks that were being processed by dead workers
5///
6/// A worker that has not sent a keep-alive signal within the heartbeat duration is considered dead
7pub async fn reenqueue_orphaned<'a, E>(
8    executor: E,
9    dead_for: i64,
10    queue: &str,
11    filter: &WorkerFilter,
12) -> Result<u64, sqlx::Error>
13where
14    E: Executor<'a, Database = sqlx::Sqlite>,
15{
16    let (exclude_id, only_id) = match filter {
17        WorkerFilter::AllExcept(id) => (Some(id), None),
18        WorkerFilter::Only(id) => (None, Some(id)),
19        WorkerFilter::None => (None, None),
20        _ => unreachable!(),
21    };
22    match sqlx::query_file!(
23        "queries/backend/reenqueue_orphaned.sql",
24        dead_for,
25        queue,
26        exclude_id,
27        only_id,
28    )
29    .execute(executor)
30    .await
31    {
32        Ok(res) => Ok(res.rows_affected()),
33        Err(e) => Err(e),
34    }
35}
36
37/// Re-enqueue tasks that were being processed by dying a worker
38///
39/// This will be invoked during `Backend::poll_close`
40pub async fn reenqueue_abandoned<'a, E>(
41    executor: E,
42    queue: &str,
43    worker: &str,
44    task_ids: &Vec<String>,
45) -> Result<u64, sqlx::Error>
46where
47    E: Executor<'a, Database = sqlx::Sqlite>,
48{
49    let task_ids = coalesce_ids(task_ids);
50    match sqlx::query_file!(
51        "queries/backend/reenqueue_abandoned.sql",
52        queue,
53        worker,
54        task_ids
55    )
56    .execute(executor)
57    .await
58    {
59        Ok(res) => Ok(res.rows_affected()),
60        Err(e) => Err(e),
61    }
62}