pub mod metadata;
pub mod scheduler;
pub mod status;
pub mod worker;
use std::sync::Arc;
pub use scheduler::{ReconcileDeps, reconcile_once, run_reconcile};
pub use status::CronStatus;
pub use worker::{WorkerDeps, run_worker};
pub struct CronDeps {
pub runtime: Arc<crate::runtime::RuntimeHandle>,
pub repo: Arc<dyn crate::storage::repositories::cron::CronRepository>,
pub trace_repo: Arc<dyn crate::storage::repositories::traces::TraceSink>,
pub persistence_queue: crate::queue::TracePersistenceQueue,
pub global_trace_storage: crate::config::TraceStorageConfig,
pub datalogic: Arc<dataflow_rs::datalogic_rs::Engine>,
pub vars: Option<Arc<serde_json::Value>>,
pub instance_id: String,
pub status: Arc<CronStatus>,
pub config: crate::config::CronConfig,
pub max_result_size_bytes: usize,
pub lease_gate: Option<Arc<crate::cluster::JobLeaseGate>>,
}
pub fn start(tasks: &crate::runtime::TaskRegistry, deps: CronDeps) {
if !deps.config.enabled {
tracing::info!("Cron scheduler disabled (cron.enabled = false)");
return;
}
let CronDeps {
runtime,
repo,
trace_repo,
persistence_queue,
global_trace_storage,
datalogic,
vars,
instance_id,
status,
config,
max_result_size_bytes,
lease_gate,
} = deps;
let reconcile = Arc::new(ReconcileDeps {
runtime: runtime.clone(),
repo: repo.clone(),
status: status.clone(),
config: config.clone(),
lease_gate,
});
tasks.supervise(
"cron_reconcile",
crate::runtime::Criticality::Required,
move |shutdown| run_reconcile(reconcile.clone(), shutdown),
);
let worker = Arc::new(WorkerDeps {
runtime,
repo,
trace_repo,
persistence_queue,
global_trace_storage,
datalogic,
vars,
instance_id,
status,
config: config.clone(),
max_result_size_bytes,
});
tasks.supervise(
"cron_worker",
crate::runtime::Criticality::Required,
move |shutdown| run_worker(worker.clone(), shutdown),
);
tracing::info!(
poll_interval_ms = config.poll_interval_ms,
workers = config.workers,
claim_batch_size = config.claim_batch_size,
"Cron scheduler started"
);
}
pub fn start_cleanup(
tasks: &crate::runtime::TaskRegistry,
retention_hours: u64,
interval_secs: u64,
repo: Arc<dyn crate::storage::repositories::cron::CronRepository>,
lease_gate: Option<Arc<crate::cluster::JobLeaseGate>>,
) {
if retention_hours == 0 {
return;
}
crate::queue::supervise_retention_job(
tasks,
"cron_cleanup",
interval_secs,
lease_gate,
move || {
let repo = repo.clone();
async move { repo.delete_terminal_older_than(retention_hours).await }
},
move |outcome| match outcome {
Ok(count) if count > 0 => {
tracing::info!(
deleted = count,
retention_hours,
"Cron occurrence cleanup completed"
)
}
Ok(_) => {}
Err(e) => tracing::error!(error = %e, "Cron occurrence cleanup failed"),
},
);
}