Skip to main content

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}