Skip to main content

apalis_sqlite/queries/
list_tasks.rs

1use apalis_core::{
2    backend::{Backend, Filter, ListAllTasks, ListTasks},
3    task::status::Status,
4};
5
6use crate::{SqliteStorage, SqliteTask};
7use crate::{error::Error, from_row::SqliteTaskRow};
8
9impl<Args> ListTasks for SqliteStorage<Args>
10where
11    Self: Backend<Error = Error>,
12    Args: 'static,
13{
14    fn list_tasks(
15        &self,
16        filter: &Filter,
17    ) -> impl Future<Output = Result<Vec<SqliteTask>, Self::Error>> + Send {
18        let queue = self.persistence.config.queue.as_ref();
19        let pool = self.persistence.pool.clone();
20        let limit = filter.limit() as i32;
21        let offset = filter.offset() as i32;
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                SqliteTaskRow,
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 SqliteStorage<Args>
47where
48    Self: Backend<Error = Error>,
49{
50    fn list_all_tasks(
51        &self,
52        filter: &Filter,
53    ) -> impl Future<Output = Result<Vec<SqliteTask>, 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 i32;
61        let offset = filter.offset() as i32;
62        async move {
63            let tasks = sqlx::query_file_as!(
64                SqliteTaskRow,
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}