1use serde::Serialize;
4
5pub const OL_SCHEMA_URL: &str =
7 "https://openlineage.io/spec/2-0-2/OpenLineage.json#/$defs/RunEvent";
8
9pub 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 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}