chronon_backend_sql_common/row/
worker.rs1use 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
11pub 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 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 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
54pub 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}