use serde::{Deserialize, Serialize};
use crate::diagnostics::{categories, Diagnostic, ValidationReport};
use crate::model::{
AnalysisContext, ContractReference, DependencyGraph, ExecutionRequirements, FailureSemantics,
PipelineContract, PipelineGraph, PipelineLineage, PipelineStep, QualityGate, SchedulingIntent,
};
use crate::resolve::{
apply_nested_provenance, resolve_contract_references, stamp_nested_parents, NestedPipeline,
ResolveOptions,
};
use crate::validation;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "jsonschema", derive(schemars::JsonSchema))]
#[serde(rename_all = "camelCase")]
pub struct PlanDependencyEdge {
pub from: String,
pub to: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[cfg_attr(feature = "jsonschema", derive(schemars::JsonSchema))]
#[serde(rename_all = "camelCase")]
pub struct PipelinePlan {
pub contract_id: String,
pub contract_version: String,
pub dpcs_version: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub version: Option<String>,
pub steps: Vec<PipelineStep>,
pub graph: PipelineGraph,
pub contract_references: Vec<ContractReference>,
pub dependency_edges: Vec<PlanDependencyEdge>,
pub step_order: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub execution: Option<ExecutionRequirements>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub scheduling: Vec<SchedulingIntent>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub quality_gates: Vec<QualityGate>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub failure_semantics: Vec<FailureSemantics>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub lineage: Option<PipelineLineage>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub nested: Vec<NestedPlanPipeline>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[cfg_attr(feature = "jsonschema", derive(schemars::JsonSchema))]
#[serde(rename_all = "camelCase")]
pub struct NestedPlanPipeline {
pub parent_step_id: String,
pub contract_ref: String,
pub contract_id: String,
pub contract_version: String,
pub dpcs_version: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub input_ports: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub output_ports: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub step_order: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub children: Vec<NestedPlanPipeline>,
}
#[derive(Debug, Clone, PartialEq)]
pub enum PlanResult {
Ok(Box<PipelinePlan>),
Err(ValidationReport),
}
impl PlanResult {
pub fn plan(self) -> Option<PipelinePlan> {
match self {
Self::Ok(plan) => Some(*plan),
Self::Err(_) => None,
}
}
pub fn as_plan(&self) -> Option<&PipelinePlan> {
match self {
Self::Ok(plan) => Some(plan),
Self::Err(_) => None,
}
}
pub fn is_ok(&self) -> bool {
matches!(self, Self::Ok(_))
}
pub fn report(&self) -> Option<&ValidationReport> {
match self {
Self::Ok(_) => None,
Self::Err(report) => Some(report),
}
}
}
pub fn plan(contract: &PipelineContract) -> PlanResult {
let opts = ResolveOptions::default_for_planning();
plan_with_resolve(contract, Some(&opts))
}
pub fn plan_with_resolve(
contract: &PipelineContract,
resolve: Option<&ResolveOptions>,
) -> PlanResult {
let ctx = AnalysisContext::build(contract);
plan_with_context_and_resolve(&ctx, resolve)
}
pub fn plan_with_context(ctx: &AnalysisContext<'_>) -> PlanResult {
let opts = ResolveOptions::default_for_planning();
plan_with_context_and_resolve(ctx, Some(&opts))
}
pub fn plan_with_context_and_resolve(
ctx: &AnalysisContext<'_>,
resolve: Option<&ResolveOptions>,
) -> PlanResult {
let contract = ctx.contract;
let mut report = validation::validate_with_context(ctx);
let owned_default = ResolveOptions::default_for_planning();
let opts = resolve.unwrap_or(&owned_default);
let mut resolution = resolve_contract_references(contract, opts);
let mut nested_loaded: Vec<NestedPipeline> = Vec::new();
if !resolution.report.is_valid() {
report.extend(resolution.report);
} else {
stamp_nested_parents(&mut resolution.nested, &contract.id);
nested_loaded = resolution.nested;
report.extend(resolution.report);
}
if !report.is_valid() {
let mut planning_report = ValidationReport::new();
planning_report.push(
Diagnostic::planning_error(
"DPCS-PLN-001",
categories::PLANNING,
"pipeline plan requires a successfully validated contract",
)
.with_remediation("Resolve validation errors before planning")
.with_related(report.errors().map(|d| d.id.clone())),
);
planning_report.extend(report);
planning_report.sort_deterministic();
return PlanResult::Err(planning_report);
}
let dependency_graph = &ctx.graph;
let step_order = match dependency_graph.topological_order() {
Ok(order) => order,
Err(cycle) => {
let mut planning_report = ValidationReport::new();
planning_report.push(
Diagnostic::planning_error(
"DPCS-PLN-002",
categories::PLANNING,
format!(
"pipeline dependency graph contains a cycle: {}",
cycle.cycle.join(" -> ")
),
)
.with_remediation("Remove cyclic graph / control-flow / data-flow dependencies"),
);
planning_report.sort_deterministic();
return PlanResult::Err(planning_report);
}
};
let dependency_edges = dependency_graph
.edges()
.into_iter()
.map(|(from, to)| PlanDependencyEdge { from, to })
.collect();
let mut lineage = contract.lineage.clone();
apply_nested_provenance(&mut lineage, &contract.id, &nested_loaded);
let nested = nested_loaded
.into_iter()
.map(nested_pipeline_to_plan)
.collect();
PlanResult::Ok(Box::new(PipelinePlan {
contract_id: contract.id.clone(),
contract_version: contract.version.clone(),
dpcs_version: contract.dpcs_version.clone(),
version: None,
steps: contract.steps.clone(),
graph: contract.graph.clone(),
contract_references: contract.contract_references.clone(),
dependency_edges,
step_order,
execution: contract.execution.clone(),
scheduling: contract.scheduling.clone(),
quality_gates: contract.quality_gates.clone(),
failure_semantics: contract.failure_semantics.clone(),
lineage,
nested,
}))
}
pub fn try_plan(contract: &PipelineContract) -> Option<PipelinePlan> {
plan(contract).plan()
}
fn nested_pipeline_to_plan(n: NestedPipeline) -> NestedPlanPipeline {
let step_order = DependencyGraph::from_contract(&n.contract)
.topological_order()
.unwrap_or_else(|_| {
n.contract
.steps
.iter()
.map(|step| step.id.clone())
.collect()
});
NestedPlanPipeline {
parent_step_id: n.parent_step_id,
contract_ref: n.contract_ref,
contract_id: n.contract.id.clone(),
contract_version: n.contract.version.clone(),
dpcs_version: n.contract.dpcs_version.clone(),
input_ports: n
.contract
.interface
.inputs
.iter()
.map(|port| port.id.clone())
.collect(),
output_ports: n
.contract
.interface
.outputs
.iter()
.map(|port| port.id.clone())
.collect(),
step_order,
children: n
.children
.into_iter()
.map(nested_pipeline_to_plan)
.collect(),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::parser::parse_yaml;
#[test]
fn refuses_invalid_contract_with_pln_001() {
let contract = parse_yaml(
r#"
dpcsVersion: "1.0.0"
id: "test"
version: "0.1.0"
interface:
inputs: []
outputs: []
steps:
- id: ""
type: "extension:noop"
graph:
edges: []
"#,
)
.unwrap();
let PlanResult::Err(report) = plan(&contract) else {
panic!("expected planning refusal");
};
assert!(report.diagnostics.iter().any(|d| d.id == "DPCS-PLN-001"));
assert!(report.diagnostics.iter().any(|d| d.id == "DPCS-COM-004"));
}
#[test]
fn plans_valid_minimal_contract() {
let contract = parse_yaml(
r#"
dpcsVersion: "1.0.0"
id: "test"
version: "0.1.0"
interface:
inputs: []
outputs: []
steps:
- id: "a"
type: "extension:noop"
- id: "b"
type: "extension:noop"
graph:
edges:
- from: "a"
to: "b"
"#,
)
.unwrap();
let PlanResult::Ok(plan) = plan(&contract) else {
panic!("expected successful plan");
};
assert_eq!(plan.step_order, vec!["a", "b"]);
assert_eq!(plan.dependency_edges.len(), 1);
assert_eq!(plan.dependency_edges[0].from, "a");
assert_eq!(plan.dependency_edges[0].to, "b");
}
#[test]
fn independent_steps_use_sorted_id_tie_break() {
let contract = parse_yaml(
r#"
dpcsVersion: "1.0.0"
id: "test"
version: "0.1.0"
interface:
inputs: []
outputs: []
steps:
- id: "z"
type: "extension:noop"
- id: "a"
type: "extension:noop"
graph:
edges: []
"#,
)
.unwrap();
let PlanResult::Ok(plan) = plan(&contract) else {
panic!("expected successful plan");
};
assert_eq!(plan.step_order, vec!["a", "z"]);
}
}