Skip to main content

faucet_lineage/
event.rs

1//! OpenLineage `RunEvent` object model (serde-faithful subset, OL 2.0.2).
2
3use serde::Serialize;
4
5/// Pinned OpenLineage RunEvent schema URL (spec 2.0.2).
6pub const OL_SCHEMA_URL: &str =
7    "https://openlineage.io/spec/2-0-2/OpenLineage.json#/$defs/RunEvent";
8
9/// Producer identifier embedded in every event.
10pub const PRODUCER: &str = concat!(
11    "https://github.com/PawanSikawat/faucet-stream/tree/v",
12    env!("CARGO_PKG_VERSION")
13);
14
15#[derive(Debug, Clone, Copy, Serialize)]
16#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
17pub enum EventType {
18    Start,
19    Running,
20    Complete,
21    Abort,
22    Fail,
23}
24
25#[derive(Debug, Clone, Serialize)]
26#[serde(rename_all = "camelCase")]
27pub struct RunEvent {
28    pub event_type: EventType,
29    pub event_time: String,
30    pub run: Run,
31    pub job: Job,
32    pub inputs: Vec<Dataset>,
33    pub outputs: Vec<Dataset>,
34    pub producer: String,
35    #[serde(rename = "schemaURL")]
36    pub schema_url: String,
37}
38
39#[derive(Debug, Clone, Serialize)]
40#[serde(rename_all = "camelCase")]
41pub struct Run {
42    pub run_id: String,
43    #[serde(skip_serializing_if = "RunFacets::is_empty")]
44    pub facets: RunFacets,
45}
46
47#[derive(Debug, Clone, Default, Serialize)]
48#[serde(rename_all = "camelCase")]
49pub struct RunFacets {
50    #[serde(skip_serializing_if = "Option::is_none")]
51    pub parent: Option<ParentRunFacet>,
52    #[serde(skip_serializing_if = "Option::is_none")]
53    pub nominal_time: Option<NominalTimeRunFacet>,
54}
55
56impl RunFacets {
57    fn is_empty(&self) -> bool {
58        self.parent.is_none() && self.nominal_time.is_none()
59    }
60}
61
62#[derive(Debug, Clone, Serialize)]
63#[serde(rename_all = "camelCase")]
64pub struct ParentRunFacet {
65    #[serde(rename = "_producer")]
66    pub producer: String,
67    #[serde(rename = "_schemaURL")]
68    pub schema_url: String,
69    pub run: ParentRunRef,
70    pub job: ParentJobRef,
71}
72
73#[derive(Debug, Clone, Serialize)]
74#[serde(rename_all = "camelCase")]
75pub struct ParentRunRef {
76    pub run_id: String,
77}
78
79#[derive(Debug, Clone, Serialize)]
80pub struct ParentJobRef {
81    pub namespace: String,
82    pub name: String,
83}
84
85#[derive(Debug, Clone, Serialize)]
86#[serde(rename_all = "camelCase")]
87pub struct NominalTimeRunFacet {
88    #[serde(rename = "_producer")]
89    pub producer: String,
90    #[serde(rename = "_schemaURL")]
91    pub schema_url: String,
92    pub nominal_start_time: String,
93    #[serde(skip_serializing_if = "Option::is_none")]
94    pub nominal_end_time: Option<String>,
95}
96
97#[derive(Debug, Clone, Serialize)]
98pub struct Job {
99    pub namespace: String,
100    pub name: String,
101    #[serde(skip_serializing_if = "JobFacets::is_empty")]
102    pub facets: JobFacets,
103}
104
105#[derive(Debug, Clone, Default, Serialize)]
106#[serde(rename_all = "camelCase")]
107pub struct JobFacets {
108    #[serde(skip_serializing_if = "Option::is_none")]
109    pub source_code: Option<SourceCodeJobFacet>,
110}
111
112impl JobFacets {
113    fn is_empty(&self) -> bool {
114        self.source_code.is_none()
115    }
116}
117
118#[derive(Debug, Clone, Serialize)]
119#[serde(rename_all = "camelCase")]
120pub struct SourceCodeJobFacet {
121    #[serde(rename = "_producer")]
122    pub producer: String,
123    #[serde(rename = "_schemaURL")]
124    pub schema_url: String,
125    pub language: String,
126    pub source_code: String,
127}
128
129#[derive(Debug, Clone, Serialize)]
130pub struct Dataset {
131    pub namespace: String,
132    pub name: String,
133    #[serde(skip_serializing_if = "DatasetFacets::is_empty")]
134    pub facets: DatasetFacets,
135}
136
137impl Dataset {
138    pub fn new(namespace: impl Into<String>, name: impl Into<String>) -> Self {
139        Self {
140            namespace: namespace.into(),
141            name: name.into(),
142            facets: DatasetFacets::default(),
143        }
144    }
145}
146
147#[derive(Debug, Clone, Default, Serialize)]
148#[serde(rename_all = "camelCase")]
149pub struct DatasetFacets {
150    #[serde(skip_serializing_if = "Option::is_none")]
151    pub schema: Option<SchemaDatasetFacet>,
152    #[serde(skip_serializing_if = "Option::is_none")]
153    pub column_lineage: Option<ColumnLineageDatasetFacet>,
154}
155
156impl DatasetFacets {
157    fn is_empty(&self) -> bool {
158        self.schema.is_none() && self.column_lineage.is_none()
159    }
160}
161
162#[derive(Debug, Clone, Serialize)]
163#[serde(rename_all = "camelCase")]
164pub struct SchemaDatasetFacet {
165    #[serde(rename = "_producer")]
166    pub producer: String,
167    #[serde(rename = "_schemaURL")]
168    pub schema_url: String,
169    pub fields: Vec<SchemaField>,
170}
171
172impl SchemaDatasetFacet {
173    pub fn new(fields: Vec<SchemaField>) -> Self {
174        Self {
175            producer: PRODUCER.into(),
176            schema_url: OL_SCHEMA_URL.into(),
177            fields,
178        }
179    }
180}
181
182#[derive(Debug, Clone, Serialize)]
183pub struct SchemaField {
184    pub name: String,
185    #[serde(rename = "type")]
186    pub type_: String,
187}
188
189#[derive(Debug, Clone, Serialize)]
190#[serde(rename_all = "camelCase")]
191pub struct ColumnLineageDatasetFacet {
192    #[serde(rename = "_producer")]
193    pub producer: String,
194    #[serde(rename = "_schemaURL")]
195    pub schema_url: String,
196    pub fields: std::collections::BTreeMap<String, ColumnLineageFieldEntry>,
197}
198
199#[derive(Debug, Clone, Serialize)]
200pub struct ColumnLineageFieldEntry {
201    pub input_fields: Vec<ColumnLineageInputField>,
202}
203
204#[derive(Debug, Clone, Serialize)]
205pub struct ColumnLineageInputField {
206    pub namespace: String,
207    pub name: String,
208    pub field: String,
209}
210
211#[cfg(test)]
212mod tests {
213    use super::*;
214
215    #[test]
216    fn serializes_minimal_start_event() {
217        let ev = RunEvent {
218            event_type: EventType::Start,
219            event_time: "2026-06-07T00:00:00Z".into(),
220            run: Run {
221                run_id: "r1".into(),
222                facets: RunFacets::default(),
223            },
224            job: Job {
225                namespace: "ns".into(),
226                name: "job1".into(),
227                facets: JobFacets::default(),
228            },
229            inputs: vec![Dataset::new("ns", "postgres://h/db?table=t")],
230            outputs: vec![Dataset::new("ns", "bigquery://p.d.t")],
231            producer: PRODUCER.into(),
232            schema_url: OL_SCHEMA_URL.into(),
233        };
234        let v = serde_json::to_value(&ev).unwrap();
235        assert_eq!(v["eventType"], "START");
236        assert_eq!(v["run"]["runId"], "r1");
237        assert_eq!(v["job"]["name"], "job1");
238        assert_eq!(v["inputs"][0]["name"], "postgres://h/db?table=t");
239        assert_eq!(v["outputs"][0]["namespace"], "ns");
240        assert_eq!(v["schemaURL"], OL_SCHEMA_URL);
241        // empty facets must not serialize as noise
242        assert!(v["run"].get("facets").is_none() || v["run"]["facets"].is_object());
243    }
244
245    #[test]
246    fn schema_facet_round_trips() {
247        let mut ds = Dataset::new("ns", "file:///x");
248        ds.facets.schema = Some(SchemaDatasetFacet::new(vec![SchemaField {
249            name: "id".into(),
250            type_: "integer".into(),
251        }]));
252        let v = serde_json::to_value(&ds).unwrap();
253        assert_eq!(v["facets"]["schema"]["fields"][0]["name"], "id");
254        assert_eq!(v["facets"]["schema"]["fields"][0]["type"], "integer");
255    }
256}