Skip to main content

apalis_postgres/queries/
list_queues.rs

1use apalis_core::backend::{Backend, ListQueues, QueueInfo};
2use serde_json::Value;
3
4use crate::{PostgresStorage, error::Error};
5
6impl<Args> ListQueues for PostgresStorage<Args>
7where
8    PostgresStorage<Args>: Backend<Error = Error>,
9{
10    fn list_queues(&self) -> impl Future<Output = Result<Vec<QueueInfo>, Self::Error>> + Send {
11        let pool = self.persistence.pool.clone();
12        struct QueueInfoRow {
13            pub name: Option<String>,
14            pub stats: Option<Value>,
15            pub workers: Option<Value>,
16            pub activity: Option<Value>,
17        }
18
19        async move {
20            let queues = sqlx::query_file_as!(QueueInfoRow, "queries/backend/list_queues.sql")
21                .fetch_all(&pool)
22                .await?
23                .into_iter()
24                .map(|row| QueueInfo {
25                    name: row.name.unwrap_or_default(),
26                    stats: serde_json::from_value(row.stats.unwrap()).unwrap_or_default(),
27                    workers: serde_json::from_value(row.workers.unwrap()).unwrap_or_default(),
28                    activity: serde_json::from_value(row.activity.unwrap()).unwrap_or_default(),
29                })
30                .collect();
31            Ok(queues)
32        }
33    }
34}