use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
#[serde(rename_all = "snake_case")]
pub enum ScheduleKind {
#[default]
Cron,
RunOnce,
Manual,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RetryPolicy {
pub max_attempts: u32,
pub base_delay_ms: u64,
pub backoff_multiplier: f64,
pub max_delay_ms: u64,
}
impl Default for RetryPolicy {
fn default() -> Self {
Self {
max_attempts: 0,
base_delay_ms: 0,
backoff_multiplier: 1.0,
max_delay_ms: 0,
}
}
}
impl RetryPolicy {
pub fn should_retry(&self, attempt: i32) -> bool {
attempt > 0 && (attempt as u32) <= self.max_attempts
}
pub fn delay_ms_after(&self, failed_attempt: i32) -> u64 {
let exp = failed_attempt.saturating_sub(1).max(0);
let raw = (self.base_delay_ms as f64) * self.backoff_multiplier.powi(exp);
let ms = if raw.is_finite() && raw > 0.0 {
raw.min(u64::MAX as f64) as u64
} else {
0
};
if self.max_delay_ms == 0 {
ms
} else {
ms.min(self.max_delay_ms)
}
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct MisfirePolicy {
pub run_immediately: bool,
pub max_misfire_window_secs: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Job {
pub job_id: String,
pub job_name: String,
pub script_name: String,
pub script_sig_hash: String,
pub enabled: bool,
pub schedule_kind: ScheduleKind,
pub cron_expr: Option<String>,
pub timezone: Option<String>,
pub run_once_at: Option<DateTime<Utc>>,
pub run_once_claimed_at: Option<DateTime<Utc>>,
pub run_once_claimed_by: Option<String>,
pub run_once_completed_at: Option<DateTime<Utc>>,
pub run_once_claim_expires_at: Option<DateTime<Utc>>,
pub partition_hash: Option<i64>,
pub claim_lease_id: Option<String>,
pub claim_lease_until: Option<DateTime<Utc>>,
pub pool: Option<String>,
pub region: Option<String>,
pub placement_json: Option<Value>,
pub actor_json: Value,
pub params_json: Value,
pub concurrency: i32,
pub timeout_ms: Option<i64>,
pub retry_policy_json: Value,
pub misfire_policy_json: Value,
pub parent_limits_json: Option<Value>,
pub next_run_at: Option<DateTime<Utc>>,
pub current_revision: i32,
pub updated_at: DateTime<Utc>,
pub created_at: DateTime<Utc>,
}
impl Job {
pub fn new(job_name: impl Into<String>, script_name: impl Into<String>) -> Self {
let now = Utc::now();
Self {
job_id: uuid::Uuid::new_v4().to_string(),
job_name: job_name.into(),
script_name: script_name.into(),
script_sig_hash: String::new(),
enabled: true,
schedule_kind: ScheduleKind::default(),
cron_expr: None,
timezone: None,
run_once_at: None,
run_once_claimed_at: None,
run_once_claimed_by: None,
run_once_completed_at: None,
run_once_claim_expires_at: None,
partition_hash: None,
claim_lease_id: None,
claim_lease_until: None,
pool: None,
region: None,
placement_json: None,
actor_json: Value::Null,
params_json: Value::Object(serde_json::Map::default()),
concurrency: 1,
timeout_ms: None,
retry_policy_json: serde_json::to_value(RetryPolicy::default()).unwrap_or_default(),
misfire_policy_json: serde_json::to_value(MisfirePolicy::default()).unwrap_or_default(),
parent_limits_json: None,
next_run_at: None,
current_revision: 1,
updated_at: now,
created_at: now,
}
}
pub fn retry_policy(&self) -> RetryPolicy {
serde_json::from_value(self.retry_policy_json.clone()).unwrap_or_default()
}
pub fn misfire_policy(&self) -> MisfirePolicy {
serde_json::from_value(self.misfire_policy_json.clone()).unwrap_or_default()
}
pub fn set_retry_policy(&mut self, policy: &RetryPolicy) {
self.retry_policy_json = serde_json::to_value(policy).unwrap_or_default();
}
pub fn set_misfire_policy(&mut self, policy: &MisfirePolicy) {
self.misfire_policy_json = serde_json::to_value(policy).unwrap_or_default();
}
}
#[cfg(test)]
mod policy_tests {
use super::*;
#[test]
fn retry_default_does_not_retry() {
let p = RetryPolicy::default();
assert!(!p.should_retry(1));
assert_eq!(p.delay_ms_after(1), 0);
}
#[test]
fn retry_backoff_and_cap() {
let p = RetryPolicy {
max_attempts: 3,
base_delay_ms: 100,
backoff_multiplier: 2.0,
max_delay_ms: 250,
};
assert!(p.should_retry(1));
assert!(p.should_retry(3));
assert!(!p.should_retry(4));
assert_eq!(p.delay_ms_after(1), 100);
assert_eq!(p.delay_ms_after(2), 200);
assert_eq!(p.delay_ms_after(3), 250);
}
#[test]
fn job_policy_roundtrip() {
let mut job = Job::new("n", "s");
let retry = RetryPolicy {
max_attempts: 2,
base_delay_ms: 50,
backoff_multiplier: 1.5,
max_delay_ms: 500,
};
let misfire = MisfirePolicy {
run_immediately: true,
max_misfire_window_secs: 3600,
};
job.set_retry_policy(&retry);
job.set_misfire_policy(&misfire);
assert_eq!(job.retry_policy().max_attempts, 2);
assert!(job.misfire_policy().run_immediately);
assert_eq!(job.misfire_policy().max_misfire_window_secs, 3600);
}
}