use serde::{ser::SerializeMap, Deserialize, Serialize, Serializer};
use super::{
is_reserved_root_field, CompatibilityPolicy, ContractReference, ControlFlow, DataFlow,
ExecutionRequirements, ExtensionMap, ExtensionValue, FailureSemantics, GovernanceMetadata,
IdentityCatalog, Metadata, PipelineGraph, PipelineIdentity, PipelineInterface, PipelineLineage,
PipelineStep, QualityGate, SchedulingIntent, SecurityMetadata,
};
use crate::diagnostics::ValidationReport;
use crate::error::Result;
use crate::validation;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[cfg_attr(feature = "jsonschema", derive(schemars::JsonSchema))]
#[serde(rename_all = "camelCase")]
pub struct PipelineContract {
pub dpcs_version: String,
pub id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
pub version: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub metadata: Option<Metadata>,
pub interface: PipelineInterface,
pub graph: PipelineGraph,
#[serde(default)]
pub steps: Vec<PipelineStep>,
#[serde(default)]
pub contract_references: Vec<ContractReference>,
#[serde(default)]
pub data_flow: Vec<DataFlow>,
#[serde(default)]
pub control_flow: Vec<ControlFlow>,
#[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)]
pub quality_gates: Vec<QualityGate>,
#[serde(default)]
pub failure_semantics: Vec<FailureSemantics>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub lineage: Option<PipelineLineage>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub compatibility: Option<CompatibilityPolicy>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub security: Option<SecurityMetadata>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub governance: Option<GovernanceMetadata>,
#[serde(default, flatten, serialize_with = "serialize_root_extensions")]
pub extensions: ExtensionMap,
}
fn serialize_root_extensions<S>(
extensions: &ExtensionMap,
serializer: S,
) -> std::result::Result<S::Ok, S::Error>
where
S: Serializer,
{
let filtered: Vec<(&String, &ExtensionValue)> = extensions
.iter()
.filter(|(key, _)| !is_reserved_root_field(key))
.collect();
let mut map = serializer.serialize_map(Some(filtered.len()))?;
for (key, value) in filtered {
map.serialize_entry(key, value)?;
}
map.end()
}
impl PipelineContract {
pub fn from_yaml_file(path: impl AsRef<std::path::Path>) -> Result<Self> {
crate::parser::parse_yaml_file(path)
}
pub fn from_json_file(path: impl AsRef<std::path::Path>) -> Result<Self> {
crate::parser::parse_json_file(path)
}
pub fn from_yaml_str(input: &str) -> Result<Self> {
crate::parser::parse_yaml(input)
}
pub fn from_json_str(input: &str) -> Result<Self> {
crate::parser::parse_json(input)
}
pub fn to_yaml_str(&self) -> Result<String> {
crate::parser::to_yaml(self)
}
pub fn to_json_str(&self) -> Result<String> {
crate::parser::to_json(self)
}
pub fn to_yaml_file(&self, path: impl AsRef<std::path::Path>) -> Result<()> {
crate::parser::to_yaml_file(self, path)
}
pub fn to_json_file(&self, path: impl AsRef<std::path::Path>) -> Result<()> {
crate::parser::to_json_file(self, path)
}
pub fn validate(&self) -> ValidationReport {
validation::validate(self)
}
pub fn identity(&self) -> PipelineIdentity {
PipelineIdentity {
id: self.id.clone().into(),
version: self.version.clone().into(),
dpcs_version: self.dpcs_version.clone().into(),
name: self.name.clone(),
}
}
pub fn identity_catalog(&self) -> IdentityCatalog {
IdentityCatalog::from_contract(self)
}
pub fn step_ids(&self) -> std::collections::BTreeSet<&str> {
self.steps.iter().map(|step| step.id.as_str()).collect()
}
}