apalis_postgres/queries/
list_queues.rs1use 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}