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