apalis_sqlite/queries/
list_workers.rs1use 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}