use std::collections::BTreeMap;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use smallvec::SmallVec;
use uuid::Uuid;
macro_rules! define_id {
($name:ident, $doc:literal) => {
#[doc = $doc]
#[derive(
Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize,
)]
pub struct $name(pub Uuid);
impl $name {
pub fn new() -> Self {
Self(Uuid::new_v4())
}
}
impl Default for $name {
fn default() -> Self {
Self::new()
}
}
impl std::fmt::Display for $name {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.0)
}
}
};
}
define_id!(JobId, "Stable id for a single Job within a Run.");
define_id!(BundleId, "Stable id for a Bundle within a Run.");
define_id!(PipelineId, "Stable id for a registered PipelineSpec.");
define_id!(RunId, "Stable id for one execution of a Pipeline.");
define_id!(EdgeId, "Stable id for a single Edge instance within a Run.");
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum JobStatus {
Pending,
Queued,
Running,
Completed,
Failed,
Timeout,
Cancelled,
Skipped,
}
impl JobStatus {
pub fn is_terminal(self) -> bool {
matches!(
self,
Self::Completed | Self::Failed | Self::Timeout | Self::Cancelled | Self::Skipped
)
}
pub fn is_success(self) -> bool {
matches!(self, Self::Completed)
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum BundleStatus {
Pending,
Queued,
Running,
Completed,
Blocked,
Cancelled,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RunStatus {
Pending,
Running,
Completed,
Failed,
Cancelled,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum OutputProjection {
Whole,
Field(String),
JsonPath(String),
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum MergePolicy {
Reject,
LastWriteWins,
AppendArray,
ObjectMerge,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum EdgeCondition {
OnSuccess,
OnFailure,
OnTerminals(SmallVec<[JobStatus; 4]>),
Always,
}
impl EdgeCondition {
pub fn matches(&self, source: JobStatus) -> bool {
debug_assert!(source.is_terminal());
match self {
Self::OnSuccess => matches!(source, JobStatus::Completed),
Self::OnFailure => matches!(
source,
JobStatus::Failed | JobStatus::Timeout | JobStatus::Cancelled | JobStatus::Skipped
),
Self::OnTerminals(set) => set.contains(&source),
Self::Always => true,
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum JoinPolicy {
#[default]
AllRequired,
AnyApplied,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RunCondition {
#[default]
AllSuccess,
AnySuccess,
AllComplete,
Always,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum EdgeResolution {
Pending,
Applied,
Unsatisfied,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum TerminalOutcome {
Success {
output: serde_json::Value,
findings_count: Option<usize>,
},
Failure {
error: String,
exit_code: Option<i32>,
},
Timeout,
Cancelled,
}
impl TerminalOutcome {
pub fn status(&self) -> JobStatus {
match self {
Self::Success { .. } => JobStatus::Completed,
Self::Failure { .. } => JobStatus::Failed,
Self::Timeout => JobStatus::Timeout,
Self::Cancelled => JobStatus::Cancelled,
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct DispatchTicket {
pub job_id: JobId,
pub run_id: RunId,
pub expected_completion_gen: u64,
pub expected_cancel_gen: u64,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Job {
pub id: JobId,
pub run_id: RunId,
pub pipeline_id: PipelineId,
pub bundle_id: Option<BundleId>,
pub kind: String,
pub inputs: serde_json::Value,
pub output: Option<serde_json::Value>,
pub status: JobStatus,
pub error: Option<String>,
pub pending_inputs: u32,
pub edge_resolutions: BTreeMap<EdgeId, EdgeResolution>,
pub completion_generation: u64,
pub cancel_generation: u64,
pub join_policy: JoinPolicy,
pub cancel_kill_pending: bool,
pub pid: Option<u32>,
pub exit_code: Option<i32>,
pub findings_count: Option<usize>,
pub created_at: DateTime<Utc>,
pub started_at: Option<DateTime<Utc>>,
pub completed_at: Option<DateTime<Utc>>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Bundle {
pub id: BundleId,
pub run_id: RunId,
pub pipeline_id: PipelineId,
pub parent: Option<BundleId>,
pub job_ids: Vec<JobId>,
pub successor_ids: Vec<BundleId>,
pub run_condition: RunCondition,
pub status: BundleStatus,
pub blocked_reason: Option<String>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Run {
pub id: RunId,
pub pipeline_id: PipelineId,
pub status: RunStatus,
pub inputs: serde_json::Value,
pub cancel_generation: u64,
pub created_at: DateTime<Utc>,
pub completed_at: Option<DateTime<Utc>>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct JobTemplate {
pub id: JobId,
pub kind: String,
#[serde(default)]
pub default_inputs: serde_json::Value,
pub bundle_id: Option<BundleId>,
#[serde(default)]
pub join_policy: JoinPolicy,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct EdgeTemplate {
pub id: EdgeId,
pub from: JobId,
pub to: JobId,
pub source: OutputProjection,
pub target: String,
pub condition: EdgeCondition,
pub merge: MergePolicy,
#[serde(default = "default_required")]
pub required: bool,
}
fn default_required() -> bool {
true
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct BundleTemplate {
pub id: BundleId,
pub parent: Option<BundleId>,
pub job_ids: Vec<JobId>,
pub successor_ids: Vec<BundleId>,
#[serde(default)]
pub run_condition: RunCondition,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct PipelineSpec {
pub id: PipelineId,
pub name: String,
pub jobs: Vec<JobTemplate>,
pub edges: Vec<EdgeTemplate>,
pub bundles: Vec<BundleTemplate>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Default)]
pub enum LogRange {
#[default]
All,
Range {
start: u64,
end: u64,
},
Tail(usize),
}
#[derive(Clone, Debug)]
pub struct RunOutcome {
pub status: RunStatus,
pub jobs: Vec<Job>,
}