mod claims;
mod coordinator;
mod jobs;
mod runs;
mod trait_impl;
use std::collections::HashMap;
use chrono::{DateTime, Utc};
use chronon_core::error::{ChrononError, Result};
use chronon_core::models::{
Job, JobRevision, PartitionAssignment, Run, SchedulerLeader, Script, Worker,
};
use parking_lot::RwLock;
pub const LEADER_ROW_ID: &str = "singleton";
#[derive(Default)]
pub struct InMemorySchedulerStore {
pub(super) jobs: RwLock<HashMap<String, Job>>,
pub(super) jobs_by_name: RwLock<HashMap<String, String>>,
pub(super) runs: RwLock<HashMap<String, Run>>,
pub(super) revisions: RwLock<HashMap<String, Vec<JobRevision>>>,
pub(super) scripts: RwLock<HashMap<String, Script>>,
pub(super) leader: RwLock<Option<SchedulerLeader>>,
pub(super) partitions: RwLock<HashMap<String, PartitionAssignment>>,
pub(super) workers: RwLock<HashMap<String, Worker>>,
}
impl InMemorySchedulerStore {
pub fn new() -> Self {
Self::default()
}
pub(super) fn write_job(&self, job: Job) -> Result<()> {
let mut jobs = self.jobs.write();
let mut by_name = self.jobs_by_name.write();
by_name.insert(job.job_name.clone(), job.job_id.clone());
jobs.insert(job.job_id.clone(), job);
Ok(())
}
pub(super) fn mutate_job<F>(&self, job_id: &str, f: F) -> Result<()>
where
F: FnOnce(&mut Job),
{
let mut jobs = self.jobs.write();
let job = jobs
.get_mut(job_id)
.ok_or_else(|| ChrononError::JobNotFound(job_id.to_string()))?;
f(job);
job.updated_at = Utc::now();
Ok(())
}
pub(super) fn mutate_job_at<F>(&self, job_id: &str, now: DateTime<Utc>, f: F) -> Result<()>
where
F: FnOnce(&mut Job),
{
let mut jobs = self.jobs.write();
let job = jobs
.get_mut(job_id)
.ok_or_else(|| ChrononError::JobNotFound(job_id.to_string()))?;
f(job);
job.updated_at = now;
Ok(())
}
}