Skip to main content

chronon_backend_sql_common/
coordinator.rs

1//! Revisions, scripts, leader election, partitions, and workers.
2//!
3//! Internal — used by [`SqlSchedulerStore`](crate::SqlSchedulerStore); not a stable public API.
4
5use 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
20/// Fixed primary key for the singleton leader election row.
21pub 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}