use std::sync::Arc;
use chrono::Utc;
use chronon_core::models::RunStatus;
use chronon_core::store::SchedulerStore;
use chronon_executor::ExecutorEvent;
use chronon_telemetry::CapturedLogs;
use tracing::{debug, warn};
use crate::retry::finalize_failed_run;
pub async fn handle_executor_event(store: &Arc<dyn SchedulerStore>, event: ExecutorEvent) {
match event {
ExecutorEvent::RunStarted { run_id } => {
if let Ok(Some(mut run)) = store.get_run(&run_id).await {
if !matches!(run.status, RunStatus::Queued | RunStatus::Claimed) {
debug!(
run_id = %run_id,
status = %run.status,
"ignoring RunStarted for non-startable run"
);
return;
}
run.started_at = Some(Utc::now());
run.status = RunStatus::Running;
if let Err(e) = store.update_run(&run).await {
warn!(
run_id = %run_id,
error = %e,
"failed to persist RunStarted status"
);
}
}
}
ExecutorEvent::RunCompleted {
run_id,
duration_ms,
logs,
} => {
if let Ok(Some(mut run)) = store.get_run(&run_id).await {
if run.status != RunStatus::Running {
debug!(
run_id = %run_id,
status = %run.status,
"ignoring RunCompleted for non-running run"
);
return;
}
run.complete();
run.duration_ms = Some(duration_ms);
apply_logs(&mut run, logs);
if let Err(e) = store.update_run(&run).await {
warn!(
run_id = %run_id,
error = %e,
"failed to persist RunCompleted status"
);
}
}
}
ExecutorEvent::RunFailed {
run_id,
error,
logs,
} => {
if let Ok(Some(run)) = store.get_run(&run_id).await {
if run.status != RunStatus::Running {
debug!(
run_id = %run_id,
status = %run.status,
"ignoring RunFailed for non-running run"
);
return;
}
let job = match run.job_id.as_deref() {
Some(id) => store.get_job(id).await.ok().flatten(),
None => None,
};
if let Some(job) = job {
finalize_failed_run(store, run, &job, RunStatus::Failed, error, Some(logs))
.await;
} else {
let safe = chronon_core::sanitize_error_message(
&chronon_core::redact_credentials_in_text(&error),
);
let mut run = run;
run.fail(safe.clone());
let mut captured = logs;
captured.ensure_stderr_message(&safe);
apply_logs(&mut run, captured);
if let Err(e) = store.update_run(&run).await {
warn!(
run_id = %run.run_id,
error = %e,
"failed to persist RunFailed status without job"
);
}
}
}
}
}
}
fn apply_logs(run: &mut chronon_core::models::Run, logs: CapturedLogs) {
if logs.stdout_text.is_some() {
run.stdout_text = logs.stdout_text;
}
if logs.stderr_text.is_some() {
run.stderr_text = logs.stderr_text;
}
}
pub fn spawn_event_handler(
store: Arc<dyn SchedulerStore>,
mut event_rx: tokio::sync::mpsc::Receiver<ExecutorEvent>,
) {
tokio::spawn(async move {
while let Some(event) = event_rx.recv().await {
handle_executor_event(&store, event).await;
}
});
}
#[cfg(test)]
mod tests {
#![allow(clippy::unwrap_used, clippy::expect_used)]
use std::sync::Arc;
use chronon_backend_mem::InMemorySchedulerStore;
use chronon_core::models::{Job, Run, RunStatus, ScheduleKind};
use chronon_core::store::SchedulerStore;
use chronon_executor::ExecutorEvent;
use chronon_telemetry::CapturedLogs;
use super::handle_executor_event;
async fn seed_queued_run(store: &Arc<dyn SchedulerStore>) -> (Job, Run) {
let mut job = Job::new("evt-job", "noop");
job.schedule_kind = ScheduleKind::Manual;
store.upsert_job(&job).await.unwrap();
let mut run = Run::for_job(&job.job_id, &job.script_name, chrono::Utc::now());
run.actor_json = serde_json::json!({"user": "alice"});
store.create_run(&run).await.unwrap();
(job, run)
}
#[tokio::test]
async fn run_started_moves_queued_to_running() {
let store: Arc<dyn SchedulerStore> = Arc::new(InMemorySchedulerStore::new());
let (_job, run) = seed_queued_run(&store).await;
handle_executor_event(
&store,
ExecutorEvent::RunStarted {
run_id: run.run_id.clone(),
},
)
.await;
let updated = store.get_run(&run.run_id).await.unwrap().unwrap();
assert_eq!(updated.status, RunStatus::Running);
}
#[tokio::test]
async fn run_completed_on_queued_is_ignored() {
let store: Arc<dyn SchedulerStore> = Arc::new(InMemorySchedulerStore::new());
let (_job, run) = seed_queued_run(&store).await;
handle_executor_event(
&store,
ExecutorEvent::RunCompleted {
run_id: run.run_id.clone(),
duration_ms: 10,
logs: CapturedLogs::default(),
},
)
.await;
let updated = store.get_run(&run.run_id).await.unwrap().unwrap();
assert_eq!(updated.status, RunStatus::Queued);
assert!(updated.finished_at.is_none());
}
#[tokio::test]
async fn run_started_from_claimed_moves_to_running() {
let store: Arc<dyn SchedulerStore> = Arc::new(InMemorySchedulerStore::new());
let (_job, mut run) = seed_queued_run(&store).await;
run.status = RunStatus::Claimed;
store.update_run(&run).await.unwrap();
handle_executor_event(
&store,
ExecutorEvent::RunStarted {
run_id: run.run_id.clone(),
},
)
.await;
let updated = store.get_run(&run.run_id).await.unwrap().unwrap();
assert_eq!(updated.status, RunStatus::Running);
}
#[tokio::test]
async fn run_completed_from_running_succeeds() {
let store: Arc<dyn SchedulerStore> = Arc::new(InMemorySchedulerStore::new());
let (_job, mut run) = seed_queued_run(&store).await;
run.status = RunStatus::Running;
store.update_run(&run).await.unwrap();
handle_executor_event(
&store,
ExecutorEvent::RunCompleted {
run_id: run.run_id.clone(),
duration_ms: 42,
logs: CapturedLogs::default(),
},
)
.await;
let updated = store.get_run(&run.run_id).await.unwrap().unwrap();
assert_eq!(updated.status, RunStatus::Success);
assert_eq!(updated.duration_ms, Some(42));
}
#[tokio::test]
async fn run_failed_from_running_marks_failed() {
let store: Arc<dyn SchedulerStore> = Arc::new(InMemorySchedulerStore::new());
let (_job, mut run) = seed_queued_run(&store).await;
run.status = RunStatus::Running;
store.update_run(&run).await.unwrap();
handle_executor_event(
&store,
ExecutorEvent::RunFailed {
run_id: run.run_id.clone(),
error: "boom".into(),
logs: CapturedLogs::default(),
},
)
.await;
let updated = store.get_run(&run.run_id).await.unwrap().unwrap();
assert_eq!(updated.status, RunStatus::Failed);
}
#[tokio::test]
async fn run_failed_on_queued_is_ignored() {
let store: Arc<dyn SchedulerStore> = Arc::new(InMemorySchedulerStore::new());
let (_job, run) = seed_queued_run(&store).await;
handle_executor_event(
&store,
ExecutorEvent::RunFailed {
run_id: run.run_id.clone(),
error: "forged".into(),
logs: CapturedLogs::default(),
},
)
.await;
let updated = store.get_run(&run.run_id).await.unwrap().unwrap();
assert_eq!(updated.status, RunStatus::Queued);
assert!(updated.finished_at.is_none());
}
#[tokio::test]
async fn run_started_on_terminal_is_ignored() {
let store: Arc<dyn SchedulerStore> = Arc::new(InMemorySchedulerStore::new());
let (_job, mut run) = seed_queued_run(&store).await;
run.status = RunStatus::Success;
store.update_run(&run).await.unwrap();
handle_executor_event(
&store,
ExecutorEvent::RunStarted {
run_id: run.run_id.clone(),
},
)
.await;
let updated = store.get_run(&run.run_id).await.unwrap().unwrap();
assert_eq!(updated.status, RunStatus::Success);
assert!(updated.started_at.is_none());
}
}