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