Skip to main content

apalis_postgres/queries/
reenqueue_orphaned.rs

1use futures::TryFutureExt;
2use sqlx::{Executor, postgres::types::PgInterval};
3
4use crate::error::Error;
5
6/// Reenqueue jobs orphaned by a dead worker
7pub async fn reenqueue_orphaned<E>(conn: &mut E, queue: &str, dead_for: u64) -> Result<u64, Error>
8where
9    for<'e> &'e mut E: Executor<'e, Database = sqlx::Postgres>,
10{
11    let dead_for = PgInterval {
12        months: 0,
13        days: 0,
14        microseconds: dead_for as i64 * 1_000_000,
15    };
16
17    match sqlx::query_file!("queries/backend/reenqueue_orphaned.sql", dead_for, queue,)
18        .execute(conn)
19        .await
20    {
21        Ok(res) => {
22            if res.rows_affected() > 0 {
23                tracing::info!(
24                    "Re-enqueued {} orphaned tasks that were being processed by dead workers",
25                    res.rows_affected()
26                );
27            }
28            Ok(res.rows_affected())
29        }
30        Err(e) => {
31            tracing::error!("Failed to re-enqueue orphaned tasks: {e}");
32            Err(e.into())
33        }
34    }
35}
36
37/// Rescues tasks that could not be executed after a worker shutdown
38pub async fn reenqueue_abandoned<E>(
39    executor: &mut E,
40    queue: &str,
41    worker: &str,
42    task_ids: &[String],
43) -> Result<u64, Error>
44where
45    for<'e> &'e mut E: Executor<'e, Database = sqlx::Postgres>,
46{
47    let res = sqlx::query_file!(
48        "queries/worker/reenqueue_abandoned.sql",
49        queue,
50        worker,
51        task_ids
52    )
53    .execute(executor)
54    .map_ok(|res| res.rows_affected())
55    .await?;
56    Ok(res)
57}