Skip to main content

apalis_sqlite/queries/
list_workers.rs

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