use std::collections::HashMap;
use std::fmt;
use std::time::Duration;
use chrono::{DateTime, TimeDelta, Utc};
use rust_decimal::Decimal;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use uuid::Uuid;
use super::{FsmState, RunActor, RunStatus, TriggerKind};
#[derive(Debug, Clone, Serialize, Deserialize)]
#[non_exhaustive]
pub struct Run {
pub id: Uuid,
pub workflow_name: String,
pub status: FsmState<RunStatus>,
pub trigger: TriggerKind,
pub payload: Value,
pub error: Option<String>,
pub retry_count: u32,
pub max_retries: u32,
pub cost_usd: Decimal,
pub duration_ms: u64,
pub created_at: DateTime<Utc>,
pub updated_at: DateTime<Utc>,
pub started_at: Option<DateTime<Utc>>,
pub completed_at: Option<DateTime<Utc>>,
pub handler_version: Option<String>,
#[serde(default)]
pub labels: HashMap<String, String>,
#[serde(default)]
pub scheduled_at: Option<DateTime<Utc>>,
#[serde(default)]
pub created_by: Option<RunActor>,
#[serde(default)]
pub created_by_label: Option<String>,
#[serde(default)]
pub idempotency_key: Option<String>,
#[serde(default)]
pub max_cost_usd: Option<Decimal>,
#[serde(default)]
pub worker_id: Option<String>,
#[serde(default)]
pub lease_expires_at: Option<DateTime<Utc>>,
}
pub const IDEMPOTENCY_WINDOW: TimeDelta = TimeDelta::hours(24);
pub const MAX_IDEMPOTENCY_KEY_LEN: usize = 255;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum RunCreation {
Created(Run),
Existing(Run),
}
impl RunCreation {
pub fn into_run(self) -> Run {
match self {
RunCreation::Created(run) | RunCreation::Existing(run) => run,
}
}
pub fn run(&self) -> &Run {
match self {
RunCreation::Created(run) | RunCreation::Existing(run) => run,
}
}
pub fn is_created(&self) -> bool {
matches!(self, RunCreation::Created(_))
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LeaseRequest {
pub worker_id: String,
pub ttl: Duration,
}
impl LeaseRequest {
pub fn expires_at(&self, from: DateTime<Utc>) -> DateTime<Utc> {
TimeDelta::from_std(self.ttl)
.ok()
.and_then(|ttl| from.checked_add_signed(ttl))
.unwrap_or(DateTime::<Utc>::MAX_UTC)
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct ReapedRun {
pub run: Run,
pub from: RunStatus,
pub to: RunStatus,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct NewRun {
pub workflow_name: String,
pub trigger: TriggerKind,
pub payload: Value,
pub max_retries: u32,
pub handler_version: Option<String>,
#[serde(default)]
pub labels: HashMap<String, String>,
#[serde(default)]
pub scheduled_at: Option<DateTime<Utc>>,
#[serde(default)]
pub created_by: Option<RunActor>,
#[serde(default)]
pub idempotency_key: Option<String>,
#[serde(default)]
pub max_cost_usd: Option<Decimal>,
}
#[derive(Debug, Clone, Default)]
pub struct RunFilter {
pub workflow_name: Option<String>,
pub status: Option<RunStatus>,
pub created_after: Option<DateTime<Utc>>,
pub created_before: Option<DateTime<Utc>>,
pub has_steps: Option<bool>,
pub labels: Option<HashMap<String, String>>,
pub created_by_user_id: Option<Uuid>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct RunUpdate {
pub status: Option<RunStatus>,
pub error: Option<String>,
pub increment_retry: bool,
pub cost_usd: Option<Decimal>,
pub duration_ms: Option<u64>,
pub started_at: Option<DateTime<Utc>>,
pub completed_at: Option<DateTime<Utc>>,
#[serde(default)]
pub scheduled_at: Option<DateTime<Utc>>,
}
#[derive(Debug, Clone)]
pub struct PurgePolicy {
pub max_age_days: u32,
pub max_runs_per_workflow: u32,
pub dry_run: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PurgeReason {
TooOld,
ExceedsWorkflowLimit,
}
impl fmt::Display for PurgeReason {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
PurgeReason::TooOld => f.write_str("too_old"),
PurgeReason::ExceedsWorkflowLimit => f.write_str("exceeds_workflow_limit"),
}
}
}
#[derive(Debug, Clone)]
pub struct PurgeableRun {
pub run_id: Uuid,
pub workflow_name: String,
pub reason: PurgeReason,
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use super::*;
use serde_json::json;
#[test]
fn newrun_serde_roundtrip() {
let new_run = NewRun {
created_by: None,
workflow_name: "deploy".to_string(),
trigger: TriggerKind::Manual,
payload: json!({"key": "value"}),
max_retries: 3,
handler_version: Some("1.2.0".to_string()),
labels: HashMap::from([("env".to_string(), "prod".to_string())]),
scheduled_at: None,
idempotency_key: None,
max_cost_usd: Some(Decimal::new(250, 2)),
};
let json = serde_json::to_string(&new_run).expect("serialize");
let back: NewRun = serde_json::from_str(&json).expect("deserialize");
assert_eq!(back.max_cost_usd, new_run.max_cost_usd);
assert_eq!(back.workflow_name, new_run.workflow_name);
assert_eq!(back.trigger, new_run.trigger);
assert_eq!(back.payload, new_run.payload);
assert_eq!(back.max_retries, new_run.max_retries);
assert_eq!(back.handler_version, new_run.handler_version);
assert_eq!(back.labels, new_run.labels);
assert_eq!(back.scheduled_at, new_run.scheduled_at);
assert_eq!(back.created_by, new_run.created_by);
assert_eq!(back.idempotency_key, new_run.idempotency_key);
}
#[test]
fn newrun_serde_roundtrip_with_actor() {
let actor = RunActor::ApiKey {
api_key_id: Uuid::now_v7(),
user_id: Uuid::now_v7(),
};
let new_run = NewRun {
workflow_name: "deploy".to_string(),
trigger: TriggerKind::Api,
payload: json!({}),
max_retries: 0,
handler_version: None,
labels: HashMap::new(),
scheduled_at: None,
created_by: Some(actor.clone()),
idempotency_key: None,
max_cost_usd: None,
};
let json = serde_json::to_string(&new_run).expect("serialize");
let back: NewRun = serde_json::from_str(&json).expect("deserialize");
assert_eq!(back.created_by, Some(actor));
}
#[test]
fn newrun_deserializes_without_created_by() {
let raw = json!({
"workflow_name": "deploy",
"trigger": {"kind": "workflow"},
"payload": {},
"max_retries": 0,
"handler_version": null,
});
let new_run: NewRun = serde_json::from_value(raw).expect("deserialize");
assert!(new_run.created_by.is_none());
}
#[test]
fn run_serde_preserves_all_fields() {
use crate::entities::FsmState;
use chrono::Utc;
use uuid::Uuid;
let now = Utc::now();
let run = Run {
id: Uuid::now_v7(),
workflow_name: "test-wf".to_string(),
status: FsmState::new(RunStatus::Running, Uuid::now_v7()),
trigger: TriggerKind::Webhook {
path: "/hooks/test".to_string(),
},
payload: json!({"data": 123}),
error: Some("test error".to_string()),
retry_count: 2,
max_retries: 5,
cost_usd: Decimal::new(1234, 2),
duration_ms: 5000,
created_at: now,
updated_at: now,
started_at: Some(now),
completed_at: Some(now),
handler_version: Some("2.0.0".to_string()),
labels: HashMap::from([
("env".to_string(), "staging".to_string()),
("team".to_string(), "platform".to_string()),
]),
scheduled_at: Some(now),
created_by: Some(RunActor::User {
user_id: Uuid::now_v7(),
}),
created_by_label: Some("alice".to_string()),
idempotency_key: Some("gh:abc-123".to_string()),
max_cost_usd: Some(Decimal::new(500, 2)),
worker_id: Some("worker-1".to_string()),
lease_expires_at: Some(now),
};
let json = serde_json::to_string(&run).expect("serialize");
let back: Run = serde_json::from_str(&json).expect("deserialize");
assert_eq!(back.id, run.id);
assert_eq!(back.workflow_name, run.workflow_name);
assert_eq!(back.status.state, run.status.state);
assert_eq!(back.trigger, run.trigger);
assert_eq!(back.payload, run.payload);
assert_eq!(back.error, run.error);
assert_eq!(back.retry_count, run.retry_count);
assert_eq!(back.max_retries, run.max_retries);
assert_eq!(back.cost_usd, run.cost_usd);
assert_eq!(back.duration_ms, run.duration_ms);
assert_eq!(back.started_at, run.started_at);
assert_eq!(back.completed_at, run.completed_at);
assert_eq!(back.handler_version, run.handler_version);
assert_eq!(back.labels, run.labels);
assert_eq!(back.scheduled_at, run.scheduled_at);
assert_eq!(back.created_by, run.created_by);
assert_eq!(back.created_by_label, run.created_by_label);
assert_eq!(back.idempotency_key, run.idempotency_key);
assert_eq!(back.max_cost_usd, run.max_cost_usd);
assert_eq!(back.worker_id, run.worker_id);
assert_eq!(back.lease_expires_at, run.lease_expires_at);
}
#[test]
fn newrun_max_cost_usd_defaults_to_none_when_absent() {
let without_cap = NewRun {
workflow_name: "deploy".to_string(),
trigger: TriggerKind::Manual,
payload: json!({}),
max_retries: 0,
handler_version: None,
labels: HashMap::new(),
scheduled_at: None,
created_by: None,
idempotency_key: None,
max_cost_usd: None,
};
let mut value = serde_json::to_value(&without_cap).expect("serialize");
value
.as_object_mut()
.expect("object")
.remove("max_cost_usd");
let parsed: NewRun = serde_json::from_value(value).expect("deserialize");
assert!(parsed.max_cost_usd.is_none());
}
#[test]
fn runupdate_serde_roundtrip() {
let update = RunUpdate {
status: Some(RunStatus::Completed),
error: Some("test error".to_string()),
increment_retry: true,
cost_usd: Some(Decimal::new(5000, 2)),
duration_ms: Some(3000),
started_at: None,
completed_at: None,
scheduled_at: Some(Utc::now()),
};
let json = serde_json::to_string(&update).expect("serialize");
let back: RunUpdate = serde_json::from_str(&json).expect("deserialize");
assert_eq!(back.status, update.status);
assert_eq!(back.error, update.error);
assert_eq!(back.increment_retry, update.increment_retry);
assert_eq!(back.cost_usd, update.cost_usd);
assert_eq!(back.duration_ms, update.duration_ms);
assert_eq!(back.scheduled_at, update.scheduled_at);
}
#[test]
fn runfilter_default_is_no_filters() {
let filter = RunFilter::default();
assert!(filter.workflow_name.is_none());
assert!(filter.status.is_none());
assert!(filter.created_after.is_none());
assert!(filter.created_before.is_none());
assert!(filter.created_by_user_id.is_none());
}
#[test]
fn runfilter_with_multiple_criteria() {
let filter = RunFilter {
workflow_name: Some("deploy".to_string()),
status: Some(RunStatus::Running),
..RunFilter::default()
};
assert_eq!(filter.workflow_name, Some("deploy".to_string()));
assert_eq!(filter.status, Some(RunStatus::Running));
assert!(filter.created_after.is_none());
assert!(filter.created_before.is_none());
}
}