apalis_sqlite/queries/
list_queues.rs1use apalis_core::backend::{Backend, ListQueues, QueueInfo};
2
3use crate::{SqliteStorage, error::Error};
4
5struct QueueInfoRow {
6 name: String,
7 stats: Option<String>, workers: Option<String>, activity: Option<String>, }
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}