Skip to main content

dpcs/model/
lineage.rs

1//! Pipeline lineage model (SPEC Ch 14).
2
3use serde::{Deserialize, Serialize};
4
5use super::ExtensionMap;
6
7/// Declared lineage information for a Pipeline Contract.
8#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
9#[cfg_attr(feature = "jsonschema", derive(schemars::JsonSchema))]
10#[serde(rename_all = "camelCase")]
11pub struct PipelineLineage {
12    /// Dataset provenance relationships.
13    #[serde(default, skip_serializing_if = "Vec::is_empty")]
14    pub datasets: Vec<DatasetLineage>,
15    /// Step provenance relationships.
16    #[serde(default, skip_serializing_if = "Vec::is_empty")]
17    pub steps: Vec<StepLineage>,
18    /// Pipeline contract provenance.
19    #[serde(default, skip_serializing_if = "Option::is_none")]
20    pub provenance: Option<PipelineProvenance>,
21    /// Optional audit metadata (non-semantic).
22    #[serde(default, skip_serializing_if = "Option::is_none")]
23    pub audit: Option<LineageAudit>,
24    /// Extension fields.
25    #[serde(default, flatten)]
26    pub extensions: ExtensionMap,
27}
28
29/// Dataset provenance within a pipeline.
30#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
31#[cfg_attr(feature = "jsonschema", derive(schemars::JsonSchema))]
32#[serde(rename_all = "camelCase")]
33pub struct DatasetLineage {
34    /// Dataset identity.
35    pub dataset: String,
36    /// Producing pipeline step identifier.
37    #[serde(default, skip_serializing_if = "Option::is_none")]
38    pub produced_by: Option<String>,
39    /// Consuming pipeline step identifiers.
40    #[serde(default, skip_serializing_if = "Vec::is_empty")]
41    pub consumed_by: Vec<String>,
42    /// Associated data contract reference id.
43    #[serde(default, skip_serializing_if = "Option::is_none")]
44    pub contract_ref: Option<String>,
45    /// Associated transformation contract reference id.
46    #[serde(default, skip_serializing_if = "Option::is_none")]
47    pub transform_ref: Option<String>,
48    /// Extension fields.
49    #[serde(default, flatten)]
50    pub extensions: ExtensionMap,
51}
52
53/// Step provenance within a pipeline.
54#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
55#[cfg_attr(feature = "jsonschema", derive(schemars::JsonSchema))]
56#[serde(rename_all = "camelCase")]
57pub struct StepLineage {
58    /// Step identity.
59    #[serde(rename = "stepId")]
60    pub step_id: String,
61    /// Predecessor step identifiers.
62    #[serde(default, skip_serializing_if = "Vec::is_empty")]
63    pub predecessors: Vec<String>,
64    /// Successor step identifiers.
65    #[serde(default, skip_serializing_if = "Vec::is_empty")]
66    pub successors: Vec<String>,
67    /// Dependency semantics description.
68    #[serde(default, skip_serializing_if = "Option::is_none")]
69    pub dependency_kind: Option<String>,
70    /// Associated contract reference id.
71    #[serde(default, skip_serializing_if = "Option::is_none")]
72    pub contract_ref: Option<String>,
73    /// Extension fields.
74    #[serde(default, flatten)]
75    pub extensions: ExtensionMap,
76}
77
78/// Provenance describing contract origin and relationships.
79#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
80#[cfg_attr(feature = "jsonschema", derive(schemars::JsonSchema))]
81#[serde(rename_all = "camelCase")]
82pub struct PipelineProvenance {
83    /// Originating pipeline contract identifier.
84    #[serde(default, skip_serializing_if = "Option::is_none")]
85    pub originating: Option<String>,
86    /// Parent pipeline contract identifiers.
87    #[serde(default, skip_serializing_if = "Vec::is_empty")]
88    pub parents: Vec<String>,
89    /// Nested pipeline contract identifiers.
90    #[serde(default, skip_serializing_if = "Vec::is_empty")]
91    pub nested: Vec<String>,
92    /// Imported pipeline contract identifiers.
93    #[serde(default, skip_serializing_if = "Vec::is_empty")]
94    pub imported: Vec<String>,
95    /// Version history entries.
96    #[serde(default, skip_serializing_if = "Vec::is_empty")]
97    pub version_history: Vec<String>,
98    /// Extension fields.
99    #[serde(default, flatten)]
100    pub extensions: ExtensionMap,
101}
102
103/// Non-semantic audit metadata for lineage.
104#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
105#[cfg_attr(feature = "jsonschema", derive(schemars::JsonSchema))]
106#[serde(rename_all = "camelCase")]
107pub struct LineageAudit {
108    /// Contract identifiers for audit.
109    #[serde(default, skip_serializing_if = "Vec::is_empty")]
110    pub contract_ids: Vec<String>,
111    /// Version identifiers for audit.
112    #[serde(default, skip_serializing_if = "Vec::is_empty")]
113    pub version_ids: Vec<String>,
114    /// Timestamps for audit.
115    #[serde(default, skip_serializing_if = "Vec::is_empty")]
116    pub timestamps: Vec<String>,
117    /// Extension fields.
118    #[serde(default, flatten)]
119    pub extensions: ExtensionMap,
120}