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 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
18/// Fixed primary key for the singleton leader election row.
19pub 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}