Skip to main content

chronon_backend_sql_common/row/
worker.rs

1//! [`Worker`] SQL row mapping.
2
3use chrono::{DateTime, Utc};
4use chronon_core::error::Result;
5use chronon_core::models::Worker;
6use sqlx::{ColumnIndex, Row};
7
8use super::{decode_json_opt, encode_json_opt, parse_worker_status, worker_status_to_str};
9use crate::error_map::map_err;
10
11/// SQL row shape for [`Worker`].
12pub struct WorkerRow {
13    pub(crate) worker_id: String,
14    pub(crate) pool_id: String,
15    pub(crate) cell_id: Option<String>,
16    pub(crate) status: String,
17    pub(crate) last_heartbeat_at: DateTime<Utc>,
18    pub(crate) capacity_json: Option<String>,
19    pub(crate) created_at: DateTime<Utc>,
20    pub(crate) updated_at: DateTime<Utc>,
21}
22
23#[allow(clippy::wrong_self_convention)]
24impl WorkerRow {
25    /// Build a row from a domain [`Worker`].
26    pub fn from_model(worker: &Worker) -> Result<Self> {
27        Ok(Self {
28            worker_id: worker.worker_id.clone(),
29            pool_id: worker.pool_id.clone(),
30            cell_id: worker.cell_id.clone(),
31            status: worker_status_to_str(worker.status).to_string(),
32            last_heartbeat_at: worker.last_heartbeat_at,
33            capacity_json: encode_json_opt(worker.capacity_json.as_ref())?,
34            created_at: worker.created_at,
35            updated_at: worker.updated_at,
36        })
37    }
38
39    /// Convert this row into a domain [`Worker`].
40    pub fn to_model(self) -> Result<Worker> {
41        Ok(Worker {
42            worker_id: self.worker_id,
43            pool_id: self.pool_id,
44            cell_id: self.cell_id,
45            status: parse_worker_status(&self.status)?,
46            last_heartbeat_at: self.last_heartbeat_at,
47            capacity_json: decode_json_opt(self.capacity_json)?,
48            created_at: self.created_at,
49            updated_at: self.updated_at,
50        })
51    }
52}
53
54/// Map a SQL row to a [`Worker`].
55pub fn row_to_worker<'r, R>(row: &'r R) -> Result<Worker>
56where
57    R: Row,
58    for<'i> &'i str: ColumnIndex<R>,
59    String: sqlx::Decode<'r, R::Database> + sqlx::Type<R::Database>,
60    DateTime<Utc>: sqlx::Decode<'r, R::Database> + sqlx::Type<R::Database>,
61    Option<String>: sqlx::Decode<'r, R::Database> + sqlx::Type<R::Database>,
62{
63    let status: String = row.try_get("status").map_err(map_err)?;
64    WorkerRow {
65        worker_id: row.try_get("worker_id").map_err(map_err)?,
66        pool_id: row.try_get("pool_id").map_err(map_err)?,
67        cell_id: row.try_get("cell_id").map_err(map_err)?,
68        status,
69        last_heartbeat_at: row.try_get("last_heartbeat_at").map_err(map_err)?,
70        capacity_json: row.try_get("capacity_json").map_err(map_err)?,
71        created_at: row.try_get("created_at").map_err(map_err)?,
72        updated_at: row.try_get("updated_at").map_err(map_err)?,
73    }
74    .to_model()
75}
76
77#[cfg(test)]
78mod tests {
79    use chrono::Utc;
80    use chronon_core::models::{Worker, WorkerStatus};
81
82    use super::WorkerRow;
83
84    #[test]
85    fn worker_row_roundtrip() {
86        let worker = Worker {
87            worker_id: "w1".into(),
88            pool_id: "general".into(),
89            cell_id: None,
90            status: WorkerStatus::Online,
91            last_heartbeat_at: Utc::now(),
92            capacity_json: None,
93            created_at: Utc::now(),
94            updated_at: Utc::now(),
95        };
96        let row = WorkerRow::from_model(&worker).expect("row");
97        let back = row.to_model().expect("model");
98        assert_eq!(back.worker_id, "w1");
99    }
100}