Skip to main content

apalis_postgres/queries/
list_workers.rs

1use apalis_core::backend::{Backend, ListWorkers, RunningWorker};
2
3use futures::TryFutureExt;
4
5use crate::{PostgresStorage, error::Error, timestamp::Timestamp};
6
7#[derive(Debug)]
8pub struct WorkerRow {
9    pub id: String,
10    pub worker_type: String,
11    pub storage_name: String,
12    pub layers: Option<String>,
13    pub last_seen: Timestamp,
14    pub started_at: Option<Timestamp>,
15}
16
17impl<Args: Sync> ListWorkers for PostgresStorage<Args>
18where
19    PostgresStorage<Args>: Backend<Error = Error>,
20{
21    fn list_workers(&self) -> impl Future<Output = Result<Vec<RunningWorker>, Self::Error>> + Send {
22        let queue = self.persistence.config.queue.to_string();
23
24        let pool = self.persistence.pool.clone();
25        let limit = 100;
26        let offset = 0;
27        async move {
28            let workers = sqlx::query_file_as!(
29                WorkerRow,
30                "queries/backend/list_workers.sql",
31                queue,
32                limit,
33                offset
34            )
35            .fetch_all(&pool)
36            .map_ok(|w| {
37                w.into_iter()
38                    .map(|w| RunningWorker {
39                        id: w.id,
40                        backend: w.storage_name,
41                        started_at: w.started_at.unwrap_or_default().0,
42                        last_heartbeat: w.last_seen.0,
43                        layers: w.layers.unwrap_or_default(),
44                        queue: w.worker_type,
45                    })
46                    .collect()
47            })
48            .await?;
49            Ok(workers)
50        }
51    }
52
53    fn list_all_workers(
54        &self,
55    ) -> impl Future<Output = Result<Vec<RunningWorker>, Self::Error>> + Send {
56        let pool = self.persistence.pool.clone();
57        let limit = 100;
58        let offset = 0;
59        async move {
60            let workers = sqlx::query_file_as!(
61                WorkerRow,
62                "queries/backend/list_all_workers.sql",
63                limit,
64                offset
65            )
66            .fetch_all(&pool)
67            .map_ok(|w| {
68                w.into_iter()
69                    .map(|w| RunningWorker {
70                        id: w.id,
71                        backend: w.storage_name,
72                        started_at: w.started_at.unwrap_or_default().0,
73                        last_heartbeat: w.last_seen.0,
74                        layers: w.layers.unwrap_or_default(),
75                        queue: w.worker_type,
76                    })
77                    .collect()
78            })
79            .await?;
80            Ok(workers)
81        }
82    }
83}