Skip to main content

apalis_sqlite/queries/
list_queues.rs

1use apalis_core::backend::{Backend, ListQueues, QueueInfo};
2
3use crate::{SqliteStorage, error::Error};
4
5struct QueueInfoRow {
6    name: String,
7    stats: Option<String>,    // JSON string
8    workers: Option<String>,  // JSON string
9    activity: Option<String>, // JSON string
10}
11
12impl From<QueueInfoRow> for QueueInfo {
13    fn from(row: QueueInfoRow) -> Self {
14        Self {
15            name: row.name,
16            stats: row
17                .stats
18                .and_then(|s| serde_json::from_str(&s).ok())
19                .unwrap_or_default(),
20            workers: row
21                .workers
22                .and_then(|s| serde_json::from_str(&s).ok())
23                .unwrap_or_default(),
24            activity: row
25                .activity
26                .and_then(|s| serde_json::from_str(&s).ok())
27                .unwrap_or_default(),
28        }
29    }
30}
31
32impl<Args> ListQueues for SqliteStorage<Args>
33where
34    Self: Backend<Error = Error>,
35{
36    fn list_queues(&self) -> impl Future<Output = Result<Vec<QueueInfo>, Self::Error>> + Send {
37        let pool = self.persistence.pool.clone();
38
39        async move {
40            let queues = sqlx::query_file_as!(QueueInfoRow, "queries/backend/list_queues.sql")
41                .fetch_all(&pool)
42                .await?
43                .into_iter()
44                .map(QueueInfo::from)
45                .collect();
46            Ok(queues)
47        }
48    }
49}