chronon_backend_mem/store/mod.rs
1//! In-memory [`SchedulerStore`] implementation split by concern.
2
3mod claims;
4mod coordinator;
5mod jobs;
6mod runs;
7mod trait_impl;
8
9use std::collections::HashMap;
10
11use chrono::{DateTime, Utc};
12use chronon_core::error::{ChrononError, Result};
13use chronon_core::models::{
14 Job, JobRevision, PartitionAssignment, Run, SchedulerLeader, Script, Worker,
15};
16use parking_lot::RwLock;
17
18/// Fixed primary key for the singleton leader election row.
19pub const LEADER_ROW_ID: &str = "singleton";
20
21/// Thread-safe in-memory persistence for jobs, runs, and coordinator metadata.
22///
23/// Process-local and **non-durable** — suitable for embedded experiments, unit tests,
24/// examples, and benchmarks. Do **not** share across processes; use SQLite or Postgres for
25/// durable / coordinator–worker topologies.
26///
27/// Enable the public crate `mem` feature to re-export this type from `chronon`. Wire with
28/// `ChrononBuilder::scheduler_store(Arc::new(InMemorySchedulerStore::new()))` or
29/// [`crate::install_default_mem_store`].
30///
31/// # Locking
32///
33/// Maps are guarded by [`parking_lot::RwLock`] for short, synchronous critical sections.
34/// **Never hold a lock across `.await`** — take the lock, clone or mutate, drop the guard,
35/// then await. This keeps Tokio worker threads unblocked under contention.
36///
37/// # Examples
38///
39/// ```
40/// use std::sync::Arc;
41/// use chronon_backend_mem::InMemorySchedulerStore;
42/// use chronon_core::{Job, SchedulerStore};
43///
44/// # #[tokio::main]
45/// # async fn main() -> chronon_core::Result<()> {
46/// let store = Arc::new(InMemorySchedulerStore::new());
47/// store.upsert_job(&Job::new("demo", "noop")).await?;
48/// assert_eq!(store.list_jobs().await?.len(), 1);
49/// # Ok(())
50/// # }
51/// ```
52#[derive(Default)]
53pub struct InMemorySchedulerStore {
54 /// Jobs keyed by [`Job::job_id`].
55 pub(super) jobs: RwLock<HashMap<String, Job>>,
56 /// Secondary index from [`Job::job_name`] to job id.
57 pub(super) jobs_by_name: RwLock<HashMap<String, String>>,
58 /// Runs keyed by [`Run::run_id`].
59 pub(super) runs: RwLock<HashMap<String, Run>>,
60 /// Job revision history keyed by job id.
61 pub(super) revisions: RwLock<HashMap<String, Vec<JobRevision>>>,
62 /// Script metadata keyed by script name.
63 pub(super) scripts: RwLock<HashMap<String, Script>>,
64 /// Singleton scheduler leader row.
65 pub(super) leader: RwLock<Option<SchedulerLeader>>,
66 /// Partition assignments keyed by partition id.
67 pub(super) partitions: RwLock<HashMap<String, PartitionAssignment>>,
68 /// Registered workers keyed by worker id.
69 pub(super) workers: RwLock<HashMap<String, Worker>>,
70}
71
72impl InMemorySchedulerStore {
73 /// Creates an empty store.
74 ///
75 /// # Examples
76 ///
77 /// ```
78 /// # #[tokio::main]
79 /// # async fn main() {
80 /// use chronon_backend_mem::InMemorySchedulerStore;
81 /// use chronon_core::SchedulerStore;
82 ///
83 /// let store = InMemorySchedulerStore::new();
84 /// assert!(store.list_jobs().await.unwrap().is_empty());
85 /// # }
86 /// ```
87 pub fn new() -> Self {
88 Self::default()
89 }
90
91 /// Insert or replace a job and update the name index.
92 pub(super) fn write_job(&self, job: Job) -> Result<()> {
93 let mut jobs = self.jobs.write();
94 let mut by_name = self.jobs_by_name.write();
95 by_name.insert(job.job_name.clone(), job.job_id.clone());
96 jobs.insert(job.job_id.clone(), job);
97 Ok(())
98 }
99
100 /// Mutate an existing job and stamp `updated_at` to now.
101 pub(super) fn mutate_job<F>(&self, job_id: &str, f: F) -> Result<()>
102 where
103 F: FnOnce(&mut Job),
104 {
105 let mut jobs = self.jobs.write();
106 let job = jobs
107 .get_mut(job_id)
108 .ok_or_else(|| ChrononError::JobNotFound(job_id.to_string()))?;
109 f(job);
110 job.updated_at = Utc::now();
111 Ok(())
112 }
113
114 /// Mutate an existing job and stamp `updated_at` to `now`.
115 pub(super) fn mutate_job_at<F>(&self, job_id: &str, now: DateTime<Utc>, f: F) -> Result<()>
116 where
117 F: FnOnce(&mut Job),
118 {
119 let mut jobs = self.jobs.write();
120 let job = jobs
121 .get_mut(job_id)
122 .ok_or_else(|| ChrononError::JobNotFound(job_id.to_string()))?;
123 f(job);
124 job.updated_at = now;
125 Ok(())
126 }
127}