apalis_postgres/queries/
list_tasks.rs1use apalis_core::{
2 backend::{Backend, Filter, ListAllTasks, ListTasks},
3 task::{Task, status::Status},
4};
5
6use crate::from_row::PgTaskRow;
7use crate::{PgTask, PostgresStorage, error::Error};
8
9impl<Args> ListTasks for PostgresStorage<Args>
10where
11 PostgresStorage<Args>: Backend<Error = Error>,
12 Args: 'static,
13{
14 fn list_tasks(
15 &self,
16 filter: &Filter,
17 ) -> impl Future<Output = Result<Vec<PgTask>, Self::Error>> + Send {
18 let queue = self.persistence.config.queue.to_string();
19 let pool = self.persistence.pool.clone();
20 let limit = filter.limit() as i64;
21 let offset = filter.offset() as i64;
22 let status = filter
23 .status
24 .as_ref()
25 .unwrap_or(&Status::Pending)
26 .to_string();
27 async move {
28 let tasks = sqlx::query_file_as!(
29 PgTaskRow,
30 "queries/backend/list_jobs.sql",
31 status,
32 queue,
33 limit,
34 offset
35 )
36 .fetch_all(&pool)
37 .await?
38 .into_iter()
39 .map(|r| r.try_into())
40 .collect::<Result<Vec<_>, _>>()?;
41 Ok(tasks)
42 }
43 }
44}
45
46impl<Args> ListAllTasks for PostgresStorage<Args>
47where
48 PostgresStorage<Args>: Backend<Error = Error>,
49{
50 fn list_all_tasks(
51 &self,
52 filter: &Filter,
53 ) -> impl Future<Output = Result<Vec<Task<Self::Compact>>, Self::Error>> + Send {
54 let status = filter
55 .status
56 .as_ref()
57 .map(|s| s.to_string())
58 .unwrap_or(Status::Pending.to_string());
59 let pool = self.persistence.pool.clone();
60 let limit = filter.limit() as i64;
61 let offset = filter.offset() as i64;
62 async move {
63 let tasks = sqlx::query_file_as!(
64 PgTaskRow,
65 "queries/backend/list_all_jobs.sql",
66 status,
67 limit,
68 offset
69 )
70 .fetch_all(&pool)
71 .await?
72 .into_iter()
73 .map(|r| r.try_into())
74 .collect::<Result<Vec<_>, _>>()?;
75 Ok(tasks)
76 }
77 }
78}