apalis_postgres/queries/
reenqueue_orphaned.rs1use futures::TryFutureExt;
2use sqlx::{Executor, postgres::types::PgInterval};
3
4use crate::error::Error;
5
6pub 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
37pub 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}