use std::collections::HashMap;
use std::sync::Mutex;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
use serde::Serialize;
const JOB_TTL: Duration = Duration::from_secs(300);
const RUNNING_JOB_MAX_AGE: Duration = Duration::from_secs(86400);
enum JobState {
Running,
Done {
output_key: String,
rows_written: u64,
},
Failed {
summary: String,
hint: String,
raw: String,
},
}
struct JobEntry {
state: JobState,
started: Instant,
finished: Option<Instant>,
}
#[derive(Debug, Serialize)]
pub struct JobStatusResponse {
pub job_id: String,
pub state: &'static str,
pub elapsed_ms: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub output_key: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub rows_written: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub summary: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub hint: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub raw: Option<String>,
}
pub struct JobRegistry {
inner: Mutex<HashMap<String, JobEntry>>,
next: AtomicU64,
}
impl Default for JobRegistry {
fn default() -> Self {
Self::new()
}
}
impl JobRegistry {
pub fn new() -> Self {
Self {
inner: Mutex::new(HashMap::new()),
next: AtomicU64::new(1),
}
}
pub fn register(&self) -> String {
let id = self.next.fetch_add(1, Ordering::Relaxed).to_string();
let mut map = self.inner.lock().unwrap_or_else(|p| p.into_inner());
let now = Instant::now();
map.retain(|_, e| match e.finished {
Some(f) => now.saturating_duration_since(f) < JOB_TTL,
None => now.saturating_duration_since(e.started) < RUNNING_JOB_MAX_AGE,
});
map.insert(
id.clone(),
JobEntry {
state: JobState::Running,
started: now,
finished: None,
},
);
id
}
pub fn complete(&self, id: &str, output_key: String, rows_written: u64) {
let mut map = self.inner.lock().unwrap_or_else(|p| p.into_inner());
if let Some(e) = map.get_mut(id) {
e.state = JobState::Done {
output_key,
rows_written,
};
e.finished = Some(Instant::now());
}
}
pub fn fail(&self, id: &str, summary: String, hint: String, raw: String) {
let mut map = self.inner.lock().unwrap_or_else(|p| p.into_inner());
if let Some(e) = map.get_mut(id) {
e.state = JobState::Failed { summary, hint, raw };
e.finished = Some(Instant::now());
}
}
pub fn status(&self, id: &str) -> Option<JobStatusResponse> {
let mut map = self.inner.lock().unwrap_or_else(|p| p.into_inner());
let now = Instant::now();
if matches!(
map.get(id),
Some(e)
if e.finished.is_none()
&& now.saturating_duration_since(e.started) >= RUNNING_JOB_MAX_AGE
) {
map.remove(id);
return None;
}
let e = map.get(id)?;
let elapsed_ms = match e.finished {
Some(f) => f.saturating_duration_since(e.started).as_millis() as u64,
None => e.started.elapsed().as_millis() as u64,
};
let resp = match &e.state {
JobState::Running => JobStatusResponse {
job_id: id.to_owned(),
state: "running",
elapsed_ms,
output_key: None,
rows_written: None,
summary: None,
hint: None,
raw: None,
},
JobState::Done {
output_key,
rows_written,
} => JobStatusResponse {
job_id: id.to_owned(),
state: "done",
elapsed_ms,
output_key: Some(output_key.clone()),
rows_written: Some(*rows_written),
summary: None,
hint: None,
raw: None,
},
JobState::Failed { summary, hint, raw } => JobStatusResponse {
job_id: id.to_owned(),
state: "failed",
elapsed_ms,
output_key: None,
rows_written: None,
summary: Some(summary.clone()),
hint: Some(hint.clone()),
raw: Some(raw.clone()),
},
};
Some(resp)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn register_returns_unique_ids() {
let reg = JobRegistry::new();
let a = reg.register();
let b = reg.register();
assert_ne!(a, b);
}
#[test]
fn new_job_is_running() {
let reg = JobRegistry::new();
let id = reg.register();
let s = reg.status(&id).unwrap();
assert_eq!(s.state, "running");
assert!(s.output_key.is_none());
assert!(s.rows_written.is_none());
assert!(s.summary.is_none());
}
#[test]
fn complete_transitions_to_done() {
let reg = JobRegistry::new();
let id = reg.register();
reg.complete(&id, "out.parquet".into(), 42);
let s = reg.status(&id).unwrap();
assert_eq!(s.state, "done");
assert_eq!(s.output_key.as_deref(), Some("out.parquet"));
assert_eq!(s.rows_written, Some(42));
}
#[test]
fn fail_transitions_to_failed() {
let reg = JobRegistry::new();
let id = reg.register();
reg.fail(&id, "oops".into(), "try X".into(), "raw err".into());
let s = reg.status(&id).unwrap();
assert_eq!(s.state, "failed");
assert_eq!(s.summary.as_deref(), Some("oops"));
assert_eq!(s.hint.as_deref(), Some("try X"));
assert_eq!(s.raw.as_deref(), Some("raw err"));
}
#[test]
fn unknown_id_returns_none() {
let reg = JobRegistry::new();
assert!(reg.status("999").is_none());
}
#[test]
fn complete_on_unknown_id_is_noop() {
let reg = JobRegistry::new();
reg.complete("999", "x.parquet".into(), 0);
}
}