chronon_backend_sql_common/
coordinator.rs1use chrono::{DateTime, Duration, Utc};
6
7use chronon_core::error::{ChrononError, Result};
8use chronon_core::models::{JobRevision, PartitionAssignment, SchedulerLeader, Script, Worker};
9use sqlx::Row;
10
11use crate::error_map::map_err;
12use crate::row::{
13 row_to_leader, row_to_partition, row_to_revision, row_to_script, JobRevisionRow,
14 PartitionAssignmentRow, ScriptRow, WorkerRow,
15};
16use crate::{bind_sql, sql_execute, sql_fetch_all_map, sql_fetch_optional_map, SqlSchedulerStore};
17
18pub const LEADER_ROW_ID: &str = "singleton";
20
21pub(crate) async fn append_revision(
22 store: &SqlSchedulerStore,
23 revision: &JobRevision,
24) -> Result<()> {
25 let row = JobRevisionRow::from_model(revision)?;
26 let sql = bind_sql(
27 store.dialect,
28 "INSERT INTO chronon_job_revision (
29 revision_id, job_id, revision_number, changed_at, changed_by_actor_json, snapshot_json
30 ) VALUES (?, ?, ?, ?, ?, ?)",
31 );
32 sql_execute!(store, &sql, |q| {
33 q.bind(&row.revision_id)
34 .bind(&row.job_id)
35 .bind(row.revision_number)
36 .bind(row.changed_at)
37 .bind(&row.changed_by_actor_json)
38 .bind(&row.snapshot_json)
39 })
40}
41
42pub(crate) async fn list_revisions(
43 store: &SqlSchedulerStore,
44 job_id: &str,
45) -> Result<Vec<JobRevision>> {
46 let sql = bind_sql(
47 store.dialect,
48 "SELECT * FROM chronon_job_revision WHERE job_id = ? ORDER BY revision_number ASC",
49 );
50 sql_fetch_all_map!(store, &sql, |q| q.bind(job_id), |r| row_to_revision(r))
51}
52
53pub(crate) async fn upsert_script(store: &SqlSchedulerStore, script: &Script) -> Result<()> {
54 let row = ScriptRow::from_model(script)?;
55 let sql = bind_sql(
56 store.dialect,
57 "INSERT INTO chronon_script (script_id, script_name, signature_json, signature_hash, created_at)
58 VALUES (?, ?, ?, ?, ?)
59 ON CONFLICT (script_name) DO UPDATE SET
60 script_id = excluded.script_id,
61 signature_json = excluded.signature_json,
62 signature_hash = excluded.signature_hash,
63 created_at = excluded.created_at",
64 );
65 sql_execute!(store, &sql, |q| {
66 q.bind(&row.script_id)
67 .bind(&row.script_name)
68 .bind(&row.signature_json)
69 .bind(&row.signature_hash)
70 .bind(row.created_at)
71 })
72}
73
74pub(crate) async fn get_script(
75 store: &SqlSchedulerStore,
76 script_name: &str,
77) -> Result<Option<Script>> {
78 let sql = bind_sql(
79 store.dialect,
80 "SELECT * FROM chronon_script WHERE script_name = ?",
81 );
82 sql_fetch_optional_map!(store, &sql, |q| q.bind(script_name), |r| row_to_script(&r))
83}
84
85pub(crate) async fn try_acquire_leader(
86 store: &SqlSchedulerStore,
87 instance_id: &str,
88 ttl_secs: i64,
89) -> Result<bool> {
90 let now = Utc::now();
91 let until = now + Duration::seconds(ttl_secs);
92 let sql = bind_sql(
93 store.dialect,
94 "INSERT INTO chronon_scheduler_leader (
95 leader_id, leader_instance_id, leader_lease_until, last_heartbeat_at
96 ) VALUES (?, ?, ?, ?)
97 ON CONFLICT (leader_id) DO UPDATE SET
98 leader_instance_id = excluded.leader_instance_id,
99 leader_lease_until = excluded.leader_lease_until,
100 last_heartbeat_at = excluded.last_heartbeat_at
101 WHERE chronon_scheduler_leader.leader_lease_until <= excluded.last_heartbeat_at
102 OR chronon_scheduler_leader.leader_instance_id = excluded.leader_instance_id
103 RETURNING leader_instance_id",
104 );
105 let acquired: Option<String> = match &store.pool {
106 crate::SqlPool::Sqlite(pool) => {
107 let q = sqlx::query(&sql)
108 .bind(LEADER_ROW_ID)
109 .bind(instance_id)
110 .bind(until)
111 .bind(now);
112 match q.fetch_optional(pool).await.map_err(map_err)? {
113 Some(row) => Some(row.try_get("leader_instance_id").map_err(map_err)?),
114 None => None,
115 }
116 }
117 crate::SqlPool::Postgres(pool) => {
118 let q = sqlx::query(&sql)
119 .bind(LEADER_ROW_ID)
120 .bind(instance_id)
121 .bind(until)
122 .bind(now);
123 match q.fetch_optional(pool).await.map_err(map_err)? {
124 Some(row) => Some(row.try_get("leader_instance_id").map_err(map_err)?),
125 None => None,
126 }
127 }
128 };
129 Ok(acquired.as_deref() == Some(instance_id))
130}
131
132pub(crate) async fn renew_leader_lease(
133 store: &SqlSchedulerStore,
134 instance_id: &str,
135 ttl_secs: i64,
136) -> Result<()> {
137 let now = Utc::now();
138 let until = now + Duration::seconds(ttl_secs);
139 let sql = bind_sql(
140 store.dialect,
141 "UPDATE chronon_scheduler_leader SET leader_lease_until = ?, last_heartbeat_at = ?
142 WHERE leader_id = ? AND leader_instance_id = ?",
143 );
144 sql_execute!(store, &sql, |q| {
145 q.bind(until)
146 .bind(now)
147 .bind(LEADER_ROW_ID)
148 .bind(instance_id)
149 })
150}
151
152pub(crate) async fn get_leader(store: &SqlSchedulerStore) -> Result<Option<SchedulerLeader>> {
153 let sql = bind_sql(
154 store.dialect,
155 "SELECT * FROM chronon_scheduler_leader WHERE leader_id = ?",
156 );
157 sql_fetch_optional_map!(store, &sql, |q| q.bind(LEADER_ROW_ID), |r| row_to_leader(
158 &r
159 ))
160}
161
162pub(crate) async fn upsert_partition_assignment(
163 store: &SqlSchedulerStore,
164 assignment: &PartitionAssignment,
165) -> Result<()> {
166 let row = PartitionAssignmentRow::from_model(assignment);
167 let sql = bind_sql(
168 store.dialect,
169 "INSERT INTO chronon_partition_assignment (
170 partition_id, owner_instance_id, lease_until, updated_at
171 ) VALUES (?, ?, ?, ?)
172 ON CONFLICT (partition_id) DO UPDATE SET
173 owner_instance_id = excluded.owner_instance_id,
174 lease_until = excluded.lease_until,
175 updated_at = excluded.updated_at",
176 );
177 sql_execute!(store, &sql, |q| {
178 q.bind(&row.partition_id)
179 .bind(&row.owner_instance_id)
180 .bind(row.lease_until)
181 .bind(row.updated_at)
182 })
183}
184
185pub(crate) async fn list_partition_assignments(
186 store: &SqlSchedulerStore,
187) -> Result<Vec<PartitionAssignment>> {
188 let sql = bind_sql(
189 store.dialect,
190 "SELECT * FROM chronon_partition_assignment ORDER BY partition_id ASC",
191 );
192 sql_fetch_all_map!(store, &sql, |q| q, |r| row_to_partition(r))
193}
194
195pub(crate) async fn register_worker(store: &SqlSchedulerStore, worker: &Worker) -> Result<()> {
196 let row = WorkerRow::from_model(worker)?;
197 let sql = bind_sql(
198 store.dialect,
199 "INSERT INTO chronon_worker (
200 worker_id, pool_id, cell_id, status, last_heartbeat_at, capacity_json, created_at, updated_at
201 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?)
202 ON CONFLICT (worker_id) DO UPDATE SET
203 pool_id = excluded.pool_id,
204 cell_id = excluded.cell_id,
205 status = excluded.status,
206 last_heartbeat_at = excluded.last_heartbeat_at,
207 capacity_json = excluded.capacity_json,
208 updated_at = excluded.updated_at",
209 );
210 sql_execute!(store, &sql, |q| {
211 q.bind(&row.worker_id)
212 .bind(&row.pool_id)
213 .bind(&row.cell_id)
214 .bind(&row.status)
215 .bind(row.last_heartbeat_at)
216 .bind(&row.capacity_json)
217 .bind(row.created_at)
218 .bind(row.updated_at)
219 })
220}
221
222pub(crate) async fn heartbeat_worker(
223 store: &SqlSchedulerStore,
224 worker_id: &str,
225 at: DateTime<Utc>,
226) -> Result<()> {
227 let sql = bind_sql(
228 store.dialect,
229 "UPDATE chronon_worker SET last_heartbeat_at = ?, updated_at = ? WHERE worker_id = ?",
230 );
231 let rows = match &store.pool {
232 crate::SqlPool::Sqlite(pool) => sqlx::query(&sql)
233 .bind(at)
234 .bind(at)
235 .bind(worker_id)
236 .execute(pool)
237 .await
238 .map_err(map_err)?
239 .rows_affected(),
240 crate::SqlPool::Postgres(pool) => sqlx::query(&sql)
241 .bind(at)
242 .bind(at)
243 .bind(worker_id)
244 .execute(pool)
245 .await
246 .map_err(map_err)?
247 .rows_affected(),
248 };
249 if rows == 0 {
250 return Err(ChrononError::Internal(format!(
251 "worker not registered: {worker_id}"
252 )));
253 }
254 Ok(())
255}